Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
444 changes: 324 additions & 120 deletions crates/dockermap-daemon/src/cache_refresh.rs

Large diffs are not rendered by default.

13 changes: 10 additions & 3 deletions crates/dockermap-daemon/src/daemon_api.rs
Original file line number Diff line number Diff line change
Expand Up @@ -85,6 +85,12 @@ pub(crate) fn daemon_router(state: AppState, daemon_token: DaemonAuthToken) -> R

async fn publication_cache(state: &AppState) -> Result<RwLockReadGuard<'_, DaemonCache>, 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,
Expand Down Expand Up @@ -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
Expand Down
29 changes: 12 additions & 17 deletions crates/dockermap-daemon/src/docker_collector.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down
118 changes: 111 additions & 7 deletions crates/dockermap-daemon/src/main.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down Expand Up @@ -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)
Expand Down Expand Up @@ -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()),
}
Expand Down Expand Up @@ -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))
Expand Down Expand Up @@ -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,
Expand Down Expand Up @@ -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]
Expand Down
9 changes: 9 additions & 0 deletions docs/architecture/ARCHITECTURE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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`.
Expand Down
43 changes: 30 additions & 13 deletions docs/testing/TIME_TO_ANSWER_EVIDENCE.md
Original file line number Diff line number Diff line change
Expand Up @@ -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.
Expand Down Expand Up @@ -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

Expand Down Expand Up @@ -358,6 +356,7 @@ npm run perf:time-to-answer -- \
--checkpoint <commit>
# 3. dedicated controlled Stage-5 capture; the only owner of the poll-phase protocol
npm run perf:stage-five -- \
--checkpoint <commit> \
--metadata /tmp/time-to-answer-metadata.json \
--output /tmp/time-to-answer-stage5.json \
--raw-dir /tmp/time-to-answer-stage5-raw
Expand All @@ -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 <candidate-commit>
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 <candidate-commit>
# 9. capture the candidate's dedicated Stage-5 section at <candidate-commit>
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 <candidate-commit>
# 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`),
Expand Down
1 change: 1 addition & 0 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -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",
Expand Down
3 changes: 2 additions & 1 deletion tests/perf/capture.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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));
Expand Down Expand Up @@ -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) }),
Expand Down
7 changes: 7 additions & 0 deletions tests/perf/captureSyntax.test.mjs
Original file line number Diff line number Diff line change
Expand Up @@ -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");
}
});
13 changes: 13 additions & 0 deletions tests/perf/listenerReadiness.mjs
Original file line number Diff line number Diff line change
@@ -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);
}
}
Loading
Loading