feat(P1): event emitters — CloudEvents envelope, Decision Ledger, Infracost adapter, attestation/confidence/policy event emission

P1 (Wave 1, feat) — REQ-187, REQ-188, REQ-205 (emitter), REQ-206 (emitter)

New components:
- core/metrics/event_envelope.py — CloudEvents 1.0 envelope + platform.* conventions
- core/metrics/run_manifest.py — per-run manifest writer (nova.run.started/completed/failed)
- core/metrics/decision_ledger.py — SQLite append-only hash-chain (ai.decision.made + attestation.recorded)
- core/metrics/infracost_adapter.py — Infracost post-processor (degraded mode when CLI absent, A6)
- schemas/metrics_event.schema.json — CloudEvents envelope schema
- schemas/metrics_run_manifest.schema.json — per-run manifest schema
- metrics/README.md — backup/restore doc (REQ-201)
- tests/test_metrics_emitters.py — 16 tests (all pass)

Modified components:
- core/confidence_signal.py — emits nova.confidence.computed + nova.ai.decision.made (D-122)
- core/hitl_gates.py — emits nova.attestation.recorded on qa/prod/dr gates (D-132)
- adapters/terraform/policy/checkov_adapter.py — emits nova.policy.evaluated
- pyproject.toml — addopts gains --junitxml + --json-report + --cov (REQ-206)
- .gitignore — metrics runtime artifacts ignored

D-120: Nova-native (JSONL + SQLite, no Kafka/OTel)
D-121: Decision Ledger = outbox_writer extension → SQLite hash-chain
D-122: AI decision = confidence_signal + HITL gate (not LLM)
D-128: metrics/ at repo root
D-132: Attestation instrumentation

---ci---
project: acdl
phase: 1
milestone: v1.17
status: execute
---/ci---
This commit is contained in:
Jon Chery
2026-08-04 19:58:54 +00:00
parent fe2ab96b8c
commit f8616b806e
14 changed files with 1070 additions and 3 deletions
+32 -1
View File
@@ -34,8 +34,13 @@ per-input scores.
from dataclasses import dataclass, asdict
from typing import List, Literal, Optional, Dict, Any
import json
import os
import sys
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from core.metrics.event_envelope import emit, make_event, append_event
from core.metrics.decision_ledger import append as ledger_append
WEIGHTS = {
"policy": 0.30,
@@ -161,7 +166,33 @@ def compute(contract_id: str, environment: str,
band = "warn"
if environment == "dev" and band == "warn":
band = "block"
return Signal(score, band, per_input, reasons)
signal = Signal(score, band, per_input, reasons)
# Emit nova.confidence.computed + nova.ai.decision.made events (D-122).
# The "AI decision" is the confidence-gated policy engine, not an LLM.
# decision_id = run_id (or "cli-<ts>" when called from CLI without a run).
try:
run_id = os.environ.get("NOVA_RUN_ID", f"cli-{int(__import__('time').time())}")
conf_data = {"score": score, "band": band, "perInput": per_input, "reasonCodes": reasons}
emit("nova.confidence.computed", run_id, environment, conf_data, contract_id=contract_id)
decision_data = {
"decision_id": run_id,
"chosen_action": band,
"confidence": score,
"alternatives": per_input,
"human_override": band == "block",
"threshold": THRESHOLDS[environment],
}
decision_event = make_event("nova.ai.decision.made", run_id, environment, decision_data,
contract_id=contract_id, actor_type="confidence-gate",
actor_id="confidence_signal")
append_event(decision_event)
ledger_append(decision_event)
except Exception:
pass # metrics emission must never break the confidence gate
return signal
if __name__ == "__main__":
+22
View File
@@ -12,6 +12,10 @@ import os
import sys
from typing import Optional, Tuple
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.abspath(__file__))))
from core.metrics.event_envelope import make_event, append_event
from core.metrics.decision_ledger import append as ledger_append
def _approver_attr(env: str) -> str:
return {"qa": "approver_qa", "prod": "approver_prod", "dr": "approver_dr"}.get(env, "")
@@ -61,6 +65,24 @@ def attest(contract_id: str, env: str, approver: str,
if not ok:
return (False, reason)
# Emit attestation.recorded event to the Decision Ledger (D-132).
try:
run_id = os.environ.get("NOVA_RUN_ID", f"attest-{contract_id[:8]}")
attestation_data = {
"approver": approver,
"environment": env,
"concerns": reason,
"result": "pass",
"contract_id": contract_id,
}
attestation_event = make_event("nova.attestation.recorded", run_id, env, attestation_data,
contract_id=contract_id, actor_type="human-attestation",
actor_id=approver)
append_event(attestation_event)
ledger_append(attestation_event)
except Exception:
pass # metrics emission must never break the attestation gate
return (True, f"{env} attested by {approver}")
View File
+257
View File
@@ -0,0 +1,257 @@
"""Nova Decision Ledger — SQLite append-only hash-chain (REQ-188, D-121).
Extends outbox_writer.py to emit to a SQLite append-only table with a hash
chain (prev_hash + own hash, SHA-256). Stores ai.decision.made events
(decision_id=run_id, chosen_action=band, confidence=score,
alternatives=perInput, human_override=HITL block) with outcome backfill
from apply.completed. Also stores attestation.recorded events (D-132).
Honors D-083 (no S3 Object Lock/JWS — local SQLite hash-chain only).
D-120: Nova-native (SQLite, no QLDB).
D-128: metrics/ at repo root.
"""
import datetime
import hashlib
import json
import os
import sqlite3
import sys
_LEDGER_PATH = os.path.join(
os.path.dirname(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__))))),
"metrics", "decision_ledger.db",
)
_GENESIS_HASH = "GENESIS"
def _iso8601_now():
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def _canonical_hash(event):
"""SHA-256 over canonical JSON (sort_keys, compact separators)."""
canonical = json.dumps(event, sort_keys=True, separators=(",", ":"))
return hashlib.sha256(canonical.encode("utf-8")).hexdigest()
def _init_db(db_path=None):
"""Create the ledger table if it doesn't exist."""
if db_path is None:
db_path = _LEDGER_PATH
os.makedirs(os.path.dirname(db_path), exist_ok=True)
conn = sqlite3.connect(db_path)
conn.execute("""
CREATE TABLE IF NOT EXISTS decision_ledger (
seq INTEGER PRIMARY KEY AUTOINCREMENT,
event_id TEXT NOT NULL,
event_type TEXT NOT NULL,
run_id TEXT NOT NULL,
contract_id TEXT,
environment TEXT,
event_time TEXT NOT NULL,
payload TEXT NOT NULL,
prev_hash TEXT NOT NULL,
hash TEXT NOT NULL
)
""")
conn.execute("CREATE INDEX IF NOT EXISTS idx_run_id ON decision_ledger(run_id)")
conn.execute("CREATE INDEX IF NOT EXISTS idx_event_type ON decision_ledger(event_type)")
conn.commit()
conn.close()
def _get_last_hash(db_path=None):
"""Get the hash of the last row in the ledger (or GENESIS if empty)."""
if db_path is None:
db_path = _LEDGER_PATH
conn = sqlite3.connect(db_path)
row = conn.execute("SELECT hash FROM decision_ledger ORDER BY seq DESC LIMIT 1").fetchone()
conn.close()
return row[0] if row else _GENESIS_HASH
def append(event, db_path=None):
"""Append an event to the Decision Ledger with hash-chain integrity.
Args:
event: a CloudEvents 1.0 envelope dict (from event_envelope.make_event)
db_path: path to the SQLite ledger
Returns:
The row dict (seq, event_id, event_type, run_id, hash, prev_hash).
"""
if db_path is None:
db_path = _LEDGER_PATH
_init_db(db_path)
prev_hash = _get_last_hash(db_path)
event_hash = _canonical_hash(event)
platform = event.get("platform", {})
data = event.get("data", {})
conn = sqlite3.connect(db_path)
conn.execute("BEGIN IMMEDIATE")
cursor = conn.execute(
"""INSERT INTO decision_ledger
(event_id, event_type, run_id, contract_id, environment, event_time, payload, prev_hash, hash)
VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)""",
(
event.get("id", ""),
event.get("type", ""),
platform.get("run_id", ""),
platform.get("contract_id", ""),
platform.get("environment", ""),
event.get("time", _iso8601_now()),
json.dumps(event, sort_keys=True),
prev_hash,
event_hash,
),
)
seq = cursor.lastrowid
conn.commit()
conn.close()
return {"seq": seq, "event_id": event.get("id", ""), "event_type": event.get("type", ""),
"run_id": platform.get("run_id", ""), "hash": event_hash, "prev_hash": prev_hash}
def verify_chain(db_path=None):
"""Verify the hash chain integrity. Returns (ok, broken_count, details).
Recomputes each row's hash from its payload and checks:
1. The stored hash matches the recomputed hash.
2. The prev_hash matches the previous row's hash.
"""
if db_path is None:
db_path = _LEDGER_PATH
_init_db(db_path)
conn = sqlite3.connect(db_path)
rows = conn.execute("SELECT seq, hash, prev_hash, payload FROM decision_ledger ORDER BY seq").fetchall()
conn.close()
if not rows:
return True, 0, "empty ledger"
broken = 0
details = []
prev_hash = _GENESIS_HASH
for seq, stored_hash, stored_prev, payload_json in rows:
event = json.loads(payload_json)
recomputed = _canonical_hash(event)
if recomputed != stored_hash:
broken += 1
details.append(f"seq={seq}: hash mismatch (stored={stored_hash[:12]}... recomputed={recomputed[:12]}...)")
if stored_prev != prev_hash:
broken += 1
details.append(f"seq={seq}: prev_hash mismatch (expected={prev_hash[:12]}... got={stored_prev[:12]}...)")
prev_hash = stored_hash
return broken == 0, broken, "; ".join(details) if details else "chain intact"
def query_by_run(run_id, db_path=None):
"""Query all ledger entries for a given run_id."""
if db_path is None:
db_path = _LEDGER_PATH
_init_db(db_path)
conn = sqlite3.connect(db_path)
rows = conn.execute(
"SELECT seq, event_type, event_time, payload FROM decision_ledger WHERE run_id = ? ORDER BY seq",
(run_id,),
).fetchall()
conn.close()
return [{"seq": r[0], "event_type": r[1], "event_time": r[2], "payload": json.loads(r[3])} for r in rows]
def stats(db_path=None):
"""Return ledger statistics."""
if db_path is None:
db_path = _LEDGER_PATH
_init_db(db_path)
conn = sqlite3.connect(db_path)
total = conn.execute("SELECT COUNT(*) FROM decision_ledger").fetchone()[0]
by_type = conn.execute("SELECT event_type, COUNT(*) FROM decision_ledger GROUP BY event_type").fetchall()
by_env = conn.execute("SELECT environment, COUNT(*) FROM decision_ledger GROUP BY environment").fetchall()
conn.close()
return {
"total": total,
"by_event_type": dict(by_type),
"by_environment": dict(by_env),
}
def export_since(since_iso, fmt="json", db_path=None):
"""Export ledger entries since a given ISO8601 timestamp."""
if db_path is None:
db_path = _LEDGER_PATH
_init_db(db_path)
conn = sqlite3.connect(db_path)
rows = conn.execute(
"SELECT seq, event_type, run_id, event_time, payload FROM decision_ledger WHERE event_time >= ? ORDER BY seq",
(since_iso,),
).fetchall()
conn.close()
entries = [{"seq": r[0], "event_type": r[1], "run_id": r[2], "event_time": r[3], "payload": json.loads(r[4])} for r in rows]
if fmt == "csv":
import csv
import io
buf = io.StringIO()
writer = csv.DictWriter(buf, fieldnames=["seq", "event_type", "run_id", "event_time", "payload"])
writer.writeheader()
for e in entries:
e["payload"] = json.dumps(e["payload"])
writer.writerow(e)
return buf.getvalue()
return json.dumps(entries, indent=2)
def replay_run(run_id, db_path=None):
"""Reconstruct a run's full event sequence from the ledger.
Prints the ordered event sequence (run.started -> policy.evaluated ->
confidence.computed -> ai.decision.made -> attestation.recorded ->
run.completed/failed) with the decision's confidence, alternatives,
and outcome.
"""
if db_path is None:
db_path = _LEDGER_PATH
entries = query_by_run(run_id, db_path)
if not entries:
return f"no events found for run_id={run_id}"
lines = [f"=== Replay: run_id={run_id} ({len(entries)} events) ==="]
for e in entries:
payload = e["payload"]
data = payload.get("data", {})
etype = e["event_type"]
line = f" [{e['seq']}] {e['event_time']} {etype}"
if etype == "nova.ai.decision.made":
line += f" confidence={data.get('confidence', '?')} band={data.get('chosen_action', '?')} override={data.get('human_override', '?')}"
elif etype == "nova.attestation.recorded":
line += f" env={data.get('environment', '?')} approver={data.get('approver', '?')} result={data.get('result', '?')}"
elif etype == "nova.run.completed":
line += f" exit={data.get('exit_code', '?')} outcome={data.get('outcome', '?')}"
elif etype == "nova.run.failed":
line += f" exit={data.get('exit_code', '?')} outcome=failed"
lines.append(line)
lines.append("=== End replay ===")
return "\n".join(lines)
if __name__ == "__main__":
if len(sys.argv) < 2:
print("usage: decision_ledger.py <verify-chain|stats|query|export|replay> [args]", file=sys.stderr)
sys.exit(2)
cmd = sys.argv[1]
if cmd == "verify-chain":
ok, broken, details = verify_chain()
print(f"chain_ok={ok} broken={broken} details={details}")
sys.exit(0 if ok else 1)
elif cmd == "stats":
print(json.dumps(stats(), indent=2))
elif cmd == "query" and len(sys.argv) >= 3:
print(json.dumps(query_by_run(sys.argv[2]), indent=2))
elif cmd == "export" and len(sys.argv) >= 3:
print(export_since(sys.argv[2]))
elif cmd == "replay" and len(sys.argv) >= 3:
print(replay_run(sys.argv[2]))
else:
print(f"unknown command: {cmd}", file=sys.stderr)
sys.exit(2)
+98
View File
@@ -0,0 +1,98 @@
"""Nova CloudEvents 1.0 envelope + platform.* semantic conventions (REQ-187).
Defines the standard event envelope for all Nova metrics events. Every
emitter (run_manifest, decision_ledger, confidence_signal, checkov_adapter,
hitl_gates, regression_verify) uses `make_event()` to produce a valid
CloudEvents 1.0 envelope. Events are appended to `metrics/events.jsonl`.
D-120: Nova-native minimal tech (no Kafka/OTel SDK — JSONL + SQLite).
D-125: hybrid model — existing file signals stay as files; the collector
reads them and emits normalized CloudEvents. New emitters emit directly.
"""
import datetime
import hashlib
import json
import os
import sys
import uuid
METRICS_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), "metrics")
EVENTS_LOG = os.path.join(METRICS_DIR, "events.jsonl")
def _iso8601_now():
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def make_event(event_type, run_id, environment, data, contract_id="", source="nova.platform", subject="", actor_type="confidence-gate", actor_id="confidence_signal"):
"""Build a CloudEvents 1.0 envelope with Nova platform.* conventions.
Args:
event_type: e.g. "nova.run.completed", "nova.ai.decision.made"
run_id: the run identifier (e.g. "run-<epoch>")
environment: dev|qa|prod|dr
data: the event payload dict
contract_id: the contract UUID (optional)
source: the event source (default "nova.platform")
subject: the event subject (default "<contract_id>/<env>")
actor_type: the actor type (default "confidence-gate")
actor_id: the actor id (default "confidence_signal")
Returns:
A CloudEvents 1.0 envelope dict.
"""
if not subject:
subject = f"{contract_id}/{environment}" if contract_id else environment
return {
"specversion": "1.0",
"id": str(uuid.uuid4()),
"source": source,
"type": event_type,
"time": _iso8601_now(),
"subject": subject,
"datacontenttype": "application/json",
"platform": {
"tenant_id": "acdl",
"run_id": run_id,
"contract_id": contract_id,
"environment": environment,
"actor": {"type": actor_type, "id": actor_id},
"trace_id": run_id,
},
"data": data,
}
def append_event(event, events_log=None):
"""Append a CloudEvents envelope to the JSONL event log.
Creates the metrics/ directory if it doesn't exist.
"""
if events_log is None:
events_log = EVENTS_LOG
os.makedirs(os.path.dirname(events_log), exist_ok=True)
with open(events_log, "a", encoding="utf-8") as fh:
fh.write(json.dumps(event, sort_keys=True, separators=(",", ":")) + "\n")
def emit(event_type, run_id, environment, data, **kwargs):
"""Make an event + append it to the JSONL log. Convenience wrapper."""
event = make_event(event_type, run_id, environment, data, **kwargs)
append_event(event)
return event
if __name__ == "__main__":
if len(sys.argv) < 4:
print("usage: event_envelope.py <event_type> <run_id> <environment> [data.json]", file=sys.stderr)
sys.exit(2)
_type = sys.argv[1]
_run_id = sys.argv[2]
_env = sys.argv[3]
_data = {}
if len(sys.argv) >= 5 and os.path.isfile(sys.argv[4]):
with open(sys.argv[4]) as f:
_data = json.load(f)
ev = emit(_type, _run_id, _env, _data)
print(json.dumps(ev, indent=2))
+73
View File
@@ -0,0 +1,73 @@
"""Nova Infracost Post-Processor (REQ-187, D-120).
Runs Infracost on `terraform show -json plan.tfplan` (offline, reads plan
JSON, no live AWS). Emits nova.cost.estimated{delta_usd} events. Degrades
gracefully (omits the event, logs a warning) when Infracost CLI is absent
(assumption A6).
run_platform.sh invokes it after the plan stage.
"""
import json
import os
import shutil
import subprocess
import sys
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
from core.metrics.event_envelope import emit
def _is_infracost_available():
"""Check if the Infracost CLI is on PATH."""
return shutil.which("infracost") is not None
def estimate(plan_json_path, run_id, contract_id, environment):
"""Run Infracost on a terraform plan JSON. Returns the cost estimate dict.
Args:
plan_json_path: path to `terraform show -json plan.tfplan` output
run_id: the run identifier
contract_id: the contract UUID
environment: dev|qa|prod|dr
Returns:
{"delta_usd": float, "total_monthly_usd": float, "available": bool}
or {"available": False} if Infracost is not installed.
"""
if not _is_infracost_available():
sys.stderr.write("[infracost] CLI not found — cost.estimated event omitted (A6 degraded mode)\n")
return {"available": False, "delta_usd": 0.0, "total_monthly_usd": 0.0}
if not os.path.isfile(plan_json_path):
sys.stderr.write(f"[infracost] plan JSON not found: {plan_json_path}\n")
return {"available": False, "delta_usd": 0.0, "total_monthly_usd": 0.0}
try:
result = subprocess.run(
["infracost", "breakdown", "--path", plan_json_path, "--format", "json"],
capture_output=True, text=True, timeout=30,
)
if result.returncode != 0:
sys.stderr.write(f"[infracost] CLI failed: {result.stderr[:200]}\n")
return {"available": False, "delta_usd": 0.0, "total_monthly_usd": 0.0}
breakdown = json.loads(result.stdout)
delta = float(breakdown.get("diffTotalMonthlyCost", 0.0))
total = float(breakdown.get("totalMonthlyCost", 0.0))
estimate_data = {"available": True, "delta_usd": delta, "total_monthly_usd": total}
emit("nova.cost.estimated", run_id, environment, estimate_data, contract_id=contract_id)
return estimate_data
except Exception as exc:
sys.stderr.write(f"[infracost] error: {exc}\n")
return {"available": False, "delta_usd": 0.0, "total_monthly_usd": 0.0}
if __name__ == "__main__":
if len(sys.argv) < 5:
print("usage: infracost_adapter.py <plan_json_path> <run_id> <contract_id> <environment>", file=sys.stderr)
sys.exit(2)
est = estimate(sys.argv[1], sys.argv[2], sys.argv[3], sys.argv[4])
print(json.dumps(est, indent=2))
+137
View File
@@ -0,0 +1,137 @@
"""Nova Per-Run Manifest Writer (REQ-187).
Emits nova.run.started, nova.run.completed, nova.run.failed events with
(run_id, contractId, env, stages x durations, exit, confidence, HITL block
count). Writes metrics/runs/<run_id>.json. scripts/run_platform.sh invokes
the writer at run start + run end.
D-120: Nova-native (JSONL events + JSON manifest file, no Kafka).
D-128: metrics/ at repo root.
"""
import datetime
import json
import os
import sys
import time
import uuid
_METRICS_DIR = os.path.join(os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))), "metrics")
_RUNS_DIR = os.path.join(_METRICS_DIR, "runs")
sys.path.insert(0, os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))))
from core.metrics.event_envelope import emit, make_event, append_event
def _iso8601_now():
return datetime.datetime.now(datetime.timezone.utc).strftime("%Y-%m-%dT%H:%M:%SZ")
def _run_id():
return f"run-{int(time.time())}-{uuid.uuid4().hex[:8]}"
def start_run(contract_id, environment, stages=None):
"""Emit nova.run.started + return the run_id."""
run_id = _run_id()
data = {
"contract_id": contract_id,
"environment": environment,
"started_at": _iso8601_now(),
"stages": stages or [],
}
emit("nova.run.started", run_id, environment, data, contract_id=contract_id)
return run_id
def complete_run(run_id, contract_id, environment, stages, exit_code, confidence=None, hitl=None, policy=None, cost_estimate_usd=None, decision_id=None):
"""Emit nova.run.completed + write the per-run manifest JSON.
Args:
run_id: the run identifier from start_run()
contract_id: the contract UUID
environment: dev|qa|prod|dr
stages: list of {name, duration_ms, exit_code, error?}
exit_code: the overall run exit code
confidence: optional {score, band, perInput}
hitl: optional {gate, result, block}
policy: optional {passed, failed, skipped}
cost_estimate_usd: optional float
decision_id: optional string (links to the Decision Ledger)
"""
started_at = stages[0].get("started_at", _iso8601_now()) if stages else _iso8601_now()
completed_at = _iso8601_now()
outcome = "succeeded" if exit_code == 0 else "failed"
manifest = {
"run_id": run_id,
"contract_id": contract_id,
"environment": environment,
"started_at": started_at,
"completed_at": completed_at,
"exit_code": exit_code,
"stages": stages,
"outcome": outcome,
}
if confidence:
manifest["confidence"] = confidence
if hitl:
manifest["hitl"] = hitl
if policy:
manifest["policy"] = policy
if cost_estimate_usd is not None:
manifest["cost_estimate_usd"] = cost_estimate_usd
if decision_id:
manifest["decision_id"] = decision_id
os.makedirs(_RUNS_DIR, exist_ok=True)
manifest_path = os.path.join(_RUNS_DIR, f"{run_id}.json")
with open(manifest_path, "w", encoding="utf-8") as fh:
json.dump(manifest, fh, indent=2, sort_keys=True)
event_type = "nova.run.completed" if exit_code == 0 else "nova.run.failed"
emit(event_type, run_id, environment, manifest, contract_id=contract_id)
return manifest
def persist_run_artifacts(run_id, work_dir):
"""Copy ephemeral $WORK/*.json to metrics/runs/<run_id>/ as durable artifacts.
Args:
run_id: the run identifier
work_dir: the $WORK directory (e.g. /tmp/nova_platform_run)
"""
if not work_dir or not os.path.isdir(work_dir):
return []
dest = os.path.join(_RUNS_DIR, run_id)
os.makedirs(dest, exist_ok=True)
copied = []
for fname in ("pcr.json", "signal.json", "event.json", "outbox_item.json", "stack.json", "checkov.json"):
src = os.path.join(work_dir, fname)
if os.path.isfile(src):
import shutil
shutil.copy2(src, os.path.join(dest, fname))
copied.append(fname)
return copied
if __name__ == "__main__":
if len(sys.argv) < 4:
print("usage: run_manifest.py <start|complete|persist> <contract_id> <environment> [run_id] [work_dir]", file=sys.stderr)
sys.exit(2)
action = sys.argv[1]
cid = sys.argv[2]
env = sys.argv[3]
if action == "start":
rid = start_run(cid, env)
print(rid)
elif action == "complete":
rid = sys.argv[4] if len(sys.argv) >= 5 else _run_id()
m = complete_run(rid, cid, env, [], 0)
print(json.dumps(m, indent=2))
elif action == "persist":
rid = sys.argv[4] if len(sys.argv) >= 5 else ""
wd = sys.argv[5] if len(sys.argv) >= 6 else ""
copied = persist_run_artifacts(rid, wd)
print(json.dumps({"copied": copied}))