agentplane 0.45.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
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
//! Batch runs — one plan, many items, per-item durability.
//!
//! # Why a batch is not a run, and not N unrelated runs
//!
//! A Jahresabrechnung over 10⁵ `Marktlokationen`, or a `MaBiS` `Clearingliste` across
//! a Bilanzierungsgebiet, is a single business act made of many independent
//! ones. Modelling it as *one run* means one journal, one budget, and one
//! failure: item 60,000 fails and the run fails, so the 59,999 settlements that
//! worked are trapped inside a failed audit record. Modelling it as *N unrelated
//! runs* loses the act — nobody can answer "did the Jahresabrechnung finish", or
//! "what did it cost", because there is no it.
//!
//! So a batch is a first-class object that owns N runs sharing one frozen plan.
//! Each item gets its own journal, its own budget, and its own outcome. The
//! batch gets a cursor, a census of what happened, and a terminal state.
//!
//! # Partial failure is a terminal state, not a degraded success
//!
//! [`BatchStatus`] has no `Succeeded`. A finished batch is
//! [`BatchStatus::Completed`] carrying counts, so a caller cannot write
//! `if status.is_ok()` and skip the 43 items that did not settle. "Mostly
//! worked" is the single most dangerous thing a batch can report, because it is
//! reported as success everywhere it is not explicitly handled — and the items
//! that failed are precisely the ones a human needed to hear about.
//!
//! Reading the outcome forces a decision:
//!
//! ```text
//! match report.status {
//!     BatchStatus::Completed { failed: 0, quarantined: 0, succeeded } => ok(succeeded),
//!     BatchStatus::Completed { failed, quarantined, .. } => escalate(failed, quarantined),
//!     BatchStatus::Running => still_going(),
//! }
//! ```
//!
//! # Resume is the effect protocol, one level up
//!
//! An item is processed exactly the way an effect is performed: **announce, act,
//! record.** Before an item runs, its run id is written to the batch store
//! ([`BatchStore::reserve`]); then the run happens; then the outcome is recorded.
//!
//! That ordering is what makes item-granular resume work, and it is worth being
//! precise about why. A crash between reserve and record leaves an item marked
//! started with a known run id. Resume finds it, and **replays that run** rather
//! than starting a new one — so the item's effects are read back from its
//! journal instead of performed again. Exactly-once for a batch item is not new
//! machinery; it is the run-level guarantee, addressed by a stored id.
//!
//! The cursor is therefore an *optimisation*, not a correctness mechanism. If it
//! were lost entirely, re-processing every item would be safe — slow, but not
//! wrong. That is the property to preserve when changing anything here.
//!
//! # Input is streamed, never materialised
//!
//! [`ItemSource`] is a cursor, not a `Vec`. A batch over 10⁵ meters must not
//! require 10⁵ items in memory, and — more subtly — must not require the caller
//! to have *produced* them all before the first one runs. A source that reads
//! from a database, a file, or a paged API is the normal case.

use std::fmt::Debug;

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

use crate::core::{BatchId, Label, RunId, Spend, StoreError};

/// One unit of work in a batch.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub struct BatchItem {
    /// Stable, unique within the batch, and **ordered**.
    ///
    /// Ordering is what lets the cursor be a single value rather than a set of
    /// seen keys, which for 10⁵ items is the difference between a resume that
    /// reads one row and one that reads the whole batch. A meter number, a
    /// document id, or a zero-padded sequence all work; a random uuid does not,
    /// because "everything after this key" then means nothing.
    pub key: String,
    /// What the plan receives as run input for this item.
    pub input: Value,
    /// The input's label, when the source vouches for one.
    ///
    /// `None` — the default — admits the item **untrusted**, from the source
    /// `batch:<batch id>`. A source reading a paged API, a file or somebody
    /// else's database is the normal case, and an item it produced is data
    /// from outside the plane: trusted input would reach every mutating sink
    /// and every protected field with no release on the record.
    #[serde(default, skip_serializing_if = "Option::is_none")]
    pub label: Option<Label>,
}

impl BatchItem {
    /// An item admitted untrusted, from the source `batch:<batch id>`.
    pub fn new(key: impl Into<String>, input: Value) -> Self {
        Self {
            key: key.into(),
            input,
            label: None,
        }
    }

    /// Admit this item under `label` instead.
    ///
    /// The operator's statement about where the source's data came from,
    /// journaled as the run's input label. `Label::trusted()` is the one to
    /// think about: it says the plane's own operator produced this input, and
    /// every taint gate downstream takes that at its word.
    #[must_use]
    pub fn labelled(mut self, label: Label) -> Self {
        self.label = Some(label);
        self
    }

    /// The label this item's run is admitted under, in batch `batch`.
    #[must_use]
    pub fn admission_label(&self, batch: BatchId) -> Label {
        self.label.clone().unwrap_or_else(|| {
            Label::untrusted(crate::core::SourceId::new(format!("batch:{batch}")))
        })
    }
}

/// A source could not produce items.
///
/// Distinct from `StoreError` because an [`ItemSource`] is the embedder's code
/// reading the embedder's system — a meter register, a paged API, a file. Making
/// it a store error would say the *plane's* storage failed, sending an operator
/// to look at the wrong thing.
#[derive(Debug, Clone, thiserror::Error)]
#[error("batch source: {0}")]
pub struct SourceError(String);

impl SourceError {
    pub fn new(detail: impl Into<String>) -> Self {
        Self(detail.into())
    }
}

/// Where a batch's work comes from.
///
/// A cursor rather than a collection — see the module docs on why input is never
/// materialised.
#[async_trait]
pub trait ItemSource: Send + Sync + Debug {
    /// The next page of items strictly after `after`, in strictly increasing
    /// byte order of [`BatchItem::key`].
    ///
    /// Byte order, not a collation or a numeric sort: the resume cursor is the
    /// greatest finished key compared as bytes, so a page in any other order
    /// resumes past items that never ran. The driver refuses a page whose keys
    /// do not strictly increase past `after`. Zero-pad numeric keys.
    ///
    /// Returning fewer than `limit` items does **not** mean the source is
    /// exhausted; returning zero does. The distinction matters for sources that
    /// page unevenly, and conflating them truncates a batch silently — which is
    /// the failure mode this whole crate exists to make loud.
    async fn next(&self, after: Option<&str>, limit: usize) -> Result<Vec<BatchItem>, SourceError>;
}

/// How one item ended.
///
/// Mirrors `RunStatus` but deliberately coarser: a batch's census is about what
/// a human must do next, and "failed" and "exhausted" both mean *this item did
/// not settle*. The item's own run journal holds the detail.
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ItemOutcome {
    Succeeded,
    /// Ran and did not settle — a failure, or a limit.
    Failed(String),
    /// Set aside for a human. Distinct from failed because the *response* is
    /// different: a failed item may be re-run, a quarantined one must not be
    /// touched until someone has looked at it.
    Quarantined(String),
    /// Waiting on something — an event, a person, a timer.
    ///
    /// Not terminal. A batch that suspends items is not finished, and counting a
    /// suspension as a failure would send someone to investigate a run that is
    /// working exactly as designed.
    Suspended(String),
    /// Paused at a ceiling. Not terminal, and deliberately not `Failed`: the
    /// item's run is intact and stays open — an operator's two honest moves
    /// are raise the ceiling and resume, or cancel — and a pause reported as
    /// a fault teaches the reader to re-run work that is standing. The string
    /// names the ceiling that was hit.
    Exhausted(String),
    /// Paused because an authority the item acts under was withdrawn. Not
    /// terminal, and not `Exhausted`: the work stands as it does at a ceiling,
    /// but the move is lifting the halt on that subject rather than raising a
    /// budget. The string names the subject and why it was withdrawn.
    Withheld(String),
}

impl ItemOutcome {
    #[must_use]
    pub const fn as_str(&self) -> &'static str {
        match self {
            Self::Succeeded => "succeeded",
            Self::Failed(_) => "failed",
            Self::Quarantined(_) => "quarantined",
            Self::Suspended(_) => "suspended",
            Self::Exhausted(_) => "exhausted",
            Self::Withheld(_) => "withheld",
        }
    }

    /// Every outcome, each carrying `detail`.
    ///
    /// The list [`parse`](Self::parse) is written over, and the list a store
    /// derives its terminal set from — so a query asking *what is still open*
    /// answers from the type rather than from a list of strings beside it.
    #[must_use]
    pub fn all(detail: String) -> [Self; 6] {
        [
            Self::Succeeded,
            Self::Failed(detail.clone()),
            Self::Quarantined(detail.clone()),
            Self::Suspended(detail.clone()),
            Self::Exhausted(detail.clone()),
            Self::Withheld(detail),
        ]
    }

    /// The inverse of [`as_str`](Self::as_str), written over
    /// [`all`](Self::all) for the reason [`CaseStatus::parse`] is.
    ///
    /// `None` is damage rather than absence: an item whose outcome cannot be
    /// read has not *failed to run*, and a census filing it as in-flight keeps
    /// the batch `Running` forever over a row nobody can explain.
    ///
    /// [`CaseStatus::parse`]: crate::core::CaseStatus::parse
    #[must_use]
    pub fn parse(tag: &str, detail: String) -> Option<Self> {
        Self::all(detail).into_iter().find(|o| o.as_str() == tag)
    }

    #[must_use]
    pub const fn is_terminal(&self) -> bool {
        !matches!(
            self,
            Self::Suspended(_) | Self::Exhausted(_) | Self::Withheld(_)
        )
    }

    /// The stored spellings of every outcome that ends an item.
    ///
    /// What a store asks for when it wants the complement: an item is open
    /// while it has no outcome or an outcome outside this set, so a
    /// non-terminal variant added later is open by construction rather than by
    /// somebody remembering to widen a list of strings.
    #[must_use]
    pub fn terminal_tags() -> Vec<&'static str> {
        Self::all(String::new())
            .iter()
            .filter(|o| o.is_terminal())
            .map(Self::as_str)
            .collect()
    }

    /// Whether this item is finished *and* needs nothing from anybody.
    ///
    /// Terminal and settled are different questions, and the difference is the
    /// backlog: `Failed` and `Quarantined` are terminal and are exactly what a
    /// person has to see. Written as a match rather than `== Succeeded` so a
    /// variant added later is unsettled by construction — the safe direction.
    #[must_use]
    pub const fn is_settled(&self) -> bool {
        match self {
            Self::Succeeded => true,
            Self::Failed(_)
            | Self::Quarantined(_)
            | Self::Suspended(_)
            | Self::Exhausted(_)
            | Self::Withheld(_) => false,
        }
    }

    /// The stored spellings of every outcome nobody has to act on.
    ///
    /// The set a store negates to build [`items_needing_attention`]. Derived
    /// from [`all`](Self::all) for the reason
    /// [`terminal_tags`](Self::terminal_tags) is: a hand-kept list of strings
    /// stops naming a class of trouble the day somebody adds one.
    ///
    /// [`items_needing_attention`]: BatchStore::items_needing_attention
    #[must_use]
    pub fn settled_tags() -> Vec<&'static str> {
        Self::all(String::new())
            .iter()
            .filter(|o| o.is_settled())
            .map(Self::as_str)
            .collect()
    }
}

/// What a batch has done, and whether it is done.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum BatchStatus {
    /// Items remain, or some are still suspended.
    Running,
    /// Every item reached a terminal outcome.
    ///
    /// Note the absence of a `Succeeded` variant: see the module docs. A caller
    /// must read the counts to know what happened, because there is no way to
    /// spell "it worked" that skips them.
    Completed {
        succeeded: u64,
        failed: u64,
        quarantined: u64,
    },
}

impl BatchStatus {
    /// Whether every item settled.
    ///
    /// Named for what it asserts rather than `is_ok`: the question a caller
    /// usually means is "is there anything to do", and a batch with 43 failures
    /// answers yes.
    #[must_use]
    pub const fn everything_settled(&self) -> bool {
        matches!(
            self,
            Self::Completed {
                failed: 0,
                quarantined: 0,
                ..
            }
        )
    }
}

/// One item's record within a batch.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct ItemRecord {
    pub key: String,
    /// The run that processed it. Stable across resume — replaying this id is
    /// what makes re-processing safe.
    pub run: RunId,
    /// `None` while the item is reserved but not yet finished.
    pub outcome: Option<ItemOutcome>,
    /// What this item consumed, so "what did the settlement run cost" is a sum
    /// rather than an estimate.
    pub spend: Spend,
}

/// What a batch cost and how it ended.
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct BatchReport {
    pub id: BatchId,
    pub status: BatchStatus,
    /// Items reserved but not yet terminal — suspended, or interrupted.
    pub in_flight: u64,
    /// Items paused at a ceiling, resumable once somebody raises it.
    ///
    /// Separate from [`in_flight`](Self::in_flight) because the move differs: an
    /// interrupted item resumes when a sweep reaches it, this one waits for a
    /// decision. Its own field because both leave the batch
    /// [`Running`](BatchStatus::Running), and so does a deliberate `max_items`
    /// window — only the counts tell those apart.
    pub exhausted: u64,
    /// Items paused because an authority they act under was withdrawn —
    /// resumable once the halt on that subject is lifted.
    pub withheld: u64,
    /// The whole batch's consumption, summed from its items.
    pub spend: Spend,
    /// The key processing stopped after, for a resume.
    pub cursor: Option<String>,
}

impl BatchReport {
    /// Whether anything here needs a person — the alerting predicate.
    ///
    /// Deliberately **not** *is this batch unfinished*. Every windowed pass
    /// returns [`Running`](BatchStatus::Running), so a status-keyed predicate
    /// would be true on every correct use of [`BatchSpec::max_items`], and a
    /// predicate that is always true is one people stop reading. It reads the
    /// facts somebody acts on instead: `in_flight`, `exhausted`, and the failed
    /// or quarantined counts.
    ///
    /// [`BatchStore::items_needing_attention`] is the listing behind the answer;
    /// a predicate that says *yes* and cannot say *which* is a counter with no
    /// alert.
    ///
    /// [`BatchSpec::max_items`]: crate::runtime::BatchSpec::max_items
    #[must_use]
    pub const fn needs_attention(&self) -> bool {
        self.in_flight > 0
            || self.exhausted > 0
            || self.withheld > 0
            || self.failed_or_quarantined() > 0
    }

    /// The terminal items that did not settle.
    #[must_use]
    pub const fn failed_or_quarantined(&self) -> u64 {
        match self.status {
            BatchStatus::Running => 0,
            BatchStatus::Completed {
                failed,
                quarantined,
                succeeded: _,
            } => failed + quarantined,
        }
    }
}

/// Durable batch state.
///
/// Separate from the journal because a batch is not a run: it has no hash chain
/// of its own, and its integrity comes from the per-item runs it points at. What
/// it needs is a cursor, a reservation, and a census — three queries, not a log.
#[async_trait]
pub trait BatchStore: Send + Sync + Debug {
    /// Which tenant this handle's batches and item reservations belong to.
    fn tenant(&self) -> &str;

    /// Register a batch. Idempotent on `id`, so a retried submission does not
    /// fork one act into two.
    ///
    /// Idempotent on the id **and held to the digest**: one batch runs one
    /// frozen plan, and that sentence is this method's to enforce, because the
    /// runner cannot — by the time it executes an item, the store's row is the
    /// only witness to what the batch was opened with.
    ///
    /// # Errors
    ///
    /// [`StoreError::BatchPlanChanged`] when the batch exists under a
    /// different `plan_digest`. A resume offering an edited plan must be
    /// refused here, in the store, or items settle under a plan the batch's
    /// record does not name.
    async fn open(&self, id: BatchId, plan_digest: &str) -> Result<(), StoreError>;

    /// The plan digest this batch was opened with, or `None` for no such batch.
    ///
    /// The existence question, answered from the batch's own row. A census
    /// cannot answer it — a batch with no items yet and a batch that does not
    /// exist both count zero rows — and the difference matters to an operator:
    /// one is work not started, the other is a mistyped id.
    async fn plan_digest(&self, id: BatchId) -> Result<Option<String>, StoreError>;

    /// Record that the source produced its last item.
    ///
    /// Without this a batch cannot tell "every item I have stored is terminal"
    /// from "I am finished" — and those differ every time processing stops
    /// early. A batch halted after 10,000 of 100,000 items has no unfinished
    /// item anywhere in its store, so a census alone would report it complete
    /// with 90,000 meters unsettled.
    ///
    /// Durable rather than in-memory because the distinction has to survive the
    /// process: a resumed batch must know whether it ever reached the end.
    ///
    /// # Errors
    ///
    /// [`StoreError::NotFound`] when no such batch exists, for the reason
    /// [`record`](Self::record) refuses an unreserved item: this mark is the
    /// one bit that lets a census read as *finished*, and reporting it written
    /// when nothing was is the quietest possible way to lose it.
    async fn mark_exhausted(&self, id: BatchId) -> Result<(), StoreError>;

    /// Whether the source has been read to the end.
    async fn is_exhausted(&self, id: BatchId) -> Result<bool, StoreError>;

    /// Claim an item and bind it to a run id, **before** the run starts.
    ///
    /// Returns the existing record if the item was already reserved — which is
    /// what makes a crashed batch resumable: the second attempt gets the first
    /// attempt's run id back and replays it rather than starting fresh.
    async fn reserve(
        &self,
        batch: BatchId,
        key: &str,
        run: RunId,
    ) -> Result<ItemRecord, StoreError>;

    /// Record how an item ended, and what it consumed.
    ///
    /// # Errors
    ///
    /// [`StoreError::NotFound`] when the item was never reserved — a refusal,
    /// never a silent no-op: `Ok` over a write that matched nothing tells the
    /// caller *recorded* about an outcome that vanished, the same lie a
    /// release that freed nothing tells. The row count the write already
    /// produces is what catches it.
    async fn record(
        &self,
        batch: BatchId,
        key: &str,
        outcome: &ItemOutcome,
        spend: Spend,
    ) -> Result<(), StoreError>;

    /// The highest key whose item reached a terminal outcome, with no
    /// unfinished item before it.
    ///
    /// The contiguous prefix, not the maximum: an item suspended at key 400 must
    /// hold the cursor at 399 even if 401 through 500 have finished, or a resume
    /// would step over it and the batch would report complete with an item still
    /// waiting.
    async fn cursor(&self, batch: BatchId) -> Result<Option<String>, StoreError>;

    /// Counts by outcome, plus reserved-but-unfinished.
    async fn census(&self, batch: BatchId) -> Result<BatchCensus, StoreError>;

    /// Every item record, oldest key first. For operators and for tests; the
    /// driver uses `cursor` and `census`.
    async fn items(&self, batch: BatchId, limit: usize) -> Result<Vec<ItemRecord>, StoreError>;

    /// The items that are not settled, oldest key first.
    ///
    /// [`census`](Self::census) answers *43 failed*; the question an operator
    /// has is **which 43**, and over the size a batch exists for the only other
    /// route is paging [`items`](Self::items) through a hundred thousand rows
    /// that are almost all successes — the finding indexed and reaching nobody.
    ///
    /// Unsettled, not un-terminal: a failed item and a suspended one are both
    /// here, because *finished* and *finished with* are different questions.
    /// Each record's [`ItemOutcome`] says whether the next move is a re-run, a
    /// raised ceiling, or a look.
    ///
    /// Oldest-first is legitimate because entries **leave**: resolving an item
    /// drops it. An ascending page over a listing nothing empties would have a
    /// permanent head and an unreachable tail.
    ///
    /// # Errors
    ///
    /// Backend failures. An unknown batch is an empty listing, not an error:
    /// [`plan_digest`](Self::plan_digest) is the existence question.
    async fn items_needing_attention(
        &self,
        batch: BatchId,
        limit: usize,
    ) -> Result<Vec<ItemRecord>, StoreError>;
}

/// A batch's tally.
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct BatchCensus {
    pub succeeded: u64,
    pub failed: u64,
    pub quarantined: u64,
    pub suspended: u64,
    /// Paused at a ceiling — resumable once somebody raises it.
    pub exhausted: u64,
    /// Paused under a withdrawn authority — resumable once the halt lifts.
    pub withheld: u64,
    /// Reserved with no outcome recorded — an item interrupted mid-flight.
    pub in_flight: u64,
    pub spend: Spend,
}

impl BatchCensus {
    /// Items that will not change without intervention.
    #[must_use]
    pub const fn terminal(&self) -> u64 {
        self.succeeded + self.failed + self.quarantined
    }
}