$ cat /tmp/claude-1000/-tmp-flows-fleet-501/144d3b43-0019-4de3-988a-7cd9ba4fc148/scratchpad/pr441-probe/src/main.rs
//! Read-only probe: an overflow fence torn between `subscription.overflow_fenced`
//! and `subscription.closed`. Uses the doc-hidden single-append fence to produce
//! the torn state, then asks what deliver / inspect / fence-retry / next do.
use relayflowd::Engine;
use relayflowd::engine::{SubscriptionNext, SubscriptionWake};
use relayflowd_core::{EntryType, RunSpec, SimClock};
use serde_json::json;

fn count(engine: &Engine<SimClock>, id: &str, t: EntryType) -> usize {
    engine.journal_entries(id, 1, 500).unwrap().iter().filter(|e| e.entry_type == t).count()
}

fn main() {
    let dir = tempfile::tempdir().unwrap();
    let engine = Engine::with_clock(dir.path(), SimClock::new(100));
    let spec = RunSpec::parse(&json!({"steps":[{"id":"body","type":"llm","prompt":"body"}]})).unwrap();
    let id = engine.start(spec, "probe", Some(0)).unwrap().run_id;
    let receipt = json!({"generation":7});
    engine.open_subscription(&id, "one", vec!["github".into()], None, 0, 100, 1000, false).unwrap();
    engine.activate_subscription(&id, "one", 0, receipt.clone()).unwrap();
    // The body parks on its first next() (no frames yet): a durable wait.event exists.
    let first = engine.next_subscription_outcome(&id, "one", None).unwrap().0;
    println!("first next() -> {:?}", matches!(first, SubscriptionNext::Suspended { .. }));
    assert!(engine.deliver_subscription_frame(&id, "one", &receipt, "event-1", json!({"type":"github","payload":{"n":1}})).unwrap());

    // TORN: only the fence lands (process died before the close command).
    engine.fence_subscription_overflow(&id, "one").unwrap();
    println!("after torn fence: fenced={} closed={}", count(&engine, &id, EntryType::SubscriptionOverflowFenced), count(&engine, &id, EntryType::SubscriptionClosed));

    // 1. deliver on the torn state
    let err = engine.deliver_subscription_frame(&id, "one", &receipt, "event-2", json!({"type":"github","payload":{"n":2}})).unwrap_err();
    println!("deliver on torn fence -> Err({err})");
    // 2. inspect on the torn state
    let snap = engine.inspect_subscriptions(&id).unwrap();
    println!("inspect on torn fence -> state={} completionReason={}", snap[0]["state"], snap[0]["completionReason"]);

    // 3. restart, then Cloud retries the fence: must complete the close once.
    drop(engine);
    let restored = Engine::with_clock(dir.path(), SimClock::new(120));
    restored.fence_router_subscription_overflow(&id, "one", &receipt).unwrap();
    println!("after fence retry: fenced={} closed={}", count(&restored, &id, EntryType::SubscriptionOverflowFenced), count(&restored, &id, EntryType::SubscriptionClosed));
    let snap = restored.inspect_subscriptions(&id).unwrap();
    println!("inspect after retry -> state={} completionReason={}", snap[0]["state"], snap[0]["completionReason"]);
    let (wake, _) = restored.next_subscription_outcome(&id, "one", None).unwrap();
    println!("next() after retry -> {:?}", wake);
    assert!(matches!(wake, SubscriptionNext::Wake(SubscriptionWake::Overflow { retained: 1, .. })));

    // 4. Alternative path: torn fence, then resume (no Cloud retry) completes it too.
    let dir2 = tempfile::tempdir().unwrap();
    let e2 = Engine::with_clock(dir2.path(), SimClock::new(100));
    let spec = RunSpec::parse(&json!({"steps":[{"id":"body","type":"llm","prompt":"body"}]})).unwrap();
    let id2 = e2.start(spec, "probe", Some(0)).unwrap().run_id;
    e2.open_subscription(&id2, "one", vec!["github".into()], None, 0, 100, 1000, false).unwrap();
    e2.activate_subscription(&id2, "one", 0, receipt.clone()).unwrap();
    e2.next_subscription_outcome(&id2, "one", None).unwrap();
    e2.fence_subscription_overflow(&id2, "one").unwrap();
    drop(e2);
    let e2 = Engine::with_clock(dir2.path(), SimClock::new(130));
    let outcome = e2.resume(&id2, None).unwrap();
    println!("resume on torn fence -> status={:?} closed={} wait.completed={}", outcome.status, count(&e2, &id2, EntryType::SubscriptionClosed), count(&e2, &id2, EntryType::WaitCompleted));
    println!("PROBE_OK");
}
$ cd /tmp/claude-1000/-tmp-flows-fleet-501/144d3b43-0019-4de3-988a-7cd9ba4fc148/scratchpad/pr441-probe && cargo run --quiet   # relayflowd path dep at /tmp/flows-pr-followup/pr441/kernel (kernel identical at 62e78e19 and a217c3f6)
first next() -> true
after torn fence: fenced=1 closed=0
deliver on torn fence -> Err(subscription_closed)
inspect on torn fence -> state="active" completionReason=null
after fence retry: fenced=1 closed=1
inspect after retry -> state="closed" completionReason="overflow"
next() after retry -> Wake(Overflow { retained: 1, bytes: 35, from: 0 })
resume on torn fence -> status=Parked closed=1 wait.completed=1
PROBE_OK
exit_code=0
