diff --git a/.github/workflows/release.yml b/.github/workflows/release.yml
index bd8e4231..ba1722be 100644
--- a/.github/workflows/release.yml
+++ b/.github/workflows/release.yml
@@ -66,9 +66,39 @@ jobs:
path: dist/OpenWorkGraph-Claude.mcpb
if-no-files-found: error
+ sdk-package:
+ runs-on: ubuntu-latest
+ steps:
+ - uses: actions/checkout@v4
+ - uses: actions/setup-python@v5
+ with:
+ python-version: "3.11"
+ - uses: actions/setup-node@v4
+ with:
+ node-version: "20"
+ - name: Validate standalone custom-agent SDKs
+ run: |
+ python -m py_compile sdk/python/openworkgraph_agent.py
+ node --check sdk/typescript/index.mjs
+ - name: Stage standalone custom-agent SDK assets
+ run: |
+ mkdir -p dist
+ cp sdk/python/openworkgraph_agent.py dist/OpenWorkGraph-Agent-Python.py
+ cp sdk/typescript/index.mjs dist/OpenWorkGraph-Agent-Node.mjs
+ cp sdk/typescript/index.d.ts dist/OpenWorkGraph-Agent-Node.d.ts
+ - name: Upload custom-agent SDK artifact
+ uses: actions/upload-artifact@v4
+ with:
+ name: pkg-agent-sdk
+ path: |
+ dist/OpenWorkGraph-Agent-Python.py
+ dist/OpenWorkGraph-Agent-Node.mjs
+ dist/OpenWorkGraph-Agent-Node.d.ts
+ if-no-files-found: error
+
publish-release:
if: github.ref == 'refs/heads/main'
- needs: [macos-package, windows-package, mcpb-package]
+ needs: [macos-package, windows-package, mcpb-package, sdk-package]
runs-on: ubuntu-latest
steps:
- uses: actions/checkout@v4
@@ -87,6 +117,11 @@ jobs:
with:
name: OpenWorkGraph-Claude
path: release-assets
+ - name: Download custom-agent SDK assets
+ uses: actions/download-artifact@v4
+ with:
+ name: pkg-agent-sdk
+ path: release-assets
- name: Publish GitHub Release when VERSION is new
env:
GH_TOKEN: ${{ github.token }}
@@ -120,6 +155,11 @@ jobs:
### Agent setup control plane
The dashboard now separates giving an AI access to OpenWorkGraph context from observing an agent's own execution. It provides reviewable setup material for Claude Code lifecycle hooks, Codex trace export, OpenAI Agents tracing and generic OpenTelemetry/custom structural adapters. OpenWorkGraph does not silently edit third-party configuration files. An integration is shown as active only when telemetry actually observed by the local evidence store supports that status.
+ ### Custom harnesses
+ Arbitrary self-built or third-party agent harnesses can now connect in either or both directions. Python and Node/TypeScript helpers, OTLP/HTTP JSON and raw structural HTTP can send privacy-safe execution telemetry through the dedicated write-only agent credential. Any MCP-capable harness can separately read only the OpenWorkGraph context the user has authorized. The setup flow keeps telemetry write permission and context/history read permission explicitly separate.
+
+ The standalone helpers do not accept or serialize prompt text, model responses, tool arguments/results, returned values, exception text or hidden reasoning. OpenWorkGraph observer failures remain fail-open for the agent. This release publishes **OpenWorkGraph-Agent-Python.py**, **OpenWorkGraph-Agent-Node.mjs** and **OpenWorkGraph-Agent-Node.d.ts** as standalone release assets.
+
### Compact MCP for new connections
New dashboard-generated connections, the Claude MCP bundle and the on-demand local HTTP bridge expose a compact MCP surface focused on current context, search, canonical evidence, repeated workflows, task context, prior-run feedback and agent execution inspection. Overlapping tools are consolidated so AI clients have fewer competing tool definitions.
@@ -158,6 +198,9 @@ jobs:
release-assets/OpenWorkGraph-Windows.zip \
release-assets/OpenWorkGraph-Windows.zip.sha256 \
release-assets/OpenWorkGraph-Claude.mcpb \
+ release-assets/OpenWorkGraph-Agent-Python.py \
+ release-assets/OpenWorkGraph-Agent-Node.mjs \
+ release-assets/OpenWorkGraph-Agent-Node.d.ts \
--target "$GITHUB_SHA" \
--title "OpenWorkGraph ${TAG}" \
--notes-file /tmp/openworkgraph-release-notes.md
diff --git a/VERSION b/VERSION
index c5c73510..5f8cbfdb 100644
--- a/VERSION
+++ b/VERSION
@@ -1 +1 @@
-0.94.0
+0.95.0
diff --git a/dashboard/custom_harness_setup.js b/dashboard/custom_harness_setup.js
new file mode 100644
index 00000000..74f492f1
--- /dev/null
+++ b/dashboard/custom_harness_setup.js
@@ -0,0 +1,90 @@
+(() => {
+ let cache=null;
+ const esc=value=>{const node=document.createElement('div');node.textContent=String(value??'');return node.innerHTML;};
+
+ function installStyle(){
+ if(document.querySelector('#owg-custom-harness-style'))return;
+ const style=document.createElement('style');style.id='owg-custom-harness-style';
+ style.textContent=`
+ .harness-card{border:1px solid var(--line);border-radius:14px;padding:15px;background:#f8faf8;display:flex;flex-direction:column;min-height:205px}
+ .harness-card h3{font-size:15px;margin:0 0 5px}.harness-card p{font-size:12.5px;color:var(--muted);line-height:1.45;margin:0 0 12px}.harness-card .actions{margin-top:auto}
+ .harness-badges{display:flex;gap:6px;flex-wrap:wrap;margin:4px 0 10px}.harness-badge{font-size:10.5px;border-radius:999px;padding:3px 7px;background:#eef3ee;color:#385143;font-weight:700}
+ .harness-methods{display:flex;gap:6px;flex-wrap:wrap;margin:10px 0}.harness-methods button{min-height:32px;padding:5px 9px;font-size:12px}.harness-methods button.active{background:#2f7d55;color:#fff}
+ .harness-code{white-space:pre-wrap;overflow-wrap:anywhere;background:#f3f4f0;border:1px solid var(--line);border-radius:9px;padding:10px;font:12px/1.45 ui-monospace,SFMono-Regular,Menlo,monospace;margin:7px 0}
+ .harness-two-way{display:grid;grid-template-columns:1fr 1fr;gap:10px;margin:12px 0}.harness-two-way>div{border:1px solid var(--line);border-radius:10px;padding:10px;font-size:12px}.harness-two-way strong{display:block;margin-bottom:3px}
+ @media(max-width:620px){.harness-two-way{grid-template-columns:1fr}}
+ `;
+ document.head.appendChild(style);
+ }
+
+ async function load(){
+ if(cache)return cache;
+ await window.__owgAuthReady;
+ const response=await fetch('/v1/custom-harness-setup',{cache:'no-store'});
+ if(!response.ok)throw new Error('custom harness setup unavailable');
+ cache=await response.json();return cache;
+ }
+
+ function injectCard(){
+ const section=document.querySelector('#agent-observation-setup');
+ if(!section||document.querySelector('#customHarnessCard'))return;
+ const grid=section.querySelector('.setup-grid');if(!grid)return;
+ const card=document.createElement('div');card.id='customHarnessCard';card.className='harness-card';
+ card.innerHTML=`
Any framework
Your own agent / harness OpenClaw / Hermes-style Internal agents Any MCP client
Connect an arbitrary harness in either or both directions: send its structural execution to OpenWorkGraph and optionally let it read the OWG context you authorize.
Connect your agent
`;
+ grid.appendChild(card);
+ card.querySelector('#customHarnessSetupButton').onclick=openSetup;
+ }
+
+ function code(text,id){return `${esc(text||'')}
Copy `;}
+ function privacy(payload){
+ const p=payload.privacy||{};
+ return `${p.prompt_content?'⚠':'✓'} No prompt content
${p.model_response_content?'⚠':'✓'} No response content
${p.tool_arguments?'⚠':'✓'} No tool arguments
${p.tool_results?'⚠':'✓'} No tool results
${p.reasoning?'⚠':'✓'} No reasoning
${p.exception_text?'⚠':'✓'} No exception text
`;
+ }
+
+ function methodBody(payload,method){
+ const write=payload.write||{};
+ if(method==='python'){
+ const x=write.python||{};
+ return `Python SDK. A dependency-free helper with context managers for runs, models and tools. It is fail-open: OWG going down never stops the agent.
Install the tiny SDK ${code(x.install,'harnessPyInstall')}Give it the local write-only credential ${code(x.environment,'harnessPyEnv')}Instrument the harness ${code(x.example,'harnessPyExample')}`;
+ }
+ if(method==='typescript'){
+ const x=write.typescript||{};
+ return `TypeScript / Node. Dependency-free ESM with TypeScript declarations. Node 18+.
Get the SDK files ${code(x.download,'harnessTsDownload')}Give it the local write-only credential ${code(x.environment,'harnessTsEnv')}Instrument the harness ${code(x.example,'harnessTsExample')}`;
+ }
+ if(method==='otel'){
+ const x=write.otel||{};
+ const token=(write.raw_http||{}).authorization||'';
+ const env=`export OTEL_EXPORTER_OTLP_TRACES_ENDPOINT="${x.endpoint||''}"\nexport OTEL_EXPORTER_OTLP_TRACES_PROTOCOL="http/json"\nexport OTEL_EXPORTER_OTLP_TRACES_HEADERS="Authorization=${token}"`;
+ return `OpenTelemetry. If the harness already emits portable GenAI spans, this is the least invasive option. OWG accepts OTLP/HTTP JSON and ignores unknown spans rather than guessing.
${code(env,'harnessOtelEnv')}`;
+ }
+ const x=write.raw_http||{};
+ const curl=`curl -X POST ${x.endpoint||''} \\\n -H 'Content-Type: application/json' \\\n -H 'Authorization: ${x.authorization||''}' \\\n --data '${String(x.example||'').replaceAll("'","'\\''")}'`;
+ return `Raw HTTP. For any language/runtime: POST the canonical structural envelope directly. The credential is write-only.
${code(curl,'harnessRawHttp')}The server validates the complete batch and rejects content-bearing fields such as prompts, responses, messages, reasoning, tool arguments and tool results.
`;
+ }
+
+ async function openSetup(){
+ try{
+ const payload=await load();
+ const body=`Use either direction independently. Observing an agent does not let it read your work history, and giving it MCP context does not automatically record its execution.
Agent → OpenWorkGraph Structural run/model/tool/handoff/approval/error telemetry through a write-only credential.
OpenWorkGraph → Agent Optional MCP context, still controlled by OWG's AI-access switch and saved-history lease.
${privacy(payload)}1. Observe this harness Python TypeScript / Node OpenTelemetry Raw HTTP
${methodBody(payload,'python')}
2. Optional: let the harness read OWG context via MCP Paste this standard stdio MCP server entry into any MCP-capable harness. Context access remains OFF unless you enable OWG's AI-access master switch; historical reads additionally require your saved-history lease.
${code(JSON.stringify(payload.read?.config||{},null,2),'harnessMcpConfig')}The telemetry token above cannot read anything. MCP is a separate connection and permission path.
`;
+ window.openModal?.('Connect your agent','Custom harness · two-way connection',body);
+ bind(payload);
+ }catch(_){window.openModal?.('Custom harness setup unavailable','Connect','Could not generate local setup material. Confirm OpenWorkGraph is running and reload the dashboard.
');}
+ }
+
+ function bind(payload){
+ document.querySelectorAll('[data-harness-copy]').forEach(button=>button.onclick=()=>copy(button.dataset.harnessCopy||'',button));
+ document.querySelectorAll('#harnessMethods [data-method]').forEach(button=>button.onclick=()=>{
+ document.querySelectorAll('#harnessMethods [data-method]').forEach(x=>x.classList.toggle('active',x===button));
+ const target=document.querySelector('#harnessMethodBody');if(target)target.innerHTML=methodBody(payload,button.dataset.method||'python');
+ document.querySelectorAll('[data-harness-copy]').forEach(copyButton=>copyButton.onclick=()=>copy(copyButton.dataset.harnessCopy||'',copyButton));
+ });
+ }
+
+ async function copy(id,button){
+ const text=document.querySelector(`#${CSS.escape(id)}`)?.textContent||'';
+ try{await navigator.clipboard.writeText(text);const before=button.textContent;button.textContent='Copied';setTimeout(()=>button.textContent=before,1100);}catch(_){window.prompt('Copy this:',text);}
+ }
+
+ function install(){installStyle();injectCard();setTimeout(injectCard,150);setTimeout(injectCard,800);}
+ if(document.readyState==='loading')document.addEventListener('DOMContentLoaded',install);else install();
+})();
diff --git a/docs/CHANGELOG_V095.md b/docs/CHANGELOG_V095.md
new file mode 100644
index 00000000..b1519570
--- /dev/null
+++ b/docs/CHANGELOG_V095.md
@@ -0,0 +1,38 @@
+# OpenWorkGraph v0.95.0
+
+## Custom agents and harnesses
+
+OpenWorkGraph now exposes a first-class, framework-neutral path for arbitrary agent runtimes. A self-built harness, an OpenClaw/Hermes-style setup, an internal company agent, or another framework can participate without OpenWorkGraph needing a named native adapter.
+
+The connection is deliberately two-way and independent:
+
+- **agent -> OpenWorkGraph:** privacy-safe structural execution telemetry through the existing write-only agent-ingest boundary;
+- **OpenWorkGraph -> agent:** optional authorized context over the compact local MCP surface.
+
+Giving a harness the telemetry token never gives it context/history read access. Giving a harness MCP context never silently enables observation of its execution.
+
+## Standalone helpers
+
+v0.95 adds:
+
+- a dependency-free Python helper (`openworkgraph-agent`) with run/model/tool context managers;
+- a dependency-free Node 18+/TypeScript helper with declarations;
+- existing OTLP/HTTP JSON integration for harnesses that already emit portable GenAI traces;
+- raw structural HTTP for any other language/runtime.
+
+The GitHub release publishes the Python and Node helper files as standalone assets in addition to the normal macOS, Windows and Claude MCP packages.
+
+## Privacy boundary
+
+The custom helpers intentionally have no API for prompt text, model-response content, tool arguments/results, returned values, exception text or hidden reasoning. Tests verify that returned secrets and exception messages remain in the harness and never enter telemetry.
+
+Observation failures remain fail-open for the agent: bounded background delivery can drop telemetry, but it does not enter the agent's control path.
+
+## Reliability
+
+- Python telemetry admission and shutdown share a synchronization boundary, preventing events from being accepted after the worker has closed.
+- Node shutdown waits for an already in-flight telemetry request before returning.
+- The local dashboard setup endpoint is authenticated and non-cacheable.
+- The generated setup keeps write-only telemetry credentials separate from MCP/history permissions.
+
+See [Custom agent harnesses](CUSTOM_HARNESSES.md) for the integration model and examples.
diff --git a/docs/CUSTOM_HARNESSES.md b/docs/CUSTOM_HARNESSES.md
new file mode 100644
index 00000000..c38dbe98
--- /dev/null
+++ b/docs/CUSTOM_HARNESSES.md
@@ -0,0 +1,120 @@
+# Custom agent harnesses
+
+OpenWorkGraph can work with an arbitrary agent runtime, including self-built harnesses and frameworks that OpenWorkGraph does not know by name.
+
+There are two independent connections:
+
+```text
+OpenWorkGraph -- MCP context --> agent / harness
+OpenWorkGraph <-- structural telemetry -- agent / harness
+```
+
+You can use either direction or both. Sending telemetry never grants read access to work history. Giving a harness MCP context never automatically enables observation of its execution.
+
+## 1. Agent -> OpenWorkGraph: structural execution
+
+OpenWorkGraph accepts the same vendor-neutral structural contract used by its native integrations. Useful operations include:
+
+- `run_started` / `run_finished`
+- `model_call`
+- `tool_call`
+- `handoff`
+- `human_approval_requested` / `human_approval_received`
+- `error`
+
+The contract can also preserve `run_id`, `trace_id`, `span_id`, `parent_span_id`, `workflow_id`, timing, a coarse tool category, model identifier and token counts when the runtime exposes them.
+
+The content boundary does not change for custom harnesses. Do not send prompts, model responses, messages, chain-of-thought/reasoning, tool arguments, tool results, returned values or exception text. The server rejects content-bearing fields in direct structural batches.
+
+### Python
+
+The repository ships a standalone, dependency-free Python helper in `sdk/python/openworkgraph_agent.py`. It can run in the harness's own environment and is also independently installable from the repository subdirectory.
+
+```python
+from openworkgraph_agent import AgentObserver
+
+owg = AgentObserver("my-agent", framework="my-harness")
+try:
+ with owg.run(workflow_id="optional-workflow-id") as run:
+ with run.model(model="model-id"):
+ call_model()
+ with run.tool("repository_search", category="search"):
+ search_repository()
+finally:
+ owg.shutdown()
+```
+
+The SDK uses a bounded background queue and is fail-open: an unavailable OpenWorkGraph observer drops telemetry rather than blocking or failing the agent. Returned values and exception messages are never serialized.
+
+### TypeScript / Node
+
+`sdk/typescript/index.mjs` is dependency-free ESM for Node 18+ and `index.d.ts` provides TypeScript declarations.
+
+```js
+import {AgentObserver} from "./openworkgraph-agent.mjs";
+
+const owg = new AgentObserver("my-agent", {framework: "my-harness"});
+try {
+ await owg.withRun(async run => {
+ await run.model({model: "model-id"}, async () => callModel());
+ await run.tool("repository_search", {category: "search"}, async () => searchRepository());
+ });
+} finally {
+ await owg.shutdown();
+}
+```
+
+### OpenTelemetry
+
+If the runtime already emits portable GenAI OpenTelemetry spans, point OTLP/HTTP JSON traces at:
+
+```text
+http://127.0.0.1:8787/agent-ingest/v1/otel
+```
+
+Use the signal-specific traces endpoint, `http/json`, and the installation's dedicated write-only agent-ingest bearer. Unknown spans are ignored rather than guessed.
+
+### Raw HTTP
+
+Any language can post canonical structural events to:
+
+```text
+POST http://127.0.0.1:8787/agent-ingest/v1/events
+Authorization: Bearer
+```
+
+The write-only token cannot read work history, exports, summaries or MCP context.
+
+## 2. OpenWorkGraph -> agent: context over MCP
+
+Any MCP-capable harness can launch the compact local MCP server using the configuration generated in the dashboard under **Connect -> Your own agent / harness**.
+
+The custom harness is still subject to OpenWorkGraph's normal disclosure controls:
+
+- the run-level AI access switch must be ON;
+- Redacted remains the default disclosure level unless the user explicitly allows Full and organization policy permits it;
+- historical reads require the user's saved-history lease and are restricted to its authorized date range.
+
+The MCP credential/path is separate from the write-only telemetry bearer.
+
+## Observation level
+
+Use the observation level that describes what the harness truly exposes:
+
+- `native_trace` — the runtime itself emits a complete structural trace;
+- `instrumented_tools` — hooks/wrappers observe structural tool execution;
+- `mcp_only` — only operations crossing an MCP boundary are visible;
+- `os_observed` — only desktop/browser structural evidence is available;
+- `outcome_only` — only externally visible results are known.
+
+Missing signals mean **not observed**, not that the agent did not perform them.
+
+## OpenClaw, Hermes and other frameworks
+
+A named adapter is not required. If a framework exposes callbacks/hooks, wrap those callbacks with the Python/TypeScript helper or translate them to raw structural events. If it already emits OpenTelemetry, use the OTLP route. If neither is available, OpenWorkGraph can still represent surface-observed activity at a lower observation level.
+
+Do not add provider-specific event vocabulary unless the information is genuinely portable. Thin adapters should translate native events into the shared OpenWorkGraph contract.
+
+## Containers and remote workers
+
+The default local observer binds to loopback. A harness running in a separate container, VM or remote machine may not be able to reach `127.0.0.1:8787` on the host. Do not expose the local observer broadly just to make telemetry convenient. Use an explicitly secured deployment/network path appropriate to the environment; until such a path exists, treat direct SDK ingestion as same-host/local integration.
diff --git a/mcpb/manifest.json b/mcpb/manifest.json
index f6c1fbc7..fb515277 100644
--- a/mcpb/manifest.json
+++ b/mcpb/manifest.json
@@ -2,9 +2,9 @@
"manifest_version": "0.3",
"name": "openworkgraph-local",
"display_name": "OpenWorkGraph",
- "version": "0.94.0",
+ "version": "0.95.0",
"description": "Connect Claude Desktop to the compact local OpenWorkGraph context surface.",
- "long_description": "Uses the OpenWorkGraph installation already running on this computer. v0.94 adds explicit local history retention, separate saved-history AI access, lightweight history navigation, and structural browser-agent lifecycle observation for major web agents. Canonical workflow evidence remains primary; Context Pulse provides incremental factual updates. Retained history is separately user-controlled: list_history can navigate saved human and agent sessions only while a time-limited saved-history lease is active, after which canonical date-ranged tools can drill into evidence. Derived task and pattern views remain optional and non-authoritative. The legacy 24-tool MCP entrypoint remains available for existing configurations while new connections use this compact surface. Prompts, model responses, tool arguments/results, typed text, clipboard contents, and hidden reasoning are not captured by these agent adapters.",
+ "long_description": "Uses the OpenWorkGraph installation already running on this computer. v0.95 adds a framework-neutral custom-harness setup flow plus standalone Python and Node helpers for privacy-safe structural agent telemetry; arbitrary MCP-capable harnesses can separately read authorized OpenWorkGraph context. v0.94 added explicit local history retention, separate saved-history AI access, lightweight history navigation, and structural browser-agent lifecycle observation. Canonical workflow evidence remains primary; Context Pulse provides incremental factual updates. Retained history is separately user-controlled: list_history can navigate saved human and agent sessions only while a time-limited saved-history lease is active. The legacy 24-tool MCP entrypoint remains available for existing configurations while new connections use this compact surface. Prompts, model responses, tool arguments/results, typed text, clipboard contents, exception text, returned values, and hidden reasoning are not captured by the custom agent helpers.",
"author": {"name": "Koyar Afrasyab / Kinvectum"},
"repository": {"type": "git", "url": "https://github.com/KAVentures/openworkgraph"},
"server": {"type": "node", "entry_point": "server/index.js", "mcp_config": {"command": "node", "args": ["${__dirname}/server/index.js"], "env": {}}},
diff --git a/openworkgraph_agent/__init__.py b/openworkgraph_agent/__init__.py
new file mode 100644
index 00000000..963a0d3b
--- /dev/null
+++ b/openworkgraph_agent/__init__.py
@@ -0,0 +1,5 @@
+"""Public OpenWorkGraph SDK for custom Python agent harnesses."""
+
+from sdk.python.openworkgraph_agent import AgentObserver, AgentRun, ObserverStats, RunIdentity
+
+__all__ = ["AgentObserver", "AgentRun", "ObserverStats", "RunIdentity"]
diff --git a/pyproject.toml b/pyproject.toml
index 3df100b1..ea4d63cb 100644
--- a/pyproject.toml
+++ b/pyproject.toml
@@ -1,6 +1,6 @@
[project]
name = "workflow-observer"
-version = "0.94.0"
+version = "0.95.0"
description = "Local-first work evidence, self-hosted organizational context gateway, REST API, and MCP access."
requires-python = ">=3.11"
license = {file = "LICENSE"}
diff --git a/sdk/__init__.py b/sdk/__init__.py
new file mode 100644
index 00000000..97d16576
--- /dev/null
+++ b/sdk/__init__.py
@@ -0,0 +1 @@
+"""Language SDK sources shipped with OpenWorkGraph."""
diff --git a/sdk/python/__init__.py b/sdk/python/__init__.py
new file mode 100644
index 00000000..5bc7d603
--- /dev/null
+++ b/sdk/python/__init__.py
@@ -0,0 +1,3 @@
+from .openworkgraph_agent import AgentObserver, AgentRun, ObserverStats, RunIdentity
+
+__all__ = ["AgentObserver", "AgentRun", "ObserverStats", "RunIdentity"]
diff --git a/sdk/python/openworkgraph_agent.py b/sdk/python/openworkgraph_agent.py
new file mode 100644
index 00000000..d8f76c5a
--- /dev/null
+++ b/sdk/python/openworkgraph_agent.py
@@ -0,0 +1,397 @@
+from __future__ import annotations
+
+"""Dependency-free OpenWorkGraph structural telemetry for custom Python agents.
+
+This file can be copied directly into another harness. Its public API is
+intentionally content-blind: no prompt, response, reasoning, tool arguments,
+tool results, returned values, or exception text are accepted or serialized.
+Delivery is bounded, asynchronous, and fail-open by default.
+"""
+
+from contextlib import AbstractContextManager
+from dataclasses import dataclass
+from datetime import datetime, timezone
+import json
+import os
+import queue
+import threading
+import time
+from typing import Any, Callable, Mapping
+from urllib import request as urllib_request
+import uuid
+
+OBSERVATION_LEVELS = frozenset({"native_trace", "instrumented_tools", "mcp_only", "os_observed", "outcome_only"})
+TOOL_CATEGORIES = frozenset({
+ "filesystem", "shell", "browser", "code", "search", "network", "database",
+ "messaging", "issue_tracker", "deployment", "mcp", "other", "none",
+})
+USAGE_KEYS = frozenset({"input_tokens", "output_tokens", "cached_input_tokens", "total_tokens"})
+
+
+def _now() -> str:
+ return datetime.now(timezone.utc).isoformat().replace("+00:00", "Z")
+
+
+def _id(prefix: str) -> str:
+ return f"{prefix}-{uuid.uuid4().hex}"
+
+
+def _clean(value: Any, *, limit: int = 240) -> str:
+ return " ".join(str(value or "").split())[:limit]
+
+
+def _usage(value: Mapping[str, int] | None) -> dict[str, int]:
+ if not value:
+ return {}
+ unknown = set(value) - USAGE_KEYS
+ if unknown:
+ raise ValueError(f"unsupported usage fields: {', '.join(sorted(unknown))}")
+ out: dict[str, int] = {}
+ for key, raw in value.items():
+ amount = int(raw)
+ if amount < 0:
+ raise ValueError(f"{key} must be non-negative")
+ out[key] = amount
+ return out
+
+
+def _default_sender(endpoint: str, token: str, events: list[dict]) -> None:
+ if not token:
+ raise RuntimeError("missing OWG agent-ingest token")
+ body = json.dumps({"events": events}, separators=(",", ":")).encode("utf-8")
+ req = urllib_request.Request(
+ endpoint,
+ data=body,
+ method="POST",
+ headers={
+ "Content-Type": "application/json",
+ "Authorization": f"Bearer {token}",
+ },
+ )
+ with urllib_request.urlopen(req, timeout=1.5) as response:
+ if not 200 <= int(response.status) < 300:
+ raise RuntimeError("OpenWorkGraph telemetry write failed")
+
+
+@dataclass(frozen=True)
+class RunIdentity:
+ run_id: str
+ trace_id: str
+ workflow_id: str
+
+
+@dataclass(frozen=True)
+class ObserverStats:
+ accepted: int
+ dropped: int
+ send_failures: int
+ batches_sent: int
+ queued: int
+
+
+class _Delivery:
+ def __init__(
+ self,
+ *,
+ endpoint: str,
+ token: str,
+ sender: Callable[[str, str, list[dict]], None] | None = None,
+ max_queue: int = 512,
+ batch_size: int = 32,
+ flush_interval: float = 0.2,
+ ) -> None:
+ self.endpoint = endpoint
+ self.token = token
+ self.sender = sender or _default_sender
+ self.queue: queue.Queue[dict] = queue.Queue(maxsize=max(1, min(int(max_queue), 10_000)))
+ self.batch_size = max(1, min(int(batch_size), 256))
+ self.flush_interval = max(0.02, min(float(flush_interval), 2.0))
+ self.stop = threading.Event()
+ self.lock = threading.Lock()
+ self.thread: threading.Thread | None = None
+ self.accepted = 0
+ self.dropped = 0
+ self.send_failures = 0
+ self.batches_sent = 0
+
+ def _ensure_worker(self) -> None:
+ with self.lock:
+ if self.thread is not None and self.thread.is_alive():
+ return
+ if self.stop.is_set():
+ return
+ self.thread = threading.Thread(target=self._run, name="openworkgraph-agent-sdk", daemon=True)
+ self.thread.start()
+
+ def emit(self, event: dict) -> bool:
+ if self.stop.is_set():
+ with self.lock:
+ self.dropped += 1
+ return False
+ self._ensure_worker()
+ # Admission and shutdown share this lock. An event either enters the
+ # queue before shutdown closes admission, or is deterministically
+ # rejected afterwards; it cannot be accepted after the worker exits.
+ with self.lock:
+ if self.stop.is_set():
+ self.dropped += 1
+ return False
+ try:
+ self.queue.put_nowait(dict(event))
+ except queue.Full:
+ self.dropped += 1
+ return False
+ self.accepted += 1
+ return True
+
+ def _take_batch(self) -> list[dict]:
+ try:
+ first = self.queue.get(timeout=self.flush_interval)
+ except queue.Empty:
+ return []
+ batch = [first]
+ while len(batch) < self.batch_size:
+ try:
+ batch.append(self.queue.get_nowait())
+ except queue.Empty:
+ break
+ return batch
+
+ def _run(self) -> None:
+ while not self.stop.is_set() or self.queue.unfinished_tasks:
+ batch = self._take_batch()
+ if not batch:
+ if self.stop.is_set():
+ break
+ continue
+ try:
+ self.sender(self.endpoint, self.token, batch)
+ except Exception:
+ with self.lock:
+ self.send_failures += 1
+ else:
+ with self.lock:
+ self.batches_sent += 1
+ finally:
+ for _ in batch:
+ self.queue.task_done()
+
+ def flush(self, timeout: float = 2.0) -> bool:
+ deadline = time.monotonic() + max(0.0, min(float(timeout), 10.0))
+ while self.queue.unfinished_tasks and time.monotonic() < deadline:
+ time.sleep(0.01)
+ return self.queue.unfinished_tasks == 0
+
+ def shutdown(self, timeout: float = 2.0) -> None:
+ with self.lock:
+ self.stop.set()
+ self.flush(timeout)
+ thread = self.thread
+ if thread is not None and thread.is_alive():
+ thread.join(timeout=max(0.0, min(float(timeout), 10.0)))
+
+ def stats(self) -> ObserverStats:
+ with self.lock:
+ return ObserverStats(
+ self.accepted,
+ self.dropped,
+ self.send_failures,
+ self.batches_sent,
+ self.queue.qsize(),
+ )
+
+
+class _TimedOperation(AbstractContextManager["_TimedOperation"]):
+ def __init__(
+ self,
+ run: "AgentRun",
+ *,
+ operation: str,
+ tool_name: str = "",
+ tool_category: str = "none",
+ model: str = "",
+ usage: Mapping[str, int] | None = None,
+ ) -> None:
+ self.run = run
+ self.operation = operation
+ self.tool_name = _clean(tool_name, limit=200)
+ self.tool_category = tool_category
+ self.model_name = _clean(model, limit=200)
+ self.token_usage = _usage(usage)
+ self.started = 0.0
+ self.span_id = _id("span")
+
+ def __enter__(self) -> "_TimedOperation":
+ self.started = time.monotonic()
+ return self
+
+ def __exit__(self, exc_type, exc, tb) -> bool:
+ self.run._emit(
+ self.operation,
+ status="error" if exc_type is not None else "success",
+ span_id=self.span_id,
+ tool_name=self.tool_name,
+ tool_category=self.tool_category,
+ model=self.model_name,
+ usage=self.token_usage,
+ duration_seconds=max(0.0, time.monotonic() - self.started),
+ )
+ return False
+
+
+class AgentRun(AbstractContextManager["AgentRun"]):
+ def __init__(
+ self,
+ observer: "AgentObserver",
+ *,
+ run_id: str | None = None,
+ trace_id: str | None = None,
+ workflow_id: str | None = None,
+ trigger_event_id: str | None = None,
+ ) -> None:
+ self.observer = observer
+ self.run_id = _clean(run_id) or _id("run")
+ self.trace_id = _clean(trace_id) or _id("trace")
+ self.workflow_id = _clean(workflow_id)
+ self.trigger_event_id = _clean(trigger_event_id)
+ self.started = 0.0
+ self.closed = False
+
+ @property
+ def identity(self) -> RunIdentity:
+ return RunIdentity(self.run_id, self.trace_id, self.workflow_id)
+
+ def __enter__(self) -> "AgentRun":
+ self.started = time.monotonic()
+ self._emit("run_started", status="running")
+ return self
+
+ def __exit__(self, exc_type, exc, tb) -> bool:
+ if self.closed:
+ return False
+ if exc_type is not None:
+ self._emit("error", status="error")
+ self._emit(
+ "run_finished",
+ status="error" if exc_type is not None else "success",
+ duration_seconds=max(0.0, time.monotonic() - self.started),
+ )
+ self.closed = True
+ return False
+
+ def tool(self, name: str, *, category: str = "other") -> _TimedOperation:
+ if category not in TOOL_CATEGORIES:
+ raise ValueError(f"unsupported tool category: {category}")
+ return _TimedOperation(self, operation="tool_call", tool_name=name, tool_category=category)
+
+ def model(self, *, model: str = "", usage: Mapping[str, int] | None = None) -> _TimedOperation:
+ return _TimedOperation(self, operation="model_call", model=model, usage=usage)
+
+ def handoff(self) -> bool:
+ return self._emit("handoff", status="success", span_id=_id("span"))
+
+ def approval_requested(self) -> bool:
+ return self._emit("human_approval_requested", status="running")
+
+ def approval_received(self, *, approved: bool) -> bool:
+ return self._emit("human_approval_received", status="success" if approved else "denied")
+
+ def record_error(self) -> bool:
+ return self._emit("error", status="error")
+
+ def _emit(self, operation: str, **fields: Any) -> bool:
+ return self.observer._emit(
+ operation,
+ run_id=self.run_id,
+ trace_id=self.trace_id,
+ workflow_id=self.workflow_id,
+ trigger_event_id=self.trigger_event_id,
+ **fields,
+ )
+
+
+class AgentObserver:
+ """Observe any Python harness without putting OWG in its control path."""
+
+ def __init__(
+ self,
+ agent_name: str,
+ *,
+ provider: str = "",
+ framework: str = "custom",
+ observation_level: str = "instrumented_tools",
+ endpoint: str | None = None,
+ token: str | None = None,
+ device_id: str = "agent-local",
+ sensor_id: str = "agent:python-sdk",
+ sender: Callable[[str, str, list[dict]], None] | None = None,
+ max_queue: int = 512,
+ batch_size: int = 32,
+ flush_interval: float = 0.2,
+ ) -> None:
+ name = _clean(agent_name, limit=160)
+ if not name:
+ raise ValueError("agent_name is required")
+ if observation_level not in OBSERVATION_LEVELS:
+ raise ValueError(f"unsupported observation level: {observation_level}")
+ self.agent_name = name
+ self.provider = _clean(provider, limit=160)
+ self.framework = _clean(framework, limit=160)
+ self.observation_level = observation_level
+ self.device_id = _clean(device_id) or "agent-local"
+ self.sensor_id = _clean(sensor_id) or "agent:python-sdk"
+ self.delivery = _Delivery(
+ endpoint=_clean(endpoint or os.getenv("OWG_AGENT_INGEST_URL") or "http://127.0.0.1:8787/agent-ingest/v1/events", limit=1000),
+ token=_clean(token or os.getenv("OWG_AGENT_INGEST_TOKEN"), limit=4000),
+ sender=sender,
+ max_queue=max_queue,
+ batch_size=batch_size,
+ flush_interval=flush_interval,
+ )
+
+ def run(
+ self,
+ *,
+ run_id: str | None = None,
+ trace_id: str | None = None,
+ workflow_id: str | None = None,
+ trigger_event_id: str | None = None,
+ ) -> AgentRun:
+ return AgentRun(
+ self,
+ run_id=run_id,
+ trace_id=trace_id,
+ workflow_id=workflow_id,
+ trigger_event_id=trigger_event_id,
+ )
+
+ def _emit(self, operation: str, **fields: Any) -> bool:
+ event = {
+ "event_id": _id("evt"),
+ "observed_at": _now(),
+ "agent_name": self.agent_name,
+ "provider": self.provider,
+ "framework": self.framework,
+ "operation": operation,
+ "status": fields.pop("status", "unknown"),
+ "observation_level": self.observation_level,
+ "device_id": self.device_id,
+ "sensor_id": self.sensor_id,
+ "tool_category": fields.get("tool_category") or "none",
+ **fields,
+ }
+ if event["tool_category"] not in TOOL_CATEGORIES:
+ return False
+ return self.delivery.emit(event)
+
+ def flush(self, *, timeout: float = 2.0) -> bool:
+ return self.delivery.flush(timeout)
+
+ def shutdown(self, *, timeout: float = 2.0) -> None:
+ self.delivery.shutdown(timeout)
+
+ def stats(self) -> ObserverStats:
+ return self.delivery.stats()
+
+
+__all__ = ["AgentObserver", "AgentRun", "RunIdentity", "ObserverStats"]
\ No newline at end of file
diff --git a/sdk/python/pyproject.toml b/sdk/python/pyproject.toml
new file mode 100644
index 00000000..0e39a70d
--- /dev/null
+++ b/sdk/python/pyproject.toml
@@ -0,0 +1,14 @@
+[build-system]
+requires = ["setuptools>=68"]
+build-backend = "setuptools.build_meta"
+
+[project]
+name = "openworkgraph-agent"
+version = "0.95.0"
+description = "Dependency-free structural telemetry helper for custom OpenWorkGraph agent harnesses"
+requires-python = ">=3.10"
+license = {text = "Apache-2.0"}
+authors = [{name = "Koyar Afrasyab / Kinvectum"}]
+
+[tool.setuptools]
+py-modules = ["openworkgraph_agent"]
diff --git a/sdk/typescript/index.d.ts b/sdk/typescript/index.d.ts
new file mode 100644
index 00000000..e1335f49
--- /dev/null
+++ b/sdk/typescript/index.d.ts
@@ -0,0 +1,57 @@
+export type ObservationLevel = 'native_trace' | 'instrumented_tools' | 'mcp_only' | 'os_observed' | 'outcome_only';
+export type ToolCategory = 'filesystem' | 'shell' | 'browser' | 'code' | 'search' | 'network' | 'database' | 'messaging' | 'issue_tracker' | 'deployment' | 'mcp' | 'other' | 'none';
+export type TokenUsage = Partial>;
+
+export interface ObserverOptions {
+ provider?: string;
+ framework?: string;
+ observationLevel?: ObservationLevel;
+ endpoint?: string;
+ token?: string;
+ deviceId?: string;
+ sensorId?: string;
+ maxQueue?: number;
+ batchSize?: number;
+ flushMs?: number;
+}
+
+export interface RunOptions {
+ runId?: string;
+ traceId?: string;
+ workflowId?: string;
+ triggerEventId?: string;
+}
+
+export interface ObserverStats {
+ accepted: number;
+ dropped: number;
+ sendFailures: number;
+ batchesSent: number;
+ queued: number;
+}
+
+export class AgentObserver {
+ constructor(agentName: string, options?: ObserverOptions);
+ startRun(options?: RunOptions): AgentRun;
+ withRun(fn: (run: AgentRun) => Promise | T, options?: RunOptions): Promise;
+ flush(): Promise;
+ shutdown(): Promise;
+ stats(): ObserverStats;
+}
+
+export class AgentRun {
+ readonly runId: string;
+ readonly traceId: string;
+ readonly workflowId: string;
+ readonly closed: boolean;
+ start(): this;
+ finish(status?: 'success' | 'error' | 'cancelled'): boolean;
+ tool(name: string, options: {category?: ToolCategory}, fn: () => Promise | T): Promise;
+ tool(name: string, fn: () => Promise | T): Promise;
+ model(options: {model?: string; usage?: TokenUsage}, fn: () => Promise | T): Promise;
+ model(fn: () => Promise | T): Promise;
+ handoff(): boolean;
+ approvalRequested(): boolean;
+ approvalReceived(approved?: boolean): boolean;
+ recordError(): boolean;
+}
diff --git a/sdk/typescript/index.mjs b/sdk/typescript/index.mjs
new file mode 100644
index 00000000..af7829a9
--- /dev/null
+++ b/sdk/typescript/index.mjs
@@ -0,0 +1,175 @@
+// OpenWorkGraph structural agent telemetry for Node.js / TypeScript projects.
+//
+// This module is deliberately content-blind. It never accepts or serializes
+// prompts, responses, tool arguments/results, reasoning, returned values, or
+// exception messages. Delivery is best-effort and fail-open by default.
+
+const OPERATIONS = new Set([
+ 'run_started','run_finished','model_call','tool_call','handoff',
+ 'human_approval_requested','human_approval_received','error'
+]);
+const LEVELS = new Set(['native_trace','instrumented_tools','mcp_only','os_observed','outcome_only']);
+const CATEGORIES = new Set([
+ 'filesystem','shell','browser','code','search','network','database',
+ 'messaging','issue_tracker','deployment','mcp','other','none'
+]);
+
+function clean(value, limit=240){return String(value??'').replace(/\s+/g,' ').trim().slice(0,limit);}
+function identifier(prefix){
+ const c=globalThis.crypto;
+ if(c?.randomUUID)return `${prefix}-${c.randomUUID().replaceAll('-','')}`;
+ return `${prefix}-${Date.now().toString(36)}${Math.random().toString(36).slice(2)}`;
+}
+function now(){return new Date().toISOString();}
+function usage(value={}){
+ const allowed=new Set(['input_tokens','output_tokens','cached_input_tokens','total_tokens']);
+ const out={};
+ for(const [key,raw] of Object.entries(value||{})){
+ if(!allowed.has(key))throw new TypeError(`unsupported usage field: ${key}`);
+ const amount=Number(raw);
+ if(!Number.isSafeInteger(amount)||amount<0)throw new TypeError(`${key} must be a non-negative integer`);
+ out[key]=amount;
+ }
+ return out;
+}
+
+export class AgentObserver {
+ constructor(agentName, options={}){
+ this.agentName=clean(agentName,160);
+ if(!this.agentName)throw new TypeError('agentName is required');
+ this.provider=clean(options.provider,160);
+ this.framework=clean(options.framework||'custom',160);
+ this.observationLevel=clean(options.observationLevel||'instrumented_tools',80);
+ if(!LEVELS.has(this.observationLevel))throw new TypeError(`unsupported observation level: ${this.observationLevel}`);
+ this.endpoint=clean(options.endpoint||process.env.OWG_AGENT_INGEST_URL||'http://127.0.0.1:8787/agent-ingest/v1/events',1000);
+ this.token=clean(options.token||process.env.OWG_AGENT_INGEST_TOKEN||'',4000);
+ this.deviceId=clean(options.deviceId||'agent-local');
+ this.sensorId=clean(options.sensorId||'agent:typescript-sdk');
+ this.maxQueue=Math.max(1,Math.min(Number(options.maxQueue||512),10000));
+ this.batchSize=Math.max(1,Math.min(Number(options.batchSize||32),256));
+ this.flushMs=Math.max(20,Math.min(Number(options.flushMs||200),2000));
+ this.queue=[];this.timer=null;this.flushPromise=null;this.closed=false;
+ this.counts={accepted:0,dropped:0,sendFailures:0,batchesSent:0};
+ }
+
+ startRun(options={}){
+ const run=new AgentRun(this,options);
+ run.start();
+ return run;
+ }
+
+ async withRun(fn, options={}){
+ const run=this.startRun(options);
+ try{
+ return await fn(run);
+ }catch(error){
+ run.recordError();
+ run.finish('error');
+ throw error;
+ }finally{
+ if(!run.closed)run.finish('success');
+ }
+ }
+
+ _enqueue(operation, fields={}){
+ if(this.closed){this.counts.dropped++;return false;}
+ if(!OPERATIONS.has(operation)){this.counts.dropped++;return false;}
+ if(this.queue.length>=this.maxQueue){this.counts.dropped++;return false;}
+ const event={
+ event_id:identifier('evt'), observed_at:now(), agent_name:this.agentName,
+ provider:this.provider, framework:this.framework, operation,
+ status:fields.status||'unknown', observation_level:this.observationLevel,
+ device_id:this.deviceId, sensor_id:this.sensorId,
+ run_id:clean(fields.run_id), trace_id:clean(fields.trace_id),
+ workflow_id:clean(fields.workflow_id), trigger_event_id:clean(fields.trigger_event_id),
+ span_id:clean(fields.span_id), parent_span_id:clean(fields.parent_span_id),
+ tool_name:clean(fields.tool_name,200), tool_category:fields.tool_category||'none',
+ model:clean(fields.model,200), duration_seconds:Math.max(0,Number(fields.duration_seconds||0)),
+ usage:usage(fields.usage||{})
+ };
+ if(!CATEGORIES.has(event.tool_category)){this.counts.dropped++;return false;}
+ this.queue.push(event);this.counts.accepted++;this._schedule();return true;
+ }
+
+ _schedule(){
+ if(this.timer||this.flushPromise||this.closed)return;
+ this.timer=setTimeout(()=>{this.timer=null;void this.flush();},this.flushMs);
+ this.timer.unref?.();
+ }
+
+ async flush(){
+ if(this.flushPromise)return this.flushPromise;
+ if(!this.queue.length)return;
+ const drain=async()=>{
+ while(this.queue.length){
+ const batch=this.queue.splice(0,this.batchSize);
+ if(!this.token){this.counts.sendFailures++;continue;}
+ const controller=new AbortController();
+ const timeout=setTimeout(()=>controller.abort(),1500);timeout.unref?.();
+ try{
+ const response=await fetch(this.endpoint,{
+ method:'POST',headers:{'content-type':'application/json','authorization':`Bearer ${this.token}`},
+ body:JSON.stringify({events:batch}),signal:controller.signal
+ });
+ if(!response.ok)this.counts.sendFailures++;else this.counts.batchesSent++;
+ }catch(_){this.counts.sendFailures++;}
+ finally{clearTimeout(timeout);}
+ }
+ };
+ this.flushPromise=drain();
+ try{await this.flushPromise;}
+ finally{
+ this.flushPromise=null;
+ if(this.queue.length&&!this.closed)this._schedule();
+ }
+ }
+
+ async shutdown(){
+ // Stop admission first, then await an already-running or newly-started
+ // drain. This prevents shutdown from returning while a fetch is in flight.
+ this.closed=true;
+ if(this.timer){clearTimeout(this.timer);this.timer=null;}
+ await this.flush();
+ }
+
+ stats(){return {...this.counts,queued:this.queue.length};}
+}
+
+export class AgentRun {
+ constructor(observer, options={}){
+ this.observer=observer;
+ this.runId=clean(options.runId)||identifier('run');
+ this.traceId=clean(options.traceId)||identifier('trace');
+ this.workflowId=clean(options.workflowId);
+ this.triggerEventId=clean(options.triggerEventId);
+ this.startedAt=0;this.closed=false;
+ }
+ _fields(extra={}){return {run_id:this.runId,trace_id:this.traceId,workflow_id:this.workflowId,trigger_event_id:this.triggerEventId,...extra};}
+ start(){if(!this.startedAt){this.startedAt=performance.now();this.observer._enqueue('run_started',this._fields({status:'running'}));}return this;}
+ finish(status='success'){
+ if(this.closed)return false;this.closed=true;
+ return this.observer._enqueue('run_finished',this._fields({status,duration_seconds:Math.max(0,(performance.now()-this.startedAt)/1000)}));
+ }
+ async tool(name, options={}, fn){
+ if(typeof options==='function'){fn=options;options={};}
+ if(typeof fn!=='function')throw new TypeError('tool requires a function');
+ const category=options.category||'other';if(!CATEGORIES.has(category))throw new TypeError(`unsupported tool category: ${category}`);
+ const started=performance.now();let failed=false;
+ try{return await fn();}
+ catch(error){failed=true;throw error;}
+ finally{
+ this.observer._enqueue('tool_call',this._fields({status:failed?'error':'success',tool_name:clean(name,200),tool_category:category,duration_seconds:(performance.now()-started)/1000}));
+ }
+ }
+ async model(options={}, fn){
+ if(typeof options==='function'){fn=options;options={};}
+ if(typeof fn!=='function')throw new TypeError('model requires a function');
+ const started=performance.now();let failed=false;
+ try{return await fn();}catch(error){failed=true;throw error;}
+ finally{this.observer._enqueue('model_call',this._fields({status:failed?'error':'success',model:clean(options.model,200),usage:usage(options.usage||{}),duration_seconds:(performance.now()-started)/1000}));}
+ }
+ handoff(){return this.observer._enqueue('handoff',this._fields({status:'success',span_id:identifier('span')}));}
+ approvalRequested(){return this.observer._enqueue('human_approval_requested',this._fields({status:'running'}));}
+ approvalReceived(approved=true){return this.observer._enqueue('human_approval_received',this._fields({status:approved?'success':'denied'}));}
+ recordError(){return this.observer._enqueue('error',this._fields({status:'error'}));}
+}
diff --git a/sdk/typescript/package.json b/sdk/typescript/package.json
new file mode 100644
index 00000000..8efd3b82
--- /dev/null
+++ b/sdk/typescript/package.json
@@ -0,0 +1,17 @@
+{
+ "name": "@openworkgraph/agent",
+ "version": "0.95.0",
+ "description": "Dependency-free structural telemetry helper for custom OpenWorkGraph agent harnesses",
+ "type": "module",
+ "exports": {
+ ".": {
+ "types": "./index.d.ts",
+ "import": "./index.mjs"
+ }
+ },
+ "types": "./index.d.ts",
+ "files": ["index.mjs", "index.d.ts"],
+ "license": "Apache-2.0",
+ "engines": {"node": ">=18"},
+ "private": true
+}
diff --git a/server/custom_harness_control_plane.py b/server/custom_harness_control_plane.py
new file mode 100644
index 00000000..c788dc51
--- /dev/null
+++ b/server/custom_harness_control_plane.py
@@ -0,0 +1,198 @@
+from __future__ import annotations
+
+"""Human-reviewed setup material for arbitrary agent harnesses.
+
+Telemetry credentials are write-only. MCP context is a separate connection and
+remains subject to the normal per-run AI-access and saved-history controls.
+"""
+
+import json
+import shlex
+import sys
+from typing import Any
+
+from fastapi import Request
+from fastapi.responses import HTMLResponse, JSONResponse, Response
+
+from server.agent_auth import ensure_agent_ingest_token
+from server.main import ROOT
+from server.secure_app import app
+
+SCRIPT = ROOT / "dashboard" / "custom_harness_setup.js"
+SCRIPT_MARKER = ''
+LAUNCHER = ROOT / "mcp_server" / "launcher.py"
+VERSION_FILE = ROOT / "VERSION"
+
+
+def _base_url(request: Request) -> str:
+ return str(request.base_url).rstrip("/")
+
+
+def _version() -> str:
+ return VERSION_FILE.read_text(encoding="utf-8").strip()
+
+
+def _mcp_config() -> dict[str, Any]:
+ return {
+ "mcpServers": {
+ "openworkgraph": {
+ "command": sys.executable,
+ "args": [str(LAUNCHER), "--client", "custom-harness"],
+ }
+ }
+ }
+
+
+def setup_payload(request: Request) -> dict[str, Any]:
+ base = _base_url(request)
+ token = ensure_agent_ingest_token()
+ version = _version()
+ endpoint = f"{base}/agent-ingest/v1/events"
+ python_install = (
+ 'pip install "openworkgraph-agent @ '
+ f'git+https://github.com/KAVentures/openworkgraph.git@v{version}#subdirectory=sdk/python"'
+ )
+ python_env = "\n".join([
+ f'export OWG_AGENT_INGEST_URL="{endpoint}"',
+ f'export OWG_AGENT_INGEST_TOKEN="{token}"',
+ ])
+ python_example = '''from openworkgraph_agent import AgentObserver
+
+owg = AgentObserver("my-agent", framework="my-harness")
+try:
+ with owg.run(workflow_id="optional-workflow-id") as run:
+ with run.model(model="model-id"):
+ call_model() # return value/content is never sent to OWG
+ with run.tool("repository_search", category="search"):
+ search_repository()
+finally:
+ owg.shutdown()
+'''
+ typescript_download = "\n".join([
+ f"curl -fsSL https://raw.githubusercontent.com/KAVentures/openworkgraph/v{version}/sdk/typescript/index.mjs -o openworkgraph-agent.mjs",
+ f"curl -fsSL https://raw.githubusercontent.com/KAVentures/openworkgraph/v{version}/sdk/typescript/index.d.ts -o openworkgraph-agent.d.ts",
+ ])
+ typescript_env = "\n".join([
+ f'export OWG_AGENT_INGEST_URL="{endpoint}"',
+ f'export OWG_AGENT_INGEST_TOKEN="{token}"',
+ ])
+ typescript_example = '''import {AgentObserver} from "./openworkgraph-agent.mjs";
+
+const owg = new AgentObserver("my-agent", {framework: "my-harness"});
+try {
+ await owg.withRun(async run => {
+ await run.model({model: "model-id"}, async () => callModel());
+ await run.tool("repository_search", {category: "search"}, async () => searchRepository());
+ }, {workflowId: "optional-workflow-id"});
+} finally {
+ await owg.shutdown();
+}
+'''
+ raw_example = {
+ "events": [{
+ "observed_at": "2026-01-01T12:00:00Z",
+ "agent_name": "my-agent",
+ "framework": "my-harness",
+ "operation": "tool_call",
+ "status": "success",
+ "observation_level": "instrumented_tools",
+ "run_id": "run-opaque-id",
+ "tool_name": "repository_search",
+ "tool_category": "search",
+ "duration_seconds": 0.42,
+ }]
+ }
+ return {
+ "version": version,
+ "write": {
+ "credential_scope": "agent_ingest_write_only",
+ "endpoint": endpoint,
+ "authorization": f"Bearer {token}",
+ "python": {
+ "install": python_install,
+ "environment": python_env,
+ "example": python_example,
+ "source": "sdk/python/openworkgraph_agent.py",
+ "dependency_free_runtime": True,
+ },
+ "typescript": {
+ "download": typescript_download,
+ "environment": typescript_env,
+ "example": typescript_example,
+ "source": "sdk/typescript/index.mjs",
+ "types_source": "sdk/typescript/index.d.ts",
+ "dependency_free_runtime": True,
+ "minimum_node": 18,
+ },
+ "raw_http": {
+ "endpoint": endpoint,
+ "authorization": f"Bearer {token}",
+ "example": json.dumps(raw_example, indent=2),
+ },
+ "otel": {
+ "endpoint": f"{base}/agent-ingest/v1/otel",
+ "protocol": "OTLP/HTTP JSON",
+ },
+ },
+ "read": {
+ "method": "MCP stdio",
+ "config": _mcp_config(),
+ "command": f"{shlex.quote(sys.executable)} {shlex.quote(str(LAUNCHER))} --client custom-harness",
+ "master_ai_access_required": True,
+ "saved_history_lease_required_for_historical_reads": True,
+ "note": "Context access is separate from telemetry. Giving a harness the write-only telemetry token never grants it MCP/history read access.",
+ },
+ "privacy": {
+ "prompt_content": False,
+ "model_response_content": False,
+ "tool_arguments": False,
+ "tool_results": False,
+ "reasoning": False,
+ "exception_text": False,
+ "returned_values": False,
+ },
+ }
+
+
+def get_setup(request: Request) -> JSONResponse:
+ return JSONResponse(setup_payload(request), headers={"Cache-Control": "no-store"})
+
+
+def script() -> Response:
+ return Response(SCRIPT.read_text(encoding="utf-8"), media_type="application/javascript")
+
+
+async def _inject(request: Request, call_next):
+ response = await call_next(request)
+ if request.method.upper() != "GET" or request.url.path != "/" or response.status_code != 200:
+ return response
+ if "text/html" not in str(response.headers.get("content-type") or ""):
+ return response
+ try:
+ if hasattr(response, "body_iterator"):
+ chunks = [chunk async for chunk in response.body_iterator]
+ body = b"".join(chunk if isinstance(chunk, bytes) else str(chunk).encode("utf-8") for chunk in chunks)
+ else:
+ body = bytes(getattr(response, "body", b""))
+ text = body.decode("utf-8")
+ except Exception:
+ return response
+ if SCRIPT_MARKER not in text:
+ text = text.replace("