From 9e2f90749fa008a3a70b02563fabda6b045701ed Mon Sep 17 00:00:00 2001 From: acamilo Date: Wed, 23 Sep 2026 15:52:10 +0000 Subject: [PATCH] flysim: FLY_TRACE's boundary half read back with fly-session-types, the journal held to the service test, rustfmt on the new files --- services/flysim/crates/flysim/src/frame.rs | 120 ++++++++++++++---- services/flysim/crates/flysim/src/journal.rs | 23 +++- services/flysim/crates/flysim/src/trace.rs | 45 +++++-- .../flysim/crates/flysim/tests/frame_trace.rs | 93 ++++++++++++++ .../flysim/crates/flysim/tests/integration.rs | 11 ++ 5 files changed, 254 insertions(+), 38 deletions(-) create mode 100644 services/flysim/crates/flysim/tests/frame_trace.rs diff --git a/services/flysim/crates/flysim/src/frame.rs b/services/flysim/crates/flysim/src/frame.rs index ebc88b9..8b023c8 100644 --- a/services/flysim/crates/flysim/src/frame.rs +++ b/services/flysim/crates/flysim/src/frame.rs @@ -34,12 +34,12 @@ use anyhow::{Result, anyhow}; use flybrain_core::agent::NeuralAgent; use flybrain_core::decoder::gameboy::to_button_mask; +use flybrain_gb::RewardEvent; use flybrain_gb::adapter::{GameAdapter, ProgressSnapshot}; use flybrain_gb::emulator::{Emulator, FRAMEBUFFER_LEN}; use flybrain_gb::macros::AdapterLedger; use flybrain_gb::ratchet::{Ratchet, Snapshot}; use flybrain_gb::recovery::{NeuralRecovery, recover_game}; -use flybrain_gb::RewardEvent; use crate::macros::{MacroEvent, MacroLayer, Silence}; use crate::trace::FrameTrace; @@ -238,14 +238,20 @@ impl LegacyFrame { /// A fresh start: one frame with no button down, then the brain's warm-up on it /// (`Environment.Initialize`, legacy-gameboy-v1 section 9). Returns that frame's audio. - pub fn initialize(&mut self, emulator: &mut Emulator, agent: &mut NeuralAgent) -> Result> { + pub fn initialize( + &mut self, + emulator: &mut Emulator, + agent: &mut NeuralAgent, + ) -> Result> { emulator .run_frame() .map_err(|error| anyhow!("running the first frame: {error}"))?; self.frame_buffer.copy_from_slice(emulator.framebuffer()); self.frame_counter = 1; let audio = emulator.take_audio_u8(); - agent.warmup(Some(&self.frame_buffer)).map_err(|error| anyhow!("{error}"))?; + agent + .warmup(Some(&self.frame_buffer)) + .map_err(|error| anyhow!("{error}"))?; Ok(audio) } @@ -255,24 +261,47 @@ impl LegacyFrame { /// The readout transient is left as a fresh process has it (`legacy-transient-reset`), which /// is what a restart of the service gives the fly. `import_state` of the agent is /// self-validating, so a refused checkpoint leaves the agent as it was. - pub fn restore(&mut self, parts: &mut Parts<'_>, checkpoint: &crate::store::Checkpoint) -> Result<()> { + pub fn restore( + &mut self, + parts: &mut Parts<'_>, + checkpoint: &crate::store::Checkpoint, + ) -> Result<()> { let runtime = &checkpoint.runtime; if runtime.framebuffer.len() != FRAMEBUFFER_LEN { - anyhow::bail!("checkpoint framebuffer is {} bytes", runtime.framebuffer.len()); + anyhow::bail!( + "checkpoint framebuffer is {} bytes", + runtime.framebuffer.len() + ); } - parts.agent.import_state(&checkpoint.agent).map_err(|error| anyhow!("{error}"))?; - parts.emulator.import_state(&runtime.emulator).map_err(|error| anyhow!("{error}"))?; + parts + .agent + .import_state(&checkpoint.agent) + .map_err(|error| anyhow!("{error}"))?; + parts + .emulator + .import_state(&runtime.emulator) + .map_err(|error| anyhow!("{error}"))?; if !runtime.reward.is_null() { - parts.adapter.import_state(&runtime.reward).map_err(|error| anyhow!("{error}"))?; + parts + .adapter + .import_state(&runtime.reward) + .map_err(|error| anyhow!("{error}"))?; } let snapshot = if runtime.ratchet_game.is_empty() { None } else { - Some(Snapshot { game: runtime.ratchet_game.clone(), frame: runtime.ratchet_frame.clone() }) + Some(Snapshot { + game: runtime.ratchet_game.clone(), + frame: runtime.ratchet_frame.clone(), + }) }; parts .ratchet - .import(Some(runtime.ratchet), snapshot, parts.adapter.rank_ladder().len()) + .import( + Some(runtime.ratchet), + snapshot, + parts.adapter.rank_ladder().len(), + ) .map_err(|error| anyhow!("{error}"))?; self.remainder = checkpoint.agent.remainder; @@ -280,7 +309,10 @@ impl LegacyFrame { self.buttons = runtime.buttons; self.frame_buffer.copy_from_slice(&runtime.framebuffer); let (width, height) = (parts.agent.frame.width, parts.agent.frame.height); - parts.agent.network.set_visual_frame(&self.frame_buffer, width, height); + parts + .agent + .network + .set_visual_frame(&self.frame_buffer, width, height); parts.emulator.set_buttons(self.buttons as u8); Ok(()) } @@ -325,11 +357,24 @@ impl LegacyFrame { let audio = self.advance(parts.emulator, parts.agent, observer)?; observer.after(FramePhase::Advanced, parts.agent); - let evaluated = self.evaluate(parts.emulator, parts.adapter, parts.macros.as_deref_mut(), ms); + let evaluated = self.evaluate( + parts.emulator, + parts.adapter, + parts.macros.as_deref_mut(), + ms, + ); self.commit(parts.agent, &evaluated.rewards, ms); observer.after(FramePhase::Committed, parts.agent); - Ok(Transition { ticks, ms, bound, active, executed, audio, evaluated }) + Ok(Transition { + ticks, + ms, + bound, + active, + executed, + audio, + evaluated, + }) } /// The rest of phase B and phase C behind a stub readout, for the drivers that measure the @@ -379,7 +424,12 @@ impl LegacyFrame { /// Phase A, the readout: decode the rates with the blocked direction and the bound channels, /// and restart the blocked window when the held channel changes. - pub fn decode(&mut self, agent: &mut NeuralAgent, boot: bool, bound: Option<&[String]>) -> Vec { + pub fn decode( + &mut self, + agent: &mut NeuralAgent, + boot: bool, + bound: Option<&[String]>, + ) -> Vec { let ms = agent.network.ms; let rates = agent.network.rates.clone(); // The readout's blocked-direction cooldown (`docs/readout.md`): the direction the group @@ -392,7 +442,9 @@ impl LegacyFrame { .then(|| agent.decoder.current()) .flatten() .map(str::to_string); - let active = agent.decoder.decode_bound(&rates, ms, boot, blocked.as_deref(), bound); + let active = agent + .decoder + .decode_bound(&rates, ms, boot, blocked.as_deref(), bound); // A new winner starts its own window: it has not had a hold to move in yet. let held = agent.decoder.current().map(str::to_string); if held != self.held_channel { @@ -423,9 +475,16 @@ impl LegacyFrame { let ledger = AdapterLedger(adapter); let decision = layer.decide(active, self.buttons, ms, emulator, &ledger); self.buttons = decision.mask; - Executed { mask: decision.mask, events: decision.events, silence: decision.silence } + Executed { + mask: decision.mask, + events: decision.events, + silence: decision.silence, + } } - None => Executed { mask: self.buttons, ..Executed::default() }, + None => Executed { + mask: self.buttons, + ..Executed::default() + }, }; if let Some(trace) = self.trace.as_mut() { trace.decided(active); @@ -496,14 +555,20 @@ impl LegacyFrame { if let Some(trace) = self.trace.as_mut() { trace.evaluated(&rewards, &abandoned, progress.rank); } - Evaluated { rewards, abandoned, progress } + Evaluated { + rewards, + abandoned, + progress, + } } /// Phase D: the frame just produced becomes the next ticks' visual drive, each reward event /// stimulates once, and the summed value reinforces once. pub fn commit(&mut self, agent: &mut NeuralAgent, rewards: &[RewardEvent], ms: f64) { let (width, height) = (agent.frame.width, agent.frame.height); - agent.network.set_visual_frame(&self.frame_buffer, width, height); + agent + .network + .set_visual_frame(&self.frame_buffer, width, height); let mut total = 0.0; for event in rewards { agent.network.stimulate(f64::from(event.stimulation_ms)); @@ -541,7 +606,10 @@ impl LegacyFrame { // rule as amended 2026-09-22): the macro layer answers "nearer the objective" with the map // graph it already walks (`docs/design/macros.md` section 12.15); in raw mode there is no // layer and no objective, and the answer is false. - let nearer = parts.macros.as_deref().is_some_and(MacroLayer::nearer_the_objective); + let nearer = parts + .macros + .as_deref() + .is_some_and(MacroLayer::nearer_the_objective); let trace = &mut self.trace; let mut saved = false; let recover = parts.ratchet.observe_with_progress( @@ -563,14 +631,20 @@ impl LegacyFrame { ); let rollback = if recover { // Two triggers, two stories on the ticker: a game over ended the run, a stall did not. - let trigger = - if parts.adapter.game_over() { RollbackTrigger::GameOver } else { RollbackTrigger::Stall }; + let trigger = if parts.adapter.game_over() { + RollbackTrigger::GameOver + } else { + RollbackTrigger::Stall + }; let events = self.rollback(parts)?; Some(Rollback { trigger, events }) } else { None }; - Ok(Boundary { captured: saved, rollback }) + Ok(Boundary { + captured: saved, + rollback, + }) } /// The ratchet's game-only rollback (`legacy-ratchet-rollback-v1`): the slot is restored, the diff --git a/services/flysim/crates/flysim/src/journal.rs b/services/flysim/crates/flysim/src/journal.rs index 472c2b7..019c54a 100644 --- a/services/flysim/crates/flysim/src/journal.rs +++ b/services/flysim/crates/flysim/src/journal.rs @@ -79,7 +79,10 @@ pub struct SugarJournal { impl SugarJournal { pub fn new(hot_dir: &Path) -> Self { - Self { path: hot_dir.join(FILE_NAME), file: None } + Self { + path: hot_dir.join(FILE_NAME), + file: None, + } } pub fn path(&self) -> &Path { @@ -90,7 +93,11 @@ impl SugarJournal { /// torn last line, which a reader skips. pub fn record(&mut self, entry: &Entry<'_>) { if self.file.is_none() { - match OpenOptions::new().create(true).append(true).open(&self.path) { + match OpenOptions::new() + .create(true) + .append(true) + .open(&self.path) + { Ok(file) => self.file = Some(file), Err(error) => { tracing::warn!(%error, path = %self.path.display(), "could not open the sugar journal"); @@ -128,7 +135,10 @@ mod tests { fn inputs_are_appended_with_their_frame_and_survive_a_reopen() { let dir = tempfile::tempdir().expect("a temp dir"); let mut journal = SugarJournal::new(dir.path()); - assert!(!journal.path().exists(), "nothing is written before an input"); + assert!( + !journal.path().exists(), + "nothing is written before an input" + ); let sugar = Entry { frame: 42, brain_ms: 703.0, @@ -141,7 +151,12 @@ mod tests { journal.record(&sugar); drop(journal); let mut journal = SugarJournal::new(dir.path()); - journal.record(&Entry { frame: 43, input: Input::Reward { value: 0.5 }, event_id: 8, ..sugar }); + journal.record(&Entry { + frame: 43, + input: Input::Reward { value: 0.5 }, + event_id: 8, + ..sugar + }); // A torn tail from a crash is skipped, not fatal. std::fs::OpenOptions::new() .append(true) diff --git a/services/flysim/crates/flysim/src/trace.rs b/services/flysim/crates/flysim/src/trace.rs index c7274b0..3cd9abf 100644 --- a/services/flysim/crates/flysim/src/trace.rs +++ b/services/flysim/crates/flysim/src/trace.rs @@ -170,12 +170,14 @@ impl FrameTrace { /// Sugar admitted at the top of the frame, before the brain ticks. pub fn sugar(&mut self, duration_ms: f64) { - self.admissions.push(json!({ "kind": "sugar", "durationMs": duration_ms })); + self.admissions + .push(json!({ "kind": "sugar", "durationMs": duration_ms })); } /// An operator reward pulse, applied at the top of the frame. pub fn reward_pulse(&mut self, value: f64) { - self.admissions.push(json!({ "kind": "reward", "value": value })); + self.admissions + .push(json!({ "kind": "reward", "value": value })); } /// A checkpoint capture at the current boundary: the open transition's, or the start's. @@ -213,7 +215,9 @@ impl FrameTrace { /// Phase A's ticks, and the network they left behind. pub fn ticked(&mut self, ticks: u64, remainder_ms: f64, network: &LifNetwork) { - let Some(record) = self.open.as_mut() else { return }; + let Some(record) = self.open.as_mut() else { + return; + }; record.ticks = ticks; record.brain_ticks = network.ms as u64; record.remainder = remainder_ns(remainder_ms); @@ -241,7 +245,9 @@ impl FrameTrace { pub fn executed(&mut self, mask: u32, events: &[MacroEvent]) { if let Some(record) = self.open.as_mut() { record.mask = mask; - record.macros.extend(events.iter().map(|event| macro_json("execute", event))); + record + .macros + .extend(events.iter().map(|event| macro_json("execute", event))); } } @@ -250,7 +256,9 @@ impl FrameTrace { pub fn advanced(&mut self, framebuffer: &[u8], emulator: &Emulator) { if let Some(record) = self.open.as_mut() { record.framebuffer = sha256_hex(framebuffer); - let wram: Vec = WRAM.map(|address| emulator.read_uncached(address)).collect(); + let wram: Vec = WRAM + .map(|address| emulator.read_uncached(address)) + .collect(); record.wram = sha256_hex(&wram); } } @@ -265,7 +273,9 @@ impl FrameTrace { "stimulationMs": event.stimulation_ms, }) })); - record.macros.extend(abandoned.iter().map(|event| macro_json("evaluate", event))); + record + .macros + .extend(abandoned.iter().map(|event| macro_json("evaluate", event))); record.rank = rank; } } @@ -289,12 +299,16 @@ impl FrameTrace { "slotId": SLOT, "stateDigest": null, })); - record.macros.extend(events.iter().map(|event| macro_json("rollback", event))); + record + .macros + .extend(events.iter().map(|event| macro_json("rollback", event))); } } fn flush_open(&mut self) { - let Some(record) = self.open.take() else { return }; + let Some(record) = self.open.take() else { + return; + }; let line = json!({ "behaviour": { "step": record.step.to_string(), @@ -348,11 +362,20 @@ mod tests { remainder += per_frame; remainder -= remainder.floor(); let value = remainder_ns(remainder); - assert!(value.get("numerator").is_some(), "{remainder} was not exact: {value}"); + assert!( + value.get("numerator").is_some(), + "{remainder} was not exact: {value}" + ); } - assert_eq!(remainder_ns(0.0), json!({ "numerator": "0", "denominator": "1" })); + assert_eq!( + remainder_ns(0.0), + json!({ "numerator": "0", "denominator": "1" }) + ); // 0.5 ms is 500000 ns. - assert_eq!(remainder_ns(0.5), json!({ "numerator": "500000", "denominator": "1" })); + assert_eq!( + remainder_ns(0.5), + json!({ "numerator": "500000", "denominator": "1" }) + ); // One 2^-15 ms step is 15625/512 ns. assert_eq!( remainder_ns(1.0 / 32_768.0), diff --git a/services/flysim/crates/flysim/tests/frame_trace.rs b/services/flysim/crates/flysim/tests/frame_trace.rs new file mode 100644 index 0000000..c92c0f2 --- /dev/null +++ b/services/flysim/crates/flysim/tests/frame_trace.rs @@ -0,0 +1,93 @@ +//! `FLY_TRACE`'s boundary half is the session framework's, field for field. +//! +//! The recorder writes `behaviour.boundaryActions` and `operational.captures` in the shapes of +//! `TraceBehaviour.boundaryActions` and `TraceOperational.captures` (step-v1 section 8, amended +//! 2026-09-23). This test records the two boundaries the legacy loop produces and reads them with +//! `fly-session-types` itself, by grafting them onto the shared baseline trace: +//! +//! - a rollback followed by the post-recovery checkpoint is a valid transition trace; +//! - a rung climb -- the milestone archive, *then* the ratchet's slot save -- is refused with the +//! capture-before-save error. That is the legacy order legacy-gameboy-v1 section 4 declares, and +//! the trace has to show it as it happens rather than tidy it away. + +use fly_session_types::fixtures; +use fly_session_types::scalar::DomainType; +use fly_session_types::trace::TransitionTrace; +use flysim::trace::FrameTrace; +use serde_json::Value; + +fn lines(path: &std::path::Path) -> Vec { + std::fs::read_to_string(path) + .expect("the trace") + .lines() + .map(|line| serde_json::from_str(line).expect("one JSON object per line")) + .collect() +} + +/// The shared baseline transition with this boundary's actions and captures in place of its own. +fn grafted(record: &Value) -> Result { + let file = fixtures::load("traces.json").expect("traces.json"); + let mut baseline = file.get("baseline").expect("baseline").clone(); + baseline["behaviour"]["boundaryActions"] = record["behaviour"]["boundaryActions"].clone(); + baseline["operational"]["captures"] = record["operational"]["captures"].clone(); + TransitionTrace::from_json(&baseline).map_err(|error| error.to_string()) +} + +#[test] +fn the_boundary_half_of_the_trace_is_the_session_frameworks() { + let dir = tempfile::tempdir().expect("a temp dir"); + let path = dir.path().join("trace.jsonl"); + let mut trace = FrameTrace::create(&path).expect("the trace file"); + + // The start: the boot checkpoint, before any transition. + trace.capture(1, 100); + // Transition 100 -> 101 ends in a stall rollback, then the post-recovery checkpoint. + trace.sugar(400.0); + trace.begin(100, 1_000.0); + trace.rolled_back(&[]); + trace.capture(2, 101); + // Transition 101 -> 102 climbs a rung: the archive first, then the ratchet captures. + trace.begin(101, 1_017.0); + trace.capture(3, 102); + trace.slot_saved(b"emulator state"); + trace.finish(); + drop(trace); + + let lines = lines(&path); + assert_eq!(lines.len(), 4, "format, start, two transitions: {lines:?}"); + assert_eq!(lines[0]["format"], flysim::trace::FORMAT); + assert_eq!(lines[1]["boundary"], "100"); + assert_eq!(lines[1]["operational"]["captures"][0]["checkpointId"], "g1"); + + let rollback = &lines[2]; + assert_eq!(rollback["behaviour"]["step"], "100"); + assert_eq!(rollback["behaviour"]["admissions"][0]["kind"], "sugar"); + assert_eq!( + rollback["behaviour"]["boundaryActions"][0]["kind"], + "rollback" + ); + assert_eq!(rollback["operational"]["captures"][0]["afterActions"], 1); + let parsed = grafted(rollback).expect("a rollback and then its checkpoint is a valid trace"); + assert_eq!(parsed.behaviour.boundary_actions.len(), 1); + assert_eq!(parsed.operational.captures.len(), 1); + + let climb = &lines[3]; + assert_eq!( + climb["behaviour"]["boundaryActions"][0]["kind"], + "save-slot" + ); + assert_eq!( + climb["behaviour"]["boundaryActions"][0]["slotId"], + flysim::trace::SLOT + ); + assert_eq!( + climb["behaviour"]["boundaryActions"][0]["stateDigest"], + flysim::trace::sha256_hex(b"emulator state") + ); + assert_eq!(climb["operational"]["captures"][0]["afterActions"], 0); + let refused = grafted(climb).expect_err("the legacy archive precedes the slot save"); + assert!( + refused.contains("before this boundary's slot saves"), + "refused for the declared reason: {refused}" + ); +} diff --git a/services/flysim/crates/flysim/tests/integration.rs b/services/flysim/crates/flysim/tests/integration.rs index 2c6c36e..afe2061 100644 --- a/services/flysim/crates/flysim/tests/integration.rs +++ b/services/flysim/crates/flysim/tests/integration.rs @@ -437,6 +437,17 @@ async fn the_service_streams_takes_sugar_checkpoints_and_resumes_after_being_kil assert_eq!(logged[0]["by"], json!("integration-test")); let on_disk = std::fs::read_to_string(dir.path().join("state/events.jsonl")).unwrap(); assert!(on_disk.contains("integration-test fed the fly sugar"), "{on_disk}"); + // And in the hot directory's journal, stamped with the frame it was applied before, which + // is what a shadow run replays (`flysim::journal`). + let journal_path = dir.path().join("hot").join(flysim::journal::FILE_NAME); + let journal = flysim::journal::read(&journal_path).expect("the sugar journal"); + assert_eq!(journal.len(), 1, "{journal:?}"); + assert_eq!(journal[0]["kind"], json!("sugar")); + assert_eq!(journal[0]["eventId"], json!(event_id)); + assert_eq!(journal[0]["durationMs"], json!(400.0)); + let stamped: u64 = + journal[0]["frame"].as_str().and_then(|frame| frame.parse().ok()).expect("a frame"); + assert!(stamped >= 1, "stamped with a frame the emulator has run: {stamped}"); // A second pulse while the first is still being applied is refused, not stacked. let (status, body) = service.post(