From fb7cfe385111dd4d91ae9f81024450ed728ca111 Mon Sep 17 00:00:00 2001 From: "glm-5.3-flash" Date: Mon, 28 Sep 2026 10:06:00 +0000 Subject: [PATCH] feat(fuzz): session_opseq target - stateful op-sequence campaign against the session pump MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Target 4 from the fuzzing plan (the alkcall-bug-class hunt): a #[derive(Arbitrary)] op sequence drives the producer session pump (drive_session_pre_negotiated) over a duplex pair against a fuzz-local mock backend. - fuzz/shared/src/session_opseq.rs: op vocabulary (client writes incl. invalid-type/oversize-length barrier probes, backend production/EOF, exit resolve/fail/drop, control messages, bounded reads, yields), kill-guard exit future (ADR-005 observability), delivery model, and the invariants: kill-on-Drop (kill_fired == !exit_resolved), exit chunk is last with first-issued code -1 on wait-failure (ADR-004), lossless per-stream FIFO + single stdout sentinel, stdin prefix-losslessness (post-shutdown chunks never delivered), control dispatch bounds, take-once allocation, teardown termination. - Harness-model correction caught by a 60s smoke run before the campaign (the alkcall §7.7 pattern, third time): a failed client write can mean the session completed and dropped the server duplex half while the write was in flight - the harness now drains and requires the exit chunk on any write/shutdown error (premature close without it is a finding). No adapter code changed. - 12 committed seeds via the deterministic generator, hand-encoded against the arbitrary 1.4.x derive layout; the encoding is pinned by a seed-decode test. Existing chunk_frame/negotiation_frame/ control_json seed corpora are byte-identical to before. - fuzz/.gitignore: the per-target seed negations never matched (corpus/* excludes the parent dir; git cannot re-include beneath an excluded dir). Fixed with the !corpus/*/ + corpus/*/* + !corpus/*/seed-* recipe - new seed files were silently ignored until now. - Campaign: 10 min detached via run-detached.sh, exited 0, artifact dir empty - 42.3k execs across 33 fork jobs, 0 crash/oom/timeout, cov 3934 -> 4075 edges. Grown corpus excised per the corpus policy. Verification: corpus replay green on stable (9 tests), cargo test 113 + --all-features 156, clippy stable/wasm/fuzz-shared, fmt, doc, publish dry-run, wasm check - all pass. --- docs/plans/fuzzing.md | 64 +- fuzz/.gitignore | 8 +- fuzz/Cargo.lock | 18 + fuzz/Cargo.toml | 7 + fuzz/README.md | 15 +- fuzz/corpus/session_opseq/seed-000 | Bin 0 -> 15 bytes fuzz/corpus/session_opseq/seed-001 | Bin 0 -> 21 bytes fuzz/corpus/session_opseq/seed-002 | Bin 0 -> 57 bytes fuzz/corpus/session_opseq/seed-003 | Bin 0 -> 41 bytes fuzz/corpus/session_opseq/seed-004 | Bin 0 -> 22 bytes fuzz/corpus/session_opseq/seed-005 | Bin 0 -> 40 bytes fuzz/corpus/session_opseq/seed-006 | Bin 0 -> 105 bytes fuzz/corpus/session_opseq/seed-007 | Bin 0 -> 50 bytes fuzz/corpus/session_opseq/seed-008 | Bin 0 -> 47 bytes fuzz/corpus/session_opseq/seed-009 | Bin 0 -> 22 bytes fuzz/corpus/session_opseq/seed-010 | Bin 0 -> 35 bytes fuzz/corpus/session_opseq/seed-011 | Bin 0 -> 48 bytes fuzz/fuzz_targets/session_opseq.rs | 9 + fuzz/gen_fuzz_seeds.py | 267 +++++++ fuzz/shared/Cargo.toml | 9 +- fuzz/shared/src/lib.rs | 1 + fuzz/shared/src/session_opseq.rs | 1162 ++++++++++++++++++++++++++++ 22 files changed, 1549 insertions(+), 11 deletions(-) create mode 100644 fuzz/corpus/session_opseq/seed-000 create mode 100644 fuzz/corpus/session_opseq/seed-001 create mode 100644 fuzz/corpus/session_opseq/seed-002 create mode 100644 fuzz/corpus/session_opseq/seed-003 create mode 100644 fuzz/corpus/session_opseq/seed-004 create mode 100644 fuzz/corpus/session_opseq/seed-005 create mode 100644 fuzz/corpus/session_opseq/seed-006 create mode 100644 fuzz/corpus/session_opseq/seed-007 create mode 100644 fuzz/corpus/session_opseq/seed-008 create mode 100644 fuzz/corpus/session_opseq/seed-009 create mode 100644 fuzz/corpus/session_opseq/seed-010 create mode 100644 fuzz/corpus/session_opseq/seed-011 create mode 100644 fuzz/fuzz_targets/session_opseq.rs create mode 100644 fuzz/shared/src/session_opseq.rs diff --git a/docs/plans/fuzzing.md b/docs/plans/fuzzing.md index 66acbcc..50f8132 100644 --- a/docs/plans/fuzzing.md +++ b/docs/plans/fuzzing.md @@ -144,7 +144,7 @@ alkcall §7.4). | 1 | `chunk_frame` | raw `&[u8]` | `ChunkReader`/`ChunkWriter` over `Cursor` + `block_on` | no-panic; error-shape partition (`ConnectionClosed` only on truncation; `InvalidStreamType` iff `stream_type > 4`; `ChunkTooLarge` iff `length > MAX_CHUNK_LEN`; `Io` impossible on a Cursor); Ok iff header complete + payload complete; round-trip (`write_chunk` then `read_chunk` reproduces stream_type + bytes); 5-byte consumption accounting; peek interleave (`peek_stream_type` → `read_chunk_after_peek` == `read_chunk`, idempotent peek, peek state resets per chunk); never allocates for out-of-range headers | | 2 | `negotiation_frame` | raw `&[u8]` | `NegotiationReader`/`NegotiationWriter` + `NegotiateRequest` serde parse + `error_response_bytes` | no-panic; error-shape partition (`ConnectionClosed` only on truncation; `FrameTooLarge` iff `length > MAX_CHUNK_LEN`; `Json` only when the full body was present; `Io` impossible); consumption accounting (4 + length); negotiation JSON round-trip (`NegotiateRequest` → `to_json` → frame → `from_slice` → equal struct); **cross-codec disambiguation** (see §4.1) | | 3 | `control_json` | raw `&[u8]` | `ControlMessage::from_slice` / `to_json`, `signal_from_name` | no-panic on any bytes; clean serde errors (never silent); `to_json` → `from_slice` structural round-trip; `signal_from_name` is total (returns `Option`, never panics) | -| 4 | `session_opseq` (second wave) | `#[derive(Arbitrary)]` op sequence | the producer session pump + mock backend (`src/testing.rs`, `drive_session`) | no-panic; take-once slot semantics; duplicate/replayed session ops leave the live session intact — the alkcall-bug-class target | +| 4 | `session_opseq` | `#[derive(Arbitrary)]` op sequence | the producer session pump (`drive_session_pre_negotiated`) + a fuzz-local mock backend (`fuzz/shared/src/session_opseq.rs`; the in-crate `MockBackend` is `#[cfg(test)]`-only) | no-panic (driver, pump tasks via the session `JoinError`, harness model); **kill-on-Drop (ADR-005)** — after teardown `kill_fired == !exit_resolved`, and with no exit op issued the kill MUST have fired; **exit chunk is last (ADR-004)** — at most one ctrl_out chunk, `type == "exit"`, `code` matches the first-issued exit op (`-1` on wait-failure or dropped exit sender), nothing observed after it, server never writes stream types 0/3; lossless per-stream FIFO of produced `(len, byte)` patterns with exactly one stdout sentinel (only after the stdout-EOF op) and no stderr sentinel; stdin prefix-losslessness (post-shutdown chunks never delivered); control dispatch never over-counts and dispatches promptly pre-exit; take-once allocation; teardown terminates within the bound | ### 4.1 Cross-codec disambiguation invariant (target 2, the subtle one) @@ -182,9 +182,10 @@ fuzz/ ├── fuzz_targets/ │ ├── chunk_frame.rs thin fuzz_target! wrapper │ ├── negotiation_frame.rs thin wrapper -│ └── control_json.rs thin wrapper +│ ├── control_json.rs thin wrapper +│ └── session_opseq.rs thin wrapper (typed SessionSequence input) ├── shared/ alktty-fuzz-shared — STABLE-toolchain library -│ └── src/{chunk_frame,negotiation_frame,control_json}.rs +│ └── src/{chunk_frame,negotiation_frame,control_json,session_opseq}.rs ├── corpus// committed seeds ├── gen_fuzz_seeds.py deterministic seed generator ├── run-detached.sh §3 detached runner (CWD-independent) @@ -200,11 +201,17 @@ Root `Cargo.toml` changes (the alkcall §7.7 footguns, both required): No `#[cfg(fuzzing)]` exposure expected — `ChunkReader`/`ChunkWriter`, `NegotiationReader`/`NegotiationWriter`, `NegotiateRequest`, -`ControlMessage`, and `error_response_bytes` are all already `pub` +`ControlMessage`, `error_response_bytes`, `drive_session_pre_negotiated`, +and the `TtyBackend`/`TtyHandle`/`TtyParams` shapes are all already `pub` (alkcall needed none of its five either). The one private constant, `MAX_CHUNK_LEN`, is wire-stable (ADR-001); the shared crate carries its own copy, and the shape invariants assert the rejection boundary so a -drift is caught. +drift is caught. The `session_opseq` harness is fuzz-local for the same +reason `MockBackend`/`TestBackend` are `#[cfg(test)]`-only: a test +backend in `src/` would either leak into the public API or force the +whole mock apparatus behind `#[cfg(any(test, fuzzing))]`; the shared +crate is the fuzz-side home for it (mirrors alkcall, whose `manager_routing` +harness also lives in `fuzz/shared/`). ## 6. Corpus policy @@ -232,6 +239,7 @@ deterministic (no randomness) so seeds are reproducible. 5. Detached campaigns (10 min per target via `run-detached.sh`), triage anything found, record results in §8. 6. Decide on target 4 (`session_opseq`) after 1–3 are clean. + **Status: landed — see §8. All four targets are in scope and live.** ## 8. Progress log (append-only across sessions) @@ -273,6 +281,52 @@ deterministic (no randomness) so seeds are reproducible. §6 policy; committed seeds unchanged. - §7 step 6 (target 4 `session_opseq`) remains open for a follow-up session. +- **2026-09-28 (later still)** — **§7 step 6 landed: target 4 + `session_opseq`. All four targets are now live.** + - `fuzz/shared/src/session_opseq.rs`: `SessionOp`/`SessionSequence` + (`#[derive(Arbitrary)]`, 12-op vocabulary — client writes incl. + barrier probes, backend production/EOF, exit resolve/fail/drop, + control messages, bounded reads, yields) + a fuzz-local + `FuzzBackend` (the in-crate `MockBackend` is `#[cfg(test)]`-only) + + the harness driving `drive_session_pre_negotiated` over a + `tokio::io::duplex` pair, with a delivery model mirroring §4's + invariants. 12 committed seeds via the deterministic generator + (hand-encoded against the `arbitrary` 1.4.x derive layout — variant + tag `(u32_le * 12) >> 32`, fields LE zero-filled, keep-going byte + gating each `Vec` element; the encoding is pinned by a decode test + asserting the exact ops each seed decodes to). + - Invariants encoded: kill-on-Drop (ADR-005), exit-chunk-is-last + (ADR-004), lossless per-stream FIFO + sentinel semantics, stdin + prefix-losslessness (post-shutdown chunks never delivered), control + dispatch bounds, take-once allocation, teardown termination. + Harness bounds: 128 ops, 4 KiB/chunk, 32 KiB produce/write caps, + unbounded backend stdin (control dispatch stays prompt). + - **Harness-model correction before the campaign** (the alkcall §7.7 + pattern repeating, third time): a 60s smoke run crashed on + `ClientWrite` after the session had completed — a failed client + write can mean the *server half of the duplex dropped* because the + session finished (exit resolved → pumps join → exit chunk → + drainer exits → duplex halves drop) while the write was in flight. + The harness now treats any write/shutdown error as "the session + closed": it drains the read side and requires the exit chunk (a + close without it is a premature close — a real bug). No adapter + code changed; the finding was a harness-model gap. + - Corpus replay green on stable (9 tests, incl. the seed-decode + pin). Full checklist green: 113 + 156 all-features tests, clippy + stable/wasm + fuzz-shared clippy, fmt, doc, publish dry-run, wasm + check. + - Campaign: **10 min detached via `run-detached.sh`, exited 0 with an + empty artifact directory — no crashes, hangs, OOMs, or leaks** (all + 33 fork jobs `oom/timeout/crash: 0/0/0`). 42.3k execs at ~74/s per + job (the op-seq target is stateful — each input runs a full tokio + session with bounded reads/yields, so exec/s is inherently tens, not + millions; the value is interleaving exploration). Cov grew + 3934 → 4075 edges over the run, still finding coverage at budget + end. Grown corpus excised per the §6 policy; 12 committed seeds + unchanged; corpus replay green after. + - §7 fuzzing plan complete: all four targets live, corpus replay is + the standing gate, campaigns pre-release and on adapter/pump/wire + changes. ## 9. Relationship to existing tests diff --git a/fuzz/.gitignore b/fuzz/.gitignore index 342eb81..631b906 100644 --- a/fuzz/.gitignore +++ b/fuzz/.gitignore @@ -1,7 +1,7 @@ target corpus/* -!corpus/chunk_frame/seed-* -!corpus/negotiation_frame/seed-* -!corpus/control_json/seed-* +!corpus/*/ +corpus/*/* +!corpus/*/seed-* artifacts -coverage \ No newline at end of file +coverage diff --git a/fuzz/Cargo.lock b/fuzz/Cargo.lock index e78b278..00bab49 100644 --- a/fuzz/Cargo.lock +++ b/fuzz/Cargo.lock @@ -78,9 +78,13 @@ name = "alktty-fuzz-shared" version = "0.0.0" dependencies = [ "alktty", + "arbitrary", + "async-trait", "bytes", + "futures-core", "serde_json", "tokio", + "tokio-stream", ] [[package]] @@ -94,6 +98,9 @@ name = "arbitrary" version = "1.4.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "c3d036a3c4ab069c7b410a2ce876bd74808d2d0888a82667669f8e783a898bf1" +dependencies = [ + "derive_arbitrary", +] [[package]] name = "async-trait" @@ -181,6 +188,17 @@ version = "2.11.1" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "4583a4551df46e2792f82ceeac45e850d2e2d5debba0b91f102385cda5b11f06" +[[package]] +name = "derive_arbitrary" +version = "1.4.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1e567bd82dcff979e4b03460c307b3cdc9e96fde3d73bed1496d2bc75d9dd62a" +dependencies = [ + "proc-macro2", + "quote", + "syn 2.0.119", +] + [[package]] name = "displaydoc" version = "0.2.7" diff --git a/fuzz/Cargo.toml b/fuzz/Cargo.toml index 5f3adc9..f56ad8b 100644 --- a/fuzz/Cargo.toml +++ b/fuzz/Cargo.toml @@ -35,4 +35,11 @@ test = false doc = false bench = false +[[bin]] +name = "session_opseq" +path = "fuzz_targets/session_opseq.rs" +test = false +doc = false +bench = false + [workspace] \ No newline at end of file diff --git a/fuzz/README.md b/fuzz/README.md index 7c07255..1f6fad6 100644 --- a/fuzz/README.md +++ b/fuzz/README.md @@ -20,6 +20,7 @@ operating rules live in `docs/plans/fuzzing.md` (adopted from alkcall's | `chunk_frame` | `ChunkReader`/`ChunkWriter` (5-byte chunk codec, ADR-001) | | `negotiation_frame` | `NegotiationReader`/`NegotiationWriter` + `NegotiateRequest` + `error_response_bytes` + the cross-codec peek-disambiguation seam | | `control_json` | `ControlMessage::from_slice`/`to_json` + `signal_from_name` | +| `session_opseq` | the producer session pump (`drive_session_pre_negotiated`) + a fuzz-local mock backend — the stateful op-sequence target (ADR-004/ADR-005 invariants; the alkcall-bug-class hunt) | ## Running a campaign — always detached @@ -57,4 +58,16 @@ required by cargo-fuzz); the crate itself stays on stable at MSRV 1.88. toolchain file (rustup resolves per directory). Run campaigns before releases and after touching `src/wire.rs`, -`src/negotiation.rs`, `src/control.rs`, the adapter, or the session pump. \ No newline at end of file +`src/negotiation.rs`, `src/control.rs`, the adapter, or the session pump. + +## session_opseq notes + +The op-sequence target is inherently slower than the byte-format targets +(each input drives a full tokio session with bounded reads and yields — +tens of exec/s per worker, not millions). That is expected: the value is +state-space exploration, not throughput. It uses `tokio::time::timeout` +internally, so the nightly target needs the `--features time` tokio +surface already present via the shared crate's deps — no extra flags. +Its invariants (kill-on-Drop, exit-chunk-is-last, sentinel/FIFO +semantics) are the ones unit tests assert individually; the fuzzer +explores interleavings no example test can enumerate. \ No newline at end of file diff --git a/fuzz/corpus/session_opseq/seed-000 b/fuzz/corpus/session_opseq/seed-000 new file mode 100644 index 0000000000000000000000000000000000000000..4f8128cc7d3fce0f7caf720fa92adf482588e70e GIT binary patch literal 15 UcmZQzT)k?QmHTYSk(&7KVHV1_lQp_o_Vu08d#4g8%>k literal 0 HcmV?d00001 diff --git a/fuzz/corpus/session_opseq/seed-005 b/fuzz/corpus/session_opseq/seed-005 new file mode 100644 index 0000000000000000000000000000000000000000..f2ee45033c9c9c22e01219de3c0bc993a9bea6d8 GIT binary patch literal 40 ncmZQ%0D==hylT}dEf$7+5SJB1FfcfThK63X2QnJ||7QRIhENDQ literal 0 HcmV?d00001 diff --git a/fuzz/corpus/session_opseq/seed-006 b/fuzz/corpus/session_opseq/seed-006 new file mode 100644 index 0000000000000000000000000000000000000000..4d376a18ecff10bd3d511f2df33b7c87bc44428f GIT binary patch literal 105 xcmZQzU|=}F5WpY-q!|t{g2>R&(5oyU9*6@1=)(Vj0E=)1g9fJR28em~3;_2s6#M`H literal 0 HcmV?d00001 diff --git a/fuzz/corpus/session_opseq/seed-007 b/fuzz/corpus/session_opseq/seed-007 new file mode 100644 index 0000000000000000000000000000000000000000..f1939a1d4fc712986abd2ed0f7eaeb8fda186a7c GIT binary patch literal 50 vcmZQzU|=}FoD3vE3>FB@z`(-9P_=5+DlHa%}klVn(0OVh_X8-`_MhhYU literal 0 HcmV?d00001 diff --git a/fuzz/corpus/session_opseq/seed-009 b/fuzz/corpus/session_opseq/seed-009 new file mode 100644 index 0000000000000000000000000000000000000000..991bfeb7fd8f744dc22e5e8caa5dea5eaa10dd37 GIT binary patch literal 22 YcmZQr1p!(t4EYQU3=W~8p;zr00A#-hmjD0& literal 0 HcmV?d00001 diff --git a/fuzz/corpus/session_opseq/seed-010 b/fuzz/corpus/session_opseq/seed-010 new file mode 100644 index 0000000000000000000000000000000000000000..b260c80ac13f3da9a2e04232a65fce72f1b08fbe GIT binary patch literal 35 lcmZP!1p*d^)K#liWdX4k3qw8w149ERkYI2CN?o;Q006uj37!A| literal 0 HcmV?d00001 diff --git a/fuzz/corpus/session_opseq/seed-011 b/fuzz/corpus/session_opseq/seed-011 new file mode 100644 index 0000000000000000000000000000000000000000..01f0f220317d75974054910dd322ecd434a3bfa5 GIT binary patch literal 48 scmZQz00Ab3L=eTu#E=A{f&643*}%xKYSk(&7KVHv$00N{^r}4r0ARlcod5s; literal 0 HcmV?d00001 diff --git a/fuzz/fuzz_targets/session_opseq.rs b/fuzz/fuzz_targets/session_opseq.rs new file mode 100644 index 0000000..d4d2252 --- /dev/null +++ b/fuzz/fuzz_targets/session_opseq.rs @@ -0,0 +1,9 @@ +#![no_main] + +use libfuzzer_sys::fuzz_target; + +fuzz_target!( + |seq: alktty_fuzz_shared::session_opseq::SessionSequence| { + alktty_fuzz_shared::session_opseq::fuzz_session_opseq(&seq); + } +); \ No newline at end of file diff --git a/fuzz/gen_fuzz_seeds.py b/fuzz/gen_fuzz_seeds.py index fff60a6..c692f44 100644 --- a/fuzz/gen_fuzz_seeds.py +++ b/fuzz/gen_fuzz_seeds.py @@ -241,10 +241,277 @@ def control_json_seeds(): return seeds +# session_opseq seeds are arbitrary-encoded `SessionSequence` inputs, so +# they are hand-encoded against the `arbitrary` 1.4.x derive layout the +# same way alkcall's manager_routing seeds are: each element of a Vec is +# gated by a keep-going byte (odd = another element follows; the final +# gated read before data exhaustion yields true via zero-fill), enum +# variant tags are 4-byte LE u32s with `(tag * 12) >> 32` selecting by +# declaration order, and fields follow in declaration order +# little-endian, zero-filled on short data. The encoding is pinned by +# the seed-decode test in fuzz/shared/src/session_opseq.rs. +SESSION_OP_VARIANTS = [ + "ClientWrite", + "ClientCloseWrite", + "Stdout", + "StdoutEof", + "Stderr", + "StderrEof", + "Exit", + "ExitFail", + "DropExitTx", + "CtrlIn", + "ClientReadAll", + "Yield", +] + + +def _le32(v): + return struct.pack("> 32` selecting the +/// variant by declaration order; fields follow in declaration order, +/// little-endian, zero-filled on short data; each element of the ops +/// vector is gated by a keep-going byte (odd = another element follows). +#[derive(Debug, Arbitrary, Clone)] +pub enum SessionOp { + /// Write one chunk-framed record on the client→server direction. + /// `st <= 4` is a real stream type (0 stdin, 1/2 ignored by the + /// input pump, 3 control, 4 protocol violation); `st == 5` writes an + /// oversize length header (input-pump barrier); `st >= 6` writes an + /// invalid stream type (input-pump barrier). + ClientWrite { st: u8, len: u16, byte: u8 }, + /// Close the client write half (input pump observes EOF). + ClientCloseWrite, + /// Backend produces one stdout payload (`len % 4096 + 1` bytes, + /// saturating at the 32 KiB production cap). + Stdout { len: u16, byte: u8 }, + /// Backend closes stdout (the pump then emits the sentinel). + StdoutEof, + /// Backend produces one stderr payload (no-op when the backend was + /// built without a stderr stream). + Stderr { len: u16, byte: u8 }, + /// Backend closes stderr. + StderrEof, + /// Resolve the exit future with `Ok(code)` (first resolution wins). + Exit { code: i16 }, + /// Resolve the exit future with `Err(WaitFailed)` (adapter sends -1). + ExitFail, + /// Drop the exit sender without resolving (future resolves to a + /// wait-failure; adapter sends -1). + DropExitTx, + /// Client→server control chunk; `kind % 5` selects resize / signal / + /// eof / exit-on-ctrl-in (protocol violation, ignored) / garbage. + CtrlIn { kind: u8, a: u16, b: u16 }, + /// Read up to `max % 64 + 1` chunks from the client side, feeding + /// the delivery model. + ClientReadAll { max: u8 }, + /// Let the spawned pumps make progress. + Yield { ticks: u8 }, +} + +#[derive(Debug, Arbitrary, Clone)] +pub struct SessionSequence { + pub stderr_enabled: bool, + pub ops: Vec, +} + +#[derive(Default)] +struct GuardState { + resolved: AtomicBool, + kill_fired: AtomicBool, +} + +struct FuzzExitFuture { + rx: Option>>, + guard: Arc, +} + +impl Future for FuzzExitFuture { + type Output = Result; + + fn poll(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll { + if let Some(rx) = self.rx.as_mut() { + if let Poll::Ready(v) = Pin::new(rx).poll(cx) { + self.rx = None; + self.guard.resolved.store(true, Ordering::SeqCst); + return Poll::Ready( + v.map_err(|_| TtyError::WaitFailed { + message: "exit sender dropped".to_string(), + }) + .and_then(|r| r), + ); + } + } + Poll::Pending + } +} + +impl Drop for FuzzExitFuture { + fn drop(&mut self) { + if self.rx.is_some() { + self.guard.kill_fired.store(true, Ordering::SeqCst); + } + } +} + +#[derive(Default)] +struct FuzzControl { + resize_calls: AtomicU32, + signal_calls: AtomicU32, +} + +impl TtyControl for FuzzControl { + fn resize(&self, _: u16, _: u16, _: u16, _: u16) { + self.resize_calls.fetch_add(1, Ordering::SeqCst); + } + + fn signal(&self, _: &str) { + self.signal_calls.fetch_add(1, Ordering::SeqCst); + } +} + +struct FuzzStdinSink { + tx: Option>, +} + +impl AsyncWrite for FuzzStdinSink { + fn poll_write( + self: Pin<&mut Self>, + _: &mut Context<'_>, + buf: &[u8], + ) -> Poll> { + match self.tx.as_ref() { + Some(tx) => match tx.send(Bytes::copy_from_slice(buf)) { + Ok(()) => Poll::Ready(Ok(buf.len())), + Err(_) => Poll::Ready(Err(std::io::Error::new( + std::io::ErrorKind::BrokenPipe, + "stdin closed", + ))), + }, + None => Poll::Ready(Err(std::io::Error::new( + std::io::ErrorKind::BrokenPipe, + "stdin shut down", + ))), + } + } + + fn poll_flush(self: Pin<&mut Self>, _: &mut Context<'_>) -> Poll> { + Poll::Ready(Ok(())) + } + + fn poll_shutdown( + self: Pin<&mut Self>, + _: &mut Context<'_>, + ) -> Poll> { + let this = self.get_mut(); + this.tx = None; + Poll::Ready(Ok(())) + } +} + +struct Slots { + stdout_tx: Option>, + stderr_tx: Option>, + stdin_rx: mpsc::UnboundedReceiver, + exit_tx: Option>>, + control: Arc, + guard: Arc, +} + +struct FuzzBackend { + stderr_enabled: bool, + slots: Mutex>, + allocations: AtomicU32, +} + +impl FuzzBackend { + fn new(stderr_enabled: bool) -> Self { + Self { + stderr_enabled, + slots: Mutex::new(Vec::new()), + allocations: AtomicU32::new(0), + } + } + + fn has_slots(&self) -> bool { + !self + .slots + .lock() + .unwrap_or_else(|e| e.into_inner()) + .is_empty() + } + + fn take_slots(&self) -> Option { + let mut slots = self.slots.lock().unwrap_or_else(|e| e.into_inner()); + if slots.is_empty() { + None + } else { + Some(slots.remove(0)) + } + } +} + +#[async_trait] +impl TtyBackend for FuzzBackend { + async fn allocate(&self, _params: &TtyParams) -> Result { + self.allocations.fetch_add(1, Ordering::SeqCst); + let (stdout_tx, stdout_rx) = mpsc::channel::(8); + let (stdin_tx, stdin_rx) = mpsc::unbounded_channel::(); + let (exit_tx, exit_rx) = oneshot::channel::>(); + let guard = Arc::new(GuardState::default()); + let control = Arc::new(FuzzControl::default()); + + let stderr_tx = if self.stderr_enabled { + Some(mpsc::channel::(8)) + } else { + None + }; + if let Some((tx, _)) = &stderr_tx { + self.slots + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(Slots { + stdout_tx: Some(stdout_tx.clone()), + stderr_tx: Some(tx.clone()), + stdin_rx, + exit_tx: Some(exit_tx), + control: Arc::clone(&control), + guard: Arc::clone(&guard), + }); + } else { + self.slots + .lock() + .unwrap_or_else(|e| e.into_inner()) + .push(Slots { + stdout_tx: Some(stdout_tx.clone()), + stderr_tx: None, + stdin_rx, + exit_tx: Some(exit_tx), + control: Arc::clone(&control), + guard: Arc::clone(&guard), + }); + } + + let stdout: Pin + Send>> = + Box::pin(ReceiverStream::new(stdout_rx)); + let stderr: Option + Send>>> = + stderr_tx.map(|(_, rx)| Box::pin(ReceiverStream::new(rx)) as _); + let stdin: Box = + Box::new(FuzzStdinSink { tx: Some(stdin_tx) }); + let exit_code: BoxFuture> = Box::pin(FuzzExitFuture { + rx: Some(exit_rx), + guard, + }); + + Ok(TtyHandle { + stdin, + stdout, + stderr, + exit_code, + control: Some(TtyControlHandle::new(control)), + }) + } +} + +type ClientWriteHalf = tokio::io::WriteHalf; +type ClientReadHalf = tokio::io::ReadHalf; + +struct Harness { + backend: Arc, + slots: Option, + client_write: ClientWriteHalf, + reader: ChunkReader, + stdout_pending: VecDeque<(usize, u8)>, + stderr_pending: VecDeque<(usize, u8)>, + produced_stdout: usize, + produced_stderr: usize, + produced_bytes: usize, + client_bytes: usize, + stdout_eof: bool, + stderr_eof: bool, + sentinel_seen: bool, + exit_seen: bool, + client_closed: bool, + stdin_expected: Vec, + stdin_shutdown: bool, + pump_dead: bool, + write_closed: bool, + exit_tx_taken: bool, + exit_op_issued: bool, + exit_code_want: i32, + issued_resizes: u32, + issued_signals: u32, +} + +impl Harness { + async fn ensure_slots(&mut self) -> Result<(), String> { + if self.slots.is_some() { + return Ok(()); + } + let mut ticks = 0usize; + while !self.backend.has_slots() { + tokio::task::yield_now().await; + ticks += 1; + if ticks > 50_000 { + return Err("the valid pre-negotiated session never allocated".to_string()); + } + } + self.slots = self.backend.take_slots(); + Ok(()) + } + + async fn write_all_bounded(&mut self, bytes: &[u8]) -> Result<(), String> { + match tokio::time::timeout(Duration::from_secs(2), self.client_write.write_all(bytes)).await + { + Err(_) => Err("client write_all parked past 2s".to_string()), + Ok(Err(_)) => { + // The server half of the duplex is dropped only when the + // session task returns, and the session returns only after + // the drainer flushed the exit chunk — so a failed write + // means the session completed while this write was in + // flight. Drain and verify; a close without the exit chunk + // is a premature close (a bug). + self.note_session_closed("client write failed").await + } + Ok(Ok(())) => Ok(()), + } + } + + /// Drain the client read side to EOF, feeding the delivery model, and + /// require the exit chunk: the session may only close after the + /// exit-chunk sequence completed (ADR-004). + async fn note_session_closed(&mut self, context: &str) -> Result<(), String> { + for _ in 0..4096 { + if self.client_closed { + break; + } + match tokio::time::timeout(Duration::from_millis(10), self.reader.read_chunk()).await { + Err(_) => break, + Ok(Err(RawError::ConnectionClosed)) => { + self.client_closed = true; + break; + } + Ok(Err(e)) => { + return Err(format!( + "server wrote undecodable chunk while closing ({context}): {e:?}" + )); + } + Ok(Ok(chunk)) => self.feed_chunk(&chunk)?, + } + } + if !self.exit_seen { + return Err(format!( + "server closed the session without the exit chunk ({context})" + )); + } + self.write_closed = true; + Ok(()) + } + + fn client_write_fits(&self, len: usize) -> bool { + self.client_bytes + len <= CLIENT_WRITE_CAP + } + + async fn step(&mut self, op: &SessionOp) -> Result<(), String> { + match op { + SessionOp::ClientWrite { st, len, byte } => { + if self.pump_dead || self.write_closed || self.exit_seen { + return Ok(()); + } + let payload_len = (*len as usize) % (MAX_PAYLOAD + 1); + match *st { + 5 => { + if !self.client_write_fits(5) { + return Ok(()); + } + self.client_bytes += 5; + self.write_all_bounded(&[1u8, 0xFF, 0xFF, 0xFF, 0xFF]) + .await?; + self.pump_dead = true; + } + 6..=255 => { + if !self.client_write_fits(5) { + return Ok(()); + } + self.client_bytes += 5; + self.write_all_bounded(&[*st, 0, 0, 0, 0]).await?; + self.pump_dead = true; + } + _ => { + let st = *st; + if st == STREAM_STDIN && payload_len == 0 { + if !self.client_write_fits(5) { + return Ok(()); + } + self.client_bytes += 5; + self.write_all_bounded(&[st, 0, 0, 0, 0]).await?; + self.stdin_shutdown = true; + return Ok(()); + } + if st == STREAM_STDIN && self.stdin_shutdown { + if !self.client_write_fits(5 + payload_len) { + return Ok(()); + } + self.client_bytes += 5 + payload_len; + let mut wire = Vec::with_capacity(5 + payload_len); + wire.push(st); + wire.extend_from_slice(&(payload_len as u32).to_be_bytes()); + wire.extend(std::iter::repeat_n(*byte, payload_len)); + self.write_all_bounded(&wire).await?; + self.pump_dead = true; + return Ok(()); + } + if !self.client_write_fits(5 + payload_len) { + return Ok(()); + } + self.client_bytes += 5 + payload_len; + let mut wire = Vec::with_capacity(5 + payload_len); + wire.push(st); + wire.extend_from_slice(&(payload_len as u32).to_be_bytes()); + wire.extend(std::iter::repeat_n(*byte, payload_len)); + self.write_all_bounded(&wire).await?; + if st == STREAM_STDIN { + self.stdin_expected + .extend(std::iter::repeat_n(*byte, payload_len)); + } + } + } + Ok(()) + } + SessionOp::ClientCloseWrite => { + if self.pump_dead || self.write_closed || self.exit_seen { + return Ok(()); + } + match tokio::time::timeout(Duration::from_secs(2), self.client_write.shutdown()) + .await + { + Err(_) => return Err("client shutdown parked past 2s".to_string()), + Ok(Err(_)) => { + // Same close-race as a failed write: the session + // completed while the shutdown was in flight. + self.note_session_closed("client shutdown failed").await?; + } + Ok(Ok(())) => {} + } + self.write_closed = true; + self.stdin_shutdown = true; + Ok(()) + } + SessionOp::Stdout { len, byte } => self.produce(StreamKind::Stdout, *len, *byte).await, + SessionOp::StdoutEof => { + self.ensure_slots().await?; + let slots = self.slots.as_mut().expect("slots taken"); + let _ = slots.stdout_tx.take(); + self.stdout_eof = true; + Ok(()) + } + SessionOp::Stderr { len, byte } => { + if !self.stderr_enabled() { + return Ok(()); + } + self.produce(StreamKind::Stderr, *len, *byte).await + } + SessionOp::StderrEof => { + if !self.stderr_enabled() { + return Ok(()); + } + self.ensure_slots().await?; + let slots = self.slots.as_mut().expect("slots taken"); + let _ = slots.stderr_tx.take(); + self.stderr_eof = true; + Ok(()) + } + SessionOp::Exit { code } => self.resolve_exit(Ok(i32::from(*code))).await, + SessionOp::ExitFail => { + self.resolve_exit(Err(TtyError::WaitFailed { + message: "fuzz wait failure".to_string(), + })) + .await + } + SessionOp::DropExitTx => { + if self.exit_tx_taken { + return Ok(()); + } + self.ensure_slots().await?; + let slots = self.slots.as_mut().expect("slots taken"); + let _ = slots.exit_tx.take(); + self.exit_tx_taken = true; + self.exit_op_issued = true; + self.exit_code_want = -1; + Ok(()) + } + SessionOp::CtrlIn { kind, a, b } => { + if self.pump_dead || self.write_closed || self.exit_seen { + return Ok(()); + } + self.ensure_slots().await?; + let kind = *kind % 5; + let json: Vec = match kind { + 0 => serde_json::to_vec(&serde_json::json!({ + "type": "resize", "cols": *a, "rows": *b + })) + .expect("resize json builds"), + 1 => serde_json::to_vec(&serde_json::json!({ + "type": "signal", "name": SIGNAL_NAMES[(*a as usize) % SIGNAL_NAMES.len()] + })) + .expect("signal json builds"), + 2 => br#"{"type":"eof"}"#.to_vec(), + 3 => serde_json::to_vec(&serde_json::json!({ + "type": "exit", "code": *a + })) + .expect("exit json builds"), + _ => b"[1,2,".to_vec(), + }; + if !self.client_write_fits(5 + json.len()) { + return Ok(()); + } + self.client_bytes += 5 + json.len(); + let mut wire = Vec::with_capacity(5 + json.len()); + wire.push(STREAM_CTRL_IN); + wire.extend_from_slice(&(json.len() as u32).to_be_bytes()); + wire.extend_from_slice(&json); + self.write_all_bounded(&wire).await?; + if kind == 2 { + self.stdin_shutdown = true; + } + if kind == 0 { + self.issued_resizes += 1; + self.poll_ctrl( + |c| c.resize_calls.load(Ordering::SeqCst), + self.issued_resizes, + ) + .await?; + } + if kind == 1 { + self.issued_signals += 1; + self.poll_ctrl( + |c| c.signal_calls.load(Ordering::SeqCst), + self.issued_signals, + ) + .await?; + } + Ok(()) + } + SessionOp::ClientReadAll { max } => { + let budget = (*max as usize) % 64 + 1; + for _ in 0..budget { + if self.client_closed { + break; + } + match tokio::time::timeout(Duration::from_millis(10), self.reader.read_chunk()) + .await + { + Err(_) => break, + Ok(Err(RawError::ConnectionClosed)) => { + self.client_closed = true; + break; + } + Ok(Err(e)) => { + return Err(format!("server wrote undecodable chunk: {e:?}")); + } + Ok(Ok(chunk)) => self.feed_chunk(&chunk)?, + } + } + Ok(()) + } + SessionOp::Yield { ticks } => { + for _ in 0..(*ticks as usize % 64) { + tokio::task::yield_now().await; + } + Ok(()) + } + } + } + + fn stderr_enabled(&self) -> bool { + self.slots + .as_ref() + .map(|s| s.stderr_tx.is_some()) + .unwrap_or_else(|| { + self.backend + .slots + .lock() + .unwrap_or_else(|e| e.into_inner()) + .first() + .map(|s| s.stderr_tx.is_some()) + .unwrap_or(true) + }) + } + + async fn produce(&mut self, kind: StreamKind, len: u16, byte: u8) -> Result<(), String> { + let payload_len = len as usize % (MAX_PAYLOAD + 1); + if payload_len == 0 { + return Ok(()); + } + if self.produced_bytes + payload_len > TOTAL_PRODUCE_CAP { + return Ok(()); + } + self.ensure_slots().await?; + let tx = { + let slots = self.slots.as_mut().expect("slots taken"); + match kind { + StreamKind::Stdout => slots.stdout_tx.as_ref(), + StreamKind::Stderr => slots.stderr_tx.as_ref(), + } + }; + let Some(tx) = tx else { + return Ok(()); + }; + let payload = vec![byte; payload_len]; + tokio::time::timeout(Duration::from_secs(2), tx.send(Bytes::from(payload))) + .await + .map_err(|_| "backend send parked past 2s (writer-channel backpressure)".to_string())? + .map_err(|_| "backend send failed (receiver dropped)".to_string())?; + self.produced_bytes += payload_len; + match kind { + StreamKind::Stdout => { + self.produced_stdout += 1; + self.stdout_pending.push_back((payload_len, byte)); + } + StreamKind::Stderr => { + self.produced_stderr += 1; + self.stderr_pending.push_back((payload_len, byte)); + } + } + Ok(()) + } + + async fn resolve_exit(&mut self, value: Result) -> Result<(), String> { + if self.exit_tx_taken { + return Ok(()); + } + self.ensure_slots().await?; + let slots = self.slots.as_mut().expect("slots taken"); + let Some(tx) = slots.exit_tx.take() else { + return Ok(()); + }; + self.exit_tx_taken = true; + self.exit_op_issued = true; + self.exit_code_want = value.as_ref().copied().unwrap_or(-1); + let _ = tx.send(value); + Ok(()) + } + + async fn poll_ctrl( + &self, + counter: impl Fn(&FuzzControl) -> u32, + want: u32, + ) -> Result<(), String> { + if self.exit_op_issued { + return Ok(()); + } + let Some(slots) = self.slots.as_ref() else { + return Err("control dispatch polled before allocation".to_string()); + }; + let mut ticks = 0usize; + while counter(&slots.control) < want { + tokio::task::yield_now().await; + ticks += 1; + if ticks > 200_000 { + return Err( + "control message not dispatched while the input pump was alive".to_string(), + ); + } + } + Ok(()) + } + + fn feed_chunk(&mut self, chunk: &Chunk) -> Result<(), String> { + if self.exit_seen { + return Err(format!( + "chunk of stream_type {} observed after the exit chunk", + chunk.stream_type + )); + } + match chunk.stream_type { + STREAM_STDOUT => { + if chunk.bytes.is_empty() { + if self.sentinel_seen { + return Err("duplicate stdout sentinel".to_string()); + } + if !self.stdout_eof { + return Err("sentinel observed before the stdout EOF op".to_string()); + } + if !self.stdout_pending.is_empty() { + return Err("sentinel observed before all stdout data chunks".to_string()); + } + self.sentinel_seen = true; + } else { + let Some((len, byte)) = self.stdout_pending.pop_front() else { + return Err(format!( + "unexpected stdout data chunk of {} bytes (produced {} total)", + chunk.bytes.len(), + self.produced_stdout + )); + }; + if chunk.bytes.len() != len || chunk.bytes.iter().any(|b| *b != byte) { + return Err("stdout payload deviates from the produced pattern".to_string()); + } + } + Ok(()) + } + STREAM_STDERR => { + if chunk.bytes.is_empty() { + return Err( + "stderr sentinel observed (the stderr pump never sends one)".to_string() + ); + } + let Some((len, byte)) = self.stderr_pending.pop_front() else { + return Err(format!( + "unexpected stderr data chunk (nothing pending, produced {} total)", + self.produced_stderr + )); + }; + if chunk.bytes.len() != len || chunk.bytes.iter().any(|b| *b != byte) { + return Err("stderr payload deviates from the produced pattern".to_string()); + } + Ok(()) + } + STREAM_CTRL_OUT => { + if !self.exit_op_issued { + return Err("exit chunk observed without an issued exit op".to_string()); + } + let v: serde_json::Value = serde_json::from_slice(&chunk.bytes) + .map_err(|e| format!("exit chunk is not JSON: {e}"))?; + if v.get("type").and_then(|t| t.as_str()) != Some("exit") { + return Err(format!("ctrl_out chunk is not an exit message: {v}")); + } + let code = v + .get("code") + .and_then(|c| c.as_i64()) + .ok_or_else(|| format!("exit chunk has no code: {v}"))?; + if code != i64::from(self.exit_code_want) { + return Err(format!( + "exit code {code} != issued {}", + self.exit_code_want + )); + } + if !self.sentinel_seen { + return Err("exit chunk observed before the stdout sentinel".to_string()); + } + if !self.stdout_pending.is_empty() { + return Err("exit chunk observed before all stdout data chunks".to_string()); + } + if !self.stderr_pending.is_empty() { + return Err("exit chunk observed before all stderr data chunks".to_string()); + } + self.exit_seen = true; + Ok(()) + } + other => Err(format!("server wrote stream_type {other}")), + } + } + + async fn teardown(mut self, session: tokio::task::JoinHandle<()>) { + self.ensure_slots() + .await + .expect("allocation observable at teardown"); + session.abort(); + let join = tokio::time::timeout(Duration::from_secs(10), session).await; + match join { + Err(_) => panic!( + "session_opseq invariant violated: session task did not terminate after abort (deadlock)" + ), + Ok(Err(e)) if e.is_panic() => { + panic!("session_opseq invariant violated: session task panicked: {e}") + } + Ok(_) => {} + } + for _ in 0..4096 { + if self.client_closed { + break; + } + match tokio::time::timeout(Duration::from_millis(10), self.reader.read_chunk()).await { + Err(_) => break, + Ok(Err(RawError::ConnectionClosed)) => { + self.client_closed = true; + break; + } + Ok(Err(e)) => { + panic!("session_opseq invariant violated: server wrote undecodable chunk {e:?}") + } + Ok(Ok(chunk)) => { + if let Err(msg) = self.feed_chunk(&chunk) { + panic!("session_opseq invariant violated during teardown drain: {msg}"); + } + } + } + } + + assert_eq!( + self.backend.allocations.load(Ordering::SeqCst), + 1, + "session_opseq invariant violated: the pre-negotiated driver must allocate exactly once" + ); + + let slots = self.slots.as_mut().expect("slots taken at teardown"); + let mut received = Vec::new(); + let mut over = 0usize; + while let Ok(b) = slots.stdin_rx.try_recv() { + received.extend_from_slice(&b); + over += 1; + assert!( + over <= 1 << 20, + "session_opseq invariant violated: stdin received exceeds any issued volume" + ); + } + assert!( + received.len() <= self.stdin_expected.len() + && received[..] == self.stdin_expected[..received.len()], + "session_opseq invariant violated: stdin received deviates from the issued prefix (expected {} bytes, got {})", + self.stdin_expected.len(), + received.len() + ); + + assert!( + slots.control.resize_calls.load(Ordering::SeqCst) <= self.issued_resizes, + "session_opseq invariant violated: resize dispatched more often than issued" + ); + assert!( + slots.control.signal_calls.load(Ordering::SeqCst) <= self.issued_signals, + "session_opseq invariant violated: signal dispatched more often than issued" + ); + + let resolved = slots.guard.resolved.load(Ordering::SeqCst); + let killed = slots.guard.kill_fired.load(Ordering::SeqCst); + assert!( + !(resolved && killed), + "session_opseq invariant violated: kill guard fired on a resolved exit future (ADR-005 §2)" + ); + if !self.exit_op_issued { + assert!( + killed, + "session_opseq invariant violated: cancelled session did not fire the kill guard (ADR-005)" + ); + assert!( + !resolved, + "session_opseq invariant violated: exit resolved without an exit op" + ); + } + if self.exit_seen { + assert!( + resolved, + "session_opseq invariant violated: exit chunk observed but the exit future never resolved" + ); + } + } +} + +enum StreamKind { + Stdout, + Stderr, +} + +/// The whole-sequence driver. +pub fn fuzz_session_opseq(seq: &SessionSequence) { + tokio::runtime::Builder::new_current_thread() + .enable_all() + .build() + .expect("current-thread runtime builds") + .block_on(async { + let stderr_enabled = seq.stderr_enabled; + let backend = Arc::new(FuzzBackend::new(stderr_enabled)); + let mut map: HashMap> = HashMap::new(); + map.insert( + "fuzz".to_string(), + Arc::clone(&backend) as Arc, + ); + let req = NegotiateRequest { + carriage: "raw".to_string(), + backend: "fuzz".to_string(), + tty: None, + cmd: vec!["fuzz".to_string()], + cwd: None, + env: HashMap::new(), + backend_params: serde_json::Map::new(), + }; + let (client, server) = tokio::io::duplex(64 * 1024); + let (client_read, client_write) = tokio::io::split(client); + let (server_read, server_write) = tokio::io::split(server); + let backends = Arc::new(map); + let session = tokio::spawn(drive_session_pre_negotiated( + server_write, + server_read, + req, + backends, + None, + None, + )); + let mut harness = Harness { + backend, + slots: None, + client_write, + reader: ChunkReader::new(client_read), + stdout_pending: VecDeque::new(), + stderr_pending: VecDeque::new(), + produced_stdout: 0, + produced_stderr: 0, + produced_bytes: 0, + client_bytes: 0, + stdout_eof: false, + stderr_eof: false, + sentinel_seen: false, + exit_seen: false, + client_closed: false, + stdin_expected: Vec::new(), + stdin_shutdown: false, + pump_dead: false, + write_closed: false, + exit_tx_taken: false, + exit_op_issued: false, + exit_code_want: 0, + issued_resizes: 0, + issued_signals: 0, + }; + for (i, op) in seq.ops.iter().take(MAX_OPS).enumerate() { + if let Err(msg) = harness.step(op).await { + panic!("session_opseq invariant violated at op {i} ({op:?}): {msg}"); + } + } + harness.teardown(session).await; + }); +} + +#[cfg(test)] +mod corpus_replay { + use super::*; + + fn decode(data: &[u8]) -> SessionSequence { + let u = arbitrary::Unstructured::new(data); + arbitrary::Arbitrary::arbitrary_take_rest(u).expect("seed decodes") + } + + #[test] + fn committed_corpus_replays_through_the_invariants() { + let corpus = + std::path::PathBuf::from(env!("CARGO_MANIFEST_DIR")).join("../corpus/session_opseq"); + let mut count = 0usize; + for entry in std::fs::read_dir(&corpus).expect("corpus directory is committed") { + let path = entry.expect("corpus entry readable").path(); + let data = std::fs::read(&path).expect("corpus entry readable"); + fuzz_session_opseq(&decode(&data)); + count += 1; + } + assert!( + count >= 8, + "committed seed corpus is present, found {count}" + ); + } + + /// The hand-encoded seed bytes decode to the intended op sequences + /// (arbitrary 1.4.x layout: variant tag `(u32_le * 12) >> 32`, fields + /// LE zero-filled, keep-going byte before each element). + #[test] + fn seed_bytes_decode_to_the_intended_sequences() { + let seq = decode(&[ + 0x00, 0x01, 0xAB, 0xAA, 0xAA, 0x2A, 0x10, 0x00, 0x61, 0x01, 0xAB, 0xAA, 0xAA, 0xEA, + 0x04, + ]); + assert!(!seq.stderr_enabled); + assert_eq!(seq.ops.len(), 2); + assert!(matches!( + seq.ops[0], + SessionOp::Stdout { + len: 16, + byte: 0x61 + } + )); + assert!(matches!(seq.ops[1], SessionOp::Yield { ticks: 4 })); + + // The seed-002 fixture: stdin flow, EOF sentinel, post-EOF trap + // chunk, Eof control, then the full data/exit/read tail. + let seq = decode(&[ + 0x81, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x03, 0x00, 0x78, 0x01, 0x00, 0x00, 0x00, + 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x00, 0x00, 0x00, 0x00, 0x01, 0x00, 0x79, + 0x01, 0x00, 0x00, 0x00, 0xC0, 0x02, 0x00, 0x00, 0x00, 0x00, 0x01, 0xAB, 0xAA, 0xAA, + 0x2A, 0x04, 0x00, 0x6F, 0x01, 0x01, 0x00, 0x00, 0x40, 0x01, 0x00, 0x00, 0x00, 0x80, + 0x03, 0x00, 0x01, 0x56, 0x55, 0x55, 0xD5, 0x3F, 0x00, + ]); + assert!(seq.stderr_enabled); + assert_eq!(seq.ops.len(), 8); + assert!(matches!( + seq.ops[0], + SessionOp::ClientWrite { + st: 0, + len: 3, + byte: 0x78 + } + )); + assert!(matches!( + seq.ops[1], + SessionOp::ClientWrite { st: 0, len: 0, .. } + )); + assert!(matches!( + seq.ops[2], + SessionOp::ClientWrite { + st: 0, + len: 1, + byte: 0x79 + } + )); + assert!(matches!(seq.ops[3], SessionOp::CtrlIn { kind: 2, .. })); + assert!(matches!( + seq.ops[4], + SessionOp::Stdout { len: 4, byte: 0x6F } + )); + assert!(matches!(seq.ops[5], SessionOp::StdoutEof)); + assert!(matches!(seq.ops[6], SessionOp::Exit { code: 3 })); + assert!(matches!(seq.ops[7], SessionOp::ClientReadAll { max: 63 })); + } + + /// Semantic fixture (alkcall §7.7 pattern): cancel mid-session with an + /// unresolved exit — the kill guard MUST fire (ADR-005). + #[test] + fn cancel_with_unresolved_exit_fires_the_kill_guard() { + fuzz_session_opseq(&SessionSequence { + stderr_enabled: false, + ops: vec![ + SessionOp::Stdout { + len: 16, + byte: b'a', + }, + SessionOp::Yield { ticks: 4 }, + ], + }); + } + + /// Semantic fixture: stdin data, then the EOF sentinel, then an `Eof` + /// control, then a further stdin chunk (the trap) — the post-EOF chunk + /// must never reach the backend. + #[test] + fn stdin_after_eof_never_reaches_the_backend() { + fuzz_session_opseq(&SessionSequence { + stderr_enabled: true, + ops: vec![ + SessionOp::ClientWrite { + st: STREAM_STDIN, + len: 3, + byte: b'x', + }, + SessionOp::ClientWrite { + st: STREAM_STDIN, + len: 0, + byte: 0, + }, + SessionOp::CtrlIn { + kind: 2, + a: 0, + b: 0, + }, + SessionOp::ClientWrite { + st: STREAM_STDIN, + len: 1, + byte: b'y', + }, + SessionOp::Stdout { len: 4, byte: b'o' }, + SessionOp::StdoutEof, + SessionOp::Stderr { len: 4, byte: b'e' }, + SessionOp::StderrEof, + SessionOp::Exit { code: 3 }, + SessionOp::ClientReadAll { max: 63 }, + ], + }); + } + + /// Semantic fixture: duplicate exit attempts and a wait-failure must + /// yield exactly one exit chunk with the first-issued code (-1 on + /// failure), after all data and the sentinel. + #[test] + fn duplicate_exit_ops_yield_one_exit_chunk() { + fuzz_session_opseq(&SessionSequence { + stderr_enabled: true, + ops: vec![ + SessionOp::Exit { code: 7 }, + SessionOp::Exit { code: 43 }, + SessionOp::Stdout { len: 4, byte: b'o' }, + SessionOp::StdoutEof, + SessionOp::Stderr { len: 4, byte: b'e' }, + SessionOp::StderrEof, + SessionOp::ClientReadAll { max: 63 }, + ], + }); + fuzz_session_opseq(&SessionSequence { + stderr_enabled: false, + ops: vec![ + SessionOp::ExitFail, + SessionOp::Stdout { len: 4, byte: b'o' }, + SessionOp::StdoutEof, + SessionOp::ClientReadAll { max: 63 }, + ], + }); + } +}