aion-server 0.21.0

Aion workflow server library: HTTP, gRPC, WebSocket, and worker endpoints. Run it with the `aion` binary from the aion-cli crate.
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
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
//! Tests for the NOI-5 transcript sequencer + fan-out, including the two
//! mandatory negative controls (concurrent-writer monotonicity, failover
//! dedup) and the transcript retention bounds (per-event truncation, the
//! per-stream cap marker, past-cap live-only delivery).

use std::num::NonZeroUsize;

use aion_core::{ActivityEventKind, ActivityId, MessageRole, RunId, WorkflowId};
use aion_store::InMemoryObservabilityStore;
use chrono::Utc;
use futures::StreamExt;
use uuid::Uuid;

use super::*;

fn capacity(value: usize) -> Result<NonZeroUsize, Box<dyn std::error::Error>> {
    NonZeroUsize::new(value).ok_or_else(|| "capacity must be non-zero".into())
}

/// These tests pin the SINGLE-EVENT sequencing contract (monotonicity, failover
/// dedup, retention bounds), so they run the sequencer under the identity flush
/// policy — one event per durable commit, no hold — which is the pre-batching
/// behaviour. Coalescing has its own tests in `activity_publisher_batching_tests`.
const UNBATCHED: TranscriptBatchPolicy = TranscriptBatchPolicy {
    max_batch_events: NonZeroUsize::MIN,
    max_hold: std::time::Duration::ZERO,
};

fn publisher(cap: usize) -> Result<ActivityEventPublisher, Box<dyn std::error::Error>> {
    let store = Arc::new(InMemoryObservabilityStore::default());
    Ok(ActivityEventPublisher::new(
        store,
        capacity(cap)?,
        UNBATCHED,
    ))
}

/// The first generation of the continue-as-new chain the fixtures live in.
fn generation_one() -> RunId {
    RunId::new(Uuid::from_u128(0x11))
}

/// The second generation: SAME workflow id, new run.
fn generation_two() -> RunId {
    RunId::new(Uuid::from_u128(0x22))
}

fn event(attempt: u32, worker_seq: u64, ephemeral: bool, text: &str) -> ActivityEvent {
    ActivityEvent {
        workflow_id: WorkflowId::new(Uuid::from_u128(1)),
        run_id: generation_one(),
        activity_id: ActivityId::from_sequence_position(3),
        attempt,
        agent_id: Uuid::from_u128(9),
        agent_role: "orchestrator".to_owned(),
        emitted_at: Utc::now(),
        worker_seq,
        store_seq: None,
        ephemeral,
        kind: if ephemeral {
            ActivityEventKind::Delta {
                message_id: "m1".to_owned(),
                text_fragment: text.to_owned(),
            }
        } else {
            ActivityEventKind::Message {
                role: MessageRole::Assistant,
                text: text.to_owned(),
            }
        },
    }
}

fn key(attempt: u32) -> ActivityStreamKey {
    ActivityStreamKey::new(
        WorkflowId::new(Uuid::from_u128(1)),
        generation_one(),
        ActivityId::from_sequence_position(3),
        attempt,
    )
}

#[tokio::test]
async fn publish_assigns_commit_allocated_monotonic_store_seq()
-> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(16)?;
    assert_eq!(publisher.publish(&event(0, 1, false, "a")).await?, Some(0));
    assert_eq!(publisher.publish(&event(0, 2, false, "b")).await?, Some(1));
    assert_eq!(publisher.publish(&event(0, 3, false, "c")).await?, Some(2));
    let tail = publisher.replay_from(&key(0), 0).await?;
    assert_eq!(
        tail.iter().map(|r| r.store_seq).collect::<Vec<_>>(),
        vec![0, 1, 2]
    );
    Ok(())
}

#[tokio::test]
async fn ephemeral_events_are_never_persisted() -> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(16)?;
    // A live subscriber (fresh, no resume cursor) must still SEE the ephemeral
    // delta AND the persisted message at store_seq 0.
    let mut live = publisher.subscribe(key(0), None);
    assert_eq!(publisher.publish(&event(0, 1, true, "wor")).await?, None);
    assert_eq!(
        publisher.publish(&event(0, 2, false, "word")).await?,
        Some(0)
    );
    // Durable tail has ONLY the non-ephemeral message.
    let tail = publisher.replay_from(&key(0), 0).await?;
    assert_eq!(tail.len(), 1);
    assert!(matches!(
        tail[0].event.kind,
        ActivityEventKind::Message { .. }
    ));

    // Live stream delivered the ephemeral delta first (store_seq None), then
    // the persisted message (store_seq Some(0)).
    let first = live.next().await.ok_or("missing ephemeral")??;
    assert!(first.ephemeral);
    assert_eq!(first.store_seq, None);
    let second = live.next().await.ok_or("missing message")??;
    assert!(!second.ephemeral);
    assert_eq!(second.store_seq, Some(0));
    Ok(())
}

/// MANDATORY NEGATIVE CONTROL (a): concurrent writers on ONE
/// `(wf,act,attempt)` stream. Many `publish` calls race the same head; the
/// read-head -> append -> on-conflict-retry loop must serialize them so the
/// final durable sequence is strictly monotonic 0..N-1 with NO gap and NO
/// duplicate — proving the retry loop (not the store) enforces monotonicity.
#[tokio::test(flavor = "multi_thread")]
async fn concurrent_writers_produce_gapless_monotonic_store_seq()
-> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(256)?;
    let writers = 32u64;
    let mut handles = Vec::new();
    for worker_seq in 0..writers {
        let publisher = publisher.clone();
        handles.push(tokio::spawn(async move {
            publisher
                .publish(&event(0, worker_seq, false, "concurrent"))
                .await
        }));
    }
    let mut assigned = Vec::new();
    for handle in handles {
        if let Some(store_seq) = handle.await?? {
            assigned.push(store_seq);
        }
    }
    assigned.sort_unstable();
    // Every writer won exactly one distinct, contiguous store_seq.
    assert_eq!(
        assigned,
        (0..writers).collect::<Vec<_>>(),
        "concurrent writers must produce a gapless, duplicate-free monotonic sequence"
    );
    // The durable tail agrees: exactly `writers` records, seqs 0..N-1 in order.
    let tail = publisher.replay_from(&key(0), 0).await?;
    assert_eq!(
        tail.iter().map(|r| r.store_seq).collect::<Vec<_>>(),
        (0..writers).collect::<Vec<_>>()
    );
    Ok(())
}

/// MANDATORY NEGATIVE CONTROL (b): failover dedup. A dying worker and an
/// adopting worker BOTH emit for one `(wf,act,attempt)` (the same session
/// resumes across failover). Modeled as two publishers sharing ONE store
/// (the server owns durability, §5.3) racing appends: the commit-allocated
/// `store_seq` collapses both emitters into one monotonic stream with no
/// gap/dup — the transcript is not duplicated or reordered by the double-emit.
#[tokio::test(flavor = "multi_thread")]
async fn failover_double_emit_dedupes_to_one_monotonic_stream()
-> Result<(), Box<dyn std::error::Error>> {
    // ONE durable store (the server), TWO publisher clones = the dying +
    // adopting worker's events both flowing through the single server writer.
    let store = Arc::new(InMemoryObservabilityStore::default());
    let dying = ActivityEventPublisher::new(store.clone(), capacity(256)?, UNBATCHED);
    let adopting = ActivityEventPublisher::new(store, capacity(256)?, UNBATCHED);

    let mut handles = Vec::new();
    for worker_seq in 0..16u64 {
        let dying = dying.clone();
        let adopting = adopting.clone();
        handles.push(tokio::spawn(async move {
            // Both survivors emit the SAME logical event for this attempt.
            let a = dying.publish(&event(0, worker_seq, false, "dying")).await;
            let b = adopting
                .publish(&event(0, worker_seq, false, "adopting"))
                .await;
            (a, b)
        }));
    }
    let mut assigned = Vec::new();
    for handle in handles {
        let (a, b) = handle.await?;
        if let Some(seq) = a? {
            assigned.push(seq);
        }
        if let Some(seq) = b? {
            assigned.push(seq);
        }
    }
    assigned.sort_unstable();
    // 16 events x 2 emitters = 32 durable records, each with a DISTINCT,
    // contiguous store_seq — the double-emit does not corrupt monotonicity.
    assert_eq!(
        assigned,
        (0..32).collect::<Vec<_>>(),
        "failover double-emit must land a gapless monotonic store_seq stream"
    );
    let tail = dying.replay_from(&key(0), 0).await?;
    let seqs: Vec<u64> = tail.iter().map(|r| r.store_seq).collect();
    assert_eq!(seqs, (0..32).collect::<Vec<_>>());
    Ok(())
}

/// MANDATORY NEGATIVE CONTROL: the WRONG-allocator case. A process-local
/// `AtomicU64` (the `ClusterEventPublisher` pattern the design forbids for
/// `store_seq`) is shown to produce COLLIDING sequences when two "survivors"
/// each start their own counter after a failover — proving the counter
/// belongs in the commit, not the process. The commit-allocated publisher
/// (above) does NOT exhibit this; this test pins the failure mode we avoid.
#[tokio::test]
async fn process_local_atomic_counter_collides_across_survivors()
-> Result<(), Box<dyn std::error::Error>> {
    use std::sync::atomic::{AtomicU64, Ordering};
    // Two survivors, each with its OWN fresh process-local counter (exactly
    // what a per-process AtomicU64 does after restart/failover).
    let dying = AtomicU64::new(0);
    let adopting = AtomicU64::new(0);
    let dying_seq = dying.fetch_add(1, Ordering::SeqCst);
    let adopting_seq = adopting.fetch_add(1, Ordering::SeqCst);
    // Both minted the SAME sequence — a collision. This is the bug the
    // commit-allocated design exists to prevent.
    assert_eq!(
        dying_seq, adopting_seq,
        "process-local counters collide across survivors (the forbidden pattern)"
    );

    // By contrast the commit-allocated publisher over a shared store never
    // collides: the two survivors get DISTINCT sequences.
    let store = Arc::new(InMemoryObservabilityStore::default());
    let a = ActivityEventPublisher::new(store.clone(), capacity(8)?, UNBATCHED);
    let b = ActivityEventPublisher::new(store, capacity(8)?, UNBATCHED);
    let sa = a.publish(&event(0, 1, false, "a")).await?;
    let sb = b.publish(&event(0, 2, false, "b")).await?;
    assert_ne!(
        sa, sb,
        "commit-allocated store_seq must be distinct across survivors"
    );
    assert_eq!((sa, sb), (Some(0), Some(1)));
    Ok(())
}

/// The live-stream + resume-by-`store_seq` gate: a subscriber tails live
/// events, and a reconnecting client resumes from the durable tail by
/// `store_seq` then splices onto the live broadcast with NO gap and NO
/// duplicate at the seam.
#[tokio::test]
async fn live_stream_then_resume_by_store_seq_has_no_gap() -> Result<(), Box<dyn std::error::Error>>
{
    let publisher = publisher(64)?;
    // Persist a few events (as if the client was connected then dropped).
    for seq in 0..3u64 {
        publisher.publish(&event(0, seq, false, "early")).await?;
    }
    // Client reconnects having last seen store_seq 1. Splice: attach the live
    // stream BEFORE reading the durable tail (gap-free splice), with the
    // resume cursor as after_seq so live re-delivery of <=cursor is suppressed.
    let mut live = publisher.subscribe(key(0), Some(1));
    let replay = publisher.replay_from(&key(0), 2).await?;
    // The durable tail from the cursor is exactly seq 2 (the one it missed).
    assert_eq!(
        replay.iter().map(|r| r.store_seq).collect::<Vec<_>>(),
        vec![2]
    );
    // Now new live events arrive.
    publisher.publish(&event(0, 10, false, "live-3")).await?;
    publisher.publish(&event(0, 11, false, "live-4")).await?;
    // The live splice yields seq 3 then 4 — no gap, no re-delivery of <=1.
    let first = live.next().await.ok_or("missing live-3")??;
    assert_eq!(first.store_seq, Some(3));
    let second = live.next().await.ok_or("missing live-4")??;
    assert_eq!(second.store_seq, Some(4));
    Ok(())
}

#[tokio::test]
async fn subscribe_filters_out_other_attempt_streams() -> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(64)?;
    let mut live = publisher.subscribe(key(0), None);
    // An event for a DIFFERENT attempt must not reach this subscriber.
    publisher
        .publish(&event(1, 1, false, "other-attempt"))
        .await?;
    publisher.publish(&event(0, 1, false, "mine")).await?;
    let received = live.next().await.ok_or("missing my event")??;
    assert_eq!(received.attempt, 0);
    assert!(matches!(
        received.kind,
        ActivityEventKind::Message { text, .. } if text == "mine"
    ));
    Ok(())
}

#[tokio::test]
async fn lagged_subscriber_yields_typed_skip_count() -> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(2)?;
    let mut live = publisher.subscribe(key(0), None);
    // Overflow the capacity-2 live buffer without consuming.
    for seq in 0..5u64 {
        publisher.publish(&event(0, seq, false, "flood")).await?;
    }
    let lagged = live.next().await.ok_or("missing lag item")?;
    assert!(matches!(lagged, Err(TranscriptStreamLagged { .. })));
    Ok(())
}

// --- transcript retention bounds -------------------------------------------

fn bounded_publisher(
    cap: usize,
    bounds: TranscriptBounds,
) -> Result<ActivityEventPublisher, Box<dyn std::error::Error>> {
    Ok(publisher(cap)?.with_bounds(bounds))
}

/// The per-stream cap: with `max_stream_events = 3`, publishes 0..2 persist
/// normally, the 4th append persists ONE marker record (a `Progress`/`Note`
/// naming the cap) and every publish past that persists nothing.
#[tokio::test]
async fn stream_cap_appends_one_marker_then_stops_persisting()
-> Result<(), Box<dyn std::error::Error>> {
    let publisher = bounded_publisher(
        64,
        TranscriptBounds {
            max_event_bytes: 256 * 1024,
            max_stream_events: 3,
        },
    )?;
    let mut assigned = Vec::new();
    for worker_seq in 0..6u64 {
        assigned.push(
            publisher
                .publish(&event(0, worker_seq, false, "chatty"))
                .await?,
        );
    }
    assert_eq!(
        assigned,
        vec![Some(0), Some(1), Some(2), None, None, None],
        "publishes past the cap return Ok(None)"
    );
    let tail = publisher.replay_from(&key(0), 0).await?;
    assert_eq!(
        tail.iter().map(|r| r.store_seq).collect::<Vec<_>>(),
        vec![0, 1, 2, 3],
        "exactly the capped records plus the one marker"
    );
    let ActivityEventKind::Progress {
        detail: ProgressDetail::Note { text },
    } = &tail[3].event.kind
    else {
        return Err("record 3 must be the retention-cap marker note".into());
    };
    assert!(
        text.contains("retention cap"),
        "the marker names the cap: {text}"
    );
    assert!(text.contains("3 events"), "the marker names the value");
    Ok(())
}

/// Past the cap the stream stays LIVE: a subscriber still receives every
/// event, with `store_seq: None` and `ephemeral == false` (real transcript,
/// just not retained).
#[tokio::test]
async fn capped_stream_still_fans_out_live_without_store_seq()
-> Result<(), Box<dyn std::error::Error>> {
    let publisher = bounded_publisher(
        64,
        TranscriptBounds {
            max_event_bytes: 256 * 1024,
            max_stream_events: 1,
        },
    )?;
    let mut live = publisher.subscribe(key(0), None);
    assert_eq!(
        publisher.publish(&event(0, 1, false, "kept")).await?,
        Some(0)
    );
    // Crosses the cap: persists the marker, fans the event out live-only.
    assert_eq!(publisher.publish(&event(0, 2, false, "over")).await?, None);
    // Fully past the cap: live-only.
    assert_eq!(
        publisher.publish(&event(0, 3, false, "way-over")).await?,
        None
    );

    let kept = live.next().await.ok_or("missing kept")??;
    assert_eq!(kept.store_seq, Some(0));
    let marker = live.next().await.ok_or("missing marker")??;
    assert_eq!(
        marker.store_seq,
        Some(1),
        "the marker carries its store_seq"
    );
    assert!(matches!(marker.kind, ActivityEventKind::Progress { .. }));
    let over = live.next().await.ok_or("missing over")??;
    assert_eq!(over.store_seq, None, "past-cap events carry no store_seq");
    assert!(!over.ephemeral, "past-cap events are NOT ephemeral");
    assert!(matches!(
        over.kind,
        ActivityEventKind::Message { text, .. } if text == "over"
    ));
    let way_over = live.next().await.ok_or("missing way-over")??;
    assert_eq!(way_over.store_seq, None);
    assert!(!way_over.ephemeral);
    Ok(())
}

/// The per-event ceiling: an oversized message is truncated BEFORE the durable
/// append, so the retained record is bounded and marked.
#[tokio::test]
async fn oversized_event_is_truncated_before_persist() -> Result<(), Box<dyn std::error::Error>> {
    let publisher = bounded_publisher(
        16,
        TranscriptBounds {
            max_event_bytes: 512,
            max_stream_events: 20_000,
        },
    )?;
    let huge = "x".repeat(10_000);
    assert_eq!(
        publisher.publish(&event(0, 1, false, &huge)).await?,
        Some(0)
    );
    let tail = publisher.replay_from(&key(0), 0).await?;
    assert_eq!(tail.len(), 1);
    let ActivityEventKind::Message { text, .. } = &tail[0].event.kind else {
        return Err("expected the truncated message".into());
    };
    assert!(
        text.ends_with("bytes by observability.max_event_bytes]"),
        "the retained text ends with the truncation marker: {text}"
    );
    assert!(
        serde_json::to_vec(&tail[0].event)?.len() <= 1024,
        "the retained record is bounded (512 + marker slack)"
    );
    assert!(
        !text.contains(&huge),
        "the original oversized text is not retained in full"
    );
    Ok(())
}

/// THE RUN-SCOPING INVARIANT at the sequencer, on the coordinates a
/// continue-as-new chain actually produces: ordinal `0`, attempt `1`, in two
/// generations of one workflow.
///
/// Both events must be sequenced as the FIRST record of their own stream
/// (`store_seq == 0` each), and a read scoped to generation two must return
/// exactly one event. Under the pre-run-axis key the second publish would have
/// seen generation one's head, retried, and landed at `store_seq == 1` on
/// generation one's stream — one fused transcript, silently.
#[tokio::test]
async fn two_generations_of_one_chain_are_sequenced_independently()
-> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(16)?;
    let ordinal_zero = ActivityId::from_sequence_position(0);
    let first = ActivityEvent {
        activity_id: ordinal_zero.clone(),
        attempt: 1,
        ..event(1, 1, false, "generation one")
    };
    let second = ActivityEvent {
        run_id: generation_two(),
        ..first.clone()
    };

    assert_eq!(publisher.publish(&first).await?, Some(0));
    assert_eq!(
        publisher.publish(&second).await?,
        Some(0),
        "generation two's first event is its own stream head, not generation one's second record"
    );

    let first_key = ActivityStreamKey::of(&first);
    let second_key = ActivityStreamKey::of(&second);
    assert_ne!(first_key, second_key);
    // The keys agree on every axis except the run — the collision under test.
    assert_eq!(first_key.workflow_id, second_key.workflow_id);
    assert_eq!(first_key.activity_id, second_key.activity_id);
    assert_eq!(first_key.attempt, second_key.attempt);

    let second_generation = publisher.replay_from(&second_key, 0).await?;
    assert_eq!(
        second_generation.len(),
        1,
        "a read scoped to generation two must return exactly its own event"
    );
    assert_eq!(second_generation[0].event.run_id, generation_two());
    let first_generation = publisher.replay_from(&first_key, 0).await?;
    assert_eq!(first_generation.len(), 1);
    assert_eq!(first_generation[0].event.run_id, generation_one());
    Ok(())
}

/// The live-tail filter does not leak across generations: a subscriber attached
/// to generation two's stream sees generation two's event and NOT generation
/// one's, even though the two share every other axis and ride one broadcast.
#[tokio::test]
async fn live_tail_does_not_leak_across_generations() -> Result<(), Box<dyn std::error::Error>> {
    let publisher = publisher(16)?;
    let ordinal_zero = ActivityId::from_sequence_position(0);
    let first = ActivityEvent {
        activity_id: ordinal_zero.clone(),
        attempt: 1,
        ..event(1, 1, false, "generation one")
    };
    let second = ActivityEvent {
        run_id: generation_two(),
        ..ActivityEvent {
            worker_seq: 2,
            ..first.clone()
        }
    };

    let mut tail = publisher.subscribe(ActivityStreamKey::of(&second), None);
    // Generation one publishes FIRST and must be suppressed for this subscriber;
    // generation two publishes second and must arrive.
    publisher.publish(&first).await?;
    publisher.publish(&second).await?;

    let delivered = tail.next().await.ok_or("the live tail closed early")??;
    assert_eq!(
        delivered.run_id,
        generation_two(),
        "the sibling generation's event must not reach a run-scoped subscriber"
    );
    assert_eq!(delivered.worker_seq, 2);
    Ok(())
}

/// The retention cap counts per `(run, activity, attempt)` stream, not per fused
/// chain: a second generation reaching the same coordinates gets its OWN budget,
/// so generation two is not silently capped by what generation one wrote.
#[tokio::test]
async fn the_retention_cap_is_measured_per_run_not_per_chain()
-> Result<(), Box<dyn std::error::Error>> {
    let publisher = bounded_publisher(
        64,
        TranscriptBounds {
            max_event_bytes: 256 * 1024,
            max_stream_events: 2,
        },
    )?;
    let ordinal_zero = ActivityId::from_sequence_position(0);
    let generation = |run_id: RunId, worker_seq: u64| ActivityEvent {
        run_id,
        activity_id: ordinal_zero.clone(),
        attempt: 1,
        ..event(1, worker_seq, false, "x")
    };

    // Generation one fills its budget: seq 0, seq 1, then the cap marker at 2.
    assert_eq!(
        publisher.publish(&generation(generation_one(), 1)).await?,
        Some(0)
    );
    assert_eq!(
        publisher.publish(&generation(generation_one(), 2)).await?,
        Some(1)
    );
    assert_eq!(
        publisher.publish(&generation(generation_one(), 3)).await?,
        None,
        "the third event crosses the cap and is live-only"
    );

    // Generation two starts at zero against its OWN cap, unaffected by the
    // exhausted sibling stream.
    assert_eq!(
        publisher.publish(&generation(generation_two(), 1)).await?,
        Some(0),
        "a new generation must not inherit its predecessor's exhausted retention budget"
    );
    assert_eq!(
        publisher.publish(&generation(generation_two(), 2)).await?,
        Some(1)
    );
    Ok(())
}