Audit log
Every step, every agent, on the record.
One append-only trail across all runs: what each agent thought, which tools it called, what came back, and every approval decision, with the risk tier stamped on each entry.
✓
Chain verified. 4 entries hash-chained and intact, 109 legacy (pre-chain). Each entry seals the one before it, so any edit, deletion, or reorder is detectable.
08-30 04:55:19
final
Here's a overview of what I can help you with:
## ๐ Knowledge & Policy
- **Search the internal knowledge base** for company policies, guidance, and documentatirun-ebc3e428
08-30 04:55:19
thought
Here's a overview of what I can help you with:
## ๐ Knowledge & Policy
- **Search the internal knowledge base** for company policies, guidance, and documentatirun-ebc3e428
08-30 04:55:19
model call
{"cost": 0.031443, "input_tokens": 8516, "latency_ms": 12621, "model": "claude-sonnet-4-6", "output_tokens": 393}run-ebc3e428
08-30 02:54:36
final
Good. Now I have the full, exact code in front of me. Let me trace through the logic precisely.
Here is the honest finding on issue #11:
---
## ๐ Issue #11 โrun-e941689e
08-30 02:54:36
thought
Good. Now I have the full, exact code in front of me. Let me trace through the logic precisely.
Here is the honest finding on issue #11:
---
## ๐ Issue #11 โrun-e941689e
08-30 02:54:36
model call
{"cost": 0.125259, "input_tokens": 37478, "latency_ms": 17656, "model": "claude-sonnet-4-6", "output_tokens": 855}run-e941689e
08-30 02:54:18
tool result
LOW
read_source → "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re, time\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nMAX_TOKENS = int(os.environ.get(\"WARDEN_MAX_TOKENS\", \"4096\"))\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\n# Estimated USD price per 1M tokens (input, output). Approximate list prices, editable;\n# used only to estimate cost for the observability view. Matched by substring of model id.\nPRICES = {\"opus\": (15.0, 75.0), \"sonnet\": (3.0, 15.0), \"haiku\": (0.80, 4.0)}\ndef _cost(model, inp, out):\n rate = (3.0, 15.0)\n m = (model or \"\").lower()\n for k, v in PRICES.items():\n if k in m:\n rate = v; break\n return round(inp / 1e6 * rate[0] + out / 1e6 * rate[1], 6)\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _repair(messages):\n \"\"\"Guarantee the API invariant: every assistant tool_use is answered by a tool_result\n in the very next message. If a tool crashed, a response was truncated, or a follow-up\n landed on a dangling turn, backfill synthetic \u0027interrupted\u0027 results so the request is\n valid. Turns a hard 400 into a graceful continuation the model can reason about.\"\"\"\n i = 0\n while i \u003c len(messages):\n m = messages[i]\n if m.get(\"role\") == \"assistant\" and isinstance(m.get(\"content\"), list):\n ids = [b[\"id\"] for b in m[\"content\"]\n if isinstance(b, dict) and b.get(\"type\") == \"tool_use\" and b.get(\"id\")]\n if ids:\n nxt = messages[i + 1] if i + 1 \u003c len(messages) else None\n answered = set()\n if nxt and nxt.get(\"role\") == \"user\" and isinstance(nxt.get(\"content\"), list):\n answered = {b.get(\"tool_use_id\") for b in nxt[\"content\"]\n if isinstance(b, dict) and b.get(\"type\") == \"tool_result\"}\n missing = [t for t in ids if t not in answered]\n if missing:\n fills = [{\"type\": \"tool_result\", \"tool_use_id\": t,\n \"content\": json.dumps({\"error\": \"interrupted\",\n \"note\": \"This tool did not complete. Do not assume it ran.\"})}\n for t in missing]\n if nxt and nxt.get(\"role\") == \"user\" and isinstance(nxt.get(\"content\"), list):\n nxt[\"content\"] = fills + nxt[\"content\"]\n else:\n messages.insert(i + 1, {\"role\": \"user\", \"content\": fills})\n i += 1\n\n_TRANSIENT = (\"rate\", \"overloaded\", \"timeout\", \"timedout\", \"internal\", \"unavailable\", \"connection\")\ndef _is_transient(ex):\n code = getattr(ex, \"status_code\", None)\n if code in (408, 409, 429, 500, 502, 503, 504, 529):\n return True\n return any(w in (type(ex).__name__ + \" \" + str(ex)).lower() for w in _TRANSIENT)\n\ndef _friendly_error(ex):\n code = getattr(ex, \"status_code\", None)\n if code == 401 or \"authentication\" in str(ex).lower():\n return \"The model rejected the API key. Check ANTHROPIC_API_KEY.\"\n if code == 429 or \"rate\" in str(ex).lower():\n return \"The model is rate-limited right now. Try again in a moment.\"\n if code and 500 \u003c= code \u003c 600:\n return \"The model service had a temporary error. Try again in a moment.\"\n if code == 400:\n return \"The model rejected the request. This run hit a malformed-request error; the transcript has been repaired, please retry.\"\n return \"The run hit an error talking to the model: \" + str(ex)[:200]\n\ndef _call_model(system, messages, tools):\n t0 = time.time()\n if SANDBOX:\n r = _sandbox_model(messages, tools)\n r[\"usage\"] = {\"input_tokens\": 0, \"output_tokens\": 0}\n r[\"model\"] = \"sandbox\"; r[\"latency_ms\"] = int((time.time() - t0) * 1000)\n return r\n import anthropic\n client = anthropic.Anthropic()\n last = None\n for attempt in range(3):\n try:\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=MAX_TOKENS,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content],\n \"usage\": {\"input_tokens\": resp.usage.input_tokens, \"output_tokens\": resp.usage.output_tokens},\n \"model\": MODEL_DEFAULT, \"latency_ms\": int((time.time() - t0) * 1000)}\n except Exception as ex:\n last = ex\n if _is_transient(ex) and attempt \u003c 2:\n time.sleep(1.5 * (attempt + 1)); continue\n raise last\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n # Intent comes from the user\u0027s actual request, not the accumulating transcript,\n # so a completed refund doesn\u0027t spuriously trigger a file-write gate.\n user_text = \"\"\n for m in messages:\n if m.get(\"role\") == \"user\" and isinstance(m.get(\"content\"), str):\n user_text = m[\"content\"].lower(); break\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", user_text)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in user_text for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in user_text for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\nimport threading\n_run_locks = {} # run_id -\u003e Lock, serializes advance() per run\n_rerun = set() # run_ids asked to advance again while already advancing\n_guard = threading.Lock()\n\ndef advance(run_id):\n \"\"\"Serialize advancing a single run. Concurrent triggers (e.g. several approvals\n decided at once) must not run the loop on the same transcript in parallel, or the\n model gets an assistant turn whose tool_use blocks aren\u0027t all answered yet (API 400).\n Only one thread advances a run at a time; triggers that arrive mid-advance cause\n exactly one more pass afterward, so the latest decisions are always picked up.\"\"\"\n with _guard:\n lock = _run_locks.setdefault(run_id, threading.Lock())\n if lock.locked():\n _rerun.add(run_id) # someone is already advancing; ask them to loop\n return store.get_run(run_id)\n with lock:\n while True:\n result = _advance_once(run_id)\n with _guard:\n if run_id in _rerun:\n _rerun.discard(run_id)\n continue # a decision landed during the pass; go again\n break\n return result\n\ndef _advance_once(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n _repair(messages) # never send an unanswered tool_use to the API\n try:\n resp = _call_model(system, messages, tools)\n except Exception as ex:\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\": _friendly_error(ex)})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n u = resp.get(\"usage\", {})\n store.audit(run_id, agent[\"id\"], \"model_call\",\n detail={\"model\": resp.get(\"model\"), \"input_tokens\": u.get(\"input_tokens\", 0),\n \"output_tokens\": u.get(\"output_tokens\", 0), \"latency_ms\": resp.get(\"latency_ms\", 0),\n \"cost\": _cost(resp.get(\"model\"), u.get(\"input_tokens\", 0), u.get(\"output_tokens\", 0))})\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \""run-e941689e
08-30 02:54:18
thought
Sure! Let me re-read the relevant section of `agent_runtime.py` carefully before forming any conclusions.run-e941689e
08-30 02:54:18
model call
{"cost": 0.102477, "input_tokens": 33724, "latency_ms": 6007, "model": "claude-sonnet-4-6", "output_tokens": 87}run-e941689e
08-30 02:54:12
Can you test #11?run-e941689e
08-30 01:03:00
final
Excellent โ I have the full picture. Here is the complete, concrete audit report.
---
# ๐ก๏ธ Warden Self-Audit Report
**Static check:** โ
12 files, 0 syntax errun-e941689e
08-30 01:03:00
thought
Excellent โ I have the full picture. Here is the complete, concrete audit report.
---
# ๐ก๏ธ Warden Self-Audit Report
**Static check:** โ
12 files, 0 syntax errun-e941689e
08-30 01:03:00
model call
{"cost": 0.132465, "input_tokens": 31105, "latency_ms": 51669, "model": "claude-sonnet-4-6", "output_tokens": 2610}run-e941689e
08-30 01:02:09
tool result
LOW
read_file → "Warden workspace. Files here are readable and writable by agents, under governance.\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nEncrypts secrets (connection tokens) at rest so they never sit in the database as\nplaintext. The key comes from WARDEN_SECRET_KEY if set (kept out of the data dir, the\nstronger option); otherwise a key is generated once and stored on the persistent disk\nnext to the data, so encryption works with zero configuration. Any string works as\nWARDEN_SECRET_KEY, it is hashed into a valid key.\n\"\"\"\nimport os\nimport base64\nimport hashlib\n\ntry:\n from cryptography.fernet import Fernet\n _OK = True\nexcept Exception:\n _OK = False\n\nimport paths\nDATA_ROOT = paths.DATA_ROOT\nKEY_FILE = os.path.join(DATA_ROOT, \".warden_key\")\n\ndef _key():\n env = os.environ.get(\"WARDEN_SECRET_KEY\")\n if env:\n return base64.urlsafe_b64encode(hashlib.sha256(env.encode()).digest())\n if os.path.exists(KEY_FILE):\n return open(KEY_FILE, \"rb\").read().strip()\n k = Fernet.generate_key()\n os.makedirs(DATA_ROOT, exist_ok=True)\n open(KEY_FILE, \"wb\").write(k)\n try:\n os.chmod(KEY_FILE, 0o600)\n except Exception:\n pass\n return k\n\n_F = None\ndef _fernet():\n global _F\n if _F is None:\n _F = Fernet(_key())\n return _F\n\ndef encrypt(s):\n \"\"\"Return an \u0027enc:\u0027-prefixed ciphertext, or the original if encryption is unavailable.\"\"\"\n if not s or not _OK:\n return s\n try:\n return \"enc:\" + _fernet().encrypt(s.encode()).decode()\n except Exception:\n return s\n\ndef decrypt(s):\n \"\"\"Reverse encrypt(); passes through anything not marked \u0027enc:\u0027 (e.g. legacy plaintext).\"\"\"\n if not s or not _OK or not isinstance(s, str) or not s.startswith(\"enc:\"):\n return s\n try:\n return _fernet().decrypt(s[4:].encode()).decode()\n except Exception:\n return s\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nTelemetry: turns a run\u0027s audit trail into an OpenTelemetry trace and exports it.\n\nA run maps onto a distributed trace: the run is the root span, and each step, model\ncall, tool call, approval gate, is a child span carrying the dimensions the runtime\nalready records (tokens, cost, latency, risk, outcome, decision). The audit log stays\nthe source of truth; this is a derived, exportable view of it.\n\nExport is over OTLP/HTTP (JSON) to whatever the customer already runs (Datadog, Grafana\nTempo, Honeycomb, ...), configured via the standard env vars:\n OTEL_EXPORTER_OTLP_ENDPOINT e.g. https://otlp.example.com\n OTEL_EXPORTER_OTLP_HEADERS e.g. x-api-key=abc,x-dataset=warden (optional auth)\nNo extra dependencies: the OTLP payload is built and POSTed with the standard library.\n\"\"\"\nimport os, json, hashlib, urllib.request\nfrom datetime import datetime, timezone\nimport store\n\nSERVICE = \"warden\"\n\ndef _ns(ts_iso):\n if not ts_iso:\n return 0\n try:\n s = ts_iso.replace(\"Z\", \"\")\n dt = datetime.fromisoformat(s)\n if dt.tzinfo is None:\n dt = dt.replace(tzinfo=timezone.utc)\n return int(dt.timestamp() * 1_000_000_000)\n except Exception:\n return 0\n\ndef _id(seed, nbytes):\n return hashlib.sha256(seed.encode()).hexdigest()[: nbytes * 2]\n\ndef _attrs(d):\n out = []\n for k, v in d.items():\n if v is None:\n continue\n if isinstance(v, bool):\n val = {\"boolValue\": v}\n elif isinstance(v, int):\n val = {\"intValue\": str(v)}\n elif isinstance(v, float):\n val = {\"doubleValue\": v}\n else:\n val = {\"stringValue\": str(v)}\n out.append({\"key\": k, \"value\": val})\n return out\n\ndef build_spans(run_id):\n \"\"\"Returns (spans, meta) where spans is a list of plain dicts we can both render as a\n waterfall and serialize to OTLP. meta carries trace-level totals.\"\"\"\n run = store.get_run(run_id)\n if not run:\n return [], {}\n events = store.audit_for_run(run_id)\n trace_id = _id(\"trace-\" + run_id, 16)\n root_id = _id(\"root-\" + run_id, 8)\n\n stamps = [_ns(e[\"ts\"]) for e in events if _ns(e[\"ts\"])]\n t0 = min(stamps) if stamps else 0\n t1 = max(stamps) if stamps else t0 + 1\n\n spans = [{\n \"trace_id\": trace_id, \"span_id\": root_id, \"parent\": None,\n \"name\": \"run: \" + (run[\"input\"] or \"\")[:60], \"start\": t0, \"end\": max(t1, t0 + 1),\n \"attrs\": {\"warden.run_id\": run_id, \"warden.agent_id\": run[\"agent_id\"],\n \"warden.status\": run[\"status\"]},\n \"error\": run[\"status\"] == \"error\", \"row\": \"run\",\n }]\n\n total_cost = 0.0; total_tokens = 0\n for i, e in enumerate(events):\n d = e.get(\"detail\") or {}\n st = _ns(e[\"ts\"]); dur_ms = d.get(\"latency_ms\") or 0\n en = st + int(dur_ms * 1_000_000) if dur_ms else st + 1\n kind = e[\"kind\"]; tool = (e[\"skill\"] or \"\").split(\"__\")[-1]\n if kind == \"model_call\":\n total_cost += d.get(\"cost\", 0) or 0\n total_tokens += (d.get(\"input_tokens\", 0) or 0) + (d.get(\"output_tokens\", 0) or 0)\n name = \"model: \" + (d.get(\"model\") or \"\")\n attrs = {\"gen_ai.system\": \"anthropic\", \"gen_ai.request.model\": d.get(\"model\"),\n \"gen_ai.usage.input_tokens\": d.get(\"input_tokens\"),\n \"gen_ai.usage.output_tokens\": d.get(\"output_tokens\"),\n \"warden.cost_usd\": d.get(\"cost\"), \"warden.latency_ms\": dur_ms}\n row = \"model\"\n elif kind in (\"tool_result\", \"tool_result_gated\"):\n name = \"tool: \" + tool\n attrs = {\"warden.tool\": tool, \"warden.risk\": e.get(\"risk\"),\n \"warden.gated\": kind == \"tool_result_gated\",\n \"warden.outcome\": d.get(\"outcome\"), \"warden.latency_ms\": dur_ms}\n row = \"tool\"\n elif kind == \"approval_request\":\n name = \"gate: \" + tool\n attrs = {\"warden.tool\": tool, \"warden.risk\": e.get(\"risk\"), \"warden.decision\": \"requested\"}\n row = \"gate\"\n elif kind == \"denied\":\n name = \"denied: \" + tool\n attrs = {\"warden.tool\": tool, \"warden.risk\": e.get(\"risk\"), \"warden.decision\": \"denied\"}\n row = \"gate\"\n elif kind == \"final\":\n name = \"final response\"; attrs = {}; row = \"final\"\n else:\n continue\n spans.append({\n \"trace_id\": trace_id, \"span_id\": _id(f\"{run_id}-{i}\", 8), \"parent\": root_id,\n \"name\": name, \"start\": st, \"end\": max(en, st + 1), \"attrs\": attrs,\n \"error\": d.get(\"outcome\") == \"error\" or kind == \"denied\", \"row\": row,\n })\n\n meta = {\"trace_id\": trace_id, \"run_id\": run_id, \"status\": run[\"status\"],\n \"spans\": len(spans), \"t0\": t0, \"total_ns\": max(t1 - t0, 1),\n \"cost\": round(total_cost, 6), \"tokens\": total_tokens}\n return spans, meta\n\ndef to_otlp(run_id):\n spans, meta = build_spans(run_id)\n otlp_spans = [{\n \"traceId\": s[\"trace_id\"], \"spanId\": s[\"span_id\"],\n **({\"parentSpanId\": s[\"parent\"]} if s[\"parent\"] else {}),\n \"name\": s[\"name\"], \"kind\": 1,\n \"startTimeUnixNano\": str(s[\"start\"]), \"endTimeUnixNano\": str(s[\"end\"]),\n \"attributes\": _attrs(s[\"attrs\"]),\n \"status\": {\"code\": 2 if s[\"error\"] else 1},\n } for s in spans]\n return {\"resourceSpans\": [{\n \"resource\": {\"attributes\": _attrs({\"service.name\": SERVICE})},\n \"scopeSpans\": [{\"scope\": {\"name\": \"warden.runtime\"}, \"spans\": otlp_spans}]}]}\n\ndef _headers():\n h = {\"Content-Type\": \"application/json\"}\n raw = os.environ.get(\"OTEL_EXPORTER_OTLP_HEADERS\", \"\")\n for pair in raw.split(\",\"):\n if \"=\" in pair:\n k, v = pair.split(\"=\", 1)\n h[k.strip()] = v.strip()\n return h\n\ndef export(run_id):\n endpoint = os.environ.get(\"OTEL_EXPORTER_OTLP_ENDPOINT\")\n payload = to_otlp(run_id)\n n = len(payload[\"resourceSpans\"][0][\"scopeSpans\"][0][\"spans\"])\n if not endpoint:\n return {\"ok\": False, \"configured\": False,\n \"reason\": \"No OTEL_EXPORTER_OTLP_ENDPOINT set. Set it (and optional \"\n \"OTEL_EXPORTER_OTLP_HEADERS) to export to your collector.\",\n \"spans\": n}\n url = endpoint.rstrip(\"/\") + \"/v1/traces\"\n try:\n req = urllib.request.Request(url, data=json.dumps(payload).encode(),\n headers=_headers(), method=\"POST\")\n with urllib.request.urlopen(req, timeout=10) as r:\n return {\"ok\": True, \"configured\": True, \"status\": r.status, \"endpoint\": url, \"spans\": n}\n except Exception as ex:\n return {\"ok\": False, \"configured\": True, \"reason\": str(ex)[:200], \"endpoint\": url, \"spans\": n}\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nPersistence for Warden: agents, runs, the audit log, and the approvals queue.\nSQLite on disk. On a platform with an ephemeral filesystem (Render) this resets on\nredeploy, which is fine for a demo; the point is that within a session every agent\naction and every approval is durably recorded and queryable.\n\"\"\"\nimport os\nimport json\nimport sqlite3\nimport datetime\nimport uuid\n\nimport paths\nDATA_ROOT = paths.DATA_ROOT\nDB = os.path.join(DATA_ROOT, \"warden.db\")\n\ndef _conn():\n c = sqlite3.connect(DB, timeout=10)\n c.row_factory = sqlite3.Row\n try:\n c.execute(\"PRAGMA journal_mode=WAL\")\n c.execute(\"PRAGMA busy_timeout=8000\")\n except Exception:\n pass\n return c\n\ndef now():\n return datetime.datetime.now(datetime.timezone.utc).isoformat()\n\ndef _id(prefix):\n return f\"{prefix}-{uuid.uuid4().hex[:8]}\"\n\ndef init():\n c = _conn()\n c.executescript(\"\"\"\n CREATE TABLE IF NOT EXISTS agents(\n id TEXT PRIMARY KEY, name TEXT, instructions TEXT, model TEXT,\n skills TEXT, created_at TEXT);\n CREATE TABLE IF NOT EXISTS runs(\n id TEXT PRIMARY KEY, agent_id TEXT, input TEXT, status TEXT,\n transcript TEXT, created_at TEXT, updated_at TEXT);\n CREATE TABLE IF NOT EXISTS audit(\n id TEXT PRIMARY KEY, run_id TEXT, agent_id TEXT, ts TEXT,\n kind TEXT, skill TEXT, risk TEXT, detail TEXT);\n CREATE TABLE IF NOT EXISTS approvals(\n id TEXT PRIMARY KEY, run_id TEXT, agent_id TEXT, skill TEXT, risk TEXT,\n arguments TEXT, status TEXT, created_at TEXT, decided_at TEXT, decided_by TEXT);\n CREATE TABLE IF NOT EXISTS connections(\n id TEXT PRIMARY KEY, transport TEXT, command TEXT, url TEXT, token TEXT,\n enabled INTEGER, created_at TEXT);\n CREATE TABLE IF NOT EXISTS tool_overrides(\n model_key TEXT PRIMARY KEY, risk TEXT);\n \"\"\")\n c.commit(); c.close()\n\n# ---- agents ----\ndef create_agent(name, instructions, model, skills):\n c = _conn(); aid = _id(\"ag\")\n c.execute(\"INSERT INTO agents VALUES(?,?,?,?,?,?)\",\n (aid, name, instructions, model, json.dumps(skills), now()))\n c.commit(); c.close(); return aid\n\ndef get_agent(aid):\n c = _conn(); r = c.execute(\"SELECT * FROM agents WHERE id=?\", (aid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"skills\"] = json.loads(d[\"skills\"] or \"[]\"); return d\n\ndef list_agents():\n c = _conn(); rows = c.execute(\"SELECT * FROM agents ORDER BY created_at DESC\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"skills\"] = json.loads(d[\"skills\"] or \"[]\"); out.append(d)\n return out\n\n# ---- runs ----\ndef create_run(agent_id, user_input):\n c = _conn(); rid = _id(\"run\")\n c.execute(\"INSERT INTO runs VALUES(?,?,?,?,?,?,?)\",\n (rid, agent_id, user_input, \"running\", json.dumps([]), now(), now()))\n c.commit(); c.close(); return rid\n\ndef update_run(rid, status=None, transcript=None):\n c = _conn()\n if status is not None:\n c.execute(\"UPDATE runs SET status=?, updated_at=? WHERE id=?\", (status, now(), rid))\n if transcript is not None:\n c.execute(\"UPDATE runs SET transcript=?, updated_at=? WHERE id=?\",\n (json.dumps(transcript), now(), rid))\n c.commit(); c.close()\n\ndef get_run(rid):\n c = _conn(); r = c.execute(\"SELECT * FROM runs WHERE id=?\", (rid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"transcript\"] = json.loads(d[\"transcript\"] or \"[]\"); return d\n\ndef list_runs(limit=50):\n c = _conn(); rows = c.execute(\"SELECT * FROM runs ORDER BY created_at DESC LIMIT ?\", (limit,)).fetchall(); c.close()\n return [dict(r) for r in rows]\n\n# ---- audit ----\ndef audit(run_id, agent_id, kind, skill=None, risk=None, detail=None):\n c = _conn()\n c.execute(\"INSERT INTO audit VALUES(?,?,?,?,?,?,?,?)\",\n (_id(\"ev\"), run_id, agent_id, now(), kind, skill, risk,\n json.dumps(detail) if detail is not None else None))\n c.commit(); c.close()\n\ndef audit_for_run(run_id):\n c = _conn(); rows = c.execute(\"SELECT * FROM audit WHERE run_id=? ORDER BY ts\", (run_id,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"detail\"] = json.loads(d[\"detail\"]) if d[\"detail\"] else None; out.append(d)\n return out\n\ndef audit_all(limit=200):\n c = _conn(); rows = c.execute(\"SELECT * FROM audit ORDER BY ts DESC LIMIT ?\", (limit,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"detail\"] = json.loads(d[\"detail\"]) if d[\"detail\"] else None; out.append(d)\n return out\n\n# ---- approvals ----\ndef create_approval(run_id, agent_id, skill, risk, arguments):\n c = _conn(); apid = _id(\"ap\")\n c.execute(\"INSERT INTO approvals VALUES(?,?,?,?,?,?,?,?,?,?)\",\n (apid, run_id, agent_id, skill, risk, json.dumps(arguments),\n \"pending\", now(), None, None))\n c.commit(); c.close(); return apid\n\ndef get_approval(apid):\n c = _conn(); r = c.execute(\"SELECT * FROM approvals WHERE id=?\", (apid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); return d\n\ndef decide_approval(apid, status, by=\"operator\"):\n c = _conn()\n c.execute(\"UPDATE approvals SET status=?, decided_at=?, decided_by=? WHERE id=?\",\n (status, now(), by, apid))\n c.commit(); c.close()\n\ndef pending_approvals():\n c = _conn(); rows = c.execute(\"SELECT * FROM approvals WHERE status=\u0027pending\u0027 ORDER BY created_at\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); out.append(d)\n return out\n\ndef approvals_for_run(run_id):\n c = _conn(); rows = c.execute(\"SELECT * FROM approvals WHERE run_id=? ORDER BY created_at\", (run_id,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); out.append(d)\n return out\n\n# ---- connections (enabled external MCP servers) ----\ndef enable_connection(cid, transport, command=None, url=None, token=None):\n import vault\n c = _conn()\n c.execute(\"INSERT OR REPLACE INTO connections VALUES(?,?,?,?,?,?,?)\",\n (cid, transport, command, url, vault.encrypt(token), 1, now()))\n c.commit(); c.close()\n\ndef disable_connection(cid):\n c = _conn(); c.execute(\"DELETE FROM connections WHERE id=?\", (cid,)); c.commit(); c.close()\n\ndef enabled_connections():\n import vault\n c = _conn(); rows = c.execute(\"SELECT * FROM connections WHERE enabled=1\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r)\n out.append({\"id\": d[\"id\"], \"transport\": d[\"transport\"], \"command\": d[\"command\"],\n \"url\": d[\"url\"], \"token\": vault.decrypt(d[\"token\"])})\n return out\n\ndef is_enabled(cid):\n c = _conn(); r = c.execute(\"SELECT 1 FROM connections WHERE id=? AND enabled=1\", (cid,)).fetchone(); c.close()\n return bool(r)\n\n# ---- tool risk overrides ----\ndef set_override(model_key, risk):\n c = _conn(); c.execute(\"INSERT OR REPLACE INTO tool_overrides VALUES(?,?)\", (model_key, risk)); c.commit(); c.close()\n\ndef get_override(model_key):\n c = _conn(); r = c.execute(\"SELECT risk FROM tool_overrides WHERE model_key=?\", (model_key,)).fetchone(); c.close()\n return r[\"risk\"] if r else None\n\ndef all_overrides():\n c = _conn(); rows = c.execute(\"SELECT * FROM tool_overrides\").fetchall(); c.close()\n return {r[\"model_key\"]: r[\"risk\"] for r in rows}\n\ndef approval_counts():\n c = _conn(); rows = c.execute(\"SELECT status, COUNT(*) n FROM approvals GROUP BY status\").fetchall(); c.close()\n return {r[\"status\"]: r[\"n\"] for r in rows}\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nResolves where Warden stores its data. Prefers WARDEN_DATA_DIR (a mounted persistent\ndisk in production). If that path can\u0027t be created or written, it falls back to a\nwritable local directory instead of crashing the app, and records that it fell back so\nthe condition is visible on /healthz. A misconfigured disk should degrade to ephemeral,\nnever take the service down.\n\"\"\"\nimport os\n\n_APPDIR = os.path.dirname(os.path.abspath(__file__))\nREQUESTED = os.environ.get(\"WARDEN_DATA_DIR\") or _APPDIR\n\ndef _writable(path):\n try:\n os.makedirs(path, exist_ok=True)\n t = os.path.join(path, \".wtest\")\n with open(t, \"w\") as f:\n f.write(\"ok\")\n os.remove(t)\n return True\n except Exception:\n return False\n\ndef _resolve():\n if _writable(REQUESTED):\n return REQUESTED, False\n fallback = os.path.join(_APPDIR, \"_localdata\")\n if _writable(fallback):\n return fallback, True\n return _APPDIR, True\n\nDATA_ROOT, FALLBACK = _resolve()\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nWarden MCP server: a real Model Context Protocol server exposing a small set of\nenterprise-flavored tools. The Warden runtime connects to this as an MCP client,\ndiscovers these tools over the protocol, and invokes them.\n\nTools deliberately span read and write so the governance layer has something to\ngovern: reads are low risk and auto-execute; writes change state and are the ones\nthe studio gates behind human approval.\n\nRun standalone for a protocol smoke test: python mcp_server.py\n(but normally it is spawned over stdio by the runtime\u0027s MCP client)\n\"\"\"\nimport json\nimport os\nimport datetime\nfrom mcp.server.fastmcp import FastMCP\n\nimport paths\nDATA_DIR = os.path.join(paths.DATA_ROOT, \"data\")\nos.makedirs(DATA_DIR, exist_ok=True)\n\nCUSTOMERS = {\n \"AC-1001\": {\"account_id\": \"AC-1001\", \"name\": \"Rivera Logistics\", \"plan\": \"Enterprise\",\n \"mrr\": 4200, \"status\": \"active\", \"last_charge\": 4200,\n \"notes\": \"Charged twice on 2026-08-03 due to a billing retry bug.\"},\n \"AC-1002\": {\"account_id\": \"AC-1002\", \"name\": \"Halcyon Health\", \"plan\": \"Premium\",\n \"mrr\": 1800, \"status\": \"active\", \"last_charge\": 1800, \"notes\": \"\"},\n \"AC-1003\": {\"account_id\": \"AC-1003\", \"name\": \"Meridian Foods\", \"plan\": \"Standard\",\n \"mrr\": 600, \"status\": \"past_due\", \"last_charge\": 0,\n \"notes\": \"Payment failed twice this month.\"},\n}\n\nKB = [\n {\"id\": \"kb-01\", \"title\": \"Refund policy\",\n \"body\": \"Duplicate charges are refunded in full once verified against the billing ledger. \"\n \"Refunds above 1000 require a human approver.\"},\n {\"id\": \"kb-02\", \"title\": \"Past-due accounts\",\n \"body\": \"Do not issue credits on past-due accounts until the balance is cleared. \"\n \"Open a billing ticket instead.\"},\n {\"id\": \"kb-03\", \"title\": \"Escalation\",\n \"body\": \"Anything touching money movement is a governed action and must be logged.\"},\n]\n\ndef _ledger_path(name):\n return os.path.join(DATA_DIR, name)\n\ndef _append(name, row):\n path = _ledger_path(name)\n rows = []\n if os.path.exists(path):\n rows = json.load(open(path))\n row[\"at\"] = datetime.datetime.now(datetime.timezone.utc).isoformat()\n rows.append(row)\n json.dump(rows, open(path, \"w\"), indent=2)\n return row\n\nmcp = FastMCP(\"warden-enterprise-tools\")\n\n@mcp.tool()\ndef lookup_customer(account_id: str) -\u003e str:\n \"\"\"Look up an enterprise customer account by id (e.g. AC-1001). Read only.\"\"\"\n rec = CUSTOMERS.get(account_id.strip().upper())\n if not rec:\n return json.dumps({\"error\": f\"no account {account_id}\"})\n return json.dumps(rec)\n\n@mcp.tool()\ndef search_knowledge(query: str) -\u003e str:\n \"\"\"Search the internal knowledge base for policy and guidance. Read only.\"\"\"\n q = query.lower()\n hits = [a for a in KB if q in a[\"title\"].lower() or q in a[\"body\"].lower()]\n if not hits:\n hits = KB # fall back to returning all short KB so the agent has context\n return json.dumps(hits)\n\n@mcp.tool()\ndef create_ticket(subject: str, body: str) -\u003e str:\n \"\"\"Open an internal billing/support ticket. Write action (changes state).\"\"\"\n row = _append(\"tickets.json\", {\"id\": f\"TK-{abs(hash(subject)) % 9000 + 1000}\",\n \"subject\": subject, \"body\": body})\n return json.dumps({\"created\": True, \"ticket\": row})\n\n@mcp.tool()\ndef issue_refund(account_id: str, amount: float, reason: str = \"\") -\u003e str:\n \"\"\"Issue a monetary refund to a customer account. High-impact write action:\n moves money, so the studio gates this behind human approval before it runs.\"\"\"\n acct = account_id.strip().upper()\n if acct not in CUSTOMERS:\n return json.dumps({\"error\": f\"no account {account_id}\"})\n row = _append(\"refunds.json\", {\"id\": f\"RF-{abs(hash(acct + str(amount))) % 9000 + 1000}\",\n \"account_id\": acct, \"amount\": float(amount), \"reason\": reason})\n return json.dumps({\"refunded\": True, \"refund\": row})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nWarden\u0027s built-in Filesystem MCP server. A real second MCP server (local, Python) so\nmulti-server connection and routing is demonstrable without needing Node. All access\nis confined to a sandbox workspace directory.\n\"\"\"\nimport os\nimport json\nfrom mcp.server.fastmcp import FastMCP\n\nimport paths\nWORKSPACE = os.path.join(paths.DATA_ROOT, \"data\", \"workspace\")\nos.makedirs(WORKSPACE, exist_ok=True)\n\n# seed one file so reads have something to find\n_seed = os.path.join(WORKSPACE, \"welcome.txt\")\nif not os.path.exists(_seed):\n open(_seed, \"w\").write(\"Warden workspace. Files here are readable and writable by agents, under governance.\\n\")\n\ndef _safe(path):\n p = os.path.abspath(os.path.join(WORKSPACE, path.lstrip(\"/\")))\n if not p.startswith(os.path.abspath(WORKSPACE)):\n raise ValueError(\"path escapes workspace\")\n return p\n\nmcp = FastMCP(\"warden-filesystem\")\n\n@mcp.tool()\ndef list_files(subdir: str = \"\") -\u003e str:\n \"\"\"List files in the workspace (or a subdirectory). Read only.\"\"\"\n base = _safe(subdir)\n if not os.path.exists(base):\n return json.dumps({\"error\": \"no such path\"})\n return json.dumps(sorted(os.listdir(base)))\n\n@mcp.tool()\ndef read_file(path: str) -\u003e str:\n \"\"\"Read a text file from the workspace. Read only.\"\"\"\n p = _safe(path)\n if not os.path.isfile(p):\n return json.dumps({\"error\": f\"no file {path}\"})\n return open(p, encoding=\"utf-8\", errors=\"replace\").read()[:8000]\n\n@mcp.tool()\ndef write_file(path: str, content: str) -\u003e str:\n \"\"\"Create or overwrite a text file in the workspace. Write action (changes state).\"\"\"\n p = _safe(path)\n os.makedirs(os.path.dirname(p), exist_ok=True)\n open(p, \"w\", encoding=\"utf-8\").write(content)\n return json.dumps({\"written\": True, \"path\": path, \"bytes\": len(content)})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nWarden self-audit MCP server. Exposes Warden\u0027s own source code to an agent so a\n\"Warden Engineer\" agent can inspect the running codebase, run a real static self-check,\nand propose fixes. Reads are safe and auto-run; proposing a patch is a gated write that\nlands in a review folder (it never overwrites the running source).\n\"\"\"\nimport os, json, ast, py_compile, tempfile\nfrom mcp.server.fastmcp import FastMCP\n\nAPP_DIR = os.path.dirname(os.path.abspath(__file__))\nimport paths\nPATCH_DIR = os.path.join(paths.DATA_ROOT, \"data\", \"patches\")\nos.makedirs(PATCH_DIR, exist_ok=True)\n\ndef _py_files():\n return sorted(f for f in os.listdir(APP_DIR) if f.endswith(\".py\"))\n\nmcp = FastMCP(\"warden-self-audit\")\n\n@mcp.tool()\ndef list_source() -\u003e str:\n \"\"\"List Warden\u0027s own Python source files with line counts. Read only.\"\"\"\n out = []\n for f in _py_files():\n n = sum(1 for _ in open(os.path.join(APP_DIR, f), encoding=\"utf-8\", errors=\"replace\"))\n out.append({\"file\": f, \"lines\": n})\n return json.dumps(out)\n\n@mcp.tool()\ndef read_source(filename: str) -\u003e str:\n \"\"\"Read one of Warden\u0027s own source files. Read only.\"\"\"\n if filename not in _py_files():\n return json.dumps({\"error\": f\"no source file {filename}\", \"available\": _py_files()})\n return open(os.path.join(APP_DIR, filename), encoding=\"utf-8\", errors=\"replace\").read()[:12000]\n\n@mcp.tool()\ndef run_selfcheck() -\u003e str:\n \"\"\"Statically check every Warden source file: byte-compile for syntax errors and\n AST-scan for bare excepts and TODO/FIXME markers. Runs real checks; no side effects.\"\"\"\n findings = []\n for f in _py_files():\n path = os.path.join(APP_DIR, f)\n try:\n py_compile.compile(path, doraise=True)\n except py_compile.PyCompileError as e:\n findings.append({\"file\": f, \"kind\": \"syntax_error\", \"detail\": str(e).splitlines()[-1][:160]})\n continue\n src = open(path, encoding=\"utf-8\", errors=\"replace\").read()\n try:\n tree = ast.parse(src)\n for node in ast.walk(tree):\n if isinstance(node, ast.ExceptHandler) and node.type is None:\n findings.append({\"file\": f, \"kind\": \"bare_except\", \"line\": node.lineno})\n except Exception:\n pass\n try:\n import tokenize, io\n for tok in tokenize.generate_tokens(io.StringIO(src).readline):\n if tok.type == tokenize.COMMENT and (\"TODO\" in tok.string.upper() or \"FIXME\" in tok.string.upper()):\n findings.append({\"file\": f, \"kind\": \"todo\", \"line\": tok.start[0], \"detail\": tok.string.strip()[:120]})\n except Exception:\n pass\n return json.dumps({\"files_checked\": len(_py_files()),\n \"findings\": findings, \"clean\": len(findings) == 0})\n\n@mcp.tool()\ndef propose_patch(filename: str, new_content: str, rationale: str = \"\") -\u003e str:\n \"\"\"Propose a fix for a source file. Gated write: saves the proposed version to a\n review folder for a human to inspect and apply. Never edits the running source.\"\"\"\n safe = os.path.basename(filename)\n out = os.path.join(PATCH_DIR, safe)\n open(out, \"w\", encoding=\"utf-8\").write(new_content)\n meta = os.path.join(PATCH_DIR, safe + \".rationale.txt\")\n open(meta, \"w\", encoding=\"utf-8\").write(rationale or \"(none)\")\n return json.dumps({\"proposed\": True, \"review_path\": f\"data/patches/{safe}\",\n \"bytes\": len(new_content), \"note\": \"saved for human review; running source unchanged\"})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-e941689e
08-30 01:02:09
tool result
LOW
read_source → "\"\"\"\nThe governance layer. Warden\u0027s point of view: tools can do things; governance decides\nwhich run on their own and which pause for a human. Built-in enterprise tools have a\nhand-set risk registry. Tools discovered from external MCP servers are classified\nautomatically, fail-closed: reads run, writes and anything unrecognized are gated.\nAn operator can override any tool\u0027s risk.\n\"\"\"\nSKILLS = {\n \"lookup_customer\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"search_knowledge\":{\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"create_ticket\": {\"risk\":\"MED\",\"gate\":\"auto\",\"kind\":\"write\"},\n \"issue_refund\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n \"list_files\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"read_file\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"write_file\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n \"list_source\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"read_source\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"run_selfcheck\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"propose_patch\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n}\n\nREAD_HINTS = (\"get\",\"list\",\"read\",\"search\",\"lookup\",\"fetch\",\"find\",\"query\",\"view\",\n \"describe\",\"show\",\"count\",\"status\",\"summary\",\"recent\",\"ask\",\"explore\",\n \"inspect\",\"check\")\nWRITE_HINTS = (\"create\",\"write\",\"update\",\"delete\",\"remove\",\"issue\",\"send\",\"post\",\"add\",\n \"set\",\"merge\",\"close\",\"open\",\"deploy\",\"execute\",\"run\",\"refund\",\"cancel\",\n \"approve\",\"edit\",\"upload\",\"move\",\"rename\",\"revoke\",\"grant\",\"pay\",\"charge\")\n\ndef classify(name, desc=\"\"):\n \"\"\"Heuristic risk for an external tool, judged by its leading verb.\n Fail-closed: unrecognized =\u003e HIGH.\"\"\"\n import re\n n = (name or \"\").lower().replace(\"-\", \"_\")\n lead = re.match(r\"[a-z]+\", n)\n lead = lead.group(0) if lead else n # leading alpha run (handles camelCase)\n first = n.split(\"_\")[0] # first snake_case token\n for h in READ_HINTS:\n if first == h or lead.startswith(h):\n return \"LOW\"\n for h in WRITE_HINTS:\n if first == h or lead.startswith(h):\n return \"HIGH\"\n return \"HIGH\"\n\ndef meta(model_key, tool_name, desc=\"\", override=None):\n \"\"\"Resolve effective governance for a tool. Precedence: override \u003e known registry \u003e classify.\"\"\"\n if override in (\"LOW\",\"MED\",\"HIGH\"):\n risk = override\n elif tool_name in SKILLS:\n risk = SKILLS[tool_name][\"risk\"]\n else:\n risk = classify(tool_name, desc)\n return {\"risk\": risk, \"gate\": \"approval\" if risk == \"HIGH\" else \"auto\"}\n\n# convenience wrappers used where only a bare name is available (built-ins)\ndef skill_meta(name):\n m = SKILLS.get(name)\n if m: return m\n r = classify(name)\n return {\"risk\": r, \"gate\": \"approval\" if r==\"HIGH\" else \"auto\", \"kind\":\"?\"}\ndef requires_approval(name):\n return skill_meta(name)[\"gate\"] == \"approval\"\ndef risk_of(name):\n return skill_meta(name)[\"risk\"]\n"run-e941689e
08-30 01:02:08
tool result
LOW
read_source → "\"\"\"\nMulti-server MCP connection manager.\n\nOwns one dedicated asyncio loop thread. Tool discovery is done once at connect time and\ncached, so listing tools never re-hits a server. For execution: stdio/builtin servers\nkeep a persistent session (subprocess spawn is expensive); HTTP servers open a fresh\nshort-lived session per call, entirely within one coroutine, because the streamable-HTTP\ntransport binds its cancel scope to the creating task and cannot be reused across tasks.\n\"\"\"\nimport os, sys, asyncio, threading, warnings, shlex\nfrom contextlib import asynccontextmanager\nfrom mcp import ClientSession, StdioServerParameters\nfrom mcp.client.stdio import stdio_client\ntry:\n from mcp.client.streamable_http import streamablehttp_client\n _HTTP_OK = True\nexcept Exception:\n _HTTP_OK = False\n\nimport catalog as catalog_mod\n\nHERE = os.path.dirname(os.path.abspath(__file__))\nBUILTINS = {\n \"builtin_enterprise\": [sys.executable, os.path.join(HERE, \"mcp_server.py\")],\n \"builtin_files\": [sys.executable, os.path.join(HERE, \"mcp_fs_server.py\")],\n \"builtin_code\": [sys.executable, os.path.join(HERE, \"mcp_code_server.py\")],\n}\n\nclass _LoopThread:\n def __init__(self):\n self.loop = asyncio.new_event_loop()\n with warnings.catch_warnings():\n warnings.simplefilter(\"ignore\")\n try:\n w = asyncio.ThreadedChildWatcher(); w.attach_loop(self.loop)\n asyncio.set_child_watcher(w)\n except Exception:\n pass\n threading.Thread(target=self._run, daemon=True).start()\n def _run(self):\n asyncio.set_event_loop(self.loop); self.loop.run_forever()\n def run(self, coro, timeout=None):\n return asyncio.run_coroutine_threadsafe(coro, self.loop).result(timeout)\n\ndef _http_params(sid, spec):\n cat = catalog_mod.BY_ID.get(sid, {})\n url = spec.get(\"url\") or cat.get(\"run\")\n headers = {}\n tok = spec.get(\"token\")\n if not tok and cat.get(\"env\"): # fall back to an environment variable\n tok = os.environ.get(cat[\"env\"])\n if tok:\n headers[\"Authorization\"] = tok if tok.lower().startswith(\"bearer\") else f\"Bearer {tok}\"\n return url, (headers or None)\n\ndef _stdio_params(sid, spec, transport):\n if transport == \"builtin\":\n cmd = BUILTINS[sid]\n else:\n run = spec.get(\"command\") or catalog_mod.BY_ID.get(sid, {}).get(\"run\", \"\")\n cmd = shlex.split(run)\n if not cmd:\n raise RuntimeError(\"no command configured\")\n return StdioServerParameters(command=cmd[0], args=cmd[1:], env=os.environ.copy())\n\n@asynccontextmanager\nasync def _http_session(sid, spec):\n url, headers = _http_params(sid, spec)\n async with streamablehttp_client(url, headers=headers) as (read, write, _):\n async with ClientSession(read, write) as session:\n await session.initialize()\n yield session\n\nclass _Manager:\n def __init__(self):\n self._lt = _LoopThread()\n self._lock = threading.Lock()\n self._sessions = {} # sid -\u003e persistent session (stdio/builtin only)\n self._keep = {} # sid -\u003e [context managers] to keep alive\n self._http = {} # sid -\u003e spec (http servers, fresh session per call)\n self._toolcache = {} # sid -\u003e [ {name,description,input_schema} ] (from connect)\n self._status = {} # sid -\u003e status dict\n self._toolmap = {} # model_name -\u003e (sid, tool)\n self._started = False\n\n def ensure_started(self, enabled_specs=None):\n with self._lock:\n if not self._started:\n for sid in BUILTINS:\n self._connect(sid, {\"id\": sid, \"transport\": \"builtin\"})\n self._started = True\n for spec in (enabled_specs or []):\n if spec[\"id\"] not in self._status:\n self._connect(spec[\"id\"], spec)\n self._rebuild_toolmap()\n\n def connect_spec(self, spec):\n with self._lock:\n self._connect(spec[\"id\"], spec); self._rebuild_toolmap()\n return self._status.get(spec[\"id\"])\n\n def disconnect(self, sid):\n with self._lock:\n self._sessions.pop(sid, None); self._keep.pop(sid, None)\n self._http.pop(sid, None); self._toolcache.pop(sid, None)\n self._status.pop(sid, None); self._rebuild_toolmap()\n\n # ---- connect (on loop thread) ----\n def _connect(self, sid, spec):\n cat = catalog_mod.BY_ID.get(sid, {})\n name = cat.get(\"name\", sid)\n transport = spec.get(\"transport\") or cat.get(\"transport\", \"stdio_node\")\n try:\n tools = self._lt.run(self._aopen(sid, spec, transport), timeout=75)\n self._toolcache[sid] = tools\n self._status[sid] = {\"status\": \"connected\", \"error\": None, \"name\": name,\n \"transport\": transport, \"tool_count\": len(tools)}\n except Exception as e:\n msg = (str(e) or e.__class__.__name__)\n low = msg.lower()\n if transport != \"http\" and (\"closed\" in low or \"exit\" in low or \"broken pipe\" in low or not msg.strip()):\n msg = (\"server exited on startup, it likely needs credentials or configuration. \"\n \"Edit the command to supply them (e.g. a real connection string or token).\")\n self._status[sid] = {\"status\": \"error\", \"error\": msg[:220],\n \"name\": name, \"transport\": transport, \"tool_count\": 0}\n\n async def _aopen(self, sid, spec, transport):\n if transport == \"http\":\n if not _HTTP_OK:\n raise RuntimeError(\"HTTP transport unavailable in this build\")\n self._http[sid] = spec\n async with _http_session(sid, spec) as session: # validate + discover, same task\n resp = await session.list_tools()\n return _tools(resp)\n # stdio / builtin: persistent session\n params = _stdio_params(sid, spec, transport)\n cm = stdio_client(params)\n read, write = await cm.__aenter__()\n sess_cm = ClientSession(read, write)\n session = await sess_cm.__aenter__()\n await session.initialize()\n self._keep[sid] = [cm, sess_cm]; self._sessions[sid] = session\n resp = await session.list_tools()\n return _tools(resp)\n\n async def _acall(self, sid, tool, args):\n if sid in self._http:\n async with _http_session(sid, self._http[sid]) as session: # fresh, same task\n result = await session.call_tool(tool, args or {})\n return _text(result)\n result = await self._sessions[sid].call_tool(tool, args or {})\n return _text(result)\n\n def _rebuild_toolmap(self):\n self._toolmap = {}\n for sid, tools in self._toolcache.items():\n for t in tools:\n self._toolmap[f\"{sid}__{t[\u0027name\u0027]}\"[:64]] = (sid, t[\"name\"])\n\n # ---- queries (use cache; no live calls) ----\n def connected_servers(self):\n return [dict(id=sid, **self._status[sid]) for sid in self._status]\n\n def all_tools(self):\n out = []\n with self._lock:\n for sid, tools in self._toolcache.items():\n sname = self._status.get(sid, {}).get(\"name\", sid)\n for t in tools:\n key = f\"{sid}__{t[\u0027name\u0027]}\"[:64]\n self._toolmap[key] = (sid, t[\"name\"])\n out.append({\"key\": key, \"server_id\": sid, \"server_name\": sname,\n \"tool\": t[\"name\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n def call_by_key(self, key, args):\n sid, tool = self._toolmap.get(key, (None, None))\n if sid is None:\n return \u0027{\"error\":\"unknown tool \u0027 + str(key) + \u0027\"}\u0027\n with self._lock:\n return self._lt.run(self._acall(sid, tool, args), timeout=90)\n\ndef _tools(resp):\n return [{\"name\": t.name, \"description\": t.description or \"\", \"input_schema\": t.inputSchema}\n for t in resp.tools]\ndef _text(result):\n return \"\\n\".join(c.text for c in result.content if getattr(c, \"type\", None) == \"text\")\n\n_CM = None\n_CM_LOCK = threading.Lock()\ndef manager():\n global _CM\n with _CM_LOCK:\n if _CM is None:\n _CM = _Manager()\n return _CM\n"run-e941689e
08-30 01:02:08
tool result
LOW
read_source → "\"\"\"\nThe connections catalog: a curated directory of common enterprise MCP servers.\nAccurate as of mid-2026. Each entry records who maintains it, how it connects\n(transport), what credentials it needs, and a default governance posture.\n\n\u0027maintainer\u0027: official = Anthropic reference (educational), vendor = product owner,\n community = third party, warden = ships built in with this app.\n\u0027transport\u0027: stdio_python (uvx ...), stdio_node (npx ...), http (remote OAuth URL),\n builtin (a local server this app runs itself).\n\u0027status\u0027: ready = connectable in this app\u0027s runtime now,\n needs_node = requires npx/Node in the runtime,\n needs_python = requires uvx/uv in the runtime,\n remote = a hosted URL you paste in (works without local runtime),\n archived = still works but no longer maintained upstream.\nRisk posture is a starting point; Warden classifies each discovered tool and lets\nyou override it.\n\"\"\"\n\nCATALOG = [\n # --- ships with Warden (always connectable) ---\n {\"id\":\"builtin_enterprise\",\"name\":\"Enterprise Tools (Warden)\",\"category\":\"Reference\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Customer lookup, knowledge search, ticketing, and refunds. The built-in demo server.\"},\n {\"id\":\"builtin_files\",\"name\":\"Filesystem (Warden)\",\"category\":\"Files \u0026 Docs\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Read, list, and write files inside a sandboxed workspace. A working second server.\"},\n {\"id\":\"builtin_code\",\"name\":\"Self-Audit (Warden)\",\"category\":\"Dev\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Reads Warden\u0027s own source, runs a static self-check, and proposes fixes (gated). Warden debugging Warden.\"},\n\n # --- live public remote server, no credentials, connect and go ---\n {\"id\":\"deepwiki\",\"name\":\"DeepWiki (GitHub repos)\",\"category\":\"Dev\",\"maintainer\":\"vendor\",\n \"transport\":\"http\",\"run\":\"https://mcp.deepwiki.com/mcp\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Ask real questions about any public GitHub repository and read its docs, live. No token needed.\"},\n\n # --- Anthropic official reference servers ---\n {\"id\":\"fetch\",\"name\":\"Fetch\",\"category\":\"Web\",\"maintainer\":\"official\",\"transport\":\"stdio_python\",\n \"run\":\"uvx mcp-server-fetch\",\"auth\":\"none\",\"status\":\"needs_python\",\n \"desc\":\"Fetch a URL and return its content as text for the agent to read.\"},\n {\"id\":\"filesystem\",\"name\":\"Filesystem (official)\",\"category\":\"Files \u0026 Docs\",\"maintainer\":\"official\",\n \"transport\":\"stdio_node\",\"run\":\"npx -y @modelcontextprotocol/server-filesystem \u003cpath\u003e\",\n \"auth\":\"none\",\"status\":\"needs_node\",\"desc\":\"Reference filesystem server. Read and write within allowed paths.\"},\n {\"id\":\"git\",\"name\":\"Git\",\"category\":\"Dev\",\"maintainer\":\"official\",\"transport\":\"stdio_python\",\n \"run\":\"uvx mcp-server-git --repository \u003cpath\u003e\",\"auth\":\"none\",\"status\":\"needs_python\",\n \"desc\":\"Read a repo: status, diff, log, branches, and commits.\"},\n {\"id\":\"memory\",\"name\":\"Memory\",\"category\":\"Reference\",\"maintainer\":\"official\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-memory\",\"auth\":\"none\",\"status\":\"needs_node\",\n \"desc\":\"A simple knowledge-graph memory the agent can write to and recall.\"},\n\n # --- vendor-maintained (the right pick for production) ---\n {\"id\":\"github\",\"env\":\"GITHUB_TOKEN\",\"name\":\"GitHub\",\"category\":\"Dev\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://api.githubcopilot.com/mcp/\",\"auth\":\"oauth_or_pat\",\"status\":\"remote\",\n \"desc\":\"Read repos and issues; create issues, branches, and pull requests. Hosted OAuth endpoint.\"},\n {\"id\":\"linear\",\"env\":\"LINEAR_API_KEY\",\"name\":\"Linear\",\"category\":\"Product\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.linear.app/sse\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Issue tracking and project planning. Read issues and create or update them.\"},\n {\"id\":\"notion\",\"env\":\"NOTION_TOKEN\",\"name\":\"Notion\",\"category\":\"Knowledge\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.notion.com/mcp\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Read and write Notion docs and databases. A common agent knowledge base.\"},\n {\"id\":\"stripe\",\"env\":\"STRIPE_API_KEY\",\"name\":\"Stripe\",\"category\":\"Payments\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.stripe.com\",\"auth\":\"api_key\",\"status\":\"remote\",\n \"desc\":\"Look up customers, invoices, and payments; issue refunds. High-impact by nature.\"},\n {\"id\":\"sentry\",\"env\":\"SENTRY_TOKEN\",\"name\":\"Sentry\",\"category\":\"Observability\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.sentry.dev/mcp\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Read issues and errors; triage and resolve. Hosted OAuth endpoint.\"},\n {\"id\":\"supabase\",\"env\":\"SUPABASE_ACCESS_TOKEN\",\"name\":\"Supabase\",\"category\":\"Data\",\"maintainer\":\"vendor\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @supabase/mcp-server-supabase\",\"auth\":\"api_key\",\"status\":\"needs_node\",\n \"desc\":\"Query and manage a Supabase Postgres project, tables, and rows.\"},\n {\"id\":\"playwright\",\"name\":\"Playwright\",\"category\":\"Web\",\"maintainer\":\"vendor\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @playwright/mcp\",\"auth\":\"none\",\"status\":\"needs_node\",\n \"desc\":\"Drive a real browser: navigate, click, fill forms, extract. Microsoft-maintained.\"},\n {\"id\":\"cloudflare\",\"env\":\"CLOUDFLARE_TOKEN\",\"name\":\"Cloudflare\",\"category\":\"Infra\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://observability.mcp.cloudflare.com/sse\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Inspect and manage Cloudflare resources over a hosted OAuth connection.\"},\n\n # --- data (archived reference, still functional) ---\n {\"id\":\"postgres\",\"name\":\"PostgreSQL\",\"category\":\"Data\",\"maintainer\":\"community\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-postgres postgresql://\u003cconn\u003e\",\"auth\":\"conn_string\",\n \"status\":\"archived\",\"desc\":\"Query a Postgres database. Reference server is archived; use read-only creds.\"},\n {\"id\":\"slack\",\"name\":\"Slack\",\"category\":\"Comms\",\"maintainer\":\"community\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-slack\",\"auth\":\"bot_token\",\"status\":\"archived\",\n \"desc\":\"Read channels and post messages. Original reference archived; community builds exist.\"},\n]\n\nBY_ID = {c[\"id\"]: c for c in CATALOG}\n\nMAINTAINER_LABEL = {\"warden\":\"Built in\",\"official\":\"Anthropic reference\",\n \"vendor\":\"Vendor-maintained\",\"community\":\"Community\"}\nSTATUS_LABEL = {\"ready\":\"Ready\",\"needs_node\":\"Needs Node runtime\",\"needs_python\":\"Needs Python runtime\",\n \"remote\":\"Remote (paste URL + token)\",\"archived\":\"Archived but works\"}\n"run-e941689e
08-30 01:02:08
tool result
LOW
read_source → "\"\"\"\nWarden, an enterprise AI agent studio where every agent is governed by default.\nConnect MCP servers, build an agent from their tools, run it against a live model, and\ngate high-risk actions behind human approval with a full audit trail.\n\"\"\"\nimport os\nimport datetime\nimport threading\nimport json as _json\nfrom flask import Flask, request, redirect, url_for, render_template, abort\nimport store, governance as gov, agent_runtime as rt\nimport connection_manager as cmod\nimport catalog as cat\n\nWARDEN_VERSION = \"0.3\"\n\ndef _build_info():\n \"\"\"Increment a build number on each new deploy. Identity comes from RENDER_GIT_COMMIT\n if Render provides it, else a BUILD_ID baked into the image at build time (see\n Dockerfile), else \u0027local\u0027. The counter (persisted on the disk) bumps whenever that\n identity changes; a plain restart of the same build does not bump it.\"\"\"\n import json\n ident = os.environ.get(\"RENDER_GIT_COMMIT\", \"\")\n src = \"commit\"\n if not ident:\n try:\n ident = open(os.path.join(os.path.dirname(os.path.abspath(__file__)), \"BUILD_ID\")).read().strip()\n src = \"build\"\n except Exception:\n ident = \"\"\n meta_path = os.path.join(store.DATA_ROOT, \"build.json\")\n try:\n meta = json.load(open(meta_path))\n except Exception:\n meta = {}\n num = meta.get(\"build\", 0)\n if not ident or ident != meta.get(\"ident\"):\n num += 1\n try:\n json.dump({\"ident\": ident or \"local\", \"build\": num}, open(meta_path, \"w\"))\n except Exception:\n pass\n if src == \"commit\":\n label = ident[:7]\n elif src == \"build\":\n label = ident[:13] # e.g. 20260828T1912\n else:\n label = \"local\"\n return num, label\n\n_BUILD_NUM, BUILD_COMMIT = _build_info()\nVERSION_FULL = f\"{WARDEN_VERSION}.{_BUILD_NUM}\"\n\ntry:\n from zoneinfo import ZoneInfo\n _PT = ZoneInfo(\"America/Los_Angeles\")\nexcept Exception:\n _PT = datetime.timezone(datetime.timedelta(hours=-7), \"PDT\")\n# captured once at process start; on Render each deploy restarts the process\nDEPLOYED_AT = datetime.datetime.now(_PT).strftime(\"%Y-%m-%d %H:%M %Z\")\n\napp = Flask(__name__)\nstore.init()\n\ndef _env_specs():\n \"\"\"Servers to auto-connect on boot, from WARDEN_AUTOCONNECT (comma-separated catalog ids).\n Tokens resolve from each server\u0027s env var, so config survives redeploys.\"\"\"\n ids = [x.strip() for x in os.environ.get(\"WARDEN_AUTOCONNECT\", \"\").split(\",\") if x.strip()]\n return [{\"id\": i, \"transport\": cat.BY_ID[i][\"transport\"]} for i in ids if i in cat.BY_ID]\n\ndef cm():\n c = cmod.manager()\n c.ensure_started(store.enabled_connections() + _env_specs())\n return c\n\n@app.context_processor\ndef inject_globals():\n return {\"pending\": store.pending_approvals(), \"mode\": rt.mode(),\n \"version\": VERSION_FULL, \"commit\": BUILD_COMMIT, \"deployed_at\": DEPLOYED_AT}\n\ndef connected_tools():\n \"\"\"All tools across connected servers, with effective governance risk.\"\"\"\n ovr = store.all_overrides()\n out = []\n for t in cm().all_tools():\n m = gov.meta(t[\"key\"], t[\"tool\"], t[\"description\"], ovr.get(t[\"key\"]))\n out.append({**t, \"risk\": m[\"risk\"], \"gate\": m[\"gate\"], \"override\": ovr.get(t[\"key\"])})\n return out\n\ndef tools_by_server():\n groups = {}\n for t in connected_tools():\n groups.setdefault(t[\"server_id\"], {\"name\": t[\"server_name\"], \"tools\": []})\n groups[t[\"server_id\"]][\"tools\"].append(t)\n return groups\n\n@app.route(\"/\")\ndef home():\n servers = cm().connected_servers()\n return render_template(\"dashboard.html\", agents=store.list_agents(), runs=store.list_runs(12),\n pending=store.pending_approvals(), servers=servers)\n\n@app.route(\"/connections\")\ndef connections():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"connections.html\", catalog=cat.CATALOG, status=status,\n enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL,\n tools=connected_tools())\n\n@app.route(\"/connections/enable\", methods=[\"POST\"])\ndef enable_connection():\n cid = request.form.get(\"id\"); entry = cat.BY_ID.get(cid)\n if not entry: abort(404)\n transport = entry[\"transport\"]\n token = request.form.get(\"token\") or None\n command = request.form.get(\"command\") or None\n url = request.form.get(\"url\") or entry.get(\"run\")\n store.enable_connection(cid, transport, command=command, url=url, token=token)\n st = cm().connect_spec({\"id\": cid, \"transport\": transport, \"command\": command, \"url\": url, \"token\": token})\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True, \"status\": (st or {}).get(\"status\"), \"error\": (st or {}).get(\"error\"),\n \"tool_count\": (st or {}).get(\"tool_count\", 0)}\n return redirect(url_for(\"connections\"))\n\n@app.route(\"/connections/disable\", methods=[\"POST\"])\ndef disable_connection():\n cid = request.form.get(\"id\")\n store.disable_connection(cid); cm().disconnect(cid)\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True}\n return redirect(url_for(\"connections\"))\n\n@app.route(\"/tool-risk\", methods=[\"POST\"])\ndef tool_risk():\n store.set_override(request.form.get(\"key\"), request.form.get(\"risk\"))\n return redirect(request.form.get(\"back\") or url_for(\"connections\"))\n\n@app.route(\"/connlist\")\ndef connlist():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"_connlist.html\", catalog=cat.CATALOG, status=status,\n enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL)\n\n@app.route(\"/tools.json\")\ndef tools_json():\n groups = {}\n for t in connected_tools():\n g = groups.setdefault(t[\"server_id\"], {\"server\": t[\"server_name\"], \"tools\": []})\n g[\"tools\"].append({\"key\": t[\"key\"], \"tool\": t[\"tool\"], \"risk\": t[\"risk\"],\n \"gate\": t[\"gate\"], \"description\": t[\"description\"]})\n return {\"groups\": list(groups.values())}\n\nAGENT_TEMPLATES = [\n {\"id\": \"billing\", \"name\": \"Billing Resolver\",\n \"instructions\": \"You resolve billing issues for enterprise customers. Look up the account, check policy, and make the customer whole. Be concise and never guess at numbers.\",\n \"tools\": [\"lookup_customer\", \"search_knowledge\", \"create_ticket\", \"issue_refund\"]},\n {\"id\": \"triage\", \"name\": \"Support Triage\",\n \"instructions\": \"You triage inbound support requests. Look up the customer, search the knowledge base for a known fix, and open a ticket with a clear summary when it needs a human. Do not promise resolutions you cannot verify.\",\n \"tools\": [\"lookup_customer\", \"search_knowledge\", \"create_ticket\"]},\n {\"id\": \"refund_audit\", \"name\": \"Refund Auditor (read-only)\",\n \"instructions\": \"You investigate refund requests but cannot issue refunds yourself. Look up the account, verify the charge against policy, and write a clear recommendation for a human to approve. State the exact amount and the policy basis.\",\n \"tools\": [\"lookup_customer\", \"search_knowledge\"]},\n {\"id\": \"repo_qa\", \"name\": \"Codebase Explainer\",\n \"instructions\": \"You answer questions about a public GitHub repository. Read its docs and structure, then explain how it works in plain language with references to the relevant files. If you are unsure, say so.\",\n \"tools\": [\"ask_question\", \"read_wiki_contents\", \"read_wiki_structure\"]},\n {\"id\": \"repo_maint\", \"name\": \"Repo Maintainer\",\n \"instructions\": \"You help maintain a GitHub repository. Read issues, pull requests, and code to understand the request, then propose changes. Any write (a branch, a commit, a pull request) is held for review before it runs. Never merge without explicit approval.\",\n \"tools\": [\"get_file_contents\", \"list_issues\", \"list_pull_requests\", \"search_code\",\n \"create_branch\", \"create_pull_request\", \"push_files\", \"merge_pull_request\"]},\n {\"id\": \"files\", \"name\": \"File Organizer\",\n \"instructions\": \"You organize a working folder. List and read files to understand what is there, then propose a tidier structure. Any file you create or overwrite is held for review first.\",\n \"tools\": [\"list_files\", \"read_file\", \"write_file\"]},\n {\"id\": \"self_audit\", \"name\": \"Warden Self-Audit\",\n \"instructions\": \"You audit Warden\u0027s own source code. List and read the source, run the self-check, and report concrete issues with file and line references. Any fix you propose is held as a patch for a human to review before anything changes.\",\n \"tools\": [\"list_source\", \"read_source\", \"run_selfcheck\", \"propose_patch\"]},\n {\"id\": \"kb\", \"name\": \"Knowledge Assistant\",\n \"instructions\": \"You answer policy and product questions from the internal knowledge base and public repo docs. Cite the source you used. If the answer is not in the sources, say you do not know rather than guessing.\",\n \"tools\": [\"search_knowledge\", \"ask_question\", \"read_wiki_contents\"]},\n]\n\n@app.route(\"/new\")\ndef new_agent():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"builder.html\", groups=tools_by_server(),\n catalog=cat.CATALOG, status=status, enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL,\n templates=AGENT_TEMPLATES)\n\n@app.route(\"/agents\", methods=[\"POST\"])\ndef create_agent():\n name = request.form.get(\"name\", \"\").strip() or \"Untitled agent\"\n instructions = request.form.get(\"instructions\", \"\").strip()\n model = request.form.get(\"model\", \"\").strip() or rt.MODEL_DEFAULT\n skills = request.form.getlist(\"skills\")\n aid = store.create_agent(name, instructions, model, skills)\n return redirect(url_for(\"agent\", aid=aid))\n\n@app.route(\"/agent/\u003caid\u003e\")\ndef agent(aid):\n ag = store.get_agent(aid)\n if not ag: abort(404)\n idx = {t[\"key\"]: t for t in connected_tools()}\n skills = ag[\"skills\"] or []\n granted = [idx[k] for k in skills if k in idx]\n groups = {}\n for t in granted:\n groups.setdefault(t[\"server_name\"], []).append(t)\n counts = {\"total\": len(granted), \"servers\": len(groups),\n \"gated\": sum(1 for t in granted if t[\"gate\"] == \"approval\"),\n \"auto\": sum(1 for t in granted if t[\"gate\"] != \"approval\")}\n missing = [k for k in skills if k not in idx]\n runs = [r for r in store.list_runs(50) if r[\"agent_id\"] == aid]\n return render_template(\"agent.html\", agent=ag, groups=groups, counts=counts,\n missing=missing, runs=runs)\n\ndef _advance_bg(rid):\n \"\"\"Run the agent loop in the background so the browser isn\u0027t blocked.\"\"\"\n def worker():\n try:\n rt.advance(rid)\n except Exception as e:\n try:\n r = store.get_run(rid)\n store.audit(rid, r[\"agent_id\"], \"error\", detail={\"text\": str(e)[:200]})\n store.update_run(rid, status=\"error\")\n except Exception:\n pass\n threading.Thread(target=worker, daemon=True).start()\n\n@app.route(\"/run\", methods=[\"POST\"])\ndef run():\n aid = request.form.get(\"agent_id\"); user_input = request.form.get(\"input\", \"\").strip()\n if not store.get_agent(aid) or not user_input: abort(400)\n rid = store.create_run(aid, user_input)\n _advance_bg(rid)\n return redirect(url_for(\"run_view\", rid=rid))\n\n@app.route(\"/run/\u003crid\u003e\")\ndef run_view(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n ag = store.get_agent(r[\"agent_id\"])\n return render_template(\"run.html\", run=r, agent=ag, audit=store.audit_for_run(rid),\n approvals=store.approvals_for_run(rid))\n\ndef _fmt_event(e):\n d = e.get(\"detail\") or {}\n kind = e[\"kind\"]\n if kind in (\"run_started\", \"user_message\"):\n text = d.get"run-e941689e
08-30 01:02:08
tool result
LOW
read_source → "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re, time\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nMAX_TOKENS = int(os.environ.get(\"WARDEN_MAX_TOKENS\", \"4096\"))\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\n# Estimated USD price per 1M tokens (input, output). Approximate list prices, editable;\n# used only to estimate cost for the observability view. Matched by substring of model id.\nPRICES = {\"opus\": (15.0, 75.0), \"sonnet\": (3.0, 15.0), \"haiku\": (0.80, 4.0)}\ndef _cost(model, inp, out):\n rate = (3.0, 15.0)\n m = (model or \"\").lower()\n for k, v in PRICES.items():\n if k in m:\n rate = v; break\n return round(inp / 1e6 * rate[0] + out / 1e6 * rate[1], 6)\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _repair(messages):\n \"\"\"Guarantee the API invariant: every assistant tool_use is answered by a tool_result\n in the very next message. If a tool crashed, a response was truncated, or a follow-up\n landed on a dangling turn, backfill synthetic \u0027interrupted\u0027 results so the request is\n valid. Turns a hard 400 into a graceful continuation the model can reason about.\"\"\"\n i = 0\n while i \u003c len(messages):\n m = messages[i]\n if m.get(\"role\") == \"assistant\" and isinstance(m.get(\"content\"), list):\n ids = [b[\"id\"] for b in m[\"content\"]\n if isinstance(b, dict) and b.get(\"type\") == \"tool_use\" and b.get(\"id\")]\n if ids:\n nxt = messages[i + 1] if i + 1 \u003c len(messages) else None\n answered = set()\n if nxt and nxt.get(\"role\") == \"user\" and isinstance(nxt.get(\"content\"), list):\n answered = {b.get(\"tool_use_id\") for b in nxt[\"content\"]\n if isinstance(b, dict) and b.get(\"type\") == \"tool_result\"}\n missing = [t for t in ids if t not in answered]\n if missing:\n fills = [{\"type\": \"tool_result\", \"tool_use_id\": t,\n \"content\": json.dumps({\"error\": \"interrupted\",\n \"note\": \"This tool did not complete. Do not assume it ran.\"})}\n for t in missing]\n if nxt and nxt.get(\"role\") == \"user\" and isinstance(nxt.get(\"content\"), list):\n nxt[\"content\"] = fills + nxt[\"content\"]\n else:\n messages.insert(i + 1, {\"role\": \"user\", \"content\": fills})\n i += 1\n\n_TRANSIENT = (\"rate\", \"overloaded\", \"timeout\", \"timedout\", \"internal\", \"unavailable\", \"connection\")\ndef _is_transient(ex):\n code = getattr(ex, \"status_code\", None)\n if code in (408, 409, 429, 500, 502, 503, 504, 529):\n return True\n return any(w in (type(ex).__name__ + \" \" + str(ex)).lower() for w in _TRANSIENT)\n\ndef _friendly_error(ex):\n code = getattr(ex, \"status_code\", None)\n if code == 401 or \"authentication\" in str(ex).lower():\n return \"The model rejected the API key. Check ANTHROPIC_API_KEY.\"\n if code == 429 or \"rate\" in str(ex).lower():\n return \"The model is rate-limited right now. Try again in a moment.\"\n if code and 500 \u003c= code \u003c 600:\n return \"The model service had a temporary error. Try again in a moment.\"\n if code == 400:\n return \"The model rejected the request. This run hit a malformed-request error; the transcript has been repaired, please retry.\"\n return \"The run hit an error talking to the model: \" + str(ex)[:200]\n\ndef _call_model(system, messages, tools):\n t0 = time.time()\n if SANDBOX:\n r = _sandbox_model(messages, tools)\n r[\"usage\"] = {\"input_tokens\": 0, \"output_tokens\": 0}\n r[\"model\"] = \"sandbox\"; r[\"latency_ms\"] = int((time.time() - t0) * 1000)\n return r\n import anthropic\n client = anthropic.Anthropic()\n last = None\n for attempt in range(3):\n try:\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=MAX_TOKENS,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content],\n \"usage\": {\"input_tokens\": resp.usage.input_tokens, \"output_tokens\": resp.usage.output_tokens},\n \"model\": MODEL_DEFAULT, \"latency_ms\": int((time.time() - t0) * 1000)}\n except Exception as ex:\n last = ex\n if _is_transient(ex) and attempt \u003c 2:\n time.sleep(1.5 * (attempt + 1)); continue\n raise last\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n # Intent comes from the user\u0027s actual request, not the accumulating transcript,\n # so a completed refund doesn\u0027t spuriously trigger a file-write gate.\n user_text = \"\"\n for m in messages:\n if m.get(\"role\") == \"user\" and isinstance(m.get(\"content\"), str):\n user_text = m[\"content\"].lower(); break\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", user_text)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in user_text for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in user_text for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\nimport threading\n_run_locks = {} # run_id -\u003e Lock, serializes advance() per run\n_rerun = set() # run_ids asked to advance again while already advancing\n_guard = threading.Lock()\n\ndef advance(run_id):\n \"\"\"Serialize advancing a single run. Concurrent triggers (e.g. several approvals\n decided at once) must not run the loop on the same transcript in parallel, or the\n model gets an assistant turn whose tool_use blocks aren\u0027t all answered yet (API 400).\n Only one thread advances a run at a time; triggers that arrive mid-advance cause\n exactly one more pass afterward, so the latest decisions are always picked up.\"\"\"\n with _guard:\n lock = _run_locks.setdefault(run_id, threading.Lock())\n if lock.locked():\n _rerun.add(run_id) # someone is already advancing; ask them to loop\n return store.get_run(run_id)\n with lock:\n while True:\n result = _advance_once(run_id)\n with _guard:\n if run_id in _rerun:\n _rerun.discard(run_id)\n continue # a decision landed during the pass; go again\n break\n return result\n\ndef _advance_once(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n _repair(messages) # never send an unanswered tool_use to the API\n try:\n resp = _call_model(system, messages, tools)\n except Exception as ex:\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\": _friendly_error(ex)})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n u = resp.get(\"usage\", {})\n store.audit(run_id, agent[\"id\"], \"model_call\",\n detail={\"model\": resp.get(\"model\"), \"input_tokens\": u.get(\"input_tokens\", 0),\n \"output_tokens\": u.get(\"output_tokens\", 0), \"latency_ms\": resp.get(\"latency_ms\", 0),\n \"cost\": _cost(resp.get(\"model\"), u.get(\"input_tokens\", 0), u.get(\"output_tokens\", 0))})\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \""run-e941689e
08-30 01:02:08
thought
Good โ static check is clean. Now I'll read every source file in parallel for a full manual review.run-e941689e
08-30 01:02:08
model call
{"cost": 0.036642, "input_tokens": 9124, "latency_ms": 7119, "model": "claude-sonnet-4-6", "output_tokens": 618}run-e941689e
08-30 01:02:01
tool result
LOW
run_selfcheck → {"clean": true, "files_checked": 12, "findings": []}run-e941689e
08-30 01:02:01
tool result
LOW
list_source → [{"file": "agent_runtime.py", "lines": 300}, {"file": "app.py", "lines": 471}, {"file": "catalog.py", "lines": 90}, {"file": "connection_manager.py", "lines": 197}, {"file": "governance.py", "lines": 64}, {"file": "mcp_code_server.py", "lines": 79}, {"file": "mcp_fs_server.py", "lines": 52}, {"file": "mcp_server.py", "lines": 95}, {"file": "paths.py", "lines": 32}, {"file": "store.py", "lines": 194}, {"file": "telemetry.py", "lines": 157}, {"file": "vault.py", "lines": 60}]run-e941689e
08-30 01:02:01
thought
Sure! I'll kick off the full self-audit in parallel โ listing source files, running the static self-check, and reading any available workspace files all at oncerun-e941689e
08-30 01:02:01
model call
{"cost": 0.027945, "input_tokens": 8685, "latency_ms": 3192, "model": "claude-sonnet-4-6", "output_tokens": 126}run-e941689e
08-30 01:00:15
final
The patch proposal itself was **denied by Warden's governance layer** โ which is exactly the system working as designed. The `propose_patch` tool is a gated HIGrun-36099a48
08-30 01:00:15
thought
The patch proposal itself was **denied by Warden's governance layer** โ which is exactly the system working as designed. The `propose_patch` tool is a gated HIGrun-36099a48
08-30 01:00:15
model call
{"cost": 0.136806, "input_tokens": 43677, "latency_ms": 15073, "model": "claude-sonnet-4-6", "output_tokens": 385}run-36099a48
08-30 01:00:00
denied
HIGH
propose_patch {"filename": "agent_runtime.py", "new_content": "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _call_model(system, messages, tools):\n if SANDBOX:\n return _sandbox_model(messages, tools)\n import anthropic\n client = anthropic.Anthropic()\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=1024,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content]}\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", text_in)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in text_in for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in text_in for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\ndef advance(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n resp = _call_model(system, messages, tools)\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \"thought\", detail={\"text\":blk[\"text\"]})\n if resp[\"stop_reason\"]!=\"tool_use\":\n final=\" \".join(b.get(\"text\",\"\") for b in resp[\"content\"] if b.get(\"type\")==\"text\").strip()\n store.audit(run_id, agent[\"id\"], \"final\", detail={\"text\":final})\n store.update_run(run_id, status=\"done\", transcript=messages)\n return store.get_run(run_id)\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\":\"loop bound reached\"})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n\ndef _has_tool_use(msg):\n return any(isinstance(b,dict) and b.get(\"type\")==\"tool_use\" for b in msg[\"content\"])\n\ndef _execute_tool_turn(run_id, agent, assistant_msg, messages, idx):\n blocks=[b for b in assistant_msg[\"content\"] if b.get(\"type\")==\"tool_use\"]\n for b in blocks:\n m = risk_for(b[\"name\"], idx)\n if m[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap is None:\n store.create_approval(run_id, agent[\"id\"], b[\"name\"], m[\"risk\"],\n {\"tool_use_id\":b[\"id\"],\"input\":b[\"input\"]})\n store.audit(run_id, agent[\"id\"], \"approval_request\", skill=b[\"name\"],\n risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n for b in blocks:\n if risk_for(b[\"name\"], idx)[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap and ap[\"status\"]==\"pending\":\n return \"paused\"\n results=[]\n for b in blocks:\n m=risk_for(b[\"name\"], idx); gated=m[\"gate\"]==\"approval\"\n ap=_approval_for(run_id,b[\"id\"]) if gated else None\n # FIX (Bug 1): initialise rtext before the if/else so it is always defined.\n # Previously, if gated=True but ap=None (missing approval record due to a race\n # or replay), neither branch executed and results.append() raised NameError.\n rtext = json.dumps({\"error\": \"tool result unavailable\"})\n if gated and ap and ap[\"status\"]==\"denied\":\n rtext=json.dumps({\"denied\":True,\"note\":\"A human approver denied this action. Do not retry; explain and stop.\"})\n store.audit(run_id, agent[\"id\"], \"denied\", skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n else:\n rtext=_cm().call_by_key(b[\"name\"], b[\"input\"])\n store.audit(run_id, agent[\"id\"], \"tool_result_gated\" if gated else \"tool_result\",\n skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"],\"result\":_safe(rtext)})\n results.append({\"type\":\"tool_result\",\"tool_use_id\":b[\"id\"],\"content\":rtext})\n messages.append({\"role\":\"user\",\"content\":results})\n store.update_run(run_id, transcript=messages)\n return \"executed\"\n\ndef _approval_for(run_id, tool_use_id):\n for ap in store.approvals_for_run(run_id):\n if ap[\"arguments\"].get(\"tool_use_id\")==tool_use_id:\n return ap\n return None\n\ndef _safe(t):\n try: return json.loads(t)\n except Exception: return t\n", "rationale": "Bug 1 of 4: NameError in _execute_tool_turn when a gated tool call has no approval record (ap is None). The if/else branches both assign rtext, but if gated=True and ap=None neither branch executes, leaving rtext undefined when results.append() runs. Fix: initialise rtext to a safe JSON fallback at the top of the loop body so it is always defined regardless of which branch is taken."}run-36099a48
08-29 22:46:20
approval request
HIGH
propose_patch {"filename": "agent_runtime.py", "new_content": "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _call_model(system, messages, tools):\n if SANDBOX:\n return _sandbox_model(messages, tools)\n import anthropic\n client = anthropic.Anthropic()\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=1024,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content]}\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", text_in)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in text_in for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in text_in for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\ndef advance(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n resp = _call_model(system, messages, tools)\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \"thought\", detail={\"text\":blk[\"text\"]})\n if resp[\"stop_reason\"]!=\"tool_use\":\n final=\" \".join(b.get(\"text\",\"\") for b in resp[\"content\"] if b.get(\"type\")==\"text\").strip()\n store.audit(run_id, agent[\"id\"], \"final\", detail={\"text\":final})\n store.update_run(run_id, status=\"done\", transcript=messages)\n return store.get_run(run_id)\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\":\"loop bound reached\"})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n\ndef _has_tool_use(msg):\n return any(isinstance(b,dict) and b.get(\"type\")==\"tool_use\" for b in msg[\"content\"])\n\ndef _execute_tool_turn(run_id, agent, assistant_msg, messages, idx):\n blocks=[b for b in assistant_msg[\"content\"] if b.get(\"type\")==\"tool_use\"]\n for b in blocks:\n m = risk_for(b[\"name\"], idx)\n if m[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap is None:\n store.create_approval(run_id, agent[\"id\"], b[\"name\"], m[\"risk\"],\n {\"tool_use_id\":b[\"id\"],\"input\":b[\"input\"]})\n store.audit(run_id, agent[\"id\"], \"approval_request\", skill=b[\"name\"],\n risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n for b in blocks:\n if risk_for(b[\"name\"], idx)[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap and ap[\"status\"]==\"pending\":\n return \"paused\"\n results=[]\n for b in blocks:\n m=risk_for(b[\"name\"], idx); gated=m[\"gate\"]==\"approval\"\n ap=_approval_for(run_id,b[\"id\"]) if gated else None\n # FIX (Bug 1): initialise rtext before the if/else so it is always defined.\n # Previously, if gated=True but ap=None (missing approval record due to a race\n # or replay), neither branch executed and results.append() raised NameError.\n rtext = json.dumps({\"error\": \"tool result unavailable\"})\n if gated and ap and ap[\"status\"]==\"denied\":\n rtext=json.dumps({\"denied\":True,\"note\":\"A human approver denied this action. Do not retry; explain and stop.\"})\n store.audit(run_id, agent[\"id\"], \"denied\", skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n else:\n rtext=_cm().call_by_key(b[\"name\"], b[\"input\"])\n store.audit(run_id, agent[\"id\"], \"tool_result_gated\" if gated else \"tool_result\",\n skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"],\"result\":_safe(rtext)})\n results.append({\"type\":\"tool_result\",\"tool_use_id\":b[\"id\"],\"content\":rtext})\n messages.append({\"role\":\"user\",\"content\":results})\n store.update_run(run_id, transcript=messages)\n return \"executed\"\n\ndef _approval_for(run_id, tool_use_id):\n for ap in store.approvals_for_run(run_id):\n if ap[\"arguments\"].get(\"tool_use_id\")==tool_use_id:\n return ap\n return None\n\ndef _safe(t):\n try: return json.loads(t)\n except Exception: return t\n", "rationale": "Bug 1 of 4: NameError in _execute_tool_turn when a gated tool call has no approval record (ap is None). The if/else branches both assign rtext, but if gated=True and ap=None neither branch executes, leaving rtext undefined when results.append() runs. Fix: initialise rtext to a safe JSON fallback at the top of the loop body so it is always defined regardless of which branch is taken."}run-36099a48
08-29 22:46:20
thought
The patch proposal was interrupted before it completed โ the tool requires the full `new_content` of the file, and I didn't supply it. The running source was **run-36099a48
08-29 22:46:20
model call
{"cost": 0.167472, "input_tokens": 40594, "latency_ms": 53121, "model": "claude-sonnet-4-6", "output_tokens": 3046}run-36099a48
08-29 22:45:27
what happened?run-36099a48
08-29 22:28:51
error
Error code: 400 - {'type': 'error', 'error': {'type': 'invalid_request_error', 'message': 'messages.6: `tool_use` ids were found without `tool_result` blocks imrun-fd49149d
08-29 22:28:51
Whats the next step?run-fd49149d
08-29 22:16:09
final
Here's what `welcome.txt` says:
> *"Warden workspace. Files here are readable and writable by agents, under governance."*
It's essentially a **descriptor filerun-7498eb2c
08-29 22:16:09
thought
Here's what `welcome.txt` says:
> *"Warden workspace. Files here are readable and writable by agents, under governance."*
It's essentially a **descriptor filerun-7498eb2c
08-29 22:16:08
model call
{"cost": 0.029673, "input_tokens": 8951, "latency_ms": 5291, "model": "claude-sonnet-4-6", "output_tokens": 188}run-7498eb2c
08-29 22:16:03
tool result
LOW
read_file → "Warden workspace. Files here are readable and writable by agents, under governance.\n"run-7498eb2c
08-29 22:16:03
model call
{"cost": 0.027477, "input_tokens": 8859, "latency_ms": 1812, "model": "claude-sonnet-4-6", "output_tokens": 60}run-7498eb2c
08-29 22:16:01
Yes. read itrun-7498eb2c
08-29 22:14:27
final
The workspace currently contains just **one file** at the root level:
| Name | Type |
|------|------|
| `welcome.txt` | File |
There are **no folders** yet โ run-7498eb2c
08-29 22:14:27
thought
The workspace currently contains just **one file** at the root level:
| Name | Type |
|------|------|
| `welcome.txt` | File |
There are **no folders** yet โ run-7498eb2c
08-29 22:14:27
model call
{"cost": 0.02856, "input_tokens": 8685, "latency_ms": 4465, "model": "claude-sonnet-4-6", "output_tokens": 167}run-7498eb2c
08-29 22:14:23
model call
{"cost": 0.026628, "input_tokens": 8616, "latency_ms": 4187, "model": "claude-sonnet-4-6", "output_tokens": 52}run-7498eb2c
08-29 09:33:03
error
Error code: 400 - {'type': 'error', 'error': {'type': 'invalid_request_error', 'message': 'messages.18: `tool_use` ids were found without `tool_result` blocks irun-36099a48
08-29 09:33:00
what happened?run-36099a48
08-29 05:30:04
final
Excellent โ I've read all 11 files thoroughly. Now let me compile the full bug report and propose patches for each confirmed bug.
---
## ๐ Bug Report โ Wardenrun-fd49149d
08-29 05:30:04
thought
Excellent โ I've read all 11 files thoroughly. Now let me compile the full bug report and propose patches for each confirmed bug.
---
## ๐ Bug Report โ Wardenrun-fd49149d
08-29 05:30:04
model call
{"cost": 0.123372, "input_tokens": 36004, "latency_ms": 23692, "model": "claude-sonnet-4-6", "output_tokens": 1024}run-fd49149d
08-29 05:29:41
tool result
LOW
read_source → "\"\"\"\nEncrypts secrets (connection tokens) at rest so they never sit in the database as\nplaintext. The key comes from WARDEN_SECRET_KEY if set (kept out of the data dir, the\nstronger option); otherwise a key is generated once and stored on the persistent disk\nnext to the data, so encryption works with zero configuration. Any string works as\nWARDEN_SECRET_KEY, it is hashed into a valid key.\n\"\"\"\nimport os\nimport base64\nimport hashlib\n\ntry:\n from cryptography.fernet import Fernet\n _OK = True\nexcept Exception:\n _OK = False\n\nimport paths\nDATA_ROOT = paths.DATA_ROOT\nKEY_FILE = os.path.join(DATA_ROOT, \".warden_key\")\n\ndef _key():\n env = os.environ.get(\"WARDEN_SECRET_KEY\")\n if env:\n return base64.urlsafe_b64encode(hashlib.sha256(env.encode()).digest())\n if os.path.exists(KEY_FILE):\n return open(KEY_FILE, \"rb\").read().strip()\n k = Fernet.generate_key()\n os.makedirs(DATA_ROOT, exist_ok=True)\n open(KEY_FILE, \"wb\").write(k)\n try:\n os.chmod(KEY_FILE, 0o600)\n except Exception:\n pass\n return k\n\n_F = None\ndef _fernet():\n global _F\n if _F is None:\n _F = Fernet(_key())\n return _F\n\ndef encrypt(s):\n \"\"\"Return an \u0027enc:\u0027-prefixed ciphertext, or the original if encryption is unavailable.\"\"\"\n if not s or not _OK:\n return s\n try:\n return \"enc:\" + _fernet().encrypt(s.encode()).decode()\n except Exception:\n return s\n\ndef decrypt(s):\n \"\"\"Reverse encrypt(); passes through anything not marked \u0027enc:\u0027 (e.g. legacy plaintext).\"\"\"\n if not s or not _OK or not isinstance(s, str) or not s.startswith(\"enc:\"):\n return s\n try:\n return _fernet().decrypt(s[4:].encode()).decode()\n except Exception:\n return s\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nPersistence for Warden: agents, runs, the audit log, and the approvals queue.\nSQLite on disk. On a platform with an ephemeral filesystem (Render) this resets on\nredeploy, which is fine for a demo; the point is that within a session every agent\naction and every approval is durably recorded and queryable.\n\"\"\"\nimport os\nimport json\nimport sqlite3\nimport datetime\nimport uuid\n\nimport paths\nDATA_ROOT = paths.DATA_ROOT\nDB = os.path.join(DATA_ROOT, \"warden.db\")\n\ndef _conn():\n c = sqlite3.connect(DB, timeout=10)\n c.row_factory = sqlite3.Row\n try:\n c.execute(\"PRAGMA journal_mode=WAL\")\n c.execute(\"PRAGMA busy_timeout=8000\")\n except Exception:\n pass\n return c\n\ndef now():\n return datetime.datetime.now(datetime.timezone.utc).isoformat()\n\ndef _id(prefix):\n return f\"{prefix}-{uuid.uuid4().hex[:8]}\"\n\ndef init():\n c = _conn()\n c.executescript(\"\"\"\n CREATE TABLE IF NOT EXISTS agents(\n id TEXT PRIMARY KEY, name TEXT, instructions TEXT, model TEXT,\n skills TEXT, created_at TEXT);\n CREATE TABLE IF NOT EXISTS runs(\n id TEXT PRIMARY KEY, agent_id TEXT, input TEXT, status TEXT,\n transcript TEXT, created_at TEXT, updated_at TEXT);\n CREATE TABLE IF NOT EXISTS audit(\n id TEXT PRIMARY KEY, run_id TEXT, agent_id TEXT, ts TEXT,\n kind TEXT, skill TEXT, risk TEXT, detail TEXT);\n CREATE TABLE IF NOT EXISTS approvals(\n id TEXT PRIMARY KEY, run_id TEXT, agent_id TEXT, skill TEXT, risk TEXT,\n arguments TEXT, status TEXT, created_at TEXT, decided_at TEXT, decided_by TEXT);\n CREATE TABLE IF NOT EXISTS connections(\n id TEXT PRIMARY KEY, transport TEXT, command TEXT, url TEXT, token TEXT,\n enabled INTEGER, created_at TEXT);\n CREATE TABLE IF NOT EXISTS tool_overrides(\n model_key TEXT PRIMARY KEY, risk TEXT);\n \"\"\")\n c.commit(); c.close()\n\n# ---- agents ----\ndef create_agent(name, instructions, model, skills):\n c = _conn(); aid = _id(\"ag\")\n c.execute(\"INSERT INTO agents VALUES(?,?,?,?,?,?)\",\n (aid, name, instructions, model, json.dumps(skills), now()))\n c.commit(); c.close(); return aid\n\ndef get_agent(aid):\n c = _conn(); r = c.execute(\"SELECT * FROM agents WHERE id=?\", (aid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"skills\"] = json.loads(d[\"skills\"] or \"[]\"); return d\n\ndef list_agents():\n c = _conn(); rows = c.execute(\"SELECT * FROM agents ORDER BY created_at DESC\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"skills\"] = json.loads(d[\"skills\"] or \"[]\"); out.append(d)\n return out\n\n# ---- runs ----\ndef create_run(agent_id, user_input):\n c = _conn(); rid = _id(\"run\")\n c.execute(\"INSERT INTO runs VALUES(?,?,?,?,?,?,?)\",\n (rid, agent_id, user_input, \"running\", json.dumps([]), now(), now()))\n c.commit(); c.close(); return rid\n\ndef update_run(rid, status=None, transcript=None):\n c = _conn()\n if status is not None:\n c.execute(\"UPDATE runs SET status=?, updated_at=? WHERE id=?\", (status, now(), rid))\n if transcript is not None:\n c.execute(\"UPDATE runs SET transcript=?, updated_at=? WHERE id=?\",\n (json.dumps(transcript), now(), rid))\n c.commit(); c.close()\n\ndef get_run(rid):\n c = _conn(); r = c.execute(\"SELECT * FROM runs WHERE id=?\", (rid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"transcript\"] = json.loads(d[\"transcript\"] or \"[]\"); return d\n\ndef list_runs(limit=50):\n c = _conn(); rows = c.execute(\"SELECT * FROM runs ORDER BY created_at DESC LIMIT ?\", (limit,)).fetchall(); c.close()\n return [dict(r) for r in rows]\n\n# ---- audit ----\ndef audit(run_id, agent_id, kind, skill=None, risk=None, detail=None):\n c = _conn()\n c.execute(\"INSERT INTO audit VALUES(?,?,?,?,?,?,?,?)\",\n (_id(\"ev\"), run_id, agent_id, now(), kind, skill, risk,\n json.dumps(detail) if detail is not None else None))\n c.commit(); c.close()\n\ndef audit_for_run(run_id):\n c = _conn(); rows = c.execute(\"SELECT * FROM audit WHERE run_id=? ORDER BY ts\", (run_id,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"detail\"] = json.loads(d[\"detail\"]) if d[\"detail\"] else None; out.append(d)\n return out\n\ndef audit_all(limit=200):\n c = _conn(); rows = c.execute(\"SELECT * FROM audit ORDER BY ts DESC LIMIT ?\", (limit,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"detail\"] = json.loads(d[\"detail\"]) if d[\"detail\"] else None; out.append(d)\n return out\n\n# ---- approvals ----\ndef create_approval(run_id, agent_id, skill, risk, arguments):\n c = _conn(); apid = _id(\"ap\")\n c.execute(\"INSERT INTO approvals VALUES(?,?,?,?,?,?,?,?,?,?)\",\n (apid, run_id, agent_id, skill, risk, json.dumps(arguments),\n \"pending\", now(), None, None))\n c.commit(); c.close(); return apid\n\ndef get_approval(apid):\n c = _conn(); r = c.execute(\"SELECT * FROM approvals WHERE id=?\", (apid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); return d\n\ndef decide_approval(apid, status, by=\"operator\"):\n c = _conn()\n c.execute(\"UPDATE approvals SET status=?, decided_at=?, decided_by=? WHERE id=?\",\n (status, now(), by, apid))\n c.commit(); c.close()\n\ndef pending_approvals():\n c = _conn(); rows = c.execute(\"SELECT * FROM approvals WHERE status=\u0027pending\u0027 ORDER BY created_at\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); out.append(d)\n return out\n\ndef approvals_for_run(run_id):\n c = _conn(); rows = c.execute(\"SELECT * FROM approvals WHERE run_id=? ORDER BY created_at\", (run_id,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); out.append(d)\n return out\n\n# ---- connections (enabled external MCP servers) ----\ndef enable_connection(cid, transport, command=None, url=None, token=None):\n import vault\n c = _conn()\n c.execute(\"INSERT OR REPLACE INTO connections VALUES(?,?,?,?,?,?,?)\",\n (cid, transport, command, url, vault.encrypt(token), 1, now()))\n c.commit(); c.close()\n\ndef disable_connection(cid):\n c = _conn(); c.execute(\"DELETE FROM connections WHERE id=?\", (cid,)); c.commit(); c.close()\n\ndef enabled_connections():\n import vault\n c = _conn(); rows = c.execute(\"SELECT * FROM connections WHERE enabled=1\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r)\n out.append({\"id\": d[\"id\"], \"transport\": d[\"transport\"], \"command\": d[\"command\"],\n \"url\": d[\"url\"], \"token\": vault.decrypt(d[\"token\"])})\n return out\n\ndef is_enabled(cid):\n c = _conn(); r = c.execute(\"SELECT 1 FROM connections WHERE id=? AND enabled=1\", (cid,)).fetchone(); c.close()\n return bool(r)\n\n# ---- tool risk overrides ----\ndef set_override(model_key, risk):\n c = _conn(); c.execute(\"INSERT OR REPLACE INTO tool_overrides VALUES(?,?)\", (model_key, risk)); c.commit(); c.close()\n\ndef get_override(model_key):\n c = _conn(); r = c.execute(\"SELECT risk FROM tool_overrides WHERE model_key=?\", (model_key,)).fetchone(); c.close()\n return r[\"risk\"] if r else None\n\ndef all_overrides():\n c = _conn(); rows = c.execute(\"SELECT * FROM tool_overrides\").fetchall(); c.close()\n return {r[\"model_key\"]: r[\"risk\"] for r in rows}\n\ndef approval_counts():\n c = _conn(); rows = c.execute(\"SELECT status, COUNT(*) n FROM approvals GROUP BY status\").fetchall(); c.close()\n return {r[\"status\"]: r[\"n\"] for r in rows}\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nResolves where Warden stores its data. Prefers WARDEN_DATA_DIR (a mounted persistent\ndisk in production). If that path can\u0027t be created or written, it falls back to a\nwritable local directory instead of crashing the app, and records that it fell back so\nthe condition is visible on /healthz. A misconfigured disk should degrade to ephemeral,\nnever take the service down.\n\"\"\"\nimport os\n\n_APPDIR = os.path.dirname(os.path.abspath(__file__))\nREQUESTED = os.environ.get(\"WARDEN_DATA_DIR\") or _APPDIR\n\ndef _writable(path):\n try:\n os.makedirs(path, exist_ok=True)\n t = os.path.join(path, \".wtest\")\n with open(t, \"w\") as f:\n f.write(\"ok\")\n os.remove(t)\n return True\n except Exception:\n return False\n\ndef _resolve():\n if _writable(REQUESTED):\n return REQUESTED, False\n fallback = os.path.join(_APPDIR, \"_localdata\")\n if _writable(fallback):\n return fallback, True\n return _APPDIR, True\n\nDATA_ROOT, FALLBACK = _resolve()\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nWarden MCP server: a real Model Context Protocol server exposing a small set of\nenterprise-flavored tools. The Warden runtime connects to this as an MCP client,\ndiscovers these tools over the protocol, and invokes them.\n\nTools deliberately span read and write so the governance layer has something to\ngovern: reads are low risk and auto-execute; writes change state and are the ones\nthe studio gates behind human approval.\n\nRun standalone for a protocol smoke test: python mcp_server.py\n(but normally it is spawned over stdio by the runtime\u0027s MCP client)\n\"\"\"\nimport json\nimport os\nimport datetime\nfrom mcp.server.fastmcp import FastMCP\n\nimport paths\nDATA_DIR = os.path.join(paths.DATA_ROOT, \"data\")\nos.makedirs(DATA_DIR, exist_ok=True)\n\nCUSTOMERS = {\n \"AC-1001\": {\"account_id\": \"AC-1001\", \"name\": \"Rivera Logistics\", \"plan\": \"Enterprise\",\n \"mrr\": 4200, \"status\": \"active\", \"last_charge\": 4200,\n \"notes\": \"Charged twice on 2026-08-03 due to a billing retry bug.\"},\n \"AC-1002\": {\"account_id\": \"AC-1002\", \"name\": \"Halcyon Health\", \"plan\": \"Premium\",\n \"mrr\": 1800, \"status\": \"active\", \"last_charge\": 1800, \"notes\": \"\"},\n \"AC-1003\": {\"account_id\": \"AC-1003\", \"name\": \"Meridian Foods\", \"plan\": \"Standard\",\n \"mrr\": 600, \"status\": \"past_due\", \"last_charge\": 0,\n \"notes\": \"Payment failed twice this month.\"},\n}\n\nKB = [\n {\"id\": \"kb-01\", \"title\": \"Refund policy\",\n \"body\": \"Duplicate charges are refunded in full once verified against the billing ledger. \"\n \"Refunds above 1000 require a human approver.\"},\n {\"id\": \"kb-02\", \"title\": \"Past-due accounts\",\n \"body\": \"Do not issue credits on past-due accounts until the balance is cleared. \"\n \"Open a billing ticket instead.\"},\n {\"id\": \"kb-03\", \"title\": \"Escalation\",\n \"body\": \"Anything touching money movement is a governed action and must be logged.\"},\n]\n\ndef _ledger_path(name):\n return os.path.join(DATA_DIR, name)\n\ndef _append(name, row):\n path = _ledger_path(name)\n rows = []\n if os.path.exists(path):\n rows = json.load(open(path))\n row[\"at\"] = datetime.datetime.now(datetime.timezone.utc).isoformat()\n rows.append(row)\n json.dump(rows, open(path, \"w\"), indent=2)\n return row\n\nmcp = FastMCP(\"warden-enterprise-tools\")\n\n@mcp.tool()\ndef lookup_customer(account_id: str) -\u003e str:\n \"\"\"Look up an enterprise customer account by id (e.g. AC-1001). Read only.\"\"\"\n rec = CUSTOMERS.get(account_id.strip().upper())\n if not rec:\n return json.dumps({\"error\": f\"no account {account_id}\"})\n return json.dumps(rec)\n\n@mcp.tool()\ndef search_knowledge(query: str) -\u003e str:\n \"\"\"Search the internal knowledge base for policy and guidance. Read only.\"\"\"\n q = query.lower()\n hits = [a for a in KB if q in a[\"title\"].lower() or q in a[\"body\"].lower()]\n if not hits:\n hits = KB # fall back to returning all short KB so the agent has context\n return json.dumps(hits)\n\n@mcp.tool()\ndef create_ticket(subject: str, body: str) -\u003e str:\n \"\"\"Open an internal billing/support ticket. Write action (changes state).\"\"\"\n row = _append(\"tickets.json\", {\"id\": f\"TK-{abs(hash(subject)) % 9000 + 1000}\",\n \"subject\": subject, \"body\": body})\n return json.dumps({\"created\": True, \"ticket\": row})\n\n@mcp.tool()\ndef issue_refund(account_id: str, amount: float, reason: str = \"\") -\u003e str:\n \"\"\"Issue a monetary refund to a customer account. High-impact write action:\n moves money, so the studio gates this behind human approval before it runs.\"\"\"\n acct = account_id.strip().upper()\n if acct not in CUSTOMERS:\n return json.dumps({\"error\": f\"no account {account_id}\"})\n row = _append(\"refunds.json\", {\"id\": f\"RF-{abs(hash(acct + str(amount))) % 9000 + 1000}\",\n \"account_id\": acct, \"amount\": float(amount), \"reason\": reason})\n return json.dumps({\"refunded\": True, \"refund\": row})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nWarden\u0027s built-in Filesystem MCP server. A real second MCP server (local, Python) so\nmulti-server connection and routing is demonstrable without needing Node. All access\nis confined to a sandbox workspace directory.\n\"\"\"\nimport os\nimport json\nfrom mcp.server.fastmcp import FastMCP\n\nimport paths\nWORKSPACE = os.path.join(paths.DATA_ROOT, \"data\", \"workspace\")\nos.makedirs(WORKSPACE, exist_ok=True)\n\n# seed one file so reads have something to find\n_seed = os.path.join(WORKSPACE, \"welcome.txt\")\nif not os.path.exists(_seed):\n open(_seed, \"w\").write(\"Warden workspace. Files here are readable and writable by agents, under governance.\\n\")\n\ndef _safe(path):\n p = os.path.abspath(os.path.join(WORKSPACE, path.lstrip(\"/\")))\n if not p.startswith(os.path.abspath(WORKSPACE)):\n raise ValueError(\"path escapes workspace\")\n return p\n\nmcp = FastMCP(\"warden-filesystem\")\n\n@mcp.tool()\ndef list_files(subdir: str = \"\") -\u003e str:\n \"\"\"List files in the workspace (or a subdirectory). Read only.\"\"\"\n base = _safe(subdir)\n if not os.path.exists(base):\n return json.dumps({\"error\": \"no such path\"})\n return json.dumps(sorted(os.listdir(base)))\n\n@mcp.tool()\ndef read_file(path: str) -\u003e str:\n \"\"\"Read a text file from the workspace. Read only.\"\"\"\n p = _safe(path)\n if not os.path.isfile(p):\n return json.dumps({\"error\": f\"no file {path}\"})\n return open(p, encoding=\"utf-8\", errors=\"replace\").read()[:8000]\n\n@mcp.tool()\ndef write_file(path: str, content: str) -\u003e str:\n \"\"\"Create or overwrite a text file in the workspace. Write action (changes state).\"\"\"\n p = _safe(path)\n os.makedirs(os.path.dirname(p), exist_ok=True)\n open(p, \"w\", encoding=\"utf-8\").write(content)\n return json.dumps({\"written\": True, \"path\": path, \"bytes\": len(content)})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nWarden self-audit MCP server. Exposes Warden\u0027s own source code to an agent so a\n\"Warden Engineer\" agent can inspect the running codebase, run a real static self-check,\nand propose fixes. Reads are safe and auto-run; proposing a patch is a gated write that\nlands in a review folder (it never overwrites the running source).\n\"\"\"\nimport os, json, ast, py_compile, tempfile\nfrom mcp.server.fastmcp import FastMCP\n\nAPP_DIR = os.path.dirname(os.path.abspath(__file__))\nimport paths\nPATCH_DIR = os.path.join(paths.DATA_ROOT, \"data\", \"patches\")\nos.makedirs(PATCH_DIR, exist_ok=True)\n\ndef _py_files():\n return sorted(f for f in os.listdir(APP_DIR) if f.endswith(\".py\"))\n\nmcp = FastMCP(\"warden-self-audit\")\n\n@mcp.tool()\ndef list_source() -\u003e str:\n \"\"\"List Warden\u0027s own Python source files with line counts. Read only.\"\"\"\n out = []\n for f in _py_files():\n n = sum(1 for _ in open(os.path.join(APP_DIR, f), encoding=\"utf-8\", errors=\"replace\"))\n out.append({\"file\": f, \"lines\": n})\n return json.dumps(out)\n\n@mcp.tool()\ndef read_source(filename: str) -\u003e str:\n \"\"\"Read one of Warden\u0027s own source files. Read only.\"\"\"\n if filename not in _py_files():\n return json.dumps({\"error\": f\"no source file {filename}\", \"available\": _py_files()})\n return open(os.path.join(APP_DIR, filename), encoding=\"utf-8\", errors=\"replace\").read()[:12000]\n\n@mcp.tool()\ndef run_selfcheck() -\u003e str:\n \"\"\"Statically check every Warden source file: byte-compile for syntax errors and\n AST-scan for bare excepts and TODO/FIXME markers. Runs real checks; no side effects.\"\"\"\n findings = []\n for f in _py_files():\n path = os.path.join(APP_DIR, f)\n try:\n py_compile.compile(path, doraise=True)\n except py_compile.PyCompileError as e:\n findings.append({\"file\": f, \"kind\": \"syntax_error\", \"detail\": str(e).splitlines()[-1][:160]})\n continue\n src = open(path, encoding=\"utf-8\", errors=\"replace\").read()\n try:\n tree = ast.parse(src)\n for node in ast.walk(tree):\n if isinstance(node, ast.ExceptHandler) and node.type is None:\n findings.append({\"file\": f, \"kind\": \"bare_except\", \"line\": node.lineno})\n except Exception:\n pass\n try:\n import tokenize, io\n for tok in tokenize.generate_tokens(io.StringIO(src).readline):\n if tok.type == tokenize.COMMENT and (\"TODO\" in tok.string.upper() or \"FIXME\" in tok.string.upper()):\n findings.append({\"file\": f, \"kind\": \"todo\", \"line\": tok.start[0], \"detail\": tok.string.strip()[:120]})\n except Exception:\n pass\n return json.dumps({\"files_checked\": len(_py_files()),\n \"findings\": findings, \"clean\": len(findings) == 0})\n\n@mcp.tool()\ndef propose_patch(filename: str, new_content: str, rationale: str = \"\") -\u003e str:\n \"\"\"Propose a fix for a source file. Gated write: saves the proposed version to a\n review folder for a human to inspect and apply. Never edits the running source.\"\"\"\n safe = os.path.basename(filename)\n out = os.path.join(PATCH_DIR, safe)\n open(out, \"w\", encoding=\"utf-8\").write(new_content)\n meta = os.path.join(PATCH_DIR, safe + \".rationale.txt\")\n open(meta, \"w\", encoding=\"utf-8\").write(rationale or \"(none)\")\n return json.dumps({\"proposed\": True, \"review_path\": f\"data/patches/{safe}\",\n \"bytes\": len(new_content), \"note\": \"saved for human review; running source unchanged\"})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nThe governance layer. Warden\u0027s point of view: tools can do things; governance decides\nwhich run on their own and which pause for a human. Built-in enterprise tools have a\nhand-set risk registry. Tools discovered from external MCP servers are classified\nautomatically, fail-closed: reads run, writes and anything unrecognized are gated.\nAn operator can override any tool\u0027s risk.\n\"\"\"\nSKILLS = {\n \"lookup_customer\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"search_knowledge\":{\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"create_ticket\": {\"risk\":\"MED\",\"gate\":\"auto\",\"kind\":\"write\"},\n \"issue_refund\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n \"list_files\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"read_file\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"write_file\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n \"list_source\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"read_source\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"run_selfcheck\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"propose_patch\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n}\n\nREAD_HINTS = (\"get\",\"list\",\"read\",\"search\",\"lookup\",\"fetch\",\"find\",\"query\",\"view\",\n \"describe\",\"show\",\"count\",\"status\",\"summary\",\"recent\",\"ask\",\"explore\",\n \"inspect\",\"check\")\nWRITE_HINTS = (\"create\",\"write\",\"update\",\"delete\",\"remove\",\"issue\",\"send\",\"post\",\"add\",\n \"set\",\"merge\",\"close\",\"open\",\"deploy\",\"execute\",\"run\",\"refund\",\"cancel\",\n \"approve\",\"edit\",\"upload\",\"move\",\"rename\",\"revoke\",\"grant\",\"pay\",\"charge\")\n\ndef classify(name, desc=\"\"):\n \"\"\"Heuristic risk for an external tool, judged by its leading verb.\n Fail-closed: unrecognized =\u003e HIGH.\"\"\"\n import re\n n = (name or \"\").lower().replace(\"-\", \"_\")\n lead = re.match(r\"[a-z]+\", n)\n lead = lead.group(0) if lead else n # leading alpha run (handles camelCase)\n first = n.split(\"_\")[0] # first snake_case token\n for h in READ_HINTS:\n if first == h or lead.startswith(h):\n return \"LOW\"\n for h in WRITE_HINTS:\n if first == h or lead.startswith(h):\n return \"HIGH\"\n return \"HIGH\"\n\ndef meta(model_key, tool_name, desc=\"\", override=None):\n \"\"\"Resolve effective governance for a tool. Precedence: override \u003e known registry \u003e classify.\"\"\"\n if override in (\"LOW\",\"MED\",\"HIGH\"):\n risk = override\n elif tool_name in SKILLS:\n risk = SKILLS[tool_name][\"risk\"]\n else:\n risk = classify(tool_name, desc)\n return {\"risk\": risk, \"gate\": \"approval\" if risk == \"HIGH\" else \"auto\"}\n\n# convenience wrappers used where only a bare name is available (built-ins)\ndef skill_meta(name):\n m = SKILLS.get(name)\n if m: return m\n r = classify(name)\n return {\"risk\": r, \"gate\": \"approval\" if r==\"HIGH\" else \"auto\", \"kind\":\"?\"}\ndef requires_approval(name):\n return skill_meta(name)[\"gate\"] == \"approval\"\ndef risk_of(name):\n return skill_meta(name)[\"risk\"]\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nMulti-server MCP connection manager.\n\nOwns one dedicated asyncio loop thread. Tool discovery is done once at connect time and\ncached, so listing tools never re-hits a server. For execution: stdio/builtin servers\nkeep a persistent session (subprocess spawn is expensive); HTTP servers open a fresh\nshort-lived session per call, entirely within one coroutine, because the streamable-HTTP\ntransport binds its cancel scope to the creating task and cannot be reused across tasks.\n\"\"\"\nimport os, sys, asyncio, threading, warnings, shlex\nfrom contextlib import asynccontextmanager\nfrom mcp import ClientSession, StdioServerParameters\nfrom mcp.client.stdio import stdio_client\ntry:\n from mcp.client.streamable_http import streamablehttp_client\n _HTTP_OK = True\nexcept Exception:\n _HTTP_OK = False\n\nimport catalog as catalog_mod\n\nHERE = os.path.dirname(os.path.abspath(__file__))\nBUILTINS = {\n \"builtin_enterprise\": [sys.executable, os.path.join(HERE, \"mcp_server.py\")],\n \"builtin_files\": [sys.executable, os.path.join(HERE, \"mcp_fs_server.py\")],\n \"builtin_code\": [sys.executable, os.path.join(HERE, \"mcp_code_server.py\")],\n}\n\nclass _LoopThread:\n def __init__(self):\n self.loop = asyncio.new_event_loop()\n with warnings.catch_warnings():\n warnings.simplefilter(\"ignore\")\n try:\n w = asyncio.ThreadedChildWatcher(); w.attach_loop(self.loop)\n asyncio.set_child_watcher(w)\n except Exception:\n pass\n threading.Thread(target=self._run, daemon=True).start()\n def _run(self):\n asyncio.set_event_loop(self.loop); self.loop.run_forever()\n def run(self, coro, timeout=None):\n return asyncio.run_coroutine_threadsafe(coro, self.loop).result(timeout)\n\ndef _http_params(sid, spec):\n cat = catalog_mod.BY_ID.get(sid, {})\n url = spec.get(\"url\") or cat.get(\"run\")\n headers = {}\n tok = spec.get(\"token\")\n if not tok and cat.get(\"env\"): # fall back to an environment variable\n tok = os.environ.get(cat[\"env\"])\n if tok:\n headers[\"Authorization\"] = tok if tok.lower().startswith(\"bearer\") else f\"Bearer {tok}\"\n return url, (headers or None)\n\ndef _stdio_params(sid, spec, transport):\n if transport == \"builtin\":\n cmd = BUILTINS[sid]\n else:\n run = spec.get(\"command\") or catalog_mod.BY_ID.get(sid, {}).get(\"run\", \"\")\n cmd = shlex.split(run)\n if not cmd:\n raise RuntimeError(\"no command configured\")\n return StdioServerParameters(command=cmd[0], args=cmd[1:], env=os.environ.copy())\n\n@asynccontextmanager\nasync def _http_session(sid, spec):\n url, headers = _http_params(sid, spec)\n async with streamablehttp_client(url, headers=headers) as (read, write, _):\n async with ClientSession(read, write) as session:\n await session.initialize()\n yield session\n\nclass _Manager:\n def __init__(self):\n self._lt = _LoopThread()\n self._lock = threading.Lock()\n self._sessions = {} # sid -\u003e persistent session (stdio/builtin only)\n self._keep = {} # sid -\u003e [context managers] to keep alive\n self._http = {} # sid -\u003e spec (http servers, fresh session per call)\n self._toolcache = {} # sid -\u003e [ {name,description,input_schema} ] (from connect)\n self._status = {} # sid -\u003e status dict\n self._toolmap = {} # model_name -\u003e (sid, tool)\n self._started = False\n\n def ensure_started(self, enabled_specs=None):\n with self._lock:\n if not self._started:\n for sid in BUILTINS:\n self._connect(sid, {\"id\": sid, \"transport\": \"builtin\"})\n self._started = True\n for spec in (enabled_specs or []):\n if spec[\"id\"] not in self._status:\n self._connect(spec[\"id\"], spec)\n self._rebuild_toolmap()\n\n def connect_spec(self, spec):\n with self._lock:\n self._connect(spec[\"id\"], spec); self._rebuild_toolmap()\n return self._status.get(spec[\"id\"])\n\n def disconnect(self, sid):\n with self._lock:\n self._sessions.pop(sid, None); self._keep.pop(sid, None)\n self._http.pop(sid, None); self._toolcache.pop(sid, None)\n self._status.pop(sid, None); self._rebuild_toolmap()\n\n # ---- connect (on loop thread) ----\n def _connect(self, sid, spec):\n cat = catalog_mod.BY_ID.get(sid, {})\n name = cat.get(\"name\", sid)\n transport = spec.get(\"transport\") or cat.get(\"transport\", \"stdio_node\")\n try:\n tools = self._lt.run(self._aopen(sid, spec, transport), timeout=75)\n self._toolcache[sid] = tools\n self._status[sid] = {\"status\": \"connected\", \"error\": None, \"name\": name,\n \"transport\": transport, \"tool_count\": len(tools)}\n except Exception as e:\n msg = (str(e) or e.__class__.__name__)\n low = msg.lower()\n if transport != \"http\" and (\"closed\" in low or \"exit\" in low or \"broken pipe\" in low or not msg.strip()):\n msg = (\"server exited on startup, it likely needs credentials or configuration. \"\n \"Edit the command to supply them (e.g. a real connection string or token).\")\n self._status[sid] = {\"status\": \"error\", \"error\": msg[:220],\n \"name\": name, \"transport\": transport, \"tool_count\": 0}\n\n async def _aopen(self, sid, spec, transport):\n if transport == \"http\":\n if not _HTTP_OK:\n raise RuntimeError(\"HTTP transport unavailable in this build\")\n self._http[sid] = spec\n async with _http_session(sid, spec) as session: # validate + discover, same task\n resp = await session.list_tools()\n return _tools(resp)\n # stdio / builtin: persistent session\n params = _stdio_params(sid, spec, transport)\n cm = stdio_client(params)\n read, write = await cm.__aenter__()\n sess_cm = ClientSession(read, write)\n session = await sess_cm.__aenter__()\n await session.initialize()\n self._keep[sid] = [cm, sess_cm]; self._sessions[sid] = session\n resp = await session.list_tools()\n return _tools(resp)\n\n async def _acall(self, sid, tool, args):\n if sid in self._http:\n async with _http_session(sid, self._http[sid]) as session: # fresh, same task\n result = await session.call_tool(tool, args or {})\n return _text(result)\n result = await self._sessions[sid].call_tool(tool, args or {})\n return _text(result)\n\n def _rebuild_toolmap(self):\n self._toolmap = {}\n for sid, tools in self._toolcache.items():\n for t in tools:\n self._toolmap[f\"{sid}__{t[\u0027name\u0027]}\"[:64]] = (sid, t[\"name\"])\n\n # ---- queries (use cache; no live calls) ----\n def connected_servers(self):\n return [dict(id=sid, **self._status[sid]) for sid in self._status]\n\n def all_tools(self):\n out = []\n with self._lock:\n for sid, tools in self._toolcache.items():\n sname = self._status.get(sid, {}).get(\"name\", sid)\n for t in tools:\n key = f\"{sid}__{t[\u0027name\u0027]}\"[:64]\n self._toolmap[key] = (sid, t[\"name\"])\n out.append({\"key\": key, \"server_id\": sid, \"server_name\": sname,\n \"tool\": t[\"name\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n def call_by_key(self, key, args):\n sid, tool = self._toolmap.get(key, (None, None))\n if sid is None:\n return \u0027{\"error\":\"unknown tool \u0027 + str(key) + \u0027\"}\u0027\n with self._lock:\n return self._lt.run(self._acall(sid, tool, args), timeout=90)\n\ndef _tools(resp):\n return [{\"name\": t.name, \"description\": t.description or \"\", \"input_schema\": t.inputSchema}\n for t in resp.tools]\ndef _text(result):\n return \"\\n\".join(c.text for c in result.content if getattr(c, \"type\", None) == \"text\")\n\n_CM = None\n_CM_LOCK = threading.Lock()\ndef manager():\n global _CM\n with _CM_LOCK:\n if _CM is None:\n _CM = _Manager()\n return _CM\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nThe connections catalog: a curated directory of common enterprise MCP servers.\nAccurate as of mid-2026. Each entry records who maintains it, how it connects\n(transport), what credentials it needs, and a default governance posture.\n\n\u0027maintainer\u0027: official = Anthropic reference (educational), vendor = product owner,\n community = third party, warden = ships built in with this app.\n\u0027transport\u0027: stdio_python (uvx ...), stdio_node (npx ...), http (remote OAuth URL),\n builtin (a local server this app runs itself).\n\u0027status\u0027: ready = connectable in this app\u0027s runtime now,\n needs_node = requires npx/Node in the runtime,\n needs_python = requires uvx/uv in the runtime,\n remote = a hosted URL you paste in (works without local runtime),\n archived = still works but no longer maintained upstream.\nRisk posture is a starting point; Warden classifies each discovered tool and lets\nyou override it.\n\"\"\"\n\nCATALOG = [\n # --- ships with Warden (always connectable) ---\n {\"id\":\"builtin_enterprise\",\"name\":\"Enterprise Tools (Warden)\",\"category\":\"Reference\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Customer lookup, knowledge search, ticketing, and refunds. The built-in demo server.\"},\n {\"id\":\"builtin_files\",\"name\":\"Filesystem (Warden)\",\"category\":\"Files \u0026 Docs\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Read, list, and write files inside a sandboxed workspace. A working second server.\"},\n {\"id\":\"builtin_code\",\"name\":\"Self-Audit (Warden)\",\"category\":\"Dev\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Reads Warden\u0027s own source, runs a static self-check, and proposes fixes (gated). Warden debugging Warden.\"},\n\n # --- live public remote server, no credentials, connect and go ---\n {\"id\":\"deepwiki\",\"name\":\"DeepWiki (GitHub repos)\",\"category\":\"Dev\",\"maintainer\":\"vendor\",\n \"transport\":\"http\",\"run\":\"https://mcp.deepwiki.com/mcp\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Ask real questions about any public GitHub repository and read its docs, live. No token needed.\"},\n\n # --- Anthropic official reference servers ---\n {\"id\":\"fetch\",\"name\":\"Fetch\",\"category\":\"Web\",\"maintainer\":\"official\",\"transport\":\"stdio_python\",\n \"run\":\"uvx mcp-server-fetch\",\"auth\":\"none\",\"status\":\"needs_python\",\n \"desc\":\"Fetch a URL and return its content as text for the agent to read.\"},\n {\"id\":\"filesystem\",\"name\":\"Filesystem (official)\",\"category\":\"Files \u0026 Docs\",\"maintainer\":\"official\",\n \"transport\":\"stdio_node\",\"run\":\"npx -y @modelcontextprotocol/server-filesystem \u003cpath\u003e\",\n \"auth\":\"none\",\"status\":\"needs_node\",\"desc\":\"Reference filesystem server. Read and write within allowed paths.\"},\n {\"id\":\"git\",\"name\":\"Git\",\"category\":\"Dev\",\"maintainer\":\"official\",\"transport\":\"stdio_python\",\n \"run\":\"uvx mcp-server-git --repository \u003cpath\u003e\",\"auth\":\"none\",\"status\":\"needs_python\",\n \"desc\":\"Read a repo: status, diff, log, branches, and commits.\"},\n {\"id\":\"memory\",\"name\":\"Memory\",\"category\":\"Reference\",\"maintainer\":\"official\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-memory\",\"auth\":\"none\",\"status\":\"needs_node\",\n \"desc\":\"A simple knowledge-graph memory the agent can write to and recall.\"},\n\n # --- vendor-maintained (the right pick for production) ---\n {\"id\":\"github\",\"env\":\"GITHUB_TOKEN\",\"name\":\"GitHub\",\"category\":\"Dev\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://api.githubcopilot.com/mcp/\",\"auth\":\"oauth_or_pat\",\"status\":\"remote\",\n \"desc\":\"Read repos and issues; create issues, branches, and pull requests. Hosted OAuth endpoint.\"},\n {\"id\":\"linear\",\"env\":\"LINEAR_API_KEY\",\"name\":\"Linear\",\"category\":\"Product\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.linear.app/sse\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Issue tracking and project planning. Read issues and create or update them.\"},\n {\"id\":\"notion\",\"env\":\"NOTION_TOKEN\",\"name\":\"Notion\",\"category\":\"Knowledge\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.notion.com/mcp\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Read and write Notion docs and databases. A common agent knowledge base.\"},\n {\"id\":\"stripe\",\"env\":\"STRIPE_API_KEY\",\"name\":\"Stripe\",\"category\":\"Payments\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.stripe.com\",\"auth\":\"api_key\",\"status\":\"remote\",\n \"desc\":\"Look up customers, invoices, and payments; issue refunds. High-impact by nature.\"},\n {\"id\":\"sentry\",\"env\":\"SENTRY_TOKEN\",\"name\":\"Sentry\",\"category\":\"Observability\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.sentry.dev/mcp\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Read issues and errors; triage and resolve. Hosted OAuth endpoint.\"},\n {\"id\":\"supabase\",\"env\":\"SUPABASE_ACCESS_TOKEN\",\"name\":\"Supabase\",\"category\":\"Data\",\"maintainer\":\"vendor\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @supabase/mcp-server-supabase\",\"auth\":\"api_key\",\"status\":\"needs_node\",\n \"desc\":\"Query and manage a Supabase Postgres project, tables, and rows.\"},\n {\"id\":\"playwright\",\"name\":\"Playwright\",\"category\":\"Web\",\"maintainer\":\"vendor\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @playwright/mcp\",\"auth\":\"none\",\"status\":\"needs_node\",\n \"desc\":\"Drive a real browser: navigate, click, fill forms, extract. Microsoft-maintained.\"},\n {\"id\":\"cloudflare\",\"env\":\"CLOUDFLARE_TOKEN\",\"name\":\"Cloudflare\",\"category\":\"Infra\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://observability.mcp.cloudflare.com/sse\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Inspect and manage Cloudflare resources over a hosted OAuth connection.\"},\n\n # --- data (archived reference, still functional) ---\n {\"id\":\"postgres\",\"name\":\"PostgreSQL\",\"category\":\"Data\",\"maintainer\":\"community\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-postgres postgresql://\u003cconn\u003e\",\"auth\":\"conn_string\",\n \"status\":\"archived\",\"desc\":\"Query a Postgres database. Reference server is archived; use read-only creds.\"},\n {\"id\":\"slack\",\"name\":\"Slack\",\"category\":\"Comms\",\"maintainer\":\"community\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-slack\",\"auth\":\"bot_token\",\"status\":\"archived\",\n \"desc\":\"Read channels and post messages. Original reference archived; community builds exist.\"},\n]\n\nBY_ID = {c[\"id\"]: c for c in CATALOG}\n\nMAINTAINER_LABEL = {\"warden\":\"Built in\",\"official\":\"Anthropic reference\",\n \"vendor\":\"Vendor-maintained\",\"community\":\"Community\"}\nSTATUS_LABEL = {\"ready\":\"Ready\",\"needs_node\":\"Needs Node runtime\",\"needs_python\":\"Needs Python runtime\",\n \"remote\":\"Remote (paste URL + token)\",\"archived\":\"Archived but works\"}\n"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nWarden, an enterprise AI agent studio where every agent is governed by default.\nConnect MCP servers, build an agent from their tools, run it against a live model, and\ngate high-risk actions behind human approval with a full audit trail.\n\"\"\"\nimport os\nimport datetime\nimport threading\nimport json as _json\nfrom flask import Flask, request, redirect, url_for, render_template, abort\nimport store, governance as gov, agent_runtime as rt\nimport connection_manager as cmod\nimport catalog as cat\n\nWARDEN_VERSION = \"0.3\"\n\ndef _build_info():\n \"\"\"Increment a build number on each new deploy. Identity comes from RENDER_GIT_COMMIT\n if Render provides it, else a BUILD_ID baked into the image at build time (see\n Dockerfile), else \u0027local\u0027. The counter (persisted on the disk) bumps whenever that\n identity changes; a plain restart of the same build does not bump it.\"\"\"\n import json\n ident = os.environ.get(\"RENDER_GIT_COMMIT\", \"\")\n src = \"commit\"\n if not ident:\n try:\n ident = open(os.path.join(os.path.dirname(os.path.abspath(__file__)), \"BUILD_ID\")).read().strip()\n src = \"build\"\n except Exception:\n ident = \"\"\n meta_path = os.path.join(store.DATA_ROOT, \"build.json\")\n try:\n meta = json.load(open(meta_path))\n except Exception:\n meta = {}\n num = meta.get(\"build\", 0)\n if not ident or ident != meta.get(\"ident\"):\n num += 1\n try:\n json.dump({\"ident\": ident or \"local\", \"build\": num}, open(meta_path, \"w\"))\n except Exception:\n pass\n if src == \"commit\":\n label = ident[:7]\n elif src == \"build\":\n label = ident[:13] # e.g. 20260828T1912\n else:\n label = \"local\"\n return num, label\n\n_BUILD_NUM, BUILD_COMMIT = _build_info()\nVERSION_FULL = f\"{WARDEN_VERSION}.{_BUILD_NUM}\"\n\ntry:\n from zoneinfo import ZoneInfo\n _PT = ZoneInfo(\"America/Los_Angeles\")\nexcept Exception:\n _PT = datetime.timezone(datetime.timedelta(hours=-7), \"PDT\")\n# captured once at process start; on Render each deploy restarts the process\nDEPLOYED_AT = datetime.datetime.now(_PT).strftime(\"%Y-%m-%d %H:%M %Z\")\n\napp = Flask(__name__)\nstore.init()\n\ndef _env_specs():\n \"\"\"Servers to auto-connect on boot, from WARDEN_AUTOCONNECT (comma-separated catalog ids).\n Tokens resolve from each server\u0027s env var, so config survives redeploys.\"\"\"\n ids = [x.strip() for x in os.environ.get(\"WARDEN_AUTOCONNECT\", \"\").split(\",\") if x.strip()]\n return [{\"id\": i, \"transport\": cat.BY_ID[i][\"transport\"]} for i in ids if i in cat.BY_ID]\n\ndef cm():\n c = cmod.manager()\n c.ensure_started(store.enabled_connections() + _env_specs())\n return c\n\n@app.context_processor\ndef inject_globals():\n return {\"pending\": store.pending_approvals(), \"mode\": rt.mode(),\n \"version\": VERSION_FULL, \"commit\": BUILD_COMMIT, \"deployed_at\": DEPLOYED_AT}\n\ndef connected_tools():\n \"\"\"All tools across connected servers, with effective governance risk.\"\"\"\n ovr = store.all_overrides()\n out = []\n for t in cm().all_tools():\n m = gov.meta(t[\"key\"], t[\"tool\"], t[\"description\"], ovr.get(t[\"key\"]))\n out.append({**t, \"risk\": m[\"risk\"], \"gate\": m[\"gate\"], \"override\": ovr.get(t[\"key\"])})\n return out\n\ndef tools_by_server():\n groups = {}\n for t in connected_tools():\n groups.setdefault(t[\"server_id\"], {\"name\": t[\"server_name\"], \"tools\": []})\n groups[t[\"server_id\"]][\"tools\"].append(t)\n return groups\n\n@app.route(\"/\")\ndef home():\n servers = cm().connected_servers()\n return render_template(\"dashboard.html\", agents=store.list_agents(), runs=store.list_runs(12),\n pending=store.pending_approvals(), servers=servers)\n\n@app.route(\"/connections\")\ndef connections():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"connections.html\", catalog=cat.CATALOG, status=status,\n enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL,\n tools=connected_tools())\n\n@app.route(\"/connections/enable\", methods=[\"POST\"])\ndef enable_connection():\n cid = request.form.get(\"id\"); entry = cat.BY_ID.get(cid)\n if not entry: abort(404)\n transport = entry[\"transport\"]\n token = request.form.get(\"token\") or None\n command = request.form.get(\"command\") or None\n url = request.form.get(\"url\") or entry.get(\"run\")\n store.enable_connection(cid, transport, command=command, url=url, token=token)\n st = cm().connect_spec({\"id\": cid, \"transport\": transport, \"command\": command, \"url\": url, \"token\": token})\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True, \"status\": (st or {}).get(\"status\"), \"error\": (st or {}).get(\"error\"),\n \"tool_count\": (st or {}).get(\"tool_count\", 0)}\n return redirect(url_for(\"connections\"))\n\n@app.route(\"/connections/disable\", methods=[\"POST\"])\ndef disable_connection():\n cid = request.form.get(\"id\")\n store.disable_connection(cid); cm().disconnect(cid)\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True}\n return redirect(url_for(\"connections\"))\n\n@app.route(\"/tool-risk\", methods=[\"POST\"])\ndef tool_risk():\n store.set_override(request.form.get(\"key\"), request.form.get(\"risk\"))\n return redirect(request.form.get(\"back\") or url_for(\"connections\"))\n\n@app.route(\"/connlist\")\ndef connlist():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"_connlist.html\", catalog=cat.CATALOG, status=status,\n enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL)\n\n@app.route(\"/tools.json\")\ndef tools_json():\n groups = {}\n for t in connected_tools():\n g = groups.setdefault(t[\"server_id\"], {\"server\": t[\"server_name\"], \"tools\": []})\n g[\"tools\"].append({\"key\": t[\"key\"], \"tool\": t[\"tool\"], \"risk\": t[\"risk\"],\n \"gate\": t[\"gate\"], \"description\": t[\"description\"]})\n return {\"groups\": list(groups.values())}\n\n@app.route(\"/new\")\ndef new_agent():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"builder.html\", groups=tools_by_server(),\n catalog=cat.CATALOG, status=status, enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL)\n\n@app.route(\"/agents\", methods=[\"POST\"])\ndef create_agent():\n name = request.form.get(\"name\", \"\").strip() or \"Untitled agent\"\n instructions = request.form.get(\"instructions\", \"\").strip()\n model = request.form.get(\"model\", \"\").strip() or rt.MODEL_DEFAULT\n skills = request.form.getlist(\"skills\")\n aid = store.create_agent(name, instructions, model, skills)\n return redirect(url_for(\"agent\", aid=aid))\n\n@app.route(\"/agent/\u003caid\u003e\")\ndef agent(aid):\n ag = store.get_agent(aid)\n if not ag: abort(404)\n idx = {t[\"key\"]: t for t in connected_tools()}\n skills = ag[\"skills\"] or []\n granted = [idx[k] for k in skills if k in idx]\n groups = {}\n for t in granted:\n groups.setdefault(t[\"server_name\"], []).append(t)\n counts = {\"total\": len(granted), \"servers\": len(groups),\n \"gated\": sum(1 for t in granted if t[\"gate\"] == \"approval\"),\n \"auto\": sum(1 for t in granted if t[\"gate\"] != \"approval\")}\n missing = [k for k in skills if k not in idx]\n runs = [r for r in store.list_runs(50) if r[\"agent_id\"] == aid]\n return render_template(\"agent.html\", agent=ag, groups=groups, counts=counts,\n missing=missing, runs=runs)\n\ndef _advance_bg(rid):\n \"\"\"Run the agent loop in the background so the browser isn\u0027t blocked.\"\"\"\n def worker():\n try:\n rt.advance(rid)\n except Exception as e:\n try:\n r = store.get_run(rid)\n store.audit(rid, r[\"agent_id\"], \"error\", detail={\"text\": str(e)[:200]})\n store.update_run(rid, status=\"error\")\n except Exception:\n pass\n threading.Thread(target=worker, daemon=True).start()\n\n@app.route(\"/run\", methods=[\"POST\"])\ndef run():\n aid = request.form.get(\"agent_id\"); user_input = request.form.get(\"input\", \"\").strip()\n if not store.get_agent(aid) or not user_input: abort(400)\n rid = store.create_run(aid, user_input)\n _advance_bg(rid)\n return redirect(url_for(\"run_view\", rid=rid))\n\n@app.route(\"/run/\u003crid\u003e\")\ndef run_view(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n ag = store.get_agent(r[\"agent_id\"])\n return render_template(\"run.html\", run=r, agent=ag, audit=store.audit_for_run(rid),\n approvals=store.approvals_for_run(rid))\n\ndef _fmt_event(e):\n d = e.get(\"detail\") or {}\n kind = e[\"kind\"]\n if kind in (\"run_started\", \"user_message\"):\n text = d.get(\"input\") or d.get(\"text\") or \"\"\n elif kind in (\"final\", \"thought\", \"error\"):\n text = d.get(\"text\", \"\")\n elif \"result\" in d:\n text = \"-\u003e \" + _json.dumps(d[\"result\"])[:300]\n elif \"input\" in d:\n text = _json.dumps(d[\"input\"])[:300]\n else:\n text = \"\"\n return {\"ts\": (e[\"ts\"] or \"\")[11:19], \"kind\": kind, \"risk\": e.get(\"risk\"),\n \"tool\": (e[\"skill\"] or \"\").split(\"__\")[-1] if e.get(\"skill\") else \"\",\n \"text\": text}\n\n@app.route(\"/run/\u003crid\u003e/events\")\ndef run_events(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n audit = store.audit_for_run(rid)\n pend = [{\"id\": a[\"id\"], \"tool\": (a[\"skill\"] or \"\").split(\"__\")[-1], \"risk\": a[\"risk\"],\n \"input\": a[\"arguments\"].get(\"input\"),\n \"approve\": url_for(\"approval\", apid=a[\"id\"])}\n for a in store.approvals_for_run(rid) if a[\"status\"] == \"pending\"]\n return {\"status\": r[\"status\"], \"events\": [_fmt_event(e) for e in audit],\n \"pending\": pend, \"back\": url_for(\"run_view\", rid=rid)}\n\n@app.route(\"/run/\u003crid\u003e/say\", methods=[\"POST\"])\ndef run_say(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n if r[\"status\"] in (\"running\", \"awaiting_approval\"):\n return {\"error\": \"busy\"}, 409\n text = request.form.get(\"input\", \"\").strip()\n if not text:\n return {\"error\": \"empty\"}, 400\n tr = r[\"transcript\"]\n tr.append({\"role\": \"user\", \"content\": text})\n store.update_run(rid, status=\"running\", transcript=tr)\n store.audit(rid, r[\"agent_id\"], \"user_message\", detail={\"text\": text})\n _advance_bg(rid)\n return {\"ok\": True}\n\n@app.route(\"/approval/\u003capid\u003e\", methods=[\"POST\"])\ndef approval(apid):\n ap = store.get_approval(apid)\n if not ap: abort(404)\n decision = request.form.get(\"decision\")\n if decision in (\"approved\", \"denied\"):\n store.decide_approval(apid, decision, by=\"operator\"); _advance_bg(ap[\"run_id\"])\n # AJAX callers get JSON; form callers get a redirect\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True}\n return redirect(request.form.get(\"back\") or url_for(\"run_view\", rid=ap[\"run_id\"]))\n\n@app.route(\"/approvals\")\ndef approvals():\n return render_template(\"approvals.html\", pending=store.pending_approvals())\n\n@app.route(\"/architecture\")\ndef architecture():\n servers = cm().connected_servers()\n return render_template(\"architecture.html\", servers=servers, tools=connected_tools())\n\n@app.route(\"/audit\")\ndef audit():\n return render_template(\"audit.html\", events=store.audit_all(300))\n\n@app.route(\"/observability\")\ndef observability():\n from collections import Counter, defaultdict\n events = store.audit_all(4000)\n runs = store.list_runs(500)\n agents = {a[\"id\"]: a[\"name\"] for a in store.list_agents()}\n\n model_calls = [e for e in events if e[\"kind\"] == \"model_call\"]\n tool_calls = [e for e in events if e[\"kind\"] in (\"tool_result\", \"tool_result_gated\")]\n denied_ev = [e for e in events if e[\"kind\"] == \"denied\"]"run-fd49149d
08-29 05:29:40
tool result
LOW
read_source → "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re, time\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\n# Estimated USD price per 1M tokens (input, output). Approximate list prices, editable;\n# used only to estimate cost for the observability view. Matched by substring of model id.\nPRICES = {\"opus\": (15.0, 75.0), \"sonnet\": (3.0, 15.0), \"haiku\": (0.80, 4.0)}\ndef _cost(model, inp, out):\n rate = (3.0, 15.0)\n m = (model or \"\").lower()\n for k, v in PRICES.items():\n if k in m:\n rate = v; break\n return round(inp / 1e6 * rate[0] + out / 1e6 * rate[1], 6)\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _call_model(system, messages, tools):\n t0 = time.time()\n if SANDBOX:\n r = _sandbox_model(messages, tools)\n r[\"usage\"] = {\"input_tokens\": 0, \"output_tokens\": 0}\n r[\"model\"] = \"sandbox\"; r[\"latency_ms\"] = int((time.time() - t0) * 1000)\n return r\n import anthropic\n client = anthropic.Anthropic()\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=1024,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content],\n \"usage\": {\"input_tokens\": resp.usage.input_tokens, \"output_tokens\": resp.usage.output_tokens},\n \"model\": MODEL_DEFAULT, \"latency_ms\": int((time.time() - t0) * 1000)}\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n # Intent comes from the user\u0027s actual request, not the accumulating transcript,\n # so a completed refund doesn\u0027t spuriously trigger a file-write gate.\n user_text = \"\"\n for m in messages:\n if m.get(\"role\") == \"user\" and isinstance(m.get(\"content\"), str):\n user_text = m[\"content\"].lower(); break\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", user_text)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in user_text for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in user_text for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\nimport threading\n_run_locks = {} # run_id -\u003e Lock, serializes advance() per run\n_rerun = set() # run_ids asked to advance again while already advancing\n_guard = threading.Lock()\n\ndef advance(run_id):\n \"\"\"Serialize advancing a single run. Concurrent triggers (e.g. several approvals\n decided at once) must not run the loop on the same transcript in parallel, or the\n model gets an assistant turn whose tool_use blocks aren\u0027t all answered yet (API 400).\n Only one thread advances a run at a time; triggers that arrive mid-advance cause\n exactly one more pass afterward, so the latest decisions are always picked up.\"\"\"\n with _guard:\n lock = _run_locks.setdefault(run_id, threading.Lock())\n if lock.locked():\n _rerun.add(run_id) # someone is already advancing; ask them to loop\n return store.get_run(run_id)\n with lock:\n while True:\n result = _advance_once(run_id)\n with _guard:\n if run_id in _rerun:\n _rerun.discard(run_id)\n continue # a decision landed during the pass; go again\n break\n return result\n\ndef _advance_once(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n resp = _call_model(system, messages, tools)\n u = resp.get(\"usage\", {})\n store.audit(run_id, agent[\"id\"], \"model_call\",\n detail={\"model\": resp.get(\"model\"), \"input_tokens\": u.get(\"input_tokens\", 0),\n \"output_tokens\": u.get(\"output_tokens\", 0), \"latency_ms\": resp.get(\"latency_ms\", 0),\n \"cost\": _cost(resp.get(\"model\"), u.get(\"input_tokens\", 0), u.get(\"output_tokens\", 0))})\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \"thought\", detail={\"text\":blk[\"text\"]})\n if resp[\"stop_reason\"]!=\"tool_use\":\n final=\" \".join(b.get(\"text\",\"\") for b in resp[\"content\"] if b.get(\"type\")==\"text\").strip()\n store.audit(run_id, agent[\"id\"], \"final\", detail={\"text\":final})\n store.update_run(run_id, status=\"done\", transcript=messages)\n return store.get_run(run_id)\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\":\"loop bound reached\"})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n\ndef _has_tool_use(msg):\n return any(isinstance(b,dict) and b.get(\"type\")==\"tool_use\" for b in msg[\"content\"])\n\ndef _execute_tool_turn(run_id, agent, assistant_msg, messages, idx):\n blocks=[b for b in assistant_msg[\"content\"] if b.get(\"type\")==\"tool_use\"]\n for b in blocks:\n m = risk_for(b[\"name\"], idx)\n if m[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap is None:\n store.create_approval(run_id, agent[\"id\"], b[\"name\"], m[\"risk\"],\n {\"tool_use_id\":b[\"id\"],\"input\":b[\"input\"]})\n store.audit(run_id, agent[\"id\"], \"approval_request\", skill=b[\"name\"],\n risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n for b in blocks:\n if risk_for(b[\"name\"], idx)[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap and ap[\"status\"]==\"pending\":\n return \"paused\"\n results=[]\n for b in blocks:\n m=risk_for(b[\"name\"], idx); gated=m[\"gate\"]==\"approval\"\n ap=_approval_for(run_id,b[\"id\"]) if gated else None\n if gated and ap and ap[\"status\"]==\"denied\":\n rtext=json.dumps({\"denied\":True,\"note\":\"A human approver denied this action. Do not retry; explain and stop.\"})\n store.audit(run_id, agent[\"id\"], \"denied\", skill=b[\"name\"], risk=m[\"risk\"],\n detail={\"input\":b[\"input\"],\"outcome\":\"denied\"})\n else:\n t0=time.time()\n rtext=_cm().call_by_key(b[\"name\"], b[\"input\"])\n parsed=_safe(rtext)\n outcome=\"error\" if isinstance(parsed, dict) and parsed.get(\"error\") else \"ok\"\n store.audit(run_id, agent[\"id\"], \"tool_result_gated\" if gated else \"tool_result\",\n skill=b[\"name\"], risk=m[\"risk\"],\n detail={\"input\":b[\"input\"],\"result\":parsed,\n \"latency_ms\":int((time.time()-t0)*1000),\"outcome\":outcome})\n results.append({\"type\":\"tool_result\",\"tool_use_id\":b[\"id\"],\"content\":rtext})\n messages.append({\"role\":\"user\",\"content\":results})\n store.update_run(run_id, transcript=messages)\n return \"executed\"\n\ndef _approval_for(run_id, tool_use_id):\n for ap in store.approvals_for_run(run_id):\n if ap[\"arguments\"].get(\"tool_use_id\")==tool_use_id:\n return ap\n return None\n\ndef _safe(t):\n try: return json.loads(t)\n except Exception: return t\n"run-fd49149d
08-29 05:29:40
thought
Good โ the self-check is clean. Now let me read all 11 source files in parallel for a deep manual review.run-fd49149d
08-29 05:29:40
model call
{"cost": 0.056877, "input_tokens": 16284, "latency_ms": 6452, "model": "claude-sonnet-4-6", "output_tokens": 535}run-fd49149d
08-29 05:29:34
tool result
LOW
run_selfcheck → {"clean": true, "files_checked": 11, "findings": []}run-fd49149d
08-29 05:29:34
tool result
LOW
list_source → [{"file": "agent_runtime.py", "lines": 233}, {"file": "app.py", "lines": 378}, {"file": "catalog.py", "lines": 90}, {"file": "connection_manager.py", "lines": 197}, {"file": "governance.py", "lines": 64}, {"file": "mcp_code_server.py", "lines": 79}, {"file": "mcp_fs_server.py", "lines": 52}, {"file": "mcp_server.py", "lines": 95}, {"file": "paths.py", "lines": 32}, {"file": "store.py", "lines": 194}, {"file": "vault.py", "lines": 60}]run-fd49149d
08-29 05:29:34
thought
Sure! Let me start by gathering information about your Warden repo โ I'll read the source files and run a self-check simultaneously.run-fd49149d
08-29 05:29:34
model call
{"cost": 0.049218, "input_tokens": 15926, "latency_ms": 4691, "model": "claude-sonnet-4-6", "output_tokens": 96}run-fd49149d
08-29 04:35:56
Yesrun-36099a48
08-29 04:30:12
final
Good. Here's the exact problem, isolated:
---
## ๐ Bug 1 of 4 โ `agent_runtime.py`: `KeyError` in `_sandbox_model`
**Location:** `_sandbox_model()`, lines 64run-36099a48
08-29 04:30:12
thought
Good. Here's the exact problem, isolated:
---
## ๐ Bug 1 of 4 โ `agent_runtime.py`: `KeyError` in `_sandbox_model`
**Location:** `_sandbox_model()`, lines 64run-36099a48
08-29 04:29:57
tool result
LOW
read_source → "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _call_model(system, messages, tools):\n if SANDBOX:\n return _sandbox_model(messages, tools)\n import anthropic\n client = anthropic.Anthropic()\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=1024,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content]}\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", text_in)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in text_in for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in text_in for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\ndef advance(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n resp = _call_model(system, messages, tools)\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \"thought\", detail={\"text\":blk[\"text\"]})\n if resp[\"stop_reason\"]!=\"tool_use\":\n final=\" \".join(b.get(\"text\",\"\") for b in resp[\"content\"] if b.get(\"type\")==\"text\").strip()\n store.audit(run_id, agent[\"id\"], \"final\", detail={\"text\":final})\n store.update_run(run_id, status=\"done\", transcript=messages)\n return store.get_run(run_id)\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\":\"loop bound reached\"})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n\ndef _has_tool_use(msg):\n return any(isinstance(b,dict) and b.get(\"type\")==\"tool_use\" for b in msg[\"content\"])\n\ndef _execute_tool_turn(run_id, agent, assistant_msg, messages, idx):\n blocks=[b for b in assistant_msg[\"content\"] if b.get(\"type\")==\"tool_use\"]\n for b in blocks:\n m = risk_for(b[\"name\"], idx)\n if m[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap is None:\n store.create_approval(run_id, agent[\"id\"], b[\"name\"], m[\"risk\"],\n {\"tool_use_id\":b[\"id\"],\"input\":b[\"input\"]})\n store.audit(run_id, agent[\"id\"], \"approval_request\", skill=b[\"name\"],\n risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n for b in blocks:\n if risk_for(b[\"name\"], idx)[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap and ap[\"status\"]==\"pending\":\n return \"paused\"\n results=[]\n for b in blocks:\n m=risk_for(b[\"name\"], idx); gated=m[\"gate\"]==\"approval\"\n ap=_approval_for(run_id,b[\"id\"]) if gated else None\n if gated and ap and ap[\"status\"]==\"denied\":\n rtext=json.dumps({\"denied\":True,\"note\":\"A human approver denied this action. Do not retry; explain and stop.\"})\n store.audit(run_id, agent[\"id\"], \"denied\", skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n else:\n rtext=_cm().call_by_key(b[\"name\"], b[\"input\"])\n store.audit(run_id, agent[\"id\"], \"tool_result_gated\" if gated else \"tool_result\",\n skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"],\"result\":_safe(rtext)})\n results.append({\"type\":\"tool_result\",\"tool_use_id\":b[\"id\"],\"content\":rtext})\n messages.append({\"role\":\"user\",\"content\":results})\n store.update_run(run_id, transcript=messages)\n return \"executed\"\n\ndef _approval_for(run_id, tool_use_id):\n for ap in store.approvals_for_run(run_id):\n if ap[\"arguments\"].get(\"tool_use_id\")==tool_use_id:\n return ap\n return None\n\ndef _safe(t):\n try: return json.loads(t)\n except Exception: return t\n"run-36099a48
08-29 04:29:57
thought
Let's start with the most severe one first. Let me re-read the exact current state of `agent_runtime.py` before proposing anything.run-36099a48
08-29 04:29:54
One at a timerun-36099a48
08-29 04:29:36
final
No patches have been proposed yet โ we identified the bugs last session but didn't commit any fixes. Here's exactly where things stand and what's next:
---
##run-36099a48
08-29 04:29:36
thought
No patches have been proposed yet โ we identified the bugs last session but didn't commit any fixes. Here's exactly where things stand and what's next:
---
##run-36099a48
08-29 04:29:27
tool result
LOW
list_source → [{"file": "agent_runtime.py", "lines": 173}, {"file": "app.py", "lines": 308}, {"file": "catalog.py", "lines": 90}, {"file": "connection_manager.py", "lines": 197}, {"file": "governance.py", "lines": 64}, {"file": "mcp_code_server.py", "lines": 79}, {"file": "mcp_fs_server.py", "lines": 52}, {"file": "mcp_server.py", "lines": 95}, {"file": "paths.py", "lines": 32}, {"file": "store.py", "lines": 190}, {"file": "vault.py", "lines": 60}]run-36099a48
08-29 04:29:27
thought
Let me check the current state of things โ the open issues, PRs, and the proposed patches from our bug-finding session.run-36099a48
08-29 04:29:24
Whats next?run-36099a48
08-29 04:26:43
final
I now have a thorough picture of every file. Let me show you all the bugs found before submitting any fixes:
---
## ๐ Bugs Found
### Bug 1 โ `agent_runtime.prun-36099a48
08-29 04:26:43
thought
I now have a thorough picture of every file. Let me show you all the bugs found before submitting any fixes:
---
## ๐ Bugs Found
### Bug 1 โ `agent_runtime.prun-36099a48
08-29 04:26:19
tool result
LOW
read_source → "\"\"\"\nResolves where Warden stores its data. Prefers WARDEN_DATA_DIR (a mounted persistent\ndisk in production). If that path can\u0027t be created or written, it falls back to a\nwritable local directory instead of crashing the app, and records that it fell back so\nthe condition is visible on /healthz. A misconfigured disk should degrade to ephemeral,\nnever take the service down.\n\"\"\"\nimport os\n\n_APPDIR = os.path.dirname(os.path.abspath(__file__))\nREQUESTED = os.environ.get(\"WARDEN_DATA_DIR\") or _APPDIR\n\ndef _writable(path):\n try:\n os.makedirs(path, exist_ok=True)\n t = os.path.join(path, \".wtest\")\n with open(t, \"w\") as f:\n f.write(\"ok\")\n os.remove(t)\n return True\n except Exception:\n return False\n\ndef _resolve():\n if _writable(REQUESTED):\n return REQUESTED, False\n fallback = os.path.join(_APPDIR, \"_localdata\")\n if _writable(fallback):\n return fallback, True\n return _APPDIR, True\n\nDATA_ROOT, FALLBACK = _resolve()\n"run-36099a48
08-29 04:26:19
tool result
LOW
read_source → "\"\"\"\nEncrypts secrets (connection tokens) at rest so they never sit in the database as\nplaintext. The key comes from WARDEN_SECRET_KEY if set (kept out of the data dir, the\nstronger option); otherwise a key is generated once and stored on the persistent disk\nnext to the data, so encryption works with zero configuration. Any string works as\nWARDEN_SECRET_KEY, it is hashed into a valid key.\n\"\"\"\nimport os\nimport base64\nimport hashlib\n\ntry:\n from cryptography.fernet import Fernet\n _OK = True\nexcept Exception:\n _OK = False\n\nimport paths\nDATA_ROOT = paths.DATA_ROOT\nKEY_FILE = os.path.join(DATA_ROOT, \".warden_key\")\n\ndef _key():\n env = os.environ.get(\"WARDEN_SECRET_KEY\")\n if env:\n return base64.urlsafe_b64encode(hashlib.sha256(env.encode()).digest())\n if os.path.exists(KEY_FILE):\n return open(KEY_FILE, \"rb\").read().strip()\n k = Fernet.generate_key()\n os.makedirs(DATA_ROOT, exist_ok=True)\n open(KEY_FILE, \"wb\").write(k)\n try:\n os.chmod(KEY_FILE, 0o600)\n except Exception:\n pass\n return k\n\n_F = None\ndef _fernet():\n global _F\n if _F is None:\n _F = Fernet(_key())\n return _F\n\ndef encrypt(s):\n \"\"\"Return an \u0027enc:\u0027-prefixed ciphertext, or the original if encryption is unavailable.\"\"\"\n if not s or not _OK:\n return s\n try:\n return \"enc:\" + _fernet().encrypt(s.encode()).decode()\n except Exception:\n return s\n\ndef decrypt(s):\n \"\"\"Reverse encrypt(); passes through anything not marked \u0027enc:\u0027 (e.g. legacy plaintext).\"\"\"\n if not s or not _OK or not isinstance(s, str) or not s.startswith(\"enc:\"):\n return s\n try:\n return _fernet().decrypt(s[4:].encode()).decode()\n except Exception:\n return s\n"run-36099a48
08-29 04:26:19
tool result
LOW
read_source → "\"\"\"\nPersistence for Warden: agents, runs, the audit log, and the approvals queue.\nSQLite on disk. On a platform with an ephemeral filesystem (Render) this resets on\nredeploy, which is fine for a demo; the point is that within a session every agent\naction and every approval is durably recorded and queryable.\n\"\"\"\nimport os\nimport json\nimport sqlite3\nimport datetime\nimport uuid\n\nimport paths\nDATA_ROOT = paths.DATA_ROOT\nDB = os.path.join(DATA_ROOT, \"warden.db\")\n\ndef _conn():\n c = sqlite3.connect(DB, timeout=10)\n c.row_factory = sqlite3.Row\n try:\n c.execute(\"PRAGMA journal_mode=WAL\")\n c.execute(\"PRAGMA busy_timeout=8000\")\n except Exception:\n pass\n return c\n\ndef now():\n return datetime.datetime.now(datetime.timezone.utc).isoformat()\n\ndef _id(prefix):\n return f\"{prefix}-{uuid.uuid4().hex[:8]}\"\n\ndef init():\n c = _conn()\n c.executescript(\"\"\"\n CREATE TABLE IF NOT EXISTS agents(\n id TEXT PRIMARY KEY, name TEXT, instructions TEXT, model TEXT,\n skills TEXT, created_at TEXT);\n CREATE TABLE IF NOT EXISTS runs(\n id TEXT PRIMARY KEY, agent_id TEXT, input TEXT, status TEXT,\n transcript TEXT, created_at TEXT, updated_at TEXT);\n CREATE TABLE IF NOT EXISTS audit(\n id TEXT PRIMARY KEY, run_id TEXT, agent_id TEXT, ts TEXT,\n kind TEXT, skill TEXT, risk TEXT, detail TEXT);\n CREATE TABLE IF NOT EXISTS approvals(\n id TEXT PRIMARY KEY, run_id TEXT, agent_id TEXT, skill TEXT, risk TEXT,\n arguments TEXT, status TEXT, created_at TEXT, decided_at TEXT, decided_by TEXT);\n CREATE TABLE IF NOT EXISTS connections(\n id TEXT PRIMARY KEY, transport TEXT, command TEXT, url TEXT, token TEXT,\n enabled INTEGER, created_at TEXT);\n CREATE TABLE IF NOT EXISTS tool_overrides(\n model_key TEXT PRIMARY KEY, risk TEXT);\n \"\"\")\n c.commit(); c.close()\n\n# ---- agents ----\ndef create_agent(name, instructions, model, skills):\n c = _conn(); aid = _id(\"ag\")\n c.execute(\"INSERT INTO agents VALUES(?,?,?,?,?,?)\",\n (aid, name, instructions, model, json.dumps(skills), now()))\n c.commit(); c.close(); return aid\n\ndef get_agent(aid):\n c = _conn(); r = c.execute(\"SELECT * FROM agents WHERE id=?\", (aid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"skills\"] = json.loads(d[\"skills\"] or \"[]\"); return d\n\ndef list_agents():\n c = _conn(); rows = c.execute(\"SELECT * FROM agents ORDER BY created_at DESC\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"skills\"] = json.loads(d[\"skills\"] or \"[]\"); out.append(d)\n return out\n\n# ---- runs ----\ndef create_run(agent_id, user_input):\n c = _conn(); rid = _id(\"run\")\n c.execute(\"INSERT INTO runs VALUES(?,?,?,?,?,?,?)\",\n (rid, agent_id, user_input, \"running\", json.dumps([]), now(), now()))\n c.commit(); c.close(); return rid\n\ndef update_run(rid, status=None, transcript=None):\n c = _conn()\n if status is not None:\n c.execute(\"UPDATE runs SET status=?, updated_at=? WHERE id=?\", (status, now(), rid))\n if transcript is not None:\n c.execute(\"UPDATE runs SET transcript=?, updated_at=? WHERE id=?\",\n (json.dumps(transcript), now(), rid))\n c.commit(); c.close()\n\ndef get_run(rid):\n c = _conn(); r = c.execute(\"SELECT * FROM runs WHERE id=?\", (rid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"transcript\"] = json.loads(d[\"transcript\"] or \"[]\"); return d\n\ndef list_runs(limit=50):\n c = _conn(); rows = c.execute(\"SELECT * FROM runs ORDER BY created_at DESC LIMIT ?\", (limit,)).fetchall(); c.close()\n return [dict(r) for r in rows]\n\n# ---- audit ----\ndef audit(run_id, agent_id, kind, skill=None, risk=None, detail=None):\n c = _conn()\n c.execute(\"INSERT INTO audit VALUES(?,?,?,?,?,?,?,?)\",\n (_id(\"ev\"), run_id, agent_id, now(), kind, skill, risk,\n json.dumps(detail) if detail is not None else None))\n c.commit(); c.close()\n\ndef audit_for_run(run_id):\n c = _conn(); rows = c.execute(\"SELECT * FROM audit WHERE run_id=? ORDER BY ts\", (run_id,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"detail\"] = json.loads(d[\"detail\"]) if d[\"detail\"] else None; out.append(d)\n return out\n\ndef audit_all(limit=200):\n c = _conn(); rows = c.execute(\"SELECT * FROM audit ORDER BY ts DESC LIMIT ?\", (limit,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"detail\"] = json.loads(d[\"detail\"]) if d[\"detail\"] else None; out.append(d)\n return out\n\n# ---- approvals ----\ndef create_approval(run_id, agent_id, skill, risk, arguments):\n c = _conn(); apid = _id(\"ap\")\n c.execute(\"INSERT INTO approvals VALUES(?,?,?,?,?,?,?,?,?,?)\",\n (apid, run_id, agent_id, skill, risk, json.dumps(arguments),\n \"pending\", now(), None, None))\n c.commit(); c.close(); return apid\n\ndef get_approval(apid):\n c = _conn(); r = c.execute(\"SELECT * FROM approvals WHERE id=?\", (apid,)).fetchone(); c.close()\n if not r: return None\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); return d\n\ndef decide_approval(apid, status, by=\"operator\"):\n c = _conn()\n c.execute(\"UPDATE approvals SET status=?, decided_at=?, decided_by=? WHERE id=?\",\n (status, now(), by, apid))\n c.commit(); c.close()\n\ndef pending_approvals():\n c = _conn(); rows = c.execute(\"SELECT * FROM approvals WHERE status=\u0027pending\u0027 ORDER BY created_at\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); out.append(d)\n return out\n\ndef approvals_for_run(run_id):\n c = _conn(); rows = c.execute(\"SELECT * FROM approvals WHERE run_id=? ORDER BY created_at\", (run_id,)).fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r); d[\"arguments\"] = json.loads(d[\"arguments\"] or \"{}\"); out.append(d)\n return out\n\n# ---- connections (enabled external MCP servers) ----\ndef enable_connection(cid, transport, command=None, url=None, token=None):\n import vault\n c = _conn()\n c.execute(\"INSERT OR REPLACE INTO connections VALUES(?,?,?,?,?,?,?)\",\n (cid, transport, command, url, vault.encrypt(token), 1, now()))\n c.commit(); c.close()\n\ndef disable_connection(cid):\n c = _conn(); c.execute(\"DELETE FROM connections WHERE id=?\", (cid,)); c.commit(); c.close()\n\ndef enabled_connections():\n import vault\n c = _conn(); rows = c.execute(\"SELECT * FROM connections WHERE enabled=1\").fetchall(); c.close()\n out = []\n for r in rows:\n d = dict(r)\n out.append({\"id\": d[\"id\"], \"transport\": d[\"transport\"], \"command\": d[\"command\"],\n \"url\": d[\"url\"], \"token\": vault.decrypt(d[\"token\"])})\n return out\n\ndef is_enabled(cid):\n c = _conn(); r = c.execute(\"SELECT 1 FROM connections WHERE id=? AND enabled=1\", (cid,)).fetchone(); c.close()\n return bool(r)\n\n# ---- tool risk overrides ----\ndef set_override(model_key, risk):\n c = _conn(); c.execute(\"INSERT OR REPLACE INTO tool_overrides VALUES(?,?)\", (model_key, risk)); c.commit(); c.close()\n\ndef get_override(model_key):\n c = _conn(); r = c.execute(\"SELECT risk FROM tool_overrides WHERE model_key=?\", (model_key,)).fetchone(); c.close()\n return r[\"risk\"] if r else None\n\ndef all_overrides():\n c = _conn(); rows = c.execute(\"SELECT * FROM tool_overrides\").fetchall(); c.close()\n return {r[\"model_key\"]: r[\"risk\"] for r in rows}\n"run-36099a48
08-29 04:26:19
tool result
LOW
read_source → "\"\"\"\nWarden MCP server: a real Model Context Protocol server exposing a small set of\nenterprise-flavored tools. The Warden runtime connects to this as an MCP client,\ndiscovers these tools over the protocol, and invokes them.\n\nTools deliberately span read and write so the governance layer has something to\ngovern: reads are low risk and auto-execute; writes change state and are the ones\nthe studio gates behind human approval.\n\nRun standalone for a protocol smoke test: python mcp_server.py\n(but normally it is spawned over stdio by the runtime\u0027s MCP client)\n\"\"\"\nimport json\nimport os\nimport datetime\nfrom mcp.server.fastmcp import FastMCP\n\nimport paths\nDATA_DIR = os.path.join(paths.DATA_ROOT, \"data\")\nos.makedirs(DATA_DIR, exist_ok=True)\n\nCUSTOMERS = {\n \"AC-1001\": {\"account_id\": \"AC-1001\", \"name\": \"Rivera Logistics\", \"plan\": \"Enterprise\",\n \"mrr\": 4200, \"status\": \"active\", \"last_charge\": 4200,\n \"notes\": \"Charged twice on 2026-08-03 due to a billing retry bug.\"},\n \"AC-1002\": {\"account_id\": \"AC-1002\", \"name\": \"Halcyon Health\", \"plan\": \"Premium\",\n \"mrr\": 1800, \"status\": \"active\", \"last_charge\": 1800, \"notes\": \"\"},\n \"AC-1003\": {\"account_id\": \"AC-1003\", \"name\": \"Meridian Foods\", \"plan\": \"Standard\",\n \"mrr\": 600, \"status\": \"past_due\", \"last_charge\": 0,\n \"notes\": \"Payment failed twice this month.\"},\n}\n\nKB = [\n {\"id\": \"kb-01\", \"title\": \"Refund policy\",\n \"body\": \"Duplicate charges are refunded in full once verified against the billing ledger. \"\n \"Refunds above 1000 require a human approver.\"},\n {\"id\": \"kb-02\", \"title\": \"Past-due accounts\",\n \"body\": \"Do not issue credits on past-due accounts until the balance is cleared. \"\n \"Open a billing ticket instead.\"},\n {\"id\": \"kb-03\", \"title\": \"Escalation\",\n \"body\": \"Anything touching money movement is a governed action and must be logged.\"},\n]\n\ndef _ledger_path(name):\n return os.path.join(DATA_DIR, name)\n\ndef _append(name, row):\n path = _ledger_path(name)\n rows = []\n if os.path.exists(path):\n rows = json.load(open(path))\n row[\"at\"] = datetime.datetime.now(datetime.timezone.utc).isoformat()\n rows.append(row)\n json.dump(rows, open(path, \"w\"), indent=2)\n return row\n\nmcp = FastMCP(\"warden-enterprise-tools\")\n\n@mcp.tool()\ndef lookup_customer(account_id: str) -\u003e str:\n \"\"\"Look up an enterprise customer account by id (e.g. AC-1001). Read only.\"\"\"\n rec = CUSTOMERS.get(account_id.strip().upper())\n if not rec:\n return json.dumps({\"error\": f\"no account {account_id}\"})\n return json.dumps(rec)\n\n@mcp.tool()\ndef search_knowledge(query: str) -\u003e str:\n \"\"\"Search the internal knowledge base for policy and guidance. Read only.\"\"\"\n q = query.lower()\n hits = [a for a in KB if q in a[\"title\"].lower() or q in a[\"body\"].lower()]\n if not hits:\n hits = KB # fall back to returning all short KB so the agent has context\n return json.dumps(hits)\n\n@mcp.tool()\ndef create_ticket(subject: str, body: str) -\u003e str:\n \"\"\"Open an internal billing/support ticket. Write action (changes state).\"\"\"\n row = _append(\"tickets.json\", {\"id\": f\"TK-{abs(hash(subject)) % 9000 + 1000}\",\n \"subject\": subject, \"body\": body})\n return json.dumps({\"created\": True, \"ticket\": row})\n\n@mcp.tool()\ndef issue_refund(account_id: str, amount: float, reason: str = \"\") -\u003e str:\n \"\"\"Issue a monetary refund to a customer account. High-impact write action:\n moves money, so the studio gates this behind human approval before it runs.\"\"\"\n acct = account_id.strip().upper()\n if acct not in CUSTOMERS:\n return json.dumps({\"error\": f\"no account {account_id}\"})\n row = _append(\"refunds.json\", {\"id\": f\"RF-{abs(hash(acct + str(amount))) % 9000 + 1000}\",\n \"account_id\": acct, \"amount\": float(amount), \"reason\": reason})\n return json.dumps({\"refunded\": True, \"refund\": row})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-36099a48
08-29 04:26:19
tool result
LOW
read_source → "\"\"\"\nWarden\u0027s built-in Filesystem MCP server. A real second MCP server (local, Python) so\nmulti-server connection and routing is demonstrable without needing Node. All access\nis confined to a sandbox workspace directory.\n\"\"\"\nimport os\nimport json\nfrom mcp.server.fastmcp import FastMCP\n\nimport paths\nWORKSPACE = os.path.join(paths.DATA_ROOT, \"data\", \"workspace\")\nos.makedirs(WORKSPACE, exist_ok=True)\n\n# seed one file so reads have something to find\n_seed = os.path.join(WORKSPACE, \"welcome.txt\")\nif not os.path.exists(_seed):\n open(_seed, \"w\").write(\"Warden workspace. Files here are readable and writable by agents, under governance.\\n\")\n\ndef _safe(path):\n p = os.path.abspath(os.path.join(WORKSPACE, path.lstrip(\"/\")))\n if not p.startswith(os.path.abspath(WORKSPACE)):\n raise ValueError(\"path escapes workspace\")\n return p\n\nmcp = FastMCP(\"warden-filesystem\")\n\n@mcp.tool()\ndef list_files(subdir: str = \"\") -\u003e str:\n \"\"\"List files in the workspace (or a subdirectory). Read only.\"\"\"\n base = _safe(subdir)\n if not os.path.exists(base):\n return json.dumps({\"error\": \"no such path\"})\n return json.dumps(sorted(os.listdir(base)))\n\n@mcp.tool()\ndef read_file(path: str) -\u003e str:\n \"\"\"Read a text file from the workspace. Read only.\"\"\"\n p = _safe(path)\n if not os.path.isfile(p):\n return json.dumps({\"error\": f\"no file {path}\"})\n return open(p, encoding=\"utf-8\", errors=\"replace\").read()[:8000]\n\n@mcp.tool()\ndef write_file(path: str, content: str) -\u003e str:\n \"\"\"Create or overwrite a text file in the workspace. Write action (changes state).\"\"\"\n p = _safe(path)\n os.makedirs(os.path.dirname(p), exist_ok=True)\n open(p, \"w\", encoding=\"utf-8\").write(content)\n return json.dumps({\"written\": True, \"path\": path, \"bytes\": len(content)})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-36099a48
08-29 04:26:19
tool result
LOW
read_source → "\"\"\"\nWarden self-audit MCP server. Exposes Warden\u0027s own source code to an agent so a\n\"Warden Engineer\" agent can inspect the running codebase, run a real static self-check,\nand propose fixes. Reads are safe and auto-run; proposing a patch is a gated write that\nlands in a review folder (it never overwrites the running source).\n\"\"\"\nimport os, json, ast, py_compile, tempfile\nfrom mcp.server.fastmcp import FastMCP\n\nAPP_DIR = os.path.dirname(os.path.abspath(__file__))\nimport paths\nPATCH_DIR = os.path.join(paths.DATA_ROOT, \"data\", \"patches\")\nos.makedirs(PATCH_DIR, exist_ok=True)\n\ndef _py_files():\n return sorted(f for f in os.listdir(APP_DIR) if f.endswith(\".py\"))\n\nmcp = FastMCP(\"warden-self-audit\")\n\n@mcp.tool()\ndef list_source() -\u003e str:\n \"\"\"List Warden\u0027s own Python source files with line counts. Read only.\"\"\"\n out = []\n for f in _py_files():\n n = sum(1 for _ in open(os.path.join(APP_DIR, f), encoding=\"utf-8\", errors=\"replace\"))\n out.append({\"file\": f, \"lines\": n})\n return json.dumps(out)\n\n@mcp.tool()\ndef read_source(filename: str) -\u003e str:\n \"\"\"Read one of Warden\u0027s own source files. Read only.\"\"\"\n if filename not in _py_files():\n return json.dumps({\"error\": f\"no source file {filename}\", \"available\": _py_files()})\n return open(os.path.join(APP_DIR, filename), encoding=\"utf-8\", errors=\"replace\").read()[:12000]\n\n@mcp.tool()\ndef run_selfcheck() -\u003e str:\n \"\"\"Statically check every Warden source file: byte-compile for syntax errors and\n AST-scan for bare excepts and TODO/FIXME markers. Runs real checks; no side effects.\"\"\"\n findings = []\n for f in _py_files():\n path = os.path.join(APP_DIR, f)\n try:\n py_compile.compile(path, doraise=True)\n except py_compile.PyCompileError as e:\n findings.append({\"file\": f, \"kind\": \"syntax_error\", \"detail\": str(e).splitlines()[-1][:160]})\n continue\n src = open(path, encoding=\"utf-8\", errors=\"replace\").read()\n try:\n tree = ast.parse(src)\n for node in ast.walk(tree):\n if isinstance(node, ast.ExceptHandler) and node.type is None:\n findings.append({\"file\": f, \"kind\": \"bare_except\", \"line\": node.lineno})\n except Exception:\n pass\n try:\n import tokenize, io\n for tok in tokenize.generate_tokens(io.StringIO(src).readline):\n if tok.type == tokenize.COMMENT and (\"TODO\" in tok.string.upper() or \"FIXME\" in tok.string.upper()):\n findings.append({\"file\": f, \"kind\": \"todo\", \"line\": tok.start[0], \"detail\": tok.string.strip()[:120]})\n except Exception:\n pass\n return json.dumps({\"files_checked\": len(_py_files()),\n \"findings\": findings, \"clean\": len(findings) == 0})\n\n@mcp.tool()\ndef propose_patch(filename: str, new_content: str, rationale: str = \"\") -\u003e str:\n \"\"\"Propose a fix for a source file. Gated write: saves the proposed version to a\n review folder for a human to inspect and apply. Never edits the running source.\"\"\"\n safe = os.path.basename(filename)\n out = os.path.join(PATCH_DIR, safe)\n open(out, \"w\", encoding=\"utf-8\").write(new_content)\n meta = os.path.join(PATCH_DIR, safe + \".rationale.txt\")\n open(meta, \"w\", encoding=\"utf-8\").write(rationale or \"(none)\")\n return json.dumps({\"proposed\": True, \"review_path\": f\"data/patches/{safe}\",\n \"bytes\": len(new_content), \"note\": \"saved for human review; running source unchanged\"})\n\nif __name__ == \"__main__\":\n mcp.run()\n"run-36099a48
08-29 04:26:14
tool result
LOW
read_source → "\"\"\"\nThe governance layer. Warden\u0027s point of view: tools can do things; governance decides\nwhich run on their own and which pause for a human. Built-in enterprise tools have a\nhand-set risk registry. Tools discovered from external MCP servers are classified\nautomatically, fail-closed: reads run, writes and anything unrecognized are gated.\nAn operator can override any tool\u0027s risk.\n\"\"\"\nSKILLS = {\n \"lookup_customer\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"search_knowledge\":{\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"create_ticket\": {\"risk\":\"MED\",\"gate\":\"auto\",\"kind\":\"write\"},\n \"issue_refund\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n \"list_files\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"read_file\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"write_file\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n \"list_source\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"read_source\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"run_selfcheck\": {\"risk\":\"LOW\",\"gate\":\"auto\",\"kind\":\"read\"},\n \"propose_patch\": {\"risk\":\"HIGH\",\"gate\":\"approval\",\"kind\":\"write\"},\n}\n\nREAD_HINTS = (\"get\",\"list\",\"read\",\"search\",\"lookup\",\"fetch\",\"find\",\"query\",\"view\",\n \"describe\",\"show\",\"count\",\"status\",\"summary\",\"recent\",\"ask\",\"explore\",\n \"inspect\",\"check\")\nWRITE_HINTS = (\"create\",\"write\",\"update\",\"delete\",\"remove\",\"issue\",\"send\",\"post\",\"add\",\n \"set\",\"merge\",\"close\",\"open\",\"deploy\",\"execute\",\"run\",\"refund\",\"cancel\",\n \"approve\",\"edit\",\"upload\",\"move\",\"rename\",\"revoke\",\"grant\",\"pay\",\"charge\")\n\ndef classify(name, desc=\"\"):\n \"\"\"Heuristic risk for an external tool, judged by its leading verb.\n Fail-closed: unrecognized =\u003e HIGH.\"\"\"\n import re\n n = (name or \"\").lower().replace(\"-\", \"_\")\n lead = re.match(r\"[a-z]+\", n)\n lead = lead.group(0) if lead else n # leading alpha run (handles camelCase)\n first = n.split(\"_\")[0] # first snake_case token\n for h in READ_HINTS:\n if first == h or lead.startswith(h):\n return \"LOW\"\n for h in WRITE_HINTS:\n if first == h or lead.startswith(h):\n return \"HIGH\"\n return \"HIGH\"\n\ndef meta(model_key, tool_name, desc=\"\", override=None):\n \"\"\"Resolve effective governance for a tool. Precedence: override \u003e known registry \u003e classify.\"\"\"\n if override in (\"LOW\",\"MED\",\"HIGH\"):\n risk = override\n elif tool_name in SKILLS:\n risk = SKILLS[tool_name][\"risk\"]\n else:\n risk = classify(tool_name, desc)\n return {\"risk\": risk, \"gate\": \"approval\" if risk == \"HIGH\" else \"auto\"}\n\n# convenience wrappers used where only a bare name is available (built-ins)\ndef skill_meta(name):\n m = SKILLS.get(name)\n if m: return m\n r = classify(name)\n return {\"risk\": r, \"gate\": \"approval\" if r==\"HIGH\" else \"auto\", \"kind\":\"?\"}\ndef requires_approval(name):\n return skill_meta(name)[\"gate\"] == \"approval\"\ndef risk_of(name):\n return skill_meta(name)[\"risk\"]\n"run-36099a48
08-29 04:26:14
tool result
LOW
read_source → "\"\"\"\nMulti-server MCP connection manager.\n\nOwns one dedicated asyncio loop thread. Tool discovery is done once at connect time and\ncached, so listing tools never re-hits a server. For execution: stdio/builtin servers\nkeep a persistent session (subprocess spawn is expensive); HTTP servers open a fresh\nshort-lived session per call, entirely within one coroutine, because the streamable-HTTP\ntransport binds its cancel scope to the creating task and cannot be reused across tasks.\n\"\"\"\nimport os, sys, asyncio, threading, warnings, shlex\nfrom contextlib import asynccontextmanager\nfrom mcp import ClientSession, StdioServerParameters\nfrom mcp.client.stdio import stdio_client\ntry:\n from mcp.client.streamable_http import streamablehttp_client\n _HTTP_OK = True\nexcept Exception:\n _HTTP_OK = False\n\nimport catalog as catalog_mod\n\nHERE = os.path.dirname(os.path.abspath(__file__))\nBUILTINS = {\n \"builtin_enterprise\": [sys.executable, os.path.join(HERE, \"mcp_server.py\")],\n \"builtin_files\": [sys.executable, os.path.join(HERE, \"mcp_fs_server.py\")],\n \"builtin_code\": [sys.executable, os.path.join(HERE, \"mcp_code_server.py\")],\n}\n\nclass _LoopThread:\n def __init__(self):\n self.loop = asyncio.new_event_loop()\n with warnings.catch_warnings():\n warnings.simplefilter(\"ignore\")\n try:\n w = asyncio.ThreadedChildWatcher(); w.attach_loop(self.loop)\n asyncio.set_child_watcher(w)\n except Exception:\n pass\n threading.Thread(target=self._run, daemon=True).start()\n def _run(self):\n asyncio.set_event_loop(self.loop); self.loop.run_forever()\n def run(self, coro, timeout=None):\n return asyncio.run_coroutine_threadsafe(coro, self.loop).result(timeout)\n\ndef _http_params(sid, spec):\n cat = catalog_mod.BY_ID.get(sid, {})\n url = spec.get(\"url\") or cat.get(\"run\")\n headers = {}\n tok = spec.get(\"token\")\n if not tok and cat.get(\"env\"): # fall back to an environment variable\n tok = os.environ.get(cat[\"env\"])\n if tok:\n headers[\"Authorization\"] = tok if tok.lower().startswith(\"bearer\") else f\"Bearer {tok}\"\n return url, (headers or None)\n\ndef _stdio_params(sid, spec, transport):\n if transport == \"builtin\":\n cmd = BUILTINS[sid]\n else:\n run = spec.get(\"command\") or catalog_mod.BY_ID.get(sid, {}).get(\"run\", \"\")\n cmd = shlex.split(run)\n if not cmd:\n raise RuntimeError(\"no command configured\")\n return StdioServerParameters(command=cmd[0], args=cmd[1:], env=os.environ.copy())\n\n@asynccontextmanager\nasync def _http_session(sid, spec):\n url, headers = _http_params(sid, spec)\n async with streamablehttp_client(url, headers=headers) as (read, write, _):\n async with ClientSession(read, write) as session:\n await session.initialize()\n yield session\n\nclass _Manager:\n def __init__(self):\n self._lt = _LoopThread()\n self._lock = threading.Lock()\n self._sessions = {} # sid -\u003e persistent session (stdio/builtin only)\n self._keep = {} # sid -\u003e [context managers] to keep alive\n self._http = {} # sid -\u003e spec (http servers, fresh session per call)\n self._toolcache = {} # sid -\u003e [ {name,description,input_schema} ] (from connect)\n self._status = {} # sid -\u003e status dict\n self._toolmap = {} # model_name -\u003e (sid, tool)\n self._started = False\n\n def ensure_started(self, enabled_specs=None):\n with self._lock:\n if not self._started:\n for sid in BUILTINS:\n self._connect(sid, {\"id\": sid, \"transport\": \"builtin\"})\n self._started = True\n for spec in (enabled_specs or []):\n if spec[\"id\"] not in self._status:\n self._connect(spec[\"id\"], spec)\n self._rebuild_toolmap()\n\n def connect_spec(self, spec):\n with self._lock:\n self._connect(spec[\"id\"], spec); self._rebuild_toolmap()\n return self._status.get(spec[\"id\"])\n\n def disconnect(self, sid):\n with self._lock:\n self._sessions.pop(sid, None); self._keep.pop(sid, None)\n self._http.pop(sid, None); self._toolcache.pop(sid, None)\n self._status.pop(sid, None); self._rebuild_toolmap()\n\n # ---- connect (on loop thread) ----\n def _connect(self, sid, spec):\n cat = catalog_mod.BY_ID.get(sid, {})\n name = cat.get(\"name\", sid)\n transport = spec.get(\"transport\") or cat.get(\"transport\", \"stdio_node\")\n try:\n tools = self._lt.run(self._aopen(sid, spec, transport), timeout=75)\n self._toolcache[sid] = tools\n self._status[sid] = {\"status\": \"connected\", \"error\": None, \"name\": name,\n \"transport\": transport, \"tool_count\": len(tools)}\n except Exception as e:\n msg = (str(e) or e.__class__.__name__)\n low = msg.lower()\n if transport != \"http\" and (\"closed\" in low or \"exit\" in low or \"broken pipe\" in low or not msg.strip()):\n msg = (\"server exited on startup, it likely needs credentials or configuration. \"\n \"Edit the command to supply them (e.g. a real connection string or token).\")\n self._status[sid] = {\"status\": \"error\", \"error\": msg[:220],\n \"name\": name, \"transport\": transport, \"tool_count\": 0}\n\n async def _aopen(self, sid, spec, transport):\n if transport == \"http\":\n if not _HTTP_OK:\n raise RuntimeError(\"HTTP transport unavailable in this build\")\n self._http[sid] = spec\n async with _http_session(sid, spec) as session: # validate + discover, same task\n resp = await session.list_tools()\n return _tools(resp)\n # stdio / builtin: persistent session\n params = _stdio_params(sid, spec, transport)\n cm = stdio_client(params)\n read, write = await cm.__aenter__()\n sess_cm = ClientSession(read, write)\n session = await sess_cm.__aenter__()\n await session.initialize()\n self._keep[sid] = [cm, sess_cm]; self._sessions[sid] = session\n resp = await session.list_tools()\n return _tools(resp)\n\n async def _acall(self, sid, tool, args):\n if sid in self._http:\n async with _http_session(sid, self._http[sid]) as session: # fresh, same task\n result = await session.call_tool(tool, args or {})\n return _text(result)\n result = await self._sessions[sid].call_tool(tool, args or {})\n return _text(result)\n\n def _rebuild_toolmap(self):\n self._toolmap = {}\n for sid, tools in self._toolcache.items():\n for t in tools:\n self._toolmap[f\"{sid}__{t[\u0027name\u0027]}\"[:64]] = (sid, t[\"name\"])\n\n # ---- queries (use cache; no live calls) ----\n def connected_servers(self):\n return [dict(id=sid, **self._status[sid]) for sid in self._status]\n\n def all_tools(self):\n out = []\n with self._lock:\n for sid, tools in self._toolcache.items():\n sname = self._status.get(sid, {}).get(\"name\", sid)\n for t in tools:\n key = f\"{sid}__{t[\u0027name\u0027]}\"[:64]\n self._toolmap[key] = (sid, t[\"name\"])\n out.append({\"key\": key, \"server_id\": sid, \"server_name\": sname,\n \"tool\": t[\"name\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n def call_by_key(self, key, args):\n sid, tool = self._toolmap.get(key, (None, None))\n if sid is None:\n return \u0027{\"error\":\"unknown tool \u0027 + str(key) + \u0027\"}\u0027\n with self._lock:\n return self._lt.run(self._acall(sid, tool, args), timeout=90)\n\ndef _tools(resp):\n return [{\"name\": t.name, \"description\": t.description or \"\", \"input_schema\": t.inputSchema}\n for t in resp.tools]\ndef _text(result):\n return \"\\n\".join(c.text for c in result.content if getattr(c, \"type\", None) == \"text\")\n\n_CM = None\n_CM_LOCK = threading.Lock()\ndef manager():\n global _CM\n with _CM_LOCK:\n if _CM is None:\n _CM = _Manager()\n return _CM\n"run-36099a48
08-29 04:26:14
tool result
LOW
read_source → "\"\"\"\nThe connections catalog: a curated directory of common enterprise MCP servers.\nAccurate as of mid-2026. Each entry records who maintains it, how it connects\n(transport), what credentials it needs, and a default governance posture.\n\n\u0027maintainer\u0027: official = Anthropic reference (educational), vendor = product owner,\n community = third party, warden = ships built in with this app.\n\u0027transport\u0027: stdio_python (uvx ...), stdio_node (npx ...), http (remote OAuth URL),\n builtin (a local server this app runs itself).\n\u0027status\u0027: ready = connectable in this app\u0027s runtime now,\n needs_node = requires npx/Node in the runtime,\n needs_python = requires uvx/uv in the runtime,\n remote = a hosted URL you paste in (works without local runtime),\n archived = still works but no longer maintained upstream.\nRisk posture is a starting point; Warden classifies each discovered tool and lets\nyou override it.\n\"\"\"\n\nCATALOG = [\n # --- ships with Warden (always connectable) ---\n {\"id\":\"builtin_enterprise\",\"name\":\"Enterprise Tools (Warden)\",\"category\":\"Reference\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Customer lookup, knowledge search, ticketing, and refunds. The built-in demo server.\"},\n {\"id\":\"builtin_files\",\"name\":\"Filesystem (Warden)\",\"category\":\"Files \u0026 Docs\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Read, list, and write files inside a sandboxed workspace. A working second server.\"},\n {\"id\":\"builtin_code\",\"name\":\"Self-Audit (Warden)\",\"category\":\"Dev\",\n \"maintainer\":\"warden\",\"transport\":\"builtin\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Reads Warden\u0027s own source, runs a static self-check, and proposes fixes (gated). Warden debugging Warden.\"},\n\n # --- live public remote server, no credentials, connect and go ---\n {\"id\":\"deepwiki\",\"name\":\"DeepWiki (GitHub repos)\",\"category\":\"Dev\",\"maintainer\":\"vendor\",\n \"transport\":\"http\",\"run\":\"https://mcp.deepwiki.com/mcp\",\"auth\":\"none\",\"status\":\"ready\",\n \"desc\":\"Ask real questions about any public GitHub repository and read its docs, live. No token needed.\"},\n\n # --- Anthropic official reference servers ---\n {\"id\":\"fetch\",\"name\":\"Fetch\",\"category\":\"Web\",\"maintainer\":\"official\",\"transport\":\"stdio_python\",\n \"run\":\"uvx mcp-server-fetch\",\"auth\":\"none\",\"status\":\"needs_python\",\n \"desc\":\"Fetch a URL and return its content as text for the agent to read.\"},\n {\"id\":\"filesystem\",\"name\":\"Filesystem (official)\",\"category\":\"Files \u0026 Docs\",\"maintainer\":\"official\",\n \"transport\":\"stdio_node\",\"run\":\"npx -y @modelcontextprotocol/server-filesystem \u003cpath\u003e\",\n \"auth\":\"none\",\"status\":\"needs_node\",\"desc\":\"Reference filesystem server. Read and write within allowed paths.\"},\n {\"id\":\"git\",\"name\":\"Git\",\"category\":\"Dev\",\"maintainer\":\"official\",\"transport\":\"stdio_python\",\n \"run\":\"uvx mcp-server-git --repository \u003cpath\u003e\",\"auth\":\"none\",\"status\":\"needs_python\",\n \"desc\":\"Read a repo: status, diff, log, branches, and commits.\"},\n {\"id\":\"memory\",\"name\":\"Memory\",\"category\":\"Reference\",\"maintainer\":\"official\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-memory\",\"auth\":\"none\",\"status\":\"needs_node\",\n \"desc\":\"A simple knowledge-graph memory the agent can write to and recall.\"},\n\n # --- vendor-maintained (the right pick for production) ---\n {\"id\":\"github\",\"env\":\"GITHUB_TOKEN\",\"name\":\"GitHub\",\"category\":\"Dev\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://api.githubcopilot.com/mcp/\",\"auth\":\"oauth_or_pat\",\"status\":\"remote\",\n \"desc\":\"Read repos and issues; create issues, branches, and pull requests. Hosted OAuth endpoint.\"},\n {\"id\":\"linear\",\"env\":\"LINEAR_API_KEY\",\"name\":\"Linear\",\"category\":\"Product\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.linear.app/sse\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Issue tracking and project planning. Read issues and create or update them.\"},\n {\"id\":\"notion\",\"env\":\"NOTION_TOKEN\",\"name\":\"Notion\",\"category\":\"Knowledge\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.notion.com/mcp\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Read and write Notion docs and databases. A common agent knowledge base.\"},\n {\"id\":\"stripe\",\"env\":\"STRIPE_API_KEY\",\"name\":\"Stripe\",\"category\":\"Payments\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.stripe.com\",\"auth\":\"api_key\",\"status\":\"remote\",\n \"desc\":\"Look up customers, invoices, and payments; issue refunds. High-impact by nature.\"},\n {\"id\":\"sentry\",\"env\":\"SENTRY_TOKEN\",\"name\":\"Sentry\",\"category\":\"Observability\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://mcp.sentry.dev/mcp\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Read issues and errors; triage and resolve. Hosted OAuth endpoint.\"},\n {\"id\":\"supabase\",\"env\":\"SUPABASE_ACCESS_TOKEN\",\"name\":\"Supabase\",\"category\":\"Data\",\"maintainer\":\"vendor\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @supabase/mcp-server-supabase\",\"auth\":\"api_key\",\"status\":\"needs_node\",\n \"desc\":\"Query and manage a Supabase Postgres project, tables, and rows.\"},\n {\"id\":\"playwright\",\"name\":\"Playwright\",\"category\":\"Web\",\"maintainer\":\"vendor\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @playwright/mcp\",\"auth\":\"none\",\"status\":\"needs_node\",\n \"desc\":\"Drive a real browser: navigate, click, fill forms, extract. Microsoft-maintained.\"},\n {\"id\":\"cloudflare\",\"env\":\"CLOUDFLARE_TOKEN\",\"name\":\"Cloudflare\",\"category\":\"Infra\",\"maintainer\":\"vendor\",\"transport\":\"http\",\n \"run\":\"https://observability.mcp.cloudflare.com/sse\",\"auth\":\"oauth\",\"status\":\"remote\",\n \"desc\":\"Inspect and manage Cloudflare resources over a hosted OAuth connection.\"},\n\n # --- data (archived reference, still functional) ---\n {\"id\":\"postgres\",\"name\":\"PostgreSQL\",\"category\":\"Data\",\"maintainer\":\"community\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-postgres postgresql://\u003cconn\u003e\",\"auth\":\"conn_string\",\n \"status\":\"archived\",\"desc\":\"Query a Postgres database. Reference server is archived; use read-only creds.\"},\n {\"id\":\"slack\",\"name\":\"Slack\",\"category\":\"Comms\",\"maintainer\":\"community\",\"transport\":\"stdio_node\",\n \"run\":\"npx -y @modelcontextprotocol/server-slack\",\"auth\":\"bot_token\",\"status\":\"archived\",\n \"desc\":\"Read channels and post messages. Original reference archived; community builds exist.\"},\n]\n\nBY_ID = {c[\"id\"]: c for c in CATALOG}\n\nMAINTAINER_LABEL = {\"warden\":\"Built in\",\"official\":\"Anthropic reference\",\n \"vendor\":\"Vendor-maintained\",\"community\":\"Community\"}\nSTATUS_LABEL = {\"ready\":\"Ready\",\"needs_node\":\"Needs Node runtime\",\"needs_python\":\"Needs Python runtime\",\n \"remote\":\"Remote (paste URL + token)\",\"archived\":\"Archived but works\"}\n"run-36099a48
08-29 04:26:14
tool result
LOW
read_source → "\"\"\"\nWarden, an enterprise AI agent studio where every agent is governed by default.\nConnect MCP servers, build an agent from their tools, run it against a live model, and\ngate high-risk actions behind human approval with a full audit trail.\n\"\"\"\nimport os\nimport datetime\nimport threading\nimport json as _json\nfrom flask import Flask, request, redirect, url_for, render_template, abort\nimport store, governance as gov, agent_runtime as rt\nimport connection_manager as cmod\nimport catalog as cat\n\nWARDEN_VERSION = \"0.3\"\n\ndef _build_info():\n \"\"\"Increment a build number on each new deploy. Identity comes from RENDER_GIT_COMMIT\n if Render provides it, else a BUILD_ID baked into the image at build time (see\n Dockerfile), else \u0027local\u0027. The counter (persisted on the disk) bumps whenever that\n identity changes; a plain restart of the same build does not bump it.\"\"\"\n import json\n ident = os.environ.get(\"RENDER_GIT_COMMIT\", \"\")\n src = \"commit\"\n if not ident:\n try:\n ident = open(os.path.join(os.path.dirname(os.path.abspath(__file__)), \"BUILD_ID\")).read().strip()\n src = \"build\"\n except Exception:\n ident = \"\"\n meta_path = os.path.join(store.DATA_ROOT, \"build.json\")\n try:\n meta = json.load(open(meta_path))\n except Exception:\n meta = {}\n num = meta.get(\"build\", 0)\n if not ident or ident != meta.get(\"ident\"):\n num += 1\n try:\n json.dump({\"ident\": ident or \"local\", \"build\": num}, open(meta_path, \"w\"))\n except Exception:\n pass\n if src == \"commit\":\n label = ident[:7]\n elif src == \"build\":\n label = ident[:13] # e.g. 20260828T1912\n else:\n label = \"local\"\n return num, label\n\n_BUILD_NUM, BUILD_COMMIT = _build_info()\nVERSION_FULL = f\"{WARDEN_VERSION}.{_BUILD_NUM}\"\n\ntry:\n from zoneinfo import ZoneInfo\n _PT = ZoneInfo(\"America/Los_Angeles\")\nexcept Exception:\n _PT = datetime.timezone(datetime.timedelta(hours=-7), \"PDT\")\n# captured once at process start; on Render each deploy restarts the process\nDEPLOYED_AT = datetime.datetime.now(_PT).strftime(\"%Y-%m-%d %H:%M %Z\")\n\napp = Flask(__name__)\nstore.init()\n\ndef _env_specs():\n \"\"\"Servers to auto-connect on boot, from WARDEN_AUTOCONNECT (comma-separated catalog ids).\n Tokens resolve from each server\u0027s env var, so config survives redeploys.\"\"\"\n ids = [x.strip() for x in os.environ.get(\"WARDEN_AUTOCONNECT\", \"\").split(\",\") if x.strip()]\n return [{\"id\": i, \"transport\": cat.BY_ID[i][\"transport\"]} for i in ids if i in cat.BY_ID]\n\ndef cm():\n c = cmod.manager()\n c.ensure_started(store.enabled_connections() + _env_specs())\n return c\n\n@app.context_processor\ndef inject_globals():\n return {\"pending\": store.pending_approvals(), \"mode\": rt.mode(),\n \"version\": VERSION_FULL, \"commit\": BUILD_COMMIT, \"deployed_at\": DEPLOYED_AT}\n\ndef connected_tools():\n \"\"\"All tools across connected servers, with effective governance risk.\"\"\"\n ovr = store.all_overrides()\n out = []\n for t in cm().all_tools():\n m = gov.meta(t[\"key\"], t[\"tool\"], t[\"description\"], ovr.get(t[\"key\"]))\n out.append({**t, \"risk\": m[\"risk\"], \"gate\": m[\"gate\"], \"override\": ovr.get(t[\"key\"])})\n return out\n\ndef tools_by_server():\n groups = {}\n for t in connected_tools():\n groups.setdefault(t[\"server_id\"], {\"name\": t[\"server_name\"], \"tools\": []})\n groups[t[\"server_id\"]][\"tools\"].append(t)\n return groups\n\n@app.route(\"/\")\ndef home():\n servers = cm().connected_servers()\n return render_template(\"dashboard.html\", agents=store.list_agents(), runs=store.list_runs(12),\n pending=store.pending_approvals(), servers=servers)\n\n@app.route(\"/connections\")\ndef connections():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"connections.html\", catalog=cat.CATALOG, status=status,\n enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL,\n tools=connected_tools())\n\n@app.route(\"/connections/enable\", methods=[\"POST\"])\ndef enable_connection():\n cid = request.form.get(\"id\"); entry = cat.BY_ID.get(cid)\n if not entry: abort(404)\n transport = entry[\"transport\"]\n token = request.form.get(\"token\") or None\n command = request.form.get(\"command\") or None\n url = request.form.get(\"url\") or entry.get(\"run\")\n store.enable_connection(cid, transport, command=command, url=url, token=token)\n st = cm().connect_spec({\"id\": cid, \"transport\": transport, \"command\": command, \"url\": url, \"token\": token})\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True, \"status\": (st or {}).get(\"status\"), \"error\": (st or {}).get(\"error\"),\n \"tool_count\": (st or {}).get(\"tool_count\", 0)}\n return redirect(url_for(\"connections\"))\n\n@app.route(\"/connections/disable\", methods=[\"POST\"])\ndef disable_connection():\n cid = request.form.get(\"id\")\n store.disable_connection(cid); cm().disconnect(cid)\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True}\n return redirect(url_for(\"connections\"))\n\n@app.route(\"/tool-risk\", methods=[\"POST\"])\ndef tool_risk():\n store.set_override(request.form.get(\"key\"), request.form.get(\"risk\"))\n return redirect(request.form.get(\"back\") or url_for(\"connections\"))\n\n@app.route(\"/connlist\")\ndef connlist():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"_connlist.html\", catalog=cat.CATALOG, status=status,\n enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL)\n\n@app.route(\"/tools.json\")\ndef tools_json():\n groups = {}\n for t in connected_tools():\n g = groups.setdefault(t[\"server_id\"], {\"server\": t[\"server_name\"], \"tools\": []})\n g[\"tools\"].append({\"key\": t[\"key\"], \"tool\": t[\"tool\"], \"risk\": t[\"risk\"],\n \"gate\": t[\"gate\"], \"description\": t[\"description\"]})\n return {\"groups\": list(groups.values())}\n\n@app.route(\"/new\")\ndef new_agent():\n status = {s[\"id\"]: s for s in cm().connected_servers()}\n return render_template(\"builder.html\", groups=tools_by_server(),\n catalog=cat.CATALOG, status=status, enabled={c[\"id\"] for c in store.enabled_connections()},\n mlabel=cat.MAINTAINER_LABEL, slabel=cat.STATUS_LABEL)\n\n@app.route(\"/agents\", methods=[\"POST\"])\ndef create_agent():\n name = request.form.get(\"name\", \"\").strip() or \"Untitled agent\"\n instructions = request.form.get(\"instructions\", \"\").strip()\n model = request.form.get(\"model\", \"\").strip() or rt.MODEL_DEFAULT\n skills = request.form.getlist(\"skills\")\n aid = store.create_agent(name, instructions, model, skills)\n return redirect(url_for(\"agent\", aid=aid))\n\n@app.route(\"/agent/\u003caid\u003e\")\ndef agent(aid):\n ag = store.get_agent(aid)\n if not ag: abort(404)\n idx = {t[\"key\"]: t for t in connected_tools()}\n skills = ag[\"skills\"] or []\n granted = [idx[k] for k in skills if k in idx]\n groups = {}\n for t in granted:\n groups.setdefault(t[\"server_name\"], []).append(t)\n counts = {\"total\": len(granted), \"servers\": len(groups),\n \"gated\": sum(1 for t in granted if t[\"gate\"] == \"approval\"),\n \"auto\": sum(1 for t in granted if t[\"gate\"] != \"approval\")}\n missing = [k for k in skills if k not in idx]\n runs = [r for r in store.list_runs(50) if r[\"agent_id\"] == aid]\n return render_template(\"agent.html\", agent=ag, groups=groups, counts=counts,\n missing=missing, runs=runs)\n\ndef _advance_bg(rid):\n \"\"\"Run the agent loop in the background so the browser isn\u0027t blocked.\"\"\"\n def worker():\n try:\n rt.advance(rid)\n except Exception as e:\n try:\n r = store.get_run(rid)\n store.audit(rid, r[\"agent_id\"], \"error\", detail={\"text\": str(e)[:200]})\n store.update_run(rid, status=\"error\")\n except Exception:\n pass\n threading.Thread(target=worker, daemon=True).start()\n\n@app.route(\"/run\", methods=[\"POST\"])\ndef run():\n aid = request.form.get(\"agent_id\"); user_input = request.form.get(\"input\", \"\").strip()\n if not store.get_agent(aid) or not user_input: abort(400)\n rid = store.create_run(aid, user_input)\n _advance_bg(rid)\n return redirect(url_for(\"run_view\", rid=rid))\n\n@app.route(\"/run/\u003crid\u003e\")\ndef run_view(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n ag = store.get_agent(r[\"agent_id\"])\n return render_template(\"run.html\", run=r, agent=ag, audit=store.audit_for_run(rid),\n approvals=store.approvals_for_run(rid))\n\ndef _fmt_event(e):\n d = e.get(\"detail\") or {}\n kind = e[\"kind\"]\n if kind in (\"run_started\", \"user_message\"):\n text = d.get(\"input\") or d.get(\"text\") or \"\"\n elif kind in (\"final\", \"thought\", \"error\"):\n text = d.get(\"text\", \"\")\n elif \"result\" in d:\n text = \"-\u003e \" + _json.dumps(d[\"result\"])[:300]\n elif \"input\" in d:\n text = _json.dumps(d[\"input\"])[:300]\n else:\n text = \"\"\n return {\"ts\": (e[\"ts\"] or \"\")[11:19], \"kind\": kind, \"risk\": e.get(\"risk\"),\n \"tool\": (e[\"skill\"] or \"\").split(\"__\")[-1] if e.get(\"skill\") else \"\",\n \"text\": text}\n\n@app.route(\"/run/\u003crid\u003e/events\")\ndef run_events(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n audit = store.audit_for_run(rid)\n pend = [{\"id\": a[\"id\"], \"tool\": (a[\"skill\"] or \"\").split(\"__\")[-1], \"risk\": a[\"risk\"],\n \"input\": a[\"arguments\"].get(\"input\"),\n \"approve\": url_for(\"approval\", apid=a[\"id\"])}\n for a in store.approvals_for_run(rid) if a[\"status\"] == \"pending\"]\n return {\"status\": r[\"status\"], \"events\": [_fmt_event(e) for e in audit],\n \"pending\": pend, \"back\": url_for(\"run_view\", rid=rid)}\n\n@app.route(\"/run/\u003crid\u003e/say\", methods=[\"POST\"])\ndef run_say(rid):\n r = store.get_run(rid)\n if not r: abort(404)\n if r[\"status\"] in (\"running\", \"awaiting_approval\"):\n return {\"error\": \"busy\"}, 409\n text = request.form.get(\"input\", \"\").strip()\n if not text:\n return {\"error\": \"empty\"}, 400\n tr = r[\"transcript\"]\n tr.append({\"role\": \"user\", \"content\": text})\n store.update_run(rid, status=\"running\", transcript=tr)\n store.audit(rid, r[\"agent_id\"], \"user_message\", detail={\"text\": text})\n _advance_bg(rid)\n return {\"ok\": True}\n\n@app.route(\"/approval/\u003capid\u003e\", methods=[\"POST\"])\ndef approval(apid):\n ap = store.get_approval(apid)\n if not ap: abort(404)\n decision = request.form.get(\"decision\")\n if decision in (\"approved\", \"denied\"):\n store.decide_approval(apid, decision, by=\"operator\"); _advance_bg(ap[\"run_id\"])\n # AJAX callers get JSON; form callers get a redirect\n if request.headers.get(\"X-Requested-With\") == \"fetch\":\n return {\"ok\": True}\n return redirect(request.form.get(\"back\") or url_for(\"run_view\", rid=ap[\"run_id\"]))\n\n@app.route(\"/approvals\")\ndef approvals():\n return render_template(\"approvals.html\", pending=store.pending_approvals())\n\n@app.route(\"/architecture\")\ndef architecture():\n servers = cm().connected_servers()\n return render_template(\"architecture.html\", servers=servers, tools=connected_tools())\n\n@app.route(\"/audit\")\ndef audit():\n return render_template(\"audit.html\", events=store.audit_all(300))\n\n@app.route(\"/healthz\")\ndef healthz():\n import paths\n dd = store.DATA_ROOT\n return {\"ok\": True, \"mode\": rt.mode(),\n \"servers\": len(cm().connected_servers()),\n \"version\": VERSION_FULL, \"commit\": BUILD_COMMIT,\n \"persistence\": {\n \"WARDEN_DATA_DIR_env\": os.environ.get(\"WARDEN_DATA_DIR\", \"(unset)\"),\n \"requested_dir\": paths.REQUESTED,\n \"data_dir\": dd,\n "run-36099a48
08-29 04:26:14
tool result
LOW
read_source → "\"\"\"\nThe Warden runtime. Runs an agent in a perceive -\u003e decide -\u003e act loop against a live\nAnthropic model, invoking tools across one or more connected MCP servers via the\nconnection manager. Governance is enforced here: each tool\u0027s risk is resolved\n(override \u003e known registry \u003e auto-classification) and high-risk tools pause the run\nfor human approval. Live model calls when ANTHROPIC_API_KEY is set; otherwise a\ndeterministic sandbox planner drives the same flow offline.\n\"\"\"\nimport os, json, re\nimport connection_manager as cmod\nimport governance as gov\nimport store\n\nMODEL_DEFAULT = os.environ.get(\"WARDEN_MODEL\", \"claude-sonnet-4-5\")\nSANDBOX = not bool(os.environ.get(\"ANTHROPIC_API_KEY\"))\n\ndef mode():\n return \"sandbox\" if SANDBOX else \"live\"\n\ndef _cm():\n cm = cmod.manager()\n cm.ensure_started(store.enabled_connections())\n return cm\n\ndef tool_index():\n \"\"\"model_key -\u003e {tool, desc, server} across all connected servers.\"\"\"\n idx = {}\n for t in _cm().all_tools():\n idx[t[\"key\"]] = {\"tool\": t[\"tool\"], \"desc\": t[\"description\"], \"server\": t[\"server_name\"]}\n return idx\n\ndef risk_for(key, idx=None):\n idx = idx or tool_index()\n info = idx.get(key, {\"tool\": key, \"desc\": \"\"})\n return gov.meta(key, info[\"tool\"], info[\"desc\"], store.get_override(key))\n\ndef tools_for(agent):\n allowed = set(agent.get(\"skills\") or [])\n out = []\n for t in _cm().all_tools():\n if t[\"key\"] in allowed:\n out.append({\"name\": t[\"key\"], \"description\": t[\"description\"],\n \"input_schema\": t[\"input_schema\"]})\n return out\n\n# ---- model dispatch ----\ndef _call_model(system, messages, tools):\n if SANDBOX:\n return _sandbox_model(messages, tools)\n import anthropic\n client = anthropic.Anthropic()\n resp = client.messages.create(model=MODEL_DEFAULT, max_tokens=1024,\n system=system, messages=messages, tools=tools)\n return {\"stop_reason\": resp.stop_reason, \"content\": [_b2d(b) for b in resp.content]}\n\ndef _b2d(b):\n if b.type == \"text\": return {\"type\": \"text\", \"text\": b.text}\n if b.type == \"tool_use\": return {\"type\": \"tool_use\", \"id\": b.id, \"name\": b.name, \"input\": b.input}\n return {\"type\": b.type}\n\n# ---- sandbox planner (offline). Emits the same message shapes, using real tool keys. ----\ndef _find_key(tools, bare):\n for t in tools:\n if t[\"name\"].endswith(\"__\" + bare) or t[\"name\"] == bare:\n return t[\"name\"]\n return None\n\ndef _sandbox_model(messages, tools):\n text_in = json.dumps(messages).lower()\n called = set()\n for m in messages:\n for blk in (m.get(\"content\") or []) if isinstance(m.get(\"content\"), list) else []:\n if isinstance(blk, dict) and blk.get(\"type\") == \"tool_use\":\n called.add(blk[\"name\"])\n mm = re.search(r\"ac-?\\d{4}\", text_in)\n acct = (\"AC-\" + mm.group(0)[-4:]) if mm else None\n money = any(w in text_in for w in [\"refund\",\"charged twice\",\"double charge\",\"duplicate\",\"make it right\",\"money back\"])\n def tu(key, inp):\n return {\"stop_reason\":\"tool_use\",\"content\":[{\"type\":\"tool_use\",\"id\":\"sbx_\"+key,\"name\":key,\"input\":inp}]}\n k_lookup=_find_key(tools,\"lookup_customer\"); k_kb=_find_key(tools,\"search_knowledge\")\n k_refund=_find_key(tools,\"issue_refund\")\n if k_lookup and k_lookup not in called and acct: return tu(k_lookup,{\"account_id\":acct})\n if k_kb and k_kb not in called and money: return tu(k_kb,{\"query\":\"refund policy\"})\n if k_refund and k_refund not in called and money and acct:\n return tu(k_refund,{\"account_id\":acct,\"amount\":4200,\"reason\":\"duplicate charge verified against ledger\"})\n # filesystem scenario\n k_list=_find_key(tools,\"list_files\"); k_write=_find_key(tools,\"write_file\")\n wants_file = any(w in text_in for w in [\"file\",\"note\",\"write\",\"summary\",\"save\"])\n if k_list and k_list not in called and wants_file:\n return tu(k_list,{\"subdir\":\"\"})\n if k_write and k_write not in called and wants_file:\n fn=re.search(r\"([\\w\\-/]+\\.\\w{1,5})\", text_in)\n name=fn.group(1) if fn else \"note.txt\"\n return tu(k_write,{\"path\":name,\"content\":\"Written by a Warden agent after human approval.\"})\n final=\"[sandbox] Done. \"\n was_gated = any((\"issue_refund\" in c or \"write_file\" in c) for c in called)\n final += \"Gated action executed after human approval; every step is in the audit log.\" if was_gated else \"Reviewed and no gated action was required.\"\n return {\"stop_reason\":\"end_turn\",\"content\":[{\"type\":\"text\",\"text\":final}]}\n\n# ---- the loop ----\ndef advance(run_id):\n run = store.get_run(run_id); agent = store.get_agent(run[\"agent_id\"])\n system = (agent[\"instructions\"] or \"\") + \\\n \"\\n\\nYou operate under Warden governance. High-impact actions may require human \" \\\n \"approval before they execute; use the tools available and Warden gates what needs a human.\"\n tools = tools_for(agent); idx = tool_index(); messages = run[\"transcript\"]\n if not messages:\n messages = [{\"role\":\"user\",\"content\":run[\"input\"]}]\n store.audit(run_id, agent[\"id\"], \"run_started\", detail={\"input\":run[\"input\"],\"mode\":mode()})\n for _ in range(12):\n last = messages[-1] if messages else None\n if last and last[\"role\"]==\"assistant\" and _has_tool_use(last):\n if _execute_tool_turn(run_id, agent, last, messages, idx) == \"paused\":\n store.update_run(run_id, status=\"awaiting_approval\", transcript=messages)\n return store.get_run(run_id)\n resp = _call_model(system, messages, tools)\n messages.append({\"role\":\"assistant\",\"content\":resp[\"content\"]})\n for blk in resp[\"content\"]:\n if blk.get(\"type\")==\"text\" and blk.get(\"text\"):\n store.audit(run_id, agent[\"id\"], \"thought\", detail={\"text\":blk[\"text\"]})\n if resp[\"stop_reason\"]!=\"tool_use\":\n final=\" \".join(b.get(\"text\",\"\") for b in resp[\"content\"] if b.get(\"type\")==\"text\").strip()\n store.audit(run_id, agent[\"id\"], \"final\", detail={\"text\":final})\n store.update_run(run_id, status=\"done\", transcript=messages)\n return store.get_run(run_id)\n store.audit(run_id, agent[\"id\"], \"error\", detail={\"text\":\"loop bound reached\"})\n store.update_run(run_id, status=\"error\", transcript=messages)\n return store.get_run(run_id)\n\ndef _has_tool_use(msg):\n return any(isinstance(b,dict) and b.get(\"type\")==\"tool_use\" for b in msg[\"content\"])\n\ndef _execute_tool_turn(run_id, agent, assistant_msg, messages, idx):\n blocks=[b for b in assistant_msg[\"content\"] if b.get(\"type\")==\"tool_use\"]\n for b in blocks:\n m = risk_for(b[\"name\"], idx)\n if m[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap is None:\n store.create_approval(run_id, agent[\"id\"], b[\"name\"], m[\"risk\"],\n {\"tool_use_id\":b[\"id\"],\"input\":b[\"input\"]})\n store.audit(run_id, agent[\"id\"], \"approval_request\", skill=b[\"name\"],\n risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n for b in blocks:\n if risk_for(b[\"name\"], idx)[\"gate\"]==\"approval\":\n ap=_approval_for(run_id,b[\"id\"])\n if ap and ap[\"status\"]==\"pending\":\n return \"paused\"\n results=[]\n for b in blocks:\n m=risk_for(b[\"name\"], idx); gated=m[\"gate\"]==\"approval\"\n ap=_approval_for(run_id,b[\"id\"]) if gated else None\n if gated and ap and ap[\"status\"]==\"denied\":\n rtext=json.dumps({\"denied\":True,\"note\":\"A human approver denied this action. Do not retry; explain and stop.\"})\n store.audit(run_id, agent[\"id\"], \"denied\", skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"]})\n else:\n rtext=_cm().call_by_key(b[\"name\"], b[\"input\"])\n store.audit(run_id, agent[\"id\"], \"tool_result_gated\" if gated else \"tool_result\",\n skill=b[\"name\"], risk=m[\"risk\"], detail={\"input\":b[\"input\"],\"result\":_safe(rtext)})\n results.append({\"type\":\"tool_result\",\"tool_use_id\":b[\"id\"],\"content\":rtext})\n messages.append({\"role\":\"user\",\"content\":results})\n store.update_run(run_id, transcript=messages)\n return \"executed\"\n\ndef _approval_for(run_id, tool_use_id):\n for ap in store.approvals_for_run(run_id):\n if ap[\"arguments\"].get(\"tool_use_id\")==tool_use_id:\n return ap\n return None\n\ndef _safe(t):\n try: return json.loads(t)\n except Exception: return t\n"run-36099a48
08-29 04:26:14
thought
Good โ the automated self-check is clean (no syntax errors, no bare excepts, no TODOs). Now let me do a deep read of every source file to catch logic-level bugsrun-36099a48
08-29 04:26:10
tool result
LOW
list_source → [{"file": "agent_runtime.py", "lines": 173}, {"file": "app.py", "lines": 308}, {"file": "catalog.py", "lines": 90}, {"file": "connection_manager.py", "lines": 197}, {"file": "governance.py", "lines": 64}, {"file": "mcp_code_server.py", "lines": 79}, {"file": "mcp_fs_server.py", "lines": 52}, {"file": "mcp_server.py", "lines": 95}, {"file": "paths.py", "lines": 32}, {"file": "store.py", "lines": 190}, {"file": "vault.py", "lines": 60}]run-36099a48
08-29 04:26:10
tool result
LOW
run_selfcheck → {"clean": true, "files_checked": 11, "findings": []}run-36099a48
08-29 04:26:10
thought
Sure! Let me start by running the self-check and listing the source files simultaneously.run-36099a48
08-29 04:26:03
run started
"Find bugs, then fix them but show me the problem before committing the fix."run-36099a48