Skip to main content

agentplane/core/
id.rs

1//! Identifiers, digests, and the hash-chain primitive.
2//!
3//! ULIDs are used for run and case ids because they are lexicographically
4//! sortable: a journal scan for one run is a range scan, and cases sort by
5//! creation time for free.
6
7use std::fmt;
8
9use serde::{Deserialize, Deserializer, Serialize, Serializer, de::Error as _};
10use sha2::{Digest as _, Sha256};
11
12/// Wall-clock instant. Only ever obtained through a journaled effect
13/// (`StepCtx::now`), never by reading the ambient clock.
14///
15/// # Format it before it reaches a wire
16///
17/// This is a re-export of [`time::OffsetDateTime`], and that type's **derived**
18/// `Serialize` is not a date. Unless the `time` crate's `serde-human-readable`
19/// feature is on — it is not here, and enabling it would silently change every
20/// stored shape — a bare `Timestamp` field serialises to its **component
21/// array**:
22///
23/// ```text
24/// [2027, 15, 8, 0, 0, 0, 0, 0, 0]
25///  year  ordinal-day  h  m  s  ns  offset-h  offset-m  offset-s
26/// ```
27///
28/// It parses, it round-trips, and every consumer that expected a date gets nine
29/// numbers. A model asked to do arithmetic on one produces confident nonsense; a
30/// dashboard renders `2027` as the hour. The failure is quiet in exactly the way
31/// a format error is not.
32///
33/// So **every timestamp on a wire in this crate is RFC 3339**, and the two ways
34/// to spell that are:
35///
36/// ```
37/// # use agentplane::core::Timestamp;
38/// #[derive(serde::Serialize, serde::Deserialize)]
39/// struct Deadline {
40///     #[serde(with = "time::serde::rfc3339")]
41///     resolved_at: Timestamp,
42///     #[serde(default, with = "time::serde::rfc3339::option")]
43///     warn_at: Option<Timestamp>,
44/// }
45/// ```
46///
47/// and [`format_timestamp`] where the value is being placed into a
48/// `serde_json::json!` literal rather than a struct field — an effect
49/// descriptor, a `CloudEvents` envelope, a log line.
50///
51/// `tests/guards` walks the crate's serialized types and fails on a component
52/// array, so this is a rule the build enforces rather than one a reviewer has to
53/// remember. The hazard is worth stating here anyway, because `Timestamp` is
54/// public API: it lands in **your** tool payloads too, and the crate cannot
55/// check those.
56pub type Timestamp = time::OffsetDateTime;
57
58/// One instant, RFC 3339, for a place that is not a struct field.
59///
60/// `#[serde(with = "time::serde::rfc3339")]` is the answer wherever there is a
61/// field to attach it to. This is the answer where there is not — inside a
62/// `json!` literal, where a bare `Timestamp` would take the derived component
63/// array with nothing in the source to hint at it. See [`Timestamp`].
64///
65/// Falls back to the numeric Unix second if formatting fails, which needs an
66/// instant outside RFC 3339's representable range. That is not a case worth
67/// failing an effect key over, and a number is at least unambiguous.
68#[must_use]
69pub fn format_timestamp(at: Timestamp) -> String {
70    at.format(&time::format_description::well_known::Rfc3339)
71        .unwrap_or_else(|_| at.unix_timestamp().to_string())
72}
73
74/// Per-run monotonic journal position, starting at 1.
75pub type Seq = u64;
76
77/// Ownership fencing epoch.
78///
79/// Every journal append carries the writer's epoch and the store rejects stale
80/// epochs *in the same transaction that writes*, so a paused or partitioned
81/// instance that wakes up and keeps writing is fenced by the store rather than
82/// by timing. Split-brain cannot corrupt the chain by construction.
83pub type Epoch = u64;
84
85macro_rules! ulid_newtype {
86    ($(#[$m:meta])* $name:ident, $prefix:literal) => {
87        $(#[$m])*
88        ///
89        /// # One spelling, everywhere
90        ///
91        #[doc = concat!(" Prefixed (`", $prefix, "_01J8Z…`) in every direction:")]
92        /// [`Display`](std::fmt::Display), `Serialize`, the store key, the
93        /// column. An id is self-describing wherever it lands — a log line, a
94        /// URL, a JSON field, a SIEM — and grepping one against another finds
95        /// it whichever way round you do it.
96        ///
97        /// What this rules out is a store spelling one id two ways — prefixed in
98        /// its keys and bare in its payloads — which costs a consumer an
99        /// afternoon and looks like two ids until it does.
100        ///
101        /// [`Deserialize`] and [`parse`](Self::parse) accept a bare ULID too,
102        /// because one arrives from somewhere else often enough. The prefix is
103        /// checked rather than stripped blindly, so a `case_…` string is refused
104        /// where a [`RunId`] is wanted instead of silently becoming one.
105        #[derive(Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord)]
106        pub struct $name(pub ulid::Ulid);
107
108        impl $name {
109            /// Mint a fresh id from the ambient clock and a random tail.
110            ///
111            /// Reserved for the runtime's admission path, which journals the
112            /// result. Skills must never call this — see `clippy.toml`.
113            ///
114            /// **Sorted to the millisecond, not beyond it.** Ordering comes
115            /// from the timestamp prefix, so two ids minted in the same
116            /// millisecond order by their random tails — arbitrarily, and
117            /// differently on each instance. A range scan over one run's
118            /// records is exact because the run id is a *prefix*; "cases sort
119            /// by creation time" is true between milliseconds and a coin flip
120            /// inside one.
121            #[allow(clippy::disallowed_methods)]
122            #[must_use]
123            pub fn generate() -> Self {
124                Self(ulid::Ulid::generate())
125            }
126
127            /// Reconstruct from a stored string.
128            ///
129            /// Accepts both the prefixed form produced by [`std::fmt::Display`]
130            /// and a bare ULID, so ids round-trip through logs, URLs, and
131            /// database columns without the caller having to know which form it
132            /// is holding.
133            ///
134            /// # Errors
135            ///
136            /// If what remains after the prefix is not a ULID.
137            pub fn parse(s: &str) -> Result<Self, ulid::DecodeError> {
138                let bare = s.strip_prefix(concat!($prefix, "_")).unwrap_or(s);
139                ulid::Ulid::from_string(bare).map(Self)
140            }
141        }
142
143        /// The trait an id arrives through: an `axum` path segment, a `clap`
144        /// argument, a TOML key, a `serde` field with `#[serde(with = …)]`.
145        ///
146        /// `.parse()` is where every consumer reaches first, and its absence
147        /// made [`parse`](Self::parse) something each of them had to find.
148        impl std::str::FromStr for $name {
149            type Err = ulid::DecodeError;
150            fn from_str(s: &str) -> Result<Self, Self::Err> {
151                Self::parse(s)
152            }
153        }
154
155        /// The written form, so a record says which kind of id it holds.
156        impl Serialize for $name {
157            fn serialize<S: Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
158                s.collect_str(self)
159            }
160        }
161
162        /// Through [`parse`](Self::parse), so a bare ULID from somewhere else
163        /// still reads.
164        ///
165        /// A refusal names the type and the input: `ulid`'s own error is
166        /// `invalid length`, which says nothing about which field was wrong or
167        /// what it wanted.
168        impl<'de> Deserialize<'de> for $name {
169            fn deserialize<D: Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
170                let raw = <std::borrow::Cow<'de, str>>::deserialize(d)?;
171                Self::parse(&raw).map_err(|e| {
172                    D::Error::custom(format!(
173                        concat!("not a ", stringify!($name), " ('{}'): {}"),
174                        raw, e
175                    ))
176                })
177            }
178        }
179
180        impl fmt::Display for $name {
181            fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
182                write!(f, "{}_{}", $prefix, self.0)
183            }
184        }
185
186        impl fmt::Debug for $name {
187            fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
188                write!(f, "{self}")
189            }
190        }
191    };
192}
193
194ulid_newtype!(
195    /// A single execution: one goal, one plan, one lifetime.
196    ///
197    /// Runs stay short *by design* — longevity lives in the [`CaseId`], so a
198    /// six-week business process never pins a code version.
199    RunId, "run"
200);
201
202ulid_newtype!(
203    /// A long-lived, correlated business fact spanning many runs and weeks.
204    CaseId, "case"
205);
206
207ulid_newtype!(
208    /// One business act made of many independent ones.
209    ///
210    /// Owns N runs sharing one frozen plan — see [`crate::batch`]. Distinct from
211    /// a [`CaseId`] because the relationship is different in kind: a case is a
212    /// matter that several runs *touch* over weeks, while a batch is a single
213    /// act that several runs *constitute* in one pass.
214    BatchId, "batch"
215);
216
217/// Position of a step within a plan.
218#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
219#[serde(transparent)]
220pub struct StepId(pub u32);
221
222impl fmt::Display for StepId {
223    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
224        write!(f, "s{}", self.0)
225    }
226}
227
228/// SHA-256 over an artifact's canonical serialization.
229///
230/// Also the hash-chain link type. [`Digest::chain`] is the only way to extend a
231/// chain, and it hashes `prev ‖ bytes` — where `bytes` are the record's *wire*
232/// bytes, never a re-serialized or upcast form. Rehashing after an upcast would
233/// destroy tamper evidence for all history the moment a schema changed.
234#[derive(Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Default)]
235pub struct Digest([u8; 32]);
236
237impl Digest {
238    /// The chain's genesis link.
239    pub const ZERO: Self = Self([0u8; 32]);
240
241    /// Hash a byte string.
242    #[must_use]
243    pub fn of(bytes: &[u8]) -> Self {
244        let mut h = Sha256::new();
245        h.update(bytes);
246        Self(h.finalize().into())
247    }
248
249    /// Extend a hash chain: `H(prev ‖ bytes)`.
250    #[must_use]
251    pub fn chain(prev: Self, bytes: &[u8]) -> Self {
252        let mut h = Sha256::new();
253        h.update(prev.0);
254        h.update(bytes);
255        Self(h.finalize().into())
256    }
257
258    #[must_use]
259    pub const fn from_bytes(b: [u8; 32]) -> Self {
260        Self(b)
261    }
262
263    #[must_use]
264    pub const fn as_bytes(&self) -> &[u8; 32] {
265        &self.0
266    }
267
268    #[must_use]
269    pub fn to_hex(self) -> String {
270        hex::encode(self.0)
271    }
272
273    pub fn from_hex(s: &str) -> Result<Self, hex::FromHexError> {
274        let mut out = [0u8; 32];
275        hex::decode_to_slice(s, &mut out)?;
276        Ok(Self(out))
277    }
278
279    /// First 8 hex chars — enough to identify a record in a log line.
280    #[must_use]
281    pub fn short(self) -> String {
282        hex::encode(&self.0[..4])
283    }
284}
285
286impl fmt::Display for Digest {
287    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
288        write!(f, "{}", self.to_hex())
289    }
290}
291
292impl fmt::Debug for Digest {
293    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
294        write!(f, "{}…", self.short())
295    }
296}
297
298impl Serialize for Digest {
299    fn serialize<S: Serializer>(&self, s: S) -> Result<S::Ok, S::Error> {
300        s.serialize_str(&self.to_hex())
301    }
302}
303
304impl<'de> Deserialize<'de> for Digest {
305    fn deserialize<D: Deserializer<'de>>(d: D) -> Result<Self, D::Error> {
306        let s = String::deserialize(d)?;
307        Self::from_hex(&s).map_err(D::Error::custom)
308    }
309}
310
311/// Stable identity of one effect within one run.
312///
313/// `H(step ‖ ordinal ‖ kind ‖ canonical(args))`. Two properties follow:
314///
315/// * **Exactly-once** — the store's unique index on `(run_id, effect_key)` makes
316///   "an effect is started at most once per run" a database invariant.
317/// * **Divergence detection** — on replay the key is recomputed from the
318///   deterministic zone. A mismatch means the code took a different path than
319///   the recorded one, and the run is quarantined rather than allowed to
320///   silently diverge.
321///
322/// Skills never construct these: the runtime derives the key from the effect's
323/// descriptor plus its position, so a skill cannot forge or collide one.
324#[derive(Clone, Copy, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
325#[serde(transparent)]
326pub struct EffectKey(Digest);
327
328/// Which pass of a step an effect belongs to.
329///
330/// A step can run twice for entirely legitimate reasons: once going forward,
331/// and once again in reverse when a later step fails and the saga unwinds. The
332/// two are different work with different effects, and the journal has to be
333/// able to tell them apart.
334#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default, Serialize, Deserialize)]
335#[serde(rename_all = "snake_case")]
336pub enum Phase {
337    /// Doing the work.
338    #[default]
339    Forward,
340    /// Undoing it.
341    Compensating,
342}
343
344impl Phase {
345    #[must_use]
346    pub const fn is_forward(self) -> bool {
347        matches!(self, Self::Forward)
348    }
349
350    /// The spelling every store writes.
351    ///
352    /// One, because three stores keep a phase column and a phase that reads
353    /// back as the other half of the saga hands the unwind logic a compensating
354    /// record wearing the forward pass's name.
355    #[must_use]
356    pub const fn as_str(self) -> &'static str {
357        match self {
358            Self::Forward => "forward",
359            Self::Compensating => "compensating",
360        }
361    }
362
363    /// The inverse of [`as_str`](Self::as_str), written over
364    /// [`ALL`](Self::ALL) for the reason [`CaseStatus::parse`] is.
365    ///
366    /// [`CaseStatus::parse`]: crate::core::CaseStatus::parse
367    #[must_use]
368    pub fn parse(s: &str) -> Option<Self> {
369        Self::ALL.iter().copied().find(|c| c.as_str() == s)
370    }
371
372    /// Both passes.
373    pub const ALL: [Self; 2] = [Self::Forward, Self::Compensating];
374
375    /// By-reference form, for `skip_serializing_if`.
376    #[allow(clippy::trivially_copy_pass_by_ref)]
377    #[must_use]
378    pub const fn is_forward_ref(v: &Self) -> bool {
379        v.is_forward()
380    }
381}
382
383impl EffectKey {
384    /// `phase` separates the forward pass from the compensating one.
385    ///
386    /// Without it a step's compensation would restart its ordinal at zero and
387    /// collide with the step's own forward effects: replay would read the
388    /// forward result back as the compensation's, and the store's uniqueness
389    /// constraint would reject the second announcement.
390    ///
391    /// `attempt` is 1-based and part of the identity, so a retry is a *new*
392    /// effect in the journal rather than a second record under an existing key.
393    ///
394    /// Without it, attempt 2 would collide with attempt 1's recorded failure:
395    /// replay would read back the failure instead of the retry that followed,
396    /// and the store's uniqueness constraint on `EffectStarted` would reject
397    /// the second attempt outright.
398    pub(crate) fn derive(
399        step: StepId,
400        phase: Phase,
401        ordinal: u32,
402        attempt: u32,
403        kind: &str,
404        canonical_args: &[u8],
405    ) -> Self {
406        let mut h = Sha256::new();
407        h.update(step.0.to_be_bytes());
408        h.update([phase as u8]);
409        h.update(ordinal.to_be_bytes());
410        h.update(attempt.to_be_bytes());
411        h.update((kind.len() as u64).to_be_bytes());
412        h.update(kind.as_bytes());
413        h.update(canonical_args);
414        Self(Digest(h.finalize().into()))
415    }
416
417    #[must_use]
418    pub fn to_hex(self) -> String {
419        self.0.to_hex()
420    }
421
422    pub fn from_hex(s: &str) -> Result<Self, hex::FromHexError> {
423        Digest::from_hex(s).map(Self)
424    }
425}
426
427impl fmt::Display for EffectKey {
428    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
429        write!(f, "ek:{}", self.0.to_hex())
430    }
431}
432
433impl fmt::Debug for EffectKey {
434    fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
435        write!(f, "ek:{}…", self.0.short())
436    }
437}
438
439#[cfg(test)]
440mod tests {
441    use super::*;
442
443    #[test]
444    fn chain_is_order_sensitive() {
445        let a = Digest::chain(Digest::ZERO, b"a");
446        let b = Digest::chain(a, b"b");
447        let swapped = Digest::chain(Digest::chain(Digest::ZERO, b"b"), b"a");
448        assert_ne!(b, swapped, "chain must not be commutative");
449    }
450
451    #[test]
452    fn chain_detects_any_mutation() {
453        let genuine = Digest::chain(Digest::ZERO, b"record-1");
454        let tampered = Digest::chain(Digest::ZERO, b"record-2");
455        assert_ne!(genuine, tampered);
456    }
457
458    #[test]
459    fn digest_hex_roundtrips() {
460        let d = Digest::of(b"hello");
461        assert_eq!(Digest::from_hex(&d.to_hex()).unwrap(), d);
462    }
463
464    #[test]
465    fn effect_key_separates_step_phase_ordinal_and_attempt() {
466        let fwd = Phase::Forward;
467        let base = EffectKey::derive(StepId(0), fwd, 0, 1, "tool", b"{}");
468        let other_ordinal = EffectKey::derive(StepId(0), fwd, 1, 1, "tool", b"{}");
469        let other_step = EffectKey::derive(StepId(1), fwd, 0, 1, "tool", b"{}");
470        let other_attempt = EffectKey::derive(StepId(0), fwd, 0, 2, "tool", b"{}");
471        let compensating = EffectKey::derive(StepId(0), Phase::Compensating, 0, 1, "tool", b"{}");
472        assert_ne!(base, other_ordinal, "ordinal must be part of the key");
473        assert_ne!(base, other_step, "step must be part of the key");
474        assert_ne!(
475            base, other_attempt,
476            "attempt must be part of the key, or a retry collides with the \
477             failure it is retrying"
478        );
479        assert_ne!(
480            base, compensating,
481            "phase must be part of the key, or a step's compensation collides \
482             with its own forward pass"
483        );
484    }
485
486    /// Length-prefixing `kind` prevents a boundary-shifting collision: without
487    /// it, `("ab", "c")` and `("a", "bc")` would hash identically.
488    #[test]
489    fn effect_key_kind_is_length_prefixed() {
490        let a = EffectKey::derive(StepId(0), Phase::Forward, 0, 1, "ab", b"c");
491        let b = EffectKey::derive(StepId(0), Phase::Forward, 0, 1, "a", b"bc");
492        assert_ne!(a, b);
493    }
494
495    #[test]
496    fn ids_display_with_prefix() {
497        let r = RunId::generate();
498        assert!(r.to_string().starts_with("run_"));
499    }
500
501    /// The property that broke first time round: `Display` and `parse` must be
502    /// inverses, or every id that goes through a database column comes back
503    /// unreadable.
504    #[test]
505    fn ids_round_trip_through_their_displayed_form() {
506        let r = RunId::generate();
507        assert_eq!(RunId::parse(&r.to_string()).unwrap(), r);
508        let c = CaseId::generate();
509        assert_eq!(CaseId::parse(&c.to_string()).unwrap(), c);
510    }
511
512    /// A bare ULID is still accepted, so ids written by other tools parse.
513    #[test]
514    fn bare_ulids_still_parse() {
515        let r = RunId::generate();
516        assert_eq!(RunId::parse(&r.0.to_string()).unwrap(), r);
517    }
518
519    /// Prefixes are not interchangeable in *meaning*, but parsing is lenient by
520    /// design: a `CaseId` column holds case ids, and the prefix is a display
521    /// affordance rather than a type check.
522    #[test]
523    fn parsing_rejects_garbage() {
524        assert!(RunId::parse("run_not-a-ulid").is_err());
525    }
526
527    /// **One spelling, in every direction.**
528    ///
529    /// `Display`, `Serialize`, the store key and the column all write the
530    /// prefixed form, so the same id greps against itself wherever it is read.
531    /// They did not always: keys carried the prefix and record payloads did
532    /// not, so a consumer who stored the string an operator sees in a log could
533    /// not deserialize it — `invalid length`, from a check inside `ulid`, naming
534    /// neither the field nor the fact that there were two forms.
535    #[test]
536    fn one_spelling_is_written_and_both_are_read() {
537        let r = RunId::generate();
538
539        assert_eq!(
540            serde_json::to_string(&r).expect("serialises"),
541            format!("\"{r}\""),
542            "the written form is the one Display shows"
543        );
544
545        for text in [r.to_string(), r.0.to_string()] {
546            let json = serde_json::to_string(&text).expect("a string");
547            assert_eq!(
548                serde_json::from_str::<RunId>(&json).expect("both forms read"),
549                r,
550                "the {text} spelling did not deserialize"
551            );
552        }
553    }
554
555    /// Leniency stops at the prefix: a case id is not a run id.
556    ///
557    /// The half that makes accepting two spellings safe. `parse` strips only
558    /// *this* type's prefix, so `case_…` fails the ULID decode rather than
559    /// quietly becoming a `RunId` — which a blind `strip_prefix` up to the
560    /// underscore would have allowed.
561    #[test]
562    fn another_types_prefix_is_refused_rather_than_stripped() {
563        let case = CaseId::generate();
564        assert!(
565            serde_json::from_str::<RunId>(&format!("\"{case}\"")).is_err(),
566            "a case id deserialized as a run id"
567        );
568        assert!(RunId::parse(&case.to_string()).is_err());
569        // And the positive half, so the negative one is not vacuous.
570        assert!(CaseId::parse(&case.to_string()).is_ok());
571    }
572
573    /// `.parse()` is where an id arrives — a path segment, a CLI argument.
574    #[test]
575    fn ids_arrive_through_from_str() {
576        let r = RunId::generate();
577        assert_eq!(r.to_string().parse::<RunId>().unwrap(), r);
578        assert_eq!(r.0.to_string().parse::<RunId>().unwrap(), r);
579        assert!("nope".parse::<RunId>().is_err());
580    }
581
582    /// **The digest is SHA-256, pinned to values computed outside this crate.**
583    ///
584    /// Every effect key, chain link, Merkle leaf and blob address in the system
585    /// comes out of these two functions, so what they emit is a durable format
586    /// even though nothing declares it one: change a byte and every historical
587    /// run becomes unverifiable, silently, because the whole suite would agree
588    /// with itself under the new value.
589    ///
590    /// That is what makes a *self-consistent* test worthless here and why these
591    /// literals were produced by `shasum -a 256` and Python's `hashlib` rather
592    /// than by running this code and writing down the answer. They survived a
593    /// `sha2` 0.10 → 0.11 upgrade, which is the class of change they exist for:
594    /// the hasher was swapped underneath and the bytes had to be shown, not
595    /// assumed, to be the same ones.
596    ///
597    /// What this does not pin: the canonicalization that decides *which* bytes
598    /// reach the hasher. That is `canon`'s own contract.
599    #[test]
600    fn the_digest_matches_sha256_computed_elsewhere() {
601        assert_eq!(
602            Digest::of(b"").to_hex(),
603            "e3b0c44298fc1c149afbf4c8996fb92427ae41e4649b934ca495991b7852b855",
604            "the empty digest moved, so every digest in every journal moved with it"
605        );
606        assert_eq!(
607            Digest::of(b"agentplane").to_hex(),
608            "c0f8f77669f4860960387db0dc9984894587bcfcc75d1846ffd3574563833443"
609        );
610        // `chain` is `H(prev ‖ bytes)`, and the concatenation order is the half
611        // a reimplementation gets wrong — reversed, it still hashes, still
612        // verifies against itself, and agrees with no other reader.
613        assert_eq!(
614            Digest::chain(Digest::of(b""), b"agentplane").to_hex(),
615            "171c5eddf30189b0efde9c71eb14f77685fab77d7750152d032f9e7cf4181159"
616        );
617    }
618}