agentplane 0.47.0

Durable, replayable agent runtime — the journal is the plan of record
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
//! Inbound events and durable waits.
//!
//! # The race that makes this hard
//!
//! A run sends a request and then waits for an acknowledgement. The obvious
//! implementation registers a subscription, then suspends. It has a hole: the
//! acknowledgement can arrive *before the run reaches the wait at all* — a fast
//! counterparty, a slow first step, a retry that overtakes. With no subscription
//! yet, the event has nowhere to go, and the run then waits forever for
//! something that already happened.
//!
//! Every durable-execution system solves this the same way and it is worth
//! stating plainly: **inbound events are buffered durably on arrival, whether or
//! not anyone is waiting.** A wait first looks in the buffer, and only suspends
//! if nothing is there. Delivery and waiting meet in the store rather than in
//! time.
//!
//! Dead-lettering therefore happens on a *sweep* of the buffer, never on
//! arrival: "nobody is waiting for this yet" and "nobody will ever want this"
//! are different claims, and only the second is safe to act on.

use serde::{Deserialize, Serialize};
use serde_json::Value;

use crate::core::{CorrelationKey, EffectKey, RunId, Timestamp};

/// A message from outside, correlated by business key.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct InboundEvent {
    /// Who sent it, as a URI — [CloudEvents]' `source`.
    ///
    /// Load-bearing twice over, and neither use is decoration.
    ///
    /// **It completes the dedup identity.** `CloudEvents` defines uniqueness as
    /// `(source, id)`, and for good reason: `id` alone is unique only within one
    /// producer. Two counterparties numbering their messages from one — or
    /// minting ids from different UUID versions — collide, and the collision is
    /// silent, because the second message looks exactly like a retry of the
    /// first and is dropped as one.
    ///
    /// **It is the payload's provenance.** An inbound message is untrusted
    /// whoever sent it, but *which* untrusted party matters: a sink may name the
    /// sources an authority-bearing field may derive from, and `event:kind` says
    /// what arrived while saying nothing about who sent it.
    ///
    /// It is not authority. A producer writes this string; anyone may claim to
    /// be anyone. What it buys is a label a policy can reason about once the
    /// transport has authenticated the sender by other means — which is why an
    /// event arriving as a [`CloudEvent`](crate::core::CloudEvent) takes this
    /// from the authenticated transport and keeps the producer's own claim
    /// inside [`id`](Self::id). See
    /// [`CloudEvent::into_inbound`](crate::core::CloudEvent::into_inbound).
    ///
    /// [CloudEvents]: https://github.com/cloudevents/spec/blob/main/cloudevents/spec.md
    pub source: String,
    /// Stable identity for deduplication, unique **within a source**. A
    /// counterparty that retries must reuse this, or the same message is
    /// delivered twice.
    pub id: String,
    /// What kind of message, e.g. `"acknowledgement.received"`.
    pub kind: String,
    /// The business keys this message carries. Real messages do not know run
    /// ids; they carry document numbers and meter ids.
    pub correlation: Vec<CorrelationKey>,
    pub payload: Value,
    /// The operator who minted this message, when this plane minted it.
    ///
    /// Distinct from [`source`](Self::source), which is a producer's URI and
    /// which anyone may claim. This is set only where the runtime itself
    /// builds an event out of an operator act it already authenticated or
    /// witnessed — today that is a worklist decision, whose deciding actor
    /// would otherwise exist only inside [`payload`](Self::payload).
    ///
    /// It matters because the payload is a **sealed** journal field and this
    /// is not: an erasure takes the decision's reason with it, as it must, and
    /// must not take *who approved* with it.
    ///
    /// **Never read off a wire.** A counterparty's message leaves this `None`
    /// — a sender naming its own operator is an assertion the transport did
    /// not check, and the one thing this field must not become is a place to
    /// claim an identity.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub by: Option<crate::core::Operator>,
}

impl InboundEvent {
    pub fn new(
        source: impl Into<String>,
        id: impl Into<String>,
        kind: impl Into<String>,
        payload: Value,
    ) -> Self {
        Self {
            source: source.into(),
            id: id.into(),
            kind: kind.into(),
            correlation: Vec::new(),
            payload,
            by: None,
        }
    }

    #[must_use]
    pub fn correlate(mut self, key: CorrelationKey) -> Self {
        self.correlation.push(key);
        self
    }

    /// Name the operator this plane minted the message for.
    ///
    /// Only for events the runtime builds out of an act it authenticated or
    /// witnessed — see [`by`](Self::by). A message that arrived over a wire
    /// has no business calling this.
    #[must_use]
    pub fn minted_by(mut self, by: crate::core::Operator) -> Self {
        self.by = Some(by);
        self
    }

    /// The identity a store deduplicates on: `(source, id)`.
    ///
    /// [`origin_key`] is the construction, and this is one of its three
    /// callers — every store has to agree byte for byte about which two
    /// messages are the same message.
    #[must_use]
    pub fn dedup_key(&self) -> String {
        origin_key(&self.source, &self.id)
    }
}

/// Whether `kind` is one this plane mints for itself.
///
/// A worklist decision reaches the run waiting on it as an event of a kind in
/// the `agentplane.` namespace. Accepted from outside, a message of that kind
/// would decide a human task for whoever may post an event — no claim, no
/// eligibility, no four-eyes. So every external intake refuses the namespace,
/// and the worklist is the one door into it.
pub(crate) fn is_reserved_kind(kind: &str) -> bool {
    kind.starts_with("agentplane.")
}

/// A producer-scoped identity — `(source, id)` — as one key.
///
/// # Why a separator is not enough
///
/// `id` is unique only within one producer: two counterparties numbering their
/// messages from one collide, and the collision is silent because the second
/// message looks exactly like a retry of the first. So the key is the pair.
///
/// Joining the pair with a separator makes it *look* unforgeable and does not
/// make it so. `("a\u{1f}b", "c")` and `("a", "b\u{1f}c")` spell the same
/// bytes, so a producer that can choose either half of its own pair can spell
/// another producer's — and pre-empt that producer's next message as an
/// apparent duplicate, which is a silent suppression rather than a rejected
/// one. Refusing the separator in both halves closes it only where something
/// refuses; a `source` derived from a deployment's own
/// [`Authenticator`](crate::api::Authenticator) is not this crate's to validate,
/// and a guarantee resting on every producer of every half being well behaved
/// is a guarantee nobody can check.
///
/// So the length of `source` goes in front. A reader takes the digits, takes
/// that many bytes, and the rest is `id` — the split is *read* rather than
/// searched for, no pair can spell another whatever either half contains, and
/// the property holds without anybody validating anything.
///
/// The separator stays for legibility: a key is a store row an operator reads.
#[must_use]
pub fn origin_key(source: &str, id: &str) -> String {
    format!("{}\u{1f}{source}\u{1f}{id}", source.len())
}

/// The `source` half of an [`origin_key`], read back.
///
/// Read the way the construction says: the digits, then exactly that many
/// bytes. `None` for a key [`origin_key`] did not produce — so a surface asking
/// "which producer admitted this run" gets no answer from a key somebody
/// spelled by hand, rather than a wrong one.
#[must_use]
pub fn origin_source(key: &str) -> Option<&str> {
    let (len, rest) = key.split_once('\u{1f}')?;
    if len.is_empty() || !len.bytes().all(|b| b.is_ascii_digit()) {
        return None;
    }
    let len: usize = len.parse().ok()?;
    let source = rest.get(..len)?;
    rest[len..].starts_with('\u{1f}').then_some(source)
}

/// What a run is waiting for.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct AwaitSpec {
    pub kind: String,
    pub correlation: Vec<CorrelationKey>,
    /// The obligation that bounds this wait.
    ///
    /// Mandatory by design. An unbounded wait is a run that can hang forever
    /// with nothing to notice it — the failure mode that presents as "the
    /// process just stalled" and is invisible until someone asks.
    pub deadline: String,
    /// The one producer whose event satisfies this wait, matched against
    /// [`InboundEvent::source`]. Absent, any producer may.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub from: Option<String>,
}

impl AwaitSpec {
    pub fn new(kind: impl Into<String>, deadline: impl Into<String>) -> Self {
        Self {
            kind: kind.into(),
            correlation: Vec::new(),
            deadline: deadline.into(),
            from: None,
        }
    }

    #[must_use]
    pub fn correlate(mut self, key: CorrelationKey) -> Self {
        self.correlation.push(key);
        self
    }

    /// Accept the event only from `source`, compared verbatim with
    /// [`InboundEvent::source`].
    ///
    /// Correlation keys are business values another producer may know; a wait
    /// that names its sender cannot be consumed by an authenticated producer
    /// that merely guessed the kind and the key.
    ///
    /// An event that arrives over the served API or A2A carries
    /// `peer:<actor>`, the authenticated caller's actor as the deployment's
    /// authenticator names it — a `CloudEvent`'s own `source` attribute is not
    /// used. An event passed to [`Runtime::deliver`](crate::runtime::Runtime::deliver)
    /// in process carries whatever source its caller built it with.
    #[must_use]
    pub fn from(mut self, source: impl Into<String>) -> Self {
        self.from = Some(source.into());
        self
    }
}

/// A durable registration of interest in a future event.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Subscription {
    pub run: RunId,
    /// The matter this wait belongs to.
    ///
    /// Carried so delivery can stamp the record it writes with the case — every
    /// record of a case-bound run must — and so "what is this matter waiting
    /// for?" is one query rather than a join through runs.
    pub case: Option<crate::core::CaseId>,
    /// The effect whose output this event will become. Delivering the event
    /// journals an `EffectDone` under this key, so resuming the run replays the
    /// wait as an ordinary completed effect.
    pub effect: EffectKey,
    /// Which step is waiting, and in which pass.
    ///
    /// Delivery journals the `EffectDone` under this position, and replay
    /// verifies effects per step — so a wait recorded against the wrong step is
    /// a wait the resumed run never finds. A fixed step here would hold only
    /// for single-step plans; a wait in a later step, or in a compensation,
    /// would suspend forever.
    pub step: crate::core::StepId,
    pub phase: crate::core::Phase,
    pub kind: String,
    pub correlation: Vec<CorrelationKey>,
    /// The one [`InboundEvent::source`] this wait accepts; `None` accepts any.
    /// Every delivery path — broadcast, buffered and targeted — holds it.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub from: Option<String>,
}

impl Subscription {
    /// Whether an event from `source` may satisfy this wait.
    #[must_use]
    pub fn accepts_source(&self, source: &str) -> bool {
        self.from.as_deref().is_none_or(|from| from == source)
    }
}

/// A run's durable wake-up.
///
/// A timer is a wait whose event is the clock. It reuses the effect protocol
/// wholesale — the wake instant is journaled under an effect key, so replay
/// reads it back rather than sleeping again, and a fired timer is recorded
/// before the run is resumed.
///
/// Unlike a correlated wait it needs no case: there is nothing to correlate and
/// no business horizon to bound it, because the instant *is* the horizon. That
/// makes durable sleep available to any run, which is what lets a long retry
/// backoff release its worker instead of holding one.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct Timer {
    pub run: crate::core::RunId,
    /// Carried so a fired timer can stamp the record it writes with the case —
    /// every record of a case-bound run must.
    pub case: Option<crate::core::CaseId>,
    /// The effect whose output this wake-up becomes.
    pub effect: EffectKey,
    /// Where the sleeping step is. Replay verifies effects per step, so a
    /// wake-up journaled against the wrong step is one the resumed run never
    /// finds.
    pub step: crate::core::StepId,
    pub phase: crate::core::Phase,
    pub fire_at: Timestamp,
}

/// Why a run stopped without finishing.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize, schemars::JsonSchema)]
#[serde(tag = "reason", rename_all = "snake_case", deny_unknown_fields)]
#[non_exhaustive]
pub enum SuspendReason {
    /// Waiting for an inbound message that has not arrived.
    AwaitingEvent {
        kind: String,
        correlation: Vec<CorrelationKey>,
        #[serde(with = "time::serde::rfc3339")]
        #[schemars(with = "String")]
        until: Timestamp,
    },
    /// Waiting for an instant to arrive.
    ///
    /// Nothing to correlate: the clock is the event. Distinct from
    /// [`AwaitingEvent`](Self::AwaitingEvent) because the two fail differently —
    /// a message that never comes is a correlation bug worth alerting on, and an
    /// instant that has not arrived yet is the system working.
    AwaitingTime {
        #[serde(with = "time::serde::rfc3339")]
        #[schemars(with = "String")]
        until: Timestamp,
    },
}

impl SuspendReason {
    /// When the wait stops being the system working.
    ///
    /// Both variants carry one, and neither means the same thing by it: for a
    /// timer it is when the run is due to continue, and for an event it is when
    /// waiting stops being reasonable. Shared here because the question a
    /// listing of waiting runs is ordered by is the same either way — *which of
    /// these should have moved by now*.
    #[must_use]
    pub const fn until(&self) -> Timestamp {
        match self {
            Self::AwaitingEvent { until, .. } | Self::AwaitingTime { until } => *until,
        }
    }
}

impl std::fmt::Display for SuspendReason {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            Self::AwaitingEvent {
                kind,
                correlation,
                until,
                ..
            } => {
                write!(f, "awaiting '{kind}'")?;
                if let Some(k) = correlation.first() {
                    write!(f, " for {k}")?;
                }
                write!(f, " until {until}")
            }
            Self::AwaitingTime { until } => write!(f, "sleeping until {until}"),
        }
    }
}

/// What happened to a delivered event.
#[derive(Debug, Clone, PartialEq)]
#[non_exhaustive]
pub enum Delivery {
    /// A waiting run consumed it and ran to its next stopping point.
    Resumed { run: RunId },
    /// Stored durably, and not yet driven onward by this call.
    ///
    /// Not an error and not a dead letter. Either nobody is waiting for it
    /// *yet* — the counterpart run may not have reached its wait, or may not
    /// have started, and the event stays claimable until a wait finds it or the
    /// sweep ages it out — or it was recorded as a waiting run's answer by a
    /// plane that could not continue that run (its owner was busy, or this
    /// plane holds no agent for it), and the run's own plane, the sweep or a
    /// replay continues it.
    Buffered,
    /// This event id was already delivered. Retries are safe.
    Duplicate,
}

impl Delivery {
    #[must_use]
    pub fn resumed_run(&self) -> Option<RunId> {
        match self {
            Self::Resumed { run } => Some(*run),
            _ => None,
        }
    }
}

/// An event that aged out of the buffer without anyone claiming it.
///
/// Reaching this list means a correlation key is wrong somewhere — the message
/// arrived, was held, and no run ever asked for it. That is a bug worth paging
/// on, and it is precisely the failure that otherwise presents as a process
/// silently never completing.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct DeadLetter {
    pub event: InboundEvent,
    #[serde(with = "time::serde::rfc3339")]
    pub received_at: Timestamp,
    pub reason: String,
}

#[cfg(test)]
mod tests {
    use super::*;

    /// `resumed_run` names the run only when one actually resumed.
    ///
    /// The three delivery outcomes are easy to conflate at a call site — all
    /// three mean "the event was accepted" — and only one of them means a run
    /// moved. A caller that treated `Buffered` as a resumption would report
    /// progress for a message nobody has claimed yet, which is precisely the
    /// failure the buffered state exists to make visible.
    #[test]
    fn only_a_resumed_delivery_names_a_run() {
        let run = RunId::generate();

        assert_eq!(Delivery::Resumed { run }.resumed_run(), Some(run));
        assert_eq!(Delivery::Buffered.resumed_run(), None);
        assert_eq!(Delivery::Duplicate.resumed_run(), None);
    }

    /// The source is read back whatever either half contains.
    #[test]
    fn an_origin_keys_source_reads_back_exactly() {
        for (source, id) in [
            ("peer:a", "m1"),
            ("a\u{1f}b", "c"),
            ("a", "b\u{1f}c"),
            ("", "x"),
        ] {
            assert_eq!(origin_source(&origin_key(source, id)), Some(source));
        }
        for key in [
            "",
            "m1",
            "x\u{1f}peer:a\u{1f}m1",
            "9\u{1f}peer:a\u{1f}m1",
            "6\u{1f}peer:a",
        ] {
            assert_eq!(origin_source(key), None, "{key:?} is not an origin key");
        }
    }
}