diff --git a/.github/workflows/ci.yml b/.github/workflows/ci.yml index b0094f8..539cd59 100644 --- a/.github/workflows/ci.yml +++ b/.github/workflows/ci.yml @@ -39,6 +39,8 @@ jobs: run: npm run typecheck - name: Unit tests run: npm test + - name: OpenCode plugin regression tests + run: node --test bridge/tests/opencode-plugin.test.mjs - name: Build run: npm run build diff --git a/README.md b/README.md index 7f8c4a1..500e810 100644 --- a/README.md +++ b/README.md @@ -1,6 +1,6 @@
-cot. Self-hosted observability for Claude Code, Cursor, and Codex +cot. Self-hosted observability for Claude Code, Cursor, Codex, and OpenCode
@@ -37,7 +37,7 @@ One command: 1. starts the collector in Docker, bound to localhost 2. installs the local bridge -3. wires up hooks for the agents you pick: **Claude Code**, **Cursor**, or **Codex** +3. wires up the agents you pick: **Claude Code**, **Cursor**, **Codex**, or **OpenCode** Cot runs in the background and collects traces while you build. Open the dashboard at **[http://127.0.0.1:31337](http://127.0.0.1:31337)**. @@ -47,6 +47,13 @@ Cot runs in the background and collects traces while you build. Open the dashboa The default port is **31337**. If it is busy, the installer picks the next free port and saves the URL in `~/.cot/config.json`. +OpenCode 1.18.29+ and OpenCode 2 use one local plugin file. The bridge installer +places it at `~/.config/opencode/plugins/cot.js`. To add it manually after +installing the cot bridge, download +`http://127.0.0.1:31337/opencode-plugin.js` (using your chosen port), save it +at that path, and restart +OpenCode. OpenCode sessions recorded before plugin installation are not imported. + | Command | What it does | | --- | --- | | `cot up` | Start or resume the collector | @@ -57,7 +64,7 @@ The default port is **31337**. If it is busy, the installer picks the next free ```mermaid flowchart LR - A["Claude Code
Cursor
Codex"] -- hook fires --> B["~/.cot/bin/cot
local bridge"] + A["Claude Code
Cursor
Codex
OpenCode"] -- hook/plugin event --> B["~/.cot/bin/cot
local bridge"] B -- "POST /v1/ingest" --> C["Collector
FastAPI · 127.0.0.1:31337"] C --> D[("~/.cot/cot.db
SQLite")] D --> E["Dashboard
React"] diff --git a/backend/app/db.py b/backend/app/db.py index 9e3477e..05b0896 100644 --- a/backend/app/db.py +++ b/backend/app/db.py @@ -1130,7 +1130,7 @@ def _recategorize_other_tools(conn: sqlite3.Connection) -> None: except json.JSONDecodeError: continue cat = categorize(row["source"], row["hook"], raw, row["tool"]) - if cat["category"] == "other": + if cat["category"] == "other" and not (row["source"] == "opencode" and cat["status"] == "error"): continue conn.execute( "UPDATE events SET category=?, title=?, detail=?, target=?, status=?," @@ -1158,7 +1158,7 @@ def _question_response_obj(detail: Any) -> dict[str, Any] | None: return obj if isinstance(obj, dict) else None -_MIGRATIONS_VERSION = "9" +_MIGRATIONS_VERSION = "10" _RAW_PAYLOAD_MAX_BYTES = 64 * 1024 @@ -2048,6 +2048,15 @@ def record_event( (norm["cwd"], sid), ) + if norm["source"] == "opencode" and raw.get("parent_session_id"): + parent = str(raw["parent_session_id"]) + if parent != sid: + conn.execute( + "UPDATE sessions SET parent_session_id = ?," + " subagent_label = COALESCE(subagent_label, ?) WHERE id = ?", + (parent, raw.get("session_title"), sid), + ) + # Imported events carry historical timestamps; ensure the session's # started_at reflects the earliest event we've seen. if origin == "import": diff --git a/backend/app/insights.py b/backend/app/insights.py index bfe2ffd..16c89a0 100644 --- a/backend/app/insights.py +++ b/backend/app/insights.py @@ -744,7 +744,7 @@ def _trend_anomaly(snap: Snapshot) -> list[dict[str, Any]]: # Paths agents legitimately edit outside the project (configs, temp, cot itself). _OUT_OF_CWD_EXCLUDES = re.compile( - r"/\.(?:claude|cursor|codex|config|cot|cache)(?:/|$)|^/(?:private/)?tmp(?:/|$)|^/var/folders/" + r"/\.(?:claude|cursor|codex|opencode|config|cot|cache)(?:/|$)|^/(?:private/)?tmp(?:/|$)|^/var/folders/" ) diff --git a/backend/app/main.py b/backend/app/main.py index 8aaca03..80908bf 100644 --- a/backend/app/main.py +++ b/backend/app/main.py @@ -39,7 +39,7 @@ app = FastAPI(title="cot collector", version=__version__) -AGENTS = ("claude", "cursor", "codex") +AGENTS = ("claude", "cursor", "codex", "opencode") EXPECTED_HOOKS = { "claude": [ "SessionStart", @@ -79,6 +79,11 @@ "SubagentStop", "Stop", ], + "opencode": [ + "SessionStart", "UserPromptSubmit", "PreToolUse", "PostToolUse", + "PostToolUseFailure", "afterAgentResponse", "afterAgentThought", + "OpenCodeUsage", "PermissionRequest", "PostCompact", "Stop", + ], } HOOK_LABELS = { "SessionStart": "Session start", @@ -106,6 +111,7 @@ "subagentStop": "Subagent finish", "preCompact": "Compaction start", "stop": "Session stopped", + "OpenCodeUsage": "Model usage", } # --- Local-only network hardening ------------------------------------------- @@ -522,7 +528,7 @@ def get_settings() -> dict[str, Any]: } -_ONBOARDING_AGENTS = ("claude", "cursor", "codex") +_ONBOARDING_AGENTS = AGENTS def _stored_onboarding_agents() -> list[str]: @@ -761,6 +767,14 @@ def bridge_script() -> FileResponse: return FileResponse(script, media_type="text/plain") +@app.get("/opencode-plugin.js") +def opencode_plugin() -> FileResponse: + script = _BRIDGE_DIR / "opencode-plugin.js" + if not script.exists(): + raise HTTPException(status_code=404, detail="OpenCode plugin not found") + return FileResponse(script, media_type="text/javascript") + + async def _json_body(request: Request) -> dict[str, Any]: try: body = await request.json() @@ -785,7 +799,7 @@ async def _ingest_body(request: Request) -> tuple[dict[str, Any] | None, str | N @app.post("/v1/ingest/{source}") async def ingest(source: str, request: Request) -> dict[str, Any]: - if source not in ("claude", "cursor", "codex"): + if source not in AGENTS: raise HTTPException(status_code=404, detail=f"Unknown source: {source}") body, malformed_error = await _ingest_body(request) if body is None: @@ -1218,7 +1232,13 @@ def get_hook_status() -> dict[str, Any]: connected = bool(conn.get("connected")) manifest_installed = bool(entry.get("installed")) installed = manifest_installed or (not manifest_agents and events > 0) - if not installed and events == 0: + if source == "opencode" and source not in manifest_agents and events > 0: + # A manually copied plugin has no bridge-written manifest. Its + # events prove installation only when no explicit entry exists. + installed = True + installed_hooks = expected_hooks + missing_hooks = [] + if not installed: health = "not_installed" elif missing_hooks: health = "missing_hooks" @@ -1259,7 +1279,7 @@ def get_hook_status() -> dict[str, Any]: "updated_at": manifest.get("updated_at"), "endpoint": manifest.get("endpoint"), "manifest_found": bool(manifest_agents), - "repair_all_url": "/repair.sh?agents=claude,cursor,codex", + "repair_all_url": "/repair.sh?agents=claude,cursor,codex,opencode", "agents": agents, } diff --git a/backend/app/normalize.py b/backend/app/normalize.py index 4e2200f..fecf675 100644 --- a/backend/app/normalize.py +++ b/backend/app/normalize.py @@ -1,4 +1,4 @@ -"""Translate raw Claude Code / Cursor / Codex hook payloads into a common shape + category.""" +"""Translate agent hook payloads into a common shape and category.""" from __future__ import annotations @@ -19,7 +19,7 @@ subagent_label as _subagent_label, ) -Source = str # 'claude' | 'cursor' | 'codex' +Source = str # 'claude' | 'cursor' | 'codex' | 'opencode' LifecycleBoundary = Literal["session_start", "turn_end", "session_end"] APPROVAL_REVIEW_PREFIX = "The following is the Codex agent history" @@ -66,7 +66,7 @@ def lifecycle_boundary(source: Source, hook: str) -> LifecycleBoundary | None: if hook in ("SessionEnd", "sessionEnd"): return "session_end" if hook in ("Stop", "stop"): - return "turn_end" if source == "claude" else "session_end" + return "turn_end" if source in ("claude", "opencode") else "session_end" return None @@ -159,6 +159,16 @@ def categorize(source: Source, hook: str, body: dict[str, Any], tool: str | None "duration_ms": duration_ms, } + if hook == "OpenCodeUsage": + return { + "category": "lifecycle", + "title": "Model usage", + "target": body.get("model"), + "detail": _json_detail(body.get("usage") or {}), + "status": "ok", + "duration_ms": duration_ms, + } + # --- Prompts --- if hook in ("UserPromptSubmit", "beforeSubmitPrompt"): prompt = body.get("prompt") or body.get("user_message") or "" @@ -216,12 +226,13 @@ def categorize(source: Source, hook: str, body: dict[str, Any], tool: str | None "duration_ms": duration_ms, } if boundary == "turn_end": + turn_status = "error" if body.get("error") else "interrupted" if body.get("interrupted") else "ok" return { "category": "lifecycle", - "title": "Turn ended", + "title": "Turn failed" if turn_status == "error" else "Turn interrupted" if turn_status == "interrupted" else "Turn ended", "target": None, "detail": _json_detail(body), - "status": "ok", + "status": turn_status, "duration_ms": duration_ms, } @@ -271,7 +282,7 @@ def categorize(source: Source, hook: str, body: dict[str, Any], tool: str | None } # --- Tool calls --- - tool_name = _canonical_tool(tool or body.get("tool_name") or "") + tool_name = _canonical_tool(tool or body.get("tool_name") or "", source) tool_input = _coerce_tool_input(body.get("tool_input")) tool_response = body.get("tool_response") or body.get("tool_output") or {} @@ -339,7 +350,7 @@ def categorize(source: Source, hook: str, body: dict[str, Any], tool: str | None "title": hook, "target": tool or body.get("tool_name"), "detail": _json_detail(body), - "status": "ok", + "status": "error" if _is_failure(hook) else "ok", "duration_ms": duration_ms, } @@ -412,8 +423,8 @@ def normalize(source: Source, body: dict[str, Any] | None) -> dict[str, Any]: cwd = roots[0] tool = _cursor_tool(hook, body) else: - # Claude Code and Codex both ride Claude-Code-style stdin payloads. - if source != "codex": + # Claude Code, Codex, and OpenCode use the same bridge payload shape. + if source not in ("codex", "opencode"): source = "claude" session_id = body.get("session_id") or "unknown" cwd = body.get("cwd") diff --git a/backend/app/tool_classification.py b/backend/app/tool_classification.py index 1816ad9..096c618 100644 --- a/backend/app/tool_classification.py +++ b/backend/app/tool_classification.py @@ -98,7 +98,18 @@ def _matches(path: str | None, patterns: tuple[re.Pattern[str], ...]) -> bool: return any(p.search(path.replace("\\", "/")) for p in patterns) -def canonical_tool(name: str | None) -> str: +_OPENCODE_TOOL_ALIASES = { + "bash": "Bash", "shell": "Shell", "read": "Read", "write": "Write", + "edit": "Edit", "multiedit": "MultiEdit", "glob": "Glob", "grep": "Grep", + "list": "Search", "webfetch": "WebFetch", "websearch": "WebSearch", + "task": "Task", "subagent": "Subagent", "question": "AskUserQuestion", + "todowrite": "TodoWrite", "todoread": "TodoRead", "skill": "Skill", +} + + +def canonical_tool(name: str | None, source: str | None = None) -> str: + if source == "opencode": + name = _OPENCODE_TOOL_ALIASES.get(name or "", name or "") return _TOOL_ALIASES.get(name or "", name or "") @@ -163,7 +174,7 @@ def extract_path(body: dict[str, Any], tool_input: dict[str, Any] | None = None) return val inputs = tool_input if isinstance(tool_input, dict) else body.get("tool_input") if isinstance(inputs, dict): - for key in ("file_path", "path", "notebook_path"): + for key in ("file_path", "filePath", "path", "notebook_path"): val = inputs.get(key) if isinstance(val, str) and val: return val @@ -330,7 +341,7 @@ def classify_tool(invocation: ToolInvocation) -> dict[str, Any] | None: "title": label, "target": key or label, "detail": json_detail({"input": tool_input, "response": tool_response}), - "status": "ok", + "status": status, "duration_ms": duration_ms, } diff --git a/backend/tests/test_opencode.py b/backend/tests/test_opencode.py new file mode 100644 index 0000000..3599e6a --- /dev/null +++ b/backend/tests/test_opencode.py @@ -0,0 +1,106 @@ +"""OpenCode-specific projection and installation-state regressions.""" +from __future__ import annotations + +import ast +from pathlib import Path + +import pytest + +from app import db, main, store +from app.normalize import normalize +from app.tool_classification import canonical_tool + + +@pytest.mark.parametrize("tool,category", [("read", "file_read"), ("write", "file_edit"), ("edit", "file_edit")]) +@pytest.mark.parametrize("path_key", ["path", "filePath"]) +def test_native_file_tools_preserve_targets_and_errors(tool, category, path_key): + event = normalize("opencode", { + "session_id": "oc", "hook_event_name": "PostToolUseFailure", "tool_name": tool, + "tool_input": {path_key: "/repo/π space.txt"}, "tool_response": "fixture error", + }) + assert (event["tool"], event["category"], event["target"], event["status"]) == ( + tool, category, "/repo/π space.txt", "error", + ) + + +@pytest.mark.parametrize("tool", ["bash", "shell"]) +def test_native_shell_keeps_command(tool): + event = normalize("opencode", { + "hook_event_name": "PostToolUse", "tool_name": tool, + "tool_input": {"command": "printf π"}, "tool_response": "π", + }) + assert (event["category"], event["target"], event["status"]) == ("shell", "printf π", "ok") + + +def test_aliases_are_source_specific_and_unknown_failures_remain_errors(): + assert canonical_tool("shell", "opencode") == "Shell" + assert canonical_tool("shell", "claude") == "shell" + event = normalize("opencode", {"tool_name": "new_tool", "hook_event_name": "PostToolUseFailure"}) + assert (event["category"], event["status"]) == ("other", "error") + + +@pytest.mark.parametrize("tool", ["task", "subagent"]) +def test_subagent_failure_keeps_call_link(tool): + event = normalize("opencode", { + "tool_name": tool, "hook_event_name": "PostToolUseFailure", "tool_use_id": "child-call", + "tool_input": {"description": "Inspect fixture", "prompt": "Read file"}, + }) + assert (event["category"], event["target"], event["status"]) == ("subagent", "child-call", "error") + + +@pytest.mark.parametrize("fields,status", [({"error": {"message": "provider rejected"}}, "error"), ({"interrupted": True}, "interrupted"), ({}, "ok")]) +def test_terminal_execution_status_is_visible(fields, status): + event = normalize("opencode", {"hook_event_name": "Stop", **fields}) + assert (event["category"], event["phase"], event["status"]) == ("lifecycle", "end", status) + + +def test_removed_plugin_is_not_inferred_from_historical_events(fresh_db, monkeypatch): + db.record_ingest("opencode", {"session_id": "history", "hook_event_name": "UserPromptSubmit", "prompt": "old"}) + monkeypatch.setattr(main, "_read_hook_manifest", lambda: {"agents": {"opencode": { + "installed": False, "installed_hooks": [], "missing_hooks": main.EXPECTED_HOOKS["opencode"], + }}}) + status = next(a for a in main.get_hook_status()["agents"] if a["source"] == "opencode") + assert status["installed"] is False + assert status["health"] == "not_installed" + assert status["missing_hooks"] + + +def test_manual_plugin_inferred_without_its_own_manifest_entry(fresh_db, monkeypatch): + db.record_ingest("opencode", {"session_id": "manual", "hook_event_name": "UserPromptSubmit", "prompt": "new"}) + monkeypatch.setattr(main, "_read_hook_manifest", lambda: {"agents": {"claude": {"installed": True}}}) + status = next(a for a in main.get_hook_status()["agents"] if a["source"] == "opencode") + assert status["installed"] is True + assert status["missing_hooks"] == [] + + +def test_desktop_spec_serves_the_plugin(tmp_path, monkeypatch): + repo = Path(__file__).resolve().parents[2] + tree = ast.parse((repo / "macos/packaging/cot-collector.spec").read_text()) + analysis = next(n for n in ast.walk(tree) if isinstance(n, ast.Call) and isinstance(n.func, ast.Name) and n.func.id == "Analysis") + datas = next(k.value for k in analysis.keywords if k.arg == "datas") + for entry in datas.elts: + source, destination = entry.elts + if ast.literal_eval(destination) == "bridge": + path = repo.joinpath(*(ast.literal_eval(a) for a in source.args[1:])) + (tmp_path / path.name).write_bytes(path.read_bytes()) + monkeypatch.setattr(main, "_BRIDGE_DIR", tmp_path) + response = main.opencode_plugin() + assert Path(response.path).read_bytes() == (repo / "bridge/opencode-plugin.js").read_bytes() + + +def test_upgrade_repairs_existing_opencode_tools_without_reimport(fresh_db): + for tool in ["read", "new_tool"]: + db.record_ingest("opencode", { + "session_id": "old-oc", "hook_event_name": "PostToolUseFailure", "tool_name": tool, + "tool_input": {"filePath": "/repo/a.py"}, "_dedup_key": tool, + }) + with store.write() as conn: + conn.execute("UPDATE events SET category='other', status='ok', target=tool") + conn.execute("UPDATE settings SET value='9' WHERE key='migrations_version'") + db.init_db() + db.init_db() + with store.read() as conn: + rows = conn.execute("SELECT tool, category, target, status FROM events ORDER BY id").fetchall() + assert len(rows) == 2 + assert tuple(rows[0]) == ("read", "file_read", "/repo/a.py", "error") + assert tuple(rows[1]) == ("new_tool", "other", "new_tool", "error") diff --git a/backend/tests/test_spool.py b/backend/tests/test_spool.py index cff0dc1..a3c8327 100644 --- a/backend/tests/test_spool.py +++ b/backend/tests/test_spool.py @@ -94,6 +94,38 @@ def _spool_lines() -> list[dict]: return [json.loads(ln) for ln in bridge.SPOOL_PATH.read_text().splitlines() if ln.strip()] +def test_opencode_persists_before_network_work_can_be_interrupted(): + def run(state): + payload = {"session_id": "oc-durable", "hook_event_name": "Stop", "_dedup_key": "oc:stop"} + def interrupted_send(url, body, timeout): + assert _spool_lines()[0]["payload"]["_dedup_key"] == "oc:stop" + raise InterruptedError("process interrupted during replay") + bridge._send_once = interrupted_send + try: + bridge._post(INGEST.replace("claude", "opencode"), payload) + except InterruptedError: + pass + else: + raise AssertionError("fixture did not interrupt the network send") + assert _spool_lines()[0]["payload"]["session_id"] == "oc-durable" + sink = _Sink() + bridge._send_once = sink.send + assert bridge._spool_flush() + assert sink.delivered[0]["_dedup_key"] == "oc:stop" + assert _spool_lines() == [] + _with_temp_spool(run) + + +def test_opencode_failed_spool_write_is_not_acknowledged_as_success(): + import pytest + def run(state): + # A directory at the expected file location reproduces a write error. + bridge.SPOOL_PATH.mkdir() + with pytest.raises(OSError, match="Unable to persist"): + bridge._post(INGEST.replace("claude", "opencode"), {"session_id": "oc"}) + _with_temp_spool(run) + + def test_post_spools_when_collector_down(): def body(_state): sink = _Sink(up=False) diff --git a/bridge/cot b/bridge/cot index a42e156..bc65e5d 100755 --- a/bridge/cot +++ b/bridge/cot @@ -248,7 +248,7 @@ def _url_request_target(url: str) -> str: _HOOK_CAPTURED_AT: str | None = None -def _spool_append(url: str, payload: dict) -> None: +def _spool_append(url: str, payload: dict) -> bool: """Queue an undelivered ingest event, trimming the oldest to stay under the byte cap.""" path = _url_request_target(url) @@ -278,11 +278,13 @@ def _spool_append(url: str, payload: dict) -> None: SPOOL_PATH.write_bytes(existing) with SPOOL_PATH.open("a", encoding="utf-8") as fh: fh.write(line) + return True except OSError as exc: print(f"cot: failed to spool event: {exc}", file=sys.stderr) + return False -def _spool_target(record: dict) -> str: +def _spool_target(record: dict, endpoint: str | None = None) -> str: """Resolve new and legacy spool records against the current collector.""" path = record.get("path") if not isinstance(path, str) or not path: @@ -292,10 +294,10 @@ def _spool_target(record: dict) -> str: path = _url_request_target(legacy_url) if not path.startswith("/v1/ingest/"): return "" - return f"{COT_ENDPOINT.rstrip('/')}{path}" + return f"{(endpoint or COT_ENDPOINT).rstrip('/')}{path}" -def _spool_flush(timeout: float = 3.0) -> bool: +def _spool_flush(timeout: float = 3.0, *, endpoint: str | None = None) -> bool: """Replay queued events in order. Stops at the first unreachable send and rewrites the unsent remainder. Returns True if the collector was reachable (spool empty or fully/partially drained), False if it is still down.""" @@ -320,7 +322,7 @@ def _spool_flush(timeout: float = 3.0) -> bool: rec = json.loads(lines[i]) except json.JSONDecodeError: continue # skip a corrupt line rather than wedge the queue - target = _spool_target(rec) + target = _spool_target(rec, endpoint) if not target: continue if not _send_once(target, rec.get("payload", {}), timeout): @@ -342,6 +344,14 @@ def _spool_flush(timeout: float = 3.0) -> bool: def _post(url: str, payload: dict) -> None: is_ingest = "/v1/ingest/" in url if is_ingest: + if "/v1/ingest/opencode" in url: + # The plugin bounds subprocess lifetime. Persist before HTTP work + # so a timeout during queue replay cannot lose the new event. + if not _spool_append(url, payload): + raise OSError("Unable to persist OpenCode hook event") + if _collector_reachable(): + _spool_flush(endpoint=url.split("/v1/ingest/", 1)[0]) + return if not _collector_reachable(): _spool_append(url, payload) return @@ -1402,7 +1412,13 @@ def _log_launch(body: dict) -> None: pass -AGENTS = ("claude", "cursor", "codex") +AGENTS = ("claude", "cursor", "codex", "opencode") +OPENCODE_EVENTS = ( + "SessionStart", "UserPromptSubmit", "PreToolUse", "PostToolUse", + "PostToolUseFailure", "afterAgentResponse", "afterAgentThought", + "OpenCodeUsage", "PermissionRequest", "PostCompact", "Stop", +) +OPENCODE_PLUGIN_MARKER = 'id: "cot.observability"' # Where each agent reads its hook config, and whether it is a Cursor-style file # (a top-level "version" alongside "hooks") vs. a Claude/Codex-style file. @@ -1410,6 +1426,7 @@ _CONFIG_TARGETS = { "claude": (Path.home() / ".claude" / "settings.json", False), "cursor": (Path.home() / ".cursor" / "hooks.json", True), "codex": (Path.home() / ".codex" / "hooks.json", False), + "opencode": (Path.home() / ".config" / "opencode" / "plugins" / "cot.js", False), } @@ -1616,6 +1633,28 @@ def _backup_config(agent: str, path: Path, action: str) -> dict | None: def _agent_hook_status(agent: str) -> dict: + if agent == "opencode": + path = _CONFIG_TARGETS[agent][0] + try: + content = path.read_text() + except OSError: + content = "" + installed = OPENCODE_PLUGIN_MARKER in content + expected = list(OPENCODE_EVENTS) + return { + "agent": agent, + "config_path": str(path), + "config_exists": path.exists(), + "valid_json": None, + "error": None, + "installed": installed, + "expected_hooks": expected, + "installed_hooks": expected if installed else [], + "missing_hooks": [] if installed else expected, + "backup_count": len(_backup_entries(agent)), + "latest_backup": _latest_backup(agent), + "warnings": [] if installed else ["plugin_missing"], + } path, _ = _CONFIG_TARGETS[agent] expected = _template_events(_hook_templates()[agent]) config, err = _read_json_file(path) @@ -1691,6 +1730,28 @@ def _hooks_installed(agent: str) -> bool: def _install_agent(agent: str, action: str = "install") -> dict: + if agent == "opencode": + path = _CONFIG_TARGETS[agent][0] + try: + current = path.read_text() + except FileNotFoundError: + current = None + if current is not None and OPENCODE_PLUGIN_MARKER not in current: + raise ValueError(f"{path} exists and is not cot's plugin; move it before installing") + req = urllib.request.Request(f"{COT_ENDPOINT}/opencode-plugin.js") + with urllib.request.urlopen(req, timeout=10) as response: + plugin = response.read().decode("utf-8") + if OPENCODE_PLUGIN_MARKER not in plugin: + raise ValueError("collector returned an invalid OpenCode plugin") + changed = current != plugin + backup = None + if changed: + path.parent.mkdir(parents=True, exist_ok=True) + backup = _backup_config(agent, path, action) + path.write_text(plugin) + status = _agent_hook_status(agent) + status.update({"changed": changed, "backup_created": backup, "status": "ok"}) + return status path, cursor_style = _CONFIG_TARGETS[agent] template = _hook_templates()[agent] @@ -1721,6 +1782,21 @@ def _install_agent(agent: str, action: str = "install") -> dict: def _uninstall_agent(agent: str, action: str = "uninstall") -> dict: + if agent == "opencode": + path = _CONFIG_TARGETS[agent][0] + try: + content = path.read_text() + except FileNotFoundError: + content = None + if content is not None and OPENCODE_PLUGIN_MARKER not in content: + return {"agent": agent, "config_path": str(path), "changed": False, + "backup_created": None, "status": "error", "error": "plugin file is not owned by cot"} + backup = None + if content is not None: + backup = _backup_config(agent, path, action) + path.unlink() + return {"agent": agent, "config_path": str(path), "changed": content is not None, + "backup_created": backup, "status": "ok"} path, _ = _CONFIG_TARGETS[agent] config, err = _read_json_file(path) if err == "missing": @@ -1767,8 +1843,8 @@ def _write_hook_status() -> None: """Record installed hook status where the collector can read it. The collector runs in Docker with ~/.cot mounted, so it cannot inspect - ~/.claude, ~/.cursor, or ~/.codex directly. This manifest is the bridge's - small, content-free handoff to the dashboard. + agent config directories directly. This manifest is the bridge's small, + content-free handoff to the dashboard. """ status = { "updated_at": _now_z(), @@ -1796,7 +1872,7 @@ def _cmd_install(args: list[str]) -> int: print(f"cot: backup saved -> {result['backup_created']['backup_path']}") elif not result.get("changed"): print(f"cot: {agent} hooks already up to date") - except OSError as exc: + except (OSError, ValueError, urllib.error.URLError) as exc: results.append( { "agent": agent, @@ -1839,9 +1915,14 @@ def _cmd_install(args: list[str]) -> int: if "codex" in agents: print("cot: run /hooks in Codex to review and trust the new hooks.") + if not status_written or any(r["status"] != "ok" for r in results): + return 1 + # Auto-import historical transcripts after hook setup. if os.environ.get("COT_SKIP_IMPORT") in ("1", "true", "yes"): return 0 + if all(agent == "opencode" for agent in agents): + return 0 print("cot: importing historical transcripts…") try: _reindex_if_parser_changed() @@ -1849,6 +1930,8 @@ def _cmd_install(args: list[str]) -> int: offsets = _import_offsets_load() total = 0 for agent in agents: + if agent == "opencode": + continue # Live plugin only; no historical OpenCode parser yet. count = _import_agent(agent, None, offsets, origins, verbose=False) total += count _import_offsets_save(offsets) @@ -2160,7 +2243,7 @@ def _confirm_purge(assume_yes: bool) -> bool: return False print("cot purge will remove:") print(f" - Docker collector container: {_container_name()}") - print(" - cot hooks from Claude, Cursor, and Codex") + print(" - cot hooks from Claude, Cursor, Codex, and OpenCode") print(f" - local cot state, bridge, and trace DB: {STATE_DIR}") answer = input("Type 'purge' to continue: ").strip() return answer == "purge" @@ -3940,8 +4023,8 @@ def main() -> int: if len(sys.argv) < 3 or sys.argv[1] != "hook": print( "Usage:\n" - " cot hook # reads a hook payload from stdin\n" - " cot install [claude cursor codex] # wire up agent hook configs\n" + " cot hook # reads a hook payload from stdin\n" + " cot install [claude cursor codex opencode] # wire up agent integrations\n" " cot import [--agent X] [--since] # import historical transcripts\n" " cot reimport # wipe + re-import all transcripts\n" " cot recover-questions # fill AskQuestion answers from transcripts\n" diff --git a/bridge/install.sh b/bridge/install.sh index cd594a6..9a23947 100755 --- a/bridge/install.sh +++ b/bridge/install.sh @@ -50,7 +50,7 @@ title() { else printf '%s\n' "cot bridge installer v${COT_VERSION}" fi - printf '%s\n\n' "hooks for Claude Code | Cursor | Codex" + printf '%s\n\n' "hooks for Claude Code | Cursor | Codex | OpenCode" } section() { @@ -284,6 +284,7 @@ if [ -z "${SELECTION}" ]; then printf '%s\n' " 1) claude Claude Code ~/.claude/settings.json" printf '%s\n' " 2) cursor Cursor ~/.cursor/hooks.json" printf '%s\n' " 3) codex Codex ~/.codex/hooks.json" + printf '%s\n' " 4) opencode OpenCode ~/.config/opencode/plugins/cot.js" echo "" printf "Enter names/numbers (space-separated), 'all', or 'none' [all]: " read -r SELECTION < /dev/tty || SELECTION="" @@ -292,7 +293,7 @@ if [ -z "${SELECTION}" ]; then SELECTION="none" echo "" warn "Non-interactive shell; skipping hook setup" - printf '%s\n' " Re-run with COT_AGENTS=\"claude cursor codex\" to wire up hooks," + printf '%s\n' " Re-run with COT_AGENTS=\"claude cursor codex opencode\" to wire up hooks," printf '%s\n' " or run: ${TARGET} install" fi fi @@ -301,13 +302,14 @@ fi AGENTS="" case " ${SELECTION} " in *" none "*|*" None "*|*" NONE "*) AGENTS="" ;; - *" all "*|*" All "*|*" ALL "*) AGENTS="claude cursor codex" ;; + *" all "*|*" All "*|*" ALL "*) AGENTS="claude cursor codex opencode" ;; *) for token in ${SELECTION}; do case "${token}" in 1|claude|Claude|CLAUDE) AGENTS="${AGENTS} claude" ;; 2|cursor|Cursor|CURSOR) AGENTS="${AGENTS} cursor" ;; 3|codex|Codex|CODEX) AGENTS="${AGENTS} codex" ;; + 4|opencode|OpenCode|OPENCODE) AGENTS="${AGENTS} opencode" ;; esac done ;; @@ -333,10 +335,12 @@ if [ -n "${AGENTS}" ]; then section "Import" IMPORT_ARGS="" - for a in ${AGENTS}; do IMPORT_ARGS="${IMPORT_ARGS} --agent ${a}"; done - if run_spinner "Importing historical transcripts" env COT_ENDPOINT="${COT_ENDPOINT}" "${TARGET}" import ${IMPORT_ARGS}; then + for a in ${AGENTS}; do + [ "${a}" = "opencode" ] || IMPORT_ARGS="${IMPORT_ARGS} --agent ${a}" + done + if [ -n "${IMPORT_ARGS}" ] && run_spinner "Importing historical transcripts" env COT_ENDPOINT="${COT_ENDPOINT}" "${TARGET}" import ${IMPORT_ARGS}; then [ -z "${RUN_OUTPUT}" ] || printf '%s\n' "${RUN_OUTPUT}" - else + elif [ -n "${IMPORT_ARGS}" ]; then warn "Transcript import had issues (non-fatal)." details "Output" "${RUN_OUTPUT}" fi diff --git a/bridge/opencode-plugin.js b/bridge/opencode-plugin.js new file mode 100644 index 0000000..c0c8ccd --- /dev/null +++ b/bridge/opencode-plugin.js @@ -0,0 +1,283 @@ +// cot OpenCode plugin (OpenCode 1.18.29+ and 2.x). +// Copy to ~/.config/opencode/plugins/cot.js; no opencode.json entry is needed. +import { spawn } from "node:child_process" +import { homedir } from "node:os" +import { join } from "node:path" + +const bridge = join(homedir(), ".cot", "bin", "cot") +const keyFor = (prefix, id) => id ? `${prefix}:${id}` : undefined + +function timestamp(value) { + const date = new Date(typeof value === "number" || typeof value === "string" ? value : NaN) + return Number.isFinite(date.getTime()) ? date.toISOString() : new Date().toISOString() +} + +function serialize(payload) { + const ancestors = [] + return JSON.stringify(payload, function (_, value) { + if (typeof value === "bigint") return String(value) + if (!value || typeof value !== "object") return value + while (ancestors.length && ancestors.at(-1) !== this) ancestors.pop() + if (ancestors.includes(value)) return "[Circular]" + ancestors.push(value) + return value + }) +} + +function createSender() { + const seen = new Set(), pending = new Map() + return function send(sessionID, hook, fields = {}, key) { + if (!sessionID) return Promise.resolve(false) + const scopedKey = key && `${sessionID}:${key}` + if (scopedKey && seen.has(scopedKey)) return Promise.resolve(true) + if (scopedKey && pending.has(scopedKey)) return pending.get(scopedKey) + const task = new Promise((resolve) => { + let child, timer, settled = false + const finish = (ok, reason) => { + if (settled) return + settled = true + clearTimeout(timer) + if (ok && scopedKey) { + seen.add(scopedKey) + if (seen.size > 10000) seen.delete(seen.values().next().value) + } + if (!ok) console.error(`cot: OpenCode ${hook} delivery failed (${reason}); retry remains available`) + resolve(ok) + } + try { + const body = serialize({ + session_id: sessionID, hook_event_name: hook, origin: "hook", ...fields, + timestamp: timestamp(fields.timestamp), + ...(scopedKey ? { _dedup_key: `opencode:${scopedKey}` } : {}), + }) + // OpenCode tears down its plugin worker at exit. A detached bridge + // must finish persisting the event even if that worker is terminated. + child = spawn(bridge, ["hook", "opencode"], { detached: true, stdio: ["pipe", "ignore", "ignore"] }) + child.on("error", (error) => finish(false, error.code || "spawn")) + child.on("close", (code) => finish(code === 0, `exit ${code}`)) + child.stdin.on("error", () => {}) + child.stdin.end(body) + // Longer than the bridge's individual three-second HTTP timeout. + timer = setTimeout(() => { child.kill("SIGKILL"); finish(false, "timeout") }, 5000) + } catch { finish(false, "serialization or spawn") } + }) + if (scopedKey) { + pending.set(scopedKey, task) + void task.then(() => pending.delete(scopedKey)) + } + return task + } +} + +function sessionStarted(send, info, directory) { + const sid = info?.id || info?.sessionID + return send(sid, "SessionStart", { + cwd: info?.directory || info?.location?.directory || directory, + parent_session_id: info?.parentID, session_title: info?.title, + }, keyFor("session", sid)) +} + +function usageFields(tokens, model) { + const t = tokens || {} + return { model, usage: { + input_tokens: t.input || 0, output_tokens: t.output || 0, + cache_read_tokens: t.cache?.read || 0, cache_write_tokens: t.cache?.write || 0, + } } +} + +function modelName(model) { + return model?.providerID && model?.id ? `${model.providerID}/${model.id}` : undefined +} + +function toolFailed(tool, status, metadata) { + return status === "error" || ((tool === "shell" || tool === "bash") && + typeof metadata?.exit === "number" && metadata.exit !== 0) +} + +function v1Event(send, event, directory) { + const p = event?.properties || {} + if (event?.type === "session.created") return sessionStarted(send, p.info, directory) + if (event?.type === "session.idle") return send(p.sessionID, "Stop", { cwd: directory }) + if (event?.type === "session.compacted") return send(p.sessionID, "PostCompact", { cwd: directory }) + if (event?.type === "permission.asked") return send(p.sessionID, "PermissionRequest", { + cwd: directory, tool_name: p.permission, permission: p, + }, keyFor("permission", p.id)) + if (event?.type === "message.updated") { + const info = p.info + if (info?.role !== "assistant" || !info.time?.completed) return + return send(info.sessionID, "OpenCodeUsage", { + cwd: directory, ...usageFields(info.tokens, info.modelID), + }, keyFor("usage", info.id)) + } + if (event?.type !== "message.part.updated") return + const part = p.part + if (!part?.sessionID || !part.id) return + if (part.type === "tool" && (part.state?.status === "completed" || part.state?.status === "error")) { + const failed = toolFailed(part.tool, part.state.status, part.state.metadata) + return send(part.sessionID, failed ? "PostToolUseFailure" : "PostToolUse", { + cwd: directory, tool_name: part.tool, tool_input: part.state.input, + tool_response: part.state.error ?? part.state.output, tool_metadata: part.state.metadata, tool_use_id: part.callID, + }, keyFor("tool:end", part.callID || part.id)) + } + if ((part.type === "text" || part.type === "reasoning") && part.time?.end && typeof part.text === "string" && part.text.trim()) { + const response = part.type === "text" + return send(part.sessionID, response ? "afterAgentResponse" : "afterAgentThought", { + cwd: directory, [response ? "response" : "thought"]: part.text, + }, keyFor("part", part.id)) + } +} + +async function v2Event(send, event, directory, stepModels) { + const p = event.data || {}, sid = p.sessionID + const type = event.type.replace("session.next.", "session.") + const fields = { cwd: directory, timestamp: p.timestamp ?? event.created } + const messageKey = p.assistantMessageID && `${sid}:${p.assistantMessageID}` + switch (type) { + case "session.step.started": + if (messageKey) stepModels.set(messageKey, modelName(p.model)) + return + case "session.text.ended": + case "session.reasoning.ended": { + if (typeof p.text !== "string" || !p.text.trim()) return + const response = type === "session.text.ended" + const partID = p.textID || p.reasoningID || (p.assistantMessageID && `${p.assistantMessageID}:${type}:${p.ordinal ?? 0}`) || event.id + return send(sid, response ? "afterAgentResponse" : "afterAgentThought", { + ...fields, [response ? "response" : "thought"]: p.text, + }, keyFor("part", partID)) + } + case "session.step.ended": { + const model = stepModels.get(messageKey) || modelName(p.model) + stepModels.delete(messageKey) + await send(sid, "OpenCodeUsage", { ...fields, ...usageFields(p.tokens, model) }, keyFor("usage", p.assistantMessageID || event.id)) + // Older v2 builds used step finish as their terminal notification. + if (event.type.startsWith("session.next.") && p.finish && p.finish !== "tool-calls") { + await send(sid, "Stop", fields, keyFor("stop", p.assistantMessageID || event.id)) + } + return + } + case "session.step.failed": + stepModels.delete(messageKey) + return + case "session.compaction.ended": + await send(sid, "PostCompact", { ...fields, summary: p.text }, keyFor("compact", event.id || p.messageID)) + if (p.tokens) await send(sid, "OpenCodeUsage", { + ...fields, ...usageFields(p.tokens, modelName(p.model)), + }, keyFor("usage:compact", event.id || p.messageID)) + return + case "session.execution.succeeded": + case "session.execution.failed": + case "session.execution.interrupted": + return send(sid, "Stop", { + ...fields, ...(p.error ? { error: p.error } : {}), + ...(type === "session.execution.interrupted" ? { interrupted: true, reason: p.reason } : {}), + }, keyFor("execution", event.id)) + // usage.updated is cumulative. Counting it as well would double tokens. + } +} + +function v1Hooks({ directory }) { + const send = createSender() + return { + "chat.message": async (input, output) => { + const prompt = (Array.isArray(output.parts) ? output.parts : []) + .filter((part) => part.type === "text" && !part.synthetic && typeof part.text === "string") + .map((part) => part.text).join("\n") + if (prompt) await send(input.sessionID, "UserPromptSubmit", { + cwd: directory, prompt, model: input.model?.modelID, + }, keyFor("prompt", input.messageID || output.message?.id)) + }, + "tool.execute.before": async (input, output) => { + await send(input.sessionID, "PreToolUse", { + cwd: directory, tool_name: input.tool, tool_input: output.args, tool_use_id: input.callID, + }, keyFor("tool:start", input.callID)) + }, + event: async ({ event }) => { await v1Event(send, event, directory) }, + } +} + +// OpenCode 1 calls server(); OpenCode 2 calls setup(). No SDK dependency. +export default { + id: "cot.observability", + async server(ctx) { return v1Hooks(ctx) }, + async setup(ctx) { + const directory = ctx.location.directory + const send = createSender() + const knownSessions = new Map(), stepModels = new Map(), inFlight = new Set() + const ensureSession = async (sessionID, knownInfo) => { + if (!sessionID) return directory + let task = knownSessions.get(sessionID) + if (knownInfo || !task) { + task = (async () => { + if (knownInfo) return knownInfo + try { return await ctx.session.get({ sessionID }) } + catch { return { id: sessionID } } + })() + knownSessions.set(sessionID, task) + } + const info = await task + await sessionStarted(send, { ...info, id: info?.id || sessionID }, directory) + if (!info?.directory && !info?.location?.directory) knownSessions.delete(sessionID) + return info?.directory || info?.location?.directory || directory + } + await ctx.session.hook("prompt", async (event) => { + const cwd = await ensureSession(event.sessionID) + if (typeof event.prompt?.text === "string" && event.prompt.text) await send(event.sessionID, "UserPromptSubmit", { + cwd, prompt: event.prompt.text, + }, keyFor("prompt", event.messageID)) + }) + await ctx.tool.hook("execute.before", async (event) => { + const cwd = await ensureSession(event.sessionID) + await send(event.sessionID, "PreToolUse", { + cwd, tool_name: event.tool, tool_input: event.input, tool_use_id: event.callID, + }, keyFor("tool:start", event.callID)) + }) + await ctx.tool.hook("execute.after", async (event) => { + const cwd = await ensureSession(event.sessionID) + const failed = toolFailed(event.tool, event.status, event.result?.metadata || event.result?.output) + await send(event.sessionID, failed ? "PostToolUseFailure" : "PostToolUse", { + cwd, tool_name: event.tool, tool_input: event.input, + tool_response: event.error?.message || event.error || event.result, tool_use_id: event.callID, + }, keyFor("tool:end", event.callID)) + }) + await ctx.permission.hook("evaluate", async (event) => { + if (event.effect !== "ask") return + const cwd = await ensureSession(event.sessionID) + const resources = Array.isArray(event.resources) ? event.resources : [] + await send(event.sessionID, "PermissionRequest", { + cwd, tool_name: event.action, + permission: { action: event.action, resources, metadata: event.metadata, message: event.message }, + }, keyFor("permission", event.source?.id && `${event.source.id}:${event.action}:${resources.join(",")}`)) + }) + const controller = new AbortController() + const handleEvent = async (event) => { + try { + if (!event || typeof event.type !== "string") return + const p = event.data || {} + const sessionID = p.sessionID || p.info?.id || (event.type === "session.created" ? p.id : undefined) + if (!sessionID) return + const cwd = await ensureSession(sessionID, p.info || (event.type === "session.created" ? p : undefined)) + await v2Event(send, { ...event, data: { ...p, sessionID } }, cwd, stepModels) + } catch { console.error("cot: ignored malformed OpenCode event") } + } + void (async () => { + try { + for await (const event of ctx.event.subscribe({ signal: controller.signal })) { + // Consume the stream promptly; bridge latency must not build a + // subscriber backlog that loses terminal events at CLI shutdown. + const task = handleEvent(event) + inFlight.add(task) + void task.then(() => inFlight.delete(task)) + } + } catch { + // OpenCode may close the subscription during reload or shutdown. + } + })() + return async () => { + // Let already-buffered events dispatch, then finish bounded deliveries. + await new Promise((resolve) => setTimeout(resolve, 0)) + controller.abort() + await Promise.allSettled([...inFlight]) + knownSessions.clear(); stepModels.clear() + } + }, +} diff --git a/bridge/tests/opencode-plugin.test.mjs b/bridge/tests/opencode-plugin.test.mjs new file mode 100644 index 0000000..8526874 --- /dev/null +++ b/bridge/tests/opencode-plugin.test.mjs @@ -0,0 +1,124 @@ +import assert from 'node:assert/strict' +import { mkdtemp, readFile, writeFile, chmod, rename, rm } from 'node:fs/promises' +import { tmpdir } from 'node:os' +import { join } from 'node:path' +import { pathToFileURL } from 'node:url' +import { setTimeout as delay } from 'node:timers/promises' +import test from 'node:test' + +async function fixture(t) { + const root = await mkdtemp(join(tmpdir(), 'cot-opencode-test-')) + t.after(() => rm(root, { recursive: true, force: true })) + const bridge = join(root, 'cot'), log = join(root, 'events.jsonl'), fail = join(root, 'fail') + const executable = `#!${process.execPath}\nimport {appendFileSync,existsSync} from 'node:fs'; let body=''; for await(const c of process.stdin)body+=c; if(existsSync(${JSON.stringify(fail)}))process.exit(1); appendFileSync(${JSON.stringify(log)},body+'\\n');\n` + await writeFile(bridge, executable); await chmod(bridge, 0o755) + // Use .mjs for the executable so the fixture has no package dependency. + await rename(bridge, bridge + '.mjs') + const path = bridge + '.mjs' + const source = (await readFile(new URL('../opencode-plugin.js', import.meta.url), 'utf8')) + .replace('const bridge = join(homedir(), ".cot", "bin", "cot")', 'const bridge = ' + JSON.stringify(path)) + const module = join(root, 'plugin.mjs'); await writeFile(module, source) + const plugin = (await import(pathToFileURL(module).href)).default + const records = async () => { try { return (await readFile(log, 'utf8')).trim().split('\n').filter(Boolean).map(JSON.parse) } catch { return [] } } + return { plugin, path, fail, records } +} + +async function context(plugin) { + const handlers = new Map(), queue = [] + let stopped = false + const cleanup = await plugin.setup({ + location: { directory: '/project/A' }, + session: { get: async ({ sessionID }) => ({ id: sessionID, location: { directory: '/project/B' } }), hook: async (name, fn) => handlers.set('session:' + name, fn) }, + tool: { hook: async (name, fn) => handlers.set('tool:' + name, fn) }, + permission: { hook: async (name, fn) => handlers.set('permission:' + name, fn) }, + event: { async *subscribe({ signal }) { + signal.addEventListener('abort', () => { stopped = true }) + while (!stopped) { if (queue.length) yield queue.shift(); else await delay(1) } + } }, + }) + return { handlers, emit: (type, data, id = type) => queue.push({ type, data, id, created: 1791000000000 }), finish: async () => { while(queue.length) await delay(1); await cleanup() } } +} + +test('v2 burst preserves responses, thoughts, usage, compaction and terminal events on disposal', async (t) => { + const f = await fixture(t), c = await context(f.plugin), sid = 'v2' + await c.handlers.get('session:prompt')({ sessionID: sid, messageID: 'prompt', prompt: { text: 'π🙂' } }) + c.emit('session.step.started', { sessionID: sid, assistantMessageID: 'msg', model: { providerID: 'review', id: 'model' } }) + c.emit('session.text.ended', { sessionID: sid, assistantMessageID: 'msg', ordinal: 0, text: 'π🙂' }) + c.emit('session.text.ended', { sessionID: sid, assistantMessageID: 'msg', ordinal: 1, text: 'second part' }, 'text-2') + c.emit('session.reasoning.ended', { sessionID: sid, assistantMessageID: 'msg', ordinal: 0, text: 'thought' }) + c.emit('session.usage.updated', { sessionID: sid, tokens: { input: 999 } }) + c.emit('session.step.ended', { sessionID: sid, assistantMessageID: 'msg', finish: 'stop', tokens: { input: 120, output: 25, cache: { read: 7, write: 3 } } }) + c.emit('session.compaction.ended', { sessionID: sid, text: 'summary' }) + c.emit('session.execution.succeeded', { sessionID: sid }) + await c.finish() + const events = await f.records(), hooks = events.map(e => e.hook_event_name) + for (const hook of ['SessionStart', 'UserPromptSubmit', 'afterAgentResponse', 'afterAgentThought', 'OpenCodeUsage', 'PostCompact', 'Stop']) assert(hooks.includes(hook), hook) + assert.equal(hooks.filter(h => h === 'afterAgentResponse').length, 2) + assert.equal(hooks.filter(h => h === 'OpenCodeUsage').length, 1) + assert.deepEqual(events.find(e => e.hook_event_name === 'OpenCodeUsage').usage, {input_tokens:120,output_tokens:25,cache_read_tokens:7,cache_write_tokens:3}) + assert(events.every(e => e.cwd === '/project/B')) +}) + +test('v2 malformed events do not stop subscription and failures/interruption survive', async (t) => { + const f = await fixture(t), c = await context(f.plugin), sid = 'edge' + c.emit('session.text.ended', { sessionID: sid, text: 123 }) + c.emit('session.text.ended', { sessionID: sid, text: 'valid', timestamp: NaN }, 'valid-text') + c.emit('session.execution.failed', { sessionID: sid, error: { message: 'provider rejected' } }) + c.emit('session.execution.interrupted', { sessionID: sid, reason: 'user' }) + await c.handlers.get('permission:evaluate')({ sessionID: sid, effect: 'ask', source: { id: 'perm' }, action: 'shell' }) + await c.finish() + const events = await f.records() + assert.equal(events.filter(e => e.hook_event_name === 'afterAgentResponse').length, 1) + assert(events.some(e => e.error?.message === 'provider rejected')) + assert(events.some(e => e.interrupted)) + assert(events.every(e => Number.isFinite(Date.parse(e.timestamp)))) +}) + +test('failed spawn and nonzero bridge exit allow the same prompt to be retried', async (t) => { + const f = await fixture(t), c = await context(f.plugin) + const prompt = { sessionID: 'retry', messageID: 'same', prompt: { text: 'retry' } } + await rename(f.path, f.path + '.saved') + await c.handlers.get('session:prompt')(prompt) + await rename(f.path + '.saved', f.path) + await writeFile(f.fail, 'fail') + await c.handlers.get('session:prompt')(prompt) + await rm(f.fail) + await c.handlers.get('session:prompt')(prompt) + await c.handlers.get('session:prompt')(prompt) + await c.finish() + assert.deepEqual((await f.records()).map(e => e.hook_event_name).sort(), ['SessionStart', 'UserPromptSubmit']) +}) + +test('v1 circular and BigInt tool inputs preserve delivery without rejecting', async (t) => { + const f = await fixture(t), hooks = await f.plugin.server({ directory: '/v1' }) + const input = { value: 12n }; input.self = input + await hooks['tool.execute.before']({ sessionID: 'v1', callID: 'call', tool: 'read' }, { args: input }) + const [event] = await f.records() + assert.equal(event.tool_input.value, '12') + assert.equal(event.tool_input.self, '[Circular]') +}) + +test('nonzero shell exits are failures even when OpenCode marks the tool completed', async (t) => { + const f = await fixture(t), c = await context(f.plugin) + await c.handlers.get('tool:execute.after')({ sessionID: 'v2-shell', tool: 'shell', callID: 'shell', status: 'completed', input: { command: 'exit 7' }, result: { metadata: { exit: 7 }, output: 'failed-command' } }) + await c.finish() + const hooks = await f.plugin.server({ directory: '/v1' }) + await hooks.event({ event: { type: 'message.part.updated', properties: { part: { id: 'part', sessionID: 'v1-shell', callID: 'bash', type: 'tool', tool: 'bash', state: { status: 'completed', input: { command: 'exit 7' }, output: 'failed-command', metadata: { exit: 7 } } } } } }) + const events = (await f.records()).filter(e => e.tool_name) + assert.equal(events.length, 2) + assert(events.every(e => e.hook_event_name === 'PostToolUseFailure')) + assert.equal(events[0].tool_response.metadata.exit, 7) + assert.equal(events[1].tool_metadata.exit, 7) +}) + +test('stalled bridge is killed within the bound and retry is available', async (t) => { + const f = await fixture(t), hooks = await f.plugin.server({ directory: '/v1' }) + const original = await readFile(f.path) + await writeFile(f.path, `#!${process.execPath}\nsetInterval(()=>{},1000)\n`) + const start = Date.now(), input = { sessionID: 'timeout', callID: 'call', tool: 'read' }, output = { args: { path: '/v1/a' } } + await hooks['tool.execute.before'](input, output) + assert(Date.now() - start < 7000) + await writeFile(f.path, original) + await hooks['tool.execute.before'](input, output) + assert.equal((await f.records()).length, 1) +}) diff --git a/macos/packaging/cot-collector.spec b/macos/packaging/cot-collector.spec index 2d06170..6320afc 100644 --- a/macos/packaging/cot-collector.spec +++ b/macos/packaging/cot-collector.spec @@ -11,10 +11,11 @@ a = Analysis( binaries=[], datas=[ (os.path.join(REPO, "backend", "app", "pricing.json"), "app"), - # Served by GET /install.sh and GET /cot. main.py resolves these two + # Served by the bridge download endpoints. main.py resolves these two # directories up from app/main.py, which lands on _internal/bridge here. (os.path.join(REPO, "bridge", "install.sh"), "bridge"), (os.path.join(REPO, "bridge", "cot"), "bridge"), + (os.path.join(REPO, "bridge", "opencode-plugin.js"), "bridge"), ], hiddenimports=[ "uvicorn.logging", diff --git a/src/components/dashboard/ExportModal.tsx b/src/components/dashboard/ExportModal.tsx index 3399dfd..6f153de 100644 --- a/src/components/dashboard/ExportModal.tsx +++ b/src/components/dashboard/ExportModal.tsx @@ -74,7 +74,7 @@ const INCLUDE_OPTIONS: IncludeDef[] = [ { key: 'clarifications', label: 'Clarifications', hint: 'Questions the agent asked and user answers' }, ]; -const SOURCE_OPTIONS = ['claude', 'cursor', 'codex'] as const; +const SOURCE_OPTIONS = ['claude', 'cursor', 'codex', 'opencode'] as const; function getNestedValue(obj: Record, path: string): unknown { return path.split('.').reduce((o, k) => (o as Record)?.[k], obj); diff --git a/src/components/dashboard/SessionList.tsx b/src/components/dashboard/SessionList.tsx index 87ea4ec..446356f 100644 --- a/src/components/dashboard/SessionList.tsx +++ b/src/components/dashboard/SessionList.tsx @@ -108,6 +108,7 @@ export function SessionList({ { value: 'claude', label: 'Claude' }, { value: 'cursor', label: 'Cursor' }, { value: 'codex', label: 'Codex' }, + { value: 'opencode', label: 'OpenCode' }, ]} />
diff --git a/src/components/dashboard/SessionsTable.tsx b/src/components/dashboard/SessionsTable.tsx index ef8b2c9..13e5df1 100644 --- a/src/components/dashboard/SessionsTable.tsx +++ b/src/components/dashboard/SessionsTable.tsx @@ -110,7 +110,7 @@ export function SessionsTable({ onSelect }: SessionsTableProps) { }; const sourceOptions = useMemo(() => { - const set = new Set(['claude', 'cursor', 'codex']); + const set = new Set(['claude', 'cursor', 'codex', 'opencode']); sessions.forEach((s) => set.add(s.source)); return Array.from(set); }, [sessions]); @@ -467,4 +467,3 @@ function groupOf( const cwd = s.cwd || '(unknown path)'; return { key: cwd, label: basename(cwd), title: cwd }; } - diff --git a/src/components/forest/AgentMark.tsx b/src/components/forest/AgentMark.tsx index 1dbf1c2..e6caff5 100644 --- a/src/components/forest/AgentMark.tsx +++ b/src/components/forest/AgentMark.tsx @@ -13,6 +13,7 @@ export const AGENT_LABELS: Record = { claude: 'Claude Code', cursor: 'Cursor', codex: 'Codex', + opencode: 'OpenCode', cowork: 'Cowork', }; export const agentLabel = (id: string) => AGENT_LABELS[id] ?? id[0].toUpperCase() + id.slice(1); @@ -33,6 +34,15 @@ export function AgentMark({ id, className = '', size = 14 }: { id: AgentId | str ); } + if (id === 'opencode') { + // Official mark: anomalyco/opencode packages/ui/src/components/logo.tsx. + return ( + + ); + } if (id !== 'cursor') { return (
+ + MANUAL PLUGIN + +

Download the plugin and save it as ~/.config/opencode/plugins/cot.js, then restart OpenCode.

+ + Download cot.js + +
+ )} {fileSteps.map( (step) => step.code && diff --git a/src/components/onboarding/steps/PassiveSetup.tsx b/src/components/onboarding/steps/PassiveSetup.tsx index b119265..ef05090 100644 --- a/src/components/onboarding/steps/PassiveSetup.tsx +++ b/src/components/onboarding/steps/PassiveSetup.tsx @@ -1,6 +1,5 @@ import { useEffect, useState } from 'react'; -import type { AgentId } from '../../../lib/agents'; -import { getPassive, PassiveUnsupportedError, previewSchedule, runPassiveNow, updatePassive, type PassiveStatus } from '../../../lib/api'; +import { getPassive, PassiveUnsupportedError, previewSchedule, runPassiveNow, updatePassive, type PassiveAgent, type PassiveStatus } from '../../../lib/api'; import { agentLabel } from '../../forest/AgentMark'; import { fmt } from '../../forest/format'; import { Icon } from '../../forest/icons'; @@ -20,7 +19,7 @@ export type PassiveDraft = ReturnType; export function usePassiveDraft() { const [p, setP] = useState(null); const [failed, setFailed] = useState<'unsupported' | 'offline' | false>(false); - const [agents, setAgents] = useState(null); + const [agents, setAgents] = useState(null); const [cron, setCron] = useState(null); const [started, setStarted] = useState(false); const [err, setErr] = useState(null); @@ -56,7 +55,7 @@ export function PassiveSetup({ }: { step: number; onStep: (s: number) => void; - onFinish: (agents: AgentId[], origin: { x: number; y: number }) => void; + onFinish: (agents: PassiveAgent[], origin: { x: number; y: number }) => void; draft: PassiveDraft; }) { const { p, failed, agents, setAgents, cron, setCron, started, setStarted, err, setErr } = draft; @@ -97,7 +96,7 @@ export function PassiveSetup({ const chosen = p.agents_detail.filter((d) => agents.includes(d.agent)); const total = chosen.reduce((n, d) => n + d.transcripts, 0); const running = chosen.reduce((n, d) => n + d.active, 0); - const toggle = (a: AgentId, on: boolean) => setAgents((prev) => (on ? [...(prev ?? []), a] : (prev ?? []).filter((x) => x !== a))); + const toggle = (a: PassiveAgent, on: boolean) => setAgents((prev) => (on ? [...(prev ?? []), a] : (prev ?? []).filter((x) => x !== a))); if (step === 1) { return ( diff --git a/src/components/ui/AgentMark.tsx b/src/components/ui/AgentMark.tsx index 89ba232..47d8018 100644 --- a/src/components/ui/AgentMark.tsx +++ b/src/components/ui/AgentMark.tsx @@ -16,8 +16,16 @@ const CURSOR_2D_PATH = const CODEX_PATH = 'M22.2819 9.8211a5.9847 5.9847 0 0 0-.5157-4.9108 6.0462 6.0462 0 0 0-6.5098-2.9A6.0651 6.0651 0 0 0 4.9807 4.1818a5.9847 5.9847 0 0 0-3.9977 2.9 6.0462 6.0462 0 0 0 .7427 7.0966 5.98 5.98 0 0 0 .511 4.9107 6.051 6.051 0 0 0 6.5146 2.9001A5.9847 5.9847 0 0 0 13.2599 24a6.0557 6.0557 0 0 0 5.7718-4.2058 5.9894 5.9894 0 0 0 3.9977-2.9001 6.0557 6.0557 0 0 0-.7475-7.0729zm-9.022 12.6081a4.4755 4.4755 0 0 1-2.8764-1.0408l.1419-.0804 4.7783-2.7582a.7948.7948 0 0 0 .3927-.6813v-6.7369l2.02 1.1686a.071.071 0 0 1 .038.052v5.5826a4.504 4.504 0 0 1-4.4945 4.4944zm-9.6607-4.1254a4.4708 4.4708 0 0 1-.5346-3.0137l.142.0852 4.783 2.7582a.7712.7712 0 0 0 .7806 0l5.8428-3.3685v2.3324a.0804.0804 0 0 1-.0332.0615L9.74 19.9502a4.4992 4.4992 0 0 1-6.1408-1.6464zM2.3408 7.8956a4.485 4.485 0 0 1 2.3655-1.9728V11.6a.7664.7664 0 0 0 .3879.6765l5.8144 3.3543-2.0201 1.1685a.0757.0757 0 0 1-.071 0l-4.8303-2.7865A4.504 4.504 0 0 1 2.3408 7.872zm16.5963 3.8558L13.1038 8.364 15.1192 7.2a.0757.0757 0 0 1 .071 0l4.8303 2.7913a4.4944 4.4944 0 0 1-.6765 8.1042v-5.6772a.79.79 0 0 0-.407-.667zm2.0107-3.0231l-.142-.0852-4.7735-2.7818a.7759.7759 0 0 0-.7854 0L9.409 9.2297V6.8974a.0662.0662 0 0 1 .0284-.0615l4.8303-2.7866a4.4992 4.4992 0 0 1 6.6802 4.66zM8.3065 12.863l-2.02-1.1638a.0804.0804 0 0 1-.038-.0567V6.0742a4.4992 4.4992 0 0 1 7.3757-3.4537l-.142.0805L8.704 5.459a.7948.7948 0 0 0-.3927.6813zm1.0976-2.3654l2.602-1.4998 2.6069 1.4998v2.9994l-2.5974 1.4997-2.6067-1.4997z'; -/** Official Claude / Cursor / Codex brand marks. */ +/** Agent marks used in onboarding and source badges. */ export function AgentMark({ id, className = '', variant = '2d' }: AgentMarkProps) { + if (id === 'opencode') { + return ( + + ); + } if (id === 'claude') { return ( = { subagentStop: 'Subagent finish', preCompact: 'Compaction start', stop: 'Session stopped', + OpenCodeUsage: 'Model usage', }; export function hookLabel(name: string): string { @@ -195,6 +196,23 @@ export const AGENTS: Agent[] = [ code: CODEX_HOOKS, }], }, + { + id: 'opencode', + name: 'OpenCode', + product: 'OpenCode', + tagline: + 'Trace prompts, tool calls, responses, and usage from OpenCode 1.18.29+ and 2 through one local plugin file.', + events: [ + 'SessionStart', + 'UserPromptSubmit', + 'PreToolUse', + 'PostToolUse', + 'afterAgentResponse', + 'OpenCodeUsage', + 'Stop', + ], + steps: [INSTALL_STEP], + }, ]; export function getAgent(id: AgentId): Agent { diff --git a/src/lib/api.ts b/src/lib/api.ts index 1b73636..e9689c8 100644 --- a/src/lib/api.ts +++ b/src/lib/api.ts @@ -288,7 +288,7 @@ export interface Settings { ui_sidebar_open: boolean; /** Onboarding finished on this install, and the agents picked there. */ ui_onboarded: boolean; - ui_onboarding_agents: ('claude' | 'cursor' | 'codex')[]; + ui_onboarding_agents: AgentId[]; } export async function getSettings(): Promise { @@ -1274,6 +1274,18 @@ export async function sendTestEvent(source: AgentId): Promise { for (const ev of events) { await ingest('codex', { ...base, ...ev }); } + } else if (source === 'opencode') { + const base = { session_id: sid, cwd, model: 'openai/gpt-5' }; + const events: Record[] = [ + { hook_event_name: 'SessionStart', timestamp: ts() }, + { hook_event_name: 'UserPromptSubmit', prompt: 'Inspect the parser and fix the failing test.', timestamp: ts() }, + { hook_event_name: 'PreToolUse', tool_name: 'bash', tool_input: { command: 'npm test' }, timestamp: ts() }, + { hook_event_name: 'PostToolUse', tool_name: 'bash', tool_input: { command: 'npm test' }, tool_response: 'Tests passed', timestamp: ts() }, + { hook_event_name: 'afterAgentResponse', response: 'The parser test now passes.', timestamp: ts() }, + { hook_event_name: 'OpenCodeUsage', usage: { input_tokens: 120, output_tokens: 40 }, timestamp: ts() }, + { hook_event_name: 'Stop', timestamp: ts() }, + ]; + for (const ev of events) await ingest('opencode', { ...base, ...ev }); } else { const base = { conversation_id: sid, workspace_roots: [cwd], cwd }; const events: Record[] = [ diff --git a/src/lib/settings.ts b/src/lib/settings.ts index 195765c..2c24ca7 100644 --- a/src/lib/settings.ts +++ b/src/lib/settings.ts @@ -57,20 +57,20 @@ export function writeTimelineSidebarMode(mode: TimelineSidebarMode): void { } } -export function readSavedAgents(): ('claude' | 'cursor' | 'codex')[] { +export function readSavedAgents(): ('claude' | 'cursor' | 'codex' | 'opencode')[] { try { const raw = localStorage.getItem(AGENTS_KEY); if (raw) { const parsed = JSON.parse(raw) as unknown; if (Array.isArray(parsed)) { return parsed.filter( - (x): x is 'claude' | 'cursor' | 'codex' => - x === 'claude' || x === 'cursor' || x === 'codex', + (x): x is 'claude' | 'cursor' | 'codex' | 'opencode' => + x === 'claude' || x === 'cursor' || x === 'codex' || x === 'opencode', ); } } const legacy = localStorage.getItem(LEGACY_AGENT_KEY); - if (legacy === 'claude' || legacy === 'cursor' || legacy === 'codex') return [legacy]; + if (legacy === 'claude' || legacy === 'cursor' || legacy === 'codex' || legacy === 'opencode') return [legacy]; } catch { /* ignore */ } @@ -86,7 +86,7 @@ export function readOnboarded(): boolean { } /** Local copy of the onboarding result; the collector holds the durable one. */ -export function writeOnboarded(agents: ('claude' | 'cursor' | 'codex')[]): void { +export function writeOnboarded(agents: ('claude' | 'cursor' | 'codex' | 'opencode')[]): void { try { localStorage.setItem(ONBOARDED_KEY, '1'); localStorage.setItem(AGENTS_KEY, JSON.stringify(agents)); diff --git a/src/lib/sourceLabels.ts b/src/lib/sourceLabels.ts index da4eb95..f1971fe 100644 --- a/src/lib/sourceLabels.ts +++ b/src/lib/sourceLabels.ts @@ -10,6 +10,8 @@ export function sourceLabel(source: string): string { return 'Cursor'; case 'codex': return 'Codex'; + case 'opencode': + return 'OpenCode'; case 'api': return 'API'; case 'custom': @@ -30,11 +32,13 @@ export function sourceTag(source: string): string { return 'CURSOR'; case 'codex': return 'CODEX'; + case 'opencode': + return 'OPENCODE'; default: return source.toUpperCase(); } } export function isKnownAgent(source: string): source is AgentId { - return source === 'claude' || source === 'cursor' || source === 'codex'; + return source === 'claude' || source === 'cursor' || source === 'codex' || source === 'opencode'; }