agentplane 0.39.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
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
//! Cases — durable state that outlives a run.
//!
//! A run is one goal, one plan, one lifetime. Real business processes are not
//! that. A supplier switch spans days: a request goes out, an acknowledgement
//! must arrive inside a regulatory window, a confirmation or rejection follows,
//! a cancellation may arrive later, and an invoice dispute may land weeks after
//! that. Each is a *separate inbound trigger at an unpredictable time*, and all
//! of them belong to **one business fact**.
//!
//! # Why a case rather than one very long-lived run
//!
//! Durable-execution engines usually model this as a workflow that lives for
//! weeks. That is a versioning trap: a six-week workflow pins your code version
//! for six weeks, and every deploy needs a migration story for in-flight
//! instances.
//!
//! agentplane inverts it — **runs stay short, longevity lives in the case.**
//! Runs are minutes; cases are months; deploys are free. The cost is that
//! continuity must be explicit (case state, not local variables), which is the
//! right trade when the alternative is an auditor asking about a process whose
//! code no longer exists.

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

use crate::core::{CaseId, Digest, RunId, Timestamp};

/// A business key that identifies a case from the outside.
///
/// Inbound messages do not know run ids. They carry document numbers, meter
/// ids, order references — so that is what correlation matches on.
#[derive(Debug, Clone, PartialEq, Eq, Hash, PartialOrd, Ord, Serialize, Deserialize)]
pub struct CorrelationKey {
    /// What kind of identifier this is, e.g. `"document-number"`, `"meter"`.
    pub namespace: String,
    pub value: String,
}

impl CorrelationKey {
    pub fn new(namespace: impl Into<String>, value: impl Into<String>) -> Self {
        Self {
            namespace: namespace.into(),
            value: value.into(),
        }
    }
}

impl std::fmt::Display for CorrelationKey {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        write!(f, "{}={}", self.namespace, self.value)
    }
}

/// Where a case is in its life.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum CaseStatus {
    Open,
    /// Waiting for an inbound message that has not arrived.
    AwaitingExternal,
    /// Waiting for a person.
    AwaitingHuman,
    /// An obligation was missed and someone has been told.
    Escalated,
    Closed,
}

impl CaseStatus {
    #[must_use]
    pub fn as_str(self) -> &'static str {
        match self {
            Self::Open => "open",
            Self::AwaitingExternal => "awaiting_external",
            Self::AwaitingHuman => "awaiting_human",
            Self::Escalated => "escalated",
            Self::Closed => "closed",
        }
    }

    /// The inverse of [`as_str`](Self::as_str), for a status arriving as text.
    ///
    /// **The reason every enum in this crate that reaches a stored column has
    /// this pair, and the reason nothing else may spell the vocabulary.**
    /// Written over [`ALL`](Self::ALL) rather than as a second `match`, so the
    /// two directions cannot disagree. `as_str` is exhaustive because the
    /// compiler insists; a hand-written reader is *total*, so it keeps
    /// compiling when a variant is added and starts refusing — or worse,
    /// silently defaulting — a value the writer beside it emits. Every store
    /// backend reads through here for that reason: two backends deciding a
    /// vocabulary apart is two answers to one question about the same data.
    ///
    /// Unknown text is `None` rather than a default. A status filter that
    /// silently became `Open` on a typo would answer *what is escalated* with a
    /// list of healthy cases, which reads as an empty backlog.
    #[must_use]
    pub fn parse(s: &str) -> Option<Self> {
        Self::ALL.iter().copied().find(|c| c.as_str() == s)
    }

    /// Every status, so a caller can enumerate them without matching.
    pub const ALL: [Self; 5] = [
        Self::Open,
        Self::AwaitingExternal,
        Self::AwaitingHuman,
        Self::Escalated,
        Self::Closed,
    ];

    #[must_use]
    pub fn is_closed(self) -> bool {
        self == Self::Closed
    }
}

/// An instruction from outside the plane that a matter must be preserved.
///
/// Orthogonal to [`CaseStatus`], which answers *where is this matter* and moves
/// because runs did things. A hold answers *may this be destroyed*, decided by
/// somebody outside and usually before anybody knows which matters it covers.
/// A held matter still closes: one field for both would either end the hold when
/// the work does, or make closure unreachable and strand [`Deadline`].
///
/// The reason and the [`Operator`] are both required, for the same reason every
/// erasure verb takes a reason: the runtime cannot check a preservation order,
/// so the row is worth exactly what it says. A hold nobody can account for is
/// indistinguishable from a sweep that quietly stopped working, and one nobody
/// is named on cannot be asked about. The name comes from the credential an
/// authenticator verified or from whoever opened the store — never from the
/// request body, which is the requester's own word.
///
/// [`Operator`]: crate::core::Operator
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct LegalHold {
    /// When the hold was placed. Supplied by the caller, never read from a
    /// clock here, for the reason every other lifecycle instant in this crate
    /// is: a pass that read the clock itself could not be tested against a year
    /// of ageing holds.
    #[serde(with = "time::serde::rfc3339")]
    pub placed_at: Timestamp,
    /// Why this matter may not be destroyed. Free text, and load-bearing: it is
    /// what the person reviewing the hold listing months later has to act on.
    pub reason: String,
    /// Who placed it, and what established the name.
    pub by: crate::core::Operator,
}

/// A long-lived, correlated business fact.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Case {
    pub id: CaseId,
    /// Opaque to the engine — `"gpke.supplier-switch"` means nothing here.
    pub kind: String,
    pub status: CaseStatus,
    pub correlation: Vec<CorrelationKey>,
    /// Schema-validated per kind by the adapter; opaque to the engine.
    pub state: Value,
    /// Which revision of [`state`](Self::state) this is.
    ///
    /// Every write bumps it, and a write must name the version it read. See
    /// [`CaseVersion`].
    pub version: CaseVersion,
    #[serde(with = "time::serde::rfc3339")]
    pub opened_at: Timestamp,
    pub runs: Vec<RunId>,
}

/// Which revision of a case's state a reader saw.
///
/// # Why case state needs this and a run's journal does not
///
/// A run is owned: the fencing lease means exactly one writer appends to its
/// journal, so "read, decide, append" cannot interleave with anybody. A **case**
/// is the opposite by construction — it is the thing several runs share, over
/// days, and the topology this crate exists to serve has several plane instances
/// writing to one store.
///
/// The window between reading case state and writing it back therefore contains
/// an *inference*, which is unbounded. Classical lost-update reasoning assumes a
/// read-to-write window measured in milliseconds; here it is measured in however
/// long a model takes to answer, and two runs on one case will overlap. A blind
/// `UPDATE ... SET state = ?` in that window silently discards whichever write
/// lost the race, and nothing in the record shows it happened.
///
/// So a write names the version it read, the store rejects it if the case has
/// moved on ([`StoreError::CaseConflict`](crate::core::StoreError::CaseConflict)),
/// and the caller re-reads. The check is a database predicate rather than
/// application logic for the same reason exactly-once is: application logic can
/// be bypassed by the next caller, a constraint cannot.
#[derive(
    Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord, Default, Serialize, Deserialize, Hash,
)]
pub struct CaseVersion(pub u64);

impl CaseVersion {
    /// The version a case has before anybody has written to it.
    pub const INITIAL: Self = Self(0);

    #[must_use]
    pub const fn next(self) -> Self {
        Self(self.0 + 1)
    }
}

impl std::fmt::Display for CaseVersion {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        write!(f, "v{}", self.0)
    }
}

/// A domain-specific deadline description.
///
/// Deliberately opaque: "5 working days at 17:00 Europe/Berlin, excluding
/// public holidays observed in any federal state" is domain knowledge and does
/// not belong in a domain-agnostic engine. The engine carries the spec to a
/// [`Calendar`](crate::core::Calendar) and enforces whatever instant comes back.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct DeadlineSpec {
    /// Which resolution rule to apply, e.g. `"hours"`, `"working-days"`.
    pub kind: String,
    /// Parameters for that rule.
    pub params: Value,
}

impl DeadlineSpec {
    pub fn new(kind: impl Into<String>, params: Value) -> Self {
        Self {
            kind: kind.into(),
            params,
        }
    }

    /// A plain wall-clock offset, understood by the built-in calendar.
    ///
    /// One constructor per rule `WallClock` accepts, and the set is closed for
    /// a reason: a rule with no canonical spelling is where `"minutes"`,
    /// `"minute"` and `"mins"` diverge between an application and the calendar
    /// adapter meant to resolve them — and a calendar that does not recognise a
    /// kind refuses it, so the drift surfaces as a failed obligation rather
    /// than a wrong one. `minutes` was the rule that had this hole:
    /// `WallClock` resolved it and nothing here spelled it.
    #[must_use]
    pub fn minutes(n: u32) -> Self {
        Self::new("minutes", serde_json::json!({ "n": n }))
    }

    /// A plain wall-clock offset, understood by the built-in calendar.
    #[must_use]
    pub fn hours(n: u32) -> Self {
        Self::new("hours", serde_json::json!({ "n": n }))
    }

    /// Calendar days, understood by the built-in calendar.
    #[must_use]
    pub fn days(n: u32) -> Self {
        Self::new("days", serde_json::json!({ "n": n }))
    }
}

#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum DeadlineState {
    Pending,
    /// The warning threshold passed and an alert was emitted.
    Warned,
    /// The instant passed with the obligation unmet.
    Breached,
    /// Satisfied before the instant.
    Met,
    Cancelled,
}

impl DeadlineState {
    #[must_use]
    pub fn as_str(self) -> &'static str {
        match self {
            Self::Pending => "pending",
            Self::Warned => "warned",
            Self::Breached => "breached",
            Self::Met => "met",
            Self::Cancelled => "cancelled",
        }
    }

    /// The inverse of [`as_str`](Self::as_str), written over
    /// [`ALL`](Self::ALL) for the reason [`CaseStatus::parse`] is.
    #[must_use]
    pub fn parse(s: &str) -> Option<Self> {
        Self::ALL.iter().copied().find(|c| c.as_str() == s)
    }

    /// Every state an obligation can be in.
    pub const ALL: [Self; 5] = [
        Self::Pending,
        Self::Warned,
        Self::Breached,
        Self::Met,
        Self::Cancelled,
    ];

    /// Whether this deadline still constitutes an open obligation.
    #[must_use]
    pub fn is_open(self) -> bool {
        matches!(self, Self::Pending | Self::Warned)
    }

    /// Whether an obligation in this state is a breach still on the listing.
    ///
    /// The membership rule for that listing, on the vocabulary rather than
    /// beside any one reader, because its readers hold different things: the
    /// index that maintains it holds a stored spelling and a column, the
    /// listing holds a decoded obligation, and a gauge counts rows. Each
    /// deciding for itself is how a breach ends up on one backend's page and
    /// not the other's.
    ///
    /// The state is only half of it, which is the part that is easy to lose:
    /// what ends the *question* is somebody accounting for the breach, and the
    /// breach itself stays [`Breached`](Self::Breached) forever.
    ///
    /// `PostgreSQL` cannot call this — a partial index predicate must be
    /// immutable — so the SQL spells both conjuncts and a guard holds the two
    /// spellings to each other.
    #[must_use]
    pub const fn is_unaccounted(self, accounted: bool) -> bool {
        matches!(self, Self::Breached) && !accounted
    }

    /// Whether an obligation in this state may be moved to `to`.
    ///
    /// The three ways an obligation ends — met, missed, withdrawn — are
    /// terminal, and [`Breached`](Self::Breached) is the one that has to be.
    /// It is the only record that a window closed unmet, and
    /// [`set_deadline_state`] is reachable from a skill: `cx.meet_deadline` on
    /// an obligation the sweep had already breached would take the miss off the
    /// operator's listing and out of the row at once, leaving no account of it
    /// anywhere. Answering late is a fact to record beside the breach — an
    /// [`acknowledge_breach`] note says so — not a way to unsay it.
    ///
    /// Re-applying the state already held is allowed. Every writer here is a
    /// sweep or a resumed run, and both repeat their own last write by design:
    /// refusing that would turn an idempotent retry into a failure.
    ///
    /// Fails **closed** on a state this relation does not name, so a new
    /// variant is unmovable until somebody says where it may go rather than
    /// silently movable anywhere.
    ///
    /// [`set_deadline_state`]: crate::case::CaseStore::set_deadline_state
    /// [`acknowledge_breach`]: crate::case::CaseStore::acknowledge_breach
    #[must_use]
    pub const fn may_become(self, to: Self) -> bool {
        matches!(
            (self, to),
            (Self::Pending, Self::Pending | Self::Warned)
                | (Self::Warned, Self::Warned)
                | (Self::Met, Self::Met)
                | (Self::Breached, Self::Breached)
                | (Self::Cancelled, Self::Cancelled)
                | (
                    Self::Pending | Self::Warned,
                    Self::Met | Self::Breached | Self::Cancelled
                )
        )
    }
}

/// Somebody's account of a breach.
///
/// A missed obligation is a fact and does not stop being one. What can end is
/// the *question* it puts to an operator — has anybody looked at this. The
/// obligation keeps [`Breached`](DeadlineState::Breached); the listing that asks
/// the question stops returning it. Why that listing has to drain is on
/// [`CaseStore::breached`].
///
/// [`CaseStore::breached`]: crate::case::CaseStore::breached
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BreachNote {
    /// The authenticated actor who accounted for it. Never taken from a request
    /// body: an account of a missed obligation that names whoever the caller
    /// said they were is not an account.
    pub by: String,
    /// What they had to say. May be empty — an operator who has looked and has
    /// nothing to add has still answered the question the listing asked.
    pub note: String,
    #[serde(with = "time::serde::rfc3339")]
    pub at: Timestamp,
}

/// A registered obligation with a resolved instant.
///
/// # The instant is a fact, not a formula
///
/// `resolved_at` is stored, and never recomputed. Calendars change — a
/// corrected holiday table, a new regulatory notice — and recomputing on replay
/// would silently move a legally binding instant under an audit. The
/// `calendar_digest` records which calendar version produced it, so a shifted
/// rule is *visible* rather than retroactive.
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct Deadline {
    pub case: CaseId,
    /// Unique within the case.
    pub name: String,
    #[serde(with = "time::serde::rfc3339")]
    pub resolved_at: Timestamp,
    /// Which calendar version produced `resolved_at`.
    pub calendar_digest: Digest,
    #[serde(default, with = "time::serde::rfc3339::option")]
    pub warn_at: Option<Timestamp>,
    pub state: DeadlineState,
    /// Who accounted for the breach, once somebody has.
    ///
    /// Only ever set on a [`Breached`](DeadlineState::Breached) obligation, and
    /// it is what takes the breach off the obligation listing — see
    /// [`CaseStore::acknowledge_breach`]. `None` on every other state, because
    /// there is nothing to account for.
    ///
    /// [`CaseStore::acknowledge_breach`]: crate::case::CaseStore::acknowledge_breach
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub acknowledged: Option<BreachNote>,
}

impl Deadline {
    #[must_use]
    pub fn is_due(&self, now: Timestamp) -> bool {
        self.state.is_open() && now >= self.resolved_at
    }

    /// A breach nobody has accounted for yet.
    ///
    /// The decoded form of [`DeadlineState::is_unaccounted`], which carries the
    /// rule. What a caller holding a whole obligation asks — including the
    /// embedded backend's own listing, which reads an index and then holds the
    /// row it decoded to this, because an index is a hint and the row is the
    /// record.
    #[must_use]
    pub const fn is_unaccounted(&self) -> bool {
        self.state.is_unaccounted(self.acknowledged.is_some())
    }

    #[must_use]
    pub fn needs_warning(&self, now: Timestamp) -> bool {
        self.state == DeadlineState::Pending
            && self.warn_at.is_some_and(|w| now >= w)
            && now < self.resolved_at
    }
}

/// What the sweeper did to something nobody was watching.
///
/// # Why these are on the record and not only in a log
///
/// The sweeper makes the plane's most consequential *automated* decisions: it
/// breaches an obligation, escalates a case, expires a person's task. Nothing
/// asked it to — that is the point of it — so there is no run whose history
/// explains why the state changed.
///
/// Without a record, *why is this case escalated* is answerable only from the
/// resulting state, and state cannot distinguish "the sweep breached this at
/// 02:00" from "somebody set it". That is the same distinction
/// [`RecordKind::StepCompensated`](crate::journal::RecordKind::StepCompensated)
/// exists for, and it matters more here because no human was present.
///
/// A typed enum rather than a message: an operator alerting on breaches should
/// not be matching on prose, and a variant added here is one every reader must
/// consider.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum SweptAction {
    /// An obligation is approaching and its warning instant passed.
    DeadlineWarned,
    /// An obligation passed unmet.
    DeadlineBreached,
    /// A case was escalated because one of its obligations was breached.
    CaseEscalated,
    /// A person's task window closed and the declared policy was applied.
    TaskExpired,
    /// A task's audience was widened because nobody answered in time.
    TaskEscalated,
    /// A run an instance died holding was taken over and resumed.
    ///
    /// The one action in this vocabulary that is the system healing rather
    /// than a finding — recorded anyway, because a takeover bumps the run's
    /// epoch and fences its previous owner, and *who fenced whom, and why*
    /// must be answerable from the journal rather than inferred from an epoch
    /// gap. The detail names the outcome the resume reached.
    RunRecovered,
}