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..d72e35f 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/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 \ + && 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/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..c97435d 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; +- 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/logger/README.md b/logger/README.md index ca3962f..76ffbab 100644 --- a/logger/README.md +++ b/logger/README.md @@ -54,6 +54,25 @@ 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. + +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/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/logger/trace_commitment.py b/logger/trace_commitment.py new file mode 100644 index 0000000..b7eff7b --- /dev/null +++ b/logger/trace_commitment.py @@ -0,0 +1,314 @@ +#!/usr/bin/env python3 +"""Commit an already-exported trace tree to this 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 + + +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 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 + 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 = 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") + 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}/{relative.as_posix()}", + } + ) + 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 = 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) + 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/tests/run.sh b/tests/run.sh index f24814f..e255a6e 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_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_commitment.py b/tests/test_trace_commitment.py new file mode 100644 index 0000000..34bd431 --- /dev/null +++ b/tests/test_trace_commitment.py @@ -0,0 +1,111 @@ +#!/usr/bin/env python3 + +import importlib.util +import json +import tempfile +import unittest +from pathlib import Path + + +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: + def __init__(self): + self.objects = {} + + def upload(self, source, uri): + self.objects[uri] = Path(source).read_bytes() + + 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[len(root) :] for uri in self.objects if uri.startswith(root + prefix)} + + +class TraceCommitterTest(unittest.TestCase): + def setUp(self): + self.temporary = tempfile.TemporaryDirectory() + root = Path(self.temporary.name) + 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() + self.posts = [] + + def tearDown(self): + self.temporary.cleanup() + + def committer(self, post=None): + def successful(url, token, body): + self.posts.append((url, token, json.loads(body))) + return 204 + + return TraceCommitter( + source=self.source, + destination="s3://trace-bucket/production/sessions/session-1", + session_id="session-1", + work=self.work, + status_file=self.status, + 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_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.committer().commit() + + self.assertTrue(status["ok"]) + self.assertTrue(status["loggerOk"]) + self.assertEqual(status["loggerPending"], 0) + event = self.posts[0][2] + self.assertEqual(event["eventType"], "trace.artifact_exported") + 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") + self.committer().commit() + self.committer().commit() + self.assertEqual(len(self.posts), 1) + + 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.committer(unavailable).commit() + self.assertTrue(failed["ok"]) + self.assertFalse(failed["loggerOk"]) + self.assertEqual(failed["loggerPending"], 1) + + recovered = self.committer().commit() + self.assertTrue(recovered["loggerOk"]) + self.assertEqual(recovered["loggerPending"], 0) + self.assertEqual(len(self.posts), 1) + + +if __name__ == "__main__": + unittest.main()