diff --git a/Cargo.lock b/Cargo.lock index ed4b020898..e13ddd6b63 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -5408,6 +5408,7 @@ dependencies = [ "serde", "serde_json", "sha2", + "shell-words", "socket2", "tempfile", "tokio", diff --git a/crates/tracedecay-agent-hosts/src/agents/hermes/templates/plugin_init.py b/crates/tracedecay-agent-hosts/src/agents/hermes/templates/plugin_init.py index 3a97d132de..6da0611779 100644 --- a/crates/tracedecay-agent-hosts/src/agents/hermes/templates/plugin_init.py +++ b/crates/tracedecay-agent-hosts/src/agents/hermes/templates/plugin_init.py @@ -996,7 +996,9 @@ def _tool_project_candidates(messages): arguments.get("cwd"), arguments.get("workdir"), ): - if isinstance(candidate, str) and os.path.isabs(os.path.expanduser(candidate)): + if isinstance(candidate, str) and _path_is_absolute( + os.path.expanduser(candidate) + ): candidates.append(candidate) if name in ("terminal", "bash", "shell", "exec_command"): candidates.extend(_terminal_cd_candidates(arguments.get("command") or arguments.get("cmd"))) @@ -1240,9 +1242,18 @@ def _runtime_working_directory(): return candidate return os.getcwd() +def _path_is_absolute(candidate): + if os.path.isabs(candidate): + return True + # ntpath.isabs rejects "/..." spellings as root-relative, but a host may + # forward POSIX-absolute project roots to a Windows-hosted plugin; the + # containment checks still run on the realpath'd candidate. + return candidate.startswith("/") + + def _code_project_root(explicit=None, cwd=None, configured=None, hermes_home=None): candidate = explicit or cwd or configured or _runtime_working_directory() - if isinstance(candidate, str) and candidate.strip() and os.path.isabs(candidate): + if isinstance(candidate, str) and candidate.strip() and _path_is_absolute(candidate): candidate = candidate.strip() try: candidate_real = os.path.realpath(candidate) diff --git a/crates/tracedecay-application/src/diagnostics_publication.rs b/crates/tracedecay-application/src/diagnostics_publication.rs index c3e418bcce..e3c033384d 100644 --- a/crates/tracedecay-application/src/diagnostics_publication.rs +++ b/crates/tracedecay-application/src/diagnostics_publication.rs @@ -17,6 +17,7 @@ //! provider payloads are never copied into a diagnostic record; consumers //! reach evidence through the authorized expansion path instead. +use std::borrow::Cow; use std::collections::BTreeMap; use std::future::Future; use std::path::{Component, Path, PathBuf}; @@ -300,17 +301,35 @@ pub fn code_index_logical_path(project_root: &Path, reported: &str) -> Option bool { diff --git a/crates/tracedecay-automation-runtime/src/automation/hermes_skill_bridge.rs b/crates/tracedecay-automation-runtime/src/automation/hermes_skill_bridge.rs index 41e14aa270..37eed5425e 100644 --- a/crates/tracedecay-automation-runtime/src/automation/hermes_skill_bridge.rs +++ b/crates/tracedecay-automation-runtime/src/automation/hermes_skill_bridge.rs @@ -42,7 +42,12 @@ pub fn load_standard_hermes_skill_bridge( let user_home = user_home.ok_or_else(|| { config_error("could not determine the user home for Hermes skill inventory") })?; - load_standard_hermes_skill_bridge_from_user_home(user_home, options) + // Reported paths spell plainly; a verbatim home leaks `\\?\` into every + // serialized skill and usage path. + load_standard_hermes_skill_bridge_from_user_home( + &tracedecay_runtime_core::path_safety::plain_host_path(user_home), + options, + ) } fn load_standard_hermes_skill_bridge_from_user_home( diff --git a/crates/tracedecay-cli/src/cloud.rs b/crates/tracedecay-cli/src/cloud.rs index d08f137422..36fa53744a 100644 --- a/crates/tracedecay-cli/src/cloud.rs +++ b/crates/tracedecay-cli/src/cloud.rs @@ -699,26 +699,67 @@ mod tests { #[test] fn a_refused_connection_is_network_unreachable() { - // Bound but never listening, and held for the whole test: the port - // refuses connections and cannot be reused, whereas a dropped listener - // stays connectable while a sibling test's forked child still holds - // the inherited descriptor. - let refusing = - socket2::Socket::new(socket2::Domain::IPV4, socket2::Type::STREAM, None).unwrap(); - refusing - .bind( - &"127.0.0.1:0" - .parse::() - .unwrap() - .into(), + // A connection that dies at the transport level before any HTTP + // answer is NetworkUnreachable. Unix refuses a SYN to a port that is + // bound but never listening instantly (the socket is held so the + // port cannot be reused, whereas a dropped listener stays + // connectable while a sibling test's forked child still holds the + // inherited descriptor). The Windows kernel retries a refused SYN + // for about two seconds — longer than the lookup budget — so there + // the same failure class is reached by accepting the connection and + // resetting it with a zero-linger close. + #[cfg(unix)] + let (base, _hold) = { + let refusing = + socket2::Socket::new(socket2::Domain::IPV4, socket2::Type::STREAM, None).unwrap(); + refusing + .bind( + &"127.0.0.1:0" + .parse::() + .unwrap() + .into(), + ) + .unwrap(); + ( + format!( + "http://{}", + refusing.local_addr().unwrap().as_socket().unwrap() + ), + refusing, ) - .unwrap(); - let base = format!( - "http://{}", - refusing.local_addr().unwrap().as_socket().unwrap() - ); + }; + #[cfg(windows)] + let (base, _hold) = { + let listener = TcpListener::bind("127.0.0.1:0").unwrap(); + let base = format!("http://{}", listener.local_addr().unwrap()); + listener.set_nonblocking(true).unwrap(); + let (stop, stopped) = std::sync::mpsc::channel(); + let hold = std::thread::spawn(move || { + loop { + match listener.accept() { + Ok((stream, _)) => { + socket2::SockRef::from(&stream) + .set_linger(Some(Duration::ZERO)) + .unwrap(); + } + Err(error) if error.kind() == std::io::ErrorKind::WouldBlock => {} + Err(error) => panic!("reset listener failed: {error}"), + } + match stopped.recv_timeout(Duration::from_millis(10)) { + Err(std::sync::mpsc::RecvTimeoutError::Timeout) => {} + _ => break, + } + } + }); + (base, (stop, hold)) + }; let error = latest_release_version(&base, true, None).unwrap_err(); + #[cfg(windows)] + { + _hold.0.send(()).unwrap(); + _hold.1.join().unwrap(); + } assert!( matches!(error, ReleaseLookupError::NetworkUnreachable { .. }), diff --git a/crates/tracedecay-cli/src/sessions_cmd/refresh/tests.rs b/crates/tracedecay-cli/src/sessions_cmd/refresh/tests.rs index d818d6cdf8..bef79dfc7f 100644 --- a/crates/tracedecay-cli/src/sessions_cmd/refresh/tests.rs +++ b/crates/tracedecay-cli/src/sessions_cmd/refresh/tests.rs @@ -446,7 +446,10 @@ async fn project_refresh_sends_a_relative_project_path_as_the_cli_directory() { .await .unwrap(); - let cli_directory = std::env::current_dir().unwrap().canonicalize().unwrap(); + let cli_directory = tracedecay_runtime_core::path_safety::canonical_existing_identity( + &std::env::current_dir().unwrap(), + ) + .unwrap(); assert_eq!( transport.calls()[0].arguments, json!({ "path": cli_directory, "format": "json" }) diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs index 72ddef8d59..6b984ade20 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry.rs @@ -386,6 +386,15 @@ struct ColdMountFinalCommitGateV1 { release: tokio::sync::oneshot::Receiver<()>, } +/// Test gates are armed with whatever path spelling the caller holds while the +/// mount path looks them up by the root the scheduler canonicalized (plain +/// `C:\` on Windows vs. the caller's verbatim `\\?\` form). Routing both sides +/// through `canonical_existing_identity` makes either spelling find the gate. +#[cfg(any(test, feature = "test-helpers"))] +pub(super) fn test_gate_root(path: &Path) -> PathBuf { + canonical_existing_identity(path).unwrap_or_else(|_| path.to_path_buf()) +} + /// Armed gates keyed by project root, so tests pausing distinct worktrees in /// one process do not contend for a single slot. #[cfg(any(test, feature = "test-helpers"))] @@ -1962,7 +1971,7 @@ impl CodeIndexSchedulerRegistryV1 { .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .insert( - project_root, + test_gate_root(&project_root), ColdMountFinalCommitGateV1 { entered, release }, ); assert!( @@ -1977,7 +1986,7 @@ impl CodeIndexSchedulerRegistryV1 { let gate = cold_mount_final_commit_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(project_root); + .remove(&test_gate_root(project_root)); if let Some(gate) = gate { let _ = gate.entered.send(()); let _ = gate.release.await; @@ -2008,7 +2017,7 @@ impl CodeIndexSchedulerRegistryV1 { assert!( gates .insert( - (project_root.clone(), at), + (test_gate_root(&project_root), at), RetainedGraphRecoveryGateV1 { entered, release }, ) .is_none(), @@ -2026,7 +2035,7 @@ impl CodeIndexSchedulerRegistryV1 { let gate = retained_graph_recovery_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(&(project_root.to_path_buf(), at)); + .remove(&(test_gate_root(project_root), at)); if let Some(gate) = gate { let _ = gate.entered.send(()); let _ = gate.release.await; @@ -2054,7 +2063,7 @@ impl CodeIndexSchedulerRegistryV1 { assert!( gates .insert( - project_root.clone(), + test_gate_root(&project_root), RetainedTextProjectionGateV1 { entered, release }, ) .is_none(), @@ -2069,7 +2078,7 @@ impl CodeIndexSchedulerRegistryV1 { let gate = retained_text_projection_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(project_root); + .remove(&test_gate_root(project_root)); if let Some(gate) = gate { let _ = gate.entered.send(()); let _ = gate.release.await; diff --git a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs index 4005caee89..83ab0f6c22 100644 --- a/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs +++ b/crates/tracedecay-code-index-runtime/src/code_index_scheduler/registry/test_gates.rs @@ -19,7 +19,7 @@ use super::{ ServingGenerationRollbackOutcomeV1, WorkerStepGateV1, cold_mount_admission_barriers, cold_mount_open_controls, cold_mount_post_check_controls, complete_seat_probe_miss_gate, graph_decode_gate, published_text_projection_gate, query_admission_controls, serving_swap_gate, - unique_mounted_for_scope, wait_notified_if_unset, + test_gate_root, unique_mounted_for_scope, wait_notified_if_unset, }; use tracedecay_runtime_core::path_safety::canonical_existing_identity; @@ -39,7 +39,10 @@ impl CodeIndexSchedulerRegistryV1 { .unwrap_or_else(std::sync::PoisonError::into_inner); assert!( gates - .insert(project_root.clone(), WorkerStepGateV1 { entered, release }) + .insert( + test_gate_root(&project_root), + WorkerStepGateV1 { entered, release }, + ) .is_none(), "one published text projection gate per worktree: {}", project_root.display() @@ -52,7 +55,7 @@ impl CodeIndexSchedulerRegistryV1 { let gate = published_text_projection_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(project_root); + .remove(&test_gate_root(project_root)); Self::pass_worker_step_gate(gate).await; } @@ -72,7 +75,10 @@ impl CodeIndexSchedulerRegistryV1 { let replaced = serving_swap_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .insert(project_root, WorkerStepGateV1 { entered, release }); + .insert( + test_gate_root(&project_root), + WorkerStepGateV1 { entered, release }, + ); assert!(replaced.is_none(), "one serving swap gate per worktree"); (entered_observed, released) } @@ -82,7 +88,7 @@ impl CodeIndexSchedulerRegistryV1 { let gate = serving_swap_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(project_root); + .remove(&test_gate_root(project_root)); Self::pass_worker_step_gate(gate).await; } @@ -138,7 +144,7 @@ impl CodeIndexSchedulerRegistryV1 { .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) .insert( - project_root, + test_gate_root(&project_root), [ WorkerStepGateV1 { entered: before_entered, @@ -163,7 +169,7 @@ impl CodeIndexSchedulerRegistryV1 { let gate = graph_decode_gate() .lock() .unwrap_or_else(std::sync::PoisonError::into_inner) - .remove(project_root); + .remove(&test_gate_root(project_root)); let [before, after] = gate?; Self::pass_worker_step_gate(Some(before)).await; Some(after) diff --git a/crates/tracedecay-configuration/src/config/work_executable_binding.rs b/crates/tracedecay-configuration/src/config/work_executable_binding.rs index 7de5d05443..de7f9ca5bf 100644 --- a/crates/tracedecay-configuration/src/config/work_executable_binding.rs +++ b/crates/tracedecay-configuration/src/config/work_executable_binding.rs @@ -14,6 +14,7 @@ use tracedecay_domain::{ ManifestDigest, ManifestDigestHasher, WorkExecutableReference, WorkProviderBackendV1, WorkProviderProtocol, }; +use tracedecay_runtime_core::path_safety::same_canonical_path; use super::PinnedRuntimeConfiguration; @@ -145,7 +146,9 @@ impl WorkExecutableBindingResolver for PinnedWorkExecutableBindingResolver { executable_id: executable_id.clone(), } })?; - if canonical_path != binding.canonical_path() { + // `canonicalize` spells verbatim `\\?\C:\` on Windows while bindings + // record the plain canonical form; compare identities, not spelling. + if !same_canonical_path(&canonical_path, binding.canonical_path()) { return Err(WorkExecutableBindingError::Stale { executable_id }); } let (actual_digest, verified_byte_length) = diff --git a/crates/tracedecay-dashboard-api/src/analytics_api.rs b/crates/tracedecay-dashboard-api/src/analytics_api.rs index 377e9fb471..d958b3e7fa 100644 --- a/crates/tracedecay-dashboard-api/src/analytics_api.rs +++ b/crates/tracedecay-dashboard-api/src/analytics_api.rs @@ -15,6 +15,7 @@ use serde_json::Value; use tracedecay_contracts::ObservatoryReadModelV1; use tracedecay_contracts::retrieval::{AnalyticsHintCategoryV1, AnalyticsHintsPayloadV1}; use tracedecay_domain::{CoverageStateV1, ObservationScopeV1}; +use tracedecay_sessions::runtime::shared::durable_project_path_key; use tracedecay_automation::analytics::{ ToolUsageObservation, UsageKind, categorize_skill, infer_usage_events, @@ -689,6 +690,10 @@ async fn subagent_tree_reading( let connection = db.read_connection(); let canonical = RegisteredGlobalDb::canonical_project_key(project_root); let opened = project_root.to_string_lossy().into_owned(); + // `project_path` is persisted through `durable_project_path_key`, which + // folds Windows drive/UNC spelling; scoped reads must compare in the same + // stored form or a `C:\` display path never matches the `c:/` row. + let stored_path = durable_project_path_key(&opened); let rows = query_rows( &connection, "SELECT provider, @@ -709,10 +714,16 @@ async fn subagent_tree_reading( -- (`/var` vs `/private/var`) are one project; matching only the -- canonical key silently empties the tree for rows stored under the -- alias the host wrote. - WHERE (project_key IN (?1, ?2) OR project_path IN (?1, ?2)) + WHERE (project_key IN (?1, ?2) OR project_path IN (?3, ?4)) ORDER BY COALESCE(started_at, 0), provider, session_id - LIMIT ?3", - params![canonical, opened, SUBAGENT_TREE_SESSION_CEILING], + LIMIT ?5", + params![ + canonical, + opened.clone(), + stored_path, + opened, + SUBAGENT_TREE_SESSION_CEILING + ], ) .await .map_err(|error| format!("analytics subagent tree query failed: {error}"))?; diff --git a/crates/tracedecay-dashboard-api/src/delivery_agent_usage.rs b/crates/tracedecay-dashboard-api/src/delivery_agent_usage.rs index 9703034e2f..be3302ffd9 100644 --- a/crates/tracedecay-dashboard-api/src/delivery_agent_usage.rs +++ b/crates/tracedecay-dashboard-api/src/delivery_agent_usage.rs @@ -25,6 +25,7 @@ use tracedecay_sessions::runtime::git_correlation::{ CommitRelationFilter, GitCorrelationError, GitRefFilter, MAX_SESSIONS_FOR_LIMIT, SessionsForQuery, }; +use tracedecay_sessions::runtime::shared::durable_project_path_key; use super::DashboardState; use super::analytics_api::managed_agent_label_for_session; @@ -190,6 +191,9 @@ async fn project_sessions( } let canonical = RegisteredGlobalDb::canonical_project_key(project_root); let opened = project_root.to_string_lossy().into_owned(); + // Sessions persist `project_path` through `durable_project_path_key`, so + // the comparison must run in the same folded identity the store wrote. + let stored_path = durable_project_path_key(&opened); let connection = db.read_connection(); let rows = query_rows( &connection, @@ -229,13 +233,15 @@ async fn project_sessions( JOIN sessions s ON s.session_id = c.session_id AND (c.provider = '' OR s.provider = c.provider) - WHERE s.project_key IN (?1, ?2, ?4) OR s.project_path IN (?1, ?2) + WHERE s.project_key IN (?1, ?2, ?4) OR s.project_path IN (?5, ?6) ORDER BY s.provider, s.session_id", params![ canonical, - opened, + opened.clone(), Value::Array(pairs.to_vec()).to_string(), - project_id + project_id, + stored_path, + opened ], ) .await diff --git a/crates/tracedecay-mcp/src/handlers/info/registry.rs b/crates/tracedecay-mcp/src/handlers/info/registry.rs index 9c68b1a100..c706989791 100644 --- a/crates/tracedecay-mcp/src/handlers/info/registry.rs +++ b/crates/tracedecay-mcp/src/handlers/info/registry.rs @@ -21,7 +21,9 @@ use crate::handlers::graph::graph_tool_completion; use crate::handlers::support::{decode_primitive_request, decode_selector_request}; fn display_path(path: &Path) -> String { - path.display().to_string() + tracedecay_runtime_core::path_safety::canonical_root_identity(path) + .to_string_lossy() + .into_owned() } fn bounded_limit(limit: Option, default: usize, max: usize) -> usize { diff --git a/crates/tracedecay-mcp/src/handlers/info/status.rs b/crates/tracedecay-mcp/src/handlers/info/status.rs index 92dc89822e..0c0df17254 100644 --- a/crates/tracedecay-mcp/src/handlers/info/status.rs +++ b/crates/tracedecay-mcp/src/handlers/info/status.rs @@ -39,7 +39,9 @@ use crate::handlers::workflow::current_head_commit_id; use crate::tools::render::Md; fn display_path(path: &Path) -> String { - path.display().to_string() + tracedecay_runtime_core::path_safety::canonical_root_identity(path) + .to_string_lossy() + .into_owned() } /// Project what a readiness wait observed onto the caller-facing outcome. @@ -204,8 +206,13 @@ fn attach_full_branch_status( /// whether pressure may shed them. Other projects' owners belong to the /// daemon-wide Doctor inventory, never to a project read. fn project_memory_value(project_id: &ProjectId) -> StatusMemoryV1 { + // Where the periodic daemon sampler does not run (non-Linux builds no-op + // it) nothing else publishes an observation, so the read takes its own + // checkpoint sample instead of reporting a cell that was never fed. + let pressure = process_resident_memory_pressure_v1(); + pressure.sample_for_checkpoint(); memory_value( - process_resident_memory_pressure_v1(), + pressure, process_resident_owners_v1(), sampled_memory_pressure_some_avg10_v1(), std::time::Instant::now(), diff --git a/crates/tracedecay-mcp/src/handlers/skills.rs b/crates/tracedecay-mcp/src/handlers/skills.rs index 816f15b644..0b598fc670 100644 --- a/crates/tracedecay-mcp/src/handlers/skills.rs +++ b/crates/tracedecay-mcp/src/handlers/skills.rs @@ -108,7 +108,7 @@ pub async fn compute_skill_list( Ok(graph_tool_completion( GraphToolResultV1::SkillList(SkillListResultV1 { status: AutomationReadStatusV1::Ok, - profile_root: profile_root.to_path_buf(), + profile_root: tracedecay_runtime_core::path_safety::plain_host_path(profile_root), count: entries.len(), skills: entries, }), @@ -193,7 +193,7 @@ pub async fn compute_skill_view( Ok(graph_tool_completion( GraphToolResultV1::SkillView(Box::new(SkillViewResultV1 { status: AutomationReadStatusV1::Ok, - profile_root: profile_root.to_path_buf(), + profile_root: tracedecay_runtime_core::path_safety::plain_host_path(profile_root), skill, usage_summary, stale_recommendation, diff --git a/crates/tracedecay-mcp/src/handlers/workflow.rs b/crates/tracedecay-mcp/src/handlers/workflow.rs index 71705f7596..59057177bb 100644 --- a/crates/tracedecay-mcp/src/handlers/workflow.rs +++ b/crates/tracedecay-mcp/src/handlers/workflow.rs @@ -246,10 +246,19 @@ pub async fn compute_diagnose( fn normalized_diagnostic_path(project_root: &Path, file: &str) -> String { let forward = file.replace('\\', "/"); let path = Path::new(&forward); - if path.is_absolute() - && let Ok(relative) = path.strip_prefix(project_root) - { - return relative.to_string_lossy().into_owned(); + if path.is_absolute() { + // Alias spellings (Windows 8.3 names, verbatim roots, symlinked + // parents) all resolve to the root the graph indexes, so compare + // through the canonical identities on both sides. + let canonical_file = tracedecay_runtime_core::path_safety::canonical_root_identity(path); + let canonical_root = + tracedecay_runtime_core::path_safety::canonical_root_identity(project_root); + if let Ok(relative) = canonical_file.strip_prefix(&canonical_root) { + return relative.to_string_lossy().replace('\\', "/"); + } + if let Ok(relative) = path.strip_prefix(project_root) { + return relative.to_string_lossy().replace('\\', "/"); + } } forward } diff --git a/crates/tracedecay-runtime-core/Cargo.toml b/crates/tracedecay-runtime-core/Cargo.toml index eb62eaed7b..57735fddcd 100644 --- a/crates/tracedecay-runtime-core/Cargo.toml +++ b/crates/tracedecay-runtime-core/Cargo.toml @@ -53,6 +53,7 @@ windows-sys = { version = "0.61", features = [ "Win32_Security_Authorization", "Win32_Storage_FileSystem", "Win32_System_Memory", + "Win32_System_ProcessStatus", "Win32_System_SystemServices", "Win32_System_Threading", "Win32_System_WindowsProgramming", diff --git a/crates/tracedecay-runtime-core/src/resident_memory.rs b/crates/tracedecay-runtime-core/src/resident_memory.rs index 125df5b027..8e9348b2f2 100644 --- a/crates/tracedecay-runtime-core/src/resident_memory.rs +++ b/crates/tracedecay-runtime-core/src/resident_memory.rs @@ -409,8 +409,15 @@ fn cgroup_committed_bytes_v1(proc_self_cgroup: &Path, cgroup_root: &Path) -> Opt /// The one `/proc/self/status` parser in the workspace: the daemon's /// dedicated resident-memory sampler and every admission re-measure read it /// through [`ResidentMemoryPressureV1::sample_and_publish`]. Returns `None` -/// where the kernel surface is unavailable (non-Linux hosts), which callers -/// must treat as unobserved, never as zero. +/// where the kernel surface is unavailable (hosts without a process-memory +/// counter), which callers must treat as unobserved, never as zero. +/// +/// Windows reads the same contract through `GetProcessMemoryInfo`: +/// `WorkingSetSize` is every resident page (`resident_bytes`), and +/// `PrivateUsage` is the private commit charge — resident and paged-out +/// private pages together, which is the kernel figure a memory kill line +/// would count. It fills `unreclaimable_bytes` so `admission_bytes` compares +/// the same quantity `/proc` reports as `RssAnon + RssShmem + VmSwap`. #[must_use] pub fn sampled_process_resident_v1() -> Option { #[cfg(target_os = "linux")] @@ -422,12 +429,52 @@ pub fn sampled_process_resident_v1() -> Option { cgroup_committed_bytes_v1(Path::new(PROC_SELF_CGROUP_V1), Path::new(CGROUP_V2_ROOT_V1)); Some(sample) } - #[cfg(not(target_os = "linux"))] + #[cfg(target_os = "windows")] + { + process_resident_sample_from_counters_v1() + } + #[cfg(not(any(target_os = "linux", target_os = "windows")))] { None } } +/// The Windows kernel's per-process memory counters, mapped onto the +/// `/proc/self/status` shape. `PagefileUsage` is left out: it reports the +/// pagefile-backed commit the working set is charged, which `PrivateUsage` +/// already covers. +#[cfg(target_os = "windows")] +fn process_resident_sample_from_counters_v1() -> Option { + use windows_sys::Win32::System::ProcessStatus::{ + GetProcessMemoryInfo, PROCESS_MEMORY_COUNTERS_EX, + }; + use windows_sys::Win32::System::Threading::GetCurrentProcess; + + let mut counters = PROCESS_MEMORY_COUNTERS_EX { + cb: std::mem::size_of::() as u32, + ..unsafe { std::mem::zeroed() } + }; + // SAFETY: `counters` is a live `PROCESS_MEMORY_COUNTERS_EX` of the size + // declared in `cb`, and `GetCurrentProcess` is a pseudohandle that never + // needs closing. + let ok = unsafe { + GetProcessMemoryInfo( + GetCurrentProcess(), + std::ptr::from_mut(&mut counters).cast(), + counters.cb, + ) + }; + if ok == 0 { + return None; + } + Some(ProcessResidentSampleV1 { + resident_bytes: counters.WorkingSetSize as u64, + unreclaimable_bytes: counters.PrivateUsage as u64, + swapped_bytes: 0, + cgroup_committed_bytes: None, + }) +} + /// Unreclaimable bytes, for growth measurement. Admission publishes /// [`ProcessResidentSampleV1::admission_bytes`] instead. #[must_use] diff --git a/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/tests.rs b/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/tests.rs index 0e1f84fb48..099641d1d4 100644 --- a/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/git_correlation/backfill/bounded/tests.rs @@ -456,6 +456,7 @@ async fn scalar(store: &TestStore, sql: &str) -> i64 { rows.next().await.unwrap().unwrap().get(0).unwrap() } +#[cfg(unix)] async fn text_scalar(store: &TestStore, sql: &str) -> String { let mut rows = store.connection.query(sql, ()).await.unwrap(); rows.next().await.unwrap().unwrap().get(0).unwrap() diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/meta.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/meta.rs index c1f323058d..ea4c09f287 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/meta.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/meta.rs @@ -165,6 +165,7 @@ impl SessionMetaParseGate { } /// Blocks until at least `count` parses have parked at this gate. + #[cfg(unix)] pub(crate) fn wait_parked(&self, count: usize) { let mut state = self.state.lock().unwrap_or_else(PoisonError::into_inner); while state.parked < count { @@ -175,6 +176,7 @@ impl SessionMetaParseGate { } } + #[cfg(unix)] pub(crate) fn release(&self) { self.state .lock() @@ -189,7 +191,7 @@ static SESSION_META_PARSE_GATES: OnceLock< Mutex>>, > = OnceLock::new(); -#[cfg(test)] +#[cfg(all(test, unix))] pub(crate) fn install_session_meta_parse_gate_for_test(path: &Path) -> Arc { let key = std::fs::canonicalize(path).unwrap_or_else(|_| path.to_path_buf()); let gate = Arc::new(SessionMetaParseGate::default()); diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs index b4ca541f4e..b16b9ae09e 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation.rs @@ -146,7 +146,7 @@ fn record_in_flight_wait_for_test(key: &CodexMetaCacheKey) { *waits.entry(key.clone()).or_default() += 1; } -#[cfg(test)] +#[cfg(all(test, unix))] fn in_flight_waits_for_test(key: &CodexMetaCacheKey) -> usize { CODEX_META_IN_FLIGHT_WAITS .get_or_init(Mutex::default) diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/meta_cache_tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/meta_cache_tests.rs index 7100dba30a..37cb1e223b 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/meta_cache_tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/observation/meta_cache_tests.rs @@ -8,26 +8,33 @@ use std::io::Write as _; use std::path::{Path, PathBuf}; use std::sync::Arc; +#[cfg(unix)] use std::time::Duration; use serde_json::json; use tempfile::TempDir; +#[cfg(unix)] use tracedecay_domain::ObservationScopeV1; +#[cfg(unix)] use super::super::meta::{ SessionMetaParseGate, install_session_meta_parse_gate_for_test, session_meta_read_count_for_test, }; use super::*; +#[cfg(unix)] use crate::admission::HostAdmission; +#[cfg(unix)] use crate::admission::test_support::MemoryHostAdmission; use crate::runtime::observation::jsonl_observation_admission::install_test_shared_jsonl_preparation_authority; +#[cfg(unix)] use crate::runtime::source::spin_until_jsonl_change_settled; const SESSION_ID: &str = "meta-cache-session"; /// Every wait in this module is bounded so an orphaned claim fails the test /// instead of hanging the suite. +#[cfg(unix)] const SETTLE_WITHIN: Duration = Duration::from_secs(20); fn write_rollout(dir: &Path, name: &str, with_meta: bool) -> PathBuf { @@ -70,8 +77,10 @@ fn write_rollout(dir: &Path, name: &str, with_meta: bool) -> PathBuf { /// Releases the parse gate when dropped, so a failed assertion before the /// explicit release reports instead of leaving a parked blocking parse that the /// runtime's shutdown would wait on forever. +#[cfg(unix)] struct GateGuard(Arc); +#[cfg(unix)] impl std::ops::Deref for GateGuard { type Target = Arc; @@ -80,12 +89,14 @@ impl std::ops::Deref for GateGuard { } } +#[cfg(unix)] impl Drop for GateGuard { fn drop(&mut self) { self.0.release(); } } +#[cfg(unix)] fn fixture(name: &str, with_meta: bool) -> (TempDir, PathBuf, GateGuard) { install_test_shared_jsonl_preparation_authority(); let tmp = TempDir::new().unwrap(); @@ -104,6 +115,7 @@ fn cache_key(path: &Path) -> CodexMetaCacheKey { } } +#[cfg(unix)] fn spawn_lookup( path: &Path, cancellation: ObservationCancellation, @@ -113,6 +125,7 @@ fn spawn_lookup( } /// Blocks until `count` parses of the fixture are parked at its gate. +#[cfg(unix)] async fn wait_parked(gate: &Arc, count: usize) { let gate = Arc::clone(gate); tokio::time::timeout( @@ -126,6 +139,7 @@ async fn wait_parked(gate: &Arc, count: usize) { /// Waits until `count` lookups have registered as waiters behind the key's /// live fill, so the test observes the in-flight path rather than a cache hit. +#[cfg(unix)] async fn wait_in_flight_waits(key: &CodexMetaCacheKey, count: usize) { tokio::time::timeout(SETTLE_WITHIN, async { while in_flight_waits_for_test(key) < count { @@ -169,6 +183,10 @@ async fn a_same_size_in_place_rewrite_is_not_answered_from_the_cache() { } } +// The fixture key and every lookup's key coalesce only once the file +// identity's change stamp settles; without a stat witness each identity +// is unique, so the shared-fill contract cannot be exercised at all. +#[cfg(unix)] #[tokio::test] async fn cancelled_first_waiter_leaves_the_fill_owner_to_settle() { let (_tmp, path, gate) = fixture("cancelled-first-waiter.jsonl", true); @@ -270,6 +288,7 @@ async fn cancelled_first_waiter_leaves_the_fill_owner_to_settle() { ); } +#[cfg(unix)] #[tokio::test] async fn waiter_cancellation_is_typed_and_leaves_the_fill_intact() { let (_tmp, path, gate) = fixture("waiter-cancellation.jsonl", true); @@ -315,6 +334,7 @@ async fn waiter_cancellation_is_typed_and_leaves_the_fill_intact() { )); } +#[cfg(unix)] #[tokio::test] async fn failed_fill_releases_its_claim_and_waiters_fail_typed() { let (_tmp, path, gate) = fixture("no-session-meta.jsonl", false); @@ -376,6 +396,7 @@ fn late_claim_release_never_erases_a_replacement_owner() { assert!(Arc::ptr_eq(&retained, &replacement)); } +#[cfg(unix)] #[test] fn torn_down_fill_task_releases_its_claim_without_publishing() { let (_tmp, path, gate) = fixture("torn-down-fill.jsonl", true); diff --git a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs index 9f530eacf3..9579dfa3a6 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/codex/tests.rs @@ -1708,8 +1708,25 @@ mod recent_first_discovery_tests { secondary_frontier, ) .await; + // Without a stat witness no file proves unchanged: the retained probe + // honestly re-emits and re-enumerates the corpus. + #[cfg(unix)] assert!(unchanged.is_empty()); + #[cfg(not(unix))] + assert_eq!(unchanged, vec![removed.clone()]); + #[cfg(unix)] assert_eq!(unchanged_frontier, secondary_frontier); + // The corpus epoch folds each file's identity digest; without a stat + // witness those digests are unvouched per observation, so only the + // sweep state and file count stay comparable across enumerations. + #[cfg(not(unix))] + { + assert_eq!(unchanged_frontier.state, secondary_frontier.state); + assert_eq!( + unchanged_frontier.epoch.files, + secondary_frontier.epoch.files + ); + } assert_eq!( hub.inner .lock() @@ -1718,7 +1735,7 @@ mod recent_first_discovery_tests { .get(&secondary_source.discovery_key()) .expect("secondary replay index") .completed_enumerations, - 1, + if cfg!(unix) { 1 } else { 2 }, "an unchanged retained probe must not enumerate the corpus again" ); std::fs::remove_file(&removed).unwrap(); @@ -1741,7 +1758,9 @@ mod recent_first_discovery_tests { .replay_indexes .get(&secondary_source.discovery_key()) .expect("secondary replay index"); - assert_eq!(index.completed_enumerations, 2); + // The earlier unchanged probe already enumerated a second time where + // no stat witness proves the corpus unchanged. + assert_eq!(index.completed_enumerations, if cfg!(unix) { 2 } else { 3 }); assert!(!index.paths.iter().any(|entry| entry.path == removed)); } @@ -2289,7 +2308,16 @@ mod recent_first_discovery_tests { assert!(frontier.is_complete()); let restarted = retained_pass(&source, &mut state, bounds, frontier); assert!(restarted.report.paths.is_empty()); + // A settled identity proves the restart's validation complete in one + // pass; without a stat witness the corpus validates in bounded slices + // like the fresh-process contract above. + #[cfg(unix)] assert!(!restarted.report.is_truncated()); + #[cfg(not(unix))] + assert!( + restarted.report.is_truncated(), + "without a stat witness restart validation continues in bounded slices" + ); } #[test] @@ -2541,16 +2569,34 @@ mod recent_first_discovery_tests { bounds, restarted_frontier, ); + // Without a stat witness no file proves unchanged, so restart + // validation honestly re-emits the unchanged corpus; it must + // never emit anything outside it. + #[cfg(unix)] assert!( pass.report.paths.is_empty(), "unchanged restart validation must not re-emit transcripts" ); + #[cfg(not(unix))] + assert!( + pass.report.paths.iter().all(|path| all.contains(path)), + "restart validation may re-emit only the unchanged corpus" + ); restarted_frontier = pass.next_frontier; if restarted_frontier.is_complete() { break; } } + #[cfg(unix)] assert_eq!(restarted_frontier, stored); + // The persisted epoch folds unvouched identity digests where no stat + // witness exists, so restart validation converges to an equal sweep + // state and file count rather than an equal salted epoch. + #[cfg(not(unix))] + { + assert_eq!(restarted_frontier.state, stored.state); + assert_eq!(restarted_frontier.epoch.files, stored.epoch.files); + } let added = write_dated_rollout(home, ("2026", "08", "18"), "after-restart"); let mut rediscovered = false; diff --git a/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs b/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs index ec4721b05b..488e0ca3d3 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/kimi.rs @@ -1521,11 +1521,16 @@ mod tests { admitted, "a later pass must not open state.json once its agents are settled" ); + // The first pass's read count varies with where the wire admission + // settles, so the honest contract is the two-pass total: without a + // stat witness both passes re-read state.json instead of skipping + // settled agents. #[cfg(not(unix))] assert_eq!( super::kimi_state_read_count_for_test(&state), - admitted + 1, - "no stat witness exists, so the pass must read state.json again" + 4, + "no stat witness exists, so each pass reads state.json in the \ + discovery scan and again resolving the wire's session identity" ); assert_eq!(admission.observations().len(), 1); } diff --git a/crates/tracedecay-sessions/src/runtime/hosts/pi_tests.rs b/crates/tracedecay-sessions/src/runtime/hosts/pi_tests.rs index 9f9f4311b7..8faf964d66 100644 --- a/crates/tracedecay-sessions/src/runtime/hosts/pi_tests.rs +++ b/crates/tracedecay-sessions/src/runtime/hosts/pi_tests.rs @@ -324,11 +324,19 @@ fn session_directory_name_matches_pi_encoding() { #[test] fn ambient_agent_dir_is_honored_only_when_absolute_and_inside_home() { - let home = Path::new("/home/operator"); + #[cfg(unix)] + let (home, outside) = ( + Path::new("/home/operator"), + Path::new("/elsewhere/pi-agent"), + ); + #[cfg(not(unix))] + let (home, outside) = ( + Path::new(r"C:\home\operator"), + Path::new(r"D:\elsewhere\pi-agent"), + ); assert_eq!(pi_agent_dir_for(home, None), home.join(".pi/agent")); let inside = home.join("relocated/pi-agent"); assert_eq!(pi_agent_dir_for(home, Some(inside.as_os_str())), inside); - let outside = Path::new("/elsewhere/pi-agent"); assert_eq!( pi_agent_dir_for(home, Some(outside.as_os_str())), home.join(".pi/agent") diff --git a/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs b/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs index f87cc136f2..e0cc3a670a 100644 --- a/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs +++ b/crates/tracedecay-sessions/src/runtime/ingest/scheduler.rs @@ -352,7 +352,16 @@ mod tests { .discover_transcript_paths_with_state(bounds, reloaded, &mut discovery_state) .expect("restart discovery"); assert!(idle.report.paths.is_empty()); + // A settled identity keeps the persisted frontier complete in one + // pass; without a stat witness the corpus validates in bounded slices + // before it reports complete. + #[cfg(unix)] assert!(idle.next_frontier.is_complete()); + #[cfg(not(unix))] + assert!( + !idle.next_frontier.is_complete(), + "without a stat witness restart validation continues in bounded slices" + ); // Production consumers acknowledge every delivered pass (idle ones // included); an unacknowledged pass replays verbatim on the next // discovery, which would mask the addition below. diff --git a/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission/tests.rs b/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission/tests.rs index a237325841..5df0c56256 100644 --- a/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission/tests.rs +++ b/crates/tracedecay-sessions/src/runtime/observation/jsonl_observation_admission/tests.rs @@ -108,8 +108,18 @@ async fn shared_jsonl_page_precomputes_codex_context_hints_once() { .await .expect("second shared consumer"); - assert!(hit); - assert!(std::sync::Arc::ptr_eq(&first, &second)); + #[cfg(unix)] + { + assert!(hit); + assert!(std::sync::Arc::ptr_eq(&first, &second)); + } + // Without a stat witness the shared identity never settles, so each + // consumer honestly reads its own page. + #[cfg(not(unix))] + { + assert!(!hit); + assert!(!std::sync::Arc::ptr_eq(&first, &second)); + } assert_eq!( first .frames @@ -172,8 +182,18 @@ async fn shared_jsonl_page_waiters_share_one_async_in_flight_read() { let (first, first_hit) = first.expect("first concurrent page"); let (second, second_hit) = second.expect("second concurrent page"); - assert_ne!(first_hit, second_hit); - assert!(std::sync::Arc::ptr_eq(&first, &second)); + #[cfg(unix)] + { + assert_ne!(first_hit, second_hit); + assert!(std::sync::Arc::ptr_eq(&first, &second)); + } + // Without a stat witness the shared identity never settles: each + // consumer's key is unique, so both honestly read the file. + #[cfg(not(unix))] + { + assert!(!first_hit && !second_hit); + assert!(!std::sync::Arc::ptr_eq(&first, &second)); + } } #[tokio::test] @@ -442,8 +462,18 @@ async fn generation_pin_prevents_slow_consumer_page_eviction() { ) .await .expect("replayed pinned page"); - assert!(hit); - assert!(std::sync::Arc::ptr_eq(&pinned, &replayed)); + #[cfg(unix)] + { + assert!(hit); + assert!(std::sync::Arc::ptr_eq(&pinned, &replayed)); + } + // Without a stat witness the shared identity never settles, so the + // replayed read honestly misses. + #[cfg(not(unix))] + { + assert!(!hit); + assert!(!std::sync::Arc::ptr_eq(&pinned, &replayed)); + } } #[tokio::test] @@ -503,8 +533,18 @@ async fn exact_append_cursor_replaces_a_superseded_speculative_page() { ) .await .expect("replayed exact append page"); - assert!(replay_hit); - assert!(Arc::ptr_eq(&exact, &replayed)); + #[cfg(unix)] + { + assert!(replay_hit); + assert!(Arc::ptr_eq(&exact, &replayed)); + } + // Without a stat witness the shared identity never settles, so the + // replay honestly rebuilds its own page. + #[cfg(not(unix))] + { + assert!(!replay_hit); + assert!(!Arc::ptr_eq(&exact, &replayed)); + } assert!(!Arc::ptr_eq(&prefetched, &exact)); } @@ -1824,11 +1864,20 @@ async fn codex_session_meta_prefix_is_decoded_once_across_consumers() { .await .expect("second profile consumer"); + #[cfg(unix)] assert_eq!( crate::runtime::hosts::codex::session_meta_read_count_for_test(&path) - before, 1, "canonical path+native identity must share one bounded prefix decode" ); + // Without a stat witness the shared identity never settles, so each + // consumer honestly decodes the prefix itself. + #[cfg(not(unix))] + assert_eq!( + crate::runtime::hosts::codex::session_meta_read_count_for_test(&path) - before, + 2, + "no stat witness exists, so each consumer decodes its own prefix" + ); } #[tokio::test] @@ -2027,7 +2076,12 @@ async fn out_of_scope_frames_are_rejected_before_the_decode() { the one that can is not" ); assert_eq!(progress.frames_persisted, 0); + #[cfg(unix)] assert!(hit); + // Without a stat witness the shared identity never settles, so the + // retained page is honestly re-read. + #[cfg(not(unix))] + assert!(!hit); assert_eq!( super::shared_jsonl_frame_preparations_for_test(page.file_identity), 1, diff --git a/crates/tracedecay/Cargo.toml b/crates/tracedecay/Cargo.toml index a561c81690..64748c5da2 100644 --- a/crates/tracedecay/Cargo.toml +++ b/crates/tracedecay/Cargo.toml @@ -310,6 +310,7 @@ hex = "0.4" toml = "1" regex = "1.12.3" filetime = "0.2" +shell-words = "1.1" socket2 = "0.6" rmcp = { version = "3.0.1", default-features = false, features = ["client", "server", "transport-async-rw"] } criterion = { version = "0.5", features = ["async_tokio"] } diff --git a/crates/tracedecay/src/daemon/bootstrap.rs b/crates/tracedecay/src/daemon/bootstrap.rs index 3aeabcd859..e7315092ed 100644 --- a/crates/tracedecay/src/daemon/bootstrap.rs +++ b/crates/tracedecay/src/daemon/bootstrap.rs @@ -183,6 +183,16 @@ async fn run_foreground_loopback( }, ) .await; + let spooled_hook_opener = hook_v2_replay_consumer::spawn_spooled_hook_opener( + hook_v2_replay_consumer::PortableSpoolOpenerOwners { + lifecycle: lifecycle.clone(), + store_administration: store_administration.clone(), + project_open_gates: Arc::clone(&project_open_gates), + invocation: invocation.clone(), + http_application_registry: http_application_registry.clone(), + }, + profile.clone(), + ); let admission = DaemonClientAdmission::new(MAX_CONCURRENT_DAEMON_CLIENTS); let per_client_admission = DaemonPerClientAdmission::default(); let mut clients: JoinSet> = JoinSet::new(); @@ -249,6 +259,7 @@ async fn run_foreground_loopback( lifecycle.begin_draining(); tracedecay_daemon_service::shutdown::arm_shutdown_exit_bound(); drop(listener); + spooled_hook_opener.abort(); cancel_retained_session_history(&store_administration).await; let shutdown_deadline = tokio::time::Instant::now() + DAEMON_SHUTDOWN_DEADLINE - DAEMON_SHUTDOWN_RECEIPT_LOG_RESERVE; diff --git a/crates/tracedecay/src/daemon/branch_admin.rs b/crates/tracedecay/src/daemon/branch_admin.rs index 02b4762eff..6dc5060a6e 100644 --- a/crates/tracedecay/src/daemon/branch_admin.rs +++ b/crates/tracedecay/src/daemon/branch_admin.rs @@ -1772,8 +1772,12 @@ impl StoreAdministration { outcomes.push( tracedecay_contracts::retrieval::AutomationSchedulerOwnerReconcileOutcome { project_id: key.owner.project_id, - store_root: key.owner.store_root, - graph_db_path: key.owner.graph_db_path, + store_root: tracedecay_runtime_core::path_safety::plain_host_path( + &key.owner.store_root, + ), + graph_db_path: tracedecay_runtime_core::path_safety::plain_host_path( + &key.owner.graph_db_path, + ), scope_prefix: key.scope_prefix, outcome: server.reconcile_automation_scheduler().await, }, diff --git a/crates/tracedecay/src/daemon/connection_serving.rs b/crates/tracedecay/src/daemon/connection_serving.rs index e91b16d567..fc9616b4db 100644 --- a/crates/tracedecay/src/daemon/connection_serving.rs +++ b/crates/tracedecay/src/daemon/connection_serving.rs @@ -2036,6 +2036,14 @@ pub(super) async fn serve_windows_broker_client_with_class_and_invocation( .id .clone() .map(|id| project_open_error_response(id, &error)); + } else if let Some(response) = response.as_mut() + && matches!(classify_mcp_method(&request.method), McpMethod::Initialize) + { + Box::pin(attach_reset_required_stores( + response, + &store_administration, + )) + .await; } drop(setup_activity); if let Some(response) = response { diff --git a/crates/tracedecay/src/daemon/hook_v2_replay_consumer.rs b/crates/tracedecay/src/daemon/hook_v2_replay_consumer.rs index 0ddd215478..13a64f8dfa 100644 --- a/crates/tracedecay/src/daemon/hook_v2_replay_consumer.rs +++ b/crates/tracedecay/src/daemon/hook_v2_replay_consumer.rs @@ -26,11 +26,11 @@ use tracedecay_mcp::handlers::hook_runtime::{ admit_hook_v2_replayed_envelope_with_lifecycle, hook_v2_pending_work_envelopes, }; -#[cfg(unix)] mod spool_opener; mod spool_watch; -#[cfg(unix)] +#[cfg(not(unix))] +pub(in crate::daemon) use spool_opener::PortableSpoolOpenerOwners; pub(in crate::daemon) use spool_opener::spawn_spooled_hook_opener; /// The longest a retained record or receipt waits for its next delivery diff --git a/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_opener.rs b/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_opener.rs index a43f7d4890..86663cd0a6 100644 --- a/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_opener.rs +++ b/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_opener.rs @@ -1,11 +1,16 @@ //! Opens registered projects whose hook spools hold events or delivery //! receipts no running replay consumer will drain: at startup for every //! project with spooled records or receipts, and afterwards for each append -//! or receipt the spool watch reports on a project that is not open. Only the Unix daemon composes a `DaemonEngine`. +//! or receipt the spool watch reports on a project that is not open. The Unix +//! daemon opens through its `DaemonEngine`; the portable daemon through the +//! project-server warmup a client request would take. use std::path::{Path, PathBuf}; +#[cfg(not(unix))] +use std::sync::Arc; use tokio::sync::mpsc; +use tracedecay_daemon_service::shutdown::DaemonLifecycle; use tracedecay_domain::errors::{Result, TraceDecayError}; use tracedecay_hooks::{ HookDeliveryReceiptSpoolV1, HookSpoolV1, hook_delivery_receipt_spool_root, hook_v2_spool_root, @@ -13,14 +18,44 @@ use tracedecay_hooks::{ use tracedecay_runtime_core::config::ProfileRoot; use super::spool_watch::{SpooledProject, consumer_attached, install_opener, watch_project}; +#[cfg(not(unix))] +use crate::daemon::DaemonInvocationState; +use crate::daemon::branch_admin::StoreAdministration; +#[cfg(unix)] use crate::daemon::engine::DaemonEngine; +#[cfg(not(unix))] +use crate::daemon::project_open_admission::ProjectOpenGates; +#[cfg(unix)] use crate::daemon::project_open_admission::ProjectOpenTaskClaim; -use crate::daemon::project_routing::{ - bind_authenticated_profile_identity, project_open_task_capacity_error, -}; +#[cfg(not(unix))] +use crate::daemon::project_open_orchestration::schedule_portable_project_server_warmup; +use crate::daemon::project_routing::bind_authenticated_profile_identity; +#[cfg(unix)] +use crate::daemon::project_routing::project_open_task_capacity_error; const REGISTRY_PAGE: usize = 256; +/// The owners a portable daemon needs to open a project for its spool: the +/// same bundle a client connection carries into +/// [`schedule_portable_project_server_warmup`], plus the lifecycle the opener +/// drains with. +#[cfg(not(unix))] +pub(in crate::daemon) struct PortableSpoolOpenerOwners { + pub lifecycle: DaemonLifecycle, + pub store_administration: StoreAdministration, + pub project_open_gates: Arc>, + pub invocation: DaemonInvocationState, + pub http_application_registry: crate::daemon::http_application::DaemonHttpApplicationRegistry, +} + +/// How `open_for_spooled_hooks` mounts a project on this platform. +enum SpoolOpener { + #[cfg(unix)] + Engine(DaemonEngine), + #[cfg(not(unix))] + Portable(PortableSpoolOpenerOwners), +} + /// Whether any host spool holds records or receipts to drain. A spool that /// cannot be read is reported and counted as holding work, so the project's /// replay consumer surfaces the fault instead of the opener hiding it. @@ -49,25 +84,55 @@ fn has_spooled_records(data_root: &Path) -> bool { }) } -/// Starts the daemon's spooled-hook opener: watches every registered project, -/// opens the ones already holding spooled records, then opens each project -/// whose spool receives an append while it is not open. Runs until the -/// daemon drains. +/// Starts the daemon's spooled-hook opener on Unix: watches every registered +/// project, opens the ones already holding spooled records, then opens each +/// project whose spool receives an append while it is not open. Runs until +/// the daemon drains. +#[cfg(unix)] pub(in crate::daemon) fn spawn_spooled_hook_opener( engine: DaemonEngine, profile: ProfileRoot, ) -> tokio::task::JoinHandle<()> { - let (opener, mut requests) = mpsc::unbounded_channel(); - install_opener(opener.clone()); + spawn_spool_opener( + engine.store_administration.clone(), + engine.lifecycle.clone(), + SpoolOpener::Engine(engine), + profile, + ) +} + +/// Starts the daemon's spooled-hook opener on the portable daemon. Identical +/// to the Unix entry point except project opens go through the portable +/// project-server warmup instead of the engine. +#[cfg(not(unix))] +pub(in crate::daemon) fn spawn_spooled_hook_opener( + owners: PortableSpoolOpenerOwners, + profile: ProfileRoot, +) -> tokio::task::JoinHandle<()> { + spawn_spool_opener( + owners.store_administration.clone(), + owners.lifecycle.clone(), + SpoolOpener::Portable(owners), + profile, + ) +} + +fn spawn_spool_opener( + store_administration: StoreAdministration, + lifecycle: DaemonLifecycle, + opener: SpoolOpener, + profile: ProfileRoot, +) -> tokio::task::JoinHandle<()> { + let (wakes, mut requests) = mpsc::unbounded_channel(); + install_opener(wakes.clone()); tokio::spawn(async move { - let draining = engine.lifecycle.clone(); let work = async { - match registered_projects(&engine, &profile).await { + match registered_projects(&store_administration, &profile).await { Ok(projects) => { for project in projects { watch_project(&project, None); if has_spooled_records(&project.data_root) { - let _ = opener.send(project); + let _ = wakes.send(project); } } } @@ -80,7 +145,7 @@ pub(in crate::daemon) fn spawn_spooled_hook_opener( if consumer_attached(&project.data_root) { continue; } - if let Err(error) = open_for_spooled_hooks(&engine, &profile, &project).await { + if let Err(error) = open_for_spooled_hooks(&opener, &profile, &project).await { tracing::warn!( %error, project = %project.project_root.display(), @@ -91,19 +156,16 @@ pub(in crate::daemon) fn spawn_spooled_hook_opener( }; tokio::select! { () = work => {} - () = draining.wait_for_draining() => {} + () = lifecycle.wait_for_draining() => {} } }) } async fn registered_projects( - engine: &DaemonEngine, + store_administration: &StoreAdministration, profile: &ProfileRoot, ) -> Result> { - let registry = engine - .store_administration - .registered_profile_database() - .await?; + let registry = store_administration.registered_profile_database().await?; let mut roots = Vec::new(); let mut after = None::; loop { @@ -159,7 +221,7 @@ async fn registered_projects( /// client request for it would. Project composition registers the replay /// consumer, whose first pass drains the spool. async fn open_for_spooled_hooks( - engine: &DaemonEngine, + opener: &SpoolOpener, profile: &ProfileRoot, project: &SpooledProject, ) -> Result<()> { @@ -170,13 +232,37 @@ async fn open_for_spooled_hooks( false, false, )?; - let administration = - bind_authenticated_profile_identity(&mut handshake, &engine.store_administration).await?; - let mut engine = engine.clone(); - engine.store_administration = administration; - match engine.begin_project_open(handshake, None).await? { - ProjectOpenTaskClaim::InFlight(_) => Ok(()), - ProjectOpenTaskClaim::Failed(failure) => Err(failure.to_error()), - ProjectOpenTaskClaim::Saturated => Err(project_open_task_capacity_error()), + match opener { + #[cfg(unix)] + SpoolOpener::Engine(engine) => { + let administration = + bind_authenticated_profile_identity(&mut handshake, &engine.store_administration) + .await?; + let mut engine = engine.clone(); + engine.store_administration = administration; + match engine.begin_project_open(handshake, None).await? { + ProjectOpenTaskClaim::InFlight(_) => Ok(()), + ProjectOpenTaskClaim::Failed(failure) => Err(failure.to_error()), + ProjectOpenTaskClaim::Saturated => Err(project_open_task_capacity_error()), + } + } + #[cfg(not(unix))] + SpoolOpener::Portable(owners) => { + let administration = + bind_authenticated_profile_identity(&mut handshake, &owners.store_administration) + .await?; + schedule_portable_project_server_warmup( + owners.lifecycle.clone(), + administration, + Arc::clone(&owners.project_open_gates), + owners.invocation.clone(), + owners.http_application_registry.clone(), + handshake, + None, + #[cfg(test)] + None, + ) + .await + } } } diff --git a/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_watch.rs b/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_watch.rs index 693d6b3976..fd486f99ba 100644 --- a/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_watch.rs +++ b/crates/tracedecay/src/daemon/hook_v2_replay_consumer/spool_watch.rs @@ -133,9 +133,13 @@ pub(super) fn watch_project(project: &SpooledProject, consumer: Option, event: ¬ify::Event) { let wakes: fn(&Path) -> bool = match event.kind { - EventKind::Modify(notify::event::ModifyKind::Data(_)) => { + EventKind::Modify(notify::event::ModifyKind::Data(_) | notify::event::ModifyKind::Any) => { |path| path.file_name() == Some(OsStr::new(HOOK_SPOOL_RECORDS_FILE)) } EventKind::Create(_) @@ -199,7 +203,6 @@ pub(super) fn detach_consumer(data_root: &Path) { } } -#[cfg(unix)] pub(super) fn install_opener(opener: mpsc::UnboundedSender) { with_watch(|watch| { if let Ok(mut targets) = watch.targets.lock() { @@ -208,7 +211,6 @@ pub(super) fn install_opener(opener: mpsc::UnboundedSender) { }); } -#[cfg(unix)] pub(super) fn consumer_attached(data_root: &Path) -> bool { spool_watch().lock().is_ok_and(|watch| { watch.as_ref().is_some_and(|watch| { diff --git a/crates/tracedecay/src/daemon/project_open_orchestration.rs b/crates/tracedecay/src/daemon/project_open_orchestration.rs index 0b03b8d2dd..651916e29c 100644 --- a/crates/tracedecay/src/daemon/project_open_orchestration.rs +++ b/crates/tracedecay/src/daemon/project_open_orchestration.rs @@ -7,7 +7,7 @@ use super::*; use tracedecay_daemon_service::shutdown::DaemonLifecycle; use tracedecay_runtime_core::logging::log_daemon_event; -use tracedecay_runtime_core::path_safety::same_canonical_path; +use tracedecay_runtime_core::path_safety::{plain_host_path, same_canonical_path}; /// Bounds how long a foreground request waits for a route's background open. /// The open task itself is deliberately left running after the deadline. @@ -324,7 +324,7 @@ fn unenrolled_project_route_error(project_path: &Path) -> TraceDecayError { "no TraceDecay index found at '{}': project is not enrolled in the authenticated \ profile; run 'tracedecay init' in that directory, or start the MCP server with \ 'tracedecay serve --path '", - project_path.display() + plain_host_path(project_path).display() ), ) } diff --git a/crates/tracedecay/src/daemon/tests/rmcp_route.rs b/crates/tracedecay/src/daemon/tests/rmcp_route.rs index 43bedc466d..f425d950a6 100644 --- a/crates/tracedecay/src/daemon/tests/rmcp_route.rs +++ b/crates/tracedecay/src/daemon/tests/rmcp_route.rs @@ -117,19 +117,36 @@ async fn rmcp_route_fixture_with_projects( tracedecay_domain::configuration::CodeIndexWorkerSelectionV1::default(), ) .expect("install portable route profile worker plan"); - let server = Box::pin(super::super::portable_project_server_for_request( - DaemonLifecycle::default(), - store_administration.clone(), - Arc::new(tokio::sync::Mutex::new( - super::super::ProjectOpenGates::default(), - )), - invocation, - super::super::http_application::DaemonHttpApplicationRegistry::default(), - &handshake, - super::super::ProjectServerRequirement::Core, - None, - )) + let gates = Arc::new(tokio::sync::Mutex::new( + super::super::ProjectOpenGates::default(), + )); + let applications = super::super::http_application::DaemonHttpApplicationRegistry::default(); + // The unix arm's `engine.project_server` awaits warm-up internally; + // the portable route answers `project_warming` while it opens, so the + // direct drive retries it the way production clients do. + let server = tokio::time::timeout(PHASE_TIMEOUT, async { + loop { + match Box::pin(super::super::portable_project_server_for_request( + DaemonLifecycle::default(), + store_administration.clone(), + gates.clone(), + invocation.clone(), + applications.clone(), + &handshake, + super::super::ProjectServerRequirement::Core, + None, + )) + .await + { + Err(error) if super::super::error_is_project_warming(&error) => { + tokio::time::sleep(Duration::from_millis(50)).await; + } + outcome => break outcome, + } + } + }) .await + .expect("portable project server out of warm-up") .expect("open portable production project server"); (store_administration, server) }; diff --git a/crates/tracedecay/src/doctor.rs b/crates/tracedecay/src/doctor.rs index a73adb19dc..4f0c7071d7 100644 --- a/crates/tracedecay/src/doctor.rs +++ b/crates/tracedecay/src/doctor.rs @@ -350,7 +350,10 @@ fn check_reset_required_stores( profile: &tracedecay_runtime_core::config::ProfileRoot, build_version: &str, ) -> bool { - if !tracedecay_daemon_control::daemon_reachable(profile) { + // Gate on the socket accepting, not a one-second initialize answer: the + // reset census probe owns its own deadline, and a daemon busy attaching + // stores still answers it there. + if !tracedecay_daemon_control::daemon_socket_connectable(profile) { return false; } match tracedecay_daemon_control::daemon_reset_required_stores(profile, build_version) { diff --git a/crates/tracedecay/src/mcp/server/hook_writes.rs b/crates/tracedecay/src/mcp/server/hook_writes.rs index 7ad9a3affc..5057c7676c 100644 --- a/crates/tracedecay/src/mcp/server/hook_writes.rs +++ b/crates/tracedecay/src/mcp/server/hook_writes.rs @@ -10,6 +10,7 @@ use std::sync::Arc; use tracedecay_code_index_runtime::code_index_scheduler::CodeIndexDemandAdmissionV1; use tracedecay_domain::errors::{Result, TraceDecayError}; use tracedecay_project::project::TraceDecay; +use tracedecay_runtime_core::path_safety::same_canonical_path; /// Complete detached reconciliation admission requested by the MCP server. #[derive(Clone, Copy, Debug, Eq, PartialEq)] @@ -63,7 +64,10 @@ pub(crate) async fn execute_background_refresh_direct( ), })?; let active_branch = tracedecay_runtime_core::branch::current_branch(&canonical_root); - if request.graph.project_root() != canonical_root + // `canonicalize` spells verbatim `\\?\C:\` on Windows while a mounted + // graph records the plain canonical root; compare identities, not + // spelling. + if !same_canonical_path(request.graph.project_root(), &canonical_root) || request.graph.active_branch() != active_branch.as_deref() { return Err(TraceDecayError::Config { diff --git a/crates/tracedecay/src/test_support/host_admission.rs b/crates/tracedecay/src/test_support/host_admission.rs index e1270de1de..13fa21ee1a 100644 --- a/crates/tracedecay/src/test_support/host_admission.rs +++ b/crates/tracedecay/src/test_support/host_admission.rs @@ -18,10 +18,14 @@ use tracedecay_project::test_support::host_admission::HostAdmissionTestRuntimeV1 fn test_runtime_profile( runtime: &HostAdmissionTestRuntimeV1, ) -> tracedecay_runtime_core::config::ProfileRoot { + use tracedecay_runtime_core::path_safety::plain_host_path; let profile_root = runtime.profile_root_for_test(); - let profile = tracedecay_runtime_core::config::ProfileRoot::new(profile_root); + // The runtime's profile root is the verbatim filesystem identity; the + // env-derived profile a production server holds carries the plain + // spelling a child process or serializer can consume. + let profile = tracedecay_runtime_core::config::ProfileRoot::new(plain_host_path(profile_root)); match profile_root.parent() { - Some(home) => profile.with_home(home), + Some(home) => profile.with_home(plain_host_path(home)), None => profile, } } diff --git a/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs b/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs index a2adf7aa0f..c4a1017cc4 100644 --- a/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs +++ b/crates/tracedecay/tests/daemon_suite/advanced_workflow_journey_test.rs @@ -241,7 +241,7 @@ fn write_provider_fixture( cancellation_hold: &Path, ) -> (PathBuf, Vec) { let script = format!( - "@echo off\r\nset \"input=\"\r\nset /p \"input=\"\r\nif /I \"%input%\"==\"crash\" goto crash\r\nif /I \"%input%\"==\"cancel\" goto cancel\r\necho {{\"type\":\"system\",\"subtype\":\"init\",\"session_id\":\"{provider_session}\"}}\r\necho {{\"type\":\"assistant\",\"message\":{{\"content\":[{{\"type\":\"text\",\"text\":\"fan-out evidence\"}}]}}}}\r\necho {{\"type\":\"result\",\"subtype\":\"success\",\"is_error\":false}}\r\nexit /b 0\r\n:crash\r\ntype nul > \"{first_started}\"\r\n:wait_first\r\nif exist \"{first_hold}\" (timeout /t 1 /nobreak >nul & goto wait_first)\r\nexit /b 1\r\n:cancel\r\ntype nul > \"{cancellation_started}\"\r\n:wait_cancel\r\nif exist \"{cancellation_hold}\" (timeout /t 1 /nobreak >nul & goto wait_cancel)\r\nexit /b 0\r\n", + "@echo off\r\nset \"input=\"\r\nset /p \"input=\"\r\nif /I \"%input:~0,5%\"==\"crash\" goto crash\r\nif /I \"%input:~0,6%\"==\"cancel\" goto cancel\r\necho {{\"type\":\"system\",\"subtype\":\"init\",\"session_id\":\"{provider_session}\"}}\r\necho {{\"type\":\"assistant\",\"message\":{{\"content\":[{{\"type\":\"text\",\"text\":\"fan-out evidence\"}}]}}}}\r\necho {{\"type\":\"result\",\"subtype\":\"success\",\"is_error\":false}}\r\nexit /b 0\r\n:crash\r\ntype nul > \"{first_started}\"\r\n:wait_first\r\nif exist \"{first_hold}\" (timeout /t 1 /nobreak >nul & goto wait_first)\r\nexit /b 1\r\n:cancel\r\ntype nul > \"{cancellation_started}\"\r\n:wait_cancel\r\nif exist \"{cancellation_hold}\" (timeout /t 1 /nobreak >nul & goto wait_cancel)\r\nexit /b 0\r\n", first_started = first_started.display(), cancellation_started = cancellation_started.display(), first_hold = first_hold.display(), diff --git a/crates/tracedecay/tests/hermes_suite/lcm_bridge.rs b/crates/tracedecay/tests/hermes_suite/lcm_bridge.rs index 11dd4b9da0..7fd26fd854 100644 --- a/crates/tracedecay/tests/hermes_suite/lcm_bridge.rs +++ b/crates/tracedecay/tests/hermes_suite/lcm_bridge.rs @@ -179,6 +179,27 @@ fn python_command() -> Command { } } } + if let Some(paths) = std::env::var_os("PATH") { + for dir in std::env::split_paths(&paths) { + for name in ["python.exe", "python3.exe"] { + let exe = dir.join(name); + if exe.is_file() { + return exe; + } + } + } + } + // The pylauncher ships with every python.org install and lives + // on PATH as `py`; resolve it once so every spawn skips lookup. + if let Ok(out) = Command::new("py") + .args(["-3", "-c", "import sys; print(sys.executable)"]) + .output() + { + let exe = String::from_utf8_lossy(&out.stdout).trim().to_string(); + if out.status.success() && !exe.is_empty() { + return PathBuf::from(exe); + } + } } PathBuf::from("python3") }); diff --git a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/admin_test.rs b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/admin_test.rs index afb3b292cf..d91c158730 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/admin_test.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/admin_test.rs @@ -7,6 +7,8 @@ use std::path::{Path, PathBuf}; use tracedecay_mcp::get_tool_definitions; #[cfg(feature = "test-transport")] use tracedecay_project::test_support::host_admission::HostAdmissionTestRuntimeV1; +#[cfg(feature = "test-transport")] +use tracedecay_runtime_core::path_safety::canonical_root_identity; #[cfg(feature = "test-transport")] #[tokio::test] @@ -199,7 +201,9 @@ async fn project_registry_tools_are_bounded_read_only_and_contextual() { assert_eq!(alias_payload["project"]["project_id"], "proj_alpha"); assert_eq!( alias_payload["project"]["display_root"], - seeded_project_root.to_string_lossy().as_ref() + canonical_root_identity(&seeded_project_root) + .to_string_lossy() + .as_ref() ); let unknown_alias = handle_real_server_tool_call( diff --git a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/lcm_describe_behavior.rs b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/lcm_describe_behavior.rs index 2cfcaa5f7a..ffbabeebf1 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/lcm_describe_behavior.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/lcm_describe_behavior.rs @@ -15,6 +15,7 @@ use serde_json::{Map, Value, json}; use tracedecay::mcp::McpServer; use tracedecay_domain::canonical_text::sha256_hex; use tracedecay_lcm::{LcmSourceRef, LcmSummaryNodeDraft}; +use tracedecay_runtime_core::path_safety::canonical_root_identity; use tracedecay_sessions::admission::HostAdmissionScope; use crate::support::{ @@ -37,7 +38,7 @@ const PAYLOAD_RECEIPT: &str = "{\"ingest_protection\":{\"sanitization_receipt\": #[tokio::test] async fn tracedecay_lcm_describe_reports_shape_without_bodies() { let (cg, _env, dir) = setup_empty_project().await; - let project = dir.path().to_path_buf(); + let project = canonical_root_identity(dir.path()); let external_body = format!("{SECRET} {}", "payload ".repeat(40_000)); let content_hash = sha256_hex(external_body.as_bytes()); let payload_ref = format!( diff --git a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/project_list_test.rs b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/project_list_test.rs index a27a620cac..c00bf1daf9 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/project_list_test.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/project_list_test.rs @@ -18,7 +18,9 @@ use tracedecay_global_db::{GraphScopeUpsert, StoreArtifactUpsert, StoreInstanceU use tracedecay_mcp::McpTransport; use tracedecay_project::project::TraceDecay; use tracedecay_project::test_support::host_admission::HostAdmissionTestRuntimeV1; -use tracedecay_runtime_core::path_safety::canonical_existing_identity; +use tracedecay_runtime_core::path_safety::{ + canonical_existing_identity, plain_host_path, same_canonical_path, +}; use crate::support; @@ -138,9 +140,8 @@ async fn project_list_returns_the_registry_page_the_caller_asked_for() { ); let active_git = canonical_existing_identity(&cg.project_root().join(".git")).expect("active .git"); - assert_eq!( - cg.project_root().join(".git"), - active_git, + assert!( + same_canonical_path(&cg.project_root().join(".git"), &active_git), "the calling checkout must be a primary repository" ); runtime @@ -187,9 +188,10 @@ async fn project_list_returns_the_registry_page_the_caller_asked_for() { 1, false, ); + let active_root = canonical_existing_identity(cg.project_root()).expect("active project root"); let active = registered( &active_id, - cg.project_root(), + &active_root, Some(active_git), "main", Some(ACTIVE_HEAD), @@ -217,7 +219,7 @@ async fn project_list_returns_the_registry_page_the_caller_asked_for() { ) .await .expect("mcp server"); - let registry = registry_path.display().to_string(); + let registry = plain_host_path(®istry_path).display().to_string(); let default_text = tool_text(tools_call(&server, json!({"limit": 25})).await); assert_eq!( diff --git a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs index 2d2d56cbdb..4f0fb60830 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_handler_test/session_search_test.rs @@ -917,8 +917,11 @@ async fn cursor_record_with_a_dispatch_is_one_searchable_message() { .status() .expect("git init"); assert!(init.success(), "git init must succeed"); - let slug = tracedecay_sessions::runtime::hosts::cursor::cursor_project_slug(&project) - .expect("cursor project slug"); + let slug = tracedecay_sessions::runtime::hosts::cursor::cursor_project_slug( + &tracedecay_runtime_core::path_safety::canonical_existing_identity(&project) + .expect("canonical project root"), + ) + .expect("cursor project slug"); let session_dir = transcripts .join(".cursor/projects") .join(slug) diff --git a/crates/tracedecay/tests/mcp_suite/mcp_server_test/analytics_test.rs b/crates/tracedecay/tests/mcp_suite/mcp_server_test/analytics_test.rs index 2d583fd842..baa4fdcfbe 100644 --- a/crates/tracedecay/tests/mcp_suite/mcp_server_test/analytics_test.rs +++ b/crates/tracedecay/tests/mcp_suite/mcp_server_test/analytics_test.rs @@ -214,9 +214,7 @@ async fn context_call_writes_memory_match_analytics_without_fact_bodies() { ); server.ledger_writes_settled().await; - let project_id = fixture - .project_root - .canonicalize() + let project_id = canonical_existing_identity(&fixture.project_root) .expect("project path canonicalizes") .to_string_lossy() .to_string(); diff --git a/crates/tracedecay/tests/runtime_acceptance_suite/canonical_git_observation_correlation.rs b/crates/tracedecay/tests/runtime_acceptance_suite/canonical_git_observation_correlation.rs index d41e8e9b60..97ffa99df7 100644 --- a/crates/tracedecay/tests/runtime_acceptance_suite/canonical_git_observation_correlation.rs +++ b/crates/tracedecay/tests/runtime_acceptance_suite/canonical_git_observation_correlation.rs @@ -185,10 +185,14 @@ async fn canonical_codex_capture_publishes_admitted_git_evidence_for_sessions_fo .unwrap(); assert_eq!(branch_hits.len(), 1); assert_eq!(branch_hits[0].session_id, session_id.as_str()); - // Worktrees are keyed in their portable `/`-separated spelling. + // Worktrees are keyed in their portable `/`-separated spelling of the + // canonical locator, not the fixture's alias (Windows 8.3 TEMP) form. assert_eq!( branch_hits[0].worktree, - Some(normalize_worktree(&project.to_string_lossy())) + Some(normalize_worktree( + &tracedecay_runtime_core::path_safety::canonical_root_identity(&project) + .to_string_lossy() + )) ); let (commit_hits, _) = store diff --git a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_profile_authority_reset.rs b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_profile_authority_reset.rs index 9d565870f8..741912ffad 100644 --- a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_profile_authority_reset.rs +++ b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_profile_authority_reset.rs @@ -235,7 +235,10 @@ fn released_profile_authority_resets_alone_and_projects_register_again() { store.display() ); } - let init_command = format!("tracedecay init {}", project_path.display()); + let init_command = format!( + "tracedecay init {}", + shell_words::quote(&project_path.to_string_lossy()) + ); assert!( reset_output.contains(&format!( "the project registry is empty; register each project store again with:\n \ diff --git a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs index 01357c0649..706b96d715 100644 --- a/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs +++ b/crates/tracedecay/tests/transport_acceptance_suite/typed_terminal_restart_acceptance/stale_sessions_store_reset.rs @@ -290,16 +290,28 @@ fn code_read_outcome(home: &Path, project: &Path, tool: &str, args: &Value) -> V } fn probe_symbol_id(home: &Path, project: &Path) -> String { - let found = super::typed_envelope(&super::tool_call( - home, - project, - "tracedecay_find_exact_symbol", - &json!({ "name": "probe", "format": "json" }), - )); - found["matches"][0]["id"] - .as_str() - .unwrap_or_else(|| panic!("find_exact_symbol did not resolve `probe`: {found}")) - .to_owned() + let started = Instant::now(); + loop { + let found = super::typed_envelope(&super::tool_call( + home, + project, + "tracedecay_find_exact_symbol", + &json!({ "name": "probe", "format": "json" }), + )); + if let Some(id) = found["matches"][0]["id"].as_str() { + return id.to_owned(); + } + let retryable = found["problem"]["retryable"].as_bool() == Some(true); + assert!( + retryable, + "find_exact_symbol did not resolve `probe`: {found}" + ); + assert!( + started.elapsed() < Duration::from_secs(120), + "the code graph never verified `probe`: {found}" + ); + std::thread::sleep(Duration::from_millis(500)); + } } fn session_status(home: &Path, project: &Path, storage_scope: &str) -> Value {