diff --git a/VERSION b/VERSION index 9188543e..c5c73510 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.93.0 +0.94.0 diff --git a/browser_extension/agent_surface_adapters.js b/browser_extension/agent_surface_adapters.js new file mode 100644 index 00000000..5bb1cb9c --- /dev/null +++ b/browser_extension/agent_surface_adapters.js @@ -0,0 +1,98 @@ +(() => { + 'use strict'; + + const PROVIDERS = [ + {key:'chatgpt', name:'ChatGPT', hosts:['chatgpt.com','chat.openai.com']}, + {key:'claude', name:'Claude', hosts:['claude.ai']}, + {key:'microsoft_copilot', name:'Microsoft Copilot', hosts:['copilot.microsoft.com','copilot.microsoft365.com','m365.cloud.microsoft']}, + {key:'lovable', name:'Lovable', hosts:['lovable.dev']}, + {key:'gemini', name:'Gemini', hosts:['gemini.google.com']}, + ]; + const STOP_RE=/^(stop|cancel)( generating| response| task| run)?$/i; + const SEND_RE=/^(send|submit|ask|run|build|generate)( prompt| message| request)?$/i; + const APPROVE_RE=/^(allow|approve|confirm|continue|accept)$/i; + + function providerForHost(host){ + const h=String(host||'').toLowerCase().replace(/^www\./,''); + return PROVIDERS.find(p=>p.hosts.some(x=>h===x||h.endsWith('.'+x)))||null; + } + function accessibleName(el){ + if(!el)return ''; + return String(el.getAttribute?.('aria-label')||el.getAttribute?.('title')||el.textContent||'').replace(/\s+/g,' ').trim().slice(0,80); + } + function buttonLike(el){return !!el?.closest?.('button,[role="button"],input[type="submit"]');} + function hasBusyState(doc){ + if(!doc)return false; + if(doc.querySelector('[aria-busy="true"]'))return true; + return [...doc.querySelectorAll('button,[role="button"]')].some(el=>STOP_RE.test(accessibleName(el))); + } + function hasVisibleErrorState(doc){return !!doc?.querySelector?.('[role="alert"][aria-live], [role="alert"]');} + function hasApprovalState(doc){ + if(!doc)return false; + for(const root of doc.querySelectorAll('dialog,[role="dialog"]')){ + if([...root.querySelectorAll('button,[role="button"]')].some(el=>APPROVE_RE.test(accessibleName(el))))return true; + } + return false; + } + function isSendControl(el){const control=buttonLike(el);return !!control&&SEND_RE.test(accessibleName(control));} + function isStopControl(el){const control=buttonLike(el);return !!control&&STOP_RE.test(accessibleName(control));} + function isApprovalControl(el){const control=buttonLike(el);return !!control&&APPROVE_RE.test(accessibleName(control))&&!!control.closest('dialog,[role="dialog"]');} + function safeRunId(){ + try{return 'web-'+crypto.randomUUID().replaceAll('-','');}catch(_){return 'web-'+Date.now().toString(36)+Math.random().toString(36).slice(2,12);} + } + function structuralPayload(provider, action, runId, source){ + return { + type:'workflow_observer_event', + observed_at:new Date().toISOString(), + action, + // background.js derives the real safe page URL from sender.tab. Deliberately + // do not send conversation paths, titles, prompts, responses or alert text. + page:{url:location.origin+'/',title:provider.name}, + target:{role:'agent-lifecycle',label:provider.name}, + metadata:{agent_provider:provider.key,agent_run_id:runId,state_source:String(source||'semantic_state').slice(0,40),top_frame:true} + }; + } + function sendLifecycle(provider,action,runId,source){ + // No DOM text or user/model content is included in this message. + const message=structuralPayload(provider,action,runId,source); + try{ + const runtime=(typeof browser!=='undefined'&&browser.runtime)?browser.runtime:(typeof chrome!=='undefined'&&chrome.runtime?chrome.runtime:null); + runtime?.sendMessage(message).catch?.(()=>{}); + }catch(_){} + } + + function start(doc=typeof document!=='undefined'?document:null){ + if(!doc||typeof location==='undefined')return; + const provider=providerForHost(location.hostname);if(!provider)return; + let runId='',wasBusy=false,errorSent=false,approvalSent=false,finishTimer=null; + const begin=source=>{if(runId)return;runId=safeRunId();wasBusy=false;errorSent=false;approvalSent=false;sendLifecycle(provider,'agent_run_started',runId,source);}; + const reset=()=>{runId='';wasBusy=false;errorSent=false;approvalSent=false;if(finishTimer){clearTimeout(finishTimer);finishTimer=null;}}; + const finish=(action='agent_run_finished',source='busy_cleared')=>{if(!runId)return;sendLifecycle(provider,action,runId,source);reset();}; + const sample=()=>{ + const busy=hasBusyState(doc); + if(busy&&!runId)begin('busy_state'); + if(runId){ + if(busy)wasBusy=true; + if(hasVisibleErrorState(doc)&&!errorSent){sendLifecycle(provider,'agent_error',runId,'aria_alert_present');errorSent=true;} + if(hasApprovalState(doc)&&!approvalSent){sendLifecycle(provider,'agent_approval_requested',runId,'semantic_dialog');approvalSent=true;} + if(wasBusy&&!busy&&!finishTimer){finishTimer=setTimeout(()=>{finishTimer=null;if(runId&&!hasBusyState(doc))finish();},900);} + if(busy&&finishTimer){clearTimeout(finishTimer);finishTimer=null;} + } + }; + doc.addEventListener('click',ev=>{ + if(isSendControl(ev.target)){begin('send_control');setTimeout(sample,0);return;} + if(runId&&isStopControl(ev.target)){finish('agent_run_cancelled','stop_control');return;} + if(runId&&isApprovalControl(ev.target)){sendLifecycle(provider,'agent_approval_received',runId,'approval_control');approvalSent=true;} + },true); + doc.addEventListener('submit',()=>{begin('form_submit');setTimeout(sample,0);},true); + const observer=new MutationObserver(()=>sample()); + const root=doc.documentElement||doc;observer.observe(root,{subtree:true,childList:true,attributes:true,attributeFilter:['aria-busy','aria-live','role','open','disabled']}); + setInterval(()=>{if(!doc.hidden)sample();},1500); + sample(); + } + + globalThis.__OWG_AGENT_SURFACE_ADAPTERS_FOR_TESTS__={PROVIDERS,providerForHost,accessibleName,hasBusyState,hasVisibleErrorState,hasApprovalState,isSendControl,isStopControl,isApprovalControl,structuralPayload}; + if(typeof document!=='undefined'&&typeof MutationObserver!=='undefined'){ + if(document.readyState==='loading')document.addEventListener('DOMContentLoaded',()=>start(),{once:true});else start(); + } +})(); diff --git a/browser_extension/manifest.json b/browser_extension/manifest.json index 0df38d5b..edb30b0d 100644 --- a/browser_extension/manifest.json +++ b/browser_extension/manifest.json @@ -1,48 +1,22 @@ { "manifest_version": 3, "name": "Workflow Observer Browser Sensor", - "version": "1.11.0", - "version_name": "1.11.0-v58-metadata-signals", - "description": "Local-only semantic browser telemetry for Workflow Observer. Query values/fragments, typed field values, clipboard contents, filenames, file contents, and ordinary key identities are not collected.", + "version": "1.12.0", + "version_name": "1.12.0-v94-agent-lifecycle", + "description": "Local-only structural browser telemetry for OpenWorkGraph. Query values/fragments, typed field values, clipboard contents, filenames, file contents, prompts and model responses are not collected.", "incognito": "not_allowed", - "permissions": [ - "tabs", - "webNavigation", - "storage", - "alarms" - ], - "host_permissions": [ - "http://127.0.0.1:8787/*", - "http://*/*", - "https://*/*" - ], + "permissions": ["tabs", "webNavigation", "storage", "alarms"], + "host_permissions": ["http://127.0.0.1:8787/*", "http://*/*", "https://*/*"], "background": { - "scripts": [ - "browser_auth.js", - "background.js", - "workflow_enrichment.js", - "browser_signal_enrichment.js" - ], + "scripts": ["browser_auth.js", "background.js", "workflow_enrichment.js", "browser_signal_enrichment.js"], "service_worker": "secure_background.js" }, - "content_scripts": [ - { - "matches": [ - "http://*/*", - "https://*/*" - ], - "js": [ - "content.js", - "content_enrichment.js", - "browser_signals.js" - ], - "run_at": "document_start", - "all_frames": true, - "match_about_blank": true - } - ], - "action": { - "default_title": "Workflow Observer Browser Sensor", - "default_popup": "popup.html" - } + "content_scripts": [{ + "matches": ["http://*/*", "https://*/*"], + "js": ["content.js", "content_enrichment.js", "browser_signals.js", "agent_surface_adapters.js"], + "run_at": "document_start", + "all_frames": true, + "match_about_blank": true + }], + "action": {"default_title": "Workflow Observer Browser Sensor", "default_popup": "popup.html"} } diff --git a/collector/outbox.py b/collector/outbox.py index 37cfa0db..e02971cd 100644 --- a/collector/outbox.py +++ b/collector/outbox.py @@ -8,6 +8,7 @@ from sensitive_identifiers import sanitize_event_identifiers from shared.evidence_deletion import event_overlaps_range +from shared.history_policy import event_kind class EventOutbox: @@ -56,8 +57,6 @@ def _sanitize_existing_pending(self) -> None: safe = sanitize_event_identifiers(event) payload = json.dumps(safe, ensure_ascii=False, separators=(",", ":")) except Exception: - # Do not destroy a delivery queue merely because one legacy - # payload is malformed; new events are always sanitized. continue if payload != row["payload_json"]: conn.execute( @@ -93,6 +92,18 @@ def acknowledge(self, event_ids: list[str]) -> None: with self._lock, self._connect() as conn: conn.execute(f"DELETE FROM pending_events WHERE event_id IN ({placeholders})", tuple(ids)) + def _delete_ids(self, conn: sqlite3.Connection, delete_ids: list[str]) -> int: + if not delete_ids: + return 0 + for offset in range(0, len(delete_ids), 500): + chunk = delete_ids[offset : offset + 500] + placeholders = ",".join("?" for _ in chunk) + conn.execute( + f"DELETE FROM pending_events WHERE event_id IN ({placeholders})", + tuple(chunk), + ) + return len(delete_ids) + def prune_range(self, since: str, until: str) -> int: """Remove queued evidence that overlaps a durable local deletion range.""" with self._lock, self._connect() as conn: @@ -105,16 +116,27 @@ def prune_range(self, since: str, until: str) -> int: continue if isinstance(event, dict) and event_overlaps_range(event, since, until): delete_ids.append(str(row["event_id"])) - if not delete_ids: - return 0 - for offset in range(0, len(delete_ids), 500): - chunk = delete_ids[offset : offset + 500] - placeholders = ",".join("?" for _ in chunk) - conn.execute( - f"DELETE FROM pending_events WHERE event_id IN ({placeholders})", - tuple(chunk), - ) - return len(delete_ids) + return self._delete_ids(conn, delete_ids) + + def prune_sessions(self, kind: str, session_ids: list[str]) -> int: + """Remove queued rows belonging to expired human or agent sessions.""" + selected_kind = "agent" if str(kind).lower() == "agent" else "human" + wanted = {str(value) for value in session_ids if str(value)} + if not wanted: + return 0 + with self._lock, self._connect() as conn: + rows = conn.execute("SELECT event_id, payload_json FROM pending_events").fetchall() + delete_ids: list[str] = [] + for row in rows: + try: + event = json.loads(row["payload_json"]) + except Exception: + continue + if not isinstance(event, dict): + continue + if str(event.get("session_id") or "") in wanted and event_kind(event) == selected_kind: + delete_ids.append(str(row["event_id"])) + return self._delete_ids(conn, delete_ids) def mark_failed(self, event_ids: list[str], error: str) -> None: ids = [str(x) for x in event_ids if x] diff --git a/dashboard/history_retention.js b/dashboard/history_retention.js new file mode 100644 index 00000000..eb3032a7 --- /dev/null +++ b/dashboard/history_retention.js @@ -0,0 +1,87 @@ +(() => { + 'use strict'; + const esc=v=>String(v??'').replace(/[&<>"']/g,c=>({'&':'&','<':'<','>':'>','"':'"',"'":'''}[c])); + const fmt=s=>{s=Math.max(0,Number(s)||0);const h=Math.floor(s/3600),m=Math.floor((s%3600)/60);return h?`${h}h ${m}m`:`${m}m`;}; + let policy=null, history=null, access=null; + + async function call(url,options={}){ + await window.__owgAuthReady; + const r=await fetch(url,{cache:'no-store',...options}); + const d=await r.json().catch(()=>({})); + if(!r.ok)throw new Error(d.detail||'History request failed'); + return d; + } + function parseChoice(value){ + if(value==='ephemeral')return {mode:'ephemeral',days:null}; + if(value==='forever')return {mode:'forever',days:null}; + const days=Number(String(value).replace('days:',''))||90;return {mode:'days',days}; + } + function choiceValue(v){return v?.mode==='days'?`days:${v.days}`:(v?.mode||'ephemeral');} + const options=selected=>['ephemeral','days:30','days:90','days:365','forever'].map(v=>``).join(''); + + function ensureUI(){ + const tabs=document.querySelector('.tabs'); const main=document.querySelector('main'); + if(!tabs||!main||document.querySelector('#tab-history'))return; + const tab=document.createElement('button');tab.type='button';tab.role='tab';tab.id='tab-history';tab.dataset.tab='history';tab.setAttribute('aria-controls','panel-history');tab.setAttribute('aria-selected','false');tab.textContent='History'; + const exportTab=document.querySelector('#tab-export');tabs.insertBefore(tab,exportTab||null); + const panel=document.createElement('section');panel.className='tabpanel';panel.role='tabpanel';panel.id='panel-history';panel.dataset.panel='history';panel.setAttribute('aria-labelledby','tab-history');panel.hidden=true; + panel.innerHTML=` +

History & retention

Capture, keeping history, and letting an AI read saved history are separate choices. Saved history stays on this computer unless you explicitly export or share it.
+

AI access to saved history

Current-session AI access does not automatically include older saved work.
+

Saved sessions

A human session is one recording run. Long idle breaks stay inside it as activity blocks. Agent sessions use observed execution boundaries.
`; + const exportPanel=document.querySelector('#panel-export');main.insertBefore(panel,exportPanel||null); + tab.addEventListener('click',()=>{ if(typeof window.activateTab==='function')window.activateTab('history'); else {document.querySelectorAll('.tabpanel').forEach(x=>x.hidden=x!==panel);panel.hidden=false;} refresh(); }); + panel.querySelector('#historyExportAll').addEventListener('click',()=>exportRange()); + } + + function onboarding(){ + if(policy?.onboarding_complete)return; + const overview=document.querySelector('#panel-overview');if(!overview||document.querySelector('#historyOnboarding'))return; + const box=document.createElement('div');box.id='historyOnboarding';box.className='card';box.style.border='2px solid #3159a5'; + box.innerHTML=`

How long should OpenWorkGraph remember your work?

If you don’t keep history, live context still works, but OpenWorkGraph cannot answer “last month” questions, compare before/after periods, or build findings across days. You can keep human and agent history separately and change this later.
${policy?.upgrade_preserved_existing_history?'
Existing history was preserved during this upgrade. Nothing was deleted automatically.
':''}`; + overview.insertBefore(box,overview.firstChild); + box.querySelector('#keep90').onclick=()=>quickPolicy('days',90); + box.querySelector('#keepNone').onclick=()=>quickPolicy('ephemeral',null); + box.querySelector('#openHistory').onclick=()=>document.querySelector('#tab-history')?.click(); + } + + async function quickPolicy(mode,days){ + await call('/v1/history-policy',{method:'PUT',headers:{'Content-Type':'application/json'},body:JSON.stringify({human_mode:mode,human_days:days,agent_mode:mode,agent_days:days,onboarding_complete:true})}); + document.querySelector('#historyOnboarding')?.remove();await refresh(); + } + + function renderPolicy(){ + const host=document.querySelector('#historyPolicyMount');if(!host||!policy)return; + const h=choiceValue(policy.human_retention),a=choiceValue(policy.agent_retention); + host.innerHTML=`
“Don’t keep” uses temporary local working storage while the session is active, then deletes it. If OpenWorkGraph crashes, stale ephemeral sessions are deleted on the next start.
`; + host.querySelector('#saveHistoryPolicy').onclick=async()=>{const hc=parseChoice(host.querySelector('#humanRetention').value),ac=parseChoice(host.querySelector('#agentRetention').value);await call('/v1/history-policy',{method:'PUT',headers:{'Content-Type':'application/json'},body:JSON.stringify({human_mode:hc.mode,human_days:hc.days,agent_mode:ac.mode,agent_days:ac.days,onboarding_complete:true})});if(typeof window.toast==='function')window.toast('History retention updated.');await refresh();}; + } + + function renderAccess(){ + const host=document.querySelector('#historyAiMount');if(!host||!access)return;const mode=access.mode||'off'; + host.innerHTML=`
${mode==='off'?'Saved-history AI access is OFF.':`Access expires ${esc(access.expires_at||'when revoked')}.`} Turning normal AI access off also prevents all MCP reads.
`; + const sync=()=>{const selected=host.querySelector('#historyAccessMode').value==='selected_range';host.querySelector('#historySince').disabled=!selected;host.querySelector('#historyUntil').disabled=!selected;};host.querySelector('#historyAccessMode').onchange=sync;sync(); + host.querySelector('#saveHistoryAccess').onclick=async()=>{const m=host.querySelector('#historyAccessMode').value;let since=null,until=null;if(m==='selected_range'){const a=host.querySelector('#historySince').value,b=host.querySelector('#historyUntil').value;if(!a||!b){if(typeof window.toast==='function')window.toast('Choose both history dates.');return;}since=new Date(a+'T00:00:00Z').toISOString();until=new Date(b+'T00:00:00Z').toISOString();}await call('/v1/history/ai-access',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({mode:m,since,until,expires_minutes:Number(host.querySelector('#historyExpiry').value)})});if(typeof window.toast==='function')window.toast('Saved-history AI access updated.');await refresh();}; + } + + function sessionLabel(s){if(s.kind==='agent'){const a=s.agent||{};return `${a.framework||a.provider||a.name||'Agent'} · ${s.observation_level||'partial observation'}`;}return `Human recording · ${fmt(s.engaged_seconds)}`;} + function renderSessions(){ + const host=document.querySelector('#historySessions');if(!host||!history)return;const rows=history.sessions||[]; + if(!rows.length){host.innerHTML='
No retained sessions yet.
';return;} + host.innerHTML=`
${rows.map((s,i)=>``).join('')}
SessionWhenEvidenceRetention
${esc(sessionLabel(s))}${s.surfaces?.length?`
${esc(s.surfaces.slice(0,5).join(' · '))}
`:''}
${esc(String(s.started_at||'').replace('T',' ').slice(0,16))}
to ${esc(String(s.ended_at||'').replace('T',' ').slice(0,16))}
${Number(s.event_count||0)} events${s.activity_blocks?.length?` · ${s.activity_blocks.length} blocks`:''}${esc(s.retention_mode||'')} ${s.expires_at?`
expires ${esc(String(s.expires_at).slice(0,10))}
`:''}
`; + host.querySelectorAll('[data-delete]').forEach(b=>b.onclick=async()=>{const s=rows[Number(b.dataset.delete)];if(!confirm('Delete this retained session from this computer? This cannot recall evidence already synchronized to an organization Gateway.'))return;await call('/v1/history/delete-session',{method:'POST',headers:{'Content-Type':'application/json'},body:JSON.stringify({history_session_id:s.history_session_id})});if(typeof window.toast==='function')window.toast('Session deleted.');await refresh();}); + host.querySelectorAll('[data-export]').forEach(b=>b.onclick=()=>{const s=rows[Number(b.dataset.export)];exportRange(s.started_at,s.ended_at);}); + } + + async function exportRange(since=null,until=null){ + const q=new URLSearchParams();if(since)q.set('since',since);if(until)q.set('until',until);q.set('include_raw','false'); + const r=await fetch('/v1/history/export-json?'+q.toString(),{cache:'no-store'});if(!r.ok){if(typeof window.toast==='function')window.toast('History export failed.');return;}const blob=await r.blob(),url=URL.createObjectURL(blob),a=document.createElement('a');a.href=url;a.download='openworkgraph-history.json';document.body.appendChild(a);a.click();a.remove();setTimeout(()=>URL.revokeObjectURL(url),5000); + } + + async function refresh(){ + ensureUI(); + try{[policy,history,access]=await Promise.all([call('/v1/history-policy'),call('/v1/history?limit=200'),call('/v1/history/ai-access')]);access=access.access||{mode:'off'};renderPolicy();renderAccess();renderSessions();onboarding();}catch(e){const host=document.querySelector('#historySessions');if(host)host.innerHTML=`
${esc(e.message)}
`;} + } + window.refreshHistory=refresh; + document.addEventListener('DOMContentLoaded',()=>{ensureUI();refresh();}); +})(); diff --git a/docs/CHANGELOG_V094.md b/docs/CHANGELOG_V094.md new file mode 100644 index 00000000..1f7c0f18 --- /dev/null +++ b/docs/CHANGELOG_V094.md @@ -0,0 +1,30 @@ +# OpenWorkGraph v0.94.0 + +## History and retention + +- Adds explicit local retention policy for human and agent evidence with separate modes: ephemeral, N days, or forever. +- New installs default to ephemeral until onboarding is completed. Existing installs preserve already-retained history on upgrade rather than deleting it silently. +- Adds a History view/API for browsing retained human sessions and agent executions, deleting one session, and exporting a selected retained time range. +- Session deletion and expiry install durable session tombstones before cleanup so buffered/retried events cannot recreate deleted history. +- Context Pulse cursors are invalidated when retained history or saved-history access changes. + +## Saved-history AI access + +- Retention and AI disclosure are separate controls. +- Saved history requires an explicit time-limited AI history lease (`selected_range` or `all_saved`). Current-session context continues to work without a saved-history lease. +- The MCP transport centrally intersects historical reads with the allowed date range so older tools cannot bypass the lease. +- Adds `list_history` for metadata-first navigation without eagerly loading full historical evidence. + +## Universal structural agent observation + +- Adds browser-surface lifecycle adapters for ChatGPT, Claude web, Microsoft Copilot, Lovable, and Gemini. +- Browser-observed agent runs are projected into the existing vendor-neutral agent evidence contract with `observation_level=os_observed`. +- The adapters emit structural lifecycle only: run started/finished/cancelled, approval requested/received, and visible error state. +- They do not capture prompts, model responses, DOM content, tool arguments/results, typed text, clipboard contents, filenames/file contents, or hidden reasoning. +- No new browser permissions are required. + +## Versioning + +- Root package version: `0.94.0`. +- Claude MCP bundle version: `0.94.0`. +- Browser sensor version: `1.12.0-v94-agent-lifecycle`. diff --git a/mcp_server/compact_http_app.py b/mcp_server/compact_http_app.py index cf92ea30..1b0f94bd 100644 --- a/mcp_server/compact_http_app.py +++ b/mcp_server/compact_http_app.py @@ -1,42 +1,29 @@ from __future__ import annotations import os - from starlette.responses import JSONResponse - from server.local_auth import mcp_bearer_matches from . import compact as _compact +from . import secure_runtime as _secure_runtime from .compact_hardening import apply_compact_hardening +from .history_guard import install_history_guard +from .history_tools import register_history_tools - +install_history_guard(_secure_runtime) apply_compact_hardening(_compact) mcp = _compact.mcp +register_history_tools(mcp, _secure_runtime) _inner = mcp.streamable_http_app() - class MCPBearerGuard: - def __init__(self, app): - self.app = app - + def __init__(self, app): self.app = app async def __call__(self, scope, receive, send): - if scope.get("type") != "http": - await self.app(scope, receive, send) - return - headers = {k.decode("latin1").lower(): v.decode("latin1") for k, v in scope.get("headers", [])} + if scope.get("type") != "http": await self.app(scope, receive, send); return + headers={k.decode("latin1").lower():v.decode("latin1") for k,v in scope.get("headers",[])} if not mcp_bearer_matches(headers.get("authorization")): - response = JSONResponse({"detail": "OpenWorkGraph MCP authentication required"}, status_code=401) - await response(scope, receive, send) - return - - if scope.get("method") == "GET" and scope.get("path") == "/openworkgraph-id": - response = JSONResponse({ - "server": "OpenWorkGraph", - "instance_nonce": os.getenv("WORKFLOW_OBSERVER_MCP_INSTANCE_NONCE", ""), - }) - await response(scope, receive, send) - return - - await self.app(scope, receive, send) - + await JSONResponse({"detail":"OpenWorkGraph MCP authentication required"},status_code=401)(scope,receive,send); return + if scope.get("method")=="GET" and scope.get("path")=="/openworkgraph-id": + await JSONResponse({"server":"OpenWorkGraph","instance_nonce":os.getenv("WORKFLOW_OBSERVER_MCP_INSTANCE_NONCE","")})(scope,receive,send); return + await self.app(scope,receive,send) -app = MCPBearerGuard(_inner) \ No newline at end of file +app=MCPBearerGuard(_inner) diff --git a/mcp_server/compact_stdio.py b/mcp_server/compact_stdio.py index 5907fa8f..6e4ac341 100644 --- a/mcp_server/compact_stdio.py +++ b/mcp_server/compact_stdio.py @@ -1,12 +1,17 @@ from __future__ import annotations from . import compact as _compact +from . import secure_runtime as _secure_runtime from .compact_hardening import apply_compact_hardening +from .history_guard import install_history_guard +from .history_tools import register_history_tools +install_history_guard(_secure_runtime) apply_compact_hardening(_compact) mcp = _compact.mcp +register_history_tools(mcp, _secure_runtime) if __name__ == "__main__": - mcp.run(transport="stdio") \ No newline at end of file + mcp.run(transport="stdio") diff --git a/mcp_server/history_guard.py b/mcp_server/history_guard.py new file mode 100644 index 00000000..ab029f17 --- /dev/null +++ b/mcp_server/history_guard.py @@ -0,0 +1,142 @@ +from __future__ import annotations + +"""Central saved-history permission guard for MCP transports. + +Local retention and AI disclosure are independent. The guard wraps the existing +secure local API transport so old and new MCP tools cannot bypass a user's saved- +history choice merely because a particular endpoint predates History settings. +""" + +from datetime import datetime +from typing import Any + +from mcp.server.mcpserver.exceptions import ToolError + + +_FULL_HISTORY_PATHS = { + "/v1/events", + "/v1/context-events", + "/v1/operational-events", +} +_RANGELESS_HISTORICAL_PREFIXES = ( + "/v1/procedural-memory", + "/v1/task-context", +) +_SESSION_PREFIXES = ( + "/v1/sessions/", + "/v1/context-sessions/", + "/v1/operational-sessions/", +) +_SCOPE_PATHS = { + "/v1/tasks", + "/v1/summary", + "/v1/operational-summary", + "/v1/semantic-activity", + "/v1/operational-semantic-activity", + "/v1/work-profile", + "/v1/patterns", +} + + +def _parse(value: Any) -> datetime | None: + try: + parsed = datetime.fromisoformat(str(value or "").replace("Z", "+00:00")) + return parsed + except Exception: + return None + + +def _iso_max(left: str | None, right: str | None) -> str | None: + a, b = _parse(left), _parse(right) + if a is None: + return right + if b is None: + return left + return (a if a >= b else b).isoformat() + + +def _iso_min(left: str | None, right: str | None) -> str | None: + a, b = _parse(left), _parse(right) + if a is None: + return right + if b is None: + return left + return (a if a <= b else b).isoformat() + + +def install_history_guard(runtime_module: Any) -> None: + if getattr(runtime_module, "_owg_history_guard_installed", False): + return + original_get = runtime_module.secure_get + + def access() -> dict[str, Any]: + try: + payload = original_get("/v1/history/ai-access") + value = payload.get("access") if isinstance(payload, dict) else None + return dict(value) if isinstance(value, dict) else {"mode": "off"} + except Exception: + return {"mode": "off"} + + def current_start() -> str | None: + try: + status = original_get("/v1/capture/status") + return str(status.get("run_started_at") or "") or None + except Exception: + return None + + def require_all_saved(path: str) -> None: + if str(access().get("mode") or "off") != "all_saved": + raise ToolError( + "Saved-history AI access does not cover this aggregate. In OpenWorkGraph History, grant " + "All saved history, or use list_history/get_workflow_trace/get_agent_runs with an allowed date range." + ) + + def bounded_range(params: dict[str, Any], *, use_current_when_off: bool) -> dict[str, Any]: + result = dict(params) + lease = access() + mode = str(lease.get("mode") or "off") + if mode == "all_saved": + return result + if mode == "selected_range": + result["since"] = _iso_max(str(result.get("since") or "") or None, str(lease.get("since") or "") or None) + result["until"] = _iso_min(str(result.get("until") or "") or None, str(lease.get("until") or "") or None) + start, end = _parse(result.get("since")), _parse(result.get("until")) + if start is not None and end is not None and end <= start: + raise ToolError("The requested history range is outside the saved-history range you allowed.") + return result + if not use_current_when_off: + raise ToolError("Saved-history AI access is OFF. Grant a date range in OpenWorkGraph History first.") + start = current_start() + if not start: + raise ToolError("OpenWorkGraph could not establish the current recording boundary; saved-history access remains off.") + result["since"] = _iso_max(str(result.get("since") or "") or None, start) + return result + + def guarded_get(path: str, params: dict[str, Any] | None = None) -> dict[str, Any]: + route = str(path or "") + query = dict(params or {}) + + if route == "/v1/history": + query = bounded_range(query, use_current_when_off=False) + elif route == "/v1/workflow-trace": + if str(query.get("scope") or "current") != "current": + query = bounded_range(query, use_current_when_off=True) + elif route == "/v1/agent-execution-traces": + query = bounded_range(query, use_current_when_off=True) + elif route in _SCOPE_PATHS: + scope = str(query.get("scope") or "current") + if scope != "current": + require_all_saved(route) + elif route in _FULL_HISTORY_PATHS or route.startswith(_RANGELESS_HISTORICAL_PREFIXES) or route.startswith(_SESSION_PREFIXES): + require_all_saved(route) + + return original_get(route, query if query else None) + + runtime_module.secure_get = guarded_get + # Legacy tools call the transport through mcp_server.main's injected hook. + if hasattr(runtime_module, "core"): + runtime_module.core._get = guarded_get + runtime_module._owg_history_guard_installed = True + + +__all__ = ["install_history_guard"] diff --git a/mcp_server/history_tools.py b/mcp_server/history_tools.py new file mode 100644 index 00000000..961f75ef --- /dev/null +++ b/mcp_server/history_tools.py @@ -0,0 +1,36 @@ +from __future__ import annotations + +from typing import Any + + +def register_history_tools(mcp: Any, runtime: Any) -> None: + if getattr(mcp, "_owg_history_tools_registered", False): + return + + @mcp.tool() + def list_history( + since: str | None = None, + until: str | None = None, + limit: int = 100, + ) -> dict[str, Any]: + """List retained human/agent sessions before drilling into canonical evidence. + + Saved-history access is user-controlled and time-limited. This tool returns + factual session metadata only: dates, activity blocks, surfaces, durations, + retention state and observed agent coverage. It does not infer task names. + After selecting relevant dates, use get_workflow_trace or get_agent_runs. + """ + name = "list_history" + runtime.core._begin(name) + params: dict[str, Any] = {"limit": max(1, min(int(limit), 500))} + if since: params["since"] = since + if until: params["until"] = until + result = runtime.secure_get("/v1/history", params) + result["next_step"] = "Use get_workflow_trace for canonical human/work evidence and get_agent_runs for structural agent evidence in the relevant dates." + result["task_labels_inferred"] = False + return runtime.core._finish(name, result) + + setattr(mcp, "_owg_history_tools_registered", True) + + +__all__ = ["register_history_tools"] diff --git a/mcp_server/http_app.py b/mcp_server/http_app.py index 0ac48b4d..b3e8afda 100644 --- a/mcp_server/http_app.py +++ b/mcp_server/http_app.py @@ -1,44 +1,26 @@ from __future__ import annotations import os - from starlette.responses import JSONResponse - from server.local_auth import mcp_bearer_matches -from .secure_runtime import mcp +from . import secure_runtime as _secure_runtime from .agent_tools import register_agent_tools +from .history_guard import install_history_guard - +install_history_guard(_secure_runtime) +mcp = _secure_runtime.mcp register_agent_tools(mcp) _inner = mcp.streamable_http_app() - class MCPBearerGuard: - def __init__(self, app): - self.app = app - + def __init__(self, app): self.app = app async def __call__(self, scope, receive, send): - if scope.get("type") != "http": - await self.app(scope, receive, send) - return - headers = {k.decode("latin1").lower(): v.decode("latin1") for k, v in scope.get("headers", [])} + if scope.get("type") != "http": await self.app(scope, receive, send); return + headers={k.decode("latin1").lower():v.decode("latin1") for k,v in scope.get("headers",[])} if not mcp_bearer_matches(headers.get("authorization")): - response = JSONResponse({"detail": "OpenWorkGraph MCP authentication required"}, status_code=401) - await response(scope, receive, send) - return - - # Private launch-time proof used only by OpenWorkGraph's controller before - # it advertises a newly spawned HTTP endpoint. A listener that merely wins - # the port race cannot produce this nonce. - if scope.get("method") == "GET" and scope.get("path") == "/openworkgraph-id": - response = JSONResponse({ - "server": "OpenWorkGraph", - "instance_nonce": os.getenv("WORKFLOW_OBSERVER_MCP_INSTANCE_NONCE", ""), - }) - await response(scope, receive, send) - return - - await self.app(scope, receive, send) - + await JSONResponse({"detail":"OpenWorkGraph MCP authentication required"},status_code=401)(scope,receive,send); return + if scope.get("method")=="GET" and scope.get("path")=="/openworkgraph-id": + await JSONResponse({"server":"OpenWorkGraph","instance_nonce":os.getenv("WORKFLOW_OBSERVER_MCP_INSTANCE_NONCE","")})(scope,receive,send); return + await self.app(scope,receive,send) -app = MCPBearerGuard(_inner) \ No newline at end of file +app=MCPBearerGuard(_inner) diff --git a/mcp_server/secure_stdio.py b/mcp_server/secure_stdio.py index db9c9ea1..7bc5db53 100644 --- a/mcp_server/secure_stdio.py +++ b/mcp_server/secure_stdio.py @@ -1,11 +1,14 @@ from __future__ import annotations +from . import secure_runtime as _secure_runtime from .secure_runtime import mcp from .agent_tools import register_agent_tools +from .history_guard import install_history_guard +install_history_guard(_secure_runtime) register_agent_tools(mcp) if __name__ == "__main__": - mcp.run(transport="stdio") \ No newline at end of file + mcp.run(transport="stdio") diff --git a/mcpb/manifest.json b/mcpb/manifest.json index 79f9095f..f6c1fbc7 100644 --- a/mcpb/manifest.json +++ b/mcpb/manifest.json @@ -2,41 +2,26 @@ "manifest_version": "0.3", "name": "openworkgraph-local", "display_name": "OpenWorkGraph", - "version": "0.93.0", + "version": "0.94.0", "description": "Connect Claude Desktop to the compact local OpenWorkGraph context surface.", - "long_description": "Uses the OpenWorkGraph installation already running on this computer. New connections expose canonical workflow evidence first through the compact read-oriented Context MCP surface, while Context Pulse provides incremental factual updates and changed long-horizon findings. Derived task and pattern views remain optional and non-authoritative. The legacy 24-tool stdio entrypoint remains available for existing configurations. Workflow evidence remains in the local OpenWorkGraph store until Claude requests it through MCP. AI access can be turned off instantly from the OpenWorkGraph dashboard.", - "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": {} - } - }, + "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.", + "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": {}}}, "tools": [ - {"name": "get_current_work_context", "description": "Read an optional compact derived overview with pointers to canonical evidence."}, - {"name": "get_context_pulse", "description": "Read what changed since the last check plus new or materially changed factual findings; pass the returned cursor back next time."}, - {"name": "search_work", "description": "Search prior work through evidence or semantic layers."}, - {"name": "get_workflow_trace", "description": "Read the primary canonical paginated workflow evidence, optionally for one session."}, - {"name": "get_work_profile", "description": "Read derived workflow signals without productivity scoring."}, - {"name": "find_repeated_workflows", "description": "Find repeated-work candidates and exact procedural-memory family keys for follow-up."}, - {"name": "get_task_context", "description": "Read a bounded task-context bundle with provenance and authority separation."}, - {"name": "how_did_similar_runs_go", "description": "Read descriptive prior-run feedback with privacy-safe human step labels, failure points, approval hotspots and observed next steps."}, - {"name": "get_agent_runs", "description": "List privacy-safe agent runs or inspect one structural trace by execution ID."} + {"name":"get_current_work_context","description":"Read an optional compact derived overview with pointers to canonical evidence."}, + {"name":"get_context_pulse","description":"Read what changed since the last check plus new or materially changed factual findings; pass the returned cursor back next time."}, + {"name":"list_history","description":"Navigate retained human and agent sessions within the saved-history date range the user explicitly allowed."}, + {"name":"search_work","description":"Search prior work through evidence or semantic layers."}, + {"name":"get_workflow_trace","description":"Read the primary canonical paginated workflow evidence, optionally for one session or date range."}, + {"name":"get_work_profile","description":"Read derived workflow signals without productivity scoring."}, + {"name":"find_repeated_workflows","description":"Find repeated-work candidates and exact procedural-memory family keys for follow-up."}, + {"name":"get_task_context","description":"Read a bounded task-context bundle with provenance and authority separation."}, + {"name":"how_did_similar_runs_go","description":"Read descriptive prior-run feedback with privacy-safe human step labels, failure points, approval hotspots and observed next steps."}, + {"name":"get_agent_runs","description":"List privacy-safe structural agent runs; saved-history leases bound historical date access."} ], - "compatibility": { - "platforms": ["darwin", "win32"], - "runtimes": {"node": ">=18.0.0"} - }, - "keywords": ["workflow", "context", "productivity", "local", "mcp"], - "license": "Apache-2.0", - "privacy_policies": [] + "compatibility": {"platforms":["darwin","win32"],"runtimes":{"node":">=18.0.0"}}, + "keywords":["workflow","context","local","mcp","history","agents"], + "license":"Apache-2.0", + "privacy_policies":[] } diff --git a/pyproject.toml b/pyproject.toml index f98ca3df..3df100b1 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "workflow-observer" -version = "0.93.0" +version = "0.94.0" description = "Local-first work evidence, self-hosted organizational context gateway, REST API, and MCP access." requires-python = ">=3.11" license = {file = "LICENSE"} diff --git a/server/agent_execution_trace_routes.py b/server/agent_execution_trace_routes.py index 93abb6f2..2ae6175a 100644 --- a/server/agent_execution_trace_routes.py +++ b/server/agent_execution_trace_routes.py @@ -1,5 +1,6 @@ from __future__ import annotations +from datetime import datetime import json from typing import Any @@ -11,7 +12,6 @@ from .agent_tool_labels import readable_tool_name from .procedural_memory import load_recent_evidence - router = APIRouter() @@ -21,36 +21,32 @@ def _require_agent_report_read(request: Request) -> None: def _row_framework(row: dict[str, Any]) -> str: - if str(row.get("source") or "") != "agent": - return "" + if str(row.get("source") or "") != "agent": return "" meta = row.get("metadata") if meta is None and row.get("metadata_json"): - try: - meta = json.loads(row["metadata_json"]) - except (TypeError, ValueError): - return "" + try: meta = json.loads(row["metadata_json"]) + except (TypeError, ValueError): return "" agent = meta.get("agent") if isinstance(meta, dict) else None return str((agent or {}).get("framework") or "") if isinstance(agent, dict) else "" +def _parse(value: Any) -> datetime | None: + try: return datetime.fromisoformat(str(value or "").replace("Z", "+00:00")) + except Exception: return None + + def _sanitize_tool_names(payload: dict[str, Any]) -> dict[str, Any]: - """Keep public runtime vocabulary readable while hiding custom identifiers.""" executions = payload.get("executions") if isinstance(payload, dict) else None - if not isinstance(executions, list): - return payload + if not isinstance(executions, list): return payload for execution in executions: - if not isinstance(execution, dict): - continue + if not isinstance(execution, dict): continue events = execution.get("events") - if not isinstance(events, list): - continue + if not isinstance(events, list): continue for event in events: - if not isinstance(event, dict): - continue + if not isinstance(event, dict): continue tool = event.get("tool") - if not isinstance(tool, dict) or not tool.get("name"): - continue - tool["name"] = readable_tool_name(tool.get("name")) + if isinstance(tool, dict) and tool.get("name"): + tool["name"] = readable_tool_name(tool.get("name")) return payload @@ -60,35 +56,27 @@ def get_agent_execution_traces( family_key: str = "", execution_id: str = "", since: str | None = None, + until: str | None = None, evidence_limit: int = 25_000, limit: int = 20, max_events_per_execution: int = 100, hide_disconnected: bool = False, ) -> dict[str, Any]: - """Return privacy-safe ordered structural traces for observed agent runs. - - ``hide_disconnected`` omits stored runs from agents whose Observe connection - is currently off or removed (history is kept, only the view is filtered). - """ + """Return privacy-safe ordered structural traces for observed agent runs.""" _require_agent_report_read(request) + if since and _parse(since) is None: raise HTTPException(status_code=422, detail="invalid since") + end = _parse(until) if until else None + if until and end is None: raise HTTPException(status_code=422, detail="invalid until") hidden: set[str] = set() try: - raw = load_recent_evidence( - limit=max(1, min(int(evidence_limit), 100_000)), - since=since, - ) + raw = load_recent_evidence(limit=max(1, min(int(evidence_limit), 100_000)), since=since) + if end is not None: + raw = [row for row in raw if (_parse(row.get("observed_at")) is not None and _parse(row.get("observed_at")) < end)] if hide_disconnected: from .connections import hidden_frameworks hidden = hidden_frameworks() - if hidden: - raw = [row for row in raw if _row_framework(row) not in hidden] - payload = agent_execution_traces( - raw, - family_key=family_key, - execution_id=execution_id, - limit=limit, - max_events_per_execution=max_events_per_execution, - ) + if hidden: raw = [row for row in raw if _row_framework(row) not in hidden] + payload = agent_execution_traces(raw, family_key=family_key, execution_id=execution_id, limit=limit, max_events_per_execution=max_events_per_execution) payload = enrich_agent_execution_payload(payload, raw) payload = _sanitize_tool_names(payload) except (TypeError, ValueError) as exc: @@ -96,6 +84,8 @@ def get_agent_execution_traces( return { **payload, "evidence_rows_considered": len(raw), + "since": since, + "until": until, "hidden_frameworks": sorted(hidden), "evidence_is_canonical": True, "read_only": True, diff --git a/server/agent_ingest.py b/server/agent_ingest.py index c0229eff..fa7f1bd7 100644 --- a/server/agent_ingest.py +++ b/server/agent_ingest.py @@ -10,6 +10,7 @@ from shared.capture_control import filter_recordable from shared.claude_otel_adapter import claude_otel_to_agent_events from shared.codex_otel_adapter import codex_otel_to_agent_events +from shared.history_policy import retention_for_kind from shared.otel_agent_adapter import otel_payload_to_agent_events from .db import insert_events @@ -30,18 +31,36 @@ def _bounded_json_size(value: Any, *, maximum: int = MAX_AGENT_BATCH_BYTES) -> N def _validated_events(payloads: list[dict[str, Any]]) -> list[dict[str, Any]]: - # Validate all untrusted events before projecting or writing any of them. This - # preserves the existing all-or-nothing batch behavior for malformed input. return [agent_event_to_evidence(validate_agent_ingress_event(item)) for item in payloads] +def _closed_agent_sessions(events: list[dict[str, Any]]) -> list[str]: + closed: list[str] = [] + for event in events: + metadata = event.get("metadata") if isinstance(event.get("metadata"), dict) else {} + if str(metadata.get("operation") or "") != "run_finished": + continue + session_id = str(event.get("session_id") or "").strip() + if session_id: + closed.append(session_id) + return list(dict.fromkeys(closed)) + + def _insert_recordable(events: list[dict[str, Any]]) -> int: - # Agent evidence obeys the same user-controlled Pause/Stop boundaries and - # permanent deletion tombstones as desktop/browser evidence. This is applied - # immediately before persistence so late or buffered agent delivery cannot - # recreate skipped/deleted work. + # Agent evidence obeys deletion, retention, Pause/Stop and crash-recovery + # boundaries immediately before persistence. For ephemeral agent history, an + # observed run_finished boundary closes and purges the corresponding session. recordable, _suppressed = filter_recordable(events) - return insert_events(recordable) if recordable else 0 + if not recordable: + return 0 + inserted = insert_events(recordable) + if retention_for_kind("agent").get("mode") == "ephemeral": + closed = _closed_agent_sessions(recordable) + if closed: + from .history_retention import cleanup_ephemeral_session + for session_id in closed: + cleanup_ephemeral_session("agent", session_id) + return inserted def ingest_agent_payloads(payloads: list[dict[str, Any]]) -> dict[str, int]: @@ -62,8 +81,6 @@ def ingest_otel_payload( *, defaults: dict[str, Any] | None = None, ) -> dict[str, int]: - # The advertised request bound covers both the OTLP spans and the optional - # OpenWorkGraph defaults. Do not let large defaults bypass the same limit. _bounded_json_size({"payload": payload, "defaults": defaults or {}}) projected, stats = otel_payload_to_agent_events( payload, diff --git a/server/browser_agent_projection.py b/server/browser_agent_projection.py new file mode 100644 index 00000000..b9027e8a --- /dev/null +++ b/server/browser_agent_projection.py @@ -0,0 +1,101 @@ +from __future__ import annotations + +"""Project privacy-safe browser lifecycle signals into the agent evidence model. + +The browser extension never supplies prompt/response/alert text here. Provider, +opaque run id and lifecycle state are the complete allowlist. +""" + +import re +from typing import Any, Callable + +from fastapi import FastAPI + +from .agent_ingest import ingest_agent_payloads +from .main import BrowserEvent +from .secure_app import app + +_PROVIDER_MAP = { + "chatgpt": ("ChatGPT", "openai", "chatgpt_web"), + "claude": ("Claude", "anthropic", "claude_web"), + "microsoft_copilot": ("Microsoft Copilot", "microsoft", "copilot_web"), + "lovable": ("Lovable", "lovable", "lovable_web"), + "gemini": ("Gemini", "google", "gemini_web"), +} +_ACTION_MAP = { + "agent_run_started": ("run_started", "running"), + "agent_run_finished": ("run_finished", "success"), + "agent_run_cancelled": ("run_finished", "cancelled"), + "agent_error": ("error", "error"), + "agent_approval_requested": ("human_approval_requested", "running"), + "agent_approval_received": ("human_approval_received", "running"), +} +_STRUCTURAL = re.compile(r"^[A-Za-z0-9_.:\-]{1,128}$") + + +def _take_endpoint(target: FastAPI, path: str, method: str) -> Callable[..., Any] | None: + wanted = method.upper() + for route in list(target.router.routes): + methods = getattr(route, "methods", None) or set() + if getattr(route, "path", None) == path and wanted in methods: + target.router.routes.remove(route) + return getattr(route, "endpoint", None) + return None + + +def _agent_payload(event: BrowserEvent) -> dict[str, Any] | None: + mapped = _ACTION_MAP.get(str(event.action or "")) + if mapped is None: + return None + meta = dict(event.metadata or {}) + provider_key = str(meta.get("agent_provider") or "").strip().lower() + provider = _PROVIDER_MAP.get(provider_key) + run_id = str(meta.get("agent_run_id") or "").strip() + if provider is None or not _STRUCTURAL.fullmatch(run_id): + return None + operation, status = mapped + agent_name, provider_name, framework = provider + payload: dict[str, Any] = { + "observed_at": event.observed_at, + "agent_name": agent_name, + "provider": provider_name, + "framework": framework, + "operation": operation, + "status": status, + "observation_level": "os_observed", + "run_id": run_id, + "session_id": run_id, + "sensor_id": f"browser-agent:{provider_key}", + "device_id": "browser-agent-local", + "tool_category": "none", + } + return payload + + +def install_browser_agent_projection(target: FastAPI = app) -> None: + if getattr(target.state, "owg_browser_agent_projection_installed", False): + return + original = _take_endpoint(target, "/v1/browser-events", "POST") + if original is None: + raise RuntimeError("browser event route must exist before agent projection is installed") + + @target.post("/v1/browser-events") + def browser_event_with_agent_projection(event: BrowserEvent) -> dict[str, Any]: + result = original(event) + payload = _agent_payload(event) + projected = 0 + if payload is not None: + try: + projected = int(ingest_agent_payloads([payload]).get("inserted") or 0) + except Exception: + # Browser observation must never break normal browser evidence. + # Invalid structural lifecycle data is simply not promoted. + projected = 0 + return {**dict(result or {}), "agent_projection_inserted": projected} + + target.state.owg_browser_agent_projection_installed = True + + +install_browser_agent_projection() + +__all__ = ["install_browser_agent_projection"] diff --git a/server/context_pulse_routes.py b/server/context_pulse_routes.py index 743951e3..b1f42247 100644 --- a/server/context_pulse_routes.py +++ b/server/context_pulse_routes.py @@ -1,15 +1,96 @@ from __future__ import annotations +import base64 +from datetime import datetime +import json +import os from typing import Any -from fastapi import HTTPException +from fastapi import HTTPException, Request +from shared.history_policy import active_ai_history_access, history_generation from .context_pulse import context_pulse from .secure_app import app +_CURSOR_ENVELOPE_VERSION = 1 + + +def _wrap_cursor(inner: str, generation: int) -> str: + raw = json.dumps( + {"h": _CURSOR_ENVELOPE_VERSION, "g": int(generation), "c": str(inner)}, + separators=(",", ":"), sort_keys=True, + ).encode("utf-8") + return base64.urlsafe_b64encode(raw).decode("ascii").rstrip("=") + + +def _unwrap_cursor(value: str | None, generation: int) -> tuple[str | None, bool]: + if not value: + return None, False + try: + raw = base64.urlsafe_b64decode((str(value) + "=" * (-len(str(value)) % 4)).encode("ascii")) + data = json.loads(raw.decode("utf-8")) + if not isinstance(data, dict) or int(data.get("h") or 0) != _CURSOR_ENVELOPE_VERSION: + return None, True + if int(data.get("g") or 0) != int(generation): + return None, True + inner = str(data.get("c") or "") + return (inner or None), False + except Exception: + # v0.93 and earlier cursors have no deletion-generation envelope. A + # one-time bootstrap is safer than letting stale finding versions survive. + return None, True + + +def _parse(value: Any) -> datetime | None: + try: + parsed = datetime.fromisoformat(str(value or "").replace("Z", "+00:00")) + return parsed + except Exception: + return None + + +def _ai_history_filter(result: dict[str, Any], request: Request) -> None: + if str(request.headers.get("X-OpenWorkGraph-Context") or "").strip().lower() != "ai": + return + access = active_ai_history_access() + mode = str(access.get("mode") or "off") + if mode == "all_saved": + result["saved_history_access"] = "all_saved" + return + + run_start = _parse(os.getenv("WORKFLOW_OBSERVER_RUN_STARTED_AT")) + selected_start = _parse(access.get("since")) if mode == "selected_range" else None + selected_end = _parse(access.get("until")) if mode == "selected_range" else None + + def allowed(observed_at: Any) -> bool: + observed = _parse(observed_at) + if observed is None: + return False + if run_start is not None and observed >= run_start: + return True + return bool( + selected_start is not None + and selected_end is not None + and selected_start <= observed < selected_end + ) + + rows = [row for row in list(result.get("recent_evidence") or []) if isinstance(row, dict) and allowed(row.get("observed_at"))] + result["recent_evidence"] = rows + result["recent_returned"] = len(rows) + # Long-horizon findings aggregate multiple dates, so a selected range cannot + # safely be enforced by filtering the final aggregate. Keep them off unless + # the user granted all saved history; selected-range analysis should use + # list_history + date-ranged canonical tools instead. + result["findings"] = [] + result["findings_returned"] = 0 + result["findings_has_more"] = False + result["saved_history_access"] = mode + result["saved_history_findings_omitted"] = True + @app.get("/v1/context-pulse") def get_context_pulse( + request: Request, cursor: str | None = None, recent_limit: int = 12, finding_limit: int = 6, @@ -18,14 +99,15 @@ def get_context_pulse( ) -> dict[str, Any]: """Return an incremental factual update for a connected AI client. - The caller owns the opaque cursor. Recent evidence is compact by default; - pass ``recent_detail=rich`` or use the canonical workflow trace when full - rich rows are needed. AI disclosure/redaction is still enforced by the - established Context boundary middleware. + The outer cursor includes only a local history-generation number and the + existing opaque Pulse cursor. Deletion/expiration invalidates it so deleted + evidence cannot remain hidden in stale finding-version state. """ + generation = history_generation() + inner, invalidated = _unwrap_cursor(cursor, generation) try: - return context_pulse( - cursor=cursor, + result = context_pulse( + cursor=inner, recent_limit=recent_limit, finding_limit=finding_limit, lookback_days=lookback_days, @@ -33,6 +115,13 @@ def get_context_pulse( ) except ValueError as exc: raise HTTPException(status_code=400, detail=str(exc)) from exc + result["next_cursor"] = _wrap_cursor(str(result.get("next_cursor") or ""), generation) + result["history_generation"] = generation + result["cursor_invalidated_by_history_change"] = bool(invalidated) + if invalidated: + result["cursor_notice"] = "Saved history changed; this Pulse restarted from a safe baseline." + _ai_history_filter(result, request) + return result __all__ = ["app", "get_context_pulse"] diff --git a/server/enterprise_runner.py b/server/enterprise_runner.py index 93772bd8..427b1281 100644 --- a/server/enterprise_runner.py +++ b/server/enterprise_runner.py @@ -1,41 +1,32 @@ from __future__ import annotations import argparse - import uvicorn - SECURE_APP = "server.secure_app:app" def main() -> None: - """Run the established secure local app with optional enterprise controls installed. - - Importing the additive route modules registers Gateway/capture/deletion routes, - lifecycle hooks and dashboard controls on the same FastAPI app object exported - by ``server.secure_app``. Uvicorn still serves the hardened secure-app target. - """ parser = argparse.ArgumentParser() parser.add_argument("--host", default="127.0.0.1") parser.add_argument("--port", type=int, default=8787) args = parser.parse_args() - # Registration imports: these modules extend, rather than replace, secure_app. import server.enterprise_app # noqa: F401 + import server.history_routes # noqa: F401 + import server.history_capture_integration # noqa: F401 + import server.history_ai_access_integration # noqa: F401 import server.evidence_delete_routes # noqa: F401 import server.v0571_polish # noqa: F401 import server.evidence_paging # noqa: F401 import server.work_profile_routes # noqa: F401 import server.context_pulse_routes # noqa: F401 import server.browser_signal_routes # noqa: F401 + import server.browser_agent_projection # noqa: F401 import server.agent_dashboard_control_plane # noqa: F401 import server.org_join_routes as org_join_routes - # Import last so its HTML middleware injects the privacy-safe dashboard loader - # after the other additive dashboard scripts have been installed. import server.dashboard_privacy # noqa: F401 - # Managed setup is best-effort and never blocks local capture or the dashboard. - # Its status is always exposed in the employee UI when a managed config exists. org_join_routes.start_managed_setup_in_background() uvicorn.run(SECURE_APP, host=args.host, port=args.port) diff --git a/server/evidence_delete_routes.py b/server/evidence_delete_routes.py index 726d6d0c..cd8d7e53 100644 --- a/server/evidence_delete_routes.py +++ b/server/evidence_delete_routes.py @@ -11,6 +11,7 @@ from connector.runtime import start_sync_worker, status as worker_status, stop_sync_worker from connector.state import SyncState from shared.evidence_deletion import add_tombstone, audit_deletion, normalize_range +from shared.history_policy import bump_history_generation from . import analytics from .db import DATA_DIR from .evidence_delete import delete_database_range, remove_local_screenshots, rewrite_jsonl_range @@ -44,8 +45,6 @@ def delete_local_evidence(request: EvidenceDeleteRequest) -> dict[str, Any]: if was_running: stop_sync_worker() - # Fail closed first. Any late collector/browser delivery for this range is - # rejected even if a later filesystem cleanup encounters an error. add_tombstone(since, until) result: dict[str, Any] | None = None @@ -54,12 +53,10 @@ def delete_local_evidence(request: EvidenceDeleteRequest) -> dict[str, Any]: screenshots_removed = 0 jsonl_removed = 0 gateway_cursor = 0 + history_generation = 0 try: sync_state = SyncState(DATA_DIR / "gateway_sync_state.db") gateway_cursor = sync_state.get_int("last_local_event_id", 0) - - # Collector JSONL and pending delivery are auxiliary local evidence - # copies, so delete them as part of the same user operation. jsonl_removed = rewrite_jsonl_range(since, until) outbox = EventOutbox(DATA_DIR / "collector_outbox.db") pruned_outbox = outbox.prune_range(since, until) @@ -68,11 +65,8 @@ def delete_local_evidence(request: EvidenceDeleteRequest) -> dict[str, Any]: local_ids = [int(value) for value in result.get("local_ids") or []] skipped_gateway = sync_state.add_skip_ids(local_ids, "local_evidence_deleted") screenshots_removed = remove_local_screenshots(list(result.get("screenshot_paths") or [])) - - # The summary cache historically assumes append-only local evidence. - # Deletion must invalidate it explicitly or the dashboard can display - # rows that no longer exist in SQLite. analytics.clear_summary_cache() + history_generation = bump_history_generation(reason="manual_range_deletion") audit_deletion( since=since, @@ -108,6 +102,8 @@ def delete_local_evidence(request: EvidenceDeleteRequest) -> dict[str, Any]: "gateway_local_rows_marked_never_share": int(skipped_gateway), "screenshots_removed": int(screenshots_removed), "late_delivery_suppressed": True, + "history_generation": int(history_generation), + "context_pulse_cursors_invalidated": True, "gateway_recall_performed": False, "rows_at_or_before_gateway_cursor": int(at_or_before_cursor), "notice": ( diff --git a/server/history_ai_access_integration.py b/server/history_ai_access_integration.py new file mode 100644 index 00000000..10a03a0a --- /dev/null +++ b/server/history_ai_access_integration.py @@ -0,0 +1,41 @@ +from __future__ import annotations + +from typing import Any, Callable + +from fastapi import FastAPI, Request + +from shared.history_policy import set_ai_history_access +from .secure_app import app + + +def _take_endpoint(target: FastAPI, path: str, method: str) -> Callable[..., Any] | None: + wanted = method.upper() + for route in list(target.router.routes): + methods = getattr(route, "methods", None) or set() + if getattr(route, "path", None) == path and wanted in methods: + target.router.routes.remove(route) + return getattr(route, "endpoint", None) + return None + + +def install_history_ai_access_integration(target: FastAPI = app) -> None: + if getattr(target.state, "owg_history_ai_access_integrated", False): + return + original = _take_endpoint(target, "/v1/ai-access", "POST") + if original is None: + raise RuntimeError("AI access route must exist before history integration") + + @target.post("/v1/ai-access") + async def update_ai_access_and_history(request: Request): + result = await original(request) + if isinstance(result, dict) and result.get("enabled") is False: + set_ai_history_access(mode="off", expires_minutes=None) + result["saved_history_access_revoked"] = True + return result + + target.state.owg_history_ai_access_integrated = True + + +install_history_ai_access_integration() + +__all__ = ["install_history_ai_access_integration"] diff --git a/server/history_capture_integration.py b/server/history_capture_integration.py new file mode 100644 index 00000000..07a2a1f8 --- /dev/null +++ b/server/history_capture_integration.py @@ -0,0 +1,57 @@ +from __future__ import annotations + +"""Bind human-session retention to explicit recording Stop/Start boundaries.""" + +from typing import Any, Callable + +from fastapi import FastAPI + +from shared.history_policy import retention_for_kind +from .history_retention import cleanup_ephemeral_session +from .main import COLLECTOR_STATUS +from .secure_app import app + + +def _take_endpoint(target: FastAPI, path: str, method: str) -> Callable[..., Any] | None: + wanted = method.upper() + for route in list(target.router.routes): + methods = getattr(route, "methods", None) or set() + if getattr(route, "path", None) == path and wanted in methods: + target.router.routes.remove(route) + return getattr(route, "endpoint", None) + return None + + +def _current_human_session() -> str: + return str(COLLECTOR_STATUS.get("session_id") or "").strip() + + +def install_history_capture_integration(target: FastAPI = app) -> None: + if getattr(target.state, "owg_history_capture_integrated", False): + return + original_stop = _take_endpoint(target, "/v1/capture/stop", "POST") + original_start = _take_endpoint(target, "/v1/capture/start", "POST") + if original_stop is None or original_start is None: + raise RuntimeError("capture routes must be registered before history integration") + + @target.post("/v1/capture/stop") + def stop_and_apply_history() -> dict[str, Any]: + session_id = _current_human_session() + result = original_stop() + if session_id and retention_for_kind("human").get("mode") == "ephemeral": + result["history_cleanup"] = cleanup_ephemeral_session("human", session_id) + return result + + @target.post("/v1/capture/start") + def start_and_apply_history() -> dict[str, Any]: + previous = _current_human_session() + if previous and retention_for_kind("human").get("mode") == "ephemeral": + cleanup_ephemeral_session("human", previous) + return original_start() + + target.state.owg_history_capture_integrated = True + + +install_history_capture_integration() + +__all__ = ["install_history_capture_integration"] diff --git a/server/history_retention.py b/server/history_retention.py new file mode 100644 index 00000000..5b917bb1 --- /dev/null +++ b/server/history_retention.py @@ -0,0 +1,378 @@ +from __future__ import annotations + +"""User-controlled local history catalogue and retention cleanup.""" + +import json +import os +import tempfile +from collections import defaultdict +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any + +from collector.outbox import EventOutbox +from connector.state import SyncState +from shared.history_policy import ( + add_session_tombstones, + initialize_policy, + read_policy, + retention_for_kind, +) +from . import analytics +from .agent_execution_traces import agent_execution_traces +from .context_layers import factual_context_timeline +from .db import DATA_DIR, connect +from .evidence_delete import remove_local_screenshots + +_ACTIVITY_BREAK_SECONDS = 30 * 60 + + +def _parse(value: Any) -> datetime | None: + try: + parsed = datetime.fromisoformat(str(value or "").replace("Z", "+00:00")) + if parsed.tzinfo is None: + return None + return parsed.astimezone(timezone.utc) + except Exception: + return None + + +def _event(row: Any) -> dict[str, Any]: + value = dict(row) + value.pop("id", None) + raw = value.pop("metadata_json", "{}") + try: + value["metadata"] = json.loads(raw or "{}") if isinstance(raw, str) else (raw or {}) + except Exception: + value["metadata"] = {} + return value + + +def _is_agent(row: dict[str, Any]) -> bool: + if str(row.get("source") or "").lower() == "agent": + return True + meta = row.get("metadata") if isinstance(row.get("metadata"), dict) else {} + return str(meta.get("actor_kind") or "").lower() == "agent" + + +def _all_rows() -> list[dict[str, Any]]: + with connect() as conn: + rows = conn.execute("SELECT * FROM events ORDER BY observed_at ASC, id ASC").fetchall() + return [_event(row) for row in rows] + + +def _kind_rows(kind: str, rows: list[dict[str, Any]] | None = None) -> list[dict[str, Any]]: + source = rows if rows is not None else _all_rows() + agent = str(kind).lower() == "agent" + return [row for row in source if _is_agent(row) == agent] + + +def _session_ids(kind: str, rows: list[dict[str, Any]] | None = None) -> list[str]: + return sorted({str(row.get("session_id") or "") for row in _kind_rows(kind, rows) if row.get("session_id")}) + + +def _delete_ids(conn, table: str, event_ids: list[str]) -> None: + for offset in range(0, len(event_ids), 500): + chunk = event_ids[offset : offset + 500] + if not chunk: + continue + placeholders = ",".join("?" for _ in chunk) + conn.execute(f"DELETE FROM {table} WHERE event_id IN ({placeholders})", tuple(chunk)) + + +def _rewrite_jsonl_sessions(kind: str, session_ids: set[str]) -> int: + removed = 0 + if not session_ids: + return 0 + for path in (DATA_DIR / "events.jsonl", DATA_DIR / "events.jsonl.1"): + if not path.exists() or not path.is_file(): + continue + fd, tmp_name = tempfile.mkstemp(prefix=path.name + ".history-", suffix=".tmp", dir=str(path.parent)) + try: + with path.open("r", encoding="utf-8", errors="replace") as source, os.fdopen(fd, "w", encoding="utf-8") as target: + for line in source: + try: + event = json.loads(line) + except Exception: + target.write(line) + continue + if isinstance(event, dict): + row_kind = "agent" if _is_agent(event) else "human" + if row_kind == kind and str(event.get("session_id") or "") in session_ids: + removed += 1 + continue + target.write(line) + target.flush() + os.fsync(target.fileno()) + os.replace(tmp_name, path) + try: + os.chmod(path, 0o600) + except Exception: + pass + finally: + try: + if os.path.exists(tmp_name): + os.unlink(tmp_name) + except Exception: + pass + return removed + + +def delete_sessions(kind: str, session_ids: list[str], *, reason: str) -> dict[str, Any]: + selected = "agent" if str(kind).lower() == "agent" else "human" + ids = {str(value) for value in session_ids if str(value)} + if not ids: + return {"kind": selected, "sessions_deleted": 0, "events_deleted": 0, "local_ids": []} + + # Tombstone first so late native/browser/outbox delivery cannot recreate a + # session while cleanup is in progress or after a crash between cleanup steps. + add_session_tombstones(selected, sorted(ids), reason=reason) + + with connect() as conn: + rows = conn.execute("SELECT * FROM events ORDER BY id ASC").fetchall() + chosen = [] + for raw in rows: + item = _event(raw) + if str(item.get("session_id") or "") not in ids: + continue + if ("agent" if _is_agent(item) else "human") != selected: + continue + chosen.append(raw) + event_ids = [str(row["event_id"]) for row in chosen] + local_ids = [int(row["id"]) for row in chosen] + screenshots = [str(row["screenshot_path"]) for row in chosen if row["screenshot_path"]] + if event_ids: + _delete_ids(conn, "context_events", event_ids) + _delete_ids(conn, "normalized_events", event_ids) + _delete_ids(conn, "events", event_ids) + + jsonl_removed = _rewrite_jsonl_sessions(selected, ids) + outbox = EventOutbox(DATA_DIR / "collector_outbox.db") + outbox_removed = outbox.prune_sessions(selected, sorted(ids)) + state = SyncState(DATA_DIR / "gateway_sync_state.db") + skipped_gateway = state.add_skip_ids(local_ids, f"history_{reason}") if local_ids else 0 + screenshots_removed = remove_local_screenshots(screenshots) + analytics.clear_summary_cache() + return { + "kind": selected, + "sessions_deleted": len(ids), + "events_deleted": len(event_ids), + "local_ids": local_ids, + "jsonl_events_removed": jsonl_removed, + "collector_outbox_events_removed": outbox_removed, + "gateway_local_rows_marked_never_share": skipped_gateway, + "screenshots_removed": screenshots_removed, + "late_delivery_suppressed": True, + } + + +def initialize_history_retention() -> dict[str, Any]: + with connect() as conn: + has_existing = bool(conn.execute("SELECT 1 FROM events LIMIT 1").fetchone()) + policy = initialize_policy(has_existing_evidence=has_existing) + cleanup_expired_history(startup=True) + return policy + + +def cleanup_expired_history(*, startup: bool = False) -> dict[str, Any]: + rows = _all_rows() + now = datetime.now(timezone.utc) + results: list[dict[str, Any]] = [] + for kind in ("human", "agent"): + retention = retention_for_kind(kind) + by_session: dict[str, list[dict[str, Any]]] = defaultdict(list) + for row in _kind_rows(kind, rows): + sid = str(row.get("session_id") or "") + if sid: + by_session[sid].append(row) + delete_ids: list[str] = [] + if retention.get("mode") == "ephemeral" and startup: + # No current-run evidence exists when called from lifespan startup, so + # every surviving ephemeral session is a stale/crash-recovery session. + delete_ids = sorted(by_session) + elif retention.get("mode") == "days": + cutoff = now - timedelta(days=int(retention.get("days") or 1)) + for sid, members in by_session.items(): + latest = max((_parse(row.get("observed_at")) for row in members), default=None) + if latest is not None and latest < cutoff: + delete_ids.append(sid) + if delete_ids: + results.append(delete_sessions(kind, delete_ids, reason="retention_expired")) + return {"startup": startup, "cleanups": results} + + +def cleanup_ephemeral_session(kind: str, session_id: str) -> dict[str, Any]: + retention = retention_for_kind(kind) + if retention.get("mode") != "ephemeral" or not str(session_id or ""): + return {"kind": kind, "sessions_deleted": 0, "events_deleted": 0} + return delete_sessions(kind, [session_id], reason="ephemeral_session_closed") + + +def _activity_blocks(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: + ordered = sorted(rows, key=lambda row: str(row.get("observed_at") or "")) + blocks: list[dict[str, Any]] = [] + for row in ordered: + start = _parse(row.get("observed_at")) + if start is None: + continue + duration = max(0.0, float(row.get("duration_seconds") or 0.0)) + end = start + timedelta(seconds=duration) + if not blocks: + blocks.append({"started_at": start.isoformat(), "ended_at": end.isoformat()}) + continue + previous_end = _parse(blocks[-1]["ended_at"]) or start + if (start - previous_end).total_seconds() > _ACTIVITY_BREAK_SECONDS: + blocks.append({"started_at": start.isoformat(), "ended_at": end.isoformat()}) + elif end > previous_end: + blocks[-1]["ended_at"] = end.isoformat() + return blocks + + +def _human_sessions(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: + grouped: dict[str, list[dict[str, Any]]] = defaultdict(list) + for row in _kind_rows("human", rows): + sid = str(row.get("session_id") or "") + if sid: + grouped[sid].append(row) + + # Reuse the factual focus projection for engaged time/surfaces while keeping + # session identity and raw event counts grounded in canonical rows. + timeline = factual_context_timeline(_raw_events=_kind_rows("human", rows)) + spans_by_session: dict[str, list[dict[str, Any]]] = defaultdict(list) + for span in timeline: + spans_by_session[str(span.get("session_id") or "")].append(span) + + output: list[dict[str, Any]] = [] + retention = retention_for_kind("human") + for sid, members in grouped.items(): + starts = [_parse(row.get("observed_at")) for row in members] + starts = [value for value in starts if value is not None] + if not starts: + continue + ended = [] + for row in members: + start = _parse(row.get("observed_at")) + if start is not None: + ended.append(start + timedelta(seconds=max(0.0, float(row.get("duration_seconds") or 0.0)))) + spans = spans_by_session.get(sid, []) + surfaces = sorted({str(span.get("work_surface") or span.get("container_app") or "") for span in spans if span.get("work_surface") or span.get("container_app")}) + expires_at = None + if retention.get("mode") == "days": + expires_at = (max(ended or starts) + timedelta(days=int(retention.get("days") or 1))).isoformat() + output.append({ + "history_session_id": f"human:{sid}", + "kind": "human", + "started_at": min(starts).isoformat(), + "ended_at": max(ended or starts).isoformat(), + "activity_blocks": _activity_blocks(members), + "event_count": len(members), + "engaged_seconds": round(sum(float(span.get("engaged_seconds") or 0.0) for span in spans), 3), + "surfaces": surfaces[:20], + "retention_mode": retention.get("mode"), + "expires_at": expires_at, + }) + return output + + +def _agent_sessions(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: + agents = _kind_rows("agent", rows) + if not agents: + return [] + payload = agent_execution_traces(agents, limit=1000, max_events_per_execution=1) + retention = retention_for_kind("agent") + output = [] + for execution in payload.get("executions") or []: + if not isinstance(execution, dict): + continue + expires_at = None + end = _parse(execution.get("ended_at") or execution.get("started_at")) + if retention.get("mode") == "days" and end is not None: + expires_at = (end + timedelta(days=int(retention.get("days") or 1))).isoformat() + output.append({ + "history_session_id": str(execution.get("execution_id") or ""), + "kind": "agent", + "started_at": execution.get("started_at"), + "ended_at": execution.get("ended_at"), + "agent": execution.get("agent"), + "observation_level": execution.get("observation_level"), + "outcome_status": execution.get("outcome_status"), + "event_count": int(execution.get("event_count_total") or 0), + "retention_mode": retention.get("mode"), + "expires_at": expires_at, + }) + return output + + +def list_history(*, since: str | None = None, until: str | None = None, limit: int = 200) -> dict[str, Any]: + rows = _all_rows() + start = _parse(since) if since else None + end = _parse(until) if until else None + if since and start is None: + raise ValueError("invalid history since") + if until and end is None: + raise ValueError("invalid history until") + if start and end and end <= start: + raise ValueError("history until must be after since") + if start or end: + filtered = [] + for row in rows: + observed = _parse(row.get("observed_at")) + if observed is None: + continue + if start and observed < start: + continue + if end and observed >= end: + continue + filtered.append(row) + rows = filtered + sessions = [*_human_sessions(rows), *_agent_sessions(rows)] + sessions.sort(key=lambda item: str(item.get("started_at") or ""), reverse=True) + bounded = max(1, min(int(limit), 1000)) + returned = sessions[:bounded] + policy = read_policy() + return { + "sessions": returned, + "returned": len(returned), + "total": len(sessions), + "has_more": len(sessions) > len(returned), + "since": start.isoformat() if start else None, + "until": end.isoformat() if end else None, + "history_generation": int(policy.get("history_generation") or 1), + "retention": { + "human": policy.get("human_retention"), + "agent": policy.get("agent_retention"), + }, + "session_semantics": { + "human": "one recording run; idle gaps over 30 minutes are activity blocks inside the same session", + "agent": "observed agent execution boundary; surface-only boundaries are best-effort and declare observation_level", + }, + "derived_task_labels_used": False, + } + + +def delete_history_session(history_session_id: str) -> dict[str, Any]: + raw = str(history_session_id or "").strip() + if raw.startswith("human:"): + return delete_sessions("human", [raw[len("human:"):]], reason="user_deleted_session") + + # Agent history exposes OWG execution IDs, not native run/session identifiers. + rows = _kind_rows("agent") + payload = agent_execution_traces(rows, execution_id=raw, limit=1, max_events_per_execution=1000) + executions = list(payload.get("executions") or []) + if not executions: + raise ValueError("history session not found") + event_ids = {str(event.get("event_id") or "") for event in executions[0].get("events") or [] if event.get("event_id")} + native_sessions = sorted({str(row.get("session_id") or "") for row in rows if str(row.get("event_id") or "") in event_ids and row.get("session_id")}) + if not native_sessions: + raise ValueError("agent history session has no deletable canonical session") + return delete_sessions("agent", native_sessions, reason="user_deleted_session") + + +__all__ = [ + "cleanup_ephemeral_session", + "cleanup_expired_history", + "delete_history_session", + "delete_sessions", + "initialize_history_retention", + "list_history", +] diff --git a/server/history_routes.py b/server/history_routes.py new file mode 100644 index 00000000..18e59c46 --- /dev/null +++ b/server/history_routes.py @@ -0,0 +1,120 @@ +from __future__ import annotations + +import json +from typing import Any + +from fastapi import HTTPException, Request, Response +from fastapi.responses import HTMLResponse +from pydantic import BaseModel + +from shared.evidence import rich_evidence_row +from shared.history_policy import active_ai_history_access, read_policy, set_ai_history_access, update_retention +from shared.lifespan import extend_lifespan +from .db import connect +from .history_retention import cleanup_expired_history, delete_history_session, initialize_history_retention, list_history +from .main import ROOT +from .privacy_pipeline import redact_for_display +from .secure_app import app + + +class RetentionChoice(BaseModel): + human_mode: str + human_days: int | None = None + agent_mode: str + agent_days: int | None = None + onboarding_complete: bool = True + +class HistoryDeleteRequest(BaseModel): history_session_id: str +class AIHistoryAccessRequest(BaseModel): + mode: str + since: str | None = None + until: str | None = None + expires_minutes: int | None = 60 + + +def _history_startup() -> None: initialize_history_retention() +def _history_shutdown() -> None: cleanup_expired_history(startup=True) +extend_lifespan(app, startup=_history_startup, shutdown=_history_shutdown) + + +@app.get("/v1/history-policy") +def get_history_policy() -> dict[str, Any]: + value=read_policy() + return {"onboarding_complete":bool(value.get("onboarding_complete")),"human_retention":value.get("human_retention"),"agent_retention":value.get("agent_retention"),"history_generation":int(value.get("history_generation") or 1),"upgrade_preserved_existing_history":bool(value.get("upgrade_preserved_existing_history")),"tradeoff":{"ephemeral":"Live context works, but completed sessions are deleted and cannot support later month/week comparisons.","saved":"Saved history stays on this computer until its retention period expires or you delete it.","ai_access_separate":True}} + +@app.put("/v1/history-policy") +def set_history_policy(request: RetentionChoice) -> dict[str, Any]: + try: + value=update_retention(human_mode=request.human_mode,human_days=request.human_days,agent_mode=request.agent_mode,agent_days=request.agent_days,onboarding_complete=request.onboarding_complete) + cleanup=cleanup_expired_history(startup=False) + except (TypeError,ValueError) as exc: raise HTTPException(status_code=400,detail=str(exc)) from exc + return {"status":"saved","human_retention":value.get("human_retention"),"agent_retention":value.get("agent_retention"),"onboarding_complete":bool(value.get("onboarding_complete")),"cleanup":cleanup} + +@app.get("/v1/history") +def get_history(since:str|None=None,until:str|None=None,limit:int=200)->dict[str,Any]: + try: return list_history(since=since,until=until,limit=limit) + except ValueError as exc: raise HTTPException(status_code=400,detail=str(exc)) from exc + +@app.post("/v1/history/delete-session") +def delete_session(request:HistoryDeleteRequest)->dict[str,Any]: + try: result=delete_history_session(request.history_session_id) + except ValueError as exc: raise HTTPException(status_code=404,detail=str(exc)) from exc + return {"status":"deleted",**result} + +@app.get("/v1/history/ai-access") +def get_ai_history_access()->dict[str,Any]: + return {"access":active_ai_history_access(),"meaning":"Saving history locally does not grant an AI permission to read saved history."} + +@app.post("/v1/history/ai-access") +def change_ai_history_access(request:AIHistoryAccessRequest)->dict[str,Any]: + try: access=set_ai_history_access(mode=request.mode,since=request.since,until=request.until,expires_minutes=request.expires_minutes) + except ValueError as exc: raise HTTPException(status_code=400,detail=str(exc)) from exc + return {"status":"updated","access":access} + +@app.post("/v1/history/cleanup") +def run_history_cleanup()->dict[str,Any]: return cleanup_expired_history(startup=False) + +@app.get("/v1/history/export-json") +def export_history_json(since:str|None=None,until:str|None=None,include_raw:bool=False,redact_names:bool=False)->Response: + # Validate the date window through the same catalogue parser first. + try: catalog=list_history(since=since,until=until,limit=1) + except ValueError as exc: raise HTTPException(status_code=400,detail=str(exc)) from exc + clauses=[]; params:list[Any]=[] + if catalog.get("since"): clauses.append("observed_at >= ?"); params.append(catalog["since"]) + if catalog.get("until"): clauses.append("observed_at < ?"); params.append(catalog["until"]) + where=(" WHERE "+" AND ".join(clauses)) if clauses else "" + with connect() as conn: + rows=conn.execute("SELECT * FROM events"+where+" ORDER BY observed_at ASC, id ASC LIMIT 100001",tuple(params)).fetchall() + if len(rows)>100000: raise HTTPException(status_code=413,detail="History export exceeds 100000 events; choose a smaller date range") + evidence=[rich_evidence_row(dict(row),include_identity=False) for row in rows] + payload={"export":{"kind":"retained_history","since":catalog.get("since"),"until":catalog.get("until"),"include_raw":bool(include_raw),"redact_names":bool(redact_names),"canonical_rows":len(evidence)},"history":catalog,"evidence":evidence} + if redact_names: + from .ai_context import redact_contextually + payload=redact_contextually(payload) + elif not include_raw: + payload=redact_for_display(payload) + body=(json.dumps(payload,ensure_ascii=False,indent=2)+"\n").encode("utf-8") + return Response(body,media_type="application/json",headers={"Content-Disposition":'attachment; filename="openworkgraph-history.json"',"Cache-Control":"no-store"}) + +@app.get("/history-retention.js",include_in_schema=False) +def history_retention_script()->Response: + path=ROOT/"dashboard"/"history_retention.js" + return Response(path.read_text(encoding="utf-8"),media_type="application/javascript",headers={"Cache-Control":"no-store","X-Content-Type-Options":"nosniff"}) + +@app.middleware("http") +async def inject_history_retention(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 + marker='' + if marker not in text:text=text.replace("",marker+"\n") + headers=dict(response.headers);headers.pop("content-length",None) + return HTMLResponse(text,status_code=response.status_code,headers=headers) + +__all__=["get_history","get_history_policy"] diff --git a/shared/capture_control.py b/shared/capture_control.py index ed3bf383..1ebeabf9 100644 --- a/shared/capture_control.py +++ b/shared/capture_control.py @@ -10,6 +10,7 @@ from typing import Any from .evidence_deletion import prepare_recordable_event as prepare_not_deleted +from .history_policy import prepare_recordable_event as prepare_retained _LOCK = threading.RLock() _MAX_INTERVALS = 500 @@ -190,16 +191,18 @@ def timestamp_is_skipped(observed_at: str | None, *, state: dict[str, Any] | Non def prepare_recordable_event(event: dict[str, Any], *, state: dict[str, Any] | None = None) -> dict[str, Any] | None: - """Apply permanent deletion tombstones, then capture pause/stop boundaries. + """Apply deletion, retention and capture-control boundaries before persistence. - Deleted time ranges are checked first so late queued evidence cannot recreate - locally deleted work. Capture-state clipping then protects against a collector - process that is shutting down just after Pause/Stop. + Deletion and retention are checked before pause/stop clipping so buffered or + late evidence cannot recreate a session the user already expired or deleted. """ deletion_safe = prepare_not_deleted(event) if deletion_safe is None: return None - event = deletion_safe + retention_safe = prepare_retained(deletion_safe) + if retention_safe is None: + return None + event = retention_safe value = state or read_state() start = _parse(str(event.get("observed_at") or "")) if start is None: diff --git a/shared/history_policy.py b/shared/history_policy.py new file mode 100644 index 00000000..4b5cec53 --- /dev/null +++ b/shared/history_policy.py @@ -0,0 +1,363 @@ +from __future__ import annotations + +"""Local history-retention and saved-history disclosure policy. + +Capture, retention, and AI disclosure are separate decisions. This module is in +``shared`` so ingestion can fail closed before persistence without importing the +server layer. The state file contains policy/bookkeeping only, never work text. +""" + +import copy +import json +import os +import tempfile +import threading +from datetime import datetime, timedelta, timezone +from pathlib import Path +from typing import Any + +_LOCK = threading.RLock() +_POLICY_VERSION = 1 +_MAX_SESSION_TOMBSTONES = 5000 +_RETENTION_MODES = {"ephemeral", "days", "forever"} +_AI_ACCESS_MODES = {"off", "selected_range", "all_saved"} + + +def _now_dt() -> datetime: + return datetime.now(timezone.utc) + + +def _now() -> str: + return _now_dt().isoformat() + + +def _parse(value: Any) -> datetime | None: + try: + parsed = datetime.fromisoformat(str(value or "").replace("Z", "+00:00")) + if parsed.tzinfo is None: + return None + return parsed.astimezone(timezone.utc) + except Exception: + return None + + +def data_dir() -> Path: + root = Path(__file__).resolve().parents[1] + path = Path(os.getenv("WORKFLOW_OBSERVER_DATA", root / "data")) + path.mkdir(parents=True, exist_ok=True) + return path + + +def policy_path() -> Path: + return data_dir() / "history_policy.json" + + +def _retention(mode: str, days: int | None = None) -> dict[str, Any]: + selected = str(mode or "").strip().lower() + if selected not in _RETENTION_MODES: + raise ValueError("retention mode must be ephemeral, days, or forever") + if selected == "days": + amount = int(days or 0) + if amount < 1 or amount > 3650: + raise ValueError("retention days must be between 1 and 3650") + return {"mode": "days", "days": amount} + return {"mode": selected, "days": None} + + +def _default(*, has_existing_evidence: bool) -> dict[str, Any]: + # Existing installations must never lose history merely by upgrading. New + # installations start ephemeral until the person makes an informed choice. + initial = "forever" if has_existing_evidence else "ephemeral" + return { + "version": _POLICY_VERSION, + "onboarding_complete": False, + "human_retention": _retention(initial), + "agent_retention": _retention(initial), + "history_generation": 1, + "session_tombstones": [], + "ai_history_access": { + "mode": "off", + "since": None, + "until": None, + "granted_at": None, + "expires_at": None, + }, + "created_at": _now(), + "updated_at": _now(), + "upgrade_preserved_existing_history": bool(has_existing_evidence), + } + + +def _uninitialized_fallback() -> dict[str, Any]: + """Preserve pre-History retention until the owning runtime initializes it. + + The production runner explicitly initializes History: genuinely new installs + then receive the ephemeral onboarding default, while upgrades preserve existing + history. Legacy/embedded callers that import agent ingestion or secure_app + directly do not run that initialization. Absence of a policy file there must + not silently turn completed agent runs into disposable data. + """ + value = _default(has_existing_evidence=True) + value["upgrade_preserved_existing_history"] = False + return value + + +def _normalize(value: dict[str, Any]) -> dict[str, Any]: + result = dict(value) + result["version"] = _POLICY_VERSION + result["onboarding_complete"] = bool(result.get("onboarding_complete")) + for key in ("human_retention", "agent_retention"): + raw = result.get(key) if isinstance(result.get(key), dict) else {} + try: + result[key] = _retention(str(raw.get("mode") or "ephemeral"), raw.get("days")) + except Exception: + result[key] = _retention("ephemeral") + result["history_generation"] = max(1, int(result.get("history_generation") or 1)) + tombstones = result.get("session_tombstones") if isinstance(result.get("session_tombstones"), list) else [] + safe_tombstones: list[dict[str, str]] = [] + for item in tombstones[-_MAX_SESSION_TOMBSTONES:]: + if not isinstance(item, dict): + continue + kind = str(item.get("kind") or "").strip().lower() + session_id = str(item.get("session_id") or "").strip() + if kind in {"human", "agent"} and session_id: + safe_tombstones.append({ + "kind": kind, + "session_id": session_id[:240], + "deleted_at": str(item.get("deleted_at") or _now())[:80], + "reason": str(item.get("reason") or "retention")[:120], + }) + result["session_tombstones"] = safe_tombstones + access = result.get("ai_history_access") if isinstance(result.get("ai_history_access"), dict) else {} + mode = str(access.get("mode") or "off").strip().lower() + if mode not in _AI_ACCESS_MODES: + mode = "off" + result["ai_history_access"] = { + "mode": mode, + "since": str(access.get("since") or "") or None, + "until": str(access.get("until") or "") or None, + "granted_at": str(access.get("granted_at") or "") or None, + "expires_at": str(access.get("expires_at") or "") or None, + } + result.setdefault("created_at", _now()) + result["updated_at"] = str(result.get("updated_at") or _now()) + return result + + +def _write_unlocked(value: dict[str, Any]) -> dict[str, Any]: + path = policy_path() + path.parent.mkdir(parents=True, exist_ok=True) + payload = json.dumps(_normalize(value), ensure_ascii=False, indent=2) + "\n" + fd, tmp_name = tempfile.mkstemp(prefix="history-policy-", suffix=".tmp", dir=str(path.parent)) + try: + with os.fdopen(fd, "w", encoding="utf-8") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) + os.replace(tmp_name, path) + try: + os.chmod(path, 0o600) + except Exception: + pass + finally: + try: + if os.path.exists(tmp_name): + os.unlink(tmp_name) + except Exception: + pass + return json.loads(payload) + + +def initialize_policy(*, has_existing_evidence: bool) -> dict[str, Any]: + with _LOCK: + path = policy_path() + if path.exists(): + return read_policy() + return _write_unlocked(_default(has_existing_evidence=has_existing_evidence)) + + +def read_policy() -> dict[str, Any]: + with _LOCK: + try: + value = json.loads(policy_path().read_text(encoding="utf-8")) + if not isinstance(value, dict): + raise ValueError("history policy must be an object") + return _normalize(value) + except Exception: + # History has not been initialized by the owning runtime yet. Keep + # legacy/direct callers non-destructive; production initialization + # writes the explicit new-install or upgrade policy before use. + return _uninitialized_fallback() + + +def update_retention( + *, + human_mode: str, + human_days: int | None, + agent_mode: str, + agent_days: int | None, + onboarding_complete: bool = True, +) -> dict[str, Any]: + with _LOCK: + value = read_policy() + value["human_retention"] = _retention(human_mode, human_days) + value["agent_retention"] = _retention(agent_mode, agent_days) + value["onboarding_complete"] = bool(onboarding_complete) + value["updated_at"] = _now() + return _write_unlocked(value) + + +def event_kind(event: dict[str, Any]) -> str: + if str(event.get("source") or "").strip().lower() == "agent": + return "agent" + metadata = event.get("metadata") if isinstance(event.get("metadata"), dict) else {} + if str(metadata.get("actor_kind") or "").strip().lower() == "agent": + return "agent" + return "human" + + +def retention_for_kind(kind: str, *, policy: dict[str, Any] | None = None) -> dict[str, Any]: + selected = "agent" if str(kind).lower() == "agent" else "human" + value = policy or read_policy() + return dict(value.get(f"{selected}_retention") or _retention("ephemeral")) + + +def session_is_tombstoned(kind: str, session_id: str, *, policy: dict[str, Any] | None = None) -> bool: + sid = str(session_id or "").strip() + if not sid: + return False + selected = "agent" if str(kind).lower() == "agent" else "human" + value = policy or read_policy() + return any( + str(item.get("kind") or "") == selected and str(item.get("session_id") or "") == sid + for item in value.get("session_tombstones") or [] + if isinstance(item, dict) + ) + + +def add_session_tombstones(kind: str, session_ids: list[str], *, reason: str) -> dict[str, Any]: + selected = "agent" if str(kind).lower() == "agent" else "human" + ids = [str(x).strip()[:240] for x in session_ids if str(x).strip()] + with _LOCK: + value = read_policy() + existing = { + (str(item.get("kind") or ""), str(item.get("session_id") or "")): item + for item in value.get("session_tombstones") or [] + if isinstance(item, dict) + } + for sid in ids: + existing[(selected, sid)] = { + "kind": selected, + "session_id": sid, + "deleted_at": _now(), + "reason": str(reason or "retention")[:120], + } + value["session_tombstones"] = list(existing.values())[-_MAX_SESSION_TOMBSTONES:] + value["history_generation"] = int(value.get("history_generation") or 1) + 1 + value["updated_at"] = _now() + return _write_unlocked(value) + + +def bump_history_generation(*, reason: str = "history_changed") -> int: + with _LOCK: + value = read_policy() + value["history_generation"] = int(value.get("history_generation") or 1) + 1 + value["last_history_change_reason"] = str(reason or "history_changed")[:120] + value["updated_at"] = _now() + return int(_write_unlocked(value)["history_generation"]) + + +def history_generation() -> int: + return int(read_policy().get("history_generation") or 1) + + +def prepare_recordable_event(event: dict[str, Any], *, now: datetime | None = None) -> dict[str, Any] | None: + value = read_policy() + kind = event_kind(event) + session_id = str(event.get("session_id") or "").strip() + if session_is_tombstoned(kind, session_id, policy=value): + return None + retention = retention_for_kind(kind, policy=value) + if retention.get("mode") == "days": + observed = _parse(event.get("observed_at")) + cutoff = (now or _now_dt()) - timedelta(days=int(retention.get("days") or 1)) + if observed is not None and observed < cutoff: + return None + return copy.deepcopy(event) + + +def set_ai_history_access( + *, + mode: str, + since: str | None = None, + until: str | None = None, + expires_minutes: int | None = 60, +) -> dict[str, Any]: + selected = str(mode or "off").strip().lower() + if selected not in _AI_ACCESS_MODES: + raise ValueError("AI history access mode must be off, selected_range, or all_saved") + normalized_since: str | None = None + normalized_until: str | None = None + if selected == "selected_range": + start = _parse(since) + end = _parse(until) + if start is None or end is None or end <= start: + raise ValueError("selected history access requires a valid since/until range") + normalized_since = start.isoformat() + normalized_until = end.isoformat() + minutes = None if expires_minutes is None else int(expires_minutes) + if minutes is not None and (minutes < 1 or minutes > 24 * 60): + raise ValueError("AI history access expiry must be between 1 minute and 24 hours") + granted = _now_dt() + expires = (granted + timedelta(minutes=minutes)).isoformat() if minutes is not None else None + with _LOCK: + value = read_policy() + previous = dict(value.get("ai_history_access") or {}) + changed = ( + str(previous.get("mode") or "off") != selected + or (str(previous.get("since") or "") or None) != normalized_since + or (str(previous.get("until") or "") or None) != normalized_until + ) + value["ai_history_access"] = { + "mode": selected, + "since": normalized_since, + "until": normalized_until, + "granted_at": granted.isoformat() if selected != "off" else None, + "expires_at": expires if selected != "off" else None, + } + if changed: + # Pulse cursors remember delivered finding versions. A disclosure-scope + # change must rebaseline them so findings hidden under a narrower lease + # cannot remain incorrectly marked as already delivered later. + value["history_generation"] = int(value.get("history_generation") or 1) + 1 + value["last_history_change_reason"] = "ai_history_access_changed" + value["updated_at"] = _now() + return _write_unlocked(value)["ai_history_access"] + + +def active_ai_history_access(*, now: datetime | None = None) -> dict[str, Any]: + value = read_policy() + access = dict(value.get("ai_history_access") or {}) + selected = str(access.get("mode") or "off") + expires = _parse(access.get("expires_at")) + current = now or _now_dt() + if selected != "off" and expires is not None and current >= expires: + set_ai_history_access(mode="off", expires_minutes=None) + return {"mode": "off", "since": None, "until": None, "granted_at": None, "expires_at": None} + return access + + +__all__ = [ + "active_ai_history_access", + "add_session_tombstones", + "bump_history_generation", + "event_kind", + "history_generation", + "initialize_policy", + "prepare_recordable_event", + "read_policy", + "retention_for_kind", + "set_ai_history_access", + "session_is_tombstoned", + "update_retention", +] diff --git a/tests/js/agent_surface_adapters.test.mjs b/tests/js/agent_surface_adapters.test.mjs new file mode 100644 index 00000000..1f1f4478 --- /dev/null +++ b/tests/js/agent_surface_adapters.test.mjs @@ -0,0 +1,47 @@ +import assert from 'node:assert/strict'; +import fs from 'node:fs'; + +const source=fs.readFileSync('browser_extension/agent_surface_adapters.js','utf8'); +const background=fs.readFileSync('browser_extension/background.js','utf8'); +const manifest=JSON.parse(fs.readFileSync('browser_extension/manifest.json','utf8')); +assert.doesNotThrow(()=>new Function(source)); +assert.ok(manifest.content_scripts[0].js.includes('agent_surface_adapters.js')); +assert.deepEqual(manifest.permissions,['tabs','webNavigation','storage','alarms']); + +assert.match(source,/aria-busy/); +assert.match(source,/\[role=["']alert/); +assert.match(source,/\[role=["']dialog/); +assert.doesNotMatch(source,/querySelector\([^\n]*\.[A-Za-z_-][A-Za-z0-9_-]{4,}/); +assert.doesNotMatch(source,/FileSystem|FileReader|webkitRelativePath/); + +new Function(source)(); +const api=globalThis.__OWG_AGENT_SURFACE_ADAPTERS_FOR_TESTS__; +assert.ok(api); +assert.equal(api.providerForHost('chatgpt.com')?.key,'chatgpt'); +assert.equal(api.providerForHost('claude.ai')?.key,'claude'); +assert.equal(api.providerForHost('copilot.microsoft.com')?.key,'microsoft_copilot'); +assert.equal(api.providerForHost('lovable.dev')?.key,'lovable'); +assert.equal(api.providerForHost('gemini.google.com')?.key,'gemini'); +assert.equal(api.providerForHost('example.com'),null); + +globalThis.location={origin:'https://chatgpt.com',hostname:'chatgpt.com'}; +const payload=api.structuralPayload({key:'chatgpt',name:'ChatGPT'},'agent_run_started','web-deadbeef','send_control'); +assert.equal(payload.type,'workflow_observer_event'); +assert.equal(payload.action,'agent_run_started'); +assert.equal(payload.metadata.agent_provider,'chatgpt'); +assert.equal(payload.metadata.agent_run_id,'web-deadbeef'); +assert.equal(payload.page.url,'https://chatgpt.com/'); +assert.equal(payload.page.title,'ChatGPT'); +// background.js consumes these exact top-level fields, so an accidental payload +// wrapper would silently downgrade lifecycle events to generic browser_event. +assert.match(background,/message\.observed_at/); +assert.match(background,/message\.action/); +assert.match(background,/message\.metadata/); +assert.equal('payload' in payload,false); + +const serialized=JSON.stringify(payload).toLowerCase(); +for(const forbidden of ['prompt_text','response_text','tool_arguments','tool_result','chain_of_thought','clipboard_contents','file_path']){ + assert.equal(serialized.includes(forbidden),false,`must not send ${forbidden}`); +} +assert.match(source,/No DOM text or user\/model content is included/); +console.log('browser agent adapters remain structural and content-free'); diff --git a/tests/test_ai_context_mcp_v092.py b/tests/test_ai_context_mcp_v092.py index 2ae8c102..6fd29d6a 100644 --- a/tests/test_ai_context_mcp_v092.py +++ b/tests/test_ai_context_mcp_v092.py @@ -29,6 +29,7 @@ CONTEXT_TOOL_ARGS = { "get_current_work_context": {}, "get_context_pulse": {}, + "list_history": {}, "search_work": {"query": "Contract"}, "get_workflow_trace": {"limit": 50}, "get_work_profile": {}, @@ -122,8 +123,20 @@ def test_every_context_mcp_tool_respects_detail_level(tmp_path): headers = {"Authorization": f"Bearer {token}"} try: _wait(base + "/health", api) + retention = httpx.put( + base + "/v1/history-policy", + json={"human_mode": "forever", "agent_mode": "forever", "onboarding_complete": True}, + headers=headers, + ) + assert retention.status_code == 200, retention.text assert httpx.post(base + "/v1/events", json={"events": _events(now)}, headers=headers).status_code == 200 assert httpx.post(base + "/v1/ai-access", json={"enabled": True}, headers=headers).json()["enabled"] is True + lease = httpx.post( + base + "/v1/history/ai-access", + json={"mode": "all_saved", "expires_minutes": 60}, + headers=headers, + ) + assert lease.status_code == 200 and lease.json()["access"]["mode"] == "all_saved" async def call_all() -> dict[str, tuple[str, dict]]: params = StdioServerParameters( @@ -190,7 +203,8 @@ async def call_all() -> dict[str, tuple[str, dict]]: finally: _stop(api) - # Raw local evidence is untouched by any of this. + # Raw local evidence is untouched by AI redaction. This test explicitly chose + # forever retention above, so shutdown retention cleanup must not remove it. db = sqlite3.connect(data / "workflow_observer.db") titles = {row[0] for row in db.execute("SELECT window_title FROM events")} labels = {json.loads(row[0] or "{}").get("target", {}).get("label") for row in db.execute("SELECT metadata_json FROM events")} diff --git a/tests/test_compact_mcp_v087.py b/tests/test_compact_mcp_v087.py index 09d36ad3..8084eb87 100644 --- a/tests/test_compact_mcp_v087.py +++ b/tests/test_compact_mcp_v087.py @@ -13,6 +13,7 @@ DEFAULT_TOOLS = { "get_current_work_context", "get_context_pulse", + "list_history", "search_work", "get_workflow_trace", "get_work_profile", diff --git a/tests/test_extension_capture.py b/tests/test_extension_capture.py index 1af41f7d..cef92aca 100644 --- a/tests/test_extension_capture.py +++ b/tests/test_extension_capture.py @@ -48,8 +48,8 @@ def test_browser_delivery_is_durable_idempotent_and_sanitized(): assert "flushBrowserQueue" in background assert "sanitizePendingBrowserQueue" in background assert "sensor_version" in background - assert manifest["version"] == "1.11.0" - assert manifest["version_name"] == "1.11.0-v58-metadata-signals" + assert manifest["version"] == "1.12.0" + assert manifest["version_name"] == "1.12.0-v94-agent-lifecycle" assert manifest["background"]["service_worker"] == "secure_background.js" diff --git a/tests/test_history_retention_v094.py b/tests/test_history_retention_v094.py new file mode 100644 index 00000000..f60f990b --- /dev/null +++ b/tests/test_history_retention_v094.py @@ -0,0 +1,102 @@ +from __future__ import annotations + +import os +import subprocess +import sys +from pathlib import Path + +ROOT=Path(__file__).resolve().parents[1] + +def _run(code:str,tmp_path:Path,timeout:int=120)->str: + env=os.environ.copy();env.update({"WORKFLOW_OBSERVER_DATA":str(tmp_path/'data'),"WORKFLOW_OBSERVER_AUTH_DIR":str(tmp_path/'auth'),"WORKFLOW_OBSERVER_CONFIG":str(tmp_path/'config.json'),"PYTHONPATH":str(ROOT)}) + result=subprocess.run([sys.executable,'-c',code],cwd=ROOT,env=env,text=True,capture_output=True,timeout=timeout) + assert result.returncode==0,f"stdout={result.stdout}\nstderr={result.stderr}" + return result.stdout + + +def test_new_install_ephemeral_existing_install_preserved(tmp_path): + _run(r''' +from datetime import datetime, timezone +from server.db import init_db,insert_events +from server.history_retention import initialize_history_retention +from shared.history_policy import read_policy +init_db();p=initialize_history_retention();assert p['human_retention']['mode']=='ephemeral';assert p['agent_retention']['mode']=='ephemeral';assert p['onboarding_complete'] is False +''',tmp_path/'new') + _run(r''' +from datetime import datetime, timezone +from server.db import init_db,insert_events +from server.history_retention import initialize_history_retention +init_db();insert_events([{'event_id':'old','observed_at':datetime.now(timezone.utc).isoformat(),'device_id':'d','session_id':'s','app':'Editor','event_type':'focus_span','duration_seconds':1,'metadata':{}}]);p=initialize_history_retention();assert p['human_retention']['mode']=='forever';assert p['upgrade_preserved_existing_history'] is True +''',tmp_path/'existing') + + +def test_uninitialized_legacy_runtime_does_not_delete_completed_agent_run(tmp_path): + _run(r''' +from datetime import datetime, timedelta, timezone +from server.db import init_db,connect +from server.agent_ingest import ingest_agent_payloads +from shared.history_policy import read_policy +init_db();assert read_policy()['agent_retention']['mode']=='forever' +base=datetime.now(timezone.utc)-timedelta(minutes=2) +events=[{'observed_at':base.isoformat(),'agent_name':'LegacyAgent','provider':'test','framework':'custom','operation':'run_started','status':'running','observation_level':'native_trace','run_id':'legacy-run','session_id':'legacy-run','tool_category':'none'},{'observed_at':(base+timedelta(seconds=1)).isoformat(),'agent_name':'LegacyAgent','provider':'test','framework':'custom','operation':'run_finished','status':'success','observation_level':'native_trace','run_id':'legacy-run','session_id':'legacy-run','tool_category':'none'}] +r=ingest_agent_payloads(events);assert r['inserted']==2 +with connect() as c: assert c.execute("SELECT COUNT(*) FROM events WHERE session_id='legacy-run'").fetchone()[0]==2 +''',tmp_path/'legacy') + + +def test_ephemeral_cleanup_tombstones_late_delivery_and_bumps_generation(tmp_path): + _run(r''' +from datetime import datetime, timezone +from server.db import init_db,insert_events,connect +from server.history_retention import initialize_history_retention,cleanup_expired_history +from shared.capture_control import filter_recordable +from shared.history_policy import history_generation +init_db();initialize_history_retention();now=datetime.now(timezone.utc).isoformat();event={'event_id':'e1','observed_at':now,'device_id':'d','session_id':'ephemeral-1','app':'Editor','event_type':'focus_span','duration_seconds':1,'metadata':{}};insert_events([event]);g=history_generation();cleanup_expired_history(startup=True);assert history_generation()>g +with connect() as c: assert c.execute('SELECT COUNT(*) FROM events').fetchone()[0]==0 +late=dict(event,event_id='e2');kept,suppressed=filter_recordable([late]);assert kept==[] and suppressed==1 +''',tmp_path) + + +def test_history_sessions_activity_blocks_and_agent_observation(tmp_path): + _run(r''' +from datetime import datetime,timedelta,timezone +from server.db import init_db,insert_events +from shared.history_policy import initialize_policy,update_retention +from server.history_retention import list_history +from server.agent_ingest import ingest_agent_payloads +init_db();initialize_policy(has_existing_evidence=False);update_retention(human_mode='days',human_days=90,agent_mode='days',agent_days=90) +base=datetime.now(timezone.utc)-timedelta(hours=3) +def h(eid,offset): return {'event_id':eid,'observed_at':(base+timedelta(minutes=offset)).isoformat(),'device_id':'d','session_id':'human-1','app':'Editor','window_title':'Work','event_type':'focus_span','duration_seconds':60,'metadata':{'activity':{'foreground_seconds':60,'engaged_seconds':50}}} +insert_events([h('h1',0),h('h2',10),h('h3',55)]) +ingest_agent_payloads([{'observed_at':(base+timedelta(minutes=5)).isoformat(),'agent_name':'ChatGPT','provider':'openai','framework':'chatgpt_web','operation':'run_started','status':'running','observation_level':'os_observed','run_id':'web-test1','session_id':'web-test1','tool_category':'none'},{'observed_at':(base+timedelta(minutes=6)).isoformat(),'agent_name':'ChatGPT','provider':'openai','framework':'chatgpt_web','operation':'run_finished','status':'success','observation_level':'os_observed','run_id':'web-test1','session_id':'web-test1','tool_category':'none'}]) +h=list_history(limit=20);human=next(x for x in h['sessions'] if x['kind']=='human');agent=next(x for x in h['sessions'] if x['kind']=='agent');assert len(human['activity_blocks'])==2,human;assert human['event_count']==3;assert agent['observation_level']=='os_observed';assert h['derived_task_labels_used'] is False +''',tmp_path) + + +def test_ai_saved_history_lease_is_explicit_range_bounded_and_expires(tmp_path): + _run(r''' +from datetime import datetime,timedelta,timezone +from shared.history_policy import initialize_policy,set_ai_history_access,active_ai_history_access +initialize_policy(has_existing_evidence=False);assert active_ai_history_access()['mode']=='off' +start=datetime.now(timezone.utc)-timedelta(days=5);end=start+timedelta(days=2) +a=set_ai_history_access(mode='selected_range',since=start.isoformat(),until=end.isoformat(),expires_minutes=60);assert a['mode']=='selected_range';assert a['since'] and a['until'];set_ai_history_access(mode='off',expires_minutes=None);assert active_ai_history_access()['mode']=='off' +''',tmp_path) + + +def test_browser_agent_projection_is_structural_only(tmp_path): + _run(r''' +from datetime import datetime,timezone +from server.main import BrowserEvent,BrowserPage,BrowserTarget +from server.browser_agent_projection import _agent_payload +b=BrowserEvent(observed_at=datetime.now(timezone.utc).isoformat(),action='agent_run_started',page=BrowserPage(hostname='chatgpt.com',pathname='/',title='ChatGPT'),target=BrowserTarget(role='agent-lifecycle',label='ChatGPT'),metadata={'agent_provider':'chatgpt','agent_run_id':'web-deadbeef','state_source':'send_control'}) +p=_agent_payload(b);assert p['agent_name']=='ChatGPT';assert p['provider']=='openai';assert p['framework']=='chatgpt_web';assert p['observation_level']=='os_observed';assert p['operation']=='run_started';assert not any(k in p for k in ('prompt','response','content','tool_arguments','tool_result')) +''',tmp_path) + + +def test_history_architecture_is_registered_and_no_filesystem_sensor_added(): + runner=(ROOT/'server'/'enterprise_runner.py').read_text(encoding='utf-8') + assert 'server.history_routes' in runner and 'server.browser_agent_projection' in runner + adapter=(ROOT/'browser_extension'/'agent_surface_adapters.js').read_text(encoding='utf-8') + assert 'FileSystem' not in adapter and 'FileReader' not in adapter + policy=(ROOT/'shared'/'history_policy.py').read_text(encoding='utf-8') + assert 'ai_history_access' in policy and 'session_tombstones' in policy diff --git a/tests/test_mcpb_manifest_v087.py b/tests/test_mcpb_manifest_v087.py index 07acfa26..3c949535 100644 --- a/tests/test_mcpb_manifest_v087.py +++ b/tests/test_mcpb_manifest_v087.py @@ -8,6 +8,7 @@ EXPECTED_COMPACT_TOOLS = { "get_current_work_context", "get_context_pulse", + "list_history", "search_work", "get_workflow_trace", "get_work_profile", diff --git a/tests/test_release_version_v087.py b/tests/test_release_version_v087.py index 9a1c07bc..19b00684 100644 --- a/tests/test_release_version_v087.py +++ b/tests/test_release_version_v087.py @@ -6,7 +6,7 @@ ROOT = Path(__file__).resolve().parents[1] -EXPECTED_VERSION = "0.93.0" +EXPECTED_VERSION = "0.94.0" def test_release_version_sources_are_aligned(): diff --git a/tests/test_v423_live_active_tab_surface.py b/tests/test_v423_live_active_tab_surface.py index f9cfa69d..1eacf2b5 100644 --- a/tests/test_v423_live_active_tab_surface.py +++ b/tests/test_v423_live_active_tab_surface.py @@ -11,8 +11,8 @@ def test_browser_heartbeat_source_reports_sanitized_active_tab(): assert 'ext.tabs.query({active: true, lastFocusedWindow: true})' in background assert 'page,' in background assert 'safeUrl(tab?.url || "")' in background - assert '"version": "1.11.0"' in manifest - assert '"version_name": "1.11.0-v58-metadata-signals"' in manifest + assert '"version": "1.12.0"' in manifest + assert '"version_name": "1.12.0-v94-agent-lifecycle"' in manifest def test_browser_heartbeat_keeps_safe_active_page_and_surface(): diff --git a/tests/test_v49_mcp_stdio.py b/tests/test_v49_mcp_stdio.py index 064b23ba..10f4ce65 100644 --- a/tests/test_v49_mcp_stdio.py +++ b/tests/test_v49_mcp_stdio.py @@ -59,11 +59,17 @@ def test_real_stdio_mcp_lists_tools_denies_then_reads_when_enabled(tmp_path): "PYTHONPATH": str(ROOT), }) api = subprocess.Popen( - [sys.executable, "-m", "uvicorn", "server.secure_app:app", "--host", "127.0.0.1", "--port", str(port)], - cwd=ROOT, env=env, stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, + [sys.executable, "-m", "server.enterprise_runner", "--host", "127.0.0.1", "--port", str(port)], + cwd=ROOT, env=env, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, ) _wait(base + "/health", api) headers = {"Authorization": f"Bearer {token}"} + retention = httpx.put( + base + "/v1/history-policy", + json={"human_mode": "forever", "agent_mode": "forever", "onboarding_complete": True}, + headers=headers, + ) + assert retention.status_code == 200, retention.text event = { "event_id": "stdio-real-1", "observed_at": "2026-09-18T16:00:00+00:00", "device_id": "d", "sensor_id": "browser:test", "source": "browser_extension", @@ -113,6 +119,7 @@ async def exercise() -> None: assert expected_procedural <= names assert expected_agent_inspection <= names assert len(names) == 24 + assert "list_history" not in names denied = await session.call_tool("get_workflow_trace", {"limit": 10}) assert denied.is_error is True @@ -120,6 +127,12 @@ async def exercise() -> None: enabled = httpx.post(base + "/v1/ai-access", json={"enabled": True}, headers=headers) assert enabled.status_code == 200 and enabled.json()["enabled"] is True + lease = httpx.post( + base + "/v1/history/ai-access", + json={"mode": "all_saved", "expires_minutes": 60}, + headers=headers, + ) + assert lease.status_code == 200 and lease.json()["access"]["mode"] == "all_saved" result = await session.call_tool("get_workflow_trace", {"limit": 10}) assert result.is_error is False