From 7cd8dbdb14c6542af8763a3de2b7f010ae7b089e Mon Sep 17 00:00:00 2001 From: Bo Wu Date: Sat, 5 Sep 2026 18:24:35 -0700 Subject: [PATCH 1/5] Integrate durable trace commitments with Logger --- .github/workflows/container.yml | 17 ++ control-server/src/server.mjs | 7 +- docker/runtime/Dockerfile | 3 +- docs/TODO.md | 7 +- gitops/README.md | 6 + logger/README.md | 6 + logger/src/main.rs | 2 +- logger/src/server.rs | 96 +++++++- tests/run.sh | 1 + tests/test_trace_exporter.py | 120 ++++++++++ trace_exporter/README.md | 27 +++ trace_exporter/__init__.py | 1 + trace_exporter/trace_exporter.py | 377 +++++++++++++++++++++++++++++++ 13 files changed, 656 insertions(+), 14 deletions(-) create mode 100644 tests/test_trace_exporter.py create mode 100644 trace_exporter/README.md create mode 100644 trace_exporter/__init__.py create mode 100644 trace_exporter/trace_exporter.py diff --git a/.github/workflows/container.yml b/.github/workflows/container.yml index cabc6f4..ef01123 100644 --- a/.github/workflows/container.yml +++ b/.github/workflows/container.yml @@ -56,6 +56,23 @@ jobs: platforms: linux/amd64 provenance: mode=max sbom: true + - uses: docker/metadata-action@c1e51972afc2121e065aed6d45c65596fe445f3f # v5 + id: logger-metadata + with: + images: ghcr.io/${{ github.repository_owner }}/multiagent-logger + tags: | + type=raw,value=main,enable={{is_default_branch}} + type=sha,prefix=sha- + - uses: docker/build-push-action@263435318d21b8e681c14492fe198d362a7d2c83 # v6 + with: + context: . + file: docker/logger/Dockerfile + push: true + tags: ${{ steps.logger-metadata.outputs.tags }} + labels: ${{ steps.logger-metadata.outputs.labels }} + platforms: linux/amd64 + provenance: mode=max + sbom: true publish-wiki: runs-on: ubuntu-latest diff --git a/control-server/src/server.mjs b/control-server/src/server.mjs index 51bcbbf..10a0192 100644 --- a/control-server/src/server.mjs +++ b/control-server/src/server.mjs @@ -931,10 +931,15 @@ function traceExportStatus() { lastAttemptAt: value.lastAttemptAt, lastSuccessAt: value.lastSuccessAt || null, fileCount: Number(value.fileCount || 0), + uploadedCount: Number(value.uploadedCount || 0), + loggerOk: value.loggerOk === true, + loggerPending: Number(value.loggerPending || 0), + loggerLastSuccessAt: value.loggerLastSuccessAt || null, + loggerError: value.loggerError || null, ageSeconds: Number.isFinite(ageSeconds) ? Math.floor(ageSeconds) : null, }; } catch { - return { configured: true, ready: false, ok: false, lastAttemptAt: null, lastSuccessAt: null, fileCount: 0, ageSeconds: null }; + return { configured: true, ready: false, ok: false, lastAttemptAt: null, lastSuccessAt: null, fileCount: 0, uploadedCount: 0, loggerOk: false, loggerPending: 0, loggerLastSuccessAt: null, loggerError: null, ageSeconds: null }; } } diff --git a/docker/runtime/Dockerfile b/docker/runtime/Dockerfile index 50a7f98..288e78e 100644 --- a/docker/runtime/Dockerfile +++ b/docker/runtime/Dockerfile @@ -20,9 +20,10 @@ RUN cd control-server && npm ci --omit=dev COPY . . COPY --from=multiagent-builder /src/target/release/multiagent /opt/multiagent/bin/multiagent RUN ln -s /opt/multiagent/wiki-service/bin/wiki-query.mjs /usr/local/bin/wiki-query \ + && ln -s /opt/multiagent/trace_exporter/trace_exporter.py /usr/local/bin/trace-exporter \ && install -m 0755 docker/runtime/container-entrypoint.sh /opt/multiagent/bin/container-entrypoint.sh \ && install -m 0755 control-server/bin/prepare-repository.mjs /opt/multiagent/bin/prepare-repository.mjs \ - && chmod +x launch.sh control-server/bin/*.mjs \ + && chmod +x launch.sh control-server/bin/*.mjs trace_exporter/trace_exporter.py \ && groupadd --gid 10000 multiagent-control \ && groupadd --gid 10001 multiagent-role \ && groupadd --gid 10004 multiagent-credentials \ diff --git a/docs/TODO.md b/docs/TODO.md index c427ccc..7eac650 100644 --- a/docs/TODO.md +++ b/docs/TODO.md @@ -33,8 +33,11 @@ applicable, and relevant evidence are complete. restart tests for newline-aligned tail truncation and missing checkpoints. This does not require producer-assigned sequence numbers; the Logger remains the sole sequencer. -- [ ] Integrate producer outboxes or a deployment-owned durable queue so Logger - delivery retries independently and backlog alerts are testable. +- [ ] Deploy and prove the metadata-only S3 trace commitment outbox so delivery + retries independently across exporter restarts, backlog is observable, and a + production Logger outage/recovery drains without duplicate ledger entries. + Additional structural-event producers must adopt the same durable delivery + property before they are enabled in production. - [ ] Add deployment-owned Loki/OpenTelemetry projections if operational demand justifies them; these must remain derived from the authoritative ledger. diff --git a/gitops/README.md b/gitops/README.md index 921ccfb..d101ba8 100644 --- a/gitops/README.md +++ b/gitops/README.md @@ -28,6 +28,12 @@ The application-owned Logger deployment contract is: - give producers a durable retry path or outbox and alert on delivery backlog; - keep the existing trace sidecar and S3 data path, then submit a bounded `trace.artifact_exported` commitment after a successful upload; +- run `trace-exporter` with a trace-commitment-only token. It stores only + bounded commitment JSON in the deployment-owned S3 outbox, uses deterministic + event IDs and delivered markers for restart-safe idempotency, and reports + `loggerPending`, `loggerOk`, and Logger delivery timestamps in its status + file. Logger failure must not change the S3 export `ok` result or workflow + progression; - configure at most one active Logger replica for a ledger volume. A standby must not write until deployment fencing has transferred ownership. diff --git a/logger/README.md b/logger/README.md index ca3962f..1badcea 100644 --- a/logger/README.md +++ b/logger/README.md @@ -54,6 +54,12 @@ cargo run -p multiagent-logger -- submit-trace-commitment \ --media-type application/gzip ``` +The helper emits `trace.artifact_exported`. Its payload contains only the +artifact digest, byte size, media type, and storage reference; it never sends +the trace body. Use a deterministic event ID and retain the event in a durable +outbox until the Logger returns `204`, because acknowledgement remains +idempotent transport evidence rather than workflow authority. + ## API and configuration The API exposes event append, log heads/entries/checkpoints, the public key, diff --git a/logger/src/main.rs b/logger/src/main.rs index 9c1bcbd..a6e6df2 100644 --- a/logger/src/main.rs +++ b/logger/src/main.rs @@ -100,7 +100,7 @@ async fn submit(trace: bool, args: Vec) -> Result<(), String> { Event { event_id: event_id.ok_or("--event-id is required")?, session_id: session.ok_or("--session-id is required")?, - event_type: "trace.commitment".into(), + event_type: "trace.artifact_exported".into(), payload_digest: digest.clone(), artifact_references: vec![ArtifactReference { uri: storage_reference.ok_or("--storage-reference is required")?, diff --git a/logger/src/server.rs b/logger/src/server.rs index 2d34e1e..7015c0f 100644 --- a/logger/src/server.rs +++ b/logger/src/server.rs @@ -429,14 +429,21 @@ mod tests { .unwrap(); let key_file = directory.path().join("signing-key.pem"); fs::write(&key_file, key.as_bytes()).unwrap(); - let token = "test-token-0123456789abcdef"; + let trace_token = "trace-token-0123456789abcdef"; + let reader_token = "reader-token-0123456789abcdef"; let clients_file = directory.path().join("clients.json"); fs::write( &clients_file, serde_json::to_vec(&json!({"clients":[{ - "id":"test-client", - "tokenSha256":format!("sha256:{:x}", Sha256::digest(token)), - "permissions":["append","read","verify"], + "id":"trace-producer", + "tokenSha256":format!("sha256:{:x}", Sha256::digest(trace_token)), + "permissions":["append"], + "eventTypes":["trace.artifact_exported"], + "sessions":["session-*"] + },{ + "id":"audit-reader", + "tokenSha256":format!("sha256:{:x}", Sha256::digest(reader_token)), + "permissions":["read","verify"], "eventTypes":["*"], "sessions":["session-*"] }]})) @@ -464,11 +471,16 @@ mod tests { async fn append_and_read_require_scoped_auth_and_return_no_receipt() { let (_directory, app) = application(); let event = json!({ - "eventId":"event-1", + "eventId":"trace-export-1111111111111111111111111111111111111111111111111111111111111111", "sessionId":"session-1", - "eventType":"reviewer.verdict", + "eventType":"trace.artifact_exported", "payloadDigest":format!("sha256:{}", "1".repeat(64)), - "artifactReferences":[] + "artifactReferences":[{ + "uri":"s3://audit/production/artifacts/session-1/trace.jsonl", + "digest":format!("sha256:{}", "1".repeat(64)), + "size":42, + "mediaType":"application/jsonl" + }] }); let unauthorized = app .clone() @@ -486,7 +498,7 @@ mod tests { .clone() .oneshot( Request::post("/v1/events") - .header("authorization", "Bearer test-token-0123456789abcdef") + .header("authorization", "Bearer trace-token-0123456789abcdef") .header("content-type", "application/json") .body(Body::from(event.to_string())) .unwrap(), @@ -503,10 +515,36 @@ mod tests { 0 ); + let duplicate = app + .clone() + .oneshot( + Request::post("/v1/events") + .header("authorization", "Bearer trace-token-0123456789abcdef") + .header("content-type", "application/json") + .body(Body::from(event.to_string())) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(duplicate.status(), StatusCode::NO_CONTENT); + + let producer_cannot_read = app + .clone() + .oneshot( + Request::get("/v1/logs/session-1/head") + .header("authorization", "Bearer trace-token-0123456789abcdef") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(producer_cannot_read.status(), StatusCode::FORBIDDEN); + let head = app + .clone() .oneshot( Request::get("/v1/logs/session-1/head") - .header("authorization", "Bearer test-token-0123456789abcdef") + .header("authorization", "Bearer reader-token-0123456789abcdef") .body(Body::empty()) .unwrap(), ) @@ -518,5 +556,45 @@ mod tests { serde_json::from_slice::(&body).unwrap()["sequence"], 1 ); + + let checkpoints = app + .clone() + .oneshot( + Request::get("/v1/logs/session-1/checkpoints") + .header("authorization", "Bearer reader-token-0123456789abcdef") + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(checkpoints.status(), StatusCode::OK); + let body = axum::body::to_bytes(checkpoints.into_body(), 16_384) + .await + .unwrap(); + let decoded: Value = serde_json::from_slice(&body).unwrap(); + assert_eq!(decoded["checkpoints"].as_array().unwrap().len(), 1); + assert_eq!( + decoded["checkpoints"][0]["loggerSignature"]["algorithm"], + "Ed25519" + ); + + let verified = app + .oneshot( + Request::post("/v1/verify") + .header("authorization", "Bearer reader-token-0123456789abcdef") + .header("content-type", "application/json") + .body(Body::from(r#"{"logId":"session-1"}"#)) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(verified.status(), StatusCode::OK); + let body = axum::body::to_bytes(verified.into_body(), 4096) + .await + .unwrap(); + let decoded: Value = serde_json::from_slice(&body).unwrap(); + assert_eq!(decoded["ok"], true); + assert_eq!(decoded["checkedEntries"], 1); + assert_eq!(decoded["checkedCheckpoints"], 1); } } diff --git a/tests/run.sh b/tests/run.sh index f24814f..9d39681 100755 --- a/tests/run.sh +++ b/tests/run.sh @@ -1247,6 +1247,7 @@ PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_native_s PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_swe_outcomes.py" PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_swe_provenance.py" PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_migration_contracts.py" +PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_trace_exporter.py" python3 -m evaluation.swe_bench_pro --help >"$TMPDIR/swe-bench-pro-help.out" assert_file_contains "$TMPDIR/swe-bench-pro-help.out" "Evaluate the production multiagent solver" assert_file_not_contains "$TMPDIR/swe-bench-pro-help.out" "--agent-framework" diff --git a/tests/test_trace_exporter.py b/tests/test_trace_exporter.py new file mode 100644 index 0000000..a299361 --- /dev/null +++ b/tests/test_trace_exporter.py @@ -0,0 +1,120 @@ +#!/usr/bin/env python3 + +import json +import tempfile +import unittest +from pathlib import Path + +from trace_exporter.trace_exporter import TraceExporter + + +class MemoryS3: + def __init__(self): + self.objects = {} + + def upload(self, source, uri): + self.objects[uri] = Path(source).read_bytes() + + def copy(self, source_uri, destination_uri): + self.objects[destination_uri] = self.objects[source_uri] + + def download(self, uri, destination): + destination.parent.mkdir(parents=True, exist_ok=True) + destination.write_bytes(self.objects[uri]) + + def list_keys(self, bucket, prefix): + root = f"s3://{bucket}/" + return {uri.removeprefix(root) for uri in self.objects if uri.startswith(root + prefix)} + + +class TraceExporterTest(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory() + root = Path(self.temporary.name) + self.source = root / "source" + self.source.mkdir() + self.work = root / "work" + self.status = root / "status.json" + self.token = root / "token" + self.token.write_text("trace-token-0123456789abcdef", encoding="utf-8") + self.s3 = MemoryS3() + self.posts = [] + + def tearDown(self): + self.temporary.cleanup() + + def exporter(self, post=None): + def successful(url, token, body): + self.posts.append((url, token, json.loads(body))) + return 204 + + return TraceExporter( + source=self.source, + work=self.work, + status_file=self.status, + bucket="trace-bucket", + trace_prefix="production/sessions/session-1", + outbox_prefix="production/logger-outbox", + delivered_prefix="production/logger-delivered", + session_id="session-1", + logger_url="http://logger", + logger_token_file=self.token, + s3=self.s3, + http_post=post or successful, + clock=lambda: "2026-09-05T12:00:00Z", + ) + + def test_successful_export_commits_digest_reference_and_no_body(self): + body = b"private trace body" + (self.source / "trace.jsonl").write_bytes(body) + + status = self.exporter().run_once() + + self.assertTrue(status["ok"]) + self.assertTrue(status["loggerOk"]) + self.assertEqual(status["loggerPending"], 0) + self.assertEqual(len(self.posts), 1) + event = self.posts[0][2] + self.assertEqual(event["eventType"], "trace.artifact_exported") + self.assertEqual(event["sessionId"], "session-1") + self.assertEqual(event["payloadDigest"], event["artifactReferences"][0]["digest"]) + self.assertEqual(event["artifactReferences"][0]["size"], len(body)) + self.assertNotIn(body, self.posts[0][2].__str__().encode()) + artifact_uri = event["artifactReferences"][0]["uri"] + self.assertEqual(self.s3.objects[artifact_uri], body) + + def test_delivered_marker_makes_restart_idempotent(self): + (self.source / "usage.json").write_text("{}", encoding="utf-8") + first = self.exporter() + first.run_once() + self.assertEqual(len(self.posts), 1) + # A fresh local cache re-uploads metadata, but the durable delivered + # marker prevents another Logger request. + for path in self.work.iterdir(): + if path.is_dir(): + import shutil + + shutil.rmtree(path) + self.exporter().run_once() + self.assertEqual(len(self.posts), 1) + + def test_logger_outage_leaves_durable_backlog_and_later_drains(self): + (self.source / "report.md").write_text("report", encoding="utf-8") + + def unavailable(_url, _token, _body): + raise OSError("temporary Logger outage") + + failed = self.exporter(unavailable).run_once() + self.assertTrue(failed["ok"], "Logger availability must not change S3 export success") + self.assertFalse(failed["loggerOk"]) + self.assertEqual(failed["loggerPending"], 1) + self.assertIn("temporary Logger outage", failed["loggerError"]) + + recovered = self.exporter().run_once() + self.assertTrue(recovered["loggerOk"]) + self.assertEqual(recovered["loggerPending"], 0) + self.assertEqual(len(self.posts), 1) + + +if __name__ == "__main__": + unittest.main() diff --git a/trace_exporter/README.md b/trace_exporter/README.md new file mode 100644 index 0000000..21e18a0 --- /dev/null +++ b/trace_exporter/README.md @@ -0,0 +1,27 @@ +# Trace exporter + +`python3 -m trace_exporter.trace_exporter` copies regular trace files to the +deployment-owned S3 trace bucket, preserves a content-addressed object for each +digest, and writes a bounded metadata-only Logger event to an S3 outbox. It +replays outbox objects until Logger acknowledges them, then writes an immutable +delivered marker. Deterministic event IDs make a retry after an uncertain +acknowledgement safe against Logger's exact idempotency contract. + +Required environment: + +| Variable | Purpose | +| --- | --- | +| `TRACE_SOURCE` | Read-only trace tree | +| `TRACE_BUCKET` | Deployment-owned S3 bucket | +| `TRACE_PREFIX` | Conventional and content-addressed trace prefix | +| `TRACE_SESSION_ID` | Logger session/log ID | +| `TRACE_OUTBOX_PREFIX` | Metadata-only durable outbox prefix | +| `TRACE_DELIVERED_PREFIX` | Durable delivery-marker prefix | +| `LOGGER_URL` | Private Logger base URL | +| `LOGGER_TOKEN_FILE` | Trace-commitment-only bearer token file | + +`TRACE_EXPORT_WORK_DIR`, `TRACE_EXPORT_STATUS_FILE`, and +`TRACE_EXPORT_INTERVAL_SECONDS` are optional. The status contract keeps S3 +success in `ok` and reports Logger delivery separately through `loggerOk`, +`loggerPending`, and Logger timestamps/errors. Consumers must never use Logger +availability or acknowledgement to grant or deny workflow transitions. diff --git a/trace_exporter/__init__.py b/trace_exporter/__init__.py new file mode 100644 index 0000000..32d6f66 --- /dev/null +++ b/trace_exporter/__init__.py @@ -0,0 +1 @@ +"""Durable S3 trace export and Logger commitment delivery.""" diff --git a/trace_exporter/trace_exporter.py b/trace_exporter/trace_exporter.py new file mode 100644 index 0000000..d0fd582 --- /dev/null +++ b/trace_exporter/trace_exporter.py @@ -0,0 +1,377 @@ +#!/usr/bin/env python3 +"""Export trace files to S3 and durably commit metadata to Logger. + +The S3 outbox is deliberately metadata-only. Trace bodies are uploaded to a +content-addressed artifact key and never sent to Logger. +""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import mimetypes +import os +import shutil +import signal +import stat +import subprocess +import tempfile +import threading +import urllib.error +import urllib.request +from urllib.parse import quote +from datetime import datetime, timezone +from pathlib import Path, PurePosixPath +from typing import Callable, Iterable + + +def utc_now() -> str: + return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") + + +def sha256_file(path: Path) -> tuple[str, int]: + digest = hashlib.sha256() + size = 0 + with path.open("rb") as handle: + while chunk := handle.read(1024 * 1024): + digest.update(chunk) + size += len(chunk) + return f"sha256:{digest.hexdigest()}", size + + +def safe_segment(value: str, name: str) -> str: + if ( + not value + or len(value) > 128 + or not value[0].isascii() + or not value[0].isalnum() + or any( + not character.isascii() + or not (character.isalnum() or character in "._:-") + for character in value + ) + ): + raise ValueError(f"{name} has an invalid format") + return value + + +def normalized_prefix(value: str, name: str) -> str: + parts = PurePosixPath(value.strip("/")).parts + if ( + not parts + or value != value.strip("/") + or any(part in ("", ".", "..") for part in parts) + or any(character in value for character in "\r\n?#\\") + ): + raise ValueError(f"{name} has an invalid format") + return "/".join(parts) + + +class AwsS3: + def __init__(self, runner: Callable[..., subprocess.CompletedProcess[str]] | None = None): + self.runner = runner or subprocess.run + + def _run(self, args: list[str]) -> subprocess.CompletedProcess[str]: + return self.runner(args, check=True, capture_output=True, text=True) + + def upload(self, source: Path, uri: str) -> None: + self._run(["aws", "s3", "cp", str(source), uri, "--only-show-errors"]) + + def copy(self, source_uri: str, destination_uri: str) -> None: + self._run(["aws", "s3", "cp", source_uri, destination_uri, "--only-show-errors"]) + + def download(self, uri: str, destination: Path) -> None: + destination.parent.mkdir(parents=True, exist_ok=True) + self._run(["aws", "s3", "cp", uri, str(destination), "--only-show-errors"]) + + def list_keys(self, bucket: str, prefix: str) -> set[str]: + result = self._run( + [ + "aws", + "s3api", + "list-objects-v2", + "--bucket", + bucket, + "--prefix", + prefix, + "--output", + "json", + ] + ) + decoded = json.loads(result.stdout or "{}") + return {item["Key"] for item in decoded.get("Contents", [])} + + +class TraceExporter: + def __init__( + self, + *, + source: Path, + work: Path, + status_file: Path, + bucket: str, + trace_prefix: str, + outbox_prefix: str, + delivered_prefix: str, + session_id: str, + logger_url: str, + logger_token_file: Path, + s3: AwsS3 | None = None, + http_post: Callable[[str, str, bytes], int] | None = None, + clock: Callable[[], str] = utc_now, + ): + self.source = source + self.work = work + self.status_file = status_file + self.bucket = bucket + self.trace_prefix = normalized_prefix(trace_prefix, "trace prefix") + self.outbox_prefix = normalized_prefix(outbox_prefix, "outbox prefix") + self.delivered_prefix = normalized_prefix(delivered_prefix, "delivered prefix") + self.session_id = safe_segment(session_id, "session ID") + self.logger_url = logger_url.rstrip("/") + self.logger_token_file = logger_token_file + self.s3 = s3 or AwsS3() + self.http_post = http_post or self._post + self.clock = clock + self.cache = work / "cache" + self.staging = work / "staging" + self.queue_cache = work / "queue" + + def run_once(self) -> dict[str, object]: + attempted_at = self.clock() + uploaded = 0 + scanned = 0 + export_error: str | None = None + logger_error: str | None = None + self.work.mkdir(parents=True, exist_ok=True) + self.cache.mkdir(parents=True, exist_ok=True) + shutil.rmtree(self.staging, ignore_errors=True) + self.staging.mkdir(parents=True, exist_ok=True) + + try: + for source_file, relative in self._source_files(): + scanned += 1 + staged = self.staging / relative + staged.parent.mkdir(parents=True, exist_ok=True) + try: + shutil.copy2(source_file, staged) + except (FileNotFoundError, PermissionError): + continue + digest, size = sha256_file(staged) + cached = self.cache / relative + if cached.is_file() and sha256_file(cached) == (digest, size): + continue + media_type = mimetypes.guess_type(relative.as_posix())[0] or "application/octet-stream" + encoded_relative = quote(relative.as_posix(), safe="/") + canonical_key = ( + f"{self.trace_prefix}/artifacts/{self.session_id}/" + f"{digest.removeprefix('sha256:')}/{encoded_relative}" + ) + canonical_uri = f"s3://{self.bucket}/{canonical_key}" + current_uri = f"s3://{self.bucket}/{self.trace_prefix}/{relative.as_posix()}" + self.s3.upload(staged, canonical_uri) + self.s3.copy(canonical_uri, current_uri) + event = self._event(canonical_uri, digest, size, media_type) + self._enqueue(event) + cached.parent.mkdir(parents=True, exist_ok=True) + os.replace(staged, cached) + uploaded += 1 + except Exception as error: # errors are reported through the shared status contract + export_error = str(error) + + pending = 0 + delivered = 0 + try: + pending, delivered = self._deliver_pending() + except Exception as error: + logger_error = str(error) + try: + pending = len(self._pending_keys()) + except Exception: + pending = max(pending, 1) + + previous = self._read_status() + status: dict[str, object] = { + "ok": export_error is None, + "lastAttemptAt": attempted_at, + "fileCount": scanned, + "uploadedCount": uploaded, + "loggerOk": logger_error is None and pending == 0, + "loggerPending": pending, + "loggerDeliveredCount": delivered, + } + if export_error is None: + status["lastSuccessAt"] = attempted_at + elif previous.get("lastSuccessAt"): + status["lastSuccessAt"] = previous["lastSuccessAt"] + status["error"] = export_error + else: + status["error"] = export_error + if logger_error is None and pending == 0: + status["loggerLastSuccessAt"] = attempted_at + else: + if previous.get("loggerLastSuccessAt"): + status["loggerLastSuccessAt"] = previous["loggerLastSuccessAt"] + if logger_error: + status["loggerError"] = logger_error[:512] + self._write_json_atomic(self.status_file, status) + return status + + def _source_files(self) -> Iterable[tuple[Path, Path]]: + if not self.source.is_dir(): + return + for root, directories, files in os.walk(self.source, followlinks=False): + directories[:] = sorted(name for name in directories if name != "worktrees") + root_path = Path(root) + for name in sorted(files): + if name == "orchestrator-bootstrap.sh": + continue + path = root_path / name + try: + if not stat.S_ISREG(path.lstat().st_mode): + continue + except (FileNotFoundError, PermissionError, OSError): + continue + yield path, path.relative_to(self.source) + + def _event(self, uri: str, digest: str, size: int, media_type: str) -> dict[str, object]: + identity = hashlib.sha256(f"{self.session_id}\0{uri}\0{digest}".encode()).hexdigest() + return { + "eventId": f"trace-export-{identity}", + "sessionId": self.session_id, + "eventType": "trace.artifact_exported", + "payloadDigest": digest, + "artifactReferences": [ + {"uri": uri, "digest": digest, "size": size, "mediaType": media_type} + ], + } + + def _enqueue(self, event: dict[str, object]) -> None: + event_id = str(event["eventId"]) + local = self.queue_cache / "outbox" / f"{event_id}.json" + self._write_json_atomic(local, event) + self.s3.upload(local, self._outbox_uri(event_id)) + + def _pending_keys(self) -> list[str]: + outbox_root = f"{self.outbox_prefix}/{self.session_id}/" + delivered_root = f"{self.delivered_prefix}/{self.session_id}/" + outbox = self.s3.list_keys(self.bucket, outbox_root) + delivered = self.s3.list_keys(self.bucket, delivered_root) + delivered_ids = {Path(key).stem for key in delivered if key.endswith(".json")} + return sorted( + key for key in outbox if key.endswith(".json") and Path(key).stem not in delivered_ids + ) + + def _deliver_pending(self) -> tuple[int, int]: + pending_keys = self._pending_keys() + delivered = 0 + token = self.logger_token_file.read_text(encoding="utf-8").strip() + if not token: + raise ValueError("Logger token file is empty") + for key in pending_keys: + event_id = Path(key).stem + local = self.queue_cache / "download" / f"{event_id}.json" + self.s3.download(f"s3://{self.bucket}/{key}", local) + body = local.read_bytes() + status = self.http_post(f"{self.logger_url}/v1/events", token, body) + if status != 204: + raise RuntimeError(f"Logger returned HTTP {status}") + marker = self.queue_cache / "delivered" / f"{event_id}.json" + self._write_json_atomic(marker, {"eventId": event_id, "deliveredAt": self.clock()}) + self.s3.upload(marker, self._delivered_uri(event_id)) + delivered += 1 + return len(pending_keys) - delivered, delivered + + def _outbox_uri(self, event_id: str) -> str: + return f"s3://{self.bucket}/{self.outbox_prefix}/{self.session_id}/{event_id}.json" + + def _delivered_uri(self, event_id: str) -> str: + return f"s3://{self.bucket}/{self.delivered_prefix}/{self.session_id}/{event_id}.json" + + def _post(self, url: str, token: str, body: bytes) -> int: + request = urllib.request.Request( + url, + data=body, + method="POST", + headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"}, + ) + try: + with urllib.request.urlopen(request, timeout=10) as response: + return response.status + except urllib.error.HTTPError as error: + return error.code + + def _read_status(self) -> dict[str, object]: + try: + return json.loads(self.status_file.read_text(encoding="utf-8")) + except (FileNotFoundError, json.JSONDecodeError, OSError): + return {} + + @staticmethod + def _write_json_atomic(path: Path, value: object) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + encoded = json.dumps(value, sort_keys=True, separators=(",", ":")) + "\n" + with tempfile.NamedTemporaryFile("w", dir=path.parent, delete=False, encoding="utf-8") as handle: + handle.write(encoded) + handle.flush() + os.fsync(handle.fileno()) + temporary = Path(handle.name) + os.replace(temporary, path) + + +def from_environment() -> TraceExporter: + required = [ + "TRACE_SOURCE", + "TRACE_BUCKET", + "TRACE_PREFIX", + "TRACE_SESSION_ID", + "TRACE_OUTBOX_PREFIX", + "TRACE_DELIVERED_PREFIX", + "LOGGER_URL", + "LOGGER_TOKEN_FILE", + ] + missing = [name for name in required if not os.environ.get(name)] + if missing: + raise ValueError(f"missing required environment: {', '.join(missing)}") + work = Path(os.environ.get("TRACE_EXPORT_WORK_DIR", "/tmp/trace-exporter")) + return TraceExporter( + source=Path(os.environ["TRACE_SOURCE"]), + work=work, + status_file=Path(os.environ.get("TRACE_EXPORT_STATUS_FILE", str(work / "status.json"))), + bucket=os.environ["TRACE_BUCKET"], + trace_prefix=os.environ["TRACE_PREFIX"], + outbox_prefix=os.environ["TRACE_OUTBOX_PREFIX"], + delivered_prefix=os.environ["TRACE_DELIVERED_PREFIX"], + session_id=os.environ["TRACE_SESSION_ID"], + logger_url=os.environ["LOGGER_URL"], + logger_token_file=Path(os.environ["LOGGER_TOKEN_FILE"]), + ) + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.add_argument("--once", action="store_true", help="run one export/delivery cycle") + args = parser.parse_args() + exporter = from_environment() + interval = int(os.environ.get("TRACE_EXPORT_INTERVAL_SECONDS", "60")) + stop = threading.Event() + + def request_stop(_signal: int, _frame: object) -> None: + stop.set() + + signal.signal(signal.SIGTERM, request_stop) + signal.signal(signal.SIGINT, request_stop) + while True: + status = exporter.run_once() + print(json.dumps({"traceExporter": status}, sort_keys=True), flush=True) + if args.once: + return 0 + if stop.wait(interval): + final_status = exporter.run_once() + print(json.dumps({"traceExporter": final_status}, sort_keys=True), flush=True) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) From 7a988e7eaa3054dce4b21d8ff433c6a1dc9de5f7 Mon Sep 17 00:00:00 2001 From: Bo Wu Date: Sat, 5 Sep 2026 18:29:37 -0700 Subject: [PATCH 2/5] Reuse existing trace export path for commitments --- docker/runtime/Dockerfile | 4 +- gitops/README.md | 12 +- tests/run.sh | 2 +- ...e_exporter.py => test_trace_commitment.py} | 68 ++-- trace_commitment/README.md | 18 + trace_commitment/__init__.py | 1 + trace_commitment/trace_commitment.py | 311 +++++++++++++++ trace_exporter/README.md | 27 -- trace_exporter/__init__.py | 1 - trace_exporter/trace_exporter.py | 377 ------------------ 10 files changed, 365 insertions(+), 456 deletions(-) rename tests/{test_trace_exporter.py => test_trace_commitment.py} (53%) create mode 100644 trace_commitment/README.md create mode 100644 trace_commitment/__init__.py create mode 100644 trace_commitment/trace_commitment.py delete mode 100644 trace_exporter/README.md delete mode 100644 trace_exporter/__init__.py delete mode 100644 trace_exporter/trace_exporter.py diff --git a/docker/runtime/Dockerfile b/docker/runtime/Dockerfile index 288e78e..cc12c7f 100644 --- a/docker/runtime/Dockerfile +++ b/docker/runtime/Dockerfile @@ -20,10 +20,10 @@ RUN cd control-server && npm ci --omit=dev COPY . . COPY --from=multiagent-builder /src/target/release/multiagent /opt/multiagent/bin/multiagent RUN ln -s /opt/multiagent/wiki-service/bin/wiki-query.mjs /usr/local/bin/wiki-query \ - && ln -s /opt/multiagent/trace_exporter/trace_exporter.py /usr/local/bin/trace-exporter \ + && ln -s /opt/multiagent/trace_commitment/trace_commitment.py /usr/local/bin/trace-commitment \ && install -m 0755 docker/runtime/container-entrypoint.sh /opt/multiagent/bin/container-entrypoint.sh \ && install -m 0755 control-server/bin/prepare-repository.mjs /opt/multiagent/bin/prepare-repository.mjs \ - && chmod +x launch.sh control-server/bin/*.mjs trace_exporter/trace_exporter.py \ + && chmod +x launch.sh control-server/bin/*.mjs trace_commitment/trace_commitment.py \ && groupadd --gid 10000 multiagent-control \ && groupadd --gid 10001 multiagent-role \ && groupadd --gid 10004 multiagent-credentials \ diff --git a/gitops/README.md b/gitops/README.md index d101ba8..c97435d 100644 --- a/gitops/README.md +++ b/gitops/README.md @@ -28,12 +28,12 @@ The application-owned Logger deployment contract is: - give producers a durable retry path or outbox and alert on delivery backlog; - keep the existing trace sidecar and S3 data path, then submit a bounded `trace.artifact_exported` commitment after a successful upload; -- run `trace-exporter` with a trace-commitment-only token. It stores only - bounded commitment JSON in the deployment-owned S3 outbox, uses deterministic - event IDs and delivered markers for restart-safe idempotency, and reports - `loggerPending`, `loggerOk`, and Logger delivery timestamps in its status - file. Logger failure must not change the S3 export `ok` result or workflow - progression; +- keep the existing exporter responsible for trace staging, delta detection, + retry, and S3 sync. After a successful sync, run `trace-commitment` with the + stable staged tree and a trace-commitment-only token. It writes the commitment + manifest, event outbox, and delivered marker as separate metadata objects in + the same existing S3 location, and reports `loggerPending`, `loggerOk`, and + delivery timestamps without changing S3 export `ok` or workflow progression; - configure at most one active Logger replica for a ledger volume. A standby must not write until deployment fencing has transferred ownership. diff --git a/tests/run.sh b/tests/run.sh index 9d39681..e255a6e 100755 --- a/tests/run.sh +++ b/tests/run.sh @@ -1247,7 +1247,7 @@ PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_native_s PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_swe_outcomes.py" PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_swe_provenance.py" PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_migration_contracts.py" -PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_trace_exporter.py" +PYTHONPATH="$ROOT${PYTHONPATH:+:$PYTHONPATH}" python3 "$ROOT/tests/test_trace_commitment.py" python3 -m evaluation.swe_bench_pro --help >"$TMPDIR/swe-bench-pro-help.out" assert_file_contains "$TMPDIR/swe-bench-pro-help.out" "Evaluate the production multiagent solver" assert_file_not_contains "$TMPDIR/swe-bench-pro-help.out" "--agent-framework" diff --git a/tests/test_trace_exporter.py b/tests/test_trace_commitment.py similarity index 53% rename from tests/test_trace_exporter.py rename to tests/test_trace_commitment.py index a299361..9e8236b 100644 --- a/tests/test_trace_exporter.py +++ b/tests/test_trace_commitment.py @@ -5,7 +5,7 @@ import unittest from pathlib import Path -from trace_exporter.trace_exporter import TraceExporter +from trace_commitment.trace_commitment import TraceCommitter class MemoryS3: @@ -15,9 +15,6 @@ def __init__(self): def upload(self, source, uri): self.objects[uri] = Path(source).read_bytes() - def copy(self, source_uri, destination_uri): - self.objects[destination_uri] = self.objects[source_uri] - def download(self, uri, destination): destination.parent.mkdir(parents=True, exist_ok=True) destination.write_bytes(self.objects[uri]) @@ -27,14 +24,15 @@ def list_keys(self, bucket, prefix): return {uri.removeprefix(root) for uri in self.objects if uri.startswith(root + prefix)} -class TraceExporterTest(unittest.TestCase): +class TraceCommitterTest(unittest.TestCase): def setUp(self): self.temporary = tempfile.TemporaryDirectory() root = Path(self.temporary.name) - self.source = root / "source" + self.source = root / "existing-export-stage" self.source.mkdir() self.work = root / "work" self.status = root / "status.json" + self.status.write_text('{"ok":true,"lastSuccessAt":"2026-09-05T11:59:00Z"}\n') self.token = root / "token" self.token.write_text("trace-token-0123456789abcdef", encoding="utf-8") self.s3 = MemoryS3() @@ -43,20 +41,17 @@ def setUp(self): def tearDown(self): self.temporary.cleanup() - def exporter(self, post=None): + def committer(self, post=None): def successful(url, token, body): self.posts.append((url, token, json.loads(body))) return 204 - return TraceExporter( + return TraceCommitter( source=self.source, + destination="s3://trace-bucket/production/sessions/session-1", + session_id="session-1", work=self.work, status_file=self.status, - bucket="trace-bucket", - trace_prefix="production/sessions/session-1", - outbox_prefix="production/logger-outbox", - delivered_prefix="production/logger-delivered", - session_id="session-1", logger_url="http://logger", logger_token_file=self.token, s3=self.s3, @@ -64,53 +59,42 @@ def successful(url, token, body): clock=lambda: "2026-09-05T12:00:00Z", ) - def test_successful_export_commits_digest_reference_and_no_body(self): - body = b"private trace body" - (self.source / "trace.jsonl").write_bytes(body) + def test_successful_sync_gets_separate_manifest_and_metadata_only_event(self): + trace_body = b"private trace body" + (self.source / "trace.jsonl").write_bytes(trace_body) - status = self.exporter().run_once() + status = self.committer().commit() self.assertTrue(status["ok"]) self.assertTrue(status["loggerOk"]) self.assertEqual(status["loggerPending"], 0) - self.assertEqual(len(self.posts), 1) event = self.posts[0][2] self.assertEqual(event["eventType"], "trace.artifact_exported") - self.assertEqual(event["sessionId"], "session-1") - self.assertEqual(event["payloadDigest"], event["artifactReferences"][0]["digest"]) - self.assertEqual(event["artifactReferences"][0]["size"], len(body)) - self.assertNotIn(body, self.posts[0][2].__str__().encode()) - artifact_uri = event["artifactReferences"][0]["uri"] - self.assertEqual(self.s3.objects[artifact_uri], body) - - def test_delivered_marker_makes_restart_idempotent(self): + self.assertEqual(len(event["artifactReferences"]), 1) + manifest_uri = event["artifactReferences"][0]["uri"] + self.assertIn("/commitments/", manifest_uri) + manifest = json.loads(self.s3.objects[manifest_uri]) + self.assertEqual(manifest["artifacts"][0]["size"], len(trace_body)) + self.assertNotIn(trace_body, json.dumps(event).encode()) + + def test_delivered_marker_makes_repeated_post_upload_hook_idempotent(self): (self.source / "usage.json").write_text("{}", encoding="utf-8") - first = self.exporter() - first.run_once() - self.assertEqual(len(self.posts), 1) - # A fresh local cache re-uploads metadata, but the durable delivered - # marker prevents another Logger request. - for path in self.work.iterdir(): - if path.is_dir(): - import shutil - - shutil.rmtree(path) - self.exporter().run_once() + self.committer().commit() + self.committer().commit() self.assertEqual(len(self.posts), 1) - def test_logger_outage_leaves_durable_backlog_and_later_drains(self): + def test_logger_outage_preserves_s3_success_and_durable_backlog(self): (self.source / "report.md").write_text("report", encoding="utf-8") def unavailable(_url, _token, _body): raise OSError("temporary Logger outage") - failed = self.exporter(unavailable).run_once() - self.assertTrue(failed["ok"], "Logger availability must not change S3 export success") + failed = self.committer(unavailable).commit() + self.assertTrue(failed["ok"]) self.assertFalse(failed["loggerOk"]) self.assertEqual(failed["loggerPending"], 1) - self.assertIn("temporary Logger outage", failed["loggerError"]) - recovered = self.exporter().run_once() + recovered = self.committer().commit() self.assertTrue(recovered["loggerOk"]) self.assertEqual(recovered["loggerPending"], 0) self.assertEqual(len(self.posts), 1) diff --git a/trace_commitment/README.md b/trace_commitment/README.md new file mode 100644 index 0000000..f669b99 --- /dev/null +++ b/trace_commitment/README.md @@ -0,0 +1,18 @@ +# Trace commitment delivery + +The deployment's existing trace exporter remains responsible for staging, +delta detection, retry, and uploading trace bodies. After a successful sync, +`python3 -m trace_commitment.trace_commitment` hashes that stable staged tree, +writes a deterministic commitment manifest as a separate object in the same S3 +location, and durably retries a bounded `trace.artifact_exported` Logger event. + +The Logger receives only the manifest digest, size, media type, and S3 +reference. The event and delivered marker are also separate metadata-only +objects beside the existing export. Deterministic IDs preserve idempotency after +an uncertain acknowledgement or process restart. + +Required environment is `TRACE_COMMITMENT_SOURCE`, `TRACE_EXPORT_DESTINATION`, +`TRACE_SESSION_ID`, `LOGGER_URL`, and `LOGGER_TOKEN_FILE`. +`TRACE_COMMITMENT_WORK_DIR` and `TRACE_EXPORT_STATUS_FILE` are optional. Logger +delivery updates `loggerOk`, `loggerPending`, and Logger timestamps/errors in +the existing status file without changing its S3 export `ok` field. diff --git a/trace_commitment/__init__.py b/trace_commitment/__init__.py new file mode 100644 index 0000000..aed859b --- /dev/null +++ b/trace_commitment/__init__.py @@ -0,0 +1 @@ +"""Post-upload trace commitment delivery.""" diff --git a/trace_commitment/trace_commitment.py b/trace_commitment/trace_commitment.py new file mode 100644 index 0000000..17c3395 --- /dev/null +++ b/trace_commitment/trace_commitment.py @@ -0,0 +1,311 @@ +#!/usr/bin/env python3 +"""Commit an already-exported trace tree to Logger without sending trace bodies.""" + +from __future__ import annotations + +import argparse +import hashlib +import json +import mimetypes +import os +import stat +import subprocess +import tempfile +import urllib.error +import urllib.request +from datetime import datetime, timezone +from pathlib import Path, PurePosixPath +from typing import Callable, Iterable +from urllib.parse import quote + + +def utc_now() -> str: + return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") + + +def sha256_bytes(value: bytes) -> str: + return f"sha256:{hashlib.sha256(value).hexdigest()}" + + +def sha256_file(path: Path) -> tuple[str, int]: + digest = hashlib.sha256() + size = 0 + with path.open("rb") as handle: + while chunk := handle.read(1024 * 1024): + digest.update(chunk) + size += len(chunk) + return f"sha256:{digest.hexdigest()}", size + + +def safe_identifier(value: str, name: str) -> str: + if ( + not value + or len(value) > 128 + or not value[0].isascii() + or not value[0].isalnum() + or any( + not character.isascii() + or not (character.isalnum() or character in "._:-") + for character in value + ) + ): + raise ValueError(f"{name} has an invalid format") + return value + + +def s3_destination(value: str) -> tuple[str, str]: + if not value.startswith("s3://") or any(character in value for character in "\r\n?#\\"): + raise ValueError("trace destination must be a bounded s3:// URI") + bucket, separator, raw_prefix = value.removeprefix("s3://").partition("/") + parts = PurePosixPath(raw_prefix.strip("/")).parts + if not bucket or not separator or not parts or raw_prefix != raw_prefix.strip("/"): + raise ValueError("trace destination requires a bucket and relative prefix") + if any(part in ("", ".", "..") for part in parts): + raise ValueError("trace destination prefix is invalid") + return bucket, "/".join(parts) + + +class AwsS3: + def __init__(self, runner: Callable[..., subprocess.CompletedProcess[str]] | None = None): + self.runner = runner or subprocess.run + + def _run(self, args: list[str]) -> subprocess.CompletedProcess[str]: + return self.runner(args, check=True, capture_output=True, text=True) + + def upload(self, source: Path, uri: str) -> None: + self._run(["aws", "s3", "cp", str(source), uri, "--only-show-errors"]) + + def download(self, uri: str, destination: Path) -> None: + destination.parent.mkdir(parents=True, exist_ok=True) + self._run(["aws", "s3", "cp", uri, str(destination), "--only-show-errors"]) + + def list_keys(self, bucket: str, prefix: str) -> set[str]: + result = self._run( + ["aws", "s3api", "list-objects-v2", "--bucket", bucket, "--prefix", prefix, "--output", "json"] + ) + return {item["Key"] for item in json.loads(result.stdout or "{}").get("Contents", [])} + + +class TraceCommitter: + def __init__( + self, + *, + source: Path, + destination: str, + session_id: str, + work: Path, + status_file: Path, + logger_url: str, + logger_token_file: Path, + s3: AwsS3 | None = None, + http_post: Callable[[str, str, bytes], int] | None = None, + clock: Callable[[], str] = utc_now, + ): + self.source = source + self.bucket, self.prefix = s3_destination(destination) + self.session_id = safe_identifier(session_id, "session ID") + self.work = work + self.status_file = status_file + self.logger_url = logger_url.rstrip("/") + self.logger_token_file = logger_token_file + self.s3 = s3 or AwsS3() + self.http_post = http_post or self._post + self.clock = clock + + def commit(self) -> dict[str, object]: + attempted_at = self.clock() + error: str | None = None + delivered = 0 + pending = 0 + try: + event = self._publish_manifest_and_event() + pending, delivered = self._deliver_pending() + event_id = event["eventId"] + except Exception as caught: # Logger path must not alter trace-export success + error = str(caught) + event_id = None + try: + pending = len(self._pending_keys()) + except Exception: + pending = max(pending, 1) + + status = self._read_status() + status.update( + { + "loggerOk": error is None and pending == 0, + "loggerPending": pending, + "loggerDeliveredCount": delivered, + "loggerLastAttemptAt": attempted_at, + } + ) + if event_id: + status["loggerEventId"] = event_id + if error is None and pending == 0: + status["loggerLastSuccessAt"] = attempted_at + status.pop("loggerError", None) + else: + status["loggerError"] = (error or "Logger backlog remains pending")[:512] + self._write_json_atomic(self.status_file, status) + return status + + def _publish_manifest_and_event(self) -> dict[str, object]: + artifacts = [] + for path, relative in self._source_files(): + digest, size = sha256_file(path) + artifacts.append( + { + "digest": digest, + "mediaType": mimetypes.guess_type(relative.as_posix())[0] or "application/octet-stream", + "size": size, + "storageReference": f"s3://{self.bucket}/{self.prefix}/{quote(relative.as_posix(), safe='/')}", + } + ) + manifest = { + "apiVersion": "trace.multiagent.dev/v1", + "kind": "TraceExportCommitment", + "sessionId": self.session_id, + "artifacts": artifacts, + } + encoded = json.dumps(manifest, sort_keys=True, separators=(",", ":")).encode() + b"\n" + digest = sha256_bytes(encoded) + digest_hex = digest.removeprefix("sha256:") + manifest_uri = f"s3://{self.bucket}/{self.prefix}/commitments/{digest_hex}.json" + manifest_file = self.work / "manifests" / f"{digest_hex}.json" + self._write_bytes_atomic(manifest_file, encoded) + self.s3.upload(manifest_file, manifest_uri) + + identity = hashlib.sha256(f"{self.session_id}\0{manifest_uri}\0{digest}".encode()).hexdigest() + event = { + "eventId": f"trace-export-{identity}", + "sessionId": self.session_id, + "eventType": "trace.artifact_exported", + "payloadDigest": digest, + "artifactReferences": [ + { + "uri": manifest_uri, + "digest": digest, + "size": len(encoded), + "mediaType": "application/vnd.multiagent.trace-commitment+json", + } + ], + } + event_file = self.work / "outbox" / f"{event['eventId']}.json" + self._write_json_atomic(event_file, event) + self.s3.upload(event_file, self._outbox_uri(str(event["eventId"]))) + return event + + def _source_files(self) -> Iterable[tuple[Path, Path]]: + for root, directories, files in os.walk(self.source, followlinks=False): + directories[:] = sorted(name for name in directories if name != "worktrees") + root_path = Path(root) + for name in sorted(files): + if name == "orchestrator-bootstrap.sh": + continue + path = root_path / name + try: + if not stat.S_ISREG(path.lstat().st_mode): + continue + except (FileNotFoundError, PermissionError, OSError): + continue + yield path, path.relative_to(self.source) + + def _pending_keys(self) -> list[str]: + outbox_root = f"{self.prefix}/logger-outbox/" + delivered_root = f"{self.prefix}/logger-delivered/" + outbox = self.s3.list_keys(self.bucket, outbox_root) + delivered = self.s3.list_keys(self.bucket, delivered_root) + delivered_ids = {Path(key).stem for key in delivered if key.endswith(".json")} + return sorted(key for key in outbox if key.endswith(".json") and Path(key).stem not in delivered_ids) + + def _deliver_pending(self) -> tuple[int, int]: + pending_keys = self._pending_keys() + token = self.logger_token_file.read_text(encoding="utf-8").strip() + if not token: + raise ValueError("Logger token file is empty") + delivered = 0 + for key in pending_keys: + event_id = Path(key).stem + local = self.work / "download" / f"{event_id}.json" + self.s3.download(f"s3://{self.bucket}/{key}", local) + status = self.http_post(f"{self.logger_url}/v1/events", token, local.read_bytes()) + if status != 204: + raise RuntimeError(f"Logger returned HTTP {status}") + marker = self.work / "delivered" / f"{event_id}.json" + self._write_json_atomic(marker, {"eventId": event_id, "deliveredAt": self.clock()}) + self.s3.upload(marker, self._delivered_uri(event_id)) + delivered += 1 + return len(pending_keys) - delivered, delivered + + def _outbox_uri(self, event_id: str) -> str: + return f"s3://{self.bucket}/{self.prefix}/logger-outbox/{event_id}.json" + + def _delivered_uri(self, event_id: str) -> str: + return f"s3://{self.bucket}/{self.prefix}/logger-delivered/{event_id}.json" + + def _post(self, url: str, token: str, body: bytes) -> int: + request = urllib.request.Request( + url, + data=body, + method="POST", + headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"}, + ) + try: + with urllib.request.urlopen(request, timeout=10) as response: + return response.status + except urllib.error.HTTPError as error: + return error.code + + def _read_status(self) -> dict[str, object]: + try: + return json.loads(self.status_file.read_text(encoding="utf-8")) + except (FileNotFoundError, json.JSONDecodeError, OSError): + return {} + + @staticmethod + def _write_bytes_atomic(path: Path, encoded: bytes) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + with tempfile.NamedTemporaryFile("wb", dir=path.parent, delete=False) as handle: + handle.write(encoded) + handle.flush() + os.fsync(handle.fileno()) + temporary = Path(handle.name) + os.replace(temporary, path) + + @classmethod + def _write_json_atomic(cls, path: Path, value: object) -> None: + cls._write_bytes_atomic(path, (json.dumps(value, sort_keys=True, separators=(",", ":")) + "\n").encode()) + + +def from_environment() -> TraceCommitter: + required = [ + "TRACE_COMMITMENT_SOURCE", + "TRACE_EXPORT_DESTINATION", + "TRACE_SESSION_ID", + "LOGGER_URL", + "LOGGER_TOKEN_FILE", + ] + missing = [name for name in required if not os.environ.get(name)] + if missing: + raise ValueError(f"missing required environment: {', '.join(missing)}") + work = Path(os.environ.get("TRACE_COMMITMENT_WORK_DIR", "/tmp/trace-commitment")) + return TraceCommitter( + source=Path(os.environ["TRACE_COMMITMENT_SOURCE"]), + destination=os.environ["TRACE_EXPORT_DESTINATION"], + session_id=os.environ["TRACE_SESSION_ID"], + work=work, + status_file=Path(os.environ.get("TRACE_EXPORT_STATUS_FILE", "/tmp/trace-exporter/status.json")), + logger_url=os.environ["LOGGER_URL"], + logger_token_file=Path(os.environ["LOGGER_TOKEN_FILE"]), + ) + + +def main() -> int: + parser = argparse.ArgumentParser(description=__doc__) + parser.parse_args() + status = from_environment().commit() + print(json.dumps({"traceCommitment": status}, sort_keys=True)) + return 0 + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/trace_exporter/README.md b/trace_exporter/README.md deleted file mode 100644 index 21e18a0..0000000 --- a/trace_exporter/README.md +++ /dev/null @@ -1,27 +0,0 @@ -# Trace exporter - -`python3 -m trace_exporter.trace_exporter` copies regular trace files to the -deployment-owned S3 trace bucket, preserves a content-addressed object for each -digest, and writes a bounded metadata-only Logger event to an S3 outbox. It -replays outbox objects until Logger acknowledges them, then writes an immutable -delivered marker. Deterministic event IDs make a retry after an uncertain -acknowledgement safe against Logger's exact idempotency contract. - -Required environment: - -| Variable | Purpose | -| --- | --- | -| `TRACE_SOURCE` | Read-only trace tree | -| `TRACE_BUCKET` | Deployment-owned S3 bucket | -| `TRACE_PREFIX` | Conventional and content-addressed trace prefix | -| `TRACE_SESSION_ID` | Logger session/log ID | -| `TRACE_OUTBOX_PREFIX` | Metadata-only durable outbox prefix | -| `TRACE_DELIVERED_PREFIX` | Durable delivery-marker prefix | -| `LOGGER_URL` | Private Logger base URL | -| `LOGGER_TOKEN_FILE` | Trace-commitment-only bearer token file | - -`TRACE_EXPORT_WORK_DIR`, `TRACE_EXPORT_STATUS_FILE`, and -`TRACE_EXPORT_INTERVAL_SECONDS` are optional. The status contract keeps S3 -success in `ok` and reports Logger delivery separately through `loggerOk`, -`loggerPending`, and Logger timestamps/errors. Consumers must never use Logger -availability or acknowledgement to grant or deny workflow transitions. diff --git a/trace_exporter/__init__.py b/trace_exporter/__init__.py deleted file mode 100644 index 32d6f66..0000000 --- a/trace_exporter/__init__.py +++ /dev/null @@ -1 +0,0 @@ -"""Durable S3 trace export and Logger commitment delivery.""" diff --git a/trace_exporter/trace_exporter.py b/trace_exporter/trace_exporter.py deleted file mode 100644 index d0fd582..0000000 --- a/trace_exporter/trace_exporter.py +++ /dev/null @@ -1,377 +0,0 @@ -#!/usr/bin/env python3 -"""Export trace files to S3 and durably commit metadata to Logger. - -The S3 outbox is deliberately metadata-only. Trace bodies are uploaded to a -content-addressed artifact key and never sent to Logger. -""" - -from __future__ import annotations - -import argparse -import hashlib -import json -import mimetypes -import os -import shutil -import signal -import stat -import subprocess -import tempfile -import threading -import urllib.error -import urllib.request -from urllib.parse import quote -from datetime import datetime, timezone -from pathlib import Path, PurePosixPath -from typing import Callable, Iterable - - -def utc_now() -> str: - return datetime.now(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") - - -def sha256_file(path: Path) -> tuple[str, int]: - digest = hashlib.sha256() - size = 0 - with path.open("rb") as handle: - while chunk := handle.read(1024 * 1024): - digest.update(chunk) - size += len(chunk) - return f"sha256:{digest.hexdigest()}", size - - -def safe_segment(value: str, name: str) -> str: - if ( - not value - or len(value) > 128 - or not value[0].isascii() - or not value[0].isalnum() - or any( - not character.isascii() - or not (character.isalnum() or character in "._:-") - for character in value - ) - ): - raise ValueError(f"{name} has an invalid format") - return value - - -def normalized_prefix(value: str, name: str) -> str: - parts = PurePosixPath(value.strip("/")).parts - if ( - not parts - or value != value.strip("/") - or any(part in ("", ".", "..") for part in parts) - or any(character in value for character in "\r\n?#\\") - ): - raise ValueError(f"{name} has an invalid format") - return "/".join(parts) - - -class AwsS3: - def __init__(self, runner: Callable[..., subprocess.CompletedProcess[str]] | None = None): - self.runner = runner or subprocess.run - - def _run(self, args: list[str]) -> subprocess.CompletedProcess[str]: - return self.runner(args, check=True, capture_output=True, text=True) - - def upload(self, source: Path, uri: str) -> None: - self._run(["aws", "s3", "cp", str(source), uri, "--only-show-errors"]) - - def copy(self, source_uri: str, destination_uri: str) -> None: - self._run(["aws", "s3", "cp", source_uri, destination_uri, "--only-show-errors"]) - - def download(self, uri: str, destination: Path) -> None: - destination.parent.mkdir(parents=True, exist_ok=True) - self._run(["aws", "s3", "cp", uri, str(destination), "--only-show-errors"]) - - def list_keys(self, bucket: str, prefix: str) -> set[str]: - result = self._run( - [ - "aws", - "s3api", - "list-objects-v2", - "--bucket", - bucket, - "--prefix", - prefix, - "--output", - "json", - ] - ) - decoded = json.loads(result.stdout or "{}") - return {item["Key"] for item in decoded.get("Contents", [])} - - -class TraceExporter: - def __init__( - self, - *, - source: Path, - work: Path, - status_file: Path, - bucket: str, - trace_prefix: str, - outbox_prefix: str, - delivered_prefix: str, - session_id: str, - logger_url: str, - logger_token_file: Path, - s3: AwsS3 | None = None, - http_post: Callable[[str, str, bytes], int] | None = None, - clock: Callable[[], str] = utc_now, - ): - self.source = source - self.work = work - self.status_file = status_file - self.bucket = bucket - self.trace_prefix = normalized_prefix(trace_prefix, "trace prefix") - self.outbox_prefix = normalized_prefix(outbox_prefix, "outbox prefix") - self.delivered_prefix = normalized_prefix(delivered_prefix, "delivered prefix") - self.session_id = safe_segment(session_id, "session ID") - self.logger_url = logger_url.rstrip("/") - self.logger_token_file = logger_token_file - self.s3 = s3 or AwsS3() - self.http_post = http_post or self._post - self.clock = clock - self.cache = work / "cache" - self.staging = work / "staging" - self.queue_cache = work / "queue" - - def run_once(self) -> dict[str, object]: - attempted_at = self.clock() - uploaded = 0 - scanned = 0 - export_error: str | None = None - logger_error: str | None = None - self.work.mkdir(parents=True, exist_ok=True) - self.cache.mkdir(parents=True, exist_ok=True) - shutil.rmtree(self.staging, ignore_errors=True) - self.staging.mkdir(parents=True, exist_ok=True) - - try: - for source_file, relative in self._source_files(): - scanned += 1 - staged = self.staging / relative - staged.parent.mkdir(parents=True, exist_ok=True) - try: - shutil.copy2(source_file, staged) - except (FileNotFoundError, PermissionError): - continue - digest, size = sha256_file(staged) - cached = self.cache / relative - if cached.is_file() and sha256_file(cached) == (digest, size): - continue - media_type = mimetypes.guess_type(relative.as_posix())[0] or "application/octet-stream" - encoded_relative = quote(relative.as_posix(), safe="/") - canonical_key = ( - f"{self.trace_prefix}/artifacts/{self.session_id}/" - f"{digest.removeprefix('sha256:')}/{encoded_relative}" - ) - canonical_uri = f"s3://{self.bucket}/{canonical_key}" - current_uri = f"s3://{self.bucket}/{self.trace_prefix}/{relative.as_posix()}" - self.s3.upload(staged, canonical_uri) - self.s3.copy(canonical_uri, current_uri) - event = self._event(canonical_uri, digest, size, media_type) - self._enqueue(event) - cached.parent.mkdir(parents=True, exist_ok=True) - os.replace(staged, cached) - uploaded += 1 - except Exception as error: # errors are reported through the shared status contract - export_error = str(error) - - pending = 0 - delivered = 0 - try: - pending, delivered = self._deliver_pending() - except Exception as error: - logger_error = str(error) - try: - pending = len(self._pending_keys()) - except Exception: - pending = max(pending, 1) - - previous = self._read_status() - status: dict[str, object] = { - "ok": export_error is None, - "lastAttemptAt": attempted_at, - "fileCount": scanned, - "uploadedCount": uploaded, - "loggerOk": logger_error is None and pending == 0, - "loggerPending": pending, - "loggerDeliveredCount": delivered, - } - if export_error is None: - status["lastSuccessAt"] = attempted_at - elif previous.get("lastSuccessAt"): - status["lastSuccessAt"] = previous["lastSuccessAt"] - status["error"] = export_error - else: - status["error"] = export_error - if logger_error is None and pending == 0: - status["loggerLastSuccessAt"] = attempted_at - else: - if previous.get("loggerLastSuccessAt"): - status["loggerLastSuccessAt"] = previous["loggerLastSuccessAt"] - if logger_error: - status["loggerError"] = logger_error[:512] - self._write_json_atomic(self.status_file, status) - return status - - def _source_files(self) -> Iterable[tuple[Path, Path]]: - if not self.source.is_dir(): - return - for root, directories, files in os.walk(self.source, followlinks=False): - directories[:] = sorted(name for name in directories if name != "worktrees") - root_path = Path(root) - for name in sorted(files): - if name == "orchestrator-bootstrap.sh": - continue - path = root_path / name - try: - if not stat.S_ISREG(path.lstat().st_mode): - continue - except (FileNotFoundError, PermissionError, OSError): - continue - yield path, path.relative_to(self.source) - - def _event(self, uri: str, digest: str, size: int, media_type: str) -> dict[str, object]: - identity = hashlib.sha256(f"{self.session_id}\0{uri}\0{digest}".encode()).hexdigest() - return { - "eventId": f"trace-export-{identity}", - "sessionId": self.session_id, - "eventType": "trace.artifact_exported", - "payloadDigest": digest, - "artifactReferences": [ - {"uri": uri, "digest": digest, "size": size, "mediaType": media_type} - ], - } - - def _enqueue(self, event: dict[str, object]) -> None: - event_id = str(event["eventId"]) - local = self.queue_cache / "outbox" / f"{event_id}.json" - self._write_json_atomic(local, event) - self.s3.upload(local, self._outbox_uri(event_id)) - - def _pending_keys(self) -> list[str]: - outbox_root = f"{self.outbox_prefix}/{self.session_id}/" - delivered_root = f"{self.delivered_prefix}/{self.session_id}/" - outbox = self.s3.list_keys(self.bucket, outbox_root) - delivered = self.s3.list_keys(self.bucket, delivered_root) - delivered_ids = {Path(key).stem for key in delivered if key.endswith(".json")} - return sorted( - key for key in outbox if key.endswith(".json") and Path(key).stem not in delivered_ids - ) - - def _deliver_pending(self) -> tuple[int, int]: - pending_keys = self._pending_keys() - delivered = 0 - token = self.logger_token_file.read_text(encoding="utf-8").strip() - if not token: - raise ValueError("Logger token file is empty") - for key in pending_keys: - event_id = Path(key).stem - local = self.queue_cache / "download" / f"{event_id}.json" - self.s3.download(f"s3://{self.bucket}/{key}", local) - body = local.read_bytes() - status = self.http_post(f"{self.logger_url}/v1/events", token, body) - if status != 204: - raise RuntimeError(f"Logger returned HTTP {status}") - marker = self.queue_cache / "delivered" / f"{event_id}.json" - self._write_json_atomic(marker, {"eventId": event_id, "deliveredAt": self.clock()}) - self.s3.upload(marker, self._delivered_uri(event_id)) - delivered += 1 - return len(pending_keys) - delivered, delivered - - def _outbox_uri(self, event_id: str) -> str: - return f"s3://{self.bucket}/{self.outbox_prefix}/{self.session_id}/{event_id}.json" - - def _delivered_uri(self, event_id: str) -> str: - return f"s3://{self.bucket}/{self.delivered_prefix}/{self.session_id}/{event_id}.json" - - def _post(self, url: str, token: str, body: bytes) -> int: - request = urllib.request.Request( - url, - data=body, - method="POST", - headers={"Authorization": f"Bearer {token}", "Content-Type": "application/json"}, - ) - try: - with urllib.request.urlopen(request, timeout=10) as response: - return response.status - except urllib.error.HTTPError as error: - return error.code - - def _read_status(self) -> dict[str, object]: - try: - return json.loads(self.status_file.read_text(encoding="utf-8")) - except (FileNotFoundError, json.JSONDecodeError, OSError): - return {} - - @staticmethod - def _write_json_atomic(path: Path, value: object) -> None: - path.parent.mkdir(parents=True, exist_ok=True) - encoded = json.dumps(value, sort_keys=True, separators=(",", ":")) + "\n" - with tempfile.NamedTemporaryFile("w", dir=path.parent, delete=False, encoding="utf-8") as handle: - handle.write(encoded) - handle.flush() - os.fsync(handle.fileno()) - temporary = Path(handle.name) - os.replace(temporary, path) - - -def from_environment() -> TraceExporter: - required = [ - "TRACE_SOURCE", - "TRACE_BUCKET", - "TRACE_PREFIX", - "TRACE_SESSION_ID", - "TRACE_OUTBOX_PREFIX", - "TRACE_DELIVERED_PREFIX", - "LOGGER_URL", - "LOGGER_TOKEN_FILE", - ] - missing = [name for name in required if not os.environ.get(name)] - if missing: - raise ValueError(f"missing required environment: {', '.join(missing)}") - work = Path(os.environ.get("TRACE_EXPORT_WORK_DIR", "/tmp/trace-exporter")) - return TraceExporter( - source=Path(os.environ["TRACE_SOURCE"]), - work=work, - status_file=Path(os.environ.get("TRACE_EXPORT_STATUS_FILE", str(work / "status.json"))), - bucket=os.environ["TRACE_BUCKET"], - trace_prefix=os.environ["TRACE_PREFIX"], - outbox_prefix=os.environ["TRACE_OUTBOX_PREFIX"], - delivered_prefix=os.environ["TRACE_DELIVERED_PREFIX"], - session_id=os.environ["TRACE_SESSION_ID"], - logger_url=os.environ["LOGGER_URL"], - logger_token_file=Path(os.environ["LOGGER_TOKEN_FILE"]), - ) - - -def main() -> int: - parser = argparse.ArgumentParser(description=__doc__) - parser.add_argument("--once", action="store_true", help="run one export/delivery cycle") - args = parser.parse_args() - exporter = from_environment() - interval = int(os.environ.get("TRACE_EXPORT_INTERVAL_SECONDS", "60")) - stop = threading.Event() - - def request_stop(_signal: int, _frame: object) -> None: - stop.set() - - signal.signal(signal.SIGTERM, request_stop) - signal.signal(signal.SIGINT, request_stop) - while True: - status = exporter.run_once() - print(json.dumps({"traceExporter": status}, sort_keys=True), flush=True) - if args.once: - return 0 - if stop.wait(interval): - final_status = exporter.run_once() - print(json.dumps({"traceExporter": final_status}, sort_keys=True), flush=True) - return 0 - - -if __name__ == "__main__": - raise SystemExit(main()) From 3ecba23256bfe2254fce78476dc4ae663340fdb9 Mon Sep 17 00:00:00 2001 From: Bo Wu Date: Sat, 5 Sep 2026 18:31:11 -0700 Subject: [PATCH 3/5] Keep trace commitment hook Python 3.8 compatible --- tests/test_trace_commitment.py | 2 +- trace_commitment/trace_commitment.py | 8 ++++++-- 2 files changed, 7 insertions(+), 3 deletions(-) diff --git a/tests/test_trace_commitment.py b/tests/test_trace_commitment.py index 9e8236b..3f92a78 100644 --- a/tests/test_trace_commitment.py +++ b/tests/test_trace_commitment.py @@ -21,7 +21,7 @@ def download(self, uri, destination): def list_keys(self, bucket, prefix): root = f"s3://{bucket}/" - return {uri.removeprefix(root) for uri in self.objects if uri.startswith(root + prefix)} + return {uri[len(root) :] for uri in self.objects if uri.startswith(root + prefix)} class TraceCommitterTest(unittest.TestCase): diff --git a/trace_commitment/trace_commitment.py b/trace_commitment/trace_commitment.py index 17c3395..ce03b4c 100644 --- a/trace_commitment/trace_commitment.py +++ b/trace_commitment/trace_commitment.py @@ -27,6 +27,10 @@ def sha256_bytes(value: bytes) -> str: return f"sha256:{hashlib.sha256(value).hexdigest()}" +def remove_prefix(value: str, prefix: str) -> str: + return value[len(prefix) :] if value.startswith(prefix) else value + + def sha256_file(path: Path) -> tuple[str, int]: digest = hashlib.sha256() size = 0 @@ -56,7 +60,7 @@ def safe_identifier(value: str, name: str) -> str: def s3_destination(value: str) -> tuple[str, str]: if not value.startswith("s3://") or any(character in value for character in "\r\n?#\\"): raise ValueError("trace destination must be a bounded s3:// URI") - bucket, separator, raw_prefix = value.removeprefix("s3://").partition("/") + bucket, separator, raw_prefix = remove_prefix(value, "s3://").partition("/") parts = PurePosixPath(raw_prefix.strip("/")).parts if not bucket or not separator or not parts or raw_prefix != raw_prefix.strip("/"): raise ValueError("trace destination requires a bucket and relative prefix") @@ -168,7 +172,7 @@ def _publish_manifest_and_event(self) -> dict[str, object]: } encoded = json.dumps(manifest, sort_keys=True, separators=(",", ":")).encode() + b"\n" digest = sha256_bytes(encoded) - digest_hex = digest.removeprefix("sha256:") + digest_hex = remove_prefix(digest, "sha256:") manifest_uri = f"s3://{self.bucket}/{self.prefix}/commitments/{digest_hex}.json" manifest_file = self.work / "manifests" / f"{digest_hex}.json" self._write_bytes_atomic(manifest_file, encoded) From 6da2ac96496352fdcb51355aadf69a5a8b862445 Mon Sep 17 00:00:00 2001 From: Bo Wu Date: Sat, 5 Sep 2026 18:32:02 -0700 Subject: [PATCH 4/5] Preserve exact existing S3 object references --- trace_commitment/trace_commitment.py | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/trace_commitment/trace_commitment.py b/trace_commitment/trace_commitment.py index ce03b4c..bf7528e 100644 --- a/trace_commitment/trace_commitment.py +++ b/trace_commitment/trace_commitment.py @@ -16,7 +16,6 @@ from datetime import datetime, timezone from pathlib import Path, PurePosixPath from typing import Callable, Iterable -from urllib.parse import quote def utc_now() -> str: @@ -161,7 +160,7 @@ def _publish_manifest_and_event(self) -> dict[str, object]: "digest": digest, "mediaType": mimetypes.guess_type(relative.as_posix())[0] or "application/octet-stream", "size": size, - "storageReference": f"s3://{self.bucket}/{self.prefix}/{quote(relative.as_posix(), safe='/')}", + "storageReference": f"s3://{self.bucket}/{self.prefix}/{relative.as_posix()}", } ) manifest = { From 9bb64366c1f7ba183dbfff30a498362dd044eb4a Mon Sep 17 00:00:00 2001 From: Bo Wu Date: Sat, 5 Sep 2026 18:42:03 -0700 Subject: [PATCH 5/5] Keep trace commitment delivery with Logger --- docker/runtime/Dockerfile | 4 ++-- logger/README.md | 13 +++++++++++++ .../trace_commitment.py | 2 +- tests/test_trace_commitment.py | 9 ++++++++- trace_commitment/README.md | 18 ------------------ trace_commitment/__init__.py | 1 - 6 files changed, 24 insertions(+), 23 deletions(-) rename {trace_commitment => logger}/trace_commitment.py (99%) delete mode 100644 trace_commitment/README.md delete mode 100644 trace_commitment/__init__.py diff --git a/docker/runtime/Dockerfile b/docker/runtime/Dockerfile index cc12c7f..d72e35f 100644 --- a/docker/runtime/Dockerfile +++ b/docker/runtime/Dockerfile @@ -20,10 +20,10 @@ RUN cd control-server && npm ci --omit=dev COPY . . COPY --from=multiagent-builder /src/target/release/multiagent /opt/multiagent/bin/multiagent RUN ln -s /opt/multiagent/wiki-service/bin/wiki-query.mjs /usr/local/bin/wiki-query \ - && ln -s /opt/multiagent/trace_commitment/trace_commitment.py /usr/local/bin/trace-commitment \ + && ln -s /opt/multiagent/logger/trace_commitment.py /usr/local/bin/trace-commitment \ && install -m 0755 docker/runtime/container-entrypoint.sh /opt/multiagent/bin/container-entrypoint.sh \ && install -m 0755 control-server/bin/prepare-repository.mjs /opt/multiagent/bin/prepare-repository.mjs \ - && chmod +x launch.sh control-server/bin/*.mjs trace_commitment/trace_commitment.py \ + && chmod +x launch.sh control-server/bin/*.mjs logger/trace_commitment.py \ && groupadd --gid 10000 multiagent-control \ && groupadd --gid 10001 multiagent-role \ && groupadd --gid 10004 multiagent-credentials \ diff --git a/logger/README.md b/logger/README.md index 1badcea..76ffbab 100644 --- a/logger/README.md +++ b/logger/README.md @@ -60,6 +60,19 @@ the trace body. Use a deterministic event ID and retain the event in a durable outbox until the Logger returns `204`, because acknowledgement remains idempotent transport evidence rather than workflow authority. +For the production sidecar, `trace_commitment.py` complements that low-level +CLI with durable S3 outbox delivery. The existing trace exporter remains +responsible for staging, delta detection, retry, and uploading trace bodies. +After a successful sync, `trace-commitment` hashes the stable staged tree, +writes a deterministic commitment manifest as a separate metadata object in +the same S3 location, and retries a bounded `trace.artifact_exported` event. + +Required environment is `TRACE_COMMITMENT_SOURCE`, `TRACE_EXPORT_DESTINATION`, +`TRACE_SESSION_ID`, `LOGGER_URL`, and `LOGGER_TOKEN_FILE`. +`TRACE_COMMITMENT_WORK_DIR` and `TRACE_EXPORT_STATUS_FILE` are optional. Logger +delivery updates `loggerOk`, `loggerPending`, and Logger timestamps/errors in +the existing status file without changing its S3 export `ok` field. + ## API and configuration The API exposes event append, log heads/entries/checkpoints, the public key, diff --git a/trace_commitment/trace_commitment.py b/logger/trace_commitment.py similarity index 99% rename from trace_commitment/trace_commitment.py rename to logger/trace_commitment.py index bf7528e..b7eff7b 100644 --- a/trace_commitment/trace_commitment.py +++ b/logger/trace_commitment.py @@ -1,5 +1,5 @@ #!/usr/bin/env python3 -"""Commit an already-exported trace tree to Logger without sending trace bodies.""" +"""Commit an already-exported trace tree to this Logger without sending trace bodies.""" from __future__ import annotations diff --git a/tests/test_trace_commitment.py b/tests/test_trace_commitment.py index 3f92a78..34bd431 100644 --- a/tests/test_trace_commitment.py +++ b/tests/test_trace_commitment.py @@ -1,11 +1,18 @@ #!/usr/bin/env python3 +import importlib.util import json import tempfile import unittest from pathlib import Path -from trace_commitment.trace_commitment import TraceCommitter + +MODULE_PATH = Path(__file__).parents[1] / "logger" / "trace_commitment.py" +SPEC = importlib.util.spec_from_file_location("logger_trace_commitment", MODULE_PATH) +assert SPEC is not None and SPEC.loader is not None +MODULE = importlib.util.module_from_spec(SPEC) +SPEC.loader.exec_module(MODULE) +TraceCommitter = MODULE.TraceCommitter class MemoryS3: diff --git a/trace_commitment/README.md b/trace_commitment/README.md deleted file mode 100644 index f669b99..0000000 --- a/trace_commitment/README.md +++ /dev/null @@ -1,18 +0,0 @@ -# Trace commitment delivery - -The deployment's existing trace exporter remains responsible for staging, -delta detection, retry, and uploading trace bodies. After a successful sync, -`python3 -m trace_commitment.trace_commitment` hashes that stable staged tree, -writes a deterministic commitment manifest as a separate object in the same S3 -location, and durably retries a bounded `trace.artifact_exported` Logger event. - -The Logger receives only the manifest digest, size, media type, and S3 -reference. The event and delivered marker are also separate metadata-only -objects beside the existing export. Deterministic IDs preserve idempotency after -an uncertain acknowledgement or process restart. - -Required environment is `TRACE_COMMITMENT_SOURCE`, `TRACE_EXPORT_DESTINATION`, -`TRACE_SESSION_ID`, `LOGGER_URL`, and `LOGGER_TOKEN_FILE`. -`TRACE_COMMITMENT_WORK_DIR` and `TRACE_EXPORT_STATUS_FILE` are optional. Logger -delivery updates `loggerOk`, `loggerPending`, and Logger timestamps/errors in -the existing status file without changing its S3 export `ok` field. diff --git a/trace_commitment/__init__.py b/trace_commitment/__init__.py deleted file mode 100644 index aed859b..0000000 --- a/trace_commitment/__init__.py +++ /dev/null @@ -1 +0,0 @@ -"""Post-upload trace commitment delivery."""