diff --git a/crates/dockermap-daemon/src/cache_refresh.rs b/crates/dockermap-daemon/src/cache_refresh.rs index df3422a3..7efd1c10 100644 --- a/crates/dockermap-daemon/src/cache_refresh.rs +++ b/crates/dockermap-daemon/src/cache_refresh.rs @@ -83,11 +83,16 @@ 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, pub(crate) health: HealthResponse, 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 @@ -265,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 @@ -285,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()) { @@ -411,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(|| { @@ -607,6 +631,7 @@ impl DaemonCache { let last_updated = snapshot.last_updated; let mut cache = Self { + publication_ready: false, snapshot, health, runtime_map: RuntimeMap { @@ -619,6 +644,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(), @@ -634,8 +660,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(); + 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(); @@ -666,6 +697,18 @@ impl DaemonCache { }; } + /// 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(|(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) { self.docker_observation_revision .assign(&self.snapshot, &self.health.mode); @@ -720,7 +763,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, ) @@ -730,11 +780,72 @@ 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; + } + 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(); + }); +} + fn spawn_provider_slots( state: AppState, snapshot: DockerSnapshot, @@ -773,7 +884,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). @@ -786,6 +904,10 @@ async fn publish_docker_snapshot_cache( .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; + let compose_containers = updated.compose_containers.take(); updated.runtime_providers = if same_source { cache.runtime_providers.clone() } else { @@ -809,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. @@ -825,6 +964,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, ) } @@ -858,8 +1000,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( @@ -869,7 +1011,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. @@ -886,46 +1027,13 @@ 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); - 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, @@ -937,11 +1045,13 @@ where message: Some("Docker engine connected".into()), }; Ok(DaemonCache { + publication_ready: false, snapshot, health, 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(), @@ -3077,6 +3187,7 @@ mod scheduler_tests { fn docker_cache(snapshot: DockerSnapshot) -> DaemonCache { let last_updated = snapshot.last_updated; let mut cache = DaemonCache { + publication_ready: true, snapshot, health: HealthResponse { status: HealthState::Ok, @@ -3090,6 +3201,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(), @@ -3101,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 { @@ -3131,48 +3346,62 @@ mod scheduler_tests { 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 }); @@ -3194,8 +3423,14 @@ mod scheduler_tests { 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(), @@ -3239,7 +3474,6 @@ mod scheduler_tests { 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", @@ -3307,9 +3541,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(); @@ -3317,30 +3548,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(); @@ -3360,26 +3578,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); } @@ -4267,7 +4471,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; @@ -4336,14 +4540,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); @@ -4578,7 +4782,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; diff --git a/crates/dockermap-daemon/src/daemon_api.rs b/crates/dockermap-daemon/src/daemon_api.rs index 4b7bc83a..c076feea 100644 --- a/crates/dockermap-daemon/src/daemon_api.rs +++ b/crates/dockermap-daemon/src/daemon_api.rs @@ -85,6 +85,12 @@ 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 { return Err(ApiError { status: StatusCode::SERVICE_UNAVAILABLE, @@ -423,13 +429,14 @@ mod tests { 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" + 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 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 a309054c..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_cache, 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, @@ -137,15 +137,15 @@ async fn main() { 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()); - - refresh_cache(&state).await; - tokio::spawn(refresh_loop(state.clone())); - - let app = daemon_router(state, daemon_token); + 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)); + println!("dockermap-daemon listening on http://{address}"); axum::serve(listener, app) @@ -180,13 +180,16 @@ 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(DaemonCache::mock())), + cache: Arc::new(RwLock::new(cache)), docker: Arc::new(RwLock::new(None)), provider_slot_in_flight: Arc::new(crate::cache_refresh::ProviderSlotFlights::default()), } @@ -228,6 +231,28 @@ mod tests { } } + #[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)) @@ -353,6 +378,7 @@ mod tests { }); let mut cache = DaemonCache::mock(); + cache.publication_ready = true; cache.health.docker_reachable = true; let state = AppState { allow_mock: true, @@ -499,6 +525,84 @@ mod tests { ], "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] 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..f1ca4b5f 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 @@ -358,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 @@ -375,11 +374,29 @@ 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) +# 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-metadata.json \ - --output /tmp/time-to-answer-candidate.json \ - --baseline /tmp/time-to-answer-baseline.json + --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 +# 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 ``` Prerequisites: a release daemon (`cargo build --release -p dockermap-daemon`), diff --git a/package.json b/package.json index 126e17b1..48a28cd4 100644 --- a/package.json +++ b/package.json @@ -32,6 +32,7 @@ "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:preconditioning": "tsx tests/perf/preconditioning.ts", "perf:phase-control": "tsx tests/perf/phaseControl.ts", diff --git a/tests/perf/capture.ts b/tests/perf/capture.ts index 6cf5573a..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)); @@ -1256,7 +1257,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..20475f92 100644 --- a/tests/perf/captureSyntax.test.mjs +++ b/tests/perf/captureSyntax.test.mjs @@ -18,3 +18,10 @@ test("capture harness entrypoint transforms", async () => { throw new Error(`capture.ts failed esbuild transform:\n${details || error.message}`); } }); + +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("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(); +} +}); 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..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"] + "include": ["capture.ts", "captureIndependence.ts", "summarize.ts", "preconditioning.ts", "phaseControl.ts", "promote.ts"] }