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
//! One contract, run against every quota store.
//!
//! A ceiling is arithmetic at a boundary, and boundaries are where two
//! implementations of one rule diverge. This battery exists because they already
//! did: one backend compared "am I at the limit?" *inside* its counting loop,
//! which is correct for every ceiling except zero — the loop body never runs, so
//! nothing is compared, and a tenant stopped dead is admitted instead. The other
//! backend had it right. Only a shared contract catches that shape.
//!
//! The properties are the ones a ceiling stands or falls on:
//!
//! * a limit of **zero** admits nothing, because that is how an operator stops a
//!   tenant;
//! * a run at the ceiling is **refused**, and refused with the count, so an
//!   operator can tell throttling from a fault;
//! * releasing makes room, because a ceiling is back-pressure and not a
//!   permanent verdict;
//! * reserving one run **twice takes one slot**, or a retried admission costs
//!   the tenant capacity forever;
//! * accruals **sum**, because reading, adding and writing back loses one of two
//!   concurrent updates — and what it loses is spend already incurred;
//! * periods are **independent**, since the window is a billing period.

use crate::core::{RunId, Spend, Timestamp};
use crate::quota::{QuotaError, QuotaSettlement, QuotaStore};

use super::conformance::Report;

/// Run the battery against one quota store.
///
/// The store must be scoped to a tenant with no runs and no recorded spend: this
/// reserves, releases and accrues under it.
pub async fn check(store: &dyn QuotaStore, report: &mut Report) {
    let at = Timestamp::from_unix_timestamp(1_760_000_000).expect("a valid test instant");
    zero_admits_nothing(store, at, report).await;
    ceiling_refuses_and_frees(store, at, report).await;
    reserving_twice_takes_one_slot(store, at, report).await;
    spend_accrues_per_period(store, report).await;
    halt(store, report).await;
}

/// A ceiling of zero admits nothing.
async fn zero_admits_nothing(store: &dyn QuotaStore, at: Timestamp, report: &mut Report) {
    report.checked += 1;
    let run = RunId::generate();
    match store.reserve(run, Some(0), at).await {
        Err(QuotaError::TooManyRuns { .. }) => {}
        Err(e) => report.record(
            "a ceiling of zero admits nothing",
            format!("reserving under a zero ceiling failed with `{e}` rather than a refusal"),
        ),
        Ok(()) => {
            report.record(
                "a ceiling of zero admits nothing",
                "a run was admitted under a ceiling of zero — the value an \
                 operator sets to stop a tenant dead, and the one a limit \
                 compared inside its counting loop never sees",
            );
            let _ = store.release(run).await;
        }
    }
}

/// A tenant at its ceiling is refused, and a release makes room.
async fn ceiling_refuses_and_frees(store: &dyn QuotaStore, at: Timestamp, report: &mut Report) {
    let first = RunId::generate();
    let second = RunId::generate();

    report.checked += 1;
    if let Err(e) = store.reserve(first, Some(1), at).await {
        report.record(
            "a run fits under a ceiling of one",
            format!("the first reservation failed: {e}"),
        );
        return;
    }

    report.checked += 1;
    match store.reserve(second, Some(1), at).await {
        Err(QuotaError::TooManyRuns { running, .. }) => {
            report.checked += 1;
            if running == 0 {
                report.record(
                    "a refusal reports how many runs are executing",
                    "the refusal said zero runs are executing, which tells an \
                     operator asking why they are throttled precisely nothing",
                );
            }
        }
        Err(e) => report.record(
            "a tenant at its ceiling is refused",
            format!("failed with `{e}` rather than reporting the ceiling"),
        ),
        Ok(()) => report.record(
            "a tenant at its ceiling is refused",
            "a second run was admitted past a ceiling of one, so the ceiling \
             bounds nothing",
        ),
    }

    report.checked += 1;
    if let Err(e) = store
        .settle(&QuotaSettlement {
            run: first,
            epoch: 1,
            period: None,
            spend: Spend::default(),
            release_slot: true,
        })
        .await
    {
        report.record("settling and releasing a slot", format!("{e}"));
        return;
    }
    report.checked += 1;
    match store.reserve(second, Some(1), at).await {
        Ok(()) => {
            let _ = store.release(second).await;
        }
        Err(e) => report.record(
            "releasing makes room",
            format!(
                "the slot freed by a finished run could not be reused: {e}. A \
                 ceiling is back-pressure, and one that never frees is a tenant \
                 permanently stopped by its first burst"
            ),
        ),
    }
}

/// Reserving one run twice takes one slot.
async fn reserving_twice_takes_one_slot(
    store: &dyn QuotaStore,
    at: Timestamp,
    report: &mut Report,
) {
    let run = RunId::generate();
    report.checked += 1;
    if let Err(e) = store.reserve(run, Some(1), at).await {
        report.record("reserving a run", format!("{e}"));
        return;
    }

    report.checked += 1;
    match store.reserve(run, Some(1), at).await {
        Ok(()) => {}
        Err(e) => report.record(
            "reserving one run twice is idempotent",
            format!(
                "a retried admission was refused against its own slot ({e}), so \
                 a transient error during admission costs the tenant capacity \
                 until something releases a run it never really started"
            ),
        ),
    }

    report.checked += 1;
    match store.running().await {
        Ok(1) => {}
        Ok(n) => report.record(
            "reserving one run twice takes one slot",
            format!(
                "{n} slots are held for one run, so every retry permanently shrinks the ceiling"
            ),
        ),
        Err(e) => report.record("counting running runs", format!("{e}")),
    }

    // A count says *one of one* and a throttled tenant needs the id, because
    // the slot may belong to a run that stopped existing — see
    // `QuotaStore::running_runs`.
    report.checked += 1;
    match store.running_runs(100).await {
        Ok(held) if held == vec![run] => {}
        Ok(held) => report.record(
            "naming the runs that hold slots",
            format!(
                "one run holds a slot and the listing says {held:?} — an operator \
                 told only how many cannot tell a live run from a slot a dead \
                 instance stranded, which is the case the accounting exists for"
            ),
        ),
        Err(e) => report.record("naming the runs that hold slots", format!("{e}")),
    }

    let _ = store.release(run).await;

    // And it empties, which is what makes the listing a queue rather than a
    // record of everything that ever ran.
    report.checked += 1;
    match store.running_runs(100).await {
        Ok(held) if held.is_empty() => {}
        Ok(held) => report.record(
            "releasing a slot takes the run off the listing",
            format!("the released run is still listed as holding a slot: {held:?}"),
        ),
        Err(e) => report.record(
            "releasing a slot takes the run off the listing",
            format!("{e}"),
        ),
    }
}

/// Spend sums within a period and does not cross between them.
async fn spend_accrues_per_period(store: &dyn QuotaStore, report: &mut Report) {
    let (this, next) = ("2999-01", "2999-02");
    let run = RunId::generate();
    let first = QuotaSettlement {
        run,
        epoch: 1,
        period: Some(this.to_owned()),
        spend: Spend::tokens(400),
        release_slot: false,
    };
    let second = QuotaSettlement {
        run,
        epoch: 2,
        period: Some(this.to_owned()),
        spend: Spend::tokens(600),
        release_slot: false,
    };

    report.checked += 1;
    for settlement in [&first, &second] {
        if let Err(e) = store.settle(settlement).await {
            report.record("settling spend", format!("{e}"));
            return;
        }
    }

    report.checked += 1;
    match store.spent(this).await {
        Ok(s) if s.tokens == 1_000 => {}
        Ok(s) => report.record(
            "accruals sum",
            format!(
                "two accruals of 400 and 600 totalled {} rather than 1000. \
                 Reading a total, adding to it and writing it back loses one of \
                 two concurrent updates — and what it loses is spend a tenant \
                 has already incurred, so the ceiling drifts upward under load",
                s.tokens
            ),
        ),
        Err(e) => report.record("reading spend", format!("{e}")),
    }

    // A lost acknowledgement retries the exact same receipt. It must not bill
    // the pass twice, and the positive total above means this cannot pass by
    // ignoring every settlement.
    report.checked += 1;
    if let Err(e) = store.settle(&first).await {
        report.record("retrying an identical settlement", format!("{e}"));
    }
    match store.spent(this).await {
        Ok(s) if s.tokens == 1_000 => {}
        Ok(s) => report.record(
            "an identical settlement accrues once",
            format!("retrying one pass changed the total to {} tokens", s.tokens),
        ),
        Err(e) => report.record("reading spend after a settlement retry", format!("{e}")),
    }

    // A key may not be reused to rewrite accounting. The store must compare
    // the receipt, not treat every conflict as idempotent success.
    report.checked += 1;
    let changed = QuotaSettlement {
        spend: Spend::tokens(401),
        ..first.clone()
    };
    if store.settle(&changed).await.is_ok() {
        report.record(
            "one pass key names one exact settlement",
            "the same run/epoch accepted a different spend, so a retry can rewrite the bill",
        );
    }

    report.checked += 1;
    match store.spent(next).await {
        Ok(s) if s.tokens == 0 => {}
        Ok(s) => report.record(
            "periods are independent",
            format!(
                "an untouched period already reports {} tokens, so a ceiling \
                 would never reset and a tenant is billed forever for one month",
                s.tokens
            ),
        ),
        Err(e) => report.record("reading an untouched period", format!("{e}")),
    }
}

/// The emergency stop, held to the same contract on every backend.
///
/// Four properties, and the last two are the ones an in-process flag and a
/// single overwritable row respectively fail.
#[allow(clippy::too_many_lines)]
async fn halt(store: &dyn QuotaStore, report: &mut Report) {
    use crate::quota::HaltScope;

    let tenant = HaltScope::Tenant;
    let agent = HaltScope::agent("payments-clerk");
    let revision = HaltScope::revision(crate::core::Digest::of(b"a manifest revision"));

    let standing = |halts: &[crate::quota::Halt], scope: &HaltScope| -> Option<String> {
        halts
            .iter()
            .find(|h| &h.scope == scope)
            .map(|h| h.reason.clone())
    };

    report.checked += 1;
    match store.halts().await {
        Ok(halts) if halts.is_empty() => {}
        Ok(halts) => report.record(
            "a fresh tenant is not halted",
            format!(
                "an untouched tenant reports {halts:?}, so a plane would refuse \
                 every run it was never told to refuse"
            ),
        ),
        Err(e) => report.record("reading the halts", format!("{e}")),
    }

    report.checked += 1;
    // Every throw in this battery names somebody: a halt with no operator on
    // it is the state this contract exists to make unreachable.
    let thrower =
        |actor: &str| crate::core::Operator::asserted(actor).expect("a battery names its operator");
    let at = crate::core::Timestamp::from_unix_timestamp(1_700_000_000).expect("a fixed instant");
    if let Err(e) = store
        .set_halt(&tenant, &thrower("ops-alice"), at, "incident 42")
        .await
    {
        report.record("setting the halt", format!("{e}"));
    }
    match store.halts().await {
        Ok(halts) if standing(&halts, &tenant).as_deref() == Some("incident 42") => {}
        Ok(other) => report.record(
            "the halt survives being written",
            format!(
                "after halting, the store reports {other:?} — a switch that does \
                 not read back is one an operator believes they threw"
            ),
        ),
        Err(e) => report.record("reading the halt back", format!("{e}")),
    }

    // The reason is replaced rather than appended to, so the current one is
    // always the current one.
    report.checked += 1;
    if let Err(e) = store
        .set_halt(&tenant, &thrower("ops-alice"), at, "incident 43")
        .await
    {
        report.record("re-halting", format!("{e}"));
    }
    match store.halts().await {
        Ok(halts) if standing(&halts, &tenant).as_deref() == Some("incident 43") => {}
        Ok(other) => report.record(
            "re-halting replaces the reason",
            format!("expected the newer reason, got {other:?}"),
        ),
        Err(e) => report.record("re-reading the halt", format!("{e}")),
    }

    // **Scopes are independent rows.** A narrow halt beside a broad one, and
    // lifting the narrow one, must leave the broad one standing: an incident
    // that widens and then partly resolves is the ordinary shape, and a single
    // overwritable flag gets it wrong in the direction that lets work through.
    report.checked += 1;
    if let Err(e) = store
        .set_halt(&agent, &thrower("ops-bob"), at, "agent 12 is looping")
        .await
    {
        report.record("halting one agent", format!("{e}"));
    }
    if let Err(e) = store
        .set_halt(&revision, &thrower("ops-bob"), at, "bad deploy")
        .await
    {
        report.record("halting one revision", format!("{e}"));
    }
    match store.halts().await {
        Ok(halts)
            if standing(&halts, &tenant).as_deref() == Some("incident 43")
                && standing(&halts, &agent).as_deref() == Some("agent 12 is looping")
                && standing(&halts, &revision).as_deref() == Some("bad deploy") => {}
        Ok(other) => report.record(
            "scopes are independent",
            format!(
                "a narrow halt overwrote a broader one, or was not kept: {other:?} — \
                 an incident that widens must not un-stop what was already stopped"
            ),
        ),
        Err(e) => report.record("reading several standing halts", format!("{e}")),
    }

    report.checked += 1;
    if let Err(e) = store.lift_halt(&agent).await {
        report.record("lifting one scope", format!("{e}"));
    }
    match store.halts().await {
        Ok(halts)
            if standing(&halts, &agent).is_none()
                && standing(&halts, &tenant).as_deref() == Some("incident 43") => {}
        Ok(other) => report.record(
            "lifting one scope leaves the others",
            format!(
                "after lifting the agent halt the store reports {other:?} — lifting \
                 a narrow stop must not lift the broad one it sits under"
            ),
        ),
        Err(e) => report.record("reading a partly lifted halt", format!("{e}")),
    }

    report.checked += 1;
    for scope in [&tenant, &revision] {
        if let Err(e) = store.lift_halt(scope).await {
            report.record("lifting the halt", format!("{e}"));
        }
    }
    match store.halts().await {
        Ok(halts) if halts.is_empty() => {}
        Ok(other) => report.record(
            "a lifted halt stays lifted",
            format!(
                "the tenant is still halted by {other:?} after the stop was \
                 lifted, so an incident that is over never ends"
            ),
        ),
        Err(e) => report.record("reading a lifted halt", format!("{e}")),
    }

    // Lifting a halt nobody set is a no-op, not an error: an operator clearing
    // a switch they are not sure about must not be punished for it — and the
    // answer still has to say that nothing was standing, because during an
    // incident *I cleared it* and *I cleared the wrong scope* are different
    // facts and only one of them is good news.
    report.checked += 1;
    match store.lift_halt(&tenant).await {
        Ok(false) => {}
        Ok(true) => report.record(
            "lifting an unset halt",
            "the store reported that a halt was standing when none was".to_owned(),
        ),
        Err(e) => report.record("lifting an unset halt", format!("{e}")),
    }

    // **Who threw it survives the round trip, and so does what established the
    // name.** The runtime cannot check an emergency stop, so the operator on
    // the row is the whole of its evidence — a store that keeps the reason and
    // drops the name leaves a switch nobody can be asked about.
    report.checked += 1;
    let by = crate::core::Operator::authenticated("ops-carol").expect("a name");
    if let Err(e) = store.set_halt(&tenant, &by, at, "incident 44").await {
        report.record("halting with an authenticated operator", format!("{e}"));
    }
    match store.halts().await {
        Ok(halts) => match halts.iter().find(|h| h.scope == tenant) {
            Some(h) if h.by == by && h.at == at => {}
            Some(h) => report.record(
                "a halt keeps who threw it",
                format!(
                    "the store read back {:?} at {:?} rather than {by:?} at {at:?} — an \
                     emergency stop nobody is named on cannot be asked about afterwards",
                    h.by, h.at
                ),
            ),
            None => report.record(
                "a halt keeps who threw it",
                "the halt did not read back at all".to_owned(),
            ),
        },
        Err(e) => report.record("reading an attributed halt", format!("{e}")),
    }
    if let Err(e) = store.lift_halt(&tenant).await {
        report.record("clearing the attributed halt", format!("{e}"));
    }
}