diff --git a/.repository-projection.json b/.repository-projection.json index b69417392..5d3dff09b 100644 --- a/.repository-projection.json +++ b/.repository-projection.json @@ -3,11 +3,11 @@ "projection": "deixic-code", "projectionSchemaVersion": 1, "sourceRepository": "dx-corp/mono", - "sourceSha": "ba4af7b8ee1d610e55c13ddc1f7e2f797b77dd87", + "sourceSha": "0dda4b314d332ad6f95f189d10b066d1a080ee34", "destinationRepository": "dx-corp/code", - "priorProjectedBase": "6d97910afbc6fecb2a1e96c284bb1d2b92a86b9f", + "priorProjectedBase": "8991ebcc446033952c75e0858f46245ab8b0580a", "definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f", "toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04", - "contentDigest": "e2657de302d4b845da675e7d434f3446a33f1c68902dabe7bdc2c9462c9f24ac", + "contentDigest": "93b1fddc5de9161c81eb617fc50c86e6f9bb877a46f1a1206a46d8732b6b39f8", "publicationEligible": true } diff --git a/Cargo.lock b/Cargo.lock index 456ef8889..0a59a903b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -2025,6 +2025,7 @@ dependencies = [ "futures-util", "hex", "jsonschema", + "managed-inference-contract", "serde", "serde_json", "sha2 0.10.9", @@ -4375,6 +4376,18 @@ dependencies = [ "tokio", ] +[[package]] +name = "managed-inference-contract" +version = "0.1.0" +dependencies = [ + "aws-lc-rs", + "base64 0.22.1", + "serde", + "serde_json", + "sha2 0.10.9", + "thiserror 2.0.20", +] + [[package]] name = "matchers" version = "0.2.0" diff --git a/packages/dex-host-rs/src/log.rs b/packages/dex-host-rs/src/log.rs index bdf7ea7c9..251e76487 100644 --- a/packages/dex-host-rs/src/log.rs +++ b/packages/dex-host-rs/src/log.rs @@ -286,6 +286,7 @@ mod tests { Event::UserMessage { turn: TurnId::new("t1"), message_id: None, + model_binding: None, principal: PrincipalId::new("alice"), text: text.into(), attachments: Vec::new(), diff --git a/packages/dex-host-rs/src/model.rs b/packages/dex-host-rs/src/model.rs index 5852189eb..022a1e71d 100644 --- a/packages/dex-host-rs/src/model.rs +++ b/packages/dex-host-rs/src/model.rs @@ -361,6 +361,7 @@ mod tests { dex_loop::Event::UserMessage { turn: dex_loop::TurnId::new("t1"), message_id: None, + model_binding: None, principal: PrincipalId::new("alice"), text: "read a.txt".into(), attachments: Vec::new(), diff --git a/packages/dex-host-rs/src/turn.rs b/packages/dex-host-rs/src/turn.rs index 21d835bbb..c2c253774 100644 --- a/packages/dex-host-rs/src/turn.rs +++ b/packages/dex-host-rs/src/turn.rs @@ -66,6 +66,7 @@ pub async fn run_local_turn( log.append(&[Event::UserMessage { turn: request.turn, message_id: None, + model_binding: None, principal: request.principal.clone(), text: request.text, attachments: Vec::new(), @@ -272,6 +273,7 @@ mod tests { Event::UserMessage { turn: TurnId::new("t1"), message_id: None, + model_binding: None, principal: principal.clone(), text: "Work unattended".into(), attachments: vec![], diff --git a/packages/dex-host-rs/tests/turn.rs b/packages/dex-host-rs/tests/turn.rs index 10ade3289..92032f0e7 100644 --- a/packages/dex-host-rs/tests/turn.rs +++ b/packages/dex-host-rs/tests/turn.rs @@ -85,6 +85,7 @@ async fn read_then_write_then_done_with_no_approval() { log.append(&[Event::UserMessage { turn: TurnId::new("t1"), message_id: None, + model_binding: None, principal: alice(), text: "read notes.txt, then write out.txt".into(), attachments: Vec::new(), diff --git a/packages/local-host-rs/src/tools/bash/shield.rs b/packages/local-host-rs/src/tools/bash/shield.rs index 283449f6d..b5e391705 100644 --- a/packages/local-host-rs/src/tools/bash/shield.rs +++ b/packages/local-host-rs/src/tools/bash/shield.rs @@ -880,8 +880,37 @@ mod tests { ); } + /// The process environment without the repository, index and config + /// overrides `check` refuses, restored on drop. A CI runner or sandbox may + /// inject `GIT_CONFIG_*` (for example proxy settings); the end-to-end bash + /// test reads the real process environment, so it must not inherit them. + struct AmbientGitOverrides(Vec<(String, std::ffi::OsString)>); + + impl AmbientGitOverrides { + fn remove() -> Self { + let removed: Vec<_> = std::env::vars_os() + .filter_map(|(name, value)| Some((name.into_string().ok()?, value))) + .filter(|(name, _)| name.starts_with("GIT_")) + .collect(); + for (name, _) in &removed { + std::env::remove_var(name); + } + Self(removed) + } + } + + impl Drop for AmbientGitOverrides { + fn drop(&mut self) { + for (name, value) in self.0.drain(..) { + std::env::set_var(name, value); + } + } + } + #[tokio::test] async fn bash_tool_blocks_before_creating_the_commit() { + let _lock = crate::config::test_process_env_lock_async().await; + let _ambient = AmbientGitOverrides::remove(); let root = repo(); stage(root.path(), &token()); let tool = super::super::BashTool::new(root.path().display().to_string()); diff --git a/packages/local-host-rs/src/tools/extract_document/ooxml.rs b/packages/local-host-rs/src/tools/extract_document/ooxml.rs index 60fdebf8a..f83f4df33 100644 --- a/packages/local-host-rs/src/tools/extract_document/ooxml.rs +++ b/packages/local-host-rs/src/tools/extract_document/ooxml.rs @@ -67,6 +67,20 @@ fn local_name(name: &str) -> &str { name.rsplit(':').next().unwrap_or(name) } +fn tag_end(src: &str) -> Option { + let mut quote = None; + for (index, byte) in src.bytes().enumerate() { + match quote { + Some(delimiter) if byte == delimiter => quote = None, + Some(_) => {} + None if matches!(byte, b'\'' | b'"') => quote = Some(byte), + None if byte == b'>' => return Some(index), + None => {} + } + } + None +} + impl<'a> Iterator for XmlEvents<'a> { type Item = XmlEvent<'a>; @@ -91,7 +105,7 @@ impl<'a> Iterator for XmlEvents<'a> { self.pos += 9 + len + 3.min(body.len() - len); return Some(XmlEvent::Text(&body[..len])); } - let Some(close) = rest.find('>') else { + let Some(close) = tag_end(rest) else { self.pos = self.src.len(); return None; }; @@ -1430,6 +1444,63 @@ pub(super) mod tests { assert_eq!(column_index("12"), None); } + #[test] + fn xml_tag_end_preserves_quoted_angles_and_following_attributes() { + for xml in [ + r#""#, + r"", + r#""#, + ] { + let mut events = xml_events(xml); + let Some(XmlEvent::Start { + name, + attrs, + self_closing, + }) = events.next() + else { + panic!("sheet start expected: {xml}"); + }; + assert_eq!(name, "sheet"); + assert!(self_closing); + assert_eq!(attr(attrs, "name").as_deref(), Some("Q1 > 2026")); + assert_eq!(attr(attrs, "r:id").as_deref(), Some("rId2")); + assert!(events.next().is_none()); + } + assert!(tag_end(r#""#, + ), + ( + "xl/_rels/workbook.xml.rels", + r#""#, + ), + ( + "xl/worksheets/sheet1.xml", + r#"111"#, + ), + ( + "xl/worksheets/sheet2.xml", + r#"222"#, + ), + ]); + let extracted = xlsx_to_markdown(&bytes).expect("workbook extraction"); + let (first, second) = extracted + .text + .split_once("## Sheet: Overview") + .expect("second sheet"); + assert!(first.starts_with("## Sheet: Q1 > 2026 (hidden)")); + assert!(first.contains("222")); + assert!(!first.contains("111")); + assert!(second.contains("111")); + assert_eq!(extracted.sections.expect("sheet accounting").count, 2); + } + #[test] fn decode_entities_handles_numeric_and_double_escaped_references() { assert_eq!(decode_entities("a &lt; b"), "a < b"); diff --git a/packages/local-host-rs/src/tools/mod.rs b/packages/local-host-rs/src/tools/mod.rs index 8cb92b56d..a19bbbd81 100644 --- a/packages/local-host-rs/src/tools/mod.rs +++ b/packages/local-host-rs/src/tools/mod.rs @@ -98,6 +98,8 @@ pub mod orb_delegation; pub(crate) mod orb_execution; pub mod process_registry; pub(crate) mod process_utils; +#[cfg(unix)] +pub use process_utils::{ProcessGroupGuard, set_new_process_group}; mod registry; mod shell_env; mod status; diff --git a/packages/local-host-rs/src/tools/process_utils.rs b/packages/local-host-rs/src/tools/process_utils.rs index 637acd404..cefbc7d5b 100644 --- a/packages/local-host-rs/src/tools/process_utils.rs +++ b/packages/local-host-rs/src/tools/process_utils.rs @@ -103,8 +103,10 @@ pub(crate) fn kill_process_tree(pid: u32) { let _ = kill_process_tree_tracked(pid); } +/// Establish a dedicated subprocess group and Linux orphan ownership before exec. +/// Containment failures are returned by the command's spawn operation. #[cfg(unix)] -pub(crate) fn set_new_process_group(cmd: &mut tokio::process::Command) { +pub fn set_new_process_group(cmd: &mut tokio::process::Command) { set_std_process_group(cmd.as_std_mut()); } @@ -242,19 +244,102 @@ pub(crate) fn reap_owned_process_groups(groups: &[u32], direct_child: Option) {} +/// Preserve remaining descendant identities after the bounded signal sweep. +/// A dying leader may not orphan its children until after that sweep finishes. +/// This worker only reaps recorded zombies adopted by us; it never signals a +/// process, takes the direct Child's status, or follows a reused PID/group. +#[cfg(target_os = "linux")] +fn finish_owned_process_group_reaping(groups: &[u32], direct_child: Option) { + reap_owned_process_groups(groups, direct_child); + fn identity(pid: u32) -> Option<(char, u32, u32, u64)> { + let stat = std::fs::read_to_string(format!("/proc/{pid}/stat")).ok()?; + let (_, fields) = stat.rsplit_once(") ")?; + let fields: Vec<_> = fields.split_whitespace().collect(); + Some(( + fields.first()?.chars().next()?, + fields.get(1)?.parse().ok()?, + fields.get(2)?.parse().ok()?, + fields.get(19)?.parse().ok()?, + )) + } + let Ok(processes) = std::fs::read_dir("/proc") else { + return; + }; + let mut owned: Vec<_> = processes + .flatten() + .filter_map(|entry| { + let pid: u32 = entry.file_name().to_str()?.parse().ok()?; + if Some(pid) == direct_child || pid <= 1 || i32::try_from(pid).is_err() { + return None; + } + let (_, _, group, started) = identity(pid)?; + groups.contains(&group).then_some((pid, group, started)) + }) + .collect(); + if owned.is_empty() { + return; + } + let owner = std::process::id(); + if let Err(error) = std::thread::Builder::new() + .name("maestro-owned-reaper".into()) + .spawn(move || { + for _ in 0..100 { + owned.retain(|&(pid, group, started)| { + let Some((state, parent, current_group, current_started)) = identity(pid) + else { + return false; + }; + if current_group != group || current_started != started { + return false; + } + if state != 'Z' || parent != owner { + return true; + } + // SAFETY: a recorded, adopted zombie only. Positive PID + // and WNOHANG cannot wait for/reap a live or direct child. + let result = + unsafe { libc::waitpid(pid as i32, std::ptr::null_mut(), libc::WNOHANG) }; + if result == pid as i32 { + return false; + } + if result < 0 { + let error = std::io::Error::last_os_error(); + if error.raw_os_error() == Some(libc::ECHILD) { + return false; + } + if error.raw_os_error() != Some(libc::EINTR) { + tracing::warn!(%error, pid, "cannot reap a late adopted tool child"); + } + } + true + }); + if owned.is_empty() { + break; + } + std::thread::sleep(std::time::Duration::from_millis(10)); + } + }) + { + tracing::warn!(%error, "cannot start late adopted tool child reaping"); + } +} + /// Own a group created before exec, including pipes inherited after leader exit. #[cfg(unix)] -pub(crate) struct ProcessGroupGuard(Option); +pub struct ProcessGroupGuard(Option); #[cfg(unix)] impl ProcessGroupGuard { - pub(crate) fn new(pid: Option) -> Self { + /// Retain ownership of a child spawned in a dedicated process group. + /// The direct child remains the caller's responsibility to wait. + pub fn new(pid: Option) -> Self { Self(pid.filter(|pid| *pid > 1 && i32::try_from(*pid).is_ok())) } pub(crate) fn disarm(&mut self) { self.0 = None; } - pub(crate) fn terminate(&mut self) { + /// Signal the owned group while retaining it for cleanup on drop. + pub fn terminate(&mut self) { if let Some(pid) = self.0 { kill_process_group(pid); } @@ -276,6 +361,8 @@ impl Drop for ProcessGroupGuard { } std::thread::sleep(std::time::Duration::from_millis(1)); } + #[cfg(target_os = "linux")] + finish_owned_process_group_reaping(&[pid], Some(pid)); self.disarm(); } } @@ -374,6 +461,126 @@ mod tests { } } + #[cfg(target_os = "linux")] + fn late_adoption_fixture() -> (std::process::Child, i32, tempfile::TempDir) { + const FIXTURE: &str = "MAESTRO_LATE_ADOPTION_FIXTURE"; + if let Some(path) = std::env::var_os(FIXTURE) { + let path = std::path::PathBuf::from(path); + let child = std::process::Command::new("sleep") + .arg("60") + .spawn() + .unwrap(); + std::fs::write(path.join("pid"), child.id().to_string()).unwrap(); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(5); + while !path.join("release").exists() { + if std::time::Instant::now() >= deadline { + std::process::exit(8); + } + std::thread::sleep(std::time::Duration::from_millis(1)); + } + // Exit without reaping this child, exactly as a killed shell does. + std::process::exit(7); + } + let dir = tempfile::tempdir().unwrap(); + let mut command = std::process::Command::new(std::env::current_exe().unwrap()); + command + .args([ + "--exact", + std::thread::current().name().expect("test name"), + "--nocapture", + ]) + .env(FIXTURE, dir.path()); + set_std_process_group(&mut command); + let child = command.spawn().unwrap(); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2); + let descendant = loop { + if let Some(pid) = std::fs::read_to_string(dir.path().join("pid")) + .ok() + .and_then(|text| text.parse::().ok()) + { + break pid; + } + assert!(std::time::Instant::now() < deadline, "fixture pid missing"); + std::thread::sleep(std::time::Duration::from_millis(1)); + }; + (child, descendant, dir) + } + + #[cfg(target_os = "linux")] + #[test] + fn late_adopted_zombie_is_reaped_without_stealing_direct_child_status() { + let (mut child, descendant, dir) = late_adoption_fixture(); + let root = child.id(); + let guard = ProcessGroupGuard::new(Some(root)); + // SAFETY: only this fixture's recorded descendant is signaled. + assert_eq!(unsafe { libc::kill(descendant, libc::SIGKILL) }, 0); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(2); + loop { + let stat = std::fs::read_to_string(format!("/proc/{descendant}/stat")).unwrap(); + if stat.rsplit_once(") ").is_some_and(|(_, fields)| { + let fields: Vec<_> = fields.split_whitespace().collect(); + fields[0] == "Z" && fields[1] == root.to_string() + }) { + break; + } + assert!( + std::time::Instant::now() < deadline, + "descendant did not exit" + ); + std::thread::sleep(std::time::Duration::from_millis(1)); + } + finish_owned_process_group_reaping(&[root], Some(root)); + // Adopt only after the synchronous sweep has ended. + std::fs::write(dir.path().join("release"), "").unwrap(); + assert_eq!(child.wait().unwrap().code(), Some(7)); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(1); + loop { + // SAFETY: signal 0 probes only this fixture's recorded descendant. + if unsafe { libc::kill(descendant, 0) } == -1 + && std::io::Error::last_os_error().raw_os_error() == Some(libc::ESRCH) + { + break; + } + assert!( + std::time::Instant::now() < deadline, + "late adopted zombie was not reaped" + ); + std::thread::sleep(std::time::Duration::from_millis(5)); + } + drop(guard); + } + + #[cfg(target_os = "linux")] + #[test] + fn deferred_reaping_never_terminates_a_live_descendant() { + let (mut child, descendant, dir) = late_adoption_fixture(); + let root = child.id(); + let guard = ProcessGroupGuard::new(Some(root)); + finish_owned_process_group_reaping(&[root], Some(root)); + std::fs::write(dir.path().join("release"), "").unwrap(); + assert_eq!(child.wait().unwrap().code(), Some(7)); + std::thread::sleep(std::time::Duration::from_millis(50)); + // SAFETY: signal 0 probes only this fixture's recorded descendant. + assert_eq!(unsafe { libc::kill(descendant, 0) }, 0); + let stat = std::fs::read_to_string(format!("/proc/{descendant}/stat")).unwrap(); + assert!(!stat.rsplit_once(") ").unwrap().1.starts_with("Z ")); + drop(guard); + let deadline = std::time::Instant::now() + std::time::Duration::from_secs(1); + loop { + // SAFETY: signal 0 probes only this fixture's recorded descendant. + if unsafe { libc::kill(descendant, 0) } == -1 + && std::io::Error::last_os_error().raw_os_error() == Some(libc::ESRCH) + { + break; + } + assert!( + std::time::Instant::now() < deadline, + "owned live child was not killed and reaped" + ); + std::thread::sleep(std::time::Duration::from_millis(5)); + } + } + #[test] fn terminate_keeps_ownership_until_explicit_disarm() { let mut guard = ProcessGroupGuard::new(None); diff --git a/packages/tui-rs/src/workflow_cli.rs b/packages/tui-rs/src/workflow_cli.rs index 545edcaaf..a25686ad7 100644 --- a/packages/tui-rs/src/workflow_cli.rs +++ b/packages/tui-rs/src/workflow_cli.rs @@ -11,8 +11,6 @@ use std::collections::{BTreeMap, BTreeSet, HashSet}; use std::fs; use std::future::Future; use std::io; -#[cfg(unix)] -use std::os::unix::process::CommandExt; use std::path::{Component, Path, PathBuf}; use std::pin::Pin; use std::process::Stdio; @@ -20,6 +18,8 @@ use std::sync::{Arc, Mutex}; use std::time::Duration; use anyhow::{Context, Result, anyhow, bail}; +#[cfg(unix)] +use maestro_local_host::tools::{ProcessGroupGuard, set_new_process_group}; use maestro_swarm::{ RecoveryDecision, SwarmConfig, SwarmExecutor, SwarmPlan, SwarmRecoveryHooks, SwarmSnapshot, SwarmStatus, SwarmTask, SwarmTaskContext, SwarmTaskOutcome, TaskResult, @@ -1657,7 +1657,8 @@ async fn run_verification( .stdout(Stdio::piped()) .stderr(Stdio::piped()) .kill_on_drop(true); - configure_process_group(&mut command); + #[cfg(unix)] + set_new_process_group(&mut command); let mut child = match command.spawn() { Ok(child) => child, Err(error) => { @@ -1674,7 +1675,8 @@ async fn run_verification( }); } }; - let process_group_id = child.id(); + #[cfg(unix)] + let mut process_group = ProcessGroupGuard::new(child.id()); let stdout = child.stdout.take(); let stderr = child.stderr.take(); let stdout_task = tokio::spawn(read_bounded(stdout, MAX_VERIFIER_OUTPUT_BYTES)); @@ -1684,7 +1686,8 @@ async fn run_verification( biased; () = cancellation.cancelled() => { cancelled = true; - kill_process_group(process_group_id).await; + #[cfg(unix)] + process_group.terminate(); let _ = child.wait().await; None } @@ -1703,11 +1706,13 @@ async fn run_verification( // the group is terminated; the bounded wait still guarantees // that inherited pipes cannot hold acceptance open. tokio::time::sleep(Duration::from_millis(20)).await; - kill_process_group(process_group_id).await; + #[cfg(unix)] + process_group.terminate(); (status?.code(), false) } Some(Err(_)) => { - kill_process_group(process_group_id).await; + #[cfg(unix)] + process_group.terminate(); let _ = child.wait().await; (None, true) } @@ -1779,24 +1784,6 @@ where Ok(output) } -fn configure_process_group(command: &mut TokioCommand) { - #[cfg(unix)] - { - command.as_std_mut().process_group(0); - } -} - -async fn kill_process_group(process_group_id: Option) { - #[cfg(unix)] - if let Some(process_group_id) = process_group_id.and_then(|id| i32::try_from(id).ok()) { - // The child is placed in a fresh process group before spawn. Killing - // the negative PID also terminates explicit verifier descendants. - unsafe { - libc::kill(-process_group_id, libc::SIGKILL); - } - } -} - fn verifier_worktree(cwd: &Path, revision: &str) -> Result> { let root = git_repository_root(cwd)?; let path = std::env::temp_dir().join(format!( diff --git a/packages/tui-rs/src/workflow_cli/acceptance_tests.rs b/packages/tui-rs/src/workflow_cli/acceptance_tests.rs index dd4fa2a95..d3dc94f79 100644 --- a/packages/tui-rs/src/workflow_cli/acceptance_tests.rs +++ b/packages/tui-rs/src/workflow_cli/acceptance_tests.rs @@ -747,6 +747,17 @@ fn wait_for_process_exit(pid: i32) -> bool { #[cfg(unix)] #[tokio::test(flavor = "current_thread")] async fn verifier_leader_exit_drains_background_descendant_pipes() { + #[cfg(target_os = "linux")] + { + // Earlier tool execution makes this process own orphaned descendants. + // Reproduce that production context before starting the verifier. + // SAFETY: this process-wide ownership flag takes integer arguments. + assert_eq!( + unsafe { libc::prctl(libc::PR_SET_CHILD_SUBREAPER, 1, 0, 0, 0) }, + 0, + "the verifier fixture must own its orphaned descendants" + ); + } let fixture = RepoFixture::new("verifier-descendant"); let initial_head = git_revision(&fixture.repo, "HEAD"); let pid_file = fixture._root.path().join("verifier-descendant.pid"); diff --git a/packages/tui-rs/tests/pty_e2e.rs b/packages/tui-rs/tests/pty_e2e.rs index 2b79c164e..20b2ee759 100644 --- a/packages/tui-rs/tests/pty_e2e.rs +++ b/packages/tui-rs/tests/pty_e2e.rs @@ -855,10 +855,153 @@ fn save_continuity_frame(session: &PtySession, name: &str) { } } +trait ContinuitySubmissionTransport { + fn current_frame(&self) -> Option; + fn send_submission(&mut self); + fn now(&self) -> Instant; + fn wait_for_frame(&mut self); + fn failure_text(&self) -> String; +} + +impl ContinuitySubmissionTransport for PtySession { + fn current_frame(&self) -> Option { + self.output + .lock() + .unwrap_or_else(|error| error.into_inner()) + .completed_current_text() + } + + fn send_submission(&mut self) { + self.send_bytes(b"\r"); + } + + fn now(&self) -> Instant { + Instant::now() + } + + fn wait_for_frame(&mut self) { + std::thread::sleep(Duration::from_millis(50)); + } + + fn failure_text(&self) -> String { + self.screen_text() + } +} + +fn await_continuity_submission( + session: &mut impl ContinuitySubmissionTransport, + input_marker: &str, + result: &str, +) { + let deadline = session.now() + TURN_TIMEOUT; + session.send_submission(); + let mut resend_at = session.now() + Duration::from_secs(1); + let mut input_consumed = false; + loop { + if let Some(frame) = session.current_frame() { + if frame.contains(result) { + return; + } + // Enter submits input but also approves a tool. After submission + // consumes input in a completed paint, never press Enter again. + input_consumed |= !frame.contains(input_marker); + if !input_consumed && session.now() >= resend_at { + session.send_submission(); + resend_at = session.now() + Duration::from_secs(1); + } + } + assert!( + session.now() < deadline, + "current terminal frame never contained {result:?}: {}", + session.failure_text() + ); + session.wait_for_frame(); + } +} + fn submit_continuity_prompt(session: &mut PtySession, prompt: &str, result: &str) { let input = format!("\x15{prompt}"); - send_continuity_bytes_until(session, input.as_bytes(), &format!("> {prompt}")); - send_continuity_bytes_until(session, b"\r", result); + let input_marker = format!("> {prompt}"); + send_continuity_bytes_until(session, input.as_bytes(), &input_marker); + await_continuity_submission(session, &input_marker, result); +} + +#[test] +fn continuity_submission_retries_only_unconsumed_input_and_waits_for_current_result() { + struct DelayedApproval { + capture: TerminalCapture, + started: Instant, + elapsed: Duration, + sent: Vec<(Duration, u8)>, + } + + impl ContinuitySubmissionTransport for DelayedApproval { + fn current_frame(&self) -> Option { + self.capture.completed_current_text() + } + + fn send_submission(&mut self) { + self.sent.push((self.elapsed, b'\r')); + } + + fn now(&self) -> Instant { + self.started + self.elapsed + } + + fn wait_for_frame(&mut self) { + self.elapsed += Duration::from_millis(50); + if self.elapsed == Duration::from_millis(500) { + self.capture.process(b"\x1b[?2026h\x1b[2J\x1b[H"); + } + if self.elapsed == Duration::from_millis(750) { + self.capture.process(b"> submit guarded tool\x1b[?2026l"); + } + if self.elapsed == Duration::from_millis(1500) { + self.capture.process(b"\x1b[2J\x1b[HWaiting for model"); + } + if self.elapsed == Duration::from_secs(3) { + self.capture + .process(b"\x1b[2J\x1b[H> submit guarded tool\r\nWaiting for approval"); + } + if self.elapsed == Duration::from_secs(4) { + self.capture + .process(b"\x1b[2J\x1b[HAction Approval Required"); + } + } + + fn failure_text(&self) -> String { + self.capture.text() + } + } + + let mut capture = TerminalCapture::new(4, 80); + // Both an older dialog and the consumed prompt remain in capture history; + // neither may cause completion or an extra key in the current frame. + capture.process(b"Action Approval Required\x1b[2J\x1b[H> submit guarded tool"); + let mut session = DelayedApproval { + capture, + started: Instant::now(), + elapsed: Duration::ZERO, + sent: Vec::new(), + }; + await_continuity_submission( + &mut session, + "> submit guarded tool", + "Action Approval Required", + ); + assert_eq!( + session.sent, + [(Duration::ZERO, b'\r'), (Duration::from_secs(1), b'\r')], + "retry a dropped Enter only while input remains; never approve the delayed tool" + ); + assert_eq!(session.elapsed, Duration::from_secs(4)); + assert!(session.capture.text().contains("> submit guarded tool")); + assert!( + session + .current_frame() + .unwrap() + .contains("Action Approval Required") + ); } fn close_continuity_dialog(session: &mut PtySession, key: &[u8], title: &str) { diff --git a/packages/tui-rs/tests/support/terminal_capture.rs b/packages/tui-rs/tests/support/terminal_capture.rs index cc401bd6d..c556ebbdf 100644 --- a/packages/tui-rs/tests/support/terminal_capture.rs +++ b/packages/tui-rs/tests/support/terminal_capture.rs @@ -35,6 +35,11 @@ impl TerminalCapture { screen_rows(self.parser.screen()).join("\n") } + /// Defer decisions while a synchronized paint is split across PTY reads. + pub(super) fn completed_current_text(&self) -> Option { + (!self.parser.callbacks().sync_in_progress).then(|| self.current_text()) + } + /// Formatted current frame for image evidence, preserving actual cell colors. pub(super) fn current_formatted(&self) -> Vec { self.parser.screen().contents_formatted() @@ -66,6 +71,7 @@ impl TerminalCapture { #[derive(Default)] struct ScreenHistory { snapshots: VecDeque, + sync_in_progress: bool, } impl ScreenHistory { @@ -92,8 +98,15 @@ impl vt100::Callbacks for ScreenHistory { // Preserve completed synchronized frames even when a subsequent // redraw or screen clear arrives in the same read. vt100 delegates // the unsupported DEC synchronized-output mode to this callback. - if first == Some(b'?') && command == 'l' && params.contains(&&[2026][..]) { - self.record(screen_rows(screen).join("\n")); + if first == Some(b'?') && params.contains(&&[2026][..]) { + match command { + 'h' => self.sync_in_progress = true, + 'l' => { + self.sync_in_progress = false; + self.record(screen_rows(screen).join("\n")); + } + _ => {} + } } } } @@ -159,6 +172,28 @@ mod tests { assert!(!capture.parser.screen().contents().contains("Approval")); } + #[test] + fn defers_current_frame_decisions_until_split_synchronized_paint_completes() { + let mut capture = TerminalCapture::new(2, 40); + capture.process(b"\x1b[?2026h> guarded prompt\x1b[?2026l"); + assert!( + capture + .completed_current_text() + .unwrap() + .contains("> guarded prompt") + ); + capture.process(b"\x1b[?2026h\x1b[2J\x1b[H"); + assert!(!capture.current_text().contains("> guarded prompt")); + assert_eq!(capture.completed_current_text(), None); + capture.process(b"> guarded prompt\x1b[?2026l"); + assert!( + capture + .completed_current_text() + .unwrap() + .contains("> guarded prompt") + ); + } + #[test] fn retains_trailing_blank_cells_for_input_prefixes() { let mut capture = TerminalCapture::new(2, 40); diff --git a/vendor/BUILD.bazel b/vendor/BUILD.bazel index 0ad506628..b67339a9e 100644 --- a/vendor/BUILD.bazel +++ b/vendor/BUILD.bazel @@ -1,4 +1,6 @@ exports_files([ + "email-encoding-0.4.2.patch", + "lettre-0.11.23.patch", "zstd-0.13.3.patch", "zstd-safe-7.2.4.patch", "zstd-sys-2.0.16+zstd.1.5.7.patch", diff --git a/vendor/README.md b/vendor/README.md index 676719984..a9ab09e4a 100644 --- a/vendor/README.md +++ b/vendor/README.md @@ -1,4 +1,31 @@ -# Zstandard source repairs +# Reviewed source repairs + +## SMTP and MIME header repairs + +`lettre-0.11.23` and `email-encoding-0.4.2` are maintained source repairs for +the Platform notification transport. `smtp-provenance.json` records their exact +published archive checksums, original file hashes, matching patch hashes, and +the consuming `rust` workspace. Other workspaces retain their existing source +selection. Cargo consumes these directories; Bazel applies the identical +patches to the checksum-pinned archives. Neither unmodified registry release +receives an audit approval, exemption, or new trust import. + +The SMTP reader checks the existing 1,000-byte line and 100,000-byte response +limits while reading, before appending bytes. Both synchronous and asynchronous +paths preserve valid multiline responses and UTF-8 characters split across +reads. Header names reject controls, spaces, non-ASCII bytes and colons in both +constructors. RFC2231 parameter continuations return a formatting error if +they cannot consume input; plain-path length arithmetic is checked. + +Run `python3 scripts/ci/test-vendored-smtp.py` with the normal shared Cargo +target. It creates a temporary consumer workspace rather than a lockfile in +vendored source, validates source provenance, and runs the library and encoding +documentation tests without contacting an SMTP server. The regressions bound +their readers/writers so the original defects fail safely without exhausting +memory or running forever. Source review and remaining dependency audit +limitations are recorded under `rust/supply-chain/reviews/`. + +## Zstandard source repairs These packages preserve the Zstandard format used by existing transcript spools and negotiated HTTP requests. They are locally maintained source, not an audit diff --git a/vendor/dex-loop/Cargo.toml b/vendor/dex-loop/Cargo.toml index 25db88718..9bf05184b 100644 --- a/vendor/dex-loop/Cargo.toml +++ b/vendor/dex-loop/Cargo.toml @@ -11,6 +11,7 @@ futures-util = "0.3" hex = "0.4" # Tool arguments are checked against the tool's schema before a call runs. jsonschema = { version = "0.49", default-features = false } +managed-inference-contract = { path = "crates/managed-inference-contract" } serde = { version = "1.0", features = ["derive"] } serde_json = "1.0" sha2 = "0.10" diff --git a/vendor/dex-loop/src/compaction.rs b/vendor/dex-loop/src/compaction.rs index c7b20e322..f6ef8f723 100644 --- a/vendor/dex-loop/src/compaction.rs +++ b/vendor/dex-loop/src/compaction.rs @@ -272,6 +272,7 @@ mod tests { attachments: vec![ArtifactRef::new("doc@v1")], client_tools: vec![], authorized_tools: vec![], + model_binding: None, approval_mode: ApprovalMode::Interactive, } } @@ -694,6 +695,7 @@ mod cut_tests { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: ApprovalMode::Interactive, } } diff --git a/vendor/dex-loop/src/context.rs b/vendor/dex-loop/src/context.rs index 344f1e8c1..a6994dc8d 100644 --- a/vendor/dex-loop/src/context.rs +++ b/vendor/dex-loop/src/context.rs @@ -197,6 +197,7 @@ pub struct Context { client_tools: Vec, authorized_tools: Vec, approval_mode: ApprovalMode, + model_binding: Option, /// `UserMessage`s for a later turn that arrived (by log cursor) while an /// earlier turn was still running. FIFO: applied one at a time, each when /// the turn ahead of it reaches a terminal status, so the still-running @@ -240,6 +241,7 @@ struct PendingTurn { client_tools: Vec, authorized_tools: Vec, approval_mode: ApprovalMode, + model_binding: Option, } impl Context { @@ -265,6 +267,7 @@ impl Context { client_tools: Vec::new(), authorized_tools: Vec::new(), approval_mode: ApprovalMode::Interactive, + model_binding: None, authorized_principal: None, uncertain_calls: Vec::new(), action_confirmations: Vec::new(), @@ -332,6 +335,11 @@ impl Context { self.approval_mode } + /// The current turn's exact host-resolved provider coordinates. + pub fn model_binding(&self) -> Option<&crate::ManagedInferenceProviderBinding> { + self.model_binding.as_ref() + } + /// Advisory capacity, refreshed before every model stream. It does not /// grant permission or change the engine's caps. pub fn remaining_budget(&self) -> Option { @@ -540,6 +548,7 @@ impl Context { client_tools, authorized_tools, approval_mode, + model_binding, } => { let next = PendingTurn { turn: turn.clone(), @@ -550,6 +559,7 @@ impl Context { client_tools: client_tools.clone(), authorized_tools: authorized_tools.clone(), approval_mode: *approval_mode, + model_binding: model_binding.clone(), }; if self.status == Status::Running { // This message's cursor landed while the current turn was @@ -983,6 +993,7 @@ impl Context { client_tools, authorized_tools, approval_mode, + model_binding, } = next; self.turn = Some(turn.clone()); self.acting = Some(principal.clone()); @@ -998,6 +1009,7 @@ impl Context { self.client_tools = client_tools; self.authorized_tools = authorized_tools; self.approval_mode = approval_mode; + self.model_binding = model_binding; self.authorized_principal = Some(principal.clone()); self.push( cursor, diff --git a/vendor/dex-loop/src/context/evidence_contract.rs b/vendor/dex-loop/src/context/evidence_contract.rs index ab25fa8a0..f26252b41 100644 --- a/vendor/dex-loop/src/context/evidence_contract.rs +++ b/vendor/dex-loop/src/context/evidence_contract.rs @@ -25,6 +25,7 @@ fn typed_owner_evidence_is_bounded_and_survives_compaction_and_full_replay() { attachments: vec![], client_tools: vec![], authorized_tools: vec![], + model_binding: None, approval_mode: ApprovalMode::Interactive, }, ); diff --git a/vendor/dex-loop/src/event.rs b/vendor/dex-loop/src/event.rs index b9739d3c7..c09051428 100644 --- a/vendor/dex-loop/src/event.rs +++ b/vendor/dex-loop/src/event.rs @@ -554,6 +554,11 @@ pub enum Event { /// they were written under. #[serde(default)] approval_mode: ApprovalMode, + /// Exact provider coordinates resolved by the authenticated host. + /// Older turns retain the deployment default. This carries references, + /// never credential values, and grants no Gateway authority. + #[serde(default, skip_serializing_if = "Option::is_none")] + model_binding: Option, }, /// Control: becomes a user message from `principal` before the next model /// call. Calls the model then proposes act under `principal`. @@ -837,6 +842,22 @@ impl Event { mod tests { use super::*; + #[test] + fn accepted_model_binding_survives_event_round_trip() { + let input = serde_json::json!({ + "type": "user_message", "turn": "turn-1", "principal": "user-1", + "text": "hello", "attachments": [], + "model_binding": { + "provider": "vertex-ai", "model": "gemini-selected", + "provider_environment": "production", "credential_name": "gemini-ref", + "team_id": "team-1" + } + }); + let event: Event = serde_json::from_value(input.clone()).expect("accepted message"); + let replay = serde_json::to_value(event).expect("durable event"); + assert_eq!(replay["model_binding"], input["model_binding"]); + } + #[test] fn model_attempt_failed_round_trips_as_snake_case_json() { let event = Event::ModelAttemptFailed { @@ -873,6 +894,7 @@ mod tests { text: "hi".into(), attachments: vec![ArtifactRef::new("a1")], authorized_tools: Vec::new(), + model_binding: None, approval_mode: ApprovalMode::Interactive, client_tools: vec![ClientToolSpec { name: ToolName::new("browser.read_tab"), @@ -984,6 +1006,7 @@ mod tests { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: ApprovalMode::Interactive, }; let json = serde_json::to_string(&event).expect("serialize"); @@ -1012,6 +1035,7 @@ mod tests { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: ApprovalMode::Interactive, } ); diff --git a/vendor/dex-loop/src/lib.rs b/vendor/dex-loop/src/lib.rs index 9845a219b..5cc9510a8 100644 --- a/vendor/dex-loop/src/lib.rs +++ b/vendor/dex-loop/src/lib.rs @@ -44,6 +44,7 @@ pub use event::{ ProviderReasoning, ReceiptId, ServedBy, StepTiming, ThreadId, ToolName, ToolResult, TurnId, Usage, args_digest, }; +pub use managed_inference_contract::ManagedInferenceProviderBinding; pub use ports::{ Claim, Effects, ExecutorKind, Fenced, GovernanceClass, Log, Model, ModelChunk, ModelError, ToolSpec, Tools, Verdict, model_tool_name, diff --git a/vendor/dex-loop/tests/action_confirmation.rs b/vendor/dex-loop/tests/action_confirmation.rs index a2f624577..7b6bf5d0b 100644 --- a/vendor/dex-loop/tests/action_confirmation.rs +++ b/vendor/dex-loop/tests/action_confirmation.rs @@ -203,6 +203,7 @@ fn events(decision: ConfirmationDecision, principal: &str) -> Vec<(Cursor, Event attachments: vec![], client_tools: vec![], authorized_tools: vec![], + model_binding: None, approval_mode: ApprovalMode::Interactive, }, Event::ModelStepCompleted { diff --git a/vendor/dex-loop/tests/authorized_tools.rs b/vendor/dex-loop/tests/authorized_tools.rs index c4b88fc72..0b60eaa97 100644 --- a/vendor/dex-loop/tests/authorized_tools.rs +++ b/vendor/dex-loop/tests/authorized_tools.rs @@ -11,6 +11,7 @@ fn message(turn: &str, principal: &str, tools: &[&str]) -> Event { attachments: vec![], client_tools: vec![], authorized_tools: tools.iter().map(|name| ToolName::new(*name)).collect(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, } } diff --git a/vendor/dex-loop/tests/client_tool_replay.rs b/vendor/dex-loop/tests/client_tool_replay.rs index 294568e73..561c388f5 100644 --- a/vendor/dex-loop/tests/client_tool_replay.rs +++ b/vendor/dex-loop/tests/client_tool_replay.rs @@ -54,6 +54,7 @@ fn log_up_to_model_step(log: &FakeLog, call: &ProposedCall) { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); log.host_append(Event::StepStarted { diff --git a/vendor/dex-loop/tests/model_binding.rs b/vendor/dex-loop/tests/model_binding.rs new file mode 100644 index 000000000..0408b9e16 --- /dev/null +++ b/vendor/dex-loop/tests/model_binding.rs @@ -0,0 +1,69 @@ +use dex_loop::{ + Context, Cursor, Event, ManagedInferenceProviderBinding, ThreadId, TurnId, rehydrate, +}; + +fn message(turn: &str, binding: Option) -> Event { + serde_json::from_value(serde_json::json!({ + "type": "user_message", "turn": turn, "principal": "user-1", + "text": "hello", "attachments": [], "model_binding": binding + })) + .unwrap() +} + +fn binding(provider: &str, model: &str) -> ManagedInferenceProviderBinding { + ManagedInferenceProviderBinding { + provider: provider.into(), + model: model.into(), + provider_environment: "production".into(), + credential_name: format!("{provider}-ref"), + team_id: "team-1".into(), + } +} + +#[test] +fn queued_and_replayed_turns_keep_their_own_binding_and_legacy_turn_resets_it() { + let thread = ThreadId { + org: "org-1".into(), + workspace: "ws-1".into(), + thread: "thread-1".into(), + }; + let first = binding("vertex-ai", "selected-gemini"); + let second = binding("vertex-anthropic", "selected-claude"); + let events = vec![ + (Cursor(1), message("first", Some(first.clone()))), + (Cursor(2), message("second", Some(second.clone()))), + (Cursor(3), message("third", None)), + ( + Cursor(4), + Event::Steer { + principal: dex_loop::PrincipalId::new("user-1"), + text: "model_binding: forged-provider/forged-model".into(), + }, + ), + ]; + let mut live = Context::new(thread.clone()); + for (cursor, event) in &events { + live.observe(*cursor, event); + } + assert_eq!(live.model_binding(), Some(&first)); + let mut replay = rehydrate(thread, &events); + assert_eq!(replay.model_binding(), live.model_binding()); + for context in [&mut live, &mut replay] { + context.observe( + Cursor(5), + &Event::Final { + text: "done".into(), + }, + ); + assert_eq!(context.turn(), Some(&TurnId::new("second"))); + assert_eq!(context.model_binding(), Some(&second)); + context.observe( + Cursor(6), + &Event::Final { + text: "done".into(), + }, + ); + assert_eq!(context.turn(), Some(&TurnId::new("third"))); + assert_eq!(context.model_binding(), None); + } +} diff --git a/vendor/dex-loop/tests/scenarios.rs b/vendor/dex-loop/tests/scenarios.rs index 541981c44..d6ccf2620 100644 --- a/vendor/dex-loop/tests/scenarios.rs +++ b/vendor/dex-loop/tests/scenarios.rs @@ -363,6 +363,7 @@ async fn a_legacy_parked_call_is_granted_on_rehydrate_and_the_step_continues() { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }, Event::StepStarted { @@ -959,6 +960,7 @@ async fn crash_mid_stream_abandons_the_attempt_and_reissues_it() { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); log.host_append(Event::StepStarted { @@ -1177,6 +1179,7 @@ async fn same_turn_id_in_two_threads_dispatches_under_distinct_threads() { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); let model = FakeModel::new(vec![vec![call("update", json!({}))], vec![text("done")]]); @@ -2215,6 +2218,7 @@ async fn a_stale_interrupt_excluded_from_the_rehydrated_suffix_must_not_kill_the attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); // cursor 1 log.host_append(Event::Interrupt { principal: alice() }); // cursor 2 -- excluded from the suffix below @@ -2227,6 +2231,7 @@ async fn a_stale_interrupt_excluded_from_the_rehydrated_suffix_must_not_kill_the attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); // cursor 4 @@ -2271,6 +2276,7 @@ async fn a_steer_from_before_the_rehydrate_point_is_not_carried_into_a_later_tur attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); // cursor 1 -- excluded from the suffix below log.host_append(Event::Steer { @@ -2288,6 +2294,7 @@ async fn a_steer_from_before_the_rehydrate_point_is_not_carried_into_a_later_tur attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); // cursor 4 diff --git a/vendor/dex-loop/tests/sim.rs b/vendor/dex-loop/tests/sim.rs index 0e4eb2919..d923b317a 100644 --- a/vendor/dex-loop/tests/sim.rs +++ b/vendor/dex-loop/tests/sim.rs @@ -194,6 +194,7 @@ fn stale_control_kinds_would_release_a_lease_with_unprocessed_control_event() { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }, ), @@ -306,6 +307,7 @@ async fn dst_two_replicas_racing_the_same_generation_dispatch_once() { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); @@ -365,6 +367,7 @@ async fn dst_two_replica_lease_fencing() { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); diff --git a/vendor/dex-loop/tests/sim/scenario.rs b/vendor/dex-loop/tests/sim/scenario.rs index cc8904689..8be477880 100644 --- a/vendor/dex-loop/tests/sim/scenario.rs +++ b/vendor/dex-loop/tests/sim/scenario.rs @@ -357,6 +357,7 @@ pub async fn run_actions(seed: u64, actions: &[Action]) -> Vec { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); } diff --git a/vendor/dex-loop/tests/support/mod.rs b/vendor/dex-loop/tests/support/mod.rs index 0e1bbb2e9..588b1af51 100644 --- a/vendor/dex-loop/tests/support/mod.rs +++ b/vendor/dex-loop/tests/support/mod.rs @@ -125,6 +125,7 @@ impl FakeLog { attachments: Vec::new(), client_tools, authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }); self.rehydrate() @@ -146,6 +147,7 @@ impl FakeLog { attachments: Vec::new(), client_tools: Vec::new(), authorized_tools: Vec::new(), + model_binding: None, approval_mode, }); self.rehydrate() @@ -852,6 +854,7 @@ pub fn crashed_after_start(log: &FakeLog, call: &ProposedCall) { attachments: vec![], client_tools: vec![], authorized_tools: Vec::new(), + model_binding: None, approval_mode: dex_loop::ApprovalMode::Interactive, }, Event::StepStarted { diff --git a/vendor/dex-loop/tests/tool_deadline.rs b/vendor/dex-loop/tests/tool_deadline.rs index dab8d6aa6..f29f1b99e 100644 --- a/vendor/dex-loop/tests/tool_deadline.rs +++ b/vendor/dex-loop/tests/tool_deadline.rs @@ -393,7 +393,9 @@ async fn exhausted_wall_still_adopts_an_already_started_mutations_ledger_result( assert_eq!(ctx, log.rehydrate()); } -#[tokio::test(flavor = "current_thread")] +// Exercise expiry inside claim/start persistence, rather than host scheduling +// exhausting the wall before the engine can claim the mutation. +#[tokio::test(flavor = "current_thread", start_paused = true)] async fn mutation_claim_or_start_persistence_consuming_wall_never_polls_the_effect() { struct SlowClaim { effects: FakeEffects,