"""Covenant Daily Replay — an automated daily self-check, operated by Alpha Covenant Holdings. Each day it rebuilds a fixed panel of receipts and diffs them against yesterday, recomputes yesterday's registered receipt hashes, re-reads the source board, fetches every disclosure surface, and writes the whole run into a hash-chained ledger whose entry hash is anchored outside our control (OpenTimestamps calendars, Internet Archive). A model writes a plain-English summary FROM the results; it never decides anything. This is not an independent audit and nobody involved holds an audit credential.""" import asyncio import base64 import hashlib import hmac import json import logging import os from datetime import datetime, timedelta, timezone import httpx from fastapi import APIRouter, Depends, HTTPException, Request from fastapi.encoders import jsonable_encoder from fastapi.responses import HTMLResponse, Response from deps import db, require_operator logger = logging.getLogger("selfaudit") router = APIRouter(prefix="/api/audit") cron_router = APIRouter(prefix="/api/cron") JOB = "self_audit_daily" NAME = "Covenant Daily Replay" MODEL = ("anthropic", "claude-sonnet-5-5") WHAT_THIS_IS = ("Automated daily self-check, operated by Alpha Covenant Holdings. Every check is deterministic and reproducible from the published code; " "the summary is machine-written from the check results and decides nothing. This is not an independent audit and no one involved holds an audit credential. " "Anyone can run the same script and publish a different answer.") WHAT_THIS_IS_NOT = "Not an independent audit. Not a certification. Not a statement about any name on the panel — the panel exists only so the same receipts can be rebuilt every day and compared." DEFAULT_PANEL = ["Wellpath LLC", "JT Medical LLC", "Corizon Health", "Maximus Inc", "Centene Corporation"] DISCLOSURES = ["/api/public/conflicts", "/api/public/conflicts/html", "/api/public-record/known-limits", "/api/public-record/known-limits/html", "/api/public/sealed", "/api/record/reviewer-packet", "/api/public/accuracy", "/api/legal/product", "/api/receipt-watch/plans", "/api/public-record/match-register", "/api/public/code/replay.py", "/api/public/code/selfaudit.py", "/manifest.json"] OTS_CALENDARS = ["https://a.pool.opentimestamps.org", "https://b.pool.opentimestamps.org", "https://alice.btc.calendar.opentimestamps.org", "https://bob.btc.calendar.opentimestamps.org"] OTS_HEADER = b"\x00OpenTimestamps\x00\x00Proof\x00\xbf\x89\xe2\xe8\x84\xe8\x92\x94" + b"\x01" + b"\x08" UA = {"User-Agent": "AlphaCovenantHoldings/1.0 (daily-replay; contact@alphacovenantholdings.com)"} def _base() -> str: return (os.environ.get("PUBLIC_BASE_URL") or "https://alphacovenantholdings.com").rstrip("/") def canonical(o) -> bytes: return json.dumps(jsonable_encoder(o), sort_keys=True, default=str, separators=(",", ":")).encode() def h(o) -> str: return hashlib.sha256(canonical(o)).hexdigest() def entry_hash(e: dict) -> str: return h({"date": e["date"], "at": e["at"], "prev_entry_hash": e["prev_entry_hash"], "results_hash": e["results_hash"]}) def receipt_hash(body: dict) -> str: b = {k: v for k, v in body.items() if k not in ("sha256", "receipt_registry", "_id")} return hashlib.sha256(json.dumps(jsonable_encoder(b), sort_keys=True, default=str).encode()).hexdigest() async def panel() -> list[str]: p = await db.self_audit_panel.find_one({"_id": "panel"}) if p: return p["names"] names = list(DEFAULT_PANEL) async for r in db.receipt_registry.find({"name": {"$ne": None}}, {"name": 1}).sort("created_at", -1).limit(400): n = (r.get("name") or "").strip() if n and n.lower() not in {x.lower() for x in names}: names.append(n) if len(names) >= 25: break await db.self_audit_panel.insert_one({"_id": "panel", "names": names, "frozen_at": datetime.now(timezone.utc).isoformat(), "how": "5 Receipt Watch trial names plus up to 20 distinct names from past published receipts; frozen so day-to-day comparison is like-for-like"}) return names async def _receipts(names: list[str], prev: dict | None) -> list[dict]: import screening, receiptspec, receiptwatch prev_by = {r["name"]: r for r in ((prev or {}).get("results", {}).get("receipts") or []) if "sha256" in r} sem = asyncio.Semaphore(3) async def one(n): t0 = datetime.now(timezone.utc) async with sem: try: d = await screening.build(n, "", "", "", "", "", "") reg = await receiptspec.record(d, True, "self-audit") except Exception as e: return {"name": n, "error": f"receipt could not be built ({type(e).__name__})"} snap = receiptwatch.snapshot(d) p = prev_by.get(n) return {"name": n, "receipt_id": reg["receipt_id"], "sha256": d["sha256"], "prev_sha256": p["sha256"] if p else None, "prev_receipt_id": p["receipt_id"] if p else None, "diff": receiptwatch.diff(p["snapshot"], snap) if p else [], "snapshot": snap, "did_not_answer": snap["did_not_answer"], "not_read": snap["not_read"], "build_ms": int((datetime.now(timezone.utc) - t0).total_seconds() * 1000)} return list(await asyncio.gather(*[one(n) for n in names])) async def _recompute(prev: dict | None) -> dict: ids = [r["receipt_id"] for r in ((prev or {}).get("results", {}).get("receipts") or []) if r.get("receipt_id")] out = {"checked": 0, "identical": 0, "differs": [], "missing": [], "what": "Yesterday's registered receipt bodies re-hashed with the published rule (served JSON minus sha256/receipt_registry, sorted keys)."} for rid in ids: row = await db.receipt_registry.find_one({"receipt_id": rid}, {"_id": 0, "body": 1, "sha256": 1}) if not row or not row.get("body"): out["missing"].append(rid) continue out["checked"] += 1 if receipt_hash(row["body"]) == row["sha256"]: out["identical"] += 1 else: out["differs"].append(rid) return out async def _sources(prev: dict | None) -> dict: import sourcesboard sourcesboard._cache.clear() b = await sourcesboard.screening_sources() prev_files = {f["id"]: f for f in ((prev or {}).get("results", {}).get("sources", {}).get("files") or [])} files = [{k: f.get(k) for k in ("id", "name", "state", "rows", "sha256", "fetched_at")} for f in b["files"]] live = [{k: l.get(k) for k in ("id", "name", "state", "http", "latency_ms", "error")} for l in b["live"]] changes = [{"id": f["id"], "name": f["name"], "old_rows": prev_files[f["id"]].get("rows"), "new_rows": f.get("rows"), "old_sha256": prev_files[f["id"]].get("sha256"), "new_sha256": f.get("sha256")} for f in files if f["id"] in prev_files and (prev_files[f["id"]].get("sha256") != f.get("sha256") or prev_files[f["id"]].get("rows") != f.get("rows"))] return {"checked_at": b["checked_at"], "files": files, "live": live, "files_read": b["read_now"], "files_total": len(files), "live_reachable": b["reachable_now"], "live_total": len(live), "held_file_changes": changes} async def _disclosures(prev: dict | None) -> list[dict]: prev_by = {d["path"]: d for d in ((prev or {}).get("results", {}).get("disclosures") or [])} base = _base() out = [] async with httpx.AsyncClient(timeout=60, follow_redirects=True, headers=UA) as client: for p in DISCLOSURES: t0 = datetime.now(timezone.utc) row = {"path": p, "http": None, "latency_ms": None, "bytes": None, "sha256": None} try: r = await client.get(base + p) row.update(http=r.status_code, bytes=len(r.content), sha256=hashlib.sha256(r.content).hexdigest()) except Exception as e: row["error"] = type(e).__name__ row["latency_ms"] = int((datetime.now(timezone.utc) - t0).total_seconds() * 1000) pb = (prev_by.get(p) or {}).get("bytes") row["prev_bytes"] = pb row["size_change_pct"] = round((row["bytes"] - pb) * 100 / pb, 1) if pb and row["bytes"] is not None else None out.append(row) return out def _flags(res: dict) -> list[dict]: f = [] for r in res["receipts"]: if "error" in r: f.append({"kind": "receipt_not_built", "detail": f"{r['name']}: {r['error']}"}) elif r["did_not_answer"]: f.append({"kind": "source_did_not_answer", "detail": f"{r['name']}: {', '.join(r['did_not_answer'])}"}) rc = res["registry_recompute"] for rid in rc["differs"]: f.append({"kind": "registered_hash_differs", "detail": f"{rid}: stored sha256 does not equal the recomputed hash of the stored body"}) for rid in rc["missing"]: f.append({"kind": "registered_body_missing", "detail": rid}) s = res["sources"] for x in s["files"]: if x["state"] != "read": f.append({"kind": "held_file_not_read", "detail": f"{x['name']}: {x['state']}"}) for x in s["live"]: if not str(x["state"]).startswith("reachable"): f.append({"kind": "live_source_unreachable", "detail": f"{x['name']}: {x['state']} ({x.get('error') or x.get('http')})"}) for c in s["held_file_changes"]: f.append({"kind": "held_file_changed", "detail": f"{c['name']}: rows {c['old_rows']} → {c['new_rows']}"}) for d in res["disclosures"]: if d["http"] != 200: f.append({"kind": "disclosure_not_200", "detail": f"{d['path']}: {d.get('http') or d.get('error')}"}) elif d["size_change_pct"] is not None and abs(d["size_change_pct"]) > 25: f.append({"kind": "disclosure_size_moved", "detail": f"{d['path']}: {d['prev_bytes']} → {d['bytes']} bytes ({d['size_change_pct']:+}%)"}) return f def _compact(e: dict) -> dict: r = e["results"] return {"date": e["date"], "compared_against_previous_entry": e.get("prev_date") or "none — this is the first entry, so there is no previous receipt or size to compare against", "panel_size": len(r["receipts"]), "receipts": [{"name": x["name"], **({"error": x["error"]} if "error" in x else {"changed_fields": x["diff"], "sources_did_not_answer": x["did_not_answer"], "sources_not_read_paid_or_gated": x["not_read"], "first_run": x["prev_sha256"] is None})} for x in r["receipts"]], "registry_recompute": {k: r["registry_recompute"][k] for k in ("checked", "identical", "differs", "missing")}, "sources": {k: r["sources"][k] for k in ("files_read", "files_total", "live_reachable", "live_total", "held_file_changes")}, "live_not_reachable": [x["name"] for x in r["sources"]["live"] if not str(x["state"]).startswith("reachable")], "disclosures": [{"path": d["path"], "http": d["http"], "latency_ms": d["latency_ms"], "size_change_pct": d["size_change_pct"]} for d in r["disclosures"]], "flags": r["flags"]} def _fallback_summary(e: dict) -> str: r, c = e["results"], _compact(e) built = [x for x in r["receipts"] if "error" not in x] changed = [x for x in built if x["diff"]] return (f"{e['date']}: {len(built)} of {len(r['receipts'])} panel receipts rebuilt; {len(changed)} printed a different line from the previous receipt. " f"Registry recompute: {r['registry_recompute']['identical']} of {r['registry_recompute']['checked']} hashes identical. " f"Sources: {r['sources']['files_read']}/{r['sources']['files_total']} held files read, {r['sources']['live_reachable']}/{r['sources']['live_total']} live sources reachable" + (f" (not reachable: {', '.join(c['live_not_reachable'])})" if c["live_not_reachable"] else "") + ". " f"Disclosures: {sum(1 for d in r['disclosures'] if d['http'] == 200)}/{len(r['disclosures'])} answered 200. {len(r['flags'])} flag(s) for a person to look at.") SYSTEM = ("You write the daily plain-English summary for Alpha Covenant Holdings' automated self-check. You receive JSON check results and nothing else. " "Rules, absolute: describe only what the JSON states; every sentence must trace to a specific field; never judge, grade, score, rate, reassure, speculate about causes, " "or use words like 'healthy', 'passed', 'safe', 'compliant', 'good', 'concerning'. Say 'flagged for a person to look at', not 'problem'. " "Do not characterise any company on the panel. Write 3-6 short sentences, no headings, no bullets, no markdown. If a field is missing, say it was not reported.") async def _narrate(e: dict) -> dict: key = os.environ.get("EMERGENT_LLM_KEY") now = datetime.now(timezone.utc).isoformat() if not key: return {"text": _fallback_summary(e), "written_by": "deterministic template (no model key configured)", "written_at": now} try: from emergentintegrations.llm.chat import LlmChat, UserMessage chat = LlmChat(api_key=key, session_id=f"selfaudit-{e['date']}", system_message=SYSTEM).with_model(*MODEL) text = await asyncio.wait_for(chat.send_message(UserMessage(text=json.dumps(_compact(e), default=str))), 90) return {"text": (text or "").strip() or _fallback_summary(e), "written_by": f"{MODEL[1]} from the check results (summarises; decides nothing)", "written_at": now} except Exception as ex: logger.warning(f"selfaudit narrate: {ex}") return {"text": _fallback_summary(e), "written_by": f"deterministic template (model did not answer: {type(ex).__name__})", "written_at": now} async def _anchor(digest_hex: str, date: str) -> dict: digest = bytes.fromhex(digest_hex) ots, wb = [], [] async with httpx.AsyncClient(timeout=30, follow_redirects=True) as client: for cal in OTS_CALENDARS: try: r = await client.post(f"{cal}/digest", content=digest, headers={"Accept": "application/vnd.opentimestamps.v1", "User-Agent": "AlphaCovenantHoldings/1.0 opentimestamps"}) ots.append({"calendar": cal, "http": r.status_code, "proof_b64": base64.b64encode(r.content).decode() if r.status_code == 200 and r.content else None}) except Exception as e: ots.append({"calendar": cal, "http": None, "error": type(e).__name__}) base = _base() for t in ("/audit", f"/api/audit/{date}", "/api/audit/html"): try: r = await client.get(f"https://web.archive.org/save/{base}{t}", headers=UA, timeout=90) wb.append({"target": t, "save_http": r.status_code, "browse": f"https://web.archive.org/web/*/{base.split('://', 1)[-1]}{t}"}) except Exception as e: wb.append({"target": t, "save_http": None, "error": type(e).__name__, "browse": f"https://web.archive.org/web/*/{base.split('://', 1)[-1]}{t}"}) return {"opentimestamps": ots, "wayback": wb, "what": "The entry hash was submitted to public OpenTimestamps calendars (Bitcoin-anchored; download the .ots and verify with any OTS client) and the day's pages were submitted to the Internet Archive. Neither is operated by us."} def ots_file(digest_hex: str, proof_b64: str) -> bytes: return OTS_HEADER + bytes.fromhex(digest_hex) + base64.b64decode(proof_b64) def next_fire(now: datetime | None = None) -> datetime: """Next 10:05 UTC — five minutes after the platform cron, so whichever fires first does the work and the other is skipped as 'already ran today'.""" now = now or datetime.now(timezone.utc) t = now.replace(hour=10, minute=5, second=0, microsecond=0) return t if t > now else t + timedelta(days=1) async def ensure_armed(): """Self-arming: the app never depends on an external scheduler. Called at startup; re-called by run().""" import scheduler as scheduler_mod today = datetime.now(timezone.utc).strftime("%Y-%m-%d") if await db.jobs.find_one({"type": JOB, "status": "pending"}): return "pending" if not await db.self_audit_runs.find_one({"date": today}): await scheduler_mod.enqueue(db, JOB, {}, run_at=datetime.now(timezone.utc) + timedelta(minutes=2)) return "queued-now" await scheduler_mod.enqueue(db, JOB, {}, run_at=next_fire()) return "armed-tomorrow" async def run(db_=None, payload: dict | None = None) -> dict: import scheduler as scheduler_mod try: return await _run(payload) except Exception as e: logger.exception(f"self_audit_daily failed: {e}") try: import emailer to = (os.environ.get("ALERT_EMAIL_TO", "") or os.environ.get("REPLY_TO_EMAIL", "")).strip() if to: await emailer.send_alert(to, "[Alpha Covenant Ops] Daily Replay did not complete", f"The self-check job raised {type(e).__name__}: {e}\n\nIt is re-armed for the next 10:05 UTC. Environment: {_base()}") except Exception: pass return {"error": type(e).__name__} finally: if not (payload or {}).get("date"): await scheduler_mod.enqueue(db, JOB, {}, run_at=next_fire()) async def _run(payload: dict | None = None) -> dict: payload = payload or {} now = datetime.now(timezone.utc) date = payload.get("date") or now.strftime("%Y-%m-%d") existing = await db.self_audit_runs.find_one({"date": date}, {"_id": 0}) if existing and not payload.get("force"): return {"skipped": "already ran today", "date": date, "entry_hash": existing["entry_hash"]} prev = await db.self_audit_runs.find_one({"date": {"$lt": date}}, {"_id": 0}, sort=[("date", -1)]) names = payload.get("panel") or await panel() receipts, recompute, sources, disclosures = await asyncio.gather(_receipts(names, prev), _recompute(prev), _sources(prev), _disclosures(prev)) results = {"receipts": receipts, "registry_recompute": recompute, "sources": sources, "disclosures": disclosures} results["flags"] = _flags(results) e = {"date": date, "at": now.isoformat(), "prev_entry_hash": prev["entry_hash"] if prev else None, "prev_date": prev["date"] if prev else None, "results": results, "results_hash": h(results)} e["entry_hash"] = entry_hash(e) e["summary"] = await _narrate(e) e["anchors"] = await _anchor(e["entry_hash"], date) if not payload.get("no_anchor") else {"skipped": True} e["panel"] = names e["what_this_is"], e["what_this_is_not"] = WHAT_THIS_IS, WHAT_THIS_IS_NOT if existing: await db.self_audit_runs.delete_one({"date": date}) await db.self_audit_runs.insert_one(dict(e)) logger.info(f"self_audit_daily {date}: {len(results['flags'])} flags, entry {e['entry_hash'][:16]}") if results["flags"] and not payload.get("no_email"): await _alert(e) return {k: e[k] for k in ("date", "entry_hash", "results_hash", "prev_entry_hash")} | {"flags": len(results["flags"]), "summary": e["summary"]["text"]} async def _alert(e: dict): import emailer to = (os.environ.get("ALERT_EMAIL_TO", "") or os.environ.get("REPLY_TO_EMAIL", "")).strip() if not to: return env = "production" if "alphacovenantholdings.com" in _base() else f"preview ({_base()})" lines = [f"{NAME} — {e['date']} — {len(e['results']['flags'])} item(s) flagged for a person to look at. Environment: {env}.", "", e["summary"]["text"], ""] lines += [f" · {f['kind']}: {f['detail']}" for f in e["results"]["flags"]] lines += ["", f"Full entry: {_base()}/audit · entry hash {e['entry_hash']}", "", WHAT_THIS_IS, "", "— Alpha Covenant Ops"] try: await emailer.send_alert(to, f"[Alpha Covenant Ops] Daily Replay {e['date']}: {len(e['results']['flags'])} flagged", "\n".join(lines)) except Exception as ex: logger.warning(f"selfaudit alert: {ex}") def _public(e: dict) -> dict: e = dict(e) for r in e.get("results", {}).get("receipts", []): r.pop("snapshot", None) for a in (e.get("anchors") or {}).get("opentimestamps") or []: a["proof_available"] = bool(a.pop("proof_b64", None)) return e @router.get("") async def latest(): e = await db.self_audit_runs.find_one({}, {"_id": 0}, sort=[("date", -1)]) n = await db.self_audit_runs.count_documents({}) return {"name": NAME, "operated_by": "Alpha Covenant Holdings", "what_this_is": WHAT_THIS_IS, "what_this_is_not": WHAT_THIS_IS_NOT, "cadence": "daily, 10:00 UTC", "entries": n, "checker_source": "/api/public/code/selfaudit.py", "ledger": "/api/audit/ledger", "verify_chain": "/api/audit/verify", "script_free": "/api/audit/html", "summary_model": {"provider": MODEL[0], "model": MODEL[1], "role": "writes the summary from the check results; decides nothing; its output is not part of the hashed entry"}, "latest": _public(e) if e else None} @router.get("/ledger") async def ledger(days: int = 60): rows = await db.self_audit_runs.find({}, {"_id": 0, "date": 1, "at": 1, "entry_hash": 1, "prev_entry_hash": 1, "results_hash": 1, "summary.text": 1, "results.flags": 1}).sort("date", -1).to_list(max(1, min(days, 400))) return {"name": NAME, "what_this_is": WHAT_THIS_IS, "entries": [{"date": r["date"], "at": r["at"], "entry_hash": r["entry_hash"], "prev_entry_hash": r["prev_entry_hash"], "results_hash": r["results_hash"], "flags": len(r["results"]["flags"]), "summary": r["summary"]["text"]} for r in rows]} @router.get("/verify") async def verify(): rows = await db.self_audit_runs.find({}, {"_id": 0}).sort("date", 1).to_list(5000) breaks, prev = [], None for r in rows: if h(r["results"]) != r["results_hash"]: breaks.append({"date": r["date"], "what": "results_hash does not equal sha256 of stored results"}) if entry_hash(r) != r["entry_hash"]: breaks.append({"date": r["date"], "what": "entry_hash does not equal sha256 of (date, at, prev_entry_hash, results_hash)"}) if (prev["entry_hash"] if prev else None) != r["prev_entry_hash"]: breaks.append({"date": r["date"], "what": "prev_entry_hash does not equal the previous entry's hash"}) prev = r return {"entries": len(rows), "chain_intact": not breaks, "breaks": breaks, "rule": "entry_hash = sha256(canonical JSON of {date, at, prev_entry_hash, results_hash}); results_hash = sha256(canonical JSON of results); canonical = sorted keys, compact separators. The summary and anchors are outside the hash."} @router.get("/html", response_class=HTMLResponse) async def html(): import html as H e = await db.self_audit_runs.find_one({}, {"_id": 0}, sort=[("date", -1)]) esc = lambda x: H.escape(str(x if x is not None else "—")) if not e: return HTMLResponse(f"
{esc(WHAT_THIS_IS)}
No entry yet.
") r = e["results"] rec = "".join(f"{esc((x.get('sha256') or '')[:16])}{esc(WHAT_THIS_IS)}
Entry hash {esc(e['entry_hash'])}
Previous entry {esc(e['prev_entry_hash'])}
Results hash {esc(e['results_hash'])} · run at {esc(e['at'])}
{esc(e['summary']['text'])}
Written by: {esc(e['summary']['written_by'])}
| Name | Receipt | SHA-256 | Fields changed vs. previous | Did not answer |
|---|
{esc(r['registry_recompute']['identical'])} of {esc(r['registry_recompute']['checked'])} of yesterday's registered receipts re-hash to their stored SHA-256; differs: {esc(len(r['registry_recompute']['differs']))}; missing: {esc(len(r['registry_recompute']['missing']))}.
| Path | HTTP | Latency | Bytes |
|---|
OpenTimestamps (download .ots):
Internet Archive:
{esc(WHAT_THIS_IS_NOT)} Ledger: /api/audit/ledger · chain check: /api/audit/verify · code: /api/public/code/selfaudit.py
""") @router.get("/{date}/proof.ots") async def proof(date: str): e = await db.self_audit_runs.find_one({"date": date}, {"_id": 0, "entry_hash": 1, "anchors": 1}) if not e: raise HTTPException(404, "No entry for that date") proofs = [a for a in (e.get("anchors") or {}).get("opentimestamps") or [] if a.get("proof_b64")] if not proofs: raise HTTPException(404, "No calendar accepted this entry's digest; nothing to download") return Response(ots_file(e["entry_hash"], proofs[0]["proof_b64"]), media_type="application/vnd.opentimestamps.v1", headers={"Content-Disposition": f'attachment; filename="covenant-daily-replay-{date}.ots"'}) @router.get("/{date}") async def by_date(date: str): e = await db.self_audit_runs.find_one({"date": date}, {"_id": 0}) if not e: raise HTTPException(404, "No entry for that date") return _public(e) @router.post("/run-now") async def run_now(force: int = 0, user=Depends(require_operator)): import scheduler as scheduler_mod jid = await scheduler_mod.enqueue(db, JOB, {"force": bool(force)}) return {"queued": jid, "job": JOB} @cron_router.post("/selfaudit") async def cron_selfaudit(request: Request): # Cron endpoints must ack 2xx immediately; enqueue/background the actual work. secret = os.environ.get("WEBHOOK_CRON_SECRET", "") auth = request.headers.get("authorization", "") if not secret or not auth.startswith("Bearer ") or not hmac.compare_digest(auth[7:], secret): raise HTTPException(401, "unauthorized") try: body = await request.json() if await request.body() else {} except Exception: raise HTTPException(400, "invalid body") run_id = request.headers.get("x-webhook-id") or (body or {}).get("run_id") or datetime.now(timezone.utc).isoformat() if await db.cron_deliveries.find_one({"run_id": run_id, "job": JOB}): return {"accepted": True, "duplicate": True} await db.cron_deliveries.insert_one({"run_id": run_id, "job": JOB, "at": datetime.now(timezone.utc).isoformat()}) import scheduler as scheduler_mod jid = await scheduler_mod.enqueue(db, JOB, {}) return {"accepted": True, "queued": jid}