Core value types: opts structs, Job/JobState, Schedule, StreamEvent, Wake, StopToken, payload encode/decode (ADR-008 §1/§3, ADR-010 §3, ADR-015 §3, ADR-017 §3, ADR-019 §3/§4, ADR-020 §1–§4, ADR-021 §2)
This commit is contained in:
1 parent
74105a16ae
commit
92615f7b7d
12 files changed
+627
-8
No files matched your search
Generated
+27
@@ -6,6 +6,8 @@ version = 4
|
|||||||
name = "alkstore"
|
name = "alkstore"
|
||||||
version = "0.1.0"
|
version = "0.1.0"
|
||||||
dependencies = [
|
dependencies = [
|
||||||
|
"serde",
|
||||||
|
"serde_json",
|
||||||
"thiserror",
|
"thiserror",
|
||||||
]
|
]
|
||||||
|
|
||||||
@@ -318,6 +320,12 @@ dependencies = [
|
|||||||
"typenum",
|
"typenum",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "itoa"
|
||||||
|
version = "1.0.18"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "8f42a60cbdf9a97f5d2305f08a87dc4e09308d1276d28c869c684d7777685682"
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "js-sys"
|
name = "js-sys"
|
||||||
version = "0.3.106"
|
version = "0.3.106"
|
||||||
@@ -626,6 +634,19 @@ dependencies = [
|
|||||||
"syn 3.0.6",
|
"syn 3.0.6",
|
||||||
]
|
]
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "serde_json"
|
||||||
|
version = "1.0.151"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "c841b55ecdae098c80dcae9cf767f6f8a0c2cdb3416bbef72181df4d0fe73f14"
|
||||||
|
dependencies = [
|
||||||
|
"itoa",
|
||||||
|
"memchr",
|
||||||
|
"serde",
|
||||||
|
"serde_core",
|
||||||
|
"zmij",
|
||||||
|
]
|
||||||
|
|
||||||
[[package]]
|
[[package]]
|
||||||
name = "sha2"
|
name = "sha2"
|
||||||
version = "0.11.0"
|
version = "0.11.0"
|
||||||
@@ -999,3 +1020,9 @@ name = "wit-bindgen"
|
|||||||
version = "0.57.1"
|
version = "0.57.1"
|
||||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e"
|
checksum = "1ebf944e87a7c253233ad6766e082e3cd714b5d03812acc24c318f549614536e"
|
||||||
|
|
||||||
|
[[package]]
|
||||||
|
name = "zmij"
|
||||||
|
version = "1.0.23"
|
||||||
|
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||||
|
checksum = "29666d0abbfad1e3dc4dcf6144730dd3a3ab225bbbdac83319345b1b44ccfc1b"
|
||||||
@@ -6,4 +6,6 @@ license.workspace = true
|
|||||||
repository.workspace = true
|
repository.workspace = true
|
||||||
|
|
||||||
[dependencies]
|
[dependencies]
|
||||||
|
serde = "1"
|
||||||
|
serde_json = "1"
|
||||||
thiserror = "2"
|
thiserror = "2"
|
||||||
@@ -0,0 +1,91 @@
|
|||||||
|
//! Queue value types (ADR-019 §3, ADR-021 §2): `Job` and `JobState`.
|
||||||
|
//!
|
||||||
|
//! `#[non_exhaustive]` per ADR-017 §3 — consumers read but never
|
||||||
|
//! construct; field additions ride ADR-017 class 2 (semver-minor).
|
||||||
|
|
||||||
|
use serde::de::DeserializeOwned;
|
||||||
|
|
||||||
|
use crate::payload;
|
||||||
|
|
||||||
|
/// A job's lifecycle state. Dead-letter is a move to dead storage
|
||||||
|
/// (ADR-010 §4); `get_job` sees dead rows with `last_error`/`died_at`.
|
||||||
|
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
|
||||||
|
#[non_exhaustive]
|
||||||
|
pub enum JobState {
|
||||||
|
/// Enqueued, awaiting its ready time or a claim.
|
||||||
|
Pending,
|
||||||
|
/// Currently claimed and running (or within its claim deadline).
|
||||||
|
Processing,
|
||||||
|
/// Dead-lettered.
|
||||||
|
Dead,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// A job row, as returned by `get_job` and carried inside a claim's
|
||||||
|
/// boxed `JobHandle` (`job(&self) -> &Job`).
|
||||||
|
///
|
||||||
|
/// `#[non_exhaustive]` (ADR-017 §3). The resolved values — not the raw
|
||||||
|
/// enqueue fields — are what a `Job` shows: the ready time after
|
||||||
|
/// delay-over-`run_at` resolution, the absolute `expires_at` after
|
||||||
|
/// relative-`expires` resolution (ADR-020 §1–§2), and the `QueueOpts`
|
||||||
|
/// stamps attached at enqueue (ADR-010 §3a).
|
||||||
|
///
|
||||||
|
/// Dead rows carry their stamps too (`get_job` returns the stamps with
|
||||||
|
/// the row); `last_error`/`died_at` are dead-only (`None` on live
|
||||||
|
/// rows).
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
#[non_exhaustive]
|
||||||
|
pub struct Job {
|
||||||
|
/// The row id.
|
||||||
|
pub id: i64,
|
||||||
|
/// The queue name the job was enqueued into.
|
||||||
|
pub queue: String,
|
||||||
|
/// The lifecycle state.
|
||||||
|
pub state: JobState,
|
||||||
|
/// Payload bytes — the serde_json serialization of the trait
|
||||||
|
/// `Value` at enqueue time (ADR-020 §4); decode with
|
||||||
|
/// [`Job::payload_as`].
|
||||||
|
pub payload: Vec<u8>,
|
||||||
|
/// The claim-order key (priority DESC).
|
||||||
|
pub priority: i64,
|
||||||
|
/// The resolved ready time — the claim-order key within equal
|
||||||
|
/// priority (unix seconds).
|
||||||
|
pub run_at: i64,
|
||||||
|
/// Claim count; every claim counts, reclaims included.
|
||||||
|
pub attempts: i64,
|
||||||
|
/// The row's max-attempts stamp (resolved at enqueue).
|
||||||
|
pub max_attempts: i64,
|
||||||
|
/// The claimant (`worker_id` on claims); `None` pre-claim.
|
||||||
|
pub worker_id: Option<String>,
|
||||||
|
/// Unix seconds at the most recent claim (fresh or reclaim);
|
||||||
|
/// `None` before any claim (ADR-021 §2) — the dual-execution
|
||||||
|
/// window's diagnostic field, paired with `claim_expires_at`.
|
||||||
|
pub claimed_at: Option<i64>,
|
||||||
|
/// The claim deadline (unix seconds); `None` pre-claim.
|
||||||
|
pub claim_expires_at: Option<i64>,
|
||||||
|
/// Unix seconds at enqueue.
|
||||||
|
pub created_at: i64,
|
||||||
|
/// The resolved absolute job-level expiry; `None` = never expires.
|
||||||
|
pub expires_at: Option<i64>,
|
||||||
|
/// The visibility-timeout stamp resolved at enqueue.
|
||||||
|
pub visibility_timeout_s: i64,
|
||||||
|
/// The backoff-base stamp resolved at enqueue.
|
||||||
|
pub backoff_base_s: i64,
|
||||||
|
/// The dead-letter-retention stamp resolved at enqueue; `None` =
|
||||||
|
/// dead rows live forever.
|
||||||
|
pub dead_letter_retention_s: Option<i64>,
|
||||||
|
/// Dead-only: the recorded failure reason (`None` on live rows).
|
||||||
|
pub last_error: Option<String>,
|
||||||
|
/// Dead-only: unix seconds at death (`None` on live rows).
|
||||||
|
pub died_at: Option<i64>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Job {
|
||||||
|
/// Decode the stored payload bytes into a concrete target type.
|
||||||
|
///
|
||||||
|
/// The stored bytes are the serde_json serialization of the
|
||||||
|
/// enqueue-time `Value` (ADR-020 §4); any deserialization failure
|
||||||
|
/// is [`Error::Codec`](crate::Error::Codec).
|
||||||
|
pub fn payload_as<T: DeserializeOwned>(&self) -> crate::Result<T> {
|
||||||
|
payload::payload_as(&self.payload)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -6,13 +6,29 @@
|
|||||||
//! this crate and implement its traits.
|
//! this crate and implement its traits.
|
||||||
|
|
||||||
mod error;
|
mod error;
|
||||||
|
mod job;
|
||||||
|
mod opts;
|
||||||
|
mod payload;
|
||||||
|
mod schedule;
|
||||||
|
mod stop_token;
|
||||||
|
mod stream;
|
||||||
mod validation;
|
mod validation;
|
||||||
|
mod wake;
|
||||||
|
|
||||||
pub use error::{Error, Result};
|
pub use error::{Error, Result};
|
||||||
|
pub use job::{Job, JobState};
|
||||||
|
pub use opts::{EnqueueOpts, QueueOpts, ScheduleOpts};
|
||||||
|
pub use payload::{encode_payload, payload_as};
|
||||||
|
pub use schedule::Schedule;
|
||||||
|
pub use stop_token::StopToken;
|
||||||
|
pub use stream::StreamEvent;
|
||||||
pub use validation::{
|
pub use validation::{
|
||||||
NameKind, RESERVED_LISTENER_RECONNECTED, RESERVED_PREFIX, validate_local_name, validate_name,
|
NameKind, RESERVED_LISTENER_RECONNECTED, RESERVED_PREFIX, validate_local_name, validate_name,
|
||||||
validate_shared_name,
|
validate_shared_name,
|
||||||
};
|
};
|
||||||
|
pub use wake::Wake;
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests;
|
mod tests;
|
||||||
|
#[cfg(test)]
|
||||||
|
mod value_type_tests;
|
||||||
@@ -0,0 +1,95 @@
|
|||||||
|
//! Opts structs (ADR-008 §1, ADR-010 §3/§3a, ADR-020 §1–§3).
|
||||||
|
//!
|
||||||
|
//! Pinned v1 field sets, carried whole. Deliberately *not*
|
||||||
|
//! `#[non_exhaustive]` (ADR-017 §3): consumers construct these — an
|
||||||
|
//! opts *field* addition post-release is a class 3 breaking change,
|
||||||
|
//! priced that way deliberately.
|
||||||
|
//!
|
||||||
|
//! Field semantics (delay-over-`run_at` precedence, relative `expires`)
|
||||||
|
//! are engine-side resolution (ADR-020 §1–§2) — core carries only the
|
||||||
|
//! fields and their doc text.
|
||||||
|
|
||||||
|
/// Options for enqueueing a job (`Queue::enqueue`, `enqueue_tx`,
|
||||||
|
/// `outbox_enqueue_tx`).
|
||||||
|
///
|
||||||
|
/// Every field has a default-resolved meaning, so all are
|
||||||
|
/// constructor-settable with `..Default::default()` (`Default` gives
|
||||||
|
/// all-`None` and priority 0).
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Default)]
|
||||||
|
pub struct EnqueueOpts {
|
||||||
|
/// Relative delay in seconds from the enqueue instant.
|
||||||
|
///
|
||||||
|
/// When set, the row's ready time resolves to `now + delay` —
|
||||||
|
/// **wins over `run_at`** if both are supplied (ADR-020 §1;
|
||||||
|
/// precedence is defined, not rejected). `None` leaves the
|
||||||
|
/// resolution to `run_at`, or to "now" when both are unset.
|
||||||
|
pub delay: Option<i64>,
|
||||||
|
/// Absolute ready time (unix seconds).
|
||||||
|
///
|
||||||
|
/// The absolute-time counterpart of `delay`: used literally when
|
||||||
|
/// set and `delay` is not (ADR-020 §1). `get_job` shows the
|
||||||
|
/// resolved ready time, never the raw fields.
|
||||||
|
pub run_at: Option<i64>,
|
||||||
|
/// The claim-order key (priority DESC). Default 0.
|
||||||
|
pub priority: i64,
|
||||||
|
/// Per-job override of the queue-level max-attempts stamp.
|
||||||
|
///
|
||||||
|
/// `None` = the queue-level stamp decides (ADR-010 §3a).
|
||||||
|
pub max_attempts: Option<i64>,
|
||||||
|
/// Relative job-level expiry: seconds from the enqueue instant.
|
||||||
|
///
|
||||||
|
/// Resolved to an absolute row `expires_at` at enqueue (ADR-020
|
||||||
|
/// §2); `None` = never expires.
|
||||||
|
pub expires: Option<i64>,
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Queue-level option stamps (ADR-010 §3/§3a).
|
||||||
|
///
|
||||||
|
/// Queues are names, not registered objects: these stamps resolve and
|
||||||
|
/// attach to every job row at enqueue — claims, heartbeats, retries,
|
||||||
|
/// and sweeps read the job's own stamps, never a live registry.
|
||||||
|
/// Opts changes apply to *future enqueues only*.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
pub struct QueueOpts {
|
||||||
|
/// Claim visibility timeout in seconds — the default claim
|
||||||
|
/// deadline stamped onto rows. Default 300.
|
||||||
|
pub visibility_timeout_s: i64,
|
||||||
|
/// Default max attempts before dead-letter; overridden per-job by
|
||||||
|
/// `EnqueueOpts::max_attempts` when set. Default 3.
|
||||||
|
pub max_attempts: i64,
|
||||||
|
/// Base seconds for the retry backoff curve. Default 5.
|
||||||
|
pub backoff_base_s: i64,
|
||||||
|
/// Dead-letter retention in seconds; `None` = dead rows live
|
||||||
|
/// forever (the default). Default `None`.
|
||||||
|
pub dead_letter_retention_s: Option<i64>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for QueueOpts {
|
||||||
|
/// The pinned defaults (ADR-010 §3): 300 / 3 / 5 / none.
|
||||||
|
fn default() -> Self {
|
||||||
|
QueueOpts {
|
||||||
|
visibility_timeout_s: 300,
|
||||||
|
max_attempts: 3,
|
||||||
|
backoff_base_s: 5,
|
||||||
|
dead_letter_retention_s: None,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Options for `schedule()` — types identical to their `EnqueueOpts`
|
||||||
|
/// twins (no delay/run-at fields: boundary fires are ready at fire
|
||||||
|
/// time, ADR-009 §3).
|
||||||
|
///
|
||||||
|
/// Applied over the target queue's derived engine defaults at each
|
||||||
|
/// boundary fire (ADR-020 §3).
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq, Default)]
|
||||||
|
pub struct ScheduleOpts {
|
||||||
|
/// The claim-order key for fired jobs; default 0.
|
||||||
|
pub priority: i64,
|
||||||
|
/// Per-schedule override of the queue-level max-attempts stamp;
|
||||||
|
/// `None` = the default stamp decides.
|
||||||
|
pub max_attempts: Option<i64>,
|
||||||
|
/// Relative job-level expiry in seconds from the fire instant
|
||||||
|
/// (ADR-020 §2); `None` = never expires.
|
||||||
|
pub expires: Option<i64>,
|
||||||
|
}
|
||||||
@@ -0,0 +1,33 @@
|
|||||||
|
//! Payload encoding posture (ADR-020 §4): payloads cross the trait as
|
||||||
|
//! [`serde_json::Value`]; engines serialize with serde_json and store
|
||||||
|
//! exactly those bytes — nothing re-encodes on the way out.
|
||||||
|
//!
|
||||||
|
//! Core provides the encode helper (the `Value` → stored-bytes
|
||||||
|
//! serialization the engines call) and the [`payload_as`] decode
|
||||||
|
//! convenience used by [`Job`](crate::Job) and
|
||||||
|
//! [`StreamEvent`](crate::StreamEvent). `serde`/`serde_json` are core's
|
||||||
|
//! only dependencies beyond `thiserror`.
|
||||||
|
|
||||||
|
use serde::de::DeserializeOwned;
|
||||||
|
|
||||||
|
use crate::Error;
|
||||||
|
|
||||||
|
/// Serialize a trait-crossing payload `Value` to the stored byte form.
|
||||||
|
///
|
||||||
|
/// The engines call this at enqueue/publish time and store exactly the
|
||||||
|
/// returned bytes (ADR-020 §4): the stored bytes of a job row or
|
||||||
|
/// stream event row are, contract-pinned, the serde_json serialization
|
||||||
|
/// of the `Value` the call carried. Cross-engine byte equality falls
|
||||||
|
/// out — the same `Value` serializes identically on both engines.
|
||||||
|
pub fn encode_payload(value: &serde_json::Value) -> Vec<u8> {
|
||||||
|
serde_json::to_vec(value).unwrap_or_else(|_| Vec::new())
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Decode stored payload bytes into a concrete target type.
|
||||||
|
///
|
||||||
|
/// Used by [`Job::payload_as`](crate::Job::payload_as) and
|
||||||
|
/// [`StreamEvent::payload_as`](crate::StreamEvent::payload_as). The
|
||||||
|
/// error is [`Error::Codec`] on any deserialization failure (pinned).
|
||||||
|
pub fn payload_as<T: DeserializeOwned>(payload: &[u8]) -> crate::Result<T> {
|
||||||
|
serde_json::from_slice(payload).map_err(|e| Error::Codec(e.to_string()))
|
||||||
|
}
|
||||||
@@ -0,0 +1,22 @@
|
|||||||
|
//! Schedule read-back value (ADR-019 §3).
|
||||||
|
//!
|
||||||
|
//! `#[non_exhaustive]` per ADR-017 §3 — consumers read but never
|
||||||
|
//! construct.
|
||||||
|
|
||||||
|
use crate::opts::ScheduleOpts;
|
||||||
|
|
||||||
|
/// The `schedule()` read-back value (upsert by name).
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
#[non_exhaustive]
|
||||||
|
pub struct Schedule {
|
||||||
|
/// The schedule name (shared namespace — validated at the entry
|
||||||
|
/// point with the reserved prefix rejected).
|
||||||
|
pub name: String,
|
||||||
|
/// The v1 spec grammar: `@every <n><unit>` only (`s|m|h|d`).
|
||||||
|
pub spec: String,
|
||||||
|
/// The target queue name fired jobs enqueue into.
|
||||||
|
pub queue: String,
|
||||||
|
/// The schedule's opts, applied over the target queue's derived
|
||||||
|
/// engine defaults at each boundary fire (ADR-020 §3).
|
||||||
|
pub opts: ScheduleOpts,
|
||||||
|
}
|
||||||
@@ -0,0 +1,56 @@
|
|||||||
|
//! The stop token (ADR-019 §4): core-owned, contract-visible handle
|
||||||
|
//! for `run_schedules(stop)`.
|
||||||
|
//!
|
||||||
|
//! The token is pinned in core because the parameter type is contract
|
||||||
|
//! surface (ADR-017 — a pinned method's signature); the *await
|
||||||
|
//! mechanics* are engine-internal (tokio watch per engine — engines
|
||||||
|
//! may wrap their own watch channels around this token). Core
|
||||||
|
//! implements the flip-everywhere semantics with a shared
|
||||||
|
//! `Arc<AtomicBool>`.
|
||||||
|
|
||||||
|
use std::sync::Arc;
|
||||||
|
use std::sync::atomic::{AtomicBool, Ordering};
|
||||||
|
|
||||||
|
/// A cloneable cancel handle for `run_schedules(stop)`.
|
||||||
|
///
|
||||||
|
/// `cancel()` on any clone flips every clone; the engine's
|
||||||
|
/// `run_schedules` ends its sleep early and returns `Ok(())` — the
|
||||||
|
/// clean-stop return (leadership loss still returns
|
||||||
|
/// `Err(LeadershipLost)` regardless of the token, ADR-009 §1).
|
||||||
|
#[derive(Clone)]
|
||||||
|
pub struct StopToken {
|
||||||
|
cancelled: Arc<AtomicBool>,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl StopToken {
|
||||||
|
/// Create a fresh, uncancelled token.
|
||||||
|
pub fn new() -> Self {
|
||||||
|
StopToken {
|
||||||
|
cancelled: Arc::new(AtomicBool::new(false)),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Cancel the token — flips every clone (they share one flag).
|
||||||
|
pub fn cancel(&self) {
|
||||||
|
self.cancelled.store(true, Ordering::SeqCst);
|
||||||
|
}
|
||||||
|
|
||||||
|
/// Whether the token has been cancelled (by any clone).
|
||||||
|
pub fn is_cancelled(&self) -> bool {
|
||||||
|
self.cancelled.load(Ordering::SeqCst)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl Default for StopToken {
|
||||||
|
fn default() -> Self {
|
||||||
|
StopToken::new()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl std::fmt::Debug for StopToken {
|
||||||
|
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
||||||
|
f.debug_struct("StopToken")
|
||||||
|
.field("cancelled", &self.is_cancelled())
|
||||||
|
.finish()
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,43 @@
|
|||||||
|
//! Stream event value (ADR-015 §3).
|
||||||
|
//!
|
||||||
|
//! `#[non_exhaustive]` per ADR-017 §3 — consumers read but never
|
||||||
|
//! construct.
|
||||||
|
|
||||||
|
use serde::de::DeserializeOwned;
|
||||||
|
|
||||||
|
use crate::payload;
|
||||||
|
|
||||||
|
/// A stream event — honker's shape with one rename (`topic` →
|
||||||
|
/// `stream`, ADR-015 §3).
|
||||||
|
///
|
||||||
|
/// `#[non_exhaustive]` (ADR-017 §3). Offsets are positions in the
|
||||||
|
/// stream's log — immutable, never renumbered; order is global FIFO
|
||||||
|
/// under publish sequence. The key is carried metadata only: not a
|
||||||
|
/// partitioning instruction, not a dedup key, not an ordering key.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
#[non_exhaustive]
|
||||||
|
pub struct StreamEvent {
|
||||||
|
/// Position in the stream's log; immutable.
|
||||||
|
pub offset: i64,
|
||||||
|
/// The stream name.
|
||||||
|
pub stream: String,
|
||||||
|
/// The carried key; `None` unless published keyed.
|
||||||
|
pub key: Option<String>,
|
||||||
|
/// Payload bytes — the serde_json serialization of the trait
|
||||||
|
/// `Value` at publish time (ADR-020 §4); decode with
|
||||||
|
/// [`StreamEvent::payload_as`].
|
||||||
|
pub payload: Vec<u8>,
|
||||||
|
/// Unix seconds at publish (the engines' single clock source).
|
||||||
|
pub created_at: i64,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl StreamEvent {
|
||||||
|
/// Decode the stored payload bytes into a concrete target type.
|
||||||
|
///
|
||||||
|
/// The stored bytes are the serde_json serialization of the
|
||||||
|
/// publish-time `Value` (ADR-020 §4); any deserialization failure
|
||||||
|
/// is [`Error::Codec`](crate::Error::Codec).
|
||||||
|
pub fn payload_as<T: DeserializeOwned>(&self) -> crate::Result<T> {
|
||||||
|
payload::payload_as(&self.payload)
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -0,0 +1,163 @@
|
|||||||
|
use std::collections::HashMap;
|
||||||
|
|
||||||
|
use serde_json::{Value, json};
|
||||||
|
|
||||||
|
use crate::{EnqueueOpts, Job, JobState, QueueOpts, ScheduleOpts, StopToken, StreamEvent, Wake};
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn enqueue_opts_default_all_none_priority_zero() {
|
||||||
|
let d = EnqueueOpts::default();
|
||||||
|
assert_eq!(
|
||||||
|
d,
|
||||||
|
EnqueueOpts {
|
||||||
|
delay: None,
|
||||||
|
run_at: None,
|
||||||
|
priority: 0,
|
||||||
|
max_attempts: None,
|
||||||
|
expires: None,
|
||||||
|
}
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn queue_opts_default_pinned_300_3_5_none() {
|
||||||
|
let d = QueueOpts::default();
|
||||||
|
assert_eq!(
|
||||||
|
d,
|
||||||
|
QueueOpts {
|
||||||
|
visibility_timeout_s: 300,
|
||||||
|
max_attempts: 3,
|
||||||
|
backoff_base_s: 5,
|
||||||
|
dead_letter_retention_s: None,
|
||||||
|
}
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn schedule_opts_default_matches_enqueue_opts_twin_types() {
|
||||||
|
let d = ScheduleOpts::default();
|
||||||
|
assert_eq!(
|
||||||
|
d,
|
||||||
|
ScheduleOpts {
|
||||||
|
priority: 0,
|
||||||
|
max_attempts: None,
|
||||||
|
expires: None,
|
||||||
|
}
|
||||||
|
);
|
||||||
|
let e = EnqueueOpts::default();
|
||||||
|
assert_eq!(e.priority, d.priority);
|
||||||
|
assert_eq!(e.max_attempts, d.max_attempts);
|
||||||
|
assert_eq!(e.expires, d.expires);
|
||||||
|
}
|
||||||
|
|
||||||
|
fn sample_job(payload: Vec<u8>) -> Job {
|
||||||
|
Job {
|
||||||
|
id: 42,
|
||||||
|
queue: "orders".to_string(),
|
||||||
|
state: JobState::Dead,
|
||||||
|
payload,
|
||||||
|
priority: 7,
|
||||||
|
run_at: 1000,
|
||||||
|
attempts: 3,
|
||||||
|
max_attempts: 3,
|
||||||
|
worker_id: Some("w1".to_string()),
|
||||||
|
claimed_at: Some(999),
|
||||||
|
claim_expires_at: Some(1299),
|
||||||
|
created_at: 900,
|
||||||
|
expires_at: Some(5000),
|
||||||
|
visibility_timeout_s: 300,
|
||||||
|
backoff_base_s: 5,
|
||||||
|
dead_letter_retention_s: None,
|
||||||
|
last_error: Some("max attempts exceeded".to_string()),
|
||||||
|
died_at: Some(1100),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn sample_event(payload: Vec<u8>) -> StreamEvent {
|
||||||
|
StreamEvent {
|
||||||
|
offset: 17,
|
||||||
|
stream: "events".to_string(),
|
||||||
|
key: Some("k".to_string()),
|
||||||
|
payload,
|
||||||
|
created_at: 1234,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn payload_as_round_trips_json_value_through_job() {
|
||||||
|
let v = json!({"id": 9, "nested": {"ok": [true, null, 1.5]}, "s": "ünï"});
|
||||||
|
let bytes = crate::encode_payload(&v);
|
||||||
|
assert_eq!(bytes, serde_json::to_vec(&v).unwrap_or_default());
|
||||||
|
let job = sample_job(bytes);
|
||||||
|
let decoded: Value = job
|
||||||
|
.payload_as()
|
||||||
|
.unwrap_or_else(|e| panic!("decode failed: {e}"));
|
||||||
|
assert_eq!(decoded, v);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn event_payload_as_round_trips() {
|
||||||
|
let v = json!([1, 2, 3, {"x": "y"}]);
|
||||||
|
let event = sample_event(crate::encode_payload(&v));
|
||||||
|
let decoded: Value = event
|
||||||
|
.payload_as()
|
||||||
|
.unwrap_or_else(|e| panic!("decode failed: {e}"));
|
||||||
|
assert_eq!(decoded, v);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn payload_as_concrete_typed_targets() {
|
||||||
|
let v = json!({"alpha": 1, "beta": 2});
|
||||||
|
let job = sample_job(crate::encode_payload(&v));
|
||||||
|
let decoded: HashMap<String, i64> = job.payload_as().unwrap_or_else(|e| panic!("decode: {e}"));
|
||||||
|
assert_eq!(
|
||||||
|
decoded,
|
||||||
|
HashMap::from([("alpha".into(), 1), ("beta".into(), 2)])
|
||||||
|
);
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn payload_as_undecodable_target_yields_codec() {
|
||||||
|
let v = json!({"name": "x"});
|
||||||
|
let job = sample_job(crate::encode_payload(&v));
|
||||||
|
let r: crate::Result<String> = job.payload_as();
|
||||||
|
match r {
|
||||||
|
Err(crate::Error::Codec(msg)) => assert!(!msg.is_empty()),
|
||||||
|
other => panic!("expected Codec, got {other:?}"),
|
||||||
|
}
|
||||||
|
let r2: crate::Result<bool> =
|
||||||
|
serde_json::from_slice(b"not json").map_err(|e| crate::Error::Codec(e.to_string()));
|
||||||
|
assert!(matches!(r2, Err(crate::Error::Codec(_))));
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn stop_token_cancel_flips_all_clones() {
|
||||||
|
let token = StopToken::new();
|
||||||
|
let clone = token.clone();
|
||||||
|
let deep = clone.clone();
|
||||||
|
assert!(!token.is_cancelled());
|
||||||
|
assert!(!clone.is_cancelled());
|
||||||
|
assert!(!deep.is_cancelled());
|
||||||
|
clone.cancel();
|
||||||
|
assert!(token.is_cancelled());
|
||||||
|
assert!(clone.is_cancelled());
|
||||||
|
assert!(deep.is_cancelled());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn stop_token_independent_tokens_do_not_flip_each_other() {
|
||||||
|
let a = StopToken::default();
|
||||||
|
let b = StopToken::default();
|
||||||
|
a.cancel();
|
||||||
|
assert!(a.is_cancelled());
|
||||||
|
assert!(!b.is_cancelled());
|
||||||
|
}
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn wake_carries_only_the_channel() {
|
||||||
|
let w = Wake {
|
||||||
|
channel: "ch-1".to_string(),
|
||||||
|
};
|
||||||
|
assert_eq!(w.channel, "ch-1");
|
||||||
|
assert_eq!(w.clone(), w);
|
||||||
|
}
|
||||||
@@ -0,0 +1,22 @@
|
|||||||
|
//! Wake value (ADR-008 §3): the `listen(channel)` channel-carried
|
||||||
|
//! hint.
|
||||||
|
//!
|
||||||
|
//! `#[non_exhaustive]` per ADR-017 §3 — consumers read but never
|
||||||
|
//! construct.
|
||||||
|
|
||||||
|
/// A wake: the one piece of semantic content the wake contract allows
|
||||||
|
/// — the channel name. No payload is surfaced (honker's
|
||||||
|
/// `Notification { id, channel, payload }` is not); consumers needing
|
||||||
|
/// content use streams (the guarantee split, ADR-006). Wakes are
|
||||||
|
/// hints: their multiplicity is not part of the contract (coalescing
|
||||||
|
/// on SQLite, per-notify on Postgres).
|
||||||
|
///
|
||||||
|
/// The Postgres engine's synthetic reconnect-wake arrives as
|
||||||
|
/// `Wake { channel: RESERVED_LISTENER_RECONNECTED }` on every
|
||||||
|
/// subscriber's receiver after a watcher reconnect.
|
||||||
|
#[derive(Debug, Clone, PartialEq, Eq)]
|
||||||
|
#[non_exhaustive]
|
||||||
|
pub struct Wake {
|
||||||
|
/// The channel the wake was delivered on.
|
||||||
|
pub channel: String,
|
||||||
|
}
|
||||||
@@ -1,7 +1,7 @@
|
|||||||
---
|
---
|
||||||
id: core-value-types
|
id: core-value-types
|
||||||
name: Core value types (opts, Job, Schedule, StreamEvent, Wake, StopToken)
|
name: Core value types (opts, Job, Schedule, StreamEvent, Wake, StopToken)
|
||||||
status: pending
|
status: completed
|
||||||
depends_on: [core-errors-and-validation]
|
depends_on: [core-errors-and-validation]
|
||||||
scope: narrow
|
scope: narrow
|
||||||
risk: low
|
risk: low
|
||||||
@@ -68,14 +68,14 @@ part of the contract.
|
|||||||
|
|
||||||
## Acceptance Criteria
|
## Acceptance Criteria
|
||||||
|
|
||||||
- [ ] All types above with the exact pinned fields/types/derives;
|
- [x] All types above with the exact pinned fields/types/derives;
|
||||||
`#[non_exhaustive]` on read types, absent on opts structs
|
`#[non_exhaustive]` on read types, absent on opts structs
|
||||||
- [ ] `Default` impls: `EnqueueOpts` (all-None/0), `QueueOpts`
|
- [x] `Default` impls: `EnqueueOpts` (all-None/0), `QueueOpts`
|
||||||
(300/3/5/None), `ScheduleOpts`
|
(300/3/5/None), `ScheduleOpts`
|
||||||
- [ ] `payload_as<T>` round-trips a serde_json `Value` through encode →
|
- [x] `payload_as<T>` round-trips a serde_json `Value` through encode →
|
||||||
decode; non-serializable target type yields `Codec`
|
decode; non-serializable target type yields `Codec`
|
||||||
- [ ] `StopToken`: cancel on a clone flips all clones; unit-tested
|
- [x] `StopToken`: cancel on a clone flips all clones; unit-tested
|
||||||
- [ ] `cargo test -p alkstore`, clippy `-D warnings`, fmt clean
|
- [x] `cargo test -p alkstore`, clippy `-D warnings`, fmt clean
|
||||||
|
|
||||||
## References
|
## References
|
||||||
|
|
||||||
@@ -88,8 +88,57 @@ part of the contract.
|
|||||||
|
|
||||||
## Notes
|
## Notes
|
||||||
|
|
||||||
> To be filled by implementation agent
|
- `serde`/`serde_json` added to `alkstore` (ADR-020 §4 names them the
|
||||||
|
core crate's only deps beyond thiserror; used without the `derive`
|
||||||
|
feature — the trait crosses `serde_json::Value`, never generic
|
||||||
|
trait methods, so no derive machinery lives in core).
|
||||||
|
- Module layout: `src/opts.rs` (the three opts structs),
|
||||||
|
`src/job.rs` (`Job` + `JobState`), `src/schedule.rs`,
|
||||||
|
`src/stream.rs` (`StreamEvent`), `src/wake.rs`,
|
||||||
|
`src/stop_token.rs`, `src/payload.rs` (the encode/decode helpers),
|
||||||
|
re-exported from `src/lib.rs`.
|
||||||
|
- Encode helper named `encode_payload(&Value) -> Vec<u8>` — infallible
|
||||||
|
at the type level (serializing a `serde_json::Value` cannot fail;
|
||||||
|
the `unwrap_or_else` fallback is the unreachable empty vec, and
|
||||||
|
core holds no panic posture to break: it serializes, never
|
||||||
|
constructs, Values in engine paths). `payload_as<T:
|
||||||
|
DeserializeOwned>(bytes)` is the shared decode used by
|
||||||
|
`Job::payload_as` and `StreamEvent::payload_as`;
|
||||||
|
`Error::Codec(e.to_string())` on failure.
|
||||||
|
- "Non-serializable target yields Codec" is realized as a
|
||||||
|
runtime-decode failure test (target type statically decodable but
|
||||||
|
structurally mismatched with the stored bytes, plus non-JSON
|
||||||
|
bytes) — a type with no `Deserialize` impl cannot compile through
|
||||||
|
`payload_as`'s `DeserializeOwned` bound, so the Codec error is the
|
||||||
|
runtime arm by construction.
|
||||||
|
- ADR-019 §3's `Job` field list carried *as corrected by ADR-021 §2*
|
||||||
|
(`claimed_at` present, `None` pre-claim). `JobState` derives
|
||||||
|
`Copy` (small field-set enum, read-only value). Opts structs derive
|
||||||
|
`Debug/Clone/PartialEq/Eq` (+`Default` per the pinned defaults);
|
||||||
|
read types derive `Debug/Clone/PartialEq/Eq`.
|
||||||
|
- `StopToken` = `Arc<AtomicBool>` with `Ordering::SeqCst`; `new()`
|
||||||
|
and `Default` both construct uncancelled; doc text carries the
|
||||||
|
engine-internal watch-channel note (ADR-019 §4) and the
|
||||||
|
clean-stop/LeadershipLost split.
|
||||||
|
|
||||||
## Summary
|
## Summary
|
||||||
|
|
||||||
> To be filled on completion
|
Core value types implemented in `alkstore`: `EnqueueOpts`/
|
||||||
|
`QueueOpts`/`ScheduleOpts` (exact pinned fields, hand-written
|
||||||
|
`QueueOpts` `Default` = 300/3/5/None, derived defaults for the other
|
||||||
|
two — all-`None`/priority 0); `Job` (ADR-019 §3 field list + ADR-021
|
||||||
|
§2's `claimed_at`) and `JobState::{Pending, Processing, Dead}`;
|
||||||
|
`Schedule`; `StreamEvent`; `Wake`; all read types `#[non_exhaustive]`,
|
||||||
|
none of the opts structs; `Job::payload_as`/`StreamEvent::payload_as`
|
||||||
|
over the shared `payload::payload_as` + `encode_payload` helpers
|
||||||
|
(ADR-020 §4); `StopToken` (`Arc<AtomicBool>`, clone-flips-all).
|
||||||
|
Doc comments carry the pinned semantics text (delay-over-run_at,
|
||||||
|
relative expires, stamped QueueOpts, heartbeat is engine-side and
|
||||||
|
absent here, wake = channel-only hint, reconnect-wake constant
|
||||||
|
reference). Ten new unit tests (20 total in the crate) cover the
|
||||||
|
pinned defaults, a full JSON round-trip through
|
||||||
|
`encode_payload` → `Job`/`StreamEvent` → `payload_as` for `Value` and
|
||||||
|
concrete typed targets, Codec failures, clone-flip stop semantics,
|
||||||
|
and independent-token isolation. `cargo test -p alkstore` (20 pass),
|
||||||
|
`cargo clippy --all-targets -- -D warnings`, `cargo fmt --check`, and
|
||||||
|
workspace `cargo build` all clean.
|
||||||
Reference in new issue
Block a user