From 934dcf6e096e8b1a3ad220b656acd93fee902c7b Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:13:52 +0200 Subject: [PATCH 01/30] Add public Python SDK for custom agent harnesses --- openworkgraph_agent/__init__.py | 265 ++++++++++++++++++++++++++++++++ 1 file changed, 265 insertions(+) create mode 100644 openworkgraph_agent/__init__.py diff --git a/openworkgraph_agent/__init__.py b/openworkgraph_agent/__init__.py new file mode 100644 index 00000000..62bd169d --- /dev/null +++ b/openworkgraph_agent/__init__.py @@ -0,0 +1,265 @@ +from __future__ import annotations + +"""Tiny public SDK for emitting privacy-safe structural agent telemetry. + +The SDK is intentionally content-blind: its public API has no prompt, response, +tool-argument, tool-result, or reasoning fields. Delivery is fail-open and uses +the same bounded background sink as OpenWorkGraph's native adapters. +""" + +from contextlib import AbstractContextManager +from dataclasses import dataclass +from datetime import datetime, timezone +import time +import uuid +from typing import Any, Callable, Mapping + +from adapters.sdk import BufferedAgentEventSink +from shared.agent_evidence import ( + AGENT_OBSERVATION_LEVELS, + AGENT_TOOL_CATEGORIES, + AgentEvidenceError, +) + + +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: str | None, *, 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 {} + allowed = {"input_tokens", "output_tokens", "cached_input_tokens", "total_tokens"} + unknown = set(value) - allowed + 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 + + +@dataclass(frozen=True) +class RunIdentity: + run_id: str + trace_id: str + workflow_id: str + + +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 = _clean(model, limit=200) + self._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: + duration = max(0.0, time.monotonic() - self._started) + 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, + usage=self._usage, + duration_seconds=duration, + ) + return False + + +class AgentRun(AbstractContextManager["AgentRun"]): + """One harness run. Methods emit structural facts only.""" + + 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_at = 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_at = 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: + # Error existence is useful structural evidence. Exception text is + # deliberately not accepted or transmitted by this SDK. + self._emit("error", status="error") + duration = max(0.0, time.monotonic() - self._started_at) + self._emit( + "run_finished", + status="error" if exc_type is not None else "success", + duration_seconds=duration, + ) + self._closed = True + return False + + def tool(self, name: str, *, category: str = "other") -> _TimedOperation: + if category not in AGENT_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, *, span_id: str | None = None, parent_span_id: str | None = None) -> bool: + return self._emit( + "handoff", + status="success", + span_id=_clean(span_id) or _id("span"), + parent_span_id=_clean(parent_span_id), + ) + + 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: + """Record that an error occurred without accepting error text/content.""" + 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: + """Fail-open structural telemetry for any Python agent or harness. + + A caller may inject ``sender`` in tests or advanced embeddings. Normal use + relies on the installation-local write-only agent ingest credential already + used by OpenWorkGraph's native adapters. + """ + + def __init__( + self, + agent_name: str, + *, + provider: str = "", + framework: str = "custom", + observation_level: str = "instrumented_tools", + device_id: str = "agent-local", + sensor_id: str = "agent:python-sdk", + sender: Callable[[list[dict]], dict] | None = None, + ) -> None: + name = _clean(agent_name, limit=160) + if not name: + raise ValueError("agent_name is required") + if observation_level not in AGENT_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._sink = BufferedAgentEventSink(sender=sender) + + 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, + **fields, + } + # Validate before it reaches the asynchronous sink so programming errors + # are fail-closed for telemetry but remain fail-open for the agent. + try: + return self._sink.emit(event) + except (AgentEvidenceError, TypeError, ValueError): + return False + + def flush(self, *, timeout: float = 2.0) -> bool: + return self._sink.force_flush(timeout=timeout) + + def shutdown(self, *, timeout: float = 2.0) -> None: + self._sink.shutdown(timeout=timeout) + + def stats(self): + return self._sink.stats() + + +__all__ = ["AgentObserver", "AgentRun", "RunIdentity"] From 7ec01f26e362f52ad77fa35eda80d84bf7227db9 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:14:39 +0200 Subject: [PATCH 02/30] Add dependency-free TypeScript and Node harness client --- sdk/typescript/index.mjs | 167 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 167 insertions(+) create mode 100644 sdk/typescript/index.mjs diff --git a/sdk/typescript/index.mjs b/sdk/typescript/index.mjs new file mode 100644 index 00000000..bdd266a5 --- /dev/null +++ b/sdk/typescript/index.mjs @@ -0,0 +1,167 @@ +// 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.sending=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(!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.sending)return; + this.timer=setTimeout(()=>{this.timer=null;void this.flush();},this.flushMs); + this.timer.unref?.(); + } + + async flush(){ + if(this.sending||!this.queue.length)return; + this.sending=true; + try{ + 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);} + } + }finally{this.sending=false;if(this.queue.length)this._schedule();} + } + + async shutdown(){ + 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(); + try{return await fn();} + catch(error){this.observer._enqueue('tool_call',this._fields({status:'error',tool_name:clean(name,200),tool_category:category,duration_seconds:(performance.now()-started)/1000}));throw error;} + finally{ + // Successful calls are emitted here; failed calls were already emitted in catch. + // Returned values and exception content are never serialized. + if(!arguments.callee){} // no-op kept out of payload; see success branch below + } + } + 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'}));} +} From e668d4f3e34ba082a9e27640c6ad62b0ee909d73 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:15:15 +0200 Subject: [PATCH 03/30] Fix TypeScript tool wrapper success telemetry --- sdk/typescript/index.mjs | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) diff --git a/sdk/typescript/index.mjs b/sdk/typescript/index.mjs index bdd266a5..1bb2691d 100644 --- a/sdk/typescript/index.mjs +++ b/sdk/typescript/index.mjs @@ -144,13 +144,11 @@ export class AgentRun { 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(); + const started=performance.now();let failed=false; try{return await fn();} - catch(error){this.observer._enqueue('tool_call',this._fields({status:'error',tool_name:clean(name,200),tool_category:category,duration_seconds:(performance.now()-started)/1000}));throw error;} + catch(error){failed=true;throw error;} finally{ - // Successful calls are emitted here; failed calls were already emitted in catch. - // Returned values and exception content are never serialized. - if(!arguments.callee){} // no-op kept out of payload; see success branch below + 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){ From 15db1cbff28a6c53a4d22e94e61a1b86c4e34090 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:15:35 +0200 Subject: [PATCH 04/30] Add TypeScript declarations for custom harness SDK --- sdk/typescript/index.d.ts | 57 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 57 insertions(+) create mode 100644 sdk/typescript/index.d.ts 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; +} From eba9833442706af203034c4be269c358d8c6c0a9 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:16:21 +0200 Subject: [PATCH 05/30] Make custom harness SDK sources packageable --- sdk/__init__.py | 1 + 1 file changed, 1 insertion(+) create mode 100644 sdk/__init__.py 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.""" From c2aa73933f960a4bfeb56bbd135aceb0c98dad9e Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:16:26 +0200 Subject: [PATCH 06/30] Expose standalone Python harness SDK --- sdk/python/__init__.py | 3 +++ 1 file changed, 3 insertions(+) create mode 100644 sdk/python/__init__.py diff --git a/sdk/python/__init__.py b/sdk/python/__init__.py new file mode 100644 index 00000000..0557b842 --- /dev/null +++ b/sdk/python/__init__.py @@ -0,0 +1,3 @@ +from .openworkgraph_agent import AgentObserver, AgentRun, RunIdentity + +__all__ = ["AgentObserver", "AgentRun", "RunIdentity"] From 79015c0ad7f74f9a41251bd59a99217739fa03e9 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:17:10 +0200 Subject: [PATCH 07/30] Add standalone stdlib-only Python harness SDK --- sdk/python/openworkgraph_agent.py | 390 ++++++++++++++++++++++++++++++ 1 file changed, 390 insertions(+) create mode 100644 sdk/python/openworkgraph_agent.py diff --git a/sdk/python/openworkgraph_agent.py b/sdk/python/openworkgraph_agent.py new file mode 100644 index 00000000..8f087efe --- /dev/null +++ b/sdk/python/openworkgraph_agent.py @@ -0,0 +1,390 @@ +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() + try: + self.queue.put_nowait(dict(event)) + except queue.Full: + with self.lock: + self.dropped += 1 + return False + with self.lock: + 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: + self.stop.set() + self.flush(timeout) + if self.thread is not None and self.thread.is_alive(): + self.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"] From 7d2b1946a2eb196c56e7d48c018ac1c838bbb061 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:17:21 +0200 Subject: [PATCH 08/30] Use standalone public harness SDK implementation --- openworkgraph_agent/__init__.py | 266 +------------------------------- 1 file changed, 3 insertions(+), 263 deletions(-) diff --git a/openworkgraph_agent/__init__.py b/openworkgraph_agent/__init__.py index 62bd169d..963a0d3b 100644 --- a/openworkgraph_agent/__init__.py +++ b/openworkgraph_agent/__init__.py @@ -1,265 +1,5 @@ -from __future__ import annotations +"""Public OpenWorkGraph SDK for custom Python agent harnesses.""" -"""Tiny public SDK for emitting privacy-safe structural agent telemetry. +from sdk.python.openworkgraph_agent import AgentObserver, AgentRun, ObserverStats, RunIdentity -The SDK is intentionally content-blind: its public API has no prompt, response, -tool-argument, tool-result, or reasoning fields. Delivery is fail-open and uses -the same bounded background sink as OpenWorkGraph's native adapters. -""" - -from contextlib import AbstractContextManager -from dataclasses import dataclass -from datetime import datetime, timezone -import time -import uuid -from typing import Any, Callable, Mapping - -from adapters.sdk import BufferedAgentEventSink -from shared.agent_evidence import ( - AGENT_OBSERVATION_LEVELS, - AGENT_TOOL_CATEGORIES, - AgentEvidenceError, -) - - -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: str | None, *, 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 {} - allowed = {"input_tokens", "output_tokens", "cached_input_tokens", "total_tokens"} - unknown = set(value) - allowed - 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 - - -@dataclass(frozen=True) -class RunIdentity: - run_id: str - trace_id: str - workflow_id: str - - -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 = _clean(model, limit=200) - self._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: - duration = max(0.0, time.monotonic() - self._started) - 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, - usage=self._usage, - duration_seconds=duration, - ) - return False - - -class AgentRun(AbstractContextManager["AgentRun"]): - """One harness run. Methods emit structural facts only.""" - - 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_at = 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_at = 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: - # Error existence is useful structural evidence. Exception text is - # deliberately not accepted or transmitted by this SDK. - self._emit("error", status="error") - duration = max(0.0, time.monotonic() - self._started_at) - self._emit( - "run_finished", - status="error" if exc_type is not None else "success", - duration_seconds=duration, - ) - self._closed = True - return False - - def tool(self, name: str, *, category: str = "other") -> _TimedOperation: - if category not in AGENT_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, *, span_id: str | None = None, parent_span_id: str | None = None) -> bool: - return self._emit( - "handoff", - status="success", - span_id=_clean(span_id) or _id("span"), - parent_span_id=_clean(parent_span_id), - ) - - 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: - """Record that an error occurred without accepting error text/content.""" - 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: - """Fail-open structural telemetry for any Python agent or harness. - - A caller may inject ``sender`` in tests or advanced embeddings. Normal use - relies on the installation-local write-only agent ingest credential already - used by OpenWorkGraph's native adapters. - """ - - def __init__( - self, - agent_name: str, - *, - provider: str = "", - framework: str = "custom", - observation_level: str = "instrumented_tools", - device_id: str = "agent-local", - sensor_id: str = "agent:python-sdk", - sender: Callable[[list[dict]], dict] | None = None, - ) -> None: - name = _clean(agent_name, limit=160) - if not name: - raise ValueError("agent_name is required") - if observation_level not in AGENT_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._sink = BufferedAgentEventSink(sender=sender) - - 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, - **fields, - } - # Validate before it reaches the asynchronous sink so programming errors - # are fail-closed for telemetry but remain fail-open for the agent. - try: - return self._sink.emit(event) - except (AgentEvidenceError, TypeError, ValueError): - return False - - def flush(self, *, timeout: float = 2.0) -> bool: - return self._sink.force_flush(timeout=timeout) - - def shutdown(self, *, timeout: float = 2.0) -> None: - self._sink.shutdown(timeout=timeout) - - def stats(self): - return self._sink.stats() - - -__all__ = ["AgentObserver", "AgentRun", "RunIdentity"] +__all__ = ["AgentObserver", "AgentRun", "ObserverStats", "RunIdentity"] From 606fcf0af3815812253f3bf4d7c21a19beb24906 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:17:50 +0200 Subject: [PATCH 09/30] Test standalone custom harness Python SDK --- tests/test_custom_harness_sdk_v095.py | 106 ++++++++++++++++++++++++++ 1 file changed, 106 insertions(+) create mode 100644 tests/test_custom_harness_sdk_v095.py diff --git a/tests/test_custom_harness_sdk_v095.py b/tests/test_custom_harness_sdk_v095.py new file mode 100644 index 00000000..2ed14c77 --- /dev/null +++ b/tests/test_custom_harness_sdk_v095.py @@ -0,0 +1,106 @@ +from __future__ import annotations + +import json + +import pytest + +from openworkgraph_agent import AgentObserver +from shared.agent_evidence import agent_event_to_evidence + + +def test_python_harness_sdk_emits_only_server_compatible_structural_events(): + batches: list[list[dict]] = [] + + def sender(endpoint: str, token: str, events: list[dict]) -> None: + assert endpoint.endswith("/agent-ingest/v1/events") + assert token == "write-only-test-token" + batches.append(events) + + observer = AgentObserver( + "OpenClaw-style local harness", + provider="local", + framework="custom-harness", + endpoint="http://127.0.0.1:8787/agent-ingest/v1/events", + token="write-only-test-token", + sender=sender, + flush_interval=0.02, + ) + with observer.run(run_id="run-1", workflow_id="workflow-1") as run: + with run.model(model="example-model", usage={"input_tokens": 12, "output_tokens": 4}): + pass + with run.tool("repository_search", category="search"): + returned_secret = "THIS_VALUE_MUST_NEVER_BE_TELEMETRY" + assert returned_secret + run.handoff() + run.approval_requested() + run.approval_received(approved=True) + + assert observer.flush(timeout=1.0) is True + observer.shutdown(timeout=1.0) + + events = [event for batch in batches for event in batch] + assert [event["operation"] for event in events] == [ + "run_started", + "model_call", + "tool_call", + "handoff", + "human_approval_requested", + "human_approval_received", + "run_finished", + ] + assert {event["run_id"] for event in events} == {"run-1"} + assert {event["workflow_id"] for event in events} == {"workflow-1"} + assert all(agent_event_to_evidence(event)["source"] == "agent" for event in events) + serialized = json.dumps(events) + assert "THIS_VALUE_MUST_NEVER_BE_TELEMETRY" not in serialized + assert "prompt" not in serialized.lower() + assert "response" not in serialized.lower() + + +def test_python_harness_sdk_never_serializes_exception_text_and_does_not_swallow_it(): + batches: list[list[dict]] = [] + observer = AgentObserver( + "Harness", + token="t", + sender=lambda _endpoint, _token, events: batches.append(events), + flush_interval=0.02, + ) + + with pytest.raises(RuntimeError, match="TOP_SECRET_EXCEPTION_TEXT"): + with observer.run(run_id="run-error") as run: + with run.tool("shell", category="shell"): + raise RuntimeError("TOP_SECRET_EXCEPTION_TEXT") + + assert observer.flush(timeout=1.0) is True + observer.shutdown(timeout=1.0) + events = [event for batch in batches for event in batch] + serialized = json.dumps(events) + assert "TOP_SECRET_EXCEPTION_TEXT" not in serialized + assert any(event["operation"] == "tool_call" and event["status"] == "error" for event in events) + assert any(event["operation"] == "error" for event in events) + assert any(event["operation"] == "run_finished" and event["status"] == "error" for event in events) + + +def test_python_harness_sdk_delivery_failure_is_fail_open(): + def fail(_endpoint: str, _token: str, _events: list[dict]) -> None: + raise OSError("network details must not escape") + + observer = AgentObserver("Harness", token="t", sender=fail, flush_interval=0.02) + with observer.run() as run: + with run.tool("search", category="search"): + pass + assert observer.flush(timeout=1.0) is True + stats = observer.stats() + assert stats.accepted == 3 + assert stats.send_failures >= 1 + observer.shutdown(timeout=1.0) + + +def test_python_harness_sdk_rejects_non_structural_configuration(): + with pytest.raises(ValueError): + AgentObserver("Harness", observation_level="full_transcript") + observer = AgentObserver("Harness", token="t", sender=lambda *_: None) + with observer.run() as run: + with pytest.raises(ValueError): + run.tool("anything", category="prompt_content") + observer.shutdown(timeout=1.0) From 5ce9fdcbdf5b85a76b1b04e55b0e045c9d73c314 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:18:06 +0200 Subject: [PATCH 10/30] Test TypeScript custom harness SDK privacy contract --- tests/js/custom_harness_sdk.test.mjs | 70 ++++++++++++++++++++++++++++ 1 file changed, 70 insertions(+) create mode 100644 tests/js/custom_harness_sdk.test.mjs diff --git a/tests/js/custom_harness_sdk.test.mjs b/tests/js/custom_harness_sdk.test.mjs new file mode 100644 index 00000000..28836f52 --- /dev/null +++ b/tests/js/custom_harness_sdk.test.mjs @@ -0,0 +1,70 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import {AgentObserver} from '../../sdk/typescript/index.mjs'; + +function installFetchCapture(){ + const calls=[]; + const previous=globalThis.fetch; + globalThis.fetch=async (url,options)=>{ + calls.push({url:String(url),options}); + return {ok:true,status:200}; + }; + return {calls,restore(){globalThis.fetch=previous;}}; +} + +test('TypeScript SDK emits structural run/model/tool events without returned content', async ()=>{ + const capture=installFetchCapture(); + try{ + const observer=new AgentObserver('Hermes-style harness',{ + provider:'local',framework:'custom-harness',token:'write-only-token',flushMs:20 + }); + const result=await observer.withRun(async run=>{ + await run.model({model:'example-model',usage:{input_tokens:3,output_tokens:2}},async()=> 'MODEL_SECRET_RETURN'); + return run.tool('repository_search',{category:'search'},async()=> 'TOOL_SECRET_RETURN'); + },{runId:'run-ts-1',workflowId:'workflow-ts-1'}); + assert.equal(result,'TOOL_SECRET_RETURN'); + await observer.shutdown(); + + assert.ok(capture.calls.length>=1); + const bodies=capture.calls.flatMap(call=>JSON.parse(call.options.body).events); + assert.deepEqual(bodies.map(x=>x.operation),['run_started','model_call','tool_call','run_finished']); + assert.ok(bodies.every(x=>x.run_id==='run-ts-1')); + assert.ok(bodies.every(x=>x.workflow_id==='workflow-ts-1')); + const serialized=JSON.stringify(bodies); + assert.equal(serialized.includes('MODEL_SECRET_RETURN'),false); + assert.equal(serialized.includes('TOOL_SECRET_RETURN'),false); + assert.equal(serialized.toLowerCase().includes('prompt'),false); + assert.equal(serialized.toLowerCase().includes('response'),false); + assert.equal(capture.calls[0].options.headers.authorization,'Bearer write-only-token'); + }finally{capture.restore();} +}); + +test('TypeScript SDK records failure structurally but never exception text', async ()=>{ + const capture=installFetchCapture(); + try{ + const observer=new AgentObserver('Custom harness',{token:'t',flushMs:20}); + await assert.rejects( + observer.withRun(async run=>{ + await run.tool('shell',{category:'shell'},async()=>{throw new Error('TOP_SECRET_EXCEPTION');}); + },{runId:'run-error'}), + /TOP_SECRET_EXCEPTION/ + ); + await observer.shutdown(); + const bodies=capture.calls.flatMap(call=>JSON.parse(call.options.body).events); + assert.equal(JSON.stringify(bodies).includes('TOP_SECRET_EXCEPTION'),false); + assert.ok(bodies.some(x=>x.operation==='tool_call'&&x.status==='error')); + assert.ok(bodies.some(x=>x.operation==='error')); + assert.ok(bodies.some(x=>x.operation==='run_finished'&&x.status==='error')); + }finally{capture.restore();} +}); + +test('TypeScript SDK remains fail-open when OWG is unavailable', async ()=>{ + const previous=globalThis.fetch; + globalThis.fetch=async()=>{throw new Error('observer unavailable');}; + try{ + const observer=new AgentObserver('Custom harness',{token:'t',flushMs:20}); + await observer.withRun(async run=>run.tool('search',{category:'search'},async()=>42)); + await observer.shutdown(); + assert.ok(observer.stats().sendFailures>=1); + }finally{globalThis.fetch=previous;} +}); From ff367d70e38f924bbabcd584b9afd79a940d65ff Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:19:38 +0200 Subject: [PATCH 11/30] Make Python agent helper independently installable --- sdk/python/pyproject.toml | 14 ++++++++++++++ 1 file changed, 14 insertions(+) create mode 100644 sdk/python/pyproject.toml 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"] From 7b0853547fee36cd93ead922e97511ef6dcbfb98 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:19:45 +0200 Subject: [PATCH 12/30] Add standalone Node and TypeScript SDK package metadata --- sdk/typescript/package.json | 17 +++++++++++++++++ 1 file changed, 17 insertions(+) create mode 100644 sdk/typescript/package.json 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 +} From 2cdb7e631d5d7f448e3e7cecf0043c6e32c77911 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:21:25 +0200 Subject: [PATCH 13/30] Add custom harness setup control plane --- server/custom_harness_control_plane.py | 202 +++++++++++++++++++++++++ 1 file changed, 202 insertions(+) create mode 100644 server/custom_harness_control_plane.py diff --git a/server/custom_harness_control_plane.py b/server/custom_harness_control_plane.py new file mode 100644 index 00000000..ac4d63b7 --- /dev/null +++ b/server/custom_harness_control_plane.py @@ -0,0 +1,202 @@ +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 pathlib import Path +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 = '' +PYTHON_SDK = ROOT / "sdk" / "python" / "openworkgraph_agent.py" +TYPESCRIPT_SDK = ROOT / "sdk" / "typescript" / "index.mjs" +TYPESCRIPT_TYPES = ROOT / "sdk" / "typescript" / "index.d.ts" +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, + "bundled_source": str(PYTHON_SDK), + "dependency_free_runtime": True, + }, + "typescript": { + "download": typescript_download, + "environment": typescript_env, + "example": typescript_example, + "bundled_source": str(TYPESCRIPT_SDK), + "types_source": str(TYPESCRIPT_TYPES), + "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("", SCRIPT_MARKER + "\n") + headers = dict(response.headers) + headers.pop("content-length", None) + return HTMLResponse(text, status_code=response.status_code, headers=headers) + + +def _install() -> None: + app.add_api_route("/v1/custom-harness-setup", get_setup, methods=["GET"]) + app.add_api_route("/custom-harness-setup.js", script, methods=["GET"]) + app.middleware("http")(_inject) + + +if not getattr(app.state, "owg_custom_harness_control_plane_installed", False): + _install() + app.state.owg_custom_harness_control_plane_installed = True + + +__all__ = ["setup_payload"] From 26ee2862bb109edb31da70cec43615fc7fa9d2d2 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:22:17 +0200 Subject: [PATCH 14/30] Add Connect your own agent dashboard flow --- dashboard/custom_harness_setup.js | 90 +++++++++++++++++++++++++++++++ 1 file changed, 90 insertions(+) create mode 100644 dashboard/custom_harness_setup.js 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-styleInternal agentsAny 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.

`; + grid.appendChild(card); + card.querySelector('#customHarnessSetupButton').onclick=openSetup; + } + + function code(text,id){return `
${esc(text||'')}
`;} + 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 → OpenWorkGraphStructural run/model/tool/handoff/approval/error telemetry through a write-only credential.
OpenWorkGraph → AgentOptional MCP context, still controlled by OWG's AI-access switch and saved-history lease.
${privacy(payload)}

1. Observe this harness

${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(); +})(); From f45d4a2c64d27e65eaf812012f31487b11744870 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:22:28 +0200 Subject: [PATCH 15/30] Register custom harness setup control plane --- server/enterprise_runner.py | 1 + 1 file changed, 1 insertion(+) diff --git a/server/enterprise_runner.py b/server/enterprise_runner.py index 427b1281..11a70392 100644 --- a/server/enterprise_runner.py +++ b/server/enterprise_runner.py @@ -24,6 +24,7 @@ def main() -> None: import server.browser_signal_routes # noqa: F401 import server.browser_agent_projection # noqa: F401 import server.agent_dashboard_control_plane # noqa: F401 + import server.custom_harness_control_plane # noqa: F401 import server.org_join_routes as org_join_routes import server.dashboard_privacy # noqa: F401 From 555cd5bd13e563f3f6c99f4f92e8c7e42ed44e8f Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:23:07 +0200 Subject: [PATCH 16/30] Test custom harness setup permission separation --- .../test_custom_harness_control_plane_v095.py | 87 +++++++++++++++++++ 1 file changed, 87 insertions(+) create mode 100644 tests/test_custom_harness_control_plane_v095.py diff --git a/tests/test_custom_harness_control_plane_v095.py b/tests/test_custom_harness_control_plane_v095.py new file mode 100644 index 00000000..9729f39b --- /dev/null +++ b/tests/test_custom_harness_control_plane_v095.py @@ -0,0 +1,87 @@ +from __future__ import annotations + +import os +import subprocess +import sys +from pathlib import Path + + +ROOT = Path(__file__).resolve().parents[1] + + +def _run_child(code: str, tmp_path: Path) -> None: + env = os.environ.copy() + env.update({ + "WORKFLOW_OBSERVER_DATA": str(tmp_path / "data"), + "WORKFLOW_OBSERVER_AUTH_DIR": str(tmp_path / "auth"), + "WORKFLOW_OBSERVER_MODE": "observe", + "WORKFLOW_OBSERVER_RUN_STARTED_AT": "2026-09-27T20:00:00+00:00", + }) + result = subprocess.run( + [sys.executable, "-c", code], cwd=ROOT, env=env, + text=True, capture_output=True, timeout=45, + ) + assert result.returncode == 0, f"stdout={result.stdout}\nstderr={result.stderr}" + + +def test_custom_harness_setup_has_separate_write_and_mcp_permissions(tmp_path): + code = r''' +from starlette.requests import Request +from server.agent_auth import ensure_agent_ingest_token +import server.custom_harness_control_plane as control + +scope={ + 'type':'http','http_version':'1.1','method':'GET','scheme':'http', + 'path':'/v1/custom-harness-setup','raw_path':b'/v1/custom-harness-setup','query_string':b'', + 'headers':[], 'client':('127.0.0.1',12345), 'server':('127.0.0.1',8787), +} +payload=control.setup_payload(Request(scope)) +token=ensure_agent_ingest_token() +write=payload['write'] +read=payload['read'] +assert write['credential_scope']=='agent_ingest_write_only' +assert write['raw_http']['authorization']==f'Bearer {token}' +assert write['raw_http']['endpoint'].endswith('/agent-ingest/v1/events') +assert write['otel']['endpoint'].endswith('/agent-ingest/v1/otel') +assert 'openworkgraph-agent' in write['python']['install'] +assert 'AgentObserver' in write['python']['example'] +assert 'openworkgraph-agent.mjs' in write['typescript']['download'] +assert 'AgentObserver' in write['typescript']['example'] +assert read['method']=='MCP stdio' +assert read['master_ai_access_required'] is True +assert read['saved_history_lease_required_for_historical_reads'] is True +entry=read['config']['mcpServers']['openworkgraph'] +assert '--client' in entry['args'] and 'custom-harness' in entry['args'] +assert token not in str(read) +assert all(value is False for value in payload['privacy'].values()) +''' + _run_child(code, tmp_path) + + +def test_custom_harness_setup_route_is_human_auth_protected_and_injected(tmp_path): + code = r''' +from fastapi.testclient import TestClient +from server.local_auth import ensure_api_token +import server.agent_dashboard_control_plane +import server.custom_harness_control_plane +from server.secure_app import app + +api_token=ensure_api_token() +with TestClient(app) as client: + denied=client.get('/v1/custom-harness-setup') + assert denied.status_code==401, denied.text + allowed=client.get('/v1/custom-harness-setup',headers={'Authorization':f'Bearer {api_token}'}) + assert allowed.status_code==200, allowed.text + assert allowed.headers.get('cache-control')=='no-store' + root=client.get('/') + assert root.status_code==200 + assert '' in root.text +''' + _run_child(code, tmp_path) + + +def test_enterprise_runner_registers_custom_harness_control_plane(): + source=(ROOT/'server'/'enterprise_runner.py').read_text(encoding='utf-8') + assert 'import server.custom_harness_control_plane' in source + assert source.index('import server.agent_dashboard_control_plane') < source.index('import server.custom_harness_control_plane') + assert source.index('import server.custom_harness_control_plane') < source.index('import server.dashboard_privacy') From d45c873c95221387a5e23a7b83d6eebef4286034 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:24:50 +0200 Subject: [PATCH 17/30] Document arbitrary agent and harness integration --- docs/CUSTOM_HARNESSES.md | 120 +++++++++++++++++++++++++++++++++++++++ 1 file changed, 120 insertions(+) create mode 100644 docs/CUSTOM_HARNESSES.md 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. From c5cd4a8c9985471e9293d8cdd4ef0f9c4815f89a Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:28:03 +0200 Subject: [PATCH 18/30] Make standalone Python SDK shutdown admission race-safe --- sdk/python/openworkgraph_agent.py | 27 +++++++++++++++++---------- 1 file changed, 17 insertions(+), 10 deletions(-) diff --git a/sdk/python/openworkgraph_agent.py b/sdk/python/openworkgraph_agent.py index 8f087efe..d8f76c5a 100644 --- a/sdk/python/openworkgraph_agent.py +++ b/sdk/python/openworkgraph_agent.py @@ -129,13 +129,18 @@ def emit(self, event: dict) -> bool: self.dropped += 1 return False self._ensure_worker() - try: - self.queue.put_nowait(dict(event)) - except queue.Full: - with self.lock: - self.dropped += 1 - return False + # 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 @@ -178,10 +183,12 @@ def flush(self, timeout: float = 2.0) -> bool: return self.queue.unfinished_tasks == 0 def shutdown(self, timeout: float = 2.0) -> None: - self.stop.set() + with self.lock: + self.stop.set() self.flush(timeout) - if self.thread is not None and self.thread.is_alive(): - self.thread.join(timeout=max(0.0, min(float(timeout), 10.0))) + 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: @@ -387,4 +394,4 @@ def stats(self) -> ObserverStats: return self.delivery.stats() -__all__ = ["AgentObserver", "AgentRun", "RunIdentity", "ObserverStats"] +__all__ = ["AgentObserver", "AgentRun", "RunIdentity", "ObserverStats"] \ No newline at end of file From 8edf69c9a62228244a4815c77588b0628f041fe6 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:28:45 +0200 Subject: [PATCH 19/30] Make Node SDK shutdown await in-flight delivery --- sdk/typescript/index.mjs | 22 ++++++++++++++++------ 1 file changed, 16 insertions(+), 6 deletions(-) diff --git a/sdk/typescript/index.mjs b/sdk/typescript/index.mjs index 1bb2691d..af7829a9 100644 --- a/sdk/typescript/index.mjs +++ b/sdk/typescript/index.mjs @@ -48,7 +48,7 @@ export class AgentObserver { 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.sending=false; + this.queue=[];this.timer=null;this.flushPromise=null;this.closed=false; this.counts={accepted:0,dropped:0,sendFailures:0,batchesSent:0}; } @@ -72,6 +72,7 @@ export class AgentObserver { } _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={ @@ -91,15 +92,15 @@ export class AgentObserver { } _schedule(){ - if(this.timer||this.sending)return; + 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.sending||!this.queue.length)return; - this.sending=true; - try{ + 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;} @@ -114,10 +115,19 @@ export class AgentObserver { }catch(_){this.counts.sendFailures++;} finally{clearTimeout(timeout);} } - }finally{this.sending=false;if(this.queue.length)this._schedule();} + }; + 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(); } From 252ea2acb2e8c8a5542c62f185e7d23e2df50377 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:29:22 +0200 Subject: [PATCH 20/30] Test custom Python SDK shutdown admission semantics --- tests/test_custom_harness_sdk_v095.py | 33 +++++++++++++++++++++++++++ 1 file changed, 33 insertions(+) diff --git a/tests/test_custom_harness_sdk_v095.py b/tests/test_custom_harness_sdk_v095.py index 2ed14c77..d97bc117 100644 --- a/tests/test_custom_harness_sdk_v095.py +++ b/tests/test_custom_harness_sdk_v095.py @@ -1,6 +1,7 @@ from __future__ import annotations import json +import threading import pytest @@ -96,6 +97,38 @@ def fail(_endpoint: str, _token: str, _events: list[dict]) -> None: observer.shutdown(timeout=1.0) +def test_python_harness_sdk_shutdown_never_accepts_an_event_after_admission_closes(): + batches: list[list[dict]] = [] + batches_lock = threading.Lock() + decisions: list[bool] = [] + + def sender(_endpoint: str, _token: str, events: list[dict]) -> None: + with batches_lock: + batches.append(events) + + observer = AgentObserver("Harness", token="t", sender=sender, flush_interval=0.02, batch_size=7) + start = threading.Event() + + def emitter() -> None: + start.wait() + run = observer.run(run_id="race-run") + for _ in range(400): + decisions.append(run.record_error()) + + thread = threading.Thread(target=emitter) + thread.start() + start.set() + observer.shutdown(timeout=2.0) + thread.join(timeout=2.0) + assert not thread.is_alive() + + delivered = [event for batch in batches for event in batch] + assert len(delivered) == sum(decisions) + assert observer.stats().accepted == len(delivered) + # Once shutdown returns, admission is permanently closed. + assert observer.run(run_id="after-shutdown").record_error() is False + + def test_python_harness_sdk_rejects_non_structural_configuration(): with pytest.raises(ValueError): AgentObserver("Harness", observation_level="full_transcript") From d5c7325e6e69e03bb2d9725f205ecf115c37fec8 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:29:50 +0200 Subject: [PATCH 21/30] Test Node SDK waits for in-flight delivery on shutdown --- tests/js/custom_harness_sdk.test.mjs | 31 ++++++++++++++++++++++++++++ 1 file changed, 31 insertions(+) diff --git a/tests/js/custom_harness_sdk.test.mjs b/tests/js/custom_harness_sdk.test.mjs index 28836f52..3d12cb5f 100644 --- a/tests/js/custom_harness_sdk.test.mjs +++ b/tests/js/custom_harness_sdk.test.mjs @@ -68,3 +68,34 @@ test('TypeScript SDK remains fail-open when OWG is unavailable', async ()=>{ assert.ok(observer.stats().sendFailures>=1); }finally{globalThis.fetch=previous;} }); + +test('TypeScript shutdown waits for an already in-flight telemetry request', async ()=>{ + const previous=globalThis.fetch; + let releaseFetch; + let fetchStarted; + const started=new Promise(resolve=>{fetchStarted=resolve;}); + const gate=new Promise(resolve=>{releaseFetch=resolve;}); + globalThis.fetch=async()=>{ + fetchStarted(); + await gate; + return {ok:true,status:200}; + }; + try{ + const observer=new AgentObserver('Custom harness',{token:'t',flushMs:20}); + const run=observer.startRun({runId:'shutdown-run'}); + run.finish('success'); + const firstFlush=observer.flush(); + await started; + let shutdownResolved=false; + const shutdown=observer.shutdown().then(()=>{shutdownResolved=true;}); + await new Promise(resolve=>setTimeout(resolve,10)); + assert.equal(shutdownResolved,false); + releaseFetch(); + await firstFlush; + await shutdown; + assert.equal(shutdownResolved,true); + assert.equal(observer.stats().queued,0); + assert.equal(observer.stats().batchesSent,1); + assert.equal(run.recordError(),false); + }finally{globalThis.fetch=previous;} +}); From 2652f9f6fe7d3979895e7b04113afcc2d45dbf76 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:31:54 +0200 Subject: [PATCH 22/30] Export Python SDK observer stats --- sdk/python/__init__.py | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/sdk/python/__init__.py b/sdk/python/__init__.py index 0557b842..5bc7d603 100644 --- a/sdk/python/__init__.py +++ b/sdk/python/__init__.py @@ -1,3 +1,3 @@ -from .openworkgraph_agent import AgentObserver, AgentRun, RunIdentity +from .openworkgraph_agent import AgentObserver, AgentRun, ObserverStats, RunIdentity -__all__ = ["AgentObserver", "AgentRun", "RunIdentity"] +__all__ = ["AgentObserver", "AgentRun", "ObserverStats", "RunIdentity"] From 793c4db01f031f399ecf05628869334154c48d65 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:32:21 +0200 Subject: [PATCH 23/30] Keep custom harness setup paths privacy-minimal --- server/custom_harness_control_plane.py | 10 +++------- 1 file changed, 3 insertions(+), 7 deletions(-) diff --git a/server/custom_harness_control_plane.py b/server/custom_harness_control_plane.py index ac4d63b7..c788dc51 100644 --- a/server/custom_harness_control_plane.py +++ b/server/custom_harness_control_plane.py @@ -9,7 +9,6 @@ import json import shlex import sys -from pathlib import Path from typing import Any from fastapi import Request @@ -21,9 +20,6 @@ SCRIPT = ROOT / "dashboard" / "custom_harness_setup.js" SCRIPT_MARKER = '' -PYTHON_SDK = ROOT / "sdk" / "python" / "openworkgraph_agent.py" -TYPESCRIPT_SDK = ROOT / "sdk" / "typescript" / "index.mjs" -TYPESCRIPT_TYPES = ROOT / "sdk" / "typescript" / "index.d.ts" LAUNCHER = ROOT / "mcp_server" / "launcher.py" VERSION_FILE = ROOT / "VERSION" @@ -116,15 +112,15 @@ def setup_payload(request: Request) -> dict[str, Any]: "install": python_install, "environment": python_env, "example": python_example, - "bundled_source": str(PYTHON_SDK), + "source": "sdk/python/openworkgraph_agent.py", "dependency_free_runtime": True, }, "typescript": { "download": typescript_download, "environment": typescript_env, "example": typescript_example, - "bundled_source": str(TYPESCRIPT_SDK), - "types_source": str(TYPESCRIPT_TYPES), + "source": "sdk/typescript/index.mjs", + "types_source": "sdk/typescript/index.d.ts", "dependency_free_runtime": True, "minimum_node": 18, }, From defc695a6d4894c2bc2ec1f260a1d4abee68705b Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:37:17 +0200 Subject: [PATCH 24/30] Bump OpenWorkGraph to v0.95.0 --- VERSION | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/VERSION b/VERSION index c5c73510..5f8cbfdb 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.94.0 +0.95.0 From d805be76b03ef5a752da12ea0eed98d738a1ec07 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:37:33 +0200 Subject: [PATCH 25/30] Align root package version to v0.95.0 --- pyproject.toml | 39 +++++++++++++++++++++------------------ 1 file changed, 21 insertions(+), 18 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index 3df100b1..b637c379 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"} @@ -30,20 +30,23 @@ gateway = [ "PyJWT[crypto]>=2.10,<3", ] -[tool.pytest.ini_options] -pythonpath = ["."] - -[build-system] -requires = ["setuptools>=69"] -build-backend = "setuptools.build_meta" - -[tool.setuptools.packages.find] -where = ["."] - -[tool.setuptools.package-data] -# Bundled name lists for contextual AI-context redaction (US Census 1990: public -# domain; Statistics Sweden 2022: CC0). See server/name_lexicon/README.md. -server = ["name_lexicon/*.txt", "name_lexicon/README.md"] -# The enterprise Gateway serves this shell at /admin. Every data/action request -# still requires the separately configured Gateway admin credential. -gateway = ["admin_console.html"] +[project.scripts] +workflow-observer = "server.main:main" +workflow-observer-gateway = "server.gateway_runner:main" +workflow-observer-enterprise = "server.enterprise_runner:main" +workflow-observer-enterprise-gateway = "server.enterprise_gateway_runner:main" +workflow-observer-agent-auth = "server.agent_auth:main" +workflow-observer-agent-adapter = "adapters.client:main" +workflow-observer-claude-hook = "adapters.claude_code_hook:main" +workflow-observer-codex-export = "adapters.codex_otel:main" +workflow-observer-openai-agent = "adapters.openai_agents:main" +workflow-observer-policy = "server.policy_admin:main" +workflow-observer-policy-proposal = "server.policy_proposal_cli:main" +workflow-observer-policy-source = "server.policy_source_cli:main" +workflow-observer-policy-drift = "server.policy_drift_cli:main" +workflow-observer-gateway-policy = "server.enterprise_policy_admin:main" +workflow-observer-demo = "server.demo_data:main" +workflow-observer-gateway-enroll = "server.gateway_enroll:main" +workflow-observer-mcp = "mcp_server.secure_stdio:main" +workflow-observer-mcp-compact = "mcp_server.compact_stdio:main" +workflow-observer-mcp-http = "mcp_server.compact_http:main" From ce5f102c13daa6f657aa7f907d7333c00bce3192 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:37:57 +0200 Subject: [PATCH 26/30] Preserve root packaging config while bumping v0.95.0 --- pyproject.toml | 37 +++++++++++++++++-------------------- 1 file changed, 17 insertions(+), 20 deletions(-) diff --git a/pyproject.toml b/pyproject.toml index b637c379..ea4d63cb 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -30,23 +30,20 @@ gateway = [ "PyJWT[crypto]>=2.10,<3", ] -[project.scripts] -workflow-observer = "server.main:main" -workflow-observer-gateway = "server.gateway_runner:main" -workflow-observer-enterprise = "server.enterprise_runner:main" -workflow-observer-enterprise-gateway = "server.enterprise_gateway_runner:main" -workflow-observer-agent-auth = "server.agent_auth:main" -workflow-observer-agent-adapter = "adapters.client:main" -workflow-observer-claude-hook = "adapters.claude_code_hook:main" -workflow-observer-codex-export = "adapters.codex_otel:main" -workflow-observer-openai-agent = "adapters.openai_agents:main" -workflow-observer-policy = "server.policy_admin:main" -workflow-observer-policy-proposal = "server.policy_proposal_cli:main" -workflow-observer-policy-source = "server.policy_source_cli:main" -workflow-observer-policy-drift = "server.policy_drift_cli:main" -workflow-observer-gateway-policy = "server.enterprise_policy_admin:main" -workflow-observer-demo = "server.demo_data:main" -workflow-observer-gateway-enroll = "server.gateway_enroll:main" -workflow-observer-mcp = "mcp_server.secure_stdio:main" -workflow-observer-mcp-compact = "mcp_server.compact_stdio:main" -workflow-observer-mcp-http = "mcp_server.compact_http:main" +[tool.pytest.ini_options] +pythonpath = ["."] + +[build-system] +requires = ["setuptools>=69"] +build-backend = "setuptools.build_meta" + +[tool.setuptools.packages.find] +where = ["."] + +[tool.setuptools.package-data] +# Bundled name lists for contextual AI-context redaction (US Census 1990: public +# domain; Statistics Sweden 2022: CC0). See server/name_lexicon/README.md. +server = ["name_lexicon/*.txt", "name_lexicon/README.md"] +# The enterprise Gateway serves this shell at /admin. Every data/action request +# still requires the separately configured Gateway admin credential. +gateway = ["admin_console.html"] From 30bc83f5b39f00ccfcfb7df68fff805a41514317 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:38:14 +0200 Subject: [PATCH 27/30] Align Claude MCP bundle to v0.95.0 --- mcpb/manifest.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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": {}}}, From ebb16ee3922cb6d0dcbf2a6429a6679a66bb7717 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:38:28 +0200 Subject: [PATCH 28/30] Align v0.95 release contract and SDK asset checks --- tests/test_release_version_v087.py | 11 ++++++++++- 1 file changed, 10 insertions(+), 1 deletion(-) diff --git a/tests/test_release_version_v087.py b/tests/test_release_version_v087.py index 19b00684..e1130288 100644 --- a/tests/test_release_version_v087.py +++ b/tests/test_release_version_v087.py @@ -6,17 +6,21 @@ ROOT = Path(__file__).resolve().parents[1] -EXPECTED_VERSION = "0.94.0" +EXPECTED_VERSION = "0.95.0" def test_release_version_sources_are_aligned(): version = (ROOT / "VERSION").read_text(encoding="utf-8").strip() pyproject = tomllib.loads((ROOT / "pyproject.toml").read_text(encoding="utf-8")) manifest = json.loads((ROOT / "mcpb" / "manifest.json").read_text(encoding="utf-8")) + sdk_python = tomllib.loads((ROOT / "sdk" / "python" / "pyproject.toml").read_text(encoding="utf-8")) + sdk_node = json.loads((ROOT / "sdk" / "typescript" / "package.json").read_text(encoding="utf-8")) assert version == EXPECTED_VERSION assert pyproject["project"]["version"] == EXPECTED_VERSION assert manifest["version"] == EXPECTED_VERSION + assert sdk_python["project"]["version"] == EXPECTED_VERSION + assert sdk_node["version"] == EXPECTED_VERSION def test_release_notes_are_current_and_version_driven(): @@ -35,6 +39,11 @@ def test_release_notes_are_current_and_version_driven(): assert "stable structural identity layer" in workflow assert "Agent setup control plane" in workflow assert "telemetry actually observed" in workflow + assert "Custom harnesses" in workflow + assert "pkg-agent-sdk" in workflow + assert "OpenWorkGraph-Agent-Python.py" in workflow + assert "OpenWorkGraph-Agent-Node.mjs" in workflow + assert "OpenWorkGraph-Agent-Node.d.ts" in workflow def test_generic_otel_docs_use_exact_json_trace_endpoint(): From acbf3391ab405b137f6c1114308defe2d785e306 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:39:04 +0200 Subject: [PATCH 29/30] Package custom agent SDKs in v0.95 releases --- .github/workflows/release.yml | 45 ++++++++++++++++++++++++++++++++++- 1 file changed, 44 insertions(+), 1 deletion(-) 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 From e3dd4de8e5d3c0416a4b3bd1b7c5f4f1924efdaa Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 00:40:01 +0200 Subject: [PATCH 30/30] Document v0.95 custom harness support --- docs/CHANGELOG_V095.md | 38 ++++++++++++++++++++++++++++++++++++++ 1 file changed, 38 insertions(+) create mode 100644 docs/CHANGELOG_V095.md 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.