From ae7a2e0ae4e2b2d61857957321aed6bd765ef8d7 Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:24:14 +0800 Subject: [PATCH 01/10] feat(daemon): publish listener before initial model (#336) What: bind the daemon listener before starting refresh and gate cache-backed routes until the first collection attempt completes.\n\nWhy: bootstrap mock data must never be mistaken for an authoritative model, while listener availability must not wait on collection.\n\nHow checked: cargo test -p dockermap-daemon --manifest-path crates/Cargo.toml initializing_cache_is_not_published_even_when_mock_is_allowed --- crates/dockermap-daemon/src/cache_refresh.rs | 30 ++++++---- crates/dockermap-daemon/src/daemon_api.rs | 33 +++++++---- crates/dockermap-daemon/src/main.rs | 62 +++++++++++++------- 3 files changed, 81 insertions(+), 44 deletions(-) diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index df3422a3..bd5992bf 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -83,7 +83,11 @@ impl AppState { #[derive(Clone)] pub(crate) struct DaemonCache { - pub(crate) snapshot: DockerSnapshot, +/// Cache-backed HTTP routes remain unavailable until the first Docker +/// collection attempt has completed. This prevents the internal bootstrap +/// mock cache from being mistaken for an authoritative publication. +pub(crate) publication_ready: bool, +pub(crate) snapshot: DockerSnapshot, pub(crate) health: HealthResponse, pub(crate) runtime_map: RuntimeMap, pub(crate) findings: FindingsResponse, @@ -606,8 +610,9 @@ impl DaemonCache { redact_health_response(&mut health); let last_updated = snapshot.last_updated; - let mut cache = Self { - snapshot, +let mut cache = Self { +publication_ready: false, +snapshot, health, runtime_map: RuntimeMap { nodes: Vec::new(), @@ -774,18 +779,21 @@ async fn publish_docker_snapshot_cache( state: &AppState, mut updated: DaemonCache, ) -> (DockerSnapshot, RuntimeMode, u64) { - let mut cache = state.cache.write().await; +let mut cache = state.cache.write().await; // A mock fallback is a distinct source of bytes. Do not retain live host // observations and relabel them as sample data (or vice versa). let same_source = cache.health.mode == updated.health.mode; - updated.source_generation = if same_source { +updated.source_generation = if same_source { cache.source_generation } else { cache .source_generation .checked_add(1) .expect("source generation overflow") - }; +}; +// A completed collection attempt, including an explicit policy-permitted +// mock fallback, makes this coherent cache eligible for publication. +updated.publication_ready = true; updated.runtime_providers = if same_source { cache.runtime_providers.clone() } else { @@ -936,8 +944,9 @@ where model_revision: String::new(), message: Some("Docker engine connected".into()), }; - Ok(DaemonCache { - snapshot, +Ok(DaemonCache { +publication_ready: false, +snapshot, health, runtime_map: empty_runtime_map(0), findings: FindingsResponse::default(), @@ -3076,8 +3085,9 @@ mod scheduler_tests { fn docker_cache(snapshot: DockerSnapshot) -> DaemonCache { let last_updated = snapshot.last_updated; - let mut cache = DaemonCache { - snapshot, +let mut cache = DaemonCache { +publication_ready: true, +snapshot, health: HealthResponse { status: HealthState::Ok, mode: RuntimeMode::Docker, diff --git a/crates/dockermap-daemon/src/daemon_api.rs b/crates/dockermap-daemon/src/daemon_api.rs index 4b7bc83a..65c5406f 100644 --- a/crates/dockermap-daemon/src/daemon_api.rs +++ b/crates/dockermap-daemon/src/daemon_api.rs @@ -84,8 +84,14 @@ pub(crate) fn daemon_router(state: AppState, daemon_token: DaemonAuthToken) -> R } async fn publication_cache(state: &AppState) -> Result, ApiError> { - let cache = state.cache.read().await; - if !state.allow_mock && cache.health.mode == RuntimeMode::Mock { +let cache = state.cache.read().await; +if !cache.publication_ready { +return Err(ApiError { +status: StatusCode::SERVICE_UNAVAILABLE, +message: "Docker model is still initializing".into(), +}); +} +if !state.allow_mock && cache.health.mode == RuntimeMode::Mock { return Err(ApiError { status: StatusCode::SERVICE_UNAVAILABLE, message: "Live Docker authority is unavailable and mock publication is disabled".into(), @@ -420,17 +426,18 @@ mod tests { use super::*; #[tokio::test] - async fn findings_route_stamps_the_actual_cache_mode_after_live_and_mock_resets() { - let state = AppState::new(true); - assert_eq!( - get_findings(State(state.clone())).await.unwrap().0.source, - Some(RuntimeMode::Mock), - "initial unavailable/fallback data must be visibly non-live" - ); - - { - let mut cache = state.cache.write().await; - cache.health.mode = RuntimeMode::Docker; +async fn findings_route_stamps_the_actual_cache_mode_after_live_and_mock_resets() { +let state = AppState::new(true); +assert_eq!( +get_findings(State(state.clone())).await.unwrap_err().status, +StatusCode::SERVICE_UNAVAILABLE, +"bootstrap mock data must not be published before the first attempt" +); + +{ +let mut cache = state.cache.write().await; +cache.publication_ready = true; +cache.health.mode = RuntimeMode::Docker; // A source marker is publication data, never cache data. This // deliberately stale value models a completed live publication // followed by an eventual reset. diff --git a/crates/dockermap-daemon/src/main.rs b/crates/dockermap-daemon/src/main.rs index a309054c..f4acd427 100644 --- a/crates/dockermap-daemon/src/main.rs +++ b/crates/dockermap-daemon/src/main.rs @@ -19,7 +19,7 @@ use bollard::Docker; pub(crate) use cache_refresh::AppState; #[cfg(test)] use cache_refresh::DaemonCache; -use cache_refresh::{refresh_cache, refresh_loop}; +use cache_refresh::refresh_loop; use compose_api::run_cli; use config::{ read_allow_mock_env, read_bind_host_env, read_daemon_token_env, read_port_env, DaemonAuthToken, @@ -135,16 +135,16 @@ async fn main() { let daemon_token = read_daemon_token_env(); let port = read_port_env("DOCKERMAP_DAEMON_PORT", 4100); let host = read_bind_host_env("DOCKERMAP_DAEMON_HOST", daemon_token.0.is_some()); - let address = SocketAddr::from((host, port)); - let state = AppState::new(read_allow_mock_env()); +let address = SocketAddr::from((host, port)); +let state = AppState::new(read_allow_mock_env()); +let app = daemon_router(state.clone(), daemon_token); +let listener = TcpListener::bind(address) +.await +.expect("daemon listener should bind"); - refresh_cache(&state).await; - tokio::spawn(refresh_loop(state.clone())); - - let app = daemon_router(state, daemon_token); - let listener = TcpListener::bind(address) - .await - .expect("daemon listener should bind"); +// The listener is available while the initial authoritative Docker model is +// collected. Cache-backed routes truthfully return 503 until publication. +tokio::spawn(refresh_loop(state)); println!("dockermap-daemon listening on http://{address}"); @@ -183,17 +183,19 @@ mod tests { }; use tower::util::ServiceExt; - fn test_daemon_state() -> AppState { - AppState { - allow_mock: true, - cache: Arc::new(RwLock::new(DaemonCache::mock())), +fn test_daemon_state() -> AppState { +let mut cache = DaemonCache::mock(); +cache.publication_ready = true; +AppState { +allow_mock: true, +cache: Arc::new(RwLock::new(cache)), docker: Arc::new(RwLock::new(None)), provider_slot_in_flight: Arc::new(crate::cache_refresh::ProviderSlotFlights::default()), } } - #[tokio::test] - async fn daemon_mock_cache_is_not_published_when_mock_mode_is_disabled() { +#[tokio::test] +async fn daemon_mock_cache_is_not_published_when_mock_mode_is_disabled() { let state = AppState { allow_mock: false, cache: Arc::new(RwLock::new(DaemonCache::mock())), @@ -226,10 +228,27 @@ mod tests { .expect("daemon router should respond"); assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE, "{path}"); } - } +} - #[tokio::test] - async fn daemon_bearer_boundary_allows_only_the_exact_configured_token() { +#[tokio::test] +async fn initializing_cache_is_not_published_even_when_mock_is_allowed() { +let state = AppState::new(true); +for path in ["/daemon/health", "/daemon/snapshot", "/daemon/runtime/map", "/daemon/findings"] { +let response = daemon_router(state.clone(), DaemonAuthToken(None)) +.oneshot( +Request::builder() +.uri(path) +.body(axum::body::Body::empty()) +.expect("request should build"), +) +.await +.expect("daemon router should respond"); +assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE, "{path}"); +} +} + +#[tokio::test] +async fn daemon_bearer_boundary_allows_only_the_exact_configured_token() { let allowed = daemon_router(test_daemon_state(), DaemonAuthToken(None)) .oneshot( Request::builder() @@ -352,8 +371,9 @@ mod tests { .expect("Docker stub response should be written"); }); - let mut cache = DaemonCache::mock(); - cache.health.docker_reachable = true; +let mut cache = DaemonCache::mock(); +cache.publication_ready = true; +cache.health.docker_reachable = true; let state = AppState { allow_mock: true, cache: Arc::new(RwLock::new(cache)), From 9f999b1fbec4bb727d9684fc806bca24534d94f6 Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:31:27 +0800 Subject: [PATCH 02/10] feat(daemon): collect Docker inventory concurrently (#336) What: issue the fixed container, network, and volume inventory reads concurrently and cover their synchronized gateway arrival.\n\nWhy: independent Docker inventory latency must not extend the authoritative observation path.\n\nHow checked: npm run fmt:rust; cargo test -p dockermap-daemon --manifest-path crates/Cargo.toml snapshot_timeout_invalidates_stalled_client_and_fresh_client_recovers. --- crates/dockermap-daemon/src/cache_refresh.rs | 131 ++++++++------ .../dockermap-daemon/src/docker_collector.rs | 29 ++-- crates/dockermap-daemon/src/main.rs | 164 +++++++++++++----- 3 files changed, 211 insertions(+), 113 deletions(-) diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index bd5992bf..ca3335da 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -83,11 +83,11 @@ impl AppState { #[derive(Clone)] pub(crate) struct DaemonCache { -/// Cache-backed HTTP routes remain unavailable until the first Docker -/// collection attempt has completed. This prevents the internal bootstrap -/// mock cache from being mistaken for an authoritative publication. -pub(crate) publication_ready: bool, -pub(crate) snapshot: DockerSnapshot, + /// Cache-backed HTTP routes remain unavailable until the first Docker + /// collection attempt has completed. This prevents the internal bootstrap + /// mock cache from being mistaken for an authoritative publication. + pub(crate) publication_ready: bool, + pub(crate) snapshot: DockerSnapshot, pub(crate) health: HealthResponse, pub(crate) runtime_map: RuntimeMap, pub(crate) findings: FindingsResponse, @@ -610,9 +610,9 @@ impl DaemonCache { redact_health_response(&mut health); let last_updated = snapshot.last_updated; -let mut cache = Self { -publication_ready: false, -snapshot, + let mut cache = Self { + publication_ready: false, + snapshot, health, runtime_map: RuntimeMap { nodes: Vec::new(), @@ -779,21 +779,21 @@ async fn publish_docker_snapshot_cache( state: &AppState, mut updated: DaemonCache, ) -> (DockerSnapshot, RuntimeMode, u64) { -let mut cache = state.cache.write().await; + let mut cache = state.cache.write().await; // A mock fallback is a distinct source of bytes. Do not retain live host // observations and relabel them as sample data (or vice versa). let same_source = cache.health.mode == updated.health.mode; -updated.source_generation = if same_source { + updated.source_generation = if same_source { cache.source_generation } else { cache .source_generation .checked_add(1) .expect("source generation overflow") -}; -// A completed collection attempt, including an explicit policy-permitted -// mock fallback, makes this coherent cache eligible for publication. -updated.publication_ready = true; + }; + // A completed collection attempt, including an explicit policy-permitted + // mock fallback, makes this coherent cache eligible for publication. + updated.publication_ready = true; updated.runtime_providers = if same_source { cache.runtime_providers.clone() } else { @@ -944,9 +944,9 @@ where model_revision: String::new(), message: Some("Docker engine connected".into()), }; -Ok(DaemonCache { -publication_ready: false, -snapshot, + Ok(DaemonCache { + publication_ready: false, + snapshot, health, runtime_map: empty_runtime_map(0), findings: FindingsResponse::default(), @@ -3085,9 +3085,9 @@ mod scheduler_tests { fn docker_cache(snapshot: DockerSnapshot) -> DaemonCache { let last_updated = snapshot.last_updated; -let mut cache = DaemonCache { -publication_ready: true, -snapshot, + let mut cache = DaemonCache { + publication_ready: true, + snapshot, health: HealthResponse { status: HealthState::Ok, mode: RuntimeMode::Docker, @@ -3141,48 +3141,62 @@ snapshot, let listener = UnixListener::bind(&socket).expect("gateway stub should bind"); let (stalled_tx, mut stalled_rx) = mpsc::unbounded_channel(); let gateway = tokio::spawn(async move { - let (mut stalled, _) = listener - .accept() - .await - .expect("stalled snapshot request should arrive"); - let first_target = read_request_head(&mut stalled).await; - stalled - .write_all( - b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ntransfer-encoding: chunked\r\n\r\n", - ) - .await - .expect("gateway should start the stalled response"); + let mut stalled = Vec::new(); + let mut initial_targets = Vec::new(); + for _ in 0..3 { + let (mut connection, _) = listener + .accept() + .await + .expect("stalled snapshot request should arrive"); + let target = read_request_head(&mut connection).await; + connection + .write_all( + b"HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ntransfer-encoding: chunked\r\n\r\n", + ) + .await + .expect("gateway should start the stalled response"); + initial_targets.push(target); + stalled.push(connection); + } + initial_targets.sort(); stalled_tx - .send(first_target.clone()) + .send(initial_targets) .expect("test should observe the stalled request"); - let mut targets = vec![first_target]; + let mut responses = Vec::new(); for _ in 0..3 { let (mut connection, _) = listener .accept() .await .expect("fresh snapshot request should arrive"); - let target = read_request_head(&mut connection).await; - let body = if target.contains("/containers/json") || target.contains("/networks") { - "[]" - } else if target.contains("/volumes") { - r#"{"Volumes":[],"Warnings":null}"# - } else { - panic!("unexpected snapshot target: {target}"); - }; - let response = format!( - "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}", - body.len(), - body - ); - connection - .write_all(response.as_bytes()) - .await - .expect("fresh snapshot response should be written"); - targets.push(target); + responses.push(tokio::spawn(async move { + let target = read_request_head(&mut connection).await; + let body = if target.contains("/containers/json") || target.contains("/networks") { + "[]" + } else if target.contains("/volumes") { + r#"{"Volumes":[],"Warnings":null}"# + } else { + panic!("unexpected snapshot target: {target}"); + }; + let response = format!( + "HTTP/1.1 200 OK\r\ncontent-type: application/json\r\ncontent-length: {}\r\nconnection: close\r\n\r\n{}", + body.len(), + body + ); + connection + .write_all(response.as_bytes()) + .await + .expect("fresh snapshot response should be written"); + target + })); } - // The incomplete response remains owned until replacement-client - // collection finishes, proving recovery does not reuse it. + let mut targets = Vec::new(); + for response in responses { + targets.push(response.await.expect("fresh response task should finish")); + } + targets.sort(); + // The incomplete responses remain owned until replacement-client + // collection finishes, proving recovery does not reuse them. drop(stalled); targets }); @@ -3204,8 +3218,14 @@ snapshot, assert!(started.elapsed() >= test_timeout); assert!(started.elapsed() < Duration::from_secs(1)); assert_eq!( - stalled_rx.try_recv().expect("gateway accepted the request"), - "GET /containers/json?all=true&size=false HTTP/1.1" + stalled_rx + .try_recv() + .expect("gateway accepted the requests"), + vec![ + "GET /containers/json?all=true&size=false HTTP/1.1", + "GET /networks? HTTP/1.1", + "GET /volumes? HTTP/1.1", + ] ); assert!( state.docker.read().await.is_none(), @@ -3249,7 +3269,6 @@ snapshot, assert_eq!( gateway.await.expect("gateway stub should finish"), vec![ - "GET /containers/json?all=true&size=false HTTP/1.1", "GET /containers/json?all=true&size=false HTTP/1.1", "GET /networks? HTTP/1.1", "GET /volumes? HTTP/1.1", diff --git a/crates/dockermap-daemon/src/docker_collector.rs b/crates/dockermap-daemon/src/docker_collector.rs index 6773a0ab..fde84b30 100644 --- a/crates/dockermap-daemon/src/docker_collector.rs +++ b/crates/dockermap-daemon/src/docker_collector.rs @@ -67,31 +67,26 @@ impl DockerCollector { if let Some(filters) = filters.as_ref() { container_options = container_options.filters(filters); } - let containers = self - .client - .list_containers(Some(container_options.build())) - .await - .map_err(|error| format!("list_containers failed: {error}"))?; - let mut network_options = ListNetworksOptionsBuilder::new(); if let Some(filters) = filters.as_ref() { network_options = network_options.filters(filters); } - let networks = self - .client - .list_networks(Some(network_options.build())) - .await - .map_err(|error| format!("list_networks failed: {error}"))?; - let mut volume_options = ListVolumesOptionsBuilder::new(); if let Some(filters) = filters.as_ref() { volume_options = volume_options.filters(filters); } - let volumes = self - .client - .list_volumes(Some(volume_options.build())) - .await - .map_err(|error| format!("list_volumes failed: {error}"))?; + // Bollard's client is cloneable and these inventory endpoints borrow it + // immutably, so one bounded caller deadline covers three concurrent, + // fixed gateway reads rather than serializing their independent latency. + let containers_client = self.client.clone(); + let networks_client = self.client.clone(); + let volumes_client = self.client.clone(); + let (containers, networks, volumes) = tokio::try_join!( + containers_client.list_containers(Some(container_options.build())), + networks_client.list_networks(Some(network_options.build())), + volumes_client.list_volumes(Some(volume_options.build())), + ) + .map_err(|error| format!("Docker inventory read failed: {error}"))?; let snapshot = build_snapshot(containers.clone(), networks, volumes); let mounts_by_id = snapshot diff --git a/crates/dockermap-daemon/src/main.rs b/crates/dockermap-daemon/src/main.rs index f4acd427..b9ea9da7 100644 --- a/crates/dockermap-daemon/src/main.rs +++ b/crates/dockermap-daemon/src/main.rs @@ -16,10 +16,10 @@ mod runtime_collection; use axum::{http::StatusCode, response::IntoResponse}; #[cfg(test)] use bollard::Docker; +use cache_refresh::refresh_loop; pub(crate) use cache_refresh::AppState; #[cfg(test)] use cache_refresh::DaemonCache; -use cache_refresh::refresh_loop; use compose_api::run_cli; use config::{ read_allow_mock_env, read_bind_host_env, read_daemon_token_env, read_port_env, DaemonAuthToken, @@ -135,16 +135,16 @@ async fn main() { let daemon_token = read_daemon_token_env(); let port = read_port_env("DOCKERMAP_DAEMON_PORT", 4100); let host = read_bind_host_env("DOCKERMAP_DAEMON_HOST", daemon_token.0.is_some()); -let address = SocketAddr::from((host, port)); -let state = AppState::new(read_allow_mock_env()); -let app = daemon_router(state.clone(), daemon_token); -let listener = TcpListener::bind(address) -.await -.expect("daemon listener should bind"); + let address = SocketAddr::from((host, port)); + let state = AppState::new(read_allow_mock_env()); + let app = daemon_router(state.clone(), daemon_token); + let listener = TcpListener::bind(address) + .await + .expect("daemon listener should bind"); -// The listener is available while the initial authoritative Docker model is -// collected. Cache-backed routes truthfully return 503 until publication. -tokio::spawn(refresh_loop(state)); + // The listener is available while the initial authoritative Docker model is + // collected. Cache-backed routes truthfully return 503 until publication. + tokio::spawn(refresh_loop(state)); println!("dockermap-daemon listening on http://{address}"); @@ -180,22 +180,23 @@ mod tests { use tokio::{ io::{AsyncReadExt, AsyncWriteExt}, net::UnixListener, + sync::Barrier, }; use tower::util::ServiceExt; -fn test_daemon_state() -> AppState { -let mut cache = DaemonCache::mock(); -cache.publication_ready = true; -AppState { -allow_mock: true, -cache: Arc::new(RwLock::new(cache)), + fn test_daemon_state() -> AppState { + let mut cache = DaemonCache::mock(); + cache.publication_ready = true; + AppState { + allow_mock: true, + cache: Arc::new(RwLock::new(cache)), docker: Arc::new(RwLock::new(None)), provider_slot_in_flight: Arc::new(crate::cache_refresh::ProviderSlotFlights::default()), } } -#[tokio::test] -async fn daemon_mock_cache_is_not_published_when_mock_mode_is_disabled() { + #[tokio::test] + async fn daemon_mock_cache_is_not_published_when_mock_mode_is_disabled() { let state = AppState { allow_mock: false, cache: Arc::new(RwLock::new(DaemonCache::mock())), @@ -228,27 +229,32 @@ async fn daemon_mock_cache_is_not_published_when_mock_mode_is_disabled() { .expect("daemon router should respond"); assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE, "{path}"); } -} + } -#[tokio::test] -async fn initializing_cache_is_not_published_even_when_mock_is_allowed() { -let state = AppState::new(true); -for path in ["/daemon/health", "/daemon/snapshot", "/daemon/runtime/map", "/daemon/findings"] { -let response = daemon_router(state.clone(), DaemonAuthToken(None)) -.oneshot( -Request::builder() -.uri(path) -.body(axum::body::Body::empty()) -.expect("request should build"), -) -.await -.expect("daemon router should respond"); -assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE, "{path}"); -} -} + #[tokio::test] + async fn initializing_cache_is_not_published_even_when_mock_is_allowed() { + let state = AppState::new(true); + for path in [ + "/daemon/health", + "/daemon/snapshot", + "/daemon/runtime/map", + "/daemon/findings", + ] { + let response = daemon_router(state.clone(), DaemonAuthToken(None)) + .oneshot( + Request::builder() + .uri(path) + .body(axum::body::Body::empty()) + .expect("request should build"), + ) + .await + .expect("daemon router should respond"); + assert_eq!(response.status(), StatusCode::SERVICE_UNAVAILABLE, "{path}"); + } + } -#[tokio::test] -async fn daemon_bearer_boundary_allows_only_the_exact_configured_token() { + #[tokio::test] + async fn daemon_bearer_boundary_allows_only_the_exact_configured_token() { let allowed = daemon_router(test_daemon_state(), DaemonAuthToken(None)) .oneshot( Request::builder() @@ -371,9 +377,9 @@ async fn daemon_bearer_boundary_allows_only_the_exact_configured_token() { .expect("Docker stub response should be written"); }); -let mut cache = DaemonCache::mock(); -cache.publication_ready = true; -cache.health.docker_reachable = true; + let mut cache = DaemonCache::mock(); + cache.publication_ready = true; + cache.health.docker_reachable = true; let state = AppState { allow_mock: true, cache: Arc::new(RwLock::new(cache)), @@ -519,6 +525,84 @@ cache.health.docker_reachable = true; ], "Bollard wire contract changed; update the gateway ADR and policy review before permitting a new request shape"); } + #[tokio::test] + async fn docker_inventory_reads_arrive_before_the_gateway_releases_them() { + let tempdir = tempfile::tempdir().expect("temporary Docker socket directory"); + let socket_path = tempdir.path().join("docker.sock"); + let listener = UnixListener::bind(&socket_path).expect("Docker stub should bind"); + let release = Arc::new(Barrier::new(4)); + let stub_release = release.clone(); + let stub = tokio::spawn(async move { + let mut responses = Vec::new(); + for _ in 0..3 { + let (mut stream, _) = listener + .accept() + .await + .expect("concurrent Docker request should arrive"); + let release = stub_release.clone(); + responses.push(tokio::spawn(async move { + let mut request = Vec::new(); + let mut chunk = [0_u8; 1024]; + while !request.windows(4).any(|window| window == b"\r\n\r\n") { + let read = stream + .read(&mut chunk) + .await + .expect("Docker request should be readable"); + assert!(read > 0, "Docker client should send request headers"); + request.extend_from_slice(&chunk[..read]); + } + let target = String::from_utf8(request) + .expect("Docker request should be UTF-8") + .lines() + .next() + .and_then(|line| line.split_whitespace().nth(1)) + .expect("Docker request should have a target") + .to_string(); + release.wait().await; + let (delay, body) = if target.contains("/containers/json") { + (Duration::from_millis(30), "[]") + } else if target.contains("/networks") { + (Duration::from_millis(90), "[]") + } else if target.contains("/volumes") { + (Duration::from_millis(60), r#"{"Volumes":[],"Warnings":null}"#) + } else { + panic!("unexpected Docker target: {target}"); + }; + tokio::time::sleep(delay).await; + let response = format!( + "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\nConnection: close\r\n\r\n{}", + body.len(), body + ); + stream.write_all(response.as_bytes()).await.expect("response should write"); + })); + } + for response in responses { + response.await.expect("response task should finish"); + } + }); + let collector = DockerCollector::with_client( + Docker::connect_with_unix( + socket_path.to_str().expect("socket path should be UTF-8"), + 2, + bollard::API_DEFAULT_VERSION, + ) + .expect("Bollard should connect to the Unix stub"), + None, + ); + let started = tokio::time::Instant::now(); + let collection = tokio::spawn(async move { collector.collect_snapshot().await }); + release.wait().await; + collection + .await + .expect("collection task should finish") + .expect("concurrent inventory reads should succeed"); + assert!( + started.elapsed() < Duration::from_millis(160), + "concurrent requests should take the maximum delay, not their sum" + ); + stub.await.expect("Docker stub should finish"); + } + /// Docker label filtering is part of the gateway contract, not a collector /// convenience: the proxy must fail closed if the engine-side scope changes. #[tokio::test] From 56e91c9b19ddbd76c14676de146bc2905dbf227b Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:33:12 +0800 Subject: [PATCH 03/10] feat(daemon): split Compose enrichment from the Docker critical path (#336) What: publish base Docker candidates without awaiting Compose filesystem projection.\n\nWhy: authoritative Docker topology must be available before optional Compose enrichment.\n\nHow checked: cargo test -p dockermap-daemon --manifest-path crates/Cargo.toml repeated_projection_timeouts_are_single_flight_and_keep_the_docker_client. --- crates/dockermap-daemon/src/cache_refresh.rs | 87 +++----------------- 1 file changed, 11 insertions(+), 76 deletions(-) diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index ca3335da..e74060ac 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -866,8 +866,8 @@ enum DockerReadFailure { async fn collect_docker_snapshot_candidate( collector: DockerCollector, snapshot_timeout: Duration, - projection_in_flight: Arc, - projection: F, + _projection_in_flight: Arc, + _projection: F, ) -> Result where F: FnOnce( @@ -877,7 +877,6 @@ where + Send + 'static, { - let started = tokio::time::Instant::now(); // Test-only stage attribution (#335). Disabled unless the benchmark harness // sets an absolute DOCKERMAP_BENCH_STAGE_TIMING_PATH; it records durations // only and never changes what is collected or published. @@ -896,44 +895,10 @@ where ); let mut snapshot = observation.snapshot; snapshot.images = derive_images(&snapshot); - let collected_at = snapshot.last_updated; - - let compose_runtime_binding = - if let Some(flight) = ComposeProjectionFlight::claim(projection_in_flight) { - if let Some(remaining) = snapshot_timeout - .checked_sub(started.elapsed()) - .filter(|remaining| !remaining.is_zero()) - { - // The flight token moves into the closure, so timing out this - // await cannot release it early or permit queued projections. - let projection_task = tokio::task::spawn_blocking(move || { - let _flight = flight; - projection(observation.compose_containers, collected_at) - }); - let projection_started = std::time::Instant::now(); - let binding = match tokio::time::timeout(remaining, projection_task).await { - Ok(Ok(binding)) => binding, - Ok(Err(_)) | Err(_) => None, - }; - // Attributed separately from the Docker observation so the - // baseline can show that this projection currently sits inside - // the Docker publication budget. #336 owns moving it off that - // path; nothing is decoupled here. - crate::bench_timing::record( - bench_sink.as_deref(), - crate::bench_timing::STAGE_COMPOSE_ENRICHMENT, - projection_started, - ); - binding - } else { - None - } - } else { - // A prior timed-out projection is still unwinding. The Docker - // observation is independently truthful, but Compose correlation is - // unavailable for this publication and no second task is launched. - None - }; + // Compose filesystem projection is intentionally not on the authoritative + // Docker collection path. The post-publication worker is attached in the + // following enrichment slice; this base candidate never claims a binding. + let compose_runtime_binding = None; let health = HealthResponse { status: HealthState::Ok, @@ -3336,9 +3301,6 @@ mod scheduler_tests { )))), provider_slot_in_flight: Arc::new(ProviderSlotFlights::default()), }; - let (projection_started_tx, projection_started_rx) = std::sync::mpsc::channel(); - let (release_projection_tx, release_projection_rx) = std::sync::mpsc::channel(); - let (projection_finished_tx, projection_finished_rx) = std::sync::mpsc::channel(); let projection_runs = Arc::new(std::sync::atomic::AtomicUsize::new(0)); let timeout = Duration::from_millis(150); let started = tokio::time::Instant::now(); @@ -3346,30 +3308,17 @@ mod scheduler_tests { let first = collect_snapshot_with_projection(&state, timeout, move |_containers, _collected_at| { first_runs.fetch_add(1, std::sync::atomic::Ordering::SeqCst); - projection_started_tx - .send(()) - .expect("test observes projection start"); - release_projection_rx - .recv() - .expect("test releases stalled projection"); - projection_finished_tx - .send(()) - .expect("test observes projection completion"); None }) .await; - assert!(started.elapsed() >= timeout); + assert!(started.elapsed() < timeout); assert!(started.elapsed() < Duration::from_secs(1)); - projection_started_rx - .try_recv() - .expect("Docker reads completed before projection stalled"); assert_eq!(first.health.mode, RuntimeMode::Docker); assert!(first.compose_runtime_binding.is_none()); assert!(state.docker.read().await.is_some()); - // Repeated refreshes still perform their fresh Docker reads but skip - // Compose projection while the first blocking task owns the guard. + // Repeated refreshes remain independent of Compose projection. for _ in 0..2 { let runs = projection_runs.clone(); let refresh_started = tokio::time::Instant::now(); @@ -3389,26 +3338,12 @@ mod scheduler_tests { } assert_eq!( projection_runs.load(std::sync::atomic::Ordering::SeqCst), - 1, - "only the original stalled projection may run" + 0, + "base Docker publication must not run filesystem projection" ); - // Publishing the Docker result remains an explicit caller action. The late - // read-only task owns no state and cannot replace it after release. + // Publishing the Docker result remains an explicit caller action. publish_docker_snapshot_cache(&state, first).await; - release_projection_tx - .send(()) - .expect("release private blocking task"); - projection_finished_rx - .recv_timeout(Duration::from_secs(1)) - .expect("late private projection completes"); - while state - .provider_slot_in_flight - .compose_projection - .load(std::sync::atomic::Ordering::Acquire) - { - tokio::task::yield_now().await; - } assert_eq!(state.cache.read().await.health.mode, RuntimeMode::Docker); assert_eq!(gateway.await.expect("gateway stub").len(), 9); } From 550c908830871a87cc112c4afd5d6b579595cf1f Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:36:47 +0800 Subject: [PATCH 04/10] feat(daemon): attach coherent Compose enrichment (#336) What: carry keyed Compose work after Docker publication, attach only to the matching Docker source observation, and include private binding identity in revisioning.\n\nWhy: late or timed-out filesystem work must never relabel newer Docker topology.\n\nHow checked: cargo test -p dockermap-daemon --manifest-path crates/Cargo.toml repeated_projection_timeouts_are_single_flight_and_keep_the_docker_client. --- crates/dockermap-daemon/src/cache_refresh.rs | 102 +++++++++++++++++-- 1 file changed, 94 insertions(+), 8 deletions(-) diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index e74060ac..45041082 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -92,6 +92,7 @@ pub(crate) struct DaemonCache { pub(crate) runtime_map: RuntimeMap, pub(crate) findings: FindingsResponse, compose_runtime_binding: Option<(dockermap_core::ComposeScan, ComposeRuntimeBinding)>, + compose_containers: Option>, runtime_providers: RuntimeProviderSlots, /// Increments on every Docker/mock source transition. A late worker must /// match this generation as well as evidence, so Docker→mock→Docker can @@ -269,6 +270,7 @@ impl PublicationRevision { snapshot: &mut DockerSnapshot, health: &mut HealthResponse, runtime_map: &mut RuntimeMap, + compose_binding_identity: Option<&str>, ) { // Compare precisely the model that routes can expose. Cache inventory // intentionally retains raw identities for correlation, so serializing @@ -289,6 +291,7 @@ impl PublicationRevision { &published_snapshot, &published_health, &published_runtime_map, + compose_binding_identity, )) .expect("public DockerMap models are serializable"); if self.last_observable.as_deref() != Some(observable.as_str()) { @@ -624,6 +627,7 @@ impl DaemonCache { }, findings: FindingsResponse::default(), compose_runtime_binding: None, + compose_containers: None, runtime_providers: unavailable_provider_slots(), source_generation: 0, docker_observation_revision: DockerObservationRevision::new(), @@ -639,8 +643,13 @@ impl DaemonCache { fn assign_revision(&mut self) { // All three independently routable model envelopes attest the same // publication. Provider state is runtime-topology evidence only. - self.revision - .assign(&mut self.snapshot, &mut self.health, &mut self.runtime_map); + let compose_binding_identity = self.compose_binding_identity().map(str::to_owned); + self.revision.assign( + &mut self.snapshot, + &mut self.health, + &mut self.runtime_map, + compose_binding_identity.as_deref(), + ); // Findings are a pure projection of the sanitized runtime map, so // calculate and cache them only after the publication revision exists. let bench_sink = crate::bench_timing::sink(); @@ -671,6 +680,12 @@ impl DaemonCache { }; } + fn compose_binding_identity(&self) -> Option<&str> { + self.compose_runtime_binding + .as_ref() + .map(|(_, binding)| binding.provider_revision.as_str()) + } + fn assign_docker_observation_revision(&mut self) { self.docker_observation_revision .assign(&self.snapshot, &self.health.mode); @@ -725,7 +740,14 @@ pub(crate) async fn refresh_loop(state: AppState) { } pub(crate) async fn refresh_cache(state: &AppState) { - let (snapshot, mode, source_generation) = publish_docker_snapshot_cache( + let ( + snapshot, + mode, + source_generation, + docker_observation_revision, + compose_containers, + collected_at, + ) = publish_docker_snapshot_cache( state, collect_snapshot(state, DOCKER_SNAPSHOT_COLLECTION_TIMEOUT).await, ) @@ -735,11 +757,61 @@ pub(crate) async fn refresh_cache(state: &AppState) { // Spawned slot workers use fixed per-slot guards; they have no route, // Docker client, or source-fallback authority. if mode == RuntimeMode::Docker { + spawn_compose_projection( + state.clone(), + source_generation, + docker_observation_revision, + compose_containers, + collected_at, + ); let due = claim_due_provider_slots(state, monotonic_now()).await; spawn_provider_slots(state.clone(), snapshot, mode, source_generation, due); } } +fn spawn_compose_projection( + state: AppState, + source_generation: u64, + docker_observation_revision: String, + compose_containers: Option>, + collected_at: u64, +) { + let Some(flight) = + ComposeProjectionFlight::claim(state.provider_slot_in_flight.compose_projection.clone()) + else { + return; + }; + tokio::spawn(async move { + let task = tokio::task::spawn_blocking(move || { + let _flight = flight; + bounded_compose_runtime_binding(compose_containers, collected_at) + }); + let projection_started = std::time::Instant::now(); + let binding = match tokio::time::timeout(DOCKER_SNAPSHOT_COLLECTION_TIMEOUT, task).await { + Ok(Ok(binding)) => binding, + Ok(Err(_)) | Err(_) => None, + }; + crate::bench_timing::record( + crate::bench_timing::sink().as_deref(), + crate::bench_timing::STAGE_COMPOSE_ENRICHMENT, + projection_started, + ); + let Some((scan, mut binding)) = binding else { + return; + }; + binding.provider_revision = docker_observation_revision.clone(); + let mut cache = state.cache.write().await; + if cache.health.mode != RuntimeMode::Docker + || cache.source_generation != source_generation + || cache.docker_observation_token() != docker_observation_revision + { + return; + } + cache.compose_runtime_binding = Some((scan, binding)); + cache.assign_revision(); + }); +} + fn spawn_provider_slots( state: AppState, snapshot: DockerSnapshot, @@ -778,7 +850,14 @@ fn spawn_provider_slots( async fn publish_docker_snapshot_cache( state: &AppState, mut updated: DaemonCache, -) -> (DockerSnapshot, RuntimeMode, u64) { +) -> ( + DockerSnapshot, + RuntimeMode, + u64, + String, + Option>, + u64, +) { let mut cache = state.cache.write().await; // A mock fallback is a distinct source of bytes. Do not retain live host // observations and relabel them as sample data (or vice versa). @@ -794,6 +873,7 @@ async fn publish_docker_snapshot_cache( // A completed collection attempt, including an explicit policy-permitted // mock fallback, makes this coherent cache eligible for publication. updated.publication_ready = true; + let compose_containers = updated.compose_containers.take(); updated.runtime_providers = if same_source { cache.runtime_providers.clone() } else { @@ -833,6 +913,9 @@ async fn publish_docker_snapshot_cache( cache.snapshot.clone(), cache.health.mode.clone(), cache.source_generation, + cache.docker_observation_token(), + compose_containers, + cache.snapshot.last_updated, ) } @@ -893,6 +976,7 @@ where crate::bench_timing::STAGE_DOCKER_OBSERVATION, bench_docker_started, ); + let compose_containers = observation.compose_containers; let mut snapshot = observation.snapshot; snapshot.images = derive_images(&snapshot); // Compose filesystem projection is intentionally not on the authoritative @@ -916,6 +1000,7 @@ where runtime_map: empty_runtime_map(0), findings: FindingsResponse::default(), compose_runtime_binding, + compose_containers, runtime_providers: unavailable_provider_slots(), source_generation: 0, docker_observation_revision: DockerObservationRevision::new(), @@ -3065,6 +3150,7 @@ mod scheduler_tests { runtime_map: empty_runtime_map(last_updated), findings: FindingsResponse::default(), compose_runtime_binding: None, + compose_containers: None, runtime_providers: unavailable_provider_slots(), source_generation: 0, docker_observation_revision: DockerObservationRevision::new(), @@ -4231,7 +4317,7 @@ mod scheduler_tests { docker: Arc::new(RwLock::new(None)), provider_slot_in_flight: Arc::new(ProviderSlotFlights::default()), }; - let (observed, mode, generation) = + let (observed, mode, generation, ..) = publish_docker_snapshot_cache(&state, docker_cache(mock_snapshot())).await; publish_docker_snapshot_cache(&state, DaemonCache::mock()).await; publish_docker_snapshot_cache(&state, docker_cache(mock_snapshot())).await; @@ -4300,14 +4386,14 @@ mod scheduler_tests { message: Some("controlled Docker cache".into()), }; let mut runtime_map = empty_runtime_map(10); - revision.assign(&mut snapshot, &mut health, &mut runtime_map); + revision.assign(&mut snapshot, &mut health, &mut runtime_map, None); let first = snapshot.model_revision.clone(); snapshot.last_updated = 11; health.last_updated = 11; health.snapshot_version = "11".into(); runtime_map.last_updated = 11; - revision.assign(&mut snapshot, &mut health, &mut runtime_map); + revision.assign(&mut snapshot, &mut health, &mut runtime_map, None); assert_eq!(snapshot.model_revision, first); assert_eq!(health.model_revision, first); @@ -4542,7 +4628,7 @@ mod scheduler_tests { }; let mut old = mock_snapshot(); old.last_updated = 10; - let (observed, mode, generation) = + let (observed, mode, generation, ..) = publish_docker_snapshot_cache(&state, docker_cache(old)).await; let mut newer = mock_snapshot(); newer.last_updated = 11; From 6a0467ee430836fb8bc6b5eb9e4e4cbc658451fb Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 12:47:32 +0800 Subject: [PATCH 05/10] test(perf): make the Baseline-4 promotion comparison runnable (#336) What: add the read-only composite promotion CLI, listener-boundary probe, focused spawn tests, and documentation.\n\nWhy: make the pinned Baseline-4 authority runnable while keeping Docker-model readiness distinct from listener readiness.\n\nHow checked: npm run test:perf; npm run typecheck; npm run fmt:rust:check. --- crates/dockermap-daemon/src/daemon_api.rs | 40 +++++------ docs/architecture/ARCHITECTURE.md | 9 +++ docs/testing/TIME_TO_ANSWER_EVIDENCE.md | 29 ++++---- package.json | 5 +- tests/perf/capture.ts | 16 ++++- tests/perf/captureSyntax.test.mjs | 9 +++ tests/perf/promote.test.mjs | 82 +++++++++++++++++++++++ tests/perf/promote.ts | 81 ++++++++++++++++++++++ tests/perf/tsconfig.json | 2 +- 9 files changed, 236 insertions(+), 37 deletions(-) create mode 100644 tests/perf/promote.test.mjs create mode 100644 tests/perf/promote.ts diff --git a/crates/dockermap-daemon/src/daemon_api.rs b/crates/dockermap-daemon/src/daemon_api.rs index 65c5406f..c076feea 100644 --- a/crates/dockermap-daemon/src/daemon_api.rs +++ b/crates/dockermap-daemon/src/daemon_api.rs @@ -84,14 +84,14 @@ pub(crate) fn daemon_router(state: AppState, daemon_token: DaemonAuthToken) -> R } async fn publication_cache(state: &AppState) -> Result, ApiError> { -let cache = state.cache.read().await; -if !cache.publication_ready { -return Err(ApiError { -status: StatusCode::SERVICE_UNAVAILABLE, -message: "Docker model is still initializing".into(), -}); -} -if !state.allow_mock && cache.health.mode == RuntimeMode::Mock { + let cache = state.cache.read().await; + if !cache.publication_ready { + return Err(ApiError { + status: StatusCode::SERVICE_UNAVAILABLE, + message: "Docker model is still initializing".into(), + }); + } + if !state.allow_mock && cache.health.mode == RuntimeMode::Mock { return Err(ApiError { status: StatusCode::SERVICE_UNAVAILABLE, message: "Live Docker authority is unavailable and mock publication is disabled".into(), @@ -426,18 +426,18 @@ mod tests { use super::*; #[tokio::test] -async fn findings_route_stamps_the_actual_cache_mode_after_live_and_mock_resets() { -let state = AppState::new(true); -assert_eq!( -get_findings(State(state.clone())).await.unwrap_err().status, -StatusCode::SERVICE_UNAVAILABLE, -"bootstrap mock data must not be published before the first attempt" -); - -{ -let mut cache = state.cache.write().await; -cache.publication_ready = true; -cache.health.mode = RuntimeMode::Docker; + async fn findings_route_stamps_the_actual_cache_mode_after_live_and_mock_resets() { + let state = AppState::new(true); + assert_eq!( + get_findings(State(state.clone())).await.unwrap_err().status, + StatusCode::SERVICE_UNAVAILABLE, + "bootstrap mock data must not be published before the first attempt" + ); + + { + let mut cache = state.cache.write().await; + cache.publication_ready = true; + cache.health.mode = RuntimeMode::Docker; // A source marker is publication data, never cache data. This // deliberately stale value models a completed live publication // followed by an eventual reset. diff --git a/docs/architecture/ARCHITECTURE.md b/docs/architecture/ARCHITECTURE.md index c9152550..1dccabd2 100644 --- a/docs/architecture/ARCHITECTURE.md +++ b/docs/architecture/ARCHITECTURE.md @@ -46,6 +46,15 @@ generated schemas rather than defining them. Route, request, and response associ derive the OpenAPI 3.1.1 document. The ownership map, drift checks, and remaining #65 acceptance work are recorded in [`CONTRACT_AUTHORITY.md`](CONTRACT_AUTHORITY.md). +Docker inventory is published as the first authoritative model without waiting for +Compose filesystem enrichment. Compose correlation is a bounded asynchronous follow-up +that attaches only to the matching Docker publication, so it cannot spend the Docker +publication budget or be applied to a newer model. The published performance +`daemonStartToListenerMs` stage now ends at the real HTTP listener boundary. In the +Baseline-4 capture, the listener was gated behind first collection, so that historical +stage folded first collection into its value; the Baseline-4 numbers retain their +original meaning and this semantic difference must be considered when comparing them. + ## Runtime Map `GET /daemon/runtime/map` is the backend's provider-neutral JSON graph for visualization. `apps/api` proxies it as `GET /api/runtime/map`. diff --git a/docs/testing/TIME_TO_ANSWER_EVIDENCE.md b/docs/testing/TIME_TO_ANSWER_EVIDENCE.md index 69fcb305..5674d5ef 100644 --- a/docs/testing/TIME_TO_ANSWER_EVIDENCE.md +++ b/docs/testing/TIME_TO_ANSWER_EVIDENCE.md @@ -96,10 +96,9 @@ examples of the distinction that matter most: model is observable. It does **not** prove the model is complete — optional provider evidence may still be missing, and a fast number here must never be read as "the host is fully described". -- `composeEnrichmentMs` measures Compose filesystem correlation separately from - the Docker observation. While the two remain coupled inside one publication - budget, this stage is **measured, not removed**; #336 owns moving it off the - critical path. +- `composeEnrichmentMs` measures asynchronous Compose filesystem correlation + separately from Docker publication. Compose enrichment no longer consumes the + Docker publication budget; it is measured as a later enrichment path. - `dockerObservationMs` is measured against a deterministic local fixture daemon. It is **not** a claim about a real Docker daemon's latency, host load, or image size. @@ -201,11 +200,10 @@ newline-delimited JSON records to that file. There is no route, no response field, no runtime telemetry, and no behaviour change. The hook times the Docker inventory read and the Compose filesystem projection -**separately while both still execute inside the same Docker publication -budget**. This baseline is therefore expected to show that Compose projection -currently sits inside the Docker critical path. That is the measurement, not a -fix: **nothing is decoupled here, and #336 owns moving the projection off that -path** — these are the numbers it must improve against. +separately. The Baseline-4 capture measured them while Compose still executed +inside the Docker publication budget. #336 moves Compose onto an asynchronous +enrichment path, so it no longer consumes that publication budget; the baseline +methodology and its recorded numbers remain unchanged comparison authority. ## Stage 6 and stage 7: two clocks, and why they cannot be one @@ -375,11 +373,16 @@ npx tsx tests/perf/assembleCompositeEvidence.ts \ --output /tmp/time-to-answer-baseline.json # 6. recompute summaries from the raw samples (never trust supplied aggregates) npm run perf:summarize -- --artifact /tmp/time-to-answer-baseline.json -# 7. compare a candidate against a reviewed baseline (fails closed) -npm run perf:time-to-answer -- \ - --metadata /tmp/time-to-answer-metadata.json \ +# 7. assemble the candidate's end-to-end and Stage-5 captures into its composite +npx tsx tests/perf/assembleCompositeEvidence.ts \ + --general /tmp/time-to-answer-candidate-general.json \ --output /tmp/time-to-answer-candidate.json \ - --baseline /tmp/time-to-answer-baseline.json + --stageFive /tmp/time-to-answer-candidate-stage5.json \ + --output /tmp/time-to-answer-candidate.json +# 8. compare the two composite artifacts (fails closed) +npm run perf:promote -- \ + --baseline /tmp/time-to-answer-baseline.json \ + --candidate /tmp/time-to-answer-candidate.json ``` Prerequisites: a release daemon (`cargo build --release -p dockermap-daemon`), diff --git a/package.json b/package.json index 126e17b1..86f457ef 100644 --- a/package.json +++ b/package.json @@ -31,8 +31,9 @@ "test:contracts": "npm run test --workspace @dockermap/contracts --if-present", "test:version": "node --test scripts/check-version-authority.test.mjs scripts/package-release.test.mjs", "test:perf": "node --test tests/perf/*.test.mjs", - "perf:time-to-answer": "tsx tests/perf/capture.ts", - "perf:calibrate-time-to-answer": "tsx tests/perf/capture.ts", +"perf:time-to-answer": "tsx tests/perf/capture.ts", +"perf:promote": "tsx tests/perf/promote.ts", +"perf:calibrate-time-to-answer": "tsx tests/perf/capture.ts", "perf:preconditioning": "tsx tests/perf/preconditioning.ts", "perf:phase-control": "tsx tests/perf/phaseControl.ts", "perf:stage-five": "tsx tests/perf/captureStageFive.ts", diff --git a/tests/perf/capture.ts b/tests/perf/capture.ts index 6cf5573a..d604be8a 100644 --- a/tests/perf/capture.ts +++ b/tests/perf/capture.ts @@ -401,6 +401,20 @@ async function fetchJson(url: string, timeoutMs = 5_000): Promise { } } +/** A listener is ready once it answers HTTP; a truthful 503 is not a model. */ +async function httpListenerReady(url: string, timeoutMs = 1_000): Promise { +const controller = new AbortController(); +const timer = setTimeout(() => controller.abort(), timeoutMs); +try { +await fetch(url, { signal: controller.signal }); +return true; +} catch { +return false; +} finally { +clearTimeout(timer); +} +} + async function waitForJson(url: string, predicate: (value: any) => boolean, timeoutMs: number) { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { @@ -1256,7 +1270,7 @@ let publicationTracker: PublicationTracker | null = null; ...(plan.name === "unavailable-optional-provider" ? { PATH: emptyPath } : {}) }; const healthUrl = (port: number) => `http://127.0.0.1:${port}/daemon/health`; - const daemonReady = async (port: number) => Boolean(await fetchJson(healthUrl(port), 1_000)); +const daemonReady = async (port: number) => httpListenerReady(healthUrl(port)); const startedDaemon = await startChildOnFreePort({ name: "daemon", spawnOn: (port) => spawnOwned(daemonBinary, [], { ...daemonEnv, DOCKERMAP_DAEMON_PORT: String(port) }), diff --git a/tests/perf/captureSyntax.test.mjs b/tests/perf/captureSyntax.test.mjs index 5f538032..a8ed047a 100644 --- a/tests/perf/captureSyntax.test.mjs +++ b/tests/perf/captureSyntax.test.mjs @@ -18,3 +18,12 @@ test("capture harness entrypoint transforms", async () => { throw new Error(`capture.ts failed esbuild transform:\n${details || error.message}`); } }); + +test("listener readiness accepts an HTTP response independently of Docker-model readiness", async () => { +const source = await readFile(new URL("./capture.ts", import.meta.url), "utf8"); +if (!source.includes("async function httpListenerReady")) throw new Error("capture must use an HTTP listener probe"); +if (!source.includes("ping: daemonReady")) throw new Error("daemon startup must use listener readiness"); +if (!source.includes("value.mode === \"docker\" && Boolean(value.modelRevision)")) { +throw new Error("Docker-model readiness must remain a later distinct boundary"); +} +}); diff --git a/tests/perf/promote.test.mjs b/tests/perf/promote.test.mjs new file mode 100644 index 00000000..4fb897d3 --- /dev/null +++ b/tests/perf/promote.test.mjs @@ -0,0 +1,82 @@ +import { mkdtempSync, rmSync, writeFileSync } from "node:fs"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { spawnSync } from "node:child_process"; +import assert from "node:assert/strict"; +import test from "node:test"; + +const fixtures = ["reference-25", "reference-100", "reference-250"]; +const stages = [ +["daemonStartToListenerMs", fixtures], +["listenerToFirstDockerModelMs", fixtures], +["dockerObservationMs", fixtures], +["composeEnrichmentMs", [...fixtures, "slow-bounded-compose-projection"]], +["publicationToNodeObservationMs", [...fixtures, "provider-only-revision-change", "docker-topology-change", "unavailable-optional-provider"]], +["notificationToCoherentModelMs", [...fixtures, "provider-only-revision-change", "docker-topology-change", "unavailable-optional-provider"]], +["coherentModelToUsefulRenderMs", [...fixtures, "docker-topology-change"]], +["buildModelMs", fixtures], +["findingsDerivationMs", fixtures], +["legacyTopologyLayoutMs", fixtures], +["commandQueryMs", fixtures], +["productionBundleMs", fixtures] +]; + +const environment = { +runnerClass: "linux-x86_64-dedicated", cpuClass: "cpus-16vcpu", osImage: "ubuntu-26.04", osKernel: "7.0.0-31-generic", +nodeRevision: "22.23.2", rustRevision: "1.88.0", dockerRevision: "29.8.1", ssePollIntervalMs: "2000", +daemonBinarySha256: "eeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeeee", +daemonBinaryBuild: "cargo-build-release-locked-p-dockermap-daemon", cargoRevision: "cargo-1.88.0", +harnessRevision: "dddddddddddddddddddddddddddddddddddddddd", browserEngine: "chromium", browserRevision: "1.61.0", +browserFlags: ["--disable-background-networking"], fontEnvironment: "system-default", buildMode: "production", +fixtureRevision: "dockermap-v1/time-to-answer-fixtures-1", sourceRevision: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", +methodologyVersion: "dockermap-v1/time-to-answer-methodology-8" +}; + +function artifact(value = 10, overrides = {}) { +return { +baseline: "dockermap-v1/time-to-answer-baseline-4", +environment: { ...environment, ...(overrides.environment ?? {}) }, +records: stages.flatMap(([stage, names]) => names.map((fixture) => ({ +fixture, stage, measurementProtocol: stage === "publicationToNodeObservationMs" ? "controlled-poll-phase" : "end-to-end", +sourceEvidenceFile: "synthetic.raw.json", checkpointSha: "aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa", +runs: Array.from({ length: 3 }, () => Array.from({ length: 15 }, () => value)) +}))) +}; +} + +function run(baseline, candidate) { +const directory = mkdtempSync(join(tmpdir(), "dockermap-promote-")); +try { +const baselinePath = join(directory, "baseline.json"); +const candidatePath = join(directory, "candidate.json"); +writeFileSync(baselinePath, JSON.stringify(baseline)); +writeFileSync(candidatePath, JSON.stringify(candidate)); +return spawnSync("npx", ["tsx", "tests/perf/promote.ts", "--baseline", baselinePath, "--candidate", candidatePath], { +cwd: new URL("../..", import.meta.url), encoding: "utf8", timeout: 30_000 +}); +} finally { +rmSync(directory, { recursive: true, force: true }); +} +} + +test("promotion CLI accepts a compatible candidate within the reviewed limit", () => { +const result = run(artifact(10), artifact(12)); +assert.equal(result.status, 0, result.stderr); +assert.match(result.stdout, /PROMOTION: PASS/); +}); + +test("promotion CLI names a stage beyond the reviewed limit", () => { +const candidate = artifact(10); +candidate.records[0].runs = Array.from({ length: 3 }, () => Array.from({ length: 15 }, () => 13)); +const result = run(artifact(10), candidate); +assert.notEqual(result.status, 0); +assert.match(result.stdout, /reference-25\u0000daemonStartToListenerMs/); +assert.match(result.stdout, /PROMOTION: FAIL:Time-to-answer candidate exceeds/); +}); + +test("promotion CLI rejects an incompatible environment before a timing limit", () => { +const result = run(artifact(10), artifact(99, { environment: { osImage: "different-image" } })); +assert.notEqual(result.status, 0); +assert.match(result.stdout, /candidate does not match the pinned baseline environment/); +assert.doesNotMatch(result.stdout, /candidate exceeds the reviewed promotion limit/); +}); diff --git a/tests/perf/promote.ts b/tests/perf/promote.ts new file mode 100644 index 00000000..57765b72 --- /dev/null +++ b/tests/perf/promote.ts @@ -0,0 +1,81 @@ +#!/usr/bin/env node +/** Read-only promotion comparison for two closed composite evidence artifacts. */ +import { readFileSync } from "node:fs"; +import { +assertTimeToAnswerPromotion, +derivedTimeToAnswerSummaries, +timeToAnswerLimit, +validateTimeToAnswerEvidence +} from "../../apps/web/src/lib/performance/timeToAnswerEvidence"; + +type Paths = { baseline: string; candidate: string }; + +function paths(arguments_: readonly string[]): Paths { +const values: Partial = {}; +for (let index = 0; index < arguments_.length; index += 1) { +const flag = arguments_[index]!; +if (flag !== "--baseline" && flag !== "--candidate") throw new Error(`Unknown flag: ${flag}`); +const value = arguments_[index + 1]; +if (!value || value.startsWith("--")) throw new Error(`Missing value for ${flag}`); +const key = flag.slice(2) as keyof Paths; +if (values[key]) throw new Error(`Duplicate flag: ${flag}`); +values[key] = value; +index += 1; +} +if (!values.baseline || !values.candidate) { +throw new Error("Usage: npm run perf:promote -- --baseline --candidate "); +} +return values as Paths; +} + +function readEvidence(path: string): unknown { +try { +return JSON.parse(readFileSync(path, "utf8")); +} catch (error) { +throw new Error(`Cannot read ${path}: ${error instanceof Error ? error.message : String(error)}`); +} +} + +function format(value: number): string { return value.toFixed(2); } + +function main(): void { +const input = paths(process.argv.slice(2)); +const baselineRaw = readEvidence(input.baseline); +const candidateRaw = readEvidence(input.candidate); +const baseline = validateTimeToAnswerEvidence(baselineRaw); +const candidate = validateTimeToAnswerEvidence(candidateRaw); +const baselineSummaries = derivedTimeToAnswerSummaries(baseline); +const candidateSummaries = derivedTimeToAnswerSummaries(candidate); +const regressions: string[] = []; +const violations: string[] = []; +process.stdout.write("| fixture | stage | protocol | baseline reviewed (ms) | limit (ms) | candidate reviewed (ms) | delta (ms) | verdict |\n"); +process.stdout.write("| --- | --- | --- | ---: | ---: | ---: | ---: | --- |\n"); +for (const record of candidate.records) { +const key = `${record.fixture}\u0000${record.stage}`; +const before = baselineSummaries.get(key)!; +const after = candidateSummaries.get(key)!; +const limit = timeToAnswerLimit(before.reviewedMs); +const delta = after.reviewedMs - before.reviewedMs; +if (delta > 0) regressions.push(`${record.fixture}/${record.stage}`); +if (after.reviewedMs > limit) violations.push(key); +process.stdout.write(`| ${record.fixture} | ${record.stage} | ${record.measurementProtocol} | ${format(before.reviewedMs)} | ${format(limit)} | ${format(after.reviewedMs)} | ${format(delta)} | ${after.reviewedMs <= limit ? "PASS" : "FAIL"} |\n`); +} +process.stdout.write(`Slower than baseline: ${regressions.join(", ") || "none"}\n`); +try { +assertTimeToAnswerPromotion(baselineRaw, candidateRaw); +process.stdout.write("PROMOTION: PASS\n"); +} catch (error) { +if (violations.length > 0) process.stdout.write(`Offending keys: ${violations.join(", ")}\n`); +const reason = error instanceof Error ? error.message : String(error); +process.stdout.write(`PROMOTION: FAIL:${reason}\n`); +process.exitCode = 1; +} +} + +try { +main(); +} catch (error) { +const reason = error instanceof Error ? error.message : String(error); +process.stderr.write(`PROMOTION: FAIL:${reason}\n`); +process.exitCode = 1; +} diff --git a/tests/perf/tsconfig.json b/tests/perf/tsconfig.json index 1b47df4c..aa312331 100644 --- a/tests/perf/tsconfig.json +++ b/tests/perf/tsconfig.json @@ -4,5 +4,5 @@ "types": ["node"], "allowJs": true }, - "include": ["capture.ts", "captureIndependence.ts", "summarize.ts", "preconditioning.ts", "phaseControl.ts"] +"include": ["capture.ts", "captureIndependence.ts", "summarize.ts", "preconditioning.ts", "phaseControl.ts", "promote.ts"] } From 92d7cbf82b2ff118f04c96104131fa49f8389c7a Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 13:04:30 +0800 Subject: [PATCH 06/10] test(perf): prove listener readiness with a real HTTP probe (#336) What: extract listener readiness into a reusable HTTP probe and cover successful, refused, and timed-out listener responses with real local servers.\n\nWhy: an HTTP response, including a truthful 503, proves listener availability while Docker-model readiness remains a later boundary.\n\nHow checked: npm run test:perf; npm run typecheck. --- tests/perf/capture.ts | 15 +------- tests/perf/captureSyntax.test.mjs | 4 +-- tests/perf/listenerReadiness.mjs | 13 +++++++ tests/perf/listenerReadiness.test.mjs | 49 +++++++++++++++++++++++++++ 4 files changed, 64 insertions(+), 17 deletions(-) create mode 100644 tests/perf/listenerReadiness.mjs create mode 100644 tests/perf/listenerReadiness.test.mjs diff --git a/tests/perf/capture.ts b/tests/perf/capture.ts index d604be8a..f46240f6 100644 --- a/tests/perf/capture.ts +++ b/tests/perf/capture.ts @@ -70,6 +70,7 @@ import { reservePort, startStaticServer } from "./staticServer.mjs"; import { withFreshBrowserRuns } from "./browserLifecycle.mjs"; import { startStageFivePublicationController } from "./stageFivePublicationControl.mjs"; import { armCaptureStageFivePublication } from "./stageFiveCaptureControl.mjs"; +import { httpListenerReady } from "./listenerReadiness.mjs"; const REPO_ROOT = resolve(fileURLToPath(new URL("../..", import.meta.url))); const sleep = (ms: number) => new Promise((done) => setTimeout(done, ms)); @@ -401,20 +402,6 @@ async function fetchJson(url: string, timeoutMs = 5_000): Promise { } } -/** A listener is ready once it answers HTTP; a truthful 503 is not a model. */ -async function httpListenerReady(url: string, timeoutMs = 1_000): Promise { -const controller = new AbortController(); -const timer = setTimeout(() => controller.abort(), timeoutMs); -try { -await fetch(url, { signal: controller.signal }); -return true; -} catch { -return false; -} finally { -clearTimeout(timer); -} -} - async function waitForJson(url: string, predicate: (value: any) => boolean, timeoutMs: number) { const deadline = Date.now() + timeoutMs; while (Date.now() < deadline) { diff --git a/tests/perf/captureSyntax.test.mjs b/tests/perf/captureSyntax.test.mjs index a8ed047a..20475f92 100644 --- a/tests/perf/captureSyntax.test.mjs +++ b/tests/perf/captureSyntax.test.mjs @@ -19,10 +19,8 @@ test("capture harness entrypoint transforms", async () => { } }); -test("listener readiness accepts an HTTP response independently of Docker-model readiness", async () => { +test("Docker-model readiness remains a later distinct boundary", async () => { const source = await readFile(new URL("./capture.ts", import.meta.url), "utf8"); -if (!source.includes("async function httpListenerReady")) throw new Error("capture must use an HTTP listener probe"); -if (!source.includes("ping: daemonReady")) throw new Error("daemon startup must use listener readiness"); if (!source.includes("value.mode === \"docker\" && Boolean(value.modelRevision)")) { throw new Error("Docker-model readiness must remain a later distinct boundary"); } diff --git a/tests/perf/listenerReadiness.mjs b/tests/perf/listenerReadiness.mjs new file mode 100644 index 00000000..c0a4ab55 --- /dev/null +++ b/tests/perf/listenerReadiness.mjs @@ -0,0 +1,13 @@ +/** A listener is ready once it answers HTTP; a truthful 503 is not a model. */ +export async function httpListenerReady(url, timeoutMs = 1_000) { +const controller = new AbortController(); +const timer = setTimeout(() => controller.abort(), timeoutMs); +try { +await fetch(url, { signal: controller.signal }); +return true; +} catch { +return false; +} finally { +clearTimeout(timer); +} +} diff --git a/tests/perf/listenerReadiness.test.mjs b/tests/perf/listenerReadiness.test.mjs new file mode 100644 index 00000000..12480e18 --- /dev/null +++ b/tests/perf/listenerReadiness.test.mjs @@ -0,0 +1,49 @@ +import assert from "node:assert/strict"; +import { createServer } from "node:http"; +import test from "node:test"; +import { httpListenerReady } from "./listenerReadiness.mjs"; + +async function serverFor(handler) { +const server = createServer(handler); +await new Promise((done) => server.listen(0, "127.0.0.1", done)); +return { +url: `http://127.0.0.1:${server.address().port}`, +close: () => new Promise((done) => server.close(done)) +}; +} + +test("listener readiness accepts a truthful HTTP 503", async () => { +const server = await serverFor((_request, response) => response.writeHead(503).end()); +try { +assert.equal(await httpListenerReady(server.url, 150), true); +} finally { +await server.close(); +} +}); + +test("listener readiness accepts an HTTP 200 JSON response", async () => { +const server = await serverFor((_request, response) => response.writeHead(200, { "content-type": "application/json" }).end(JSON.stringify({ ready: true }))); +try { +assert.equal(await httpListenerReady(server.url, 150), true); +} finally { +await server.close(); +} +}); + +test("listener readiness rejects a closed port", async () => { +const server = await serverFor((_request, response) => response.end()); +const url = server.url; +await server.close(); +assert.equal(await httpListenerReady(url, 150), false); +}); + +test("listener readiness rejects a listener that never responds before its timeout", async () => { +const server = await serverFor(() => {}); +try { +const startedAt = Date.now(); +assert.equal(await httpListenerReady(server.url, 150), false); +assert.ok(Date.now() - startedAt < 1_000, "listener probe should use its explicit timeout"); +} finally { +await server.close(); +} +}); From 562d4f46657047c414972e73ecbf03bfd8979cbd Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 13:07:13 +0800 Subject: [PATCH 07/10] docs(perf): correct the candidate assembly example (#336) MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit What: Removed the duplicate candidate --output flag from the Step 7 command.\n\nWhy: The example now shows the valid general → stageFive → output argument order.\n\nHow checked: Verified the Step 7 block, one-line documentation diff, whitespace check, and working-tree scope. --- docs/testing/TIME_TO_ANSWER_EVIDENCE.md | 1 - 1 file changed, 1 deletion(-) diff --git a/docs/testing/TIME_TO_ANSWER_EVIDENCE.md b/docs/testing/TIME_TO_ANSWER_EVIDENCE.md index 5674d5ef..5e95f28f 100644 --- a/docs/testing/TIME_TO_ANSWER_EVIDENCE.md +++ b/docs/testing/TIME_TO_ANSWER_EVIDENCE.md @@ -376,7 +376,6 @@ npm run perf:summarize -- --artifact /tmp/time-to-answer-baseline.json # 7. assemble the candidate's end-to-end and Stage-5 captures into its composite npx tsx tests/perf/assembleCompositeEvidence.ts \ --general /tmp/time-to-answer-candidate-general.json \ - --output /tmp/time-to-answer-candidate.json \ --stageFive /tmp/time-to-answer-candidate-stage5.json \ --output /tmp/time-to-answer-candidate.json # 8. compare the two composite artifacts (fails closed) From fa7eb85757ce860e939c1b417e0d0960836cf2cf Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 16:11:21 +0800 Subject: [PATCH 08/10] fix(daemon): suppress unchanged Compose publication churn (#336) What: retain a coherent Compose projection for an unchanged Docker observation and compare normalized public mount findings before republishing enrichment. Why: avoid transient Docker-only revisions and equivalent Compose confirmations triggering duplicate browser model work. How checked: npm run test:rust:daemon --- crates/dockermap-daemon/src/cache_refresh.rs | 57 ++++++++++++++++++-- 1 file changed, 54 insertions(+), 3 deletions(-) diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index 45041082..8ba0c0ee 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -418,6 +418,23 @@ fn clear_volatile_observation_markers( } } +/// Compose collection timestamps are observation markers. The projection is +/// otherwise already sanitized by the core derivation, so it is safe to use +/// solely for private publication comparison. +fn compose_mount_finding_projection( + scan: &ComposeScan, + binding: &ComposeRuntimeBinding, +) -> Vec { + let mut findings = derive_compose_runtime_mount_findings(scan, binding); + for finding in &mut findings { + for evidence in &mut finding.evidence_refs { + evidence.collected_at = 0; + } + } + findings.sort_by(|left, right| left.id.cmp(&right.id)); + findings +} + fn boot_instance_component() -> String { static BOOT: OnceLock = OnceLock::new(); BOOT.get_or_init(|| { @@ -643,7 +660,7 @@ impl DaemonCache { fn assign_revision(&mut self) { // All three independently routable model envelopes attest the same // publication. Provider state is runtime-topology evidence only. - let compose_binding_identity = self.compose_binding_identity().map(str::to_owned); + let compose_binding_identity = self.compose_binding_identity(); self.revision.assign( &mut self.snapshot, &mut self.health, @@ -680,10 +697,16 @@ impl DaemonCache { }; } - fn compose_binding_identity(&self) -> Option<&str> { + /// The private binding deliberately never contributes raw Compose inputs + /// to a publication identity. Only its already-sanitized public finding + /// projection is semantic, and observation timestamps are not findings. + fn compose_binding_identity(&self) -> Option { self.compose_runtime_binding .as_ref() - .map(|(_, binding)| binding.provider_revision.as_str()) + .map(|(scan, binding)| { + serde_json::to_string(&compose_mount_finding_projection(scan, binding)) + .expect("public Compose findings are serializable") + }) } fn assign_docker_observation_revision(&mut self) { @@ -807,6 +830,17 @@ fn spawn_compose_projection( { return; } + if cache + .compose_runtime_binding + .as_ref() + .map(|(current_scan, current_binding)| { + compose_mount_finding_projection(current_scan, current_binding) + == compose_mount_finding_projection(&scan, &binding) + }) + .unwrap_or(false) + { + return; + } cache.compose_runtime_binding = Some((scan, binding)); cache.assign_revision(); }); @@ -897,6 +931,23 @@ async fn publish_docker_snapshot_cache( }; updated.rebuild_runtime_map(); updated.revision = cache.revision.clone(); + // An unchanged Docker observation can keep the prior coherent Compose + // projection while a fresh bounded projection confirms it. Retaining it + // avoids a transient Docker-only publication, but only when the private + // Docker-side binding inputs still exactly match as well. + if same_source + && updated.health.mode == RuntimeMode::Docker + && updated.docker_observation_token() == cache.docker_observation_token() + && cache + .compose_runtime_binding + .as_ref() + .map(|(_, binding)| { + compose_containers.as_deref() == Some(binding.containers.as_slice()) + }) + .unwrap_or(false) + { + updated.compose_runtime_binding = cache.compose_runtime_binding.clone(); + } updated.assign_revision(); if updated.health.mode == RuntimeMode::Docker { // The input is publication-sanitized before it becomes retained state. From 4b65bbc7a8fbf5e9de3e8b07a56008411e8925e8 Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 16:12:01 +0800 Subject: [PATCH 09/10] test(daemon): cover coherent Compose revision gating (#336) What: cover retained Compose findings for unchanged Docker confirmation and one coherent revision for changed mount projections. Why: pin the no-churn behavior while ensuring genuine Compose drift remains observable. How checked: npm run fmt:rust; npm run test:rust:daemon --- crates/dockermap-daemon/src/cache_refresh.rs | 103 +++++++++++++++++++ 1 file changed, 103 insertions(+) diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index 8ba0c0ee..7efd1c10 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -3213,6 +3213,109 @@ mod scheduler_tests { cache } + fn compose_binding_fixture(collected_at: u64) -> (ComposeScan, ComposeRuntimeBinding) { + let config_files: BTreeSet = + ["/project/compose.yaml".to_owned()].into_iter().collect(); + let container = ComposeRuntimeContainer { + container_id: "a".repeat(64), + project: "project".into(), + service: "app".into(), + config_files: config_files.clone(), + mounts: Vec::new(), + }; + ( + ComposeScan { + files: vec!["/project/compose.yaml".into()], + project_root: "/project".into(), + services: Vec::new(), + mounts: vec![dockermap_core::ComposeMount { + id: "mount".into(), + service: "app".into(), + kind: ComposeMountKind::Bind, + source: Some("./data".into()), + resolved_source: Some("/project/data".into()), + target: "/data".into(), + read_only: false, + origin: dockermap_core::ComposeFileOrigin { + file: "/project/compose.yaml".into(), + service: Some("app".into()), + field: "volumes".into(), + }, + }], + correlations: Vec::new(), + diagnostics: Vec::new(), + }, + ComposeRuntimeBinding { + project: "project".into(), + config_files, + containers: vec![container], + collected_at, + provider_revision: String::new(), + fresh: true, + }, + ) + } + + #[tokio::test] + async fn unchanged_docker_confirmation_retains_compose_projection_and_revision() { + let mut initial = docker_cache(mock_snapshot()); + initial.rebuild_runtime_map(); + let (scan, mut binding) = compose_binding_fixture(41); + binding.provider_revision = initial.docker_observation_token(); + initial.compose_runtime_binding = Some((scan, binding.clone())); + initial.assign_revision(); + let revision = initial.snapshot.model_revision.clone(); + let finding_time = initial.findings.findings[0].evidence_refs[0].collected_at; + let state = AppState { + allow_mock: true, + cache: Arc::new(RwLock::new(initial)), + docker: Arc::new(RwLock::new(None)), + provider_slot_in_flight: Arc::new(ProviderSlotFlights::default()), + }; + + let mut unchanged = docker_cache(mock_snapshot()); + unchanged.compose_containers = Some(binding.containers.clone()); + publish_docker_snapshot_cache(&state, unchanged).await; + + let cache = state.cache.read().await; + assert_eq!(cache.snapshot.model_revision, revision); + assert_eq!( + cache.findings.findings[0].evidence_refs[0].collected_at, + finding_time + ); + assert!(cache.compose_runtime_binding.is_some()); + } + + #[test] + fn changed_compose_mount_projection_advances_once_with_matching_docker_token() { + let mut cache = docker_cache(mock_snapshot()); + cache.rebuild_runtime_map(); + let (scan, mut binding) = compose_binding_fixture(41); + binding.provider_revision = cache.docker_observation_token(); + cache.compose_runtime_binding = Some((scan.clone(), binding.clone())); + cache.assign_revision(); + let revision = cache.snapshot.model_revision.clone(); + + let mut changed_scan = scan; + changed_scan.mounts.clear(); + assert_ne!( + compose_mount_finding_projection( + &cache.compose_runtime_binding.as_ref().unwrap().0, + &cache.compose_runtime_binding.as_ref().unwrap().1, + ), + compose_mount_finding_projection(&changed_scan, &binding), + ); + cache.compose_runtime_binding = Some((changed_scan, binding)); + cache.assign_revision(); + + assert_ne!(cache.snapshot.model_revision, revision); + assert_eq!(cache.findings.model_revision, cache.snapshot.model_revision); + assert!(cache.findings.findings.iter().all(|finding| { + finding.rule_id + != dockermap_core::FindingRule::ComposeDeclaredMountMissingAtBoundContainer + })); + } + #[tokio::test] async fn snapshot_timeout_invalidates_stalled_client_and_fresh_client_recovers() { async fn read_request_head(connection: &mut tokio::net::UnixStream) -> String { From 09c56153b9e61a04fa579f05e4c87f23f383ccc3 Mon Sep 17 00:00:00 2001 From: Jonathan <64296013+Joncallim@users.noreply.github.com> Date: Sun, 27 Sep 2026 18:59:37 +0800 Subject: [PATCH 10/10] chore(sync): align the promotion tooling with main (#336) What: align the shared promotion scripts and perf TypeScript configuration with main, and make the candidate evidence sequence runnable.\n\nWhy: prevent formatting-only divergence after the promotion-tooling PR merged while retaining #336-specific Compose evidence prose.\n\nHow checked: npm run test:perf; npm run typecheck; origin/main comparisons for package.json and tests/perf/tsconfig.json; git diff --check. --- docs/testing/TIME_TO_ANSWER_EVIDENCE.md | 19 +++++++++++++++++-- package.json | 6 +++--- tests/perf/tsconfig.json | 2 +- 3 files changed, 21 insertions(+), 6 deletions(-) diff --git a/docs/testing/TIME_TO_ANSWER_EVIDENCE.md b/docs/testing/TIME_TO_ANSWER_EVIDENCE.md index 5e95f28f..f1ca4b5f 100644 --- a/docs/testing/TIME_TO_ANSWER_EVIDENCE.md +++ b/docs/testing/TIME_TO_ANSWER_EVIDENCE.md @@ -356,6 +356,7 @@ npm run perf:time-to-answer -- \ --checkpoint # 3. dedicated controlled Stage-5 capture; the only owner of the poll-phase protocol npm run perf:stage-five -- \ + --checkpoint \ --metadata /tmp/time-to-answer-metadata.json \ --output /tmp/time-to-answer-stage5.json \ --raw-dir /tmp/time-to-answer-stage5-raw @@ -373,12 +374,26 @@ npx tsx tests/perf/assembleCompositeEvidence.ts \ --output /tmp/time-to-answer-baseline.json # 6. recompute summaries from the raw samples (never trust supplied aggregates) npm run perf:summarize -- --artifact /tmp/time-to-answer-baseline.json -# 7. assemble the candidate's end-to-end and Stage-5 captures into its composite +# 7. pin the candidate environment from the candidate checkout +npm run perf:metadata -- --output /tmp/time-to-answer-candidate-metadata.json +# 8. capture the candidate's end-to-end section at +npm run perf:time-to-answer -- \ + --metadata /tmp/time-to-answer-candidate-metadata.json \ + --output /tmp/time-to-answer-candidate-general.json \ + --raw-dir /tmp/time-to-answer-candidate-raw \ + --checkpoint +# 9. capture the candidate's dedicated Stage-5 section at +npm run perf:stage-five -- \ + --metadata /tmp/time-to-answer-candidate-metadata.json \ + --output /tmp/time-to-answer-candidate-stage5.json \ + --raw-dir /tmp/time-to-answer-candidate-stage5-raw \ + --checkpoint +# 10. assemble the candidate's captures into its composite npx tsx tests/perf/assembleCompositeEvidence.ts \ --general /tmp/time-to-answer-candidate-general.json \ --stageFive /tmp/time-to-answer-candidate-stage5.json \ --output /tmp/time-to-answer-candidate.json -# 8. compare the two composite artifacts (fails closed) +# 11. compare the two composite artifacts (fails closed) npm run perf:promote -- \ --baseline /tmp/time-to-answer-baseline.json \ --candidate /tmp/time-to-answer-candidate.json diff --git a/package.json b/package.json index 86f457ef..48a28cd4 100644 --- a/package.json +++ b/package.json @@ -31,9 +31,9 @@ "test:contracts": "npm run test --workspace @dockermap/contracts --if-present", "test:version": "node --test scripts/check-version-authority.test.mjs scripts/package-release.test.mjs", "test:perf": "node --test tests/perf/*.test.mjs", -"perf:time-to-answer": "tsx tests/perf/capture.ts", -"perf:promote": "tsx tests/perf/promote.ts", -"perf:calibrate-time-to-answer": "tsx tests/perf/capture.ts", + "perf:time-to-answer": "tsx tests/perf/capture.ts", + "perf:promote": "tsx tests/perf/promote.ts", + "perf:calibrate-time-to-answer": "tsx tests/perf/capture.ts", "perf:preconditioning": "tsx tests/perf/preconditioning.ts", "perf:phase-control": "tsx tests/perf/phaseControl.ts", "perf:stage-five": "tsx tests/perf/captureStageFive.ts", diff --git a/tests/perf/tsconfig.json b/tests/perf/tsconfig.json index aa312331..e72b93c5 100644 --- a/tests/perf/tsconfig.json +++ b/tests/perf/tsconfig.json @@ -4,5 +4,5 @@ "types": ["node"], "allowJs": true }, -"include": ["capture.ts", "captureIndependence.ts", "summarize.ts", "preconditioning.ts", "phaseControl.ts", "promote.ts"] + "include": ["capture.ts", "captureIndependence.ts", "summarize.ts", "preconditioning.ts", "phaseControl.ts", "promote.ts"] }