From 339c386a4478cf9a396113b86789a55e2001290d Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:15:35 +0200 Subject: [PATCH 01/13] Capture correctness: one collector, away spans, document boundaries, honest health - One collector per data folder via an OS advisory lock; the supervisor backs off instead of respawning a blocked worker every second. - Away spans after 5 minutes without input (OS input clock) or when the screen locks; duration_seconds keeps its wall-clock meaning; away spans never sync to a Gateway. - Debounced, noise-normalized document boundaries within the same app (new default; "application" follows it, "application_only" keeps the old behavior); title privacy modes respected. - capture_health only for real degradation and permission changes; routine counters move to diagnostics. - macOS Accessibility/Input Monitoring preflight: blocked sensors are stated in startup output, heartbeat, capture health and the Recording pill. - Claude Code UserPromptSubmit/Stop as turn boundaries keyed by prompt_id; outdated OpenWorkGraph hook entries refreshed at launch when Observe is on. - Per-channel agent telemetry diagnostics (counts and reason codes only). - Lease-gated, bounded agent spool: hooks spool only while OpenWorkGraph is running and recording; flushed through normal ingest rules. --- adapters/_agent_client.py | 30 +- adapters/claude_code_hook.py | 2 + collector/boundaries.py | 41 ++ collector/instance_lock.py | 73 +++ collector/interactions.py | 4 + collector/main.py | 438 +++++++++++++----- collector/permissions.py | 65 +++ collector/platform.py | 60 +++ collector/secure_main.py | 13 + collector/secure_worker.py | 3 +- config.example.json | 5 +- connector/policy.py | 8 + dashboard/connections.js | 31 +- dashboard/gateway_panel.js | 6 +- docs/CHANGELOG_V098.md | 73 +++ docs/NATIVE_AGENT_ADAPTERS.md | 46 +- docs/PRIVACY_AND_DATA.md | 1 + server/agent_capture_runtime.py | 172 +++++++ server/agent_config_writer.py | 31 +- server/agent_routes.py | 75 ++- server/agent_spool.py | 222 +++++++++ server/agent_telemetry_diagnostics.py | 104 +++++ server/collector_status.py | 17 + server/enterprise_app.py | 4 + server/enterprise_runner.py | 8 +- shared/claude_code_adapter.py | 21 + tests/conftest.py | 1 + tests/js/capture_correctness.test.mjs | 21 + ...test_agent_dashboard_control_plane_v088.py | 2 +- tests/test_agent_delivery_v098.py | 288 ++++++++++++ tests/test_capture_correctness_v098.py | 329 +++++++++++++ tests/test_claude_code_adapter_v060.py | 31 +- 32 files changed, 2074 insertions(+), 151 deletions(-) create mode 100644 collector/boundaries.py create mode 100644 collector/instance_lock.py create mode 100644 collector/permissions.py create mode 100644 docs/CHANGELOG_V098.md create mode 100644 server/agent_capture_runtime.py create mode 100644 server/agent_spool.py create mode 100644 server/agent_telemetry_diagnostics.py create mode 100644 server/collector_status.py create mode 100644 tests/js/capture_correctness.test.mjs create mode 100644 tests/test_agent_delivery_v098.py create mode 100644 tests/test_capture_correctness_v098.py diff --git a/adapters/_agent_client.py b/adapters/_agent_client.py index 56fcfc3..125ffd4 100644 --- a/adapters/_agent_client.py +++ b/adapters/_agent_client.py @@ -6,6 +6,7 @@ import json import os from urllib.parse import urlparse +from urllib.error import HTTPError, URLError from urllib.request import Request, urlopen from server.agent_auth import ensure_agent_ingest_token @@ -35,7 +36,7 @@ def _token() -> str: return os.getenv("OWG_AGENT_INGEST_TOKEN", "").strip() or ensure_agent_ingest_token() -def post_json(path: str, payload: dict, *, timeout: float = 0.75) -> dict: +def post_json(path: str, payload: dict, *, timeout: float = 0.75, channel: str = "") -> dict: if not path.startswith("/agent-ingest/"): raise ValueError("native adapters may only use agent-ingest write routes") body = json.dumps(payload, ensure_ascii=False, separators=(",", ":")).encode("utf-8") @@ -45,6 +46,7 @@ def post_json(path: str, payload: dict, *, timeout: float = 0.75) -> dict: headers={ "Authorization": f"Bearer {_token()}", "Content-Type": "application/json", + **({"X-OWG-Channel": channel} if channel else {}), }, method="POST", ) @@ -54,7 +56,29 @@ def post_json(path: str, payload: dict, *, timeout: float = 0.75) -> dict: return value if isinstance(value, dict) else {} -def post_agent_events(events: list[dict], *, timeout: float = 0.75) -> dict: +def post_agent_events(events: list[dict], *, timeout: float = 0.75, spool: bool = True) -> dict: + """Deliver structural agent events; briefly spool them if OpenWorkGraph is busy. + + Spooling happens only for transport failures and server errors, and only + under a valid recording lease (see server.agent_spool). A rejection (4xx) is + final: retrying it later would not make it acceptable. + """ if not events: return {"status": "ignored", "received": 0} - return post_json("/agent-ingest/v1/events", {"events": events}, timeout=timeout) + try: + framework = str(events[0].get("framework") or "") if isinstance(events[0], dict) else "" + channel = "claude_code_hooks" if framework == "claude-code" else "agent_events" + return post_json("/agent-ingest/v1/events", {"events": events}, timeout=timeout, channel=channel) + except HTTPError as exc: + if exc.code < 500 or not spool: + raise + failure: Exception = exc + except (URLError, TimeoutError, ConnectionError, OSError) as exc: + if not spool: + raise + failure = exc + from server.agent_spool import spool_events + + if spool_events(events): + return {"status": "spooled", "received": 0, "spooled": len(events)} + raise failure diff --git a/adapters/claude_code_hook.py b/adapters/claude_code_hook.py index 62b0830..daa1270 100644 --- a/adapters/claude_code_hook.py +++ b/adapters/claude_code_hook.py @@ -20,6 +20,8 @@ SUPPORTED_EVENTS = [ "SessionStart", "SessionEnd", + "UserPromptSubmit", + "Stop", "PostToolUse", "PostToolUseFailure", "PermissionRequest", diff --git a/collector/boundaries.py b/collector/boundaries.py new file mode 100644 index 0000000..e7411c8 --- /dev/null +++ b/collector/boundaries.py @@ -0,0 +1,41 @@ +from __future__ import annotations + +"""Deciding when a same-app window title is a different document. + +Titles change for many reasons that are not a new document: unread counts, +unsaved markers, "Edited" suffixes, progress percentages, "Not Responding". +``document_key`` removes that noise so only material changes count, and the +collector additionally requires the new title to persist for a few polls. +""" + +import re + +_BADGE = r"(?:[\(\[]\s*\d{1,5}\+?\s*(?:unread|new|olästa|nya)?\s*[\)\]])" +_LEADING = re.compile(rf"^\s*(?:{_BADGE}|[•●◉∙*✱]+)\s*", re.I) +_INNER_BADGE = re.compile(rf"\s*{_BADGE}", re.I) +_TRAILING = re.compile( + r"\s*(?:" + r"[—–-]\s*(?:edited|modified|saved|saving…?|autosaved|not responding|redigerad|sparad|svarar inte)" + r"|\(\s*(?:not responding|read-only|edited|svarar inte|skrivskyddad)\s*\)" + r"|[•●*]" + r"|\d{1,3}\s?%" + r")\s*$", + re.I, +) + + +def document_key(title: str) -> str: + """A normalized title for comparing documents within one application.""" + text = re.sub(r"\s+", " ", str(title or "")).strip() + previous = None + while previous != text: + previous = text + text = _LEADING.sub("", text) + text = _TRAILING.sub("", text) + text = _INNER_BADGE.sub("", text).strip() + return text.casefold() + + +def is_material_change(current_key: str, candidate_key: str) -> bool: + """A switch to an empty title (a dialog, an untitled moment) is not a document change.""" + return bool(candidate_key) and candidate_key != current_key diff --git a/collector/instance_lock.py b/collector/instance_lock.py new file mode 100644 index 0000000..58beda0 --- /dev/null +++ b/collector/instance_lock.py @@ -0,0 +1,73 @@ +from __future__ import annotations + +"""One recording collector per data directory. + +An OS advisory lock on a file inside the data directory is held for the +collector's lifetime. The OS releases it when the process exits or crashes, so +there is no stale PID file to clean up. ``data/live`` and ``data/demo`` are +different directories and therefore different locks. +""" + +import os +from pathlib import Path + +LOCK_NAME = ".collector.lock" +# Exit code for "another collector already records this data directory". +EXIT_ALREADY_RUNNING = 75 + + +class CollectorLock: + def __init__(self, data_dir: Path) -> None: + self.path = Path(data_dir) / LOCK_NAME + self._handle = None + + def acquire(self) -> bool: + self.path.parent.mkdir(parents=True, exist_ok=True) + handle = open(self.path, "a+b") + try: + if os.name == "nt": + import msvcrt + + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_NBLCK, 1) + else: + import fcntl + + fcntl.flock(handle.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + except OSError: + handle.close() + return False + try: + handle.seek(0) + handle.truncate() + handle.write(str(os.getpid()).encode("ascii")) # informational only + handle.flush() + except OSError: + pass + self._handle = handle + return True + + def release(self) -> None: + handle, self._handle = self._handle, None + if handle is None: + return + try: + if os.name == "nt": + import msvcrt + + handle.seek(0) + msvcrt.locking(handle.fileno(), msvcrt.LK_UNLCK, 1) + else: + import fcntl + + fcntl.flock(handle.fileno(), fcntl.LOCK_UN) + except OSError: + pass + finally: + handle.close() + + def __enter__(self) -> "CollectorLock": + return self + + def __exit__(self, *_exc) -> None: + self.release() diff --git a/collector/interactions.py b/collector/interactions.py index 973bef3..f12859e 100644 --- a/collector/interactions.py +++ b/collector/interactions.py @@ -99,6 +99,10 @@ def record(self, kind: str, occurred_mono: float | None = None) -> None: while self._events and self._events[0][0] < cutoff: self._events.popleft() + def last_input_mono(self) -> float | None: + with self._lock: + return self._events[-1][0] if self._events else None + def summarize( self, start_mono: float, diff --git a/collector/main.py b/collector/main.py index d52cfc5..1ed824e 100644 --- a/collector/main.py +++ b/collector/main.py @@ -13,7 +13,7 @@ import httpx -from .platform import active_window +from .platform import active_window, screen_locked, system_idle_seconds from .accessibility import element_at_position from .interactions import ( InteractionSensor, @@ -25,6 +25,9 @@ from .privacy import should_exclude, title_for_mode from .identity import load_or_create_identity from .outbox import EventOutbox +from .boundaries import document_key, is_material_change +from .instance_lock import EXIT_ALREADY_RUNNING, CollectorLock +from .permissions import missing as missing_permissions, sensor_permissions from browser_utils import is_browser_app, normalized_browser_title from sensitive_identifiers import sanitize_event_identifiers @@ -51,7 +54,17 @@ def load_config(path: Path) -> dict: "heartbeat_seconds": 5, "focus_checkpoint_seconds": 120, "capture_gap_seconds": 30, - "change_detection": "application", + # "application_and_document": app changes plus debounced, noise-normalized + # document changes within the same app. "application_only": app changes + # only (titles of long spans go stale). "application_and_title": every raw + # title change. "application" was the shipped default copied into every + # config.json, so it follows the default rather than pinning old behavior. + "change_detection": "application_and_document", + "document_debounce_polls": 2, + # A span ends after this long without any keyboard/mouse input anywhere, + # or when the screen locks; the time until input returns is an away span. + "away_detection_enabled": True, + "away_after_seconds": 300, "screenshot_interval_seconds": 20, "screenshots_enabled": False, "upload_screenshots": False, @@ -208,16 +221,20 @@ def _public_window(w, cfg: dict) -> dict: cfg["excluded_apps"], cfg["excluded_title_patterns"], ) + mode = cfg.get("window_title_mode", "full") return { "app": "Excluded" if excluded else (w.app or "Unknown"), - "window_title": "" if excluded else title_for_mode( - w.title, cfg.get("window_title_mode", "full") - ), + "window_title": "" if excluded else title_for_mode(w.title, mode), "excluded": excluded, + # In-memory only (event builders copy app/window_title/excluded, never + # this): a noise-normalized document identity for same-app boundaries. + # Empty when titles are excluded or the user chose not to record them. + "document_key": "" if excluded or mode == "none" else document_key(w.title), } def _change_key(state: dict, cfg: dict): + """Key whose change ends a span immediately (application, browser page, raw title).""" app = state["app"] if is_browser_app(app, cfg.get("browser_app_patterns")): return (app, normalized_browser_title(state["window_title"])) @@ -226,6 +243,14 @@ def _change_key(state: dict, cfg: dict): return (app,) +def _document_tracking(state: dict, cfg: dict) -> bool: + return ( + cfg.get("change_detection", "application_and_document") in {"application_and_document", "application"} + and not is_browser_app(state["app"], cfg.get("browser_app_patterns")) + and not state.get("excluded") + ) + + def _focus_span_event( *, state: dict, @@ -280,7 +305,62 @@ def _capture_gap_event(*, started_at: str, duration_seconds: float, cfg: dict, s } -def _capture_health_event(*, capture_health: dict[str, int], cfg: dict, session_id: str) -> dict: +def _away_span_event( + *, + started_at: str, + duration_seconds: float, + cfg: dict, + session_id: str, + reason: str, + idle_source: str, + boundary_reason: str, +) -> dict: + return { + "event_id": str(uuid.uuid4()), + "observed_at": started_at, + **_identity_fields(cfg), + "session_id": session_id, + "app": "Away", + "window_title": "", + "event_type": "away_span", + "duration_seconds": max(0.0, float(duration_seconds)), + "screenshot_path": None, + "metadata": { + "source": "desktop", + "reason": reason, + "away_after_seconds": float(cfg.get("away_after_seconds", 300)), + "idle_source": idle_source, + "focus_boundary": boundary_reason, + "interpretation": ( + "No keyboard or mouse input for at least away_after_seconds, or the screen was locked. " + "The focus span before it keeps the first away_after_seconds, which covers reading " + "and thinking without input. Never shared with an organization Gateway." + ), + "privacy": {"key_identities": False, "typed_values": False}, + }, + } + + +# Counters that mean OpenWorkGraph missed or could not see something. Only these +# (and permission changes) produce capture_health evidence; routine counters stay +# in the heartbeat and in the event's diagnostics block. +DEGRADATION_KEYS = ( + "interaction_worker_errors", + "interaction_queue_dropped", + "clipboard_queue_dropped", + "active_window_unavailable", +) +HEALTH_REPEAT_SECONDS = 600.0 + + +def _capture_health_event( + *, + capture_health: dict[str, int], + cfg: dict, + session_id: str, + diagnostics: dict[str, int] | None = None, + permissions: dict[str, bool | None] | None = None, +) -> dict: return { "event_id": str(uuid.uuid4()), "observed_at": utcnow(), @@ -294,7 +374,10 @@ def _capture_health_event(*, capture_health: dict[str, int], cfg: dict, session_ "metadata": { "source": "collector", "capture_health": dict(capture_health), - "interpretation": "Nonzero counters indicate known capture degradation; zero counters are not proof that no external capture limitation existed.", + "diagnostics": dict(diagnostics or {}), + "permissions": dict(permissions or {}), + "missing_permissions": missing_permissions(dict(permissions or {})), + "interpretation": "Nonzero counters or missing permissions indicate known capture degradation; zero counters are not proof that no external capture limitation existed.", "privacy": {"key_identities": False, "typed_values": False, "clipboard_contents": False}, }, } @@ -445,33 +528,67 @@ def _interaction_worker( q.task_done() -def run(config_path: Path) -> None: +def run(config_path: Path) -> int: global STOP STOP = False + lock = CollectorLock(LOCAL_DIR) + if not lock.acquire(): + print( + f"Another OpenWorkGraph collector is already recording {LOCAL_DIR}; " + "this one is not starting, so nothing is recorded twice." + ) + return EXIT_ALREADY_RUNNING + try: + _run_locked(config_path) + finally: + lock.release() + return 0 + + +def _run_locked(config_path: Path) -> None: + global STOP cfg = load_config(config_path) identity = load_or_create_identity(LOCAL_DIR, cfg) cfg.update(identity) session_id = str(uuid.uuid4()) activity_tracker = ActivityTracker() - capture_health = { - "interaction_worker_errors": 0, - "interaction_queue_dropped": 0, - "clipboard_queue_dropped": 0, + capture_health = {key: 0 for key in DEGRADATION_KEYS} + diagnostics = { "capture_gap_count": 0, "focus_checkpoint_count": 0, + "document_boundary_count": 0, + "away_count": 0, } - last_capture_health_evidence: dict[str, int] = dict(capture_health) + permissions = sensor_permissions() + last_permission_check = time.monotonic() + last_health_evidence: dict[str, int] = dict(capture_health) + last_health_permissions: list[str] = [] + last_health_evidence_at = 0.0 _sanitize_existing_jsonl() outbox = EventOutbox(LOCAL_DIR / "collector_outbox.db") current_state: dict | None = None current_key = None + current_doc = "" current_started_wall = "" current_started_mono = 0.0 current_screenshot: str | None = None + # A same-app title change becomes a boundary only after it persists. + pending_doc = "" + pending_count = 0 + pending_wall = "" + pending_mono = 0.0 + pending_state: dict | None = None + # Away: no input for away_after_seconds, or the screen is locked. + away = False + away_started_wall = "" + away_started_mono = 0.0 + away_reason = "" + away_idle_source = "" last_heartbeat = 0.0 last_poll_mono = 0.0 last_poll_wall_epoch = 0.0 + run_started_mono = time.monotonic() signal.signal(signal.SIGINT, _stop) signal.signal(signal.SIGTERM, _stop) @@ -534,17 +651,101 @@ def enqueue_clipboard(raw: RawClipboardAction) -> None: ) keyboard_started = keyboard_sensor.start() + # What each sensor needs on macOS. A started listener is not proof: macOS + # withholds events from an untrusted process while the thread still runs. + mouse_needs = ("accessibility",) + keyboard_needs = ("accessibility", "input_monitoring") + + def receiving(started: bool, needs: tuple[str, ...]) -> bool: + return bool(started and all(permissions.get(n) is not False for n in needs)) + + def keyboard_receiving() -> bool: + return receiving(keyboard_started, keyboard_needs) + + def sensor_line(label: str, started: bool, needs: tuple[str, ...], on_text: str) -> str: + if not started: + return f"{label}: OFF / unavailable" + denied = [n.replace("_", " ") for n in needs if permissions.get(n) is False] + if denied: + return f"{label}: BLOCKED (macOS {' and '.join(denied)} permission not granted; restart after granting)" + if permissions and any(permissions.get(n) is None for n in needs): + return f"{label}: {on_text} (permission not verified)" + return f"{label}: {on_text}" + print(f"Workflow Observer collector started. session={session_id}") print(f"device={cfg['device_id']} sensor={cfg['sensor_id']}") print(f"Durable delivery queue: {outbox.count()} pending event(s) at startup.") print("Polling detects focus/tab changes; long unchanged focus is checkpointed into durable spans rather than repeated polling rows.") print(f"Focus spans checkpoint every {max(30, float(cfg.get('focus_checkpoint_seconds', 120))):g}s so crashes cannot erase arbitrarily long work periods.") - print("Screen interaction capture: " + ("ON (clicks + throttled scrolls)" if sensor_started else "OFF / permission unavailable")) - print("Keyboard activity: " + ("ON (counts only; key identities/text are discarded)" if keyboard_started else "OFF / permission unavailable")) + print(sensor_line("Screen interaction capture", sensor_started, mouse_needs, "ON (clicks + throttled scrolls)")) + print(sensor_line("Keyboard activity", keyboard_started, keyboard_needs, "ON (counts only; key identities/text are discarded)")) clipboard_started = bool(keyboard_started and cfg.get("clipboard_behavior_enabled", True)) - print("Clipboard behavior: " + ("ON (copy/cut/paste actions only; contents never read)" if clipboard_started else "OFF")) + print(sensor_line("Clipboard behavior", clipboard_started, keyboard_needs, "ON (copy/cut/paste shortcuts only; contents never read)")) + if permissions: + print(sensor_line("Window titles and UI labels", True, ("accessibility",), "ON")) print("Typed text, ordinary key identities, and clipboard contents are never stored. Ctrl+C stops collection.") + def idle_seconds(now_mono: float) -> tuple[float | None, str]: + value = system_idle_seconds() + if value is not None: + return value, "os_input_clock" + # Fall back to OpenWorkGraph's own sensors, but only when they can + # actually see input; otherwise never guess that the user is away. + if not (keyboard_receiving() or receiving(sensor_started, mouse_needs)): + return None, "unavailable" + last = activity_tracker.last_input_mono() + return max(0.0, now_mono - max(last or 0.0, run_started_mono)), "collector_input_sensors" + + def emit_focus(end_mono: float, reason: str) -> None: + activity = activity_tracker.summarize( + current_started_mono, end_mono, + active_window_seconds=float(cfg.get("activity_active_window_seconds", 5)), + engaged_grace_seconds=float(cfg.get("engaged_grace_seconds", 60)), + ) + persist_event(_focus_span_event( + state=current_state, + started_at=current_started_wall, + duration_seconds=end_mono - current_started_mono, + cfg=cfg, + session_id=session_id, + screenshot_path=current_screenshot, + activity=activity, + boundary_reason=reason, + ), outbox) + + def emit_away(end_mono: float, reason: str) -> None: + persist_event(_away_span_event( + started_at=away_started_wall, + duration_seconds=end_mono - away_started_mono, + cfg=cfg, + session_id=session_id, + reason=away_reason, + idle_source=away_idle_source, + boundary_reason=reason, + ), outbox) + + def clear_pending() -> None: + nonlocal pending_doc, pending_count, pending_wall, pending_mono, pending_state + pending_doc, pending_count, pending_wall, pending_mono, pending_state = "", 0, "", 0.0, None + + def open_span(state: dict, key, wall: str, mono: float) -> None: + nonlocal current_state, current_key, current_doc, current_started_wall, current_started_mono, current_screenshot + current_state = state + current_key = key + current_doc = state.get("document_key", "") if _document_tracking(state, cfg) else "" + current_started_wall = wall + current_started_mono = mono + current_screenshot = None + if cfg.get("screenshots_enabled") and not state["excluded"]: + current_screenshot = screenshot(str(uuid.uuid4())) + clear_pending() + + def close_span() -> None: + nonlocal current_state, current_key, current_doc, current_started_wall, current_started_mono, current_screenshot + current_state, current_key, current_doc = None, None, "" + current_started_wall, current_started_mono, current_screenshot = "", 0.0, None + clear_pending() + while not STOP: now_mono = time.monotonic() now_dt = datetime.now(timezone.utc) @@ -559,25 +760,12 @@ def enqueue_clipboard(raw: RawClipboardAction) -> None: mono_gap = max(0.0, now_mono - last_poll_mono) wall_gap = max(0.0, now_wall_epoch - last_poll_wall_epoch) if max(mono_gap, wall_gap) > gap_threshold: - if current_state is not None: - observed_end_mono = min(now_mono, last_poll_mono + expected_poll) - if observed_end_mono > current_started_mono: - activity = activity_tracker.summarize( - current_started_mono, - observed_end_mono, - active_window_seconds=float(cfg.get("activity_active_window_seconds", 5)), - engaged_grace_seconds=float(cfg.get("engaged_grace_seconds", 60)), - ) - persist_event(_focus_span_event( - state=current_state, - started_at=current_started_wall, - duration_seconds=observed_end_mono - current_started_mono, - cfg=cfg, - session_id=session_id, - screenshot_path=current_screenshot, - activity=activity, - boundary_reason="capture_gap", - ), outbox) + observed_end_mono = min(now_mono, last_poll_mono + expected_poll) + if current_state is not None and observed_end_mono > current_started_mono: + emit_focus(observed_end_mono, "capture_gap") + if away and observed_end_mono > away_started_mono: + emit_away(observed_end_mono, "capture_gap") + away = False gap_duration = max(0.0, wall_gap - expected_poll) if gap_duration > 0: gap_start = datetime.fromtimestamp( @@ -590,74 +778,76 @@ def enqueue_clipboard(raw: RawClipboardAction) -> None: cfg=cfg, session_id=session_id, ), outbox) - capture_health["capture_gap_count"] += 1 - current_state = None - current_key = None - current_started_wall = "" - current_started_mono = 0.0 - current_screenshot = None + diagnostics["capture_gap_count"] += 1 + close_span() + + checkpoint_seconds = max(30.0, float(cfg.get("focus_checkpoint_seconds", 120))) + + # ---- away / back --------------------------------------------------------------- + if cfg.get("away_detection_enabled", True): + away_after = max(60.0, float(cfg.get("away_after_seconds", 300))) + locked = screen_locked() + idle, idle_source = idle_seconds(now_mono) + away_now = bool(locked) or (idle is not None and idle >= away_after) + if away_now and not away: + if current_state is not None: + emit_focus(now_mono, "away") + close_span() + away = True + away_started_wall, away_started_mono = now_wall, now_mono + away_reason = "screen_locked" if locked else "no_input" + away_idle_source = idle_source + diagnostics["away_count"] += 1 + elif away and not away_now: + emit_away(now_mono, "input_resumed") + away = False + elif away and now_mono - away_started_mono >= checkpoint_seconds: + emit_away(now_mono, "periodic_checkpoint") + away_started_wall, away_started_mono = now_wall, now_mono w = active_window() + if (w.app or "Unknown") == "Unknown": + capture_health["active_window_unavailable"] += 1 state = _public_window(w, cfg) key = _change_key(state, cfg) - if current_state is None: - current_state = state - current_key = key - current_started_wall = now_wall - current_started_mono = now_mono - if cfg.get("screenshots_enabled") and not state["excluded"]: - current_screenshot = screenshot(str(uuid.uuid4())) - + if away: + pass + elif current_state is None: + open_span(state, key, now_wall, now_mono) elif key != current_key: - activity = activity_tracker.summarize( - current_started_mono, now_mono, - active_window_seconds=float(cfg.get("activity_active_window_seconds", 5)), - engaged_grace_seconds=float(cfg.get("engaged_grace_seconds", 60)), - ) - persist_event(_focus_span_event( - state=current_state, - started_at=current_started_wall, - duration_seconds=now_mono - current_started_mono, - cfg=cfg, - session_id=session_id, - screenshot_path=current_screenshot, - activity=activity, - boundary_reason="focus_change", - ), outbox) - - current_state = state - current_key = key - current_started_wall = now_wall - current_started_mono = now_mono - current_screenshot = None - if cfg.get("screenshots_enabled") and not state["excluded"]: - current_screenshot = screenshot(str(uuid.uuid4())) - + emit_focus(now_mono, "focus_change") + open_span(state, key, now_wall, now_mono) else: - checkpoint_seconds = max(30.0, float(cfg.get("focus_checkpoint_seconds", 120))) - if now_mono - current_started_mono >= checkpoint_seconds: - activity = activity_tracker.summarize( - current_started_mono, now_mono, - active_window_seconds=float(cfg.get("activity_active_window_seconds", 5)), - engaged_grace_seconds=float(cfg.get("engaged_grace_seconds", 60)), - ) - persist_event(_focus_span_event( - state=current_state, - started_at=current_started_wall, - duration_seconds=now_mono - current_started_mono, - cfg=cfg, - session_id=session_id, - screenshot_path=current_screenshot, - activity=activity, - boundary_reason="periodic_checkpoint", - ), outbox) - capture_health["focus_checkpoint_count"] += 1 - current_started_wall = now_wall - current_started_mono = now_mono - current_screenshot = None - if cfg.get("screenshots_enabled") and not state["excluded"]: - current_screenshot = screenshot(str(uuid.uuid4())) + candidate = state.get("document_key", "") if _document_tracking(state, cfg) else "" + if _document_tracking(state, cfg) and is_material_change(current_doc, candidate): + if candidate == pending_doc: + pending_count += 1 + else: + pending_doc, pending_count = candidate, 1 + pending_wall, pending_mono, pending_state = now_wall, now_mono, state + needed = max(1, int(cfg.get("document_debounce_polls", 2))) + if pending_count >= needed and pending_state is not None: + # The new document began when it first appeared, not when confirmed. + boundary_mono = max(pending_mono, current_started_mono) + boundary_wall = pending_wall if pending_mono >= current_started_mono else current_started_wall + first_state = pending_state + emit_focus(boundary_mono, "document_change") + diagnostics["document_boundary_count"] += 1 + open_span(first_state, key, boundary_wall, boundary_mono) + else: + clear_pending() + + if current_state is not None and now_mono - current_started_mono >= checkpoint_seconds: + emit_focus(now_mono, "periodic_checkpoint") + diagnostics["focus_checkpoint_count"] += 1 + pending = (pending_doc, pending_count, pending_wall, pending_mono, pending_state) + open_span(current_state, current_key, now_wall, now_mono) + pending_doc, pending_count, pending_wall, pending_mono, pending_state = pending + + if now_mono - last_permission_check >= 60: + permissions = sensor_permissions() + last_permission_check = now_mono if now_mono - last_heartbeat >= float(cfg.get("heartbeat_seconds", 5)): heartbeat_activity = activity_tracker.summarize( @@ -668,6 +858,9 @@ def enqueue_clipboard(raw: RawClipboardAction) -> None: # Nest inside the established `activity` field so older local APIs keep # accepting the heartbeat while new dashboards can inspect health. heartbeat_activity["capture_health"] = dict(capture_health) + heartbeat_activity["diagnostics"] = dict(diagnostics) + heartbeat_activity["permissions"] = dict(permissions) + heartbeat_activity["away"] = away post_heartbeat( { "device_id": cfg["device_id"], @@ -676,24 +869,34 @@ def enqueue_clipboard(raw: RawClipboardAction) -> None: "actor_id": cfg.get("actor_id", ""), "session_id": session_id, "observed_at": now_wall, - "app": state["app"], - "window_title": state["window_title"], - "focus_elapsed_seconds": round(max(0.0, now_mono - current_started_mono), 3), + "app": "Away" if away else state["app"], + "window_title": "" if away else state["window_title"], + "focus_elapsed_seconds": round(max(0.0, now_mono - (away_started_mono if away else current_started_mono)), 3), "activity": heartbeat_activity, - "keyboard_sensor": keyboard_started, + "keyboard_sensor": keyboard_receiving(), "outbox_pending": outbox.count(), }, cfg["backend_url"], ) - # Persist sparse quality evidence only when counters have changed and - # at least one known degradation/diagnostic counter is nonzero. - if capture_health != last_capture_health_evidence and any(capture_health.values()): + # Persist sparse quality evidence only for real degradation: a new kind + # of problem or a permission change right away, repeats at most every + # HEALTH_REPEAT_SECONDS. Routine counters never trigger it. + missing_now = missing_permissions(permissions) + new_kind = any(capture_health[k] and not last_health_evidence.get(k) for k in DEGRADATION_KEYS) + grew = capture_health != last_health_evidence + if missing_now != last_health_permissions or new_kind or ( + grew and now_mono - last_health_evidence_at >= HEALTH_REPEAT_SECONDS + ): persist_event(_capture_health_event( capture_health=capture_health, cfg=cfg, session_id=session_id, + diagnostics=diagnostics, + permissions=permissions, ), outbox) - last_capture_health_evidence = dict(capture_health) + last_health_evidence = dict(capture_health) + last_health_permissions = missing_now + last_health_evidence_at = now_mono last_heartbeat = now_mono last_poll_mono = now_mono @@ -710,29 +913,19 @@ def enqueue_clipboard(raw: RawClipboardAction) -> None: except Exception: capture_health["interaction_worker_errors"] += 1 + end_mono = time.monotonic() if current_state is not None: - end_mono = time.monotonic() - activity = activity_tracker.summarize( - current_started_mono, end_mono, - active_window_seconds=float(cfg.get("activity_active_window_seconds", 5)), - engaged_grace_seconds=float(cfg.get("engaged_grace_seconds", 60)), - ) - persist_event(_focus_span_event( - state=current_state, - started_at=current_started_wall, - duration_seconds=end_mono - current_started_mono, - cfg=cfg, - session_id=session_id, - screenshot_path=current_screenshot, - activity=activity, - boundary_reason="shutdown", - ), outbox) + emit_focus(end_mono, "shutdown") + if away: + emit_away(end_mono, "shutdown") - if capture_health != last_capture_health_evidence and any(capture_health.values()): + if capture_health != last_health_evidence: persist_event(_capture_health_event( capture_health=capture_health, cfg=cfg, session_id=session_id, + diagnostics=diagnostics, + permissions=permissions, ), outbox) for _ in range(5): @@ -749,8 +942,7 @@ def main() -> None: parser = argparse.ArgumentParser(description="Local-first workflow observation collector") parser.add_argument("--config", default=str(ROOT / "config.json")) args = parser.parse_args() - run(Path(args.config)) - + raise SystemExit(run(Path(args.config))) if __name__ == "__main__": main() diff --git a/collector/permissions.py b/collector/permissions.py new file mode 100644 index 0000000..2e50a80 --- /dev/null +++ b/collector/permissions.py @@ -0,0 +1,65 @@ +from __future__ import annotations + +"""Truthful OS permission state for the desktop sensors. + +A pynput listener thread starts even when macOS withholds input events, so +"listener started" is not evidence that OpenWorkGraph can see anything. These +preflight checks ask the OS directly and never trigger a permission prompt. + +Values: ``True`` granted, ``False`` denied, ``None`` unknown/not applicable. +""" + +import ctypes +import ctypes.util +import platform + +# IOKit: kIOHIDRequestTypeListenEvent and IOHIDAccessType values. +_LISTEN_EVENT = 1 +_GRANTED, _DENIED = 0, 1 + + +def _macos_accessibility() -> bool | None: + try: + import ApplicationServices # type: ignore + + return bool(ApplicationServices.AXIsProcessTrusted()) + except Exception: + return None + + +def _macos_input_monitoring() -> bool | None: + try: + path = ctypes.util.find_library("IOKit") + if not path: + return None + iokit = ctypes.cdll.LoadLibrary(path) + check = iokit.IOHIDCheckAccess # macOS 10.15+ + check.argtypes = [ctypes.c_uint32] + check.restype = ctypes.c_uint32 + value = int(check(_LISTEN_EVENT)) + except Exception: + return None + if value == _GRANTED: + return True + if value == _DENIED: + return False + return None # not determined yet + + +def sensor_permissions() -> dict[str, bool | None]: + """Permission state that decides what the desktop sensors can actually see. + + * accessibility: window titles and UI element labels on click (and, on + older macOS versions, input events). + * input_monitoring: keyboard activity counts and copy/cut/paste shortcuts. + """ + if platform.system() != "Darwin": + return {} + return { + "accessibility": _macos_accessibility(), + "input_monitoring": _macos_input_monitoring(), + } + + +def missing(permissions: dict[str, bool | None]) -> list[str]: + return sorted(name for name, granted in permissions.items() if granted is False) diff --git a/collector/platform.py b/collector/platform.py index 1df6d80..7f95bc5 100644 --- a/collector/platform.py +++ b/collector/platform.py @@ -70,3 +70,63 @@ def active_window() -> ActiveWindow: if system == "Windows": return _windows_active_window() return _linux_active_window() + + +def _macos_idle_seconds() -> float | None: + try: + import Quartz # type: ignore + + value = Quartz.CGEventSourceSecondsSinceLastEventType( + Quartz.kCGEventSourceStateCombinedSessionState, Quartz.kCGAnyInputEventType + ) + return max(0.0, float(value)) + except Exception: + return None + + +def _windows_idle_seconds() -> float | None: + try: + class LASTINPUTINFO(ctypes.Structure): + _fields_ = [("cbSize", ctypes.c_uint), ("dwTime", ctypes.c_uint)] + + info = LASTINPUTINFO() + info.cbSize = ctypes.sizeof(LASTINPUTINFO) + if not ctypes.windll.user32.GetLastInputInfo(ctypes.byref(info)): + return None + now = ctypes.windll.kernel32.GetTickCount() & 0xFFFFFFFF + return max(0.0, ((now - info.dwTime) & 0xFFFFFFFF) / 1000.0) + except Exception: + return None + + +def system_idle_seconds() -> float | None: + """Seconds since the last keyboard/mouse input anywhere in the session. + + Read from the OS input clock, so it needs no Input Monitoring permission and + never sees which keys were pressed. ``None`` when the platform cannot say + (Linux for now); callers must then not guess that the user is away. + """ + system = _platform.system() + if system == "Darwin": + return _macos_idle_seconds() + if system == "Windows": + return _windows_idle_seconds() + return None + + +def screen_locked() -> bool | None: + """True while the session is locked or switched away from the console.""" + if _platform.system() != "Darwin": + return None + try: + import Quartz # type: ignore + + session = Quartz.CGSessionCopyCurrentDictionary() + if not session: + return None + if bool(session.get("CGSSessionScreenIsLocked", False)): + return True + on_console = session.get("kCGSSessionOnConsoleKey") + return False if on_console is None else not bool(on_console) + except Exception: + return None diff --git a/collector/secure_main.py b/collector/secure_main.py index cdb617b..fa4631f 100644 --- a/collector/secure_main.py +++ b/collector/secure_main.py @@ -10,6 +10,12 @@ from shared.capture_control import initialize_run, read_state +from .instance_lock import EXIT_ALREADY_RUNNING + +# Another collector holds this data directory's lock (for example one left over +# from an earlier launch). Check again now and then instead of every second. +ALREADY_RUNNING_RETRY_SECONDS = 15.0 + _STOP = False @@ -42,6 +48,7 @@ def run_supervisor(config_path: Path) -> int: last_state = "" restart_not_before = 0.0 cooperative_stop_started = 0.0 + reported_other_collector = False try: while not _STOP: @@ -51,6 +58,12 @@ def run_supervisor(config_path: Path) -> int: if state == "recording": cooperative_stop_started = 0.0 + if worker is not None and worker.poll() == EXIT_ALREADY_RUNNING: + if not reported_other_collector: + print("OpenWorkGraph: another collector is already recording this data folder; not starting a second one.") + reported_other_collector = True + restart_not_before = max(restart_not_before, time.monotonic() + ALREADY_RUNNING_RETRY_SECONDS) + worker = None if worker is None or worker.poll() is not None: now = time.monotonic() if now >= restart_not_before: diff --git a/collector/secure_worker.py b/collector/secure_worker.py index eef7dc8..51671d8 100644 --- a/collector/secure_worker.py +++ b/collector/secure_worker.py @@ -92,10 +92,11 @@ def main() -> None: monitor = threading.Thread(target=_watch_capture_state, args=(generation, done), daemon=True) monitor.start() try: - base.run(Path(args.config).resolve()) + code = base.run(Path(args.config).resolve()) finally: done.set() monitor.join(timeout=1) + raise SystemExit(code or 0) if __name__ == "__main__": diff --git a/config.example.json b/config.example.json index 2116ef6..852738a 100644 --- a/config.example.json +++ b/config.example.json @@ -15,7 +15,10 @@ "heartbeat_seconds": 5, "focus_checkpoint_seconds": 120, "capture_gap_seconds": 30, - "change_detection": "application", + "change_detection": "application_and_document", + "document_debounce_polls": 2, + "away_detection_enabled": true, + "away_after_seconds": 300, "screenshot_interval_seconds": 20, "screenshots_enabled": false, "upload_screenshots": false, diff --git a/connector/policy.py b/connector/policy.py index c920c00..36fce8f 100644 --- a/connector/policy.py +++ b/connector/policy.py @@ -93,10 +93,18 @@ def _strip_keys(value: Any, blocked: set[str]) -> Any: return value +# Evidence that stays on the employee's computer whatever the organization policy +# says. Away spans record when someone stepped away; sharing them would turn +# OpenWorkGraph into presence monitoring. +LOCAL_ONLY_EVENT_TYPES = frozenset({"away_span"}) + + def prepare_event_for_gateway(event: dict[str, Any], policy: dict[str, Any]) -> dict[str, Any] | None: metadata = event.get("metadata") if not isinstance(metadata, dict): metadata = {} + if str(event.get("event_type") or "") in LOCAL_ONLY_EVENT_TYPES: + return None if str(event.get("source") or "") == "agent" and policy.get("allow_agent_events") is not True: return None if bool(metadata.get("excluded")) and not policy.get("share_excluded", False): diff --git a/dashboard/connections.js b/dashboard/connections.js index ac96da6..93fc3df 100644 --- a/dashboard/connections.js +++ b/dashboard/connections.js @@ -4,6 +4,7 @@ // /v1/connections, the same code as the owg_connect.py CLI agents can use. let state=null; let aiAccess=null; + let diagnostics=null; let busy=new Set(); let lastMessage={}; @@ -103,6 +104,29 @@ return ''; } + // Where each native telemetry signal stands since OpenWorkGraph started: + // counts and reason codes only (see /v1/agent-telemetry/diagnostics). + const CHANNELS={claude_code:[['claude_code_hooks','Hooks'],['claude_code_otel_logs','OTel logs']],codex:[['codex_otel','OTel']]}; + function channelText(label,ch){ + if(!ch||!ch.requests)return {text:`${label}: nothing received since OpenWorkGraph started`,warn:false}; + const parts=[`${label} ${relativeTime(ch.last_received_at)}`]; + if(ch.events_stored)parts.push(`${ch.events_stored} stored`); + else if(ch.processed&&ch.records_seen&&!ch.events_projected)parts.push('received, but nothing structural in it'); + const rejected=Object.entries(ch.rejected||{}); + if(rejected.length)parts.push(`${rejected.reduce((n,[,c])=>n+c,0)} rejected (${rejected.map(([r])=>r.replace(/_/g,' ')).join(', ')})`); + if(ch.observation_off)parts.push('arriving while Observe is off'); + return {text:parts.join(' · '),warn:rejected.length>0}; + } + function diagnosticsLine(client){ + const channels=CHANNELS[client.id]; + if(!channels||!diagnostics||!client.observe?.supported||!client.observe.on)return ''; + const lines=channels.map(([key,label])=>channelText(label,diagnostics.channels?.[key])); + const missing=diagnostics.configuration?.[client.id]?.missing_hook_events||[]; + if(missing.length)lines.push({text:`Hooks out of date (missing ${missing.join(', ')}); restart OpenWorkGraph to update them`,warn:true}); + const warn=lines.some(l=>l.warn); + return `
${esc(lines.map(l=>l.text).join(' · '))}
`; + } + function subline(client){ const restart=restartLine(client);if(restart)return restart; const msg=lastMessage[client.id]; @@ -110,12 +134,13 @@ const errors=[client.mcp,client.observe].map(x=>x&&x.error).filter(Boolean); if(errors.length)return `
${esc(errors[0])}
`; if(!client.detected)return '
Not found on this computer
'; + const diag=diagnosticsLine(client); // Evidence, not configuration: these only appear once data actually flowed. const parts=[]; if(client.last_used)parts.push(`Context used ${relativeTime(client.last_used)}`); if(client.last_observed)parts.push(`Telemetry observed ${relativeTime(client.last_observed)}`); else if(client.observe.supported&&client.observe.on)parts.push('No telemetry observed yet'); - return parts.length?`
${esc(parts.join(' · '))}
`:''; + return (parts.length?`
${esc(parts.join(' · '))}
`:'')+diag; } function render(){ @@ -151,10 +176,12 @@ async function refresh(){ try{ - const [connections,access,activity,traces]=await Promise.all([ + const [connections,access,activity,traces,diag]=await Promise.all([ api('/v1/connections'),api('/v1/ai-access'),api('/v1/mcp-activity?limit=200').catch(()=>({items:[]})), api('/v1/agent-execution-traces?limit=50&evidence_limit=25000&max_events_per_execution=1').catch(()=>({executions:[]})), + api('/v1/agent-telemetry/diagnostics').catch(()=>null), ]); + diagnostics=diag; const lastUsed={},lastObserved={}; for(const item of activity.items||[]){if(item.client&&item.status==='ok'&&!lastUsed[item.client])lastUsed[item.client]=item.observed_at;} for(const run of traces.executions||[]){ diff --git a/dashboard/gateway_panel.js b/dashboard/gateway_panel.js index d355779..674f495 100644 --- a/dashboard/gateway_panel.js +++ b/dashboard/gateway_panel.js @@ -120,7 +120,11 @@ } const state=String(s.state||'recording'); if (state==='recording') { - label.textContent=`Recording · ${fmtDuration(s.run_elapsed_seconds)}`; dot.style.background='#d25a5a'; + const names={accessibility:'Accessibility',input_monitoring:'Input Monitoring'}; + const missing=(s.missing_permissions||[]).map(p=>names[p]||p); + label.textContent=`Recording · ${fmtDuration(s.run_elapsed_seconds)}`+(s.away?' · away':'')+(missing.length?` · macOS ${missing.join(' and ')} permission missing`:''); + label.title=missing.length?`OpenWorkGraph cannot see everything: allow it under System Settings → Privacy & Security → ${missing.join(' / ')}, then restart OpenWorkGraph.`:''; + dot.style.background='#d25a5a'; buttons.innerHTML=''; buttons.querySelector('#pauseCapture').onclick=()=>captureAction('pause'); buttons.querySelector('#stopCapture').onclick=()=>captureAction('stop'); diff --git a/docs/CHANGELOG_V098.md b/docs/CHANGELOG_V098.md new file mode 100644 index 0000000..173cb9b --- /dev/null +++ b/docs/CHANGELOG_V098.md @@ -0,0 +1,73 @@ +# OpenWorkGraph (unreleased, planned v0.98): capture correctness + +Measured on a real install, OpenWorkGraph's evidence had four problems: +- Idle time was recorded as focus time. 198 h of foreground time had 12.3 h with any input, and single spans ran overnight. +- Two collectors sometimes recorded at once. 8% of focus time was duplicated, and some days showed more than 24 h. +- About 12% of all events were "capture health" warnings caused by routine checkpoints. +- Sensors reported "ON" on macOS while the OS was withholding their input. + +This release fixes those without changing what a field means, what is shared, or the event schema. + +## Time you were away is no longer focus time +- **When a span ends:** a focus span ends after 5 minutes (`away_after_seconds`, minimum 60) with no keyboard or mouse input anywhere, or when the screen locks. The time until input returns becomes an **away span** (`event_type: "away_span"`, app "Away"). +- **Reading time kept:** the first 5 minutes stay with the app, which covers reading and thinking without input. +- **Meanings unchanged:** `duration_seconds` still means wall-clock observed time. Effort is still `activity.engaged_seconds` / `idle_seconds`. +- **Where idle comes from:** the operating system's input clock on macOS and Windows. It needs no Input Monitoring permission and never sees keys. OpenWorkGraph's own sensors are a fallback, used only when they can actually see input. +- **No idle signal:** without one (Linux today) there is no away detection, instead of guessed absences. +- **Stays on the computer:** away spans are never sent to an organization Gateway, whatever its policy. Sharing them would make OpenWorkGraph presence monitoring. + +## One collector per data folder +- **The lock:** the collector holds an OS advisory lock on its data folder for its lifetime. The OS releases it on exit or crash, so there is no stale PID file. +- **A second collector:** another one for the same folder (an orphan from an earlier launch, or a direct launch) exits with code 75 and records nothing. The supervisor checks again every 15 s instead of respawning every second. +- **Separate folders:** `data/live` and `data/demo` are separate. + +## Document changes within one app +- **The new default:** `change_detection: "application_and_document"`. When the window title changes materially within the same app and stays changed for 2 polls, the span ends. The new span starts when the new document first appeared. +- **What counts as noise:** unread badges, unsaved markers, "Edited"/"Redigerad", progress percentages and "Not Responding" are ignored, as are one-poll dialogs and titles going empty. +- **Title privacy modes:** + - `none`: nothing about titles is used. + - `hash`: normalization happens in memory before hashing, and nothing new is stored. + - Excluded apps: never tracked. +- **Existing configs:** `"application"` was the shipped default copied into every `config.json`, so it now follows the default. `"application_only"` keeps the old behaviour. `"application_and_title"` is unchanged. + +## Capture health means degradation again +- **What triggers a warning:** only real degradation produces `capture_health` evidence, namely dropped interactions, worker errors, an unreadable foreground window, or missing permissions. A new kind of problem or a permission change is reported immediately; repeats at most every 10 minutes. +- **Routine counters:** checkpoints, capture gaps, document boundaries and away spans go to `diagnostics` in the heartbeat and event instead. + +## macOS permissions are stated, not assumed +- **How it checks:** the collector asks macOS directly with `AXIsProcessTrusted` and `IOHIDCheckAccess`, which never prompt, and rechecks every minute. +- **What it reports:** it names each sensor that is BLOCKED, and why. Clicks and titles need Accessibility; keyboard counts and copy/paste shortcuts need Accessibility and Input Monitoring. +- **Heartbeat:** `keyboard_sensor` is false when keys cannot arrive, and the permissions are included. +- **Dashboard:** the Recording pill shows "macOS … permission missing". Missing permissions are recorded as capture health. + +## Claude Code turns +- **New hooks:** `UserPromptSubmit` starts a turn and `Stop` finishes it successfully, keyed by `prompt_id`, which Claude's tool hooks and OTel share. Sessions keep their own `SessionStart`/`SessionEnd` boundary. +- **Missing `prompt_id`:** those turn hooks are ignored rather than mistaken for the session. +- **Content never read:** these hooks carry the prompt and the last reply, and the adapter reads neither. +- **Existing installs:** OpenWorkGraph adds missing hook events at startup when Claude Code Observe is on, with a backup. Your own hooks and env are kept, and hooks are never installed where Observe was not turned on. + +## Agent telemetry diagnostics +- **Where:** `GET /v1/agent-telemetry/diagnostics`, and a line in each Connections row. +- **What it shows:** per channel (Claude hooks, Claude OTel logs, Codex OTel, other agent events, generic OTel, spool), since OpenWorkGraph started: + - when requests arrived; + - rejections by reason; + - OTel records seen and ignored; + - events stored. + + It also shows which exporters and hook events are configured. +- **No content:** counts, timestamps and reason codes only. + +## Delayed delivery for agent hooks, only while recording +- **When it spools:** if a hook cannot reach OpenWorkGraph, the event may wait in a local spool (`data/auth/agent_spool`) under a 90-second recording lease. Only a running, recording, non-demo OpenWorkGraph grants the lease, and it is revoked on Pause, Stop and exit. +- **After Stop or quit:** agent events are dropped, not collected for later. +- **On delivery:** spooled events go through the normal ingest rules again (Observe switches, deletions, retention, pause windows). +- **Limits:** 2,000 files of 256 KB, 24 h maximum age. A rejection is never spooled. + +## Not in this release +- Clipboard-write detection. +- Copilot, Gemini CLI and Cursor presets. +- Multilingual web-agent detection. +- Claude metrics export. +- `PreCompact` (it needs a schema operation that older Gateways would reject). +- Platform rework: macOS without `osascript`, Windows UWP app names, Linux/Wayland. +- **Known issue, unchanged:** Claude subagent runs share their parent turn's `run_id`, so they are grouped into the parent turn rather than linked as children. diff --git a/docs/NATIVE_AGENT_ADAPTERS.md b/docs/NATIVE_AGENT_ADAPTERS.md index 2761424..6cdcf6b 100644 --- a/docs/NATIVE_AGENT_ADAPTERS.md +++ b/docs/NATIVE_AGENT_ADAPTERS.md @@ -29,16 +29,18 @@ Missing evidence is never treated as proof that an underlying agent action did n OpenWorkGraph combines two Claude Code surfaces: -1. **Lifecycle hooks** for session boundaries, completed tools, permission requests, failures, and subagent start/stop. +1. **Lifecycle hooks** for session boundaries, turn boundaries (one prompt = one turn), completed tools, permission requests, failures, and subagent start/stop. 2. **Claude Code OpenTelemetry log events** for per-prompt model calls, token usage, model identity, completed tools, permission decisions, retries and subagent completion. The hook surface remains useful as a fail-open lifecycle/fallback channel. The OTel log surface supplies the signals hooks do not expose, especially model calls and token usage. -OpenWorkGraph registers only content-safe hook points: +OpenWorkGraph registers these hook points: ```text SessionStart SessionEnd +UserPromptSubmit +Stop PostToolUse PostToolUseFailure PermissionRequest @@ -48,7 +50,18 @@ SubagentStop StopFailure ``` -Content-heavy hooks such as `UserPromptSubmit`, `PreToolUse`, `MessageDisplay`, and `Stop` are intentionally not registered. +`UserPromptSubmit` and `Stop` carry content: the prompt text and the last assistant message. OpenWorkGraph uses them only as turn boundaries. The adapter reads just `session_id`, `prompt_id` and the event name, and never copies a content field. + +| Hook | Becomes | +|---|---| +| `SessionStart` / `SessionEnd` | the session starts / finishes (`run_id` = session) | +| `UserPromptSubmit` / `Stop` | a turn starts / finishes successfully (`run_id` = `prompt_id`, which Claude's tool hooks and OTel share) | + +- **No `prompt_id`:** a turn hook without one is ignored. Falling back to the session would make a turn's `Stop` look like the whole session finishing. +- **Not registered:** `PreToolUse` and `MessageDisplay`. +- **`PreCompact`:** not mapped yet. It needs a new operation in the agent event schema, and older Gateways would reject that when syncing. + +**Existing installs:** when OpenWorkGraph starts with Claude Code Observe on, it adds any hook events its entries are missing. It keeps a backup, keeps your own hooks and env, and never installs hooks you have not turned on. The Claude OTel adapter recognizes these stable structural events: @@ -96,9 +109,34 @@ The dedicated local endpoint is: POST /agent-ingest/v1/claude-otel ``` +**Diagnostics:** `GET /v1/agent-telemetry/diagnostics` (and the Connections row in the dashboard) reports per channel since OpenWorkGraph started: + +- requests received and when; +- rejections by reason (`auth`, `invalid_payload`, `too_large`, `not_an_object`, `adapter_error`); +- OTel records seen versus ignored; +- events stored. + +It also shows which exporters and hook events are configured. So "no model calls" can be traced to one of: + +- Claude not exporting; +- a rejected request; +- records the adapter does not recognise. + +It holds counts, timestamps and reason codes only, never payload content. + The integration uses Claude's logs exporter only. OWG v0.90 does **not** require Claude's beta detailed-trace exporter. -### Hook failure behavior +### Hook failure behavior and delayed delivery + +If a hook cannot reach OpenWorkGraph (not running, restarting, busy past the 0.75 s timeout, or a server error), the event may be written to a short local spool and delivered later. This only happens under a recording lease: + +- **Who grants it:** a running, recording, non-demo OpenWorkGraph issues the lease. It lasts 90 s and is renewed while recording. +- **When it ends:** it is revoked immediately on Pause, Stop and exit. +- **What the hook checks:** the hook also reads the capture state of the data folder that issued the lease, and spools only if it says "recording". +- **What happens otherwise:** after Stop, or when OpenWorkGraph is not running, agent events are dropped, never collected for later. +- **What doesn't qualify:** a rejection (4xx) is never spooled. +- **On delivery:** spooled events go through the normal ingest path again, so Observe switches, deletions, retention and pause windows apply. +- **Limits:** at most 2,000 files of 256 KB each, and anything older than 24 h is discarded. Files live under `data/auth/agent_spool`, which only your user can read. Hooks remain asynchronous and fail-open. Invalid JSON, an unavailable OpenWorkGraph server, authentication failure, or an adapter exception cannot block or approve Claude Code execution. `OWG_AGENT_ADAPTER_DEBUG=1` prints only a fixed diagnostic notice, never native exception/payload content. diff --git a/docs/PRIVACY_AND_DATA.md b/docs/PRIVACY_AND_DATA.md index 32ddb73..9463480 100644 --- a/docs/PRIVACY_AND_DATA.md +++ b/docs/PRIVACY_AND_DATA.md @@ -52,6 +52,7 @@ Current signals include: - active application/window focus spans - browser tab activation and navigation when the optional browser sensor is installed - foreground, engaged, probable-idle and active-input timing +- away spans: time with no keyboard or mouse input for 5 minutes (configurable), or with the screen locked. They are read from the operating system's idle clock, which reports only when the last input happened and never which key was pressed. Away spans stay on this computer and are never sent to an organization Gateway, whatever its policy. - keypress counts - global clicks and throttled scrolls - safe native control identity/label metadata through macOS Accessibility or Windows UI Automation, best effort diff --git a/server/agent_capture_runtime.py b/server/agent_capture_runtime.py new file mode 100644 index 0000000..6f3a4bb --- /dev/null +++ b/server/agent_capture_runtime.py @@ -0,0 +1,172 @@ +from __future__ import annotations + +"""Agent-capture housekeeping that runs inside the real OpenWorkGraph launcher. + +* Keeps the agent spool's recording lease current while capture is recording, + and revokes it the moment capture is paused or stopped, or OpenWorkGraph exits. +* Delivers spooled agent events through the normal ingest path. +* Refreshes OpenWorkGraph's own Claude Code hook entries when a newer version + needs more hook events (only when Claude Code Observe is on and OpenWorkGraph's + hooks are already installed; a backup is written first). + +Started from ``server.enterprise_runner.main`` only, never from app startup, so +tests and embedded apps never touch the user's Claude Code settings. +""" + +import atexit +import threading +import uuid +from pathlib import Path +from typing import Any + +from . import agent_spool +from . import agent_telemetry_diagnostics as diagnostics + +TICK_SECONDS = 2.0 +LEASE_RENEW_BEFORE_SECONDS = 60.0 +FLUSH_EVERY_SECONDS = 5.0 + +_LEASE_ID = uuid.uuid4().hex +_STATE: dict[str, Any] = {"lease_active": False, "lease_valid_until": None, "last_flush": None} +_STOP = threading.Event() +_THREAD: threading.Thread | None = None + + +def _data_dir() -> Path: + from shared.capture_control import data_dir + + return data_dir() + + +def _demo() -> bool: + try: + from .enterprise_app import _demo_mode + + return bool(_demo_mode()) + except Exception: + return False + + +def _ingest_spooled(events: list[dict[str, Any]]) -> None: + from .agent_ingest import ingest_agent_payloads + from .agent_routes import _observation_enabled_for + + diagnostics.received("spool") + allowed = [event for event in events if _observation_enabled_for(event)] + if not allowed: + diagnostics.observation_off("spool") + else: + diagnostics.processed("spool", ingest_agent_payloads(allowed)) + + +def tick(now: float | None = None) -> None: + """One lease/flush step; separated from the thread for tests.""" + import time as _time + + from shared.capture_control import read_state + + current = _time.time() if now is None else now + recording = not _demo() and read_state().get("state") == "recording" + if recording: + valid_until = _STATE.get("lease_valid_until") or 0.0 + if not _STATE["lease_active"] or valid_until - current < LEASE_RENEW_BEFORE_SECONDS: + lease = agent_spool.issue_lease(_data_dir(), lease_id=_LEASE_ID, now=current) + _STATE.update(lease_active=True, lease_valid_until=lease["valid_until"]) + elif _STATE["lease_active"]: + agent_spool.revoke_lease(lease_id=_LEASE_ID) + _STATE.update(lease_active=False, lease_valid_until=None) + + last = _STATE.get("last_flush") or 0.0 + if current - last >= FLUSH_EVERY_SECONDS and agent_spool.pending_count(): + _STATE["last_flush"] = current + if not _demo(): + agent_spool.flush_spool(_data_dir(), _ingest_spooled, now=current) + + +def _loop() -> None: + while not _STOP.is_set(): + try: + tick() + except Exception: + pass + _STOP.wait(TICK_SECONDS) + + +def stop() -> None: + _STOP.set() + agent_spool.revoke_lease(lease_id=_LEASE_ID) + _STATE.update(lease_active=False, lease_valid_until=None) + + +def start() -> None: + global _THREAD + if _THREAD is not None and _THREAD.is_alive(): + return + if not _demo(): + refresh_outdated_claude_hooks() + _STOP.clear() + _THREAD = threading.Thread(target=_loop, name="owg-agent-capture-runtime", daemon=True) + _THREAD.start() + atexit.register(stop) + + +def refresh_outdated_claude_hooks() -> dict[str, Any] | None: + from adapters.claude_code_hook import SUPPORTED_EVENTS, settings_fragment + + from . import agent_config_writer as writer + from .connections import is_enabled + + try: + if not is_enabled("claude_code", "observe"): + return None + missing = writer.claude_missing_hook_events(SUPPORTED_EVENTS) + if not missing: + return None + result = writer.claude_connect(settings_fragment) + print(f"OpenWorkGraph: added Claude Code hooks {', '.join(missing)} (backup: {result.get('backup')}).") + return {"added": missing, **result} + except writer.ConfigConflict as exc: + print(f"OpenWorkGraph: Claude Code hooks are out of date but were left unchanged: {exc}") + except Exception: + pass + return None + + +def configuration_state() -> dict[str, Any]: + """What is configured, as booleans and names only (no tokens or headers).""" + from adapters.claude_code_hook import SUPPORTED_EVENTS + + from . import agent_config_writer as writer + + claude: dict[str, Any] = {} + try: + data = writer._load_claude(writer.claude_settings_path()) + env = data.get("env") if isinstance(data.get("env"), dict) else {} + present = writer.claude_owg_hook_events(data) + claude = { + "hook_events": sorted(present), + "missing_hook_events": sorted(set(SUPPORTED_EVENTS) - present) if present else [], + "telemetry_enabled": str(env.get("CLAUDE_CODE_ENABLE_TELEMETRY") or "") == "1", + "otel_logs_exporter": str(env.get("OTEL_LOGS_EXPORTER") or "") or None, + "otel_metrics_exporter": str(env.get("OTEL_METRICS_EXPORTER") or "") or None, + "otel_logs_endpoint_is_openworkgraph": "/agent-ingest/v1/claude-otel" in str(env.get("OTEL_EXPORTER_OTLP_LOGS_ENDPOINT") or ""), + } + except Exception as exc: # unreadable settings are themselves a diagnosis + claude = {"error": type(exc).__name__} + try: + codex = {"configured": bool(writer.codex_status().get("configured"))} + except Exception as exc: + codex = {"error": type(exc).__name__} + return { + "claude_code": claude, + "codex": codex, + "agent_spool": { + "lease_active": bool(_STATE["lease_active"]), + "pending_files": agent_spool.pending_count(), + "limits": { + "lease_seconds": agent_spool.LEASE_TTL_SECONDS, + "max_files": agent_spool.MAX_SPOOL_FILES, + "max_age_seconds": agent_spool.MAX_AGE_SECONDS, + }, + }, + } diff --git a/server/agent_config_writer.py b/server/agent_config_writer.py index 85ddc69..62990b5 100644 --- a/server/agent_config_writer.py +++ b/server/agent_config_writer.py @@ -21,7 +21,7 @@ import time import tomllib from pathlib import Path -from typing import Any, Callable +from typing import Any, Callable, Iterable CLAUDE_HOOK_MARKER = "adapters.claude_code_hook" CODEX_BLOCK_START = "# >>> OpenWorkGraph agent observation (managed; remove via dashboard) >>>" @@ -158,6 +158,35 @@ def _claude_hooks_configured(data: dict[str, Any]) -> bool: ) +def claude_owg_hook_events(data: dict[str, Any]) -> set[str]: + """Hook events that currently have an OpenWorkGraph handler.""" + found: set[str] = set() + for event, groups in (data.get("hooks") or {}).items(): + if not isinstance(groups, list): + continue + for group in groups: + handlers = group.get("hooks") if isinstance(group, dict) else None + if isinstance(handlers, list) and any(_is_owg_handler(h) for h in handlers): + found.add(str(event)) + return found + + +def claude_missing_hook_events(required: Iterable[str]) -> list[str]: + """Required events missing from an existing OpenWorkGraph hook install. + + Empty when OpenWorkGraph's hooks are not installed at all: that is "off", not + "outdated", and must not be turned on behind the user's back. + """ + try: + data = _load_claude(claude_settings_path()) + except ConfigConflict: + return [] + present = claude_owg_hook_events(data) + if not present: + return [] + return sorted(set(required) - present) + + def _claude_env_matches(data: dict[str, Any], managed_env: dict[str, str]) -> bool: existing = data.get("env") if isinstance(data.get("env"), dict) else {} return all(str(existing.get(key) or "") == str(value) for key, value in managed_env.items()) diff --git a/server/agent_routes.py b/server/agent_routes.py index 31dd2c5..6cc08b0 100644 --- a/server/agent_routes.py +++ b/server/agent_routes.py @@ -7,6 +7,7 @@ from pydantic import BaseModel from shared.agent_evidence import AgentEvidenceError +from . import agent_telemetry_diagnostics as diagnostics from .agent_auth import agent_bearer_matches from .agent_ingest import ( MAX_AGENT_BATCH_BYTES, @@ -89,71 +90,125 @@ async def _read_bounded_json(request: Request) -> Any: raise HTTPException(status_code=400, detail="invalid agent JSON payload") from exc +def _events_channel(payload: Any) -> str: + events = payload.get("events") if isinstance(payload, dict) else None + first = events[0] if isinstance(events, list) and events and isinstance(events[0], dict) else {} + return "claude_code_hooks" if str(first.get("framework") or "") == "claude-code" else "agent_events" + + +async def _authorized_json(request: Request, channel: str | None) -> Any: + """Authenticate and read the body, recording the outcome for diagnostics.""" + if channel: + diagnostics.received(channel) + try: + _require_agent_write_bearer(request) + except HTTPException: + if channel: + diagnostics.rejected(channel, "auth") + raise + try: + return await _read_bounded_json(request) + except HTTPException as exc: + if channel: + diagnostics.rejected(channel, "too_large" if exc.status_code == 413 else "invalid_payload") + raise + + @router.post(AGENT_EVENT_PATH) async def ingest_agent_events(request: Request) -> dict[str, int | str]: - _require_agent_write_bearer(request) - payload = await _read_bounded_json(request) + # The native client names its channel so an auth failure is attributable + # before the body is read; otherwise classify by the events themselves. + hint = str(request.headers.get("x-owg-channel") or "") + channel = hint if hint in {"claude_code_hooks", "agent_events"} else "" + payload = await _authorized_json(request, channel or None) + if not channel: + channel = _events_channel(payload) + diagnostics.received(channel) if not isinstance(payload, dict) or not isinstance(payload.get("events"), list): + diagnostics.rejected(channel, "invalid_payload") raise HTTPException(status_code=422, detail="events must be a list") events = [event for event in payload["events"] if _observation_enabled_for(event)] if not events: + diagnostics.observation_off(channel) return {"received": 0, "status": "observation_off"} try: result = ingest_agent_payloads(events) except AgentEvidenceError as exc: + diagnostics.rejected(channel, "invalid_payload") raise HTTPException(status_code=422, detail=str(exc)) from exc + diagnostics.processed(channel, result) return {**result, "status": "ok"} @router.post(AGENT_OTEL_PATH) async def ingest_agent_otel(request: Request) -> dict[str, int | str]: - _require_agent_write_bearer(request) - payload = await _read_bounded_json(request) + channel = "otel_generic" + payload = await _authorized_json(request, channel) if not isinstance(payload, dict): + diagnostics.rejected(channel, "not_an_object") raise HTTPException(status_code=422, detail="OpenTelemetry payload must be an object") defaults_raw = payload.pop("openworkgraph", {}) try: defaults = OTelDefaults.model_validate(defaults_raw if isinstance(defaults_raw, dict) else {}).model_dump() result = ingest_otel_payload(payload, defaults=defaults) except (AgentEvidenceError, ValueError) as exc: + diagnostics.rejected(channel, "adapter_error") raise HTTPException(status_code=422, detail=str(exc)) from exc + diagnostics.processed(channel, result) return {**result, "status": "ok"} @router.post(AGENT_CODEX_OTEL_PATH, status_code=202) async def ingest_codex_otel(request: Request) -> Response: - _require_agent_write_bearer(request) - payload = await _read_bounded_json(request) + channel = "codex_otel" + payload = await _authorized_json(request, channel) if not isinstance(payload, dict): + diagnostics.rejected(channel, "not_an_object") raise HTTPException(status_code=422, detail="OpenTelemetry payload must be an object") if not is_enabled("codex", "observe"): + diagnostics.observation_off(channel) return Response(status_code=202) defaults_raw = payload.pop("openworkgraph", {}) try: defaults = OTelDefaults.model_validate(defaults_raw if isinstance(defaults_raw, dict) else {}).model_dump() - ingest_codex_otel_payload(payload, defaults=defaults) + result = ingest_codex_otel_payload(payload, defaults=defaults) except (AgentEvidenceError, ValueError) as exc: + diagnostics.rejected(channel, "adapter_error") raise HTTPException(status_code=422, detail=str(exc)) from exc + diagnostics.processed(channel, result) return Response(status_code=202) @router.post(AGENT_CLAUDE_OTEL_PATH, status_code=202) async def ingest_claude_otel(request: Request) -> Response: - _require_agent_write_bearer(request) - payload = await _read_bounded_json(request) + channel = "claude_code_otel_logs" + payload = await _authorized_json(request, channel) if not isinstance(payload, dict): + diagnostics.rejected(channel, "not_an_object") raise HTTPException(status_code=422, detail="OpenTelemetry payload must be an object") if not is_enabled("claude_code", "observe"): + diagnostics.observation_off(channel) return Response(status_code=202) defaults_raw = payload.pop("openworkgraph", {}) try: defaults = OTelDefaults.model_validate(defaults_raw if isinstance(defaults_raw, dict) else {}).model_dump() - ingest_claude_otel_payload(payload, defaults=defaults) + result = ingest_claude_otel_payload(payload, defaults=defaults) except (AgentEvidenceError, ValueError) as exc: + diagnostics.rejected(channel, "adapter_error") raise HTTPException(status_code=422, detail=str(exc)) from exc + diagnostics.processed(channel, result) return Response(status_code=202) +@router.get("/v1/agent-telemetry/diagnostics") +def get_agent_telemetry_diagnostics(request: Request) -> dict[str, Any]: + """Per-channel delivery counts so a missing signal can be located, not guessed.""" + _require_api_read_bearer(request) + from .agent_capture_runtime import configuration_state + + return {**diagnostics.snapshot(), "configuration": configuration_state()} + + @router.get("/v1/agent-workflows") def get_agent_workflows(request: Request, since: str | None = None, limit: int = 5000) -> dict[str, Any]: _require_api_read_bearer(request) diff --git a/server/agent_spool.py b/server/agent_spool.py new file mode 100644 index 0000000..ba5f103 --- /dev/null +++ b/server/agent_spool.py @@ -0,0 +1,222 @@ +from __future__ import annotations + +"""A short, bounded buffer for native agent events while OpenWorkGraph is busy. + +Native agent hooks (Claude Code) run outside OpenWorkGraph and POST each event +once. If that POST fails, the event may be kept on disk and delivered later, +but only under a recording lease: + +* Only a running, recording, non-demo OpenWorkGraph issues the lease. It is short + (``LEASE_TTL_SECONDS``), renewed while recording, and revoked immediately on + Pause, Stop and shutdown. A crashed OpenWorkGraph stops granting it within the + TTL. +* A hook spools only while the lease is valid *and* the capture state file of + the data folder that issued it says "recording" (read directly, never + defaulted). So after the user presses Stop or quits OpenWorkGraph, agent + events are dropped, never quietly collected for later. +* Flushing goes through the normal ingest path, which applies the Observe + switches, deletion tombstones, retention and pause/stop windows again. +* The spool is bounded by file count, file size and age. + +Spooled events are the same allowlisted structural events the hook would have +sent; no prompt, response, argument or result content exists in them. +""" + +import json +import os +import tempfile +import time +import uuid +from datetime import datetime, timezone +from pathlib import Path +from typing import Any, Callable + +from .local_auth import auth_dir + +LEASE_NAME = "agent_spool_lease.json" +SPOOL_DIR_NAME = "agent_spool" +LEASE_TTL_SECONDS = 90.0 +MAX_SPOOL_FILES = 2000 +MAX_FILE_BYTES = 256_000 +MAX_AGE_SECONDS = 24 * 3600.0 + + +def _lease_path() -> Path: + return auth_dir() / LEASE_NAME + + +def spool_dir() -> Path: + path = auth_dir() / SPOOL_DIR_NAME + path.mkdir(parents=True, exist_ok=True) + try: + os.chmod(path, 0o700) + except Exception: + pass + return path + + +def _atomic_write(path: Path, payload: bytes) -> None: + fd, tmp = tempfile.mkstemp(prefix=".tmp-", dir=str(path.parent)) + try: + with os.fdopen(fd, "wb") as handle: + handle.write(payload) + handle.flush() + os.fsync(handle.fileno()) + try: + os.chmod(tmp, 0o600) + except Exception: + pass + os.replace(tmp, path) + finally: + if os.path.exists(tmp): + os.unlink(tmp) + + +# ---------------------------------------------------------------- lease (server side) + +def issue_lease(data_dir: Path, *, lease_id: str, ttl_seconds: float = LEASE_TTL_SECONDS, now: float | None = None) -> dict[str, Any]: + issued = time.time() if now is None else now + lease = { + "lease_id": lease_id, + "data_dir": str(Path(data_dir).resolve()), + "issued_at": issued, + "valid_until": issued + float(ttl_seconds), + } + _atomic_write(_lease_path(), (json.dumps(lease) + "\n").encode("utf-8")) + return lease + + +def revoke_lease(*, lease_id: str | None = None) -> None: + """Remove the lease (only our own, when an id is given).""" + path = _lease_path() + try: + if lease_id is not None: + current = json.loads(path.read_text(encoding="utf-8")) + if current.get("lease_id") != lease_id: + return + path.unlink() + except FileNotFoundError: + pass + except Exception: + try: + path.unlink() + except Exception: + pass + + +# ---------------------------------------------------------------- lease (hook side) + +def _capture_state(data_dir: Path) -> str: + try: + value = json.loads((data_dir / "capture_control.json").read_text(encoding="utf-8")) + return str(value.get("state") or "") if isinstance(value, dict) else "" + except Exception: + return "" + + +def valid_lease(now: float | None = None) -> dict[str, Any] | None: + try: + lease = json.loads(_lease_path().read_text(encoding="utf-8")) + except Exception: + return None + if not isinstance(lease, dict): + return None + current = time.time() if now is None else now + try: + if float(lease.get("valid_until") or 0) <= current: + return None + except Exception: + return None + data_dir = Path(str(lease.get("data_dir") or "")) + if not str(lease.get("data_dir") or "") or _capture_state(data_dir) != "recording": + return None + return lease + + +def spool_events(events: list[dict[str, Any]], *, now: float | None = None) -> bool: + """Keep events for later delivery if, and only if, recording is leased.""" + if not events: + return False + lease = valid_lease(now) + if lease is None: + return False + directory = spool_dir() + try: + if sum(1 for p in directory.iterdir() if p.suffix == ".json") >= MAX_SPOOL_FILES: + return False + except Exception: + return False + payload = json.dumps({ + "lease_id": lease["lease_id"], + "data_dir": lease["data_dir"], + "spooled_at": datetime.now(timezone.utc).isoformat(), + "spooled_epoch": time.time() if now is None else now, + "events": events, + }, ensure_ascii=False, separators=(",", ":")).encode("utf-8") + if len(payload) > MAX_FILE_BYTES: + return False + try: + _atomic_write(directory / f"{time.time_ns()}-{uuid.uuid4().hex}.json", payload) + except Exception: + return False + return True + + +# ---------------------------------------------------------------- flush (server side) + +def flush_spool( + data_dir: Path, + ingest: Callable[[list[dict[str, Any]]], Any], + *, + now: float | None = None, +) -> dict[str, int]: + """Deliver spooled events that belong to this data folder, oldest first. + + Files older than ``MAX_AGE_SECONDS`` are discarded unread. An ingest error + that is about the payload discards that file; any other error stops the + flush so the file is retried next time. + """ + from shared.agent_evidence import AgentEvidenceError + + counts = {"delivered_files": 0, "delivered_events": 0, "expired_files": 0, "rejected_files": 0} + current = time.time() if now is None else now + ours = str(Path(data_dir).resolve()) + try: + files = sorted(p for p in spool_dir().iterdir() if p.suffix == ".json" and not p.name.startswith(".tmp-")) + except Exception: + return counts + for path in files: + try: + item = json.loads(path.read_text(encoding="utf-8")) + except Exception: + path.unlink(missing_ok=True) + counts["rejected_files"] += 1 + continue + age = current - float(item.get("spooled_epoch") or 0) + if age > MAX_AGE_SECONDS: + path.unlink(missing_ok=True) + counts["expired_files"] += 1 + continue + if str(item.get("data_dir") or "") != ours: + continue # another OpenWorkGraph data folder (e.g. demo) owns it + events = [e for e in item.get("events") or [] if isinstance(e, dict)] + try: + if events: + ingest(events) + except (AgentEvidenceError, ValueError): + path.unlink(missing_ok=True) + counts["rejected_files"] += 1 + continue + except Exception: + break + path.unlink(missing_ok=True) + counts["delivered_files"] += 1 + counts["delivered_events"] += len(events) + return counts + + +def pending_count() -> int: + try: + return sum(1 for p in spool_dir().iterdir() if p.suffix == ".json" and not p.name.startswith(".tmp-")) + except Exception: + return 0 diff --git a/server/agent_telemetry_diagnostics.py b/server/agent_telemetry_diagnostics.py new file mode 100644 index 0000000..0bef6bb --- /dev/null +++ b/server/agent_telemetry_diagnostics.py @@ -0,0 +1,104 @@ +from __future__ import annotations + +"""Structural delivery diagnostics for native agent telemetry channels. + +Answers "is Claude Code sending anything, and if so what happens to it?" with +counts, timestamps and reason codes only. No payload, attribute value, header or +error text is ever stored here. + +Per channel, a request is counted when it arrives, then as exactly one of: +rejected (with a reason code), observation off, or processed. Processed OTel +requests also report records seen/ignored by the structural allowlist and how +many events were stored versus suppressed by deletion/retention/pause windows. +""" + +import threading +from datetime import datetime, timezone +from typing import Any + +CHANNELS = { + "claude_code_hooks": "Claude Code hooks", + "claude_code_otel_logs": "Claude Code OpenTelemetry logs", + "codex_otel": "Codex OpenTelemetry", + "agent_events": "Other agent events (SDKs, adapters)", + "otel_generic": "Generic OpenTelemetry", + "spool": "Delayed delivery (agent spool)", +} +REJECTION_REASONS = ("auth", "invalid_payload", "too_large", "not_an_object", "adapter_error") + +_LOCK = threading.Lock() +_STARTED_AT = datetime.now(timezone.utc).isoformat() +_STATS: dict[str, dict[str, Any]] = {} + + +def _now() -> str: + return datetime.now(timezone.utc).isoformat() + + +def _channel(name: str) -> dict[str, Any]: + return _STATS.setdefault(name, { + "requests": 0, + "processed": 0, + "observation_off": 0, + "rejected": {}, + "records_seen": 0, + "records_ignored": 0, + "events_projected": 0, + "events_stored": 0, + "last_received_at": None, + "last_stored_at": None, + "last_rejected_at": None, + "last_rejection_reason": None, + }) + + +def received(channel: str) -> None: + with _LOCK: + entry = _channel(channel) + entry["requests"] += 1 + entry["last_received_at"] = _now() + + +def rejected(channel: str, reason: str) -> None: + reason = reason if reason in REJECTION_REASONS else "adapter_error" + with _LOCK: + entry = _channel(channel) + entry["rejected"][reason] = int(entry["rejected"].get(reason, 0)) + 1 + entry["last_rejected_at"] = _now() + entry["last_rejection_reason"] = reason + + +def observation_off(channel: str) -> None: + with _LOCK: + _channel(channel)["observation_off"] += 1 + + +def processed(channel: str, result: dict[str, Any] | None = None) -> None: + result = result or {} + stored = int(result.get("inserted") or 0) + projected = int(result.get("projected", result.get("received", stored)) or 0) + seen = int(result.get("records_seen", result.get("spans_seen", 0)) or 0) + ignored = int(result.get("records_ignored", result.get("spans_ignored", 0)) or 0) + with _LOCK: + entry = _channel(channel) + entry["processed"] += 1 + entry["records_seen"] += seen + entry["records_ignored"] += ignored + entry["events_projected"] += projected + entry["events_stored"] += stored + if stored: + entry["last_stored_at"] = _now() + + +def snapshot() -> dict[str, Any]: + with _LOCK: + channels = { + name: {"label": label, **{k: (dict(v) if isinstance(v, dict) else v) for k, v in _channel(name).items()}} + for name, label in CHANNELS.items() + } + return {"since": _STARTED_AT, "channels": channels} + + +def reset_for_tests() -> None: + with _LOCK: + _STATS.clear() diff --git a/server/collector_status.py b/server/collector_status.py new file mode 100644 index 0000000..24ecf49 --- /dev/null +++ b/server/collector_status.py @@ -0,0 +1,17 @@ +from __future__ import annotations + +"""What the latest collector heartbeat says about sensor visibility.""" + +from typing import Any + + +def sensor_state(status: dict[str, Any] | None) -> dict[str, Any]: + """Missing OS permissions and away state; a blind sensor is not "recording".""" + activity = (status or {}).get("activity") + activity = activity if isinstance(activity, dict) else {} + permissions = activity.get("permissions") + permissions = permissions if isinstance(permissions, dict) else {} + return { + "missing_permissions": sorted(str(k) for k, v in permissions.items() if v is False), + "away": bool(activity.get("away")), + } diff --git a/server/enterprise_app.py b/server/enterprise_app.py index 0fa0243..0dc9dfc 100644 --- a/server/enterprise_app.py +++ b/server/enterprise_app.py @@ -20,6 +20,7 @@ set_state, ) from shared.lifespan import extend_lifespan +from .collector_status import sensor_state from . import main as main_module from .context_layers import candidate_tasks, factual_context_timeline from .main import CONFIG_PATH, ROOT @@ -129,6 +130,9 @@ def _capture_status() -> dict[str, Any]: pass value = public_status(engaged_seconds=engaged) value["collector_alive"] = bool(main_module.COLLECTOR_STATUS) + # From the collector's own OS preflight: a sensor that is "on" but blind is + # reported as such, not as recording. + value.update(sensor_state(main_module.COLLECTOR_STATUS)) value["browser_sensor_alive"] = bool(main_module.BROWSER_STATUS) value["demo"] = _demo_mode() return value diff --git a/server/enterprise_runner.py b/server/enterprise_runner.py index a96f36a..572a72f 100644 --- a/server/enterprise_runner.py +++ b/server/enterprise_runner.py @@ -29,8 +29,14 @@ def main() -> None: import server.org_join_routes as org_join_routes import server.dashboard_privacy # noqa: F401 + import server.agent_capture_runtime as agent_capture_runtime + org_join_routes.start_managed_setup_in_background() - uvicorn.run(SECURE_APP, host=args.host, port=args.port) + agent_capture_runtime.start() + try: + uvicorn.run(SECURE_APP, host=args.host, port=args.port) + finally: + agent_capture_runtime.stop() if __name__ == "__main__": diff --git a/shared/claude_code_adapter.py b/shared/claude_code_adapter.py index c882dcc..a5328e2 100644 --- a/shared/claude_code_adapter.py +++ b/shared/claude_code_adapter.py @@ -16,6 +16,8 @@ _SUPPORTED_EVENTS = frozenset({ "SessionStart", "SessionEnd", + "UserPromptSubmit", + "Stop", "PostToolUse", "PostToolUseFailure", "PermissionRequest", @@ -176,6 +178,25 @@ def claude_hook_to_agent_events( event_key=hook, )] + if hook in {"UserPromptSubmit", "Stop"}: + # One turn = one prompt: it starts when the prompt is submitted and + # finishes when Claude stops responding. These payloads carry the prompt + # text and the last assistant message; neither is ever read here. Without + # a prompt_id there is no turn identity, and falling back to the session + # would make a turn's Stop look like the whole session finishing. + prompt_id = _text(payload.get("prompt_id"), 128) + if not prompt_id: + return [] + return [_base_event( + payload, + operation="run_started" if hook == "UserPromptSubmit" else "run_finished", + status="running" if hook == "UserPromptSubmit" else "success", + observed_at=timestamp, + run_id=prompt_id, + trace_id=trace_id, + event_key=hook, + )] + if hook in {"PostToolUse", "PostToolUseFailure"}: tool_use_id = _text(payload.get("tool_use_id"), 128) tool_name = _safe_label(payload.get("tool_name"), default="unknown-tool", limit=160) diff --git a/tests/conftest.py b/tests/conftest.py index 56fc448..9a902b9 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -32,6 +32,7 @@ def _isolate_agent_ingestion_database(request): "test_agent_ingestion_v059.py", "test_agent_native_routes_v060.py", "test_agent_ingress_boundaries_v084.py", + "test_agent_delivery_v098.py", } if Path(str(request.fspath)).name not in isolated_modules: yield diff --git a/tests/js/capture_correctness.test.mjs b/tests/js/capture_correctness.test.mjs new file mode 100644 index 0000000..e72021c --- /dev/null +++ b/tests/js/capture_correctness.test.mjs @@ -0,0 +1,21 @@ +import test from 'node:test'; +import assert from 'node:assert/strict'; +import fs from 'node:fs'; +import path from 'node:path'; +import {fileURLToPath} from 'node:url'; + +const root = path.resolve(path.dirname(fileURLToPath(import.meta.url)), '..', '..'); +const read = (p) => fs.readFileSync(path.join(root, p), 'utf8'); + +test('Connections rows say where each agent telemetry signal stands', () => { + const js = read('dashboard/connections.js'); + assert.doesNotThrow(() => new Function(js)); + for (const needle of ['/v1/agent-telemetry/diagnostics', 'claude_code_hooks', 'claude_code_otel_logs', 'codex_otel', + 'nothing received since OpenWorkGraph started', 'Hooks out of date']) assert.ok(js.includes(needle), needle); +}); + +test('the Recording pill does not hide a blind sensor', () => { + const js = read('dashboard/gateway_panel.js'); + assert.doesNotThrow(() => new Function(js)); + assert.ok(js.includes('missing_permissions') && js.includes('permission missing') && js.includes("' · away'")); +}); diff --git a/tests/test_agent_dashboard_control_plane_v088.py b/tests/test_agent_dashboard_control_plane_v088.py index a55d203..bebb172 100644 --- a/tests/test_agent_dashboard_control_plane_v088.py +++ b/tests/test_agent_dashboard_control_plane_v088.py @@ -53,7 +53,7 @@ def test_agent_setup_payload_is_write_only_and_privacy_minimized(tmp_path): claude=payload['integrations']['claude_code'] blob=__import__('json').dumps(claude['settings']) assert 'SessionStart' in blob and 'PostToolUse' in blob and 'SubagentStart' in blob -assert 'UserPromptSubmit' not in blob and 'PreToolUse' not in blob +assert 'UserPromptSubmit' in blob and 'Stop' in blob and 'PreToolUse' not in blob assert '"async": true' in blob codex=payload['integrations']['codex']['config'] diff --git a/tests/test_agent_delivery_v098.py b/tests/test_agent_delivery_v098.py new file mode 100644 index 0000000..3bbd440 --- /dev/null +++ b/tests/test_agent_delivery_v098.py @@ -0,0 +1,288 @@ +from __future__ import annotations + +"""Agent delivery: lease-gated spool, telemetry diagnostics, hook refresh.""" + +import json +import time +from datetime import datetime, timezone +from pathlib import Path + +import pytest +from fastapi.testclient import TestClient + +from adapters import _agent_client +from server import agent_capture_runtime as runtime +from server import agent_spool as spool +from server import agent_telemetry_diagnostics as diagnostics +from shared.claude_code_adapter import claude_hook_to_agent_events + + +@pytest.fixture() +def dirs(tmp_path, monkeypatch): + auth = tmp_path / "auth" + data = tmp_path / "live" + auth.mkdir() + data.mkdir() + monkeypatch.setenv("WORKFLOW_OBSERVER_AUTH_DIR", str(auth)) + monkeypatch.setenv("WORKFLOW_OBSERVER_DATA", str(data)) + monkeypatch.setattr(runtime, "_demo", lambda: False) + runtime._STATE.update(lease_active=False, lease_valid_until=None, last_flush=None) + return auth, data + + +def _state(data: Path, state: str) -> None: + (data / "capture_control.json").write_text(json.dumps({"state": state, "generation": 1, "skip_intervals": []}), encoding="utf-8") + + +def _events(prompt: str = "p1", tool: str = "t1") -> list[dict]: + return claude_hook_to_agent_events({ + "session_id": "s-delivery", "prompt_id": prompt, "hook_event_name": "PostToolUse", + "tool_use_id": tool, "tool_name": "Read", + }, observed_at=datetime.now(timezone.utc).isoformat()) + + +# --- the lease ----------------------------------------------------------------------------------- + +def test_spooling_requires_a_live_recording_lease(dirs): + _auth, data = dirs + assert spool.spool_events(_events()) is False # OpenWorkGraph never granted a lease + _state(data, "recording") + spool.issue_lease(data, lease_id="L1") + assert spool.spool_events(_events()) is True + _state(data, "paused") + assert spool.spool_events(_events()) is False # paused: dropped, not held for later + _state(data, "stopped") + assert spool.spool_events(_events()) is False + _state(data, "recording") + assert spool.spool_events(_events(), now=time.time() + spool.LEASE_TTL_SECONDS + 1) is False # expired + (data / "capture_control.json").unlink() + assert spool.spool_events(_events()) is False # unknown state is never "recording" + spool.revoke_lease() + assert spool.valid_lease() is None + assert spool.pending_count() == 1 + + +def test_runtime_grants_lease_only_while_recording_and_revokes_on_pause(dirs): + _auth, data = dirs + from shared.capture_control import initialize_run, set_state + + initialize_run("2026-09-28T08:00:00+00:00") + runtime.tick() + assert spool.valid_lease() is not None + set_state("pause") + runtime.tick() + assert spool.valid_lease() is None + set_state("resume") + runtime.tick() + assert spool.valid_lease() is not None + runtime.stop() + assert spool.valid_lease() is None + + +def test_demo_mode_never_grants_a_lease(dirs, monkeypatch): + from shared.capture_control import initialize_run + + initialize_run("2026-09-28T08:00:00+00:00") + monkeypatch.setattr(runtime, "_demo", lambda: True) + runtime.tick() + assert spool.valid_lease() is None + + +# --- flushing ------------------------------------------------------------------------------------ + +def test_flush_delivers_own_folder_skips_others_and_discards_old_or_invalid(dirs, tmp_path): + _auth, data = dirs + _state(data, "recording") + spool.issue_lease(data, lease_id="L1") + assert spool.spool_events(_events("p1")) + # A file from another data folder (e.g. demo) and an old one. + other = tmp_path / "demo" + other.mkdir() + _state(other, "recording") + spool.issue_lease(other, lease_id="L2") + assert spool.spool_events(_events("p3")) + old = spool.spool_dir() / "0-old.json" + old.write_text(json.dumps({"data_dir": str(data.resolve()), "spooled_epoch": time.time() - spool.MAX_AGE_SECONDS - 5, "events": _events("p4")})) + bad = spool.spool_dir() / "1-bad.json" + bad.write_text("{not json") + + delivered: list[list[dict]] = [] + counts = spool.flush_spool(data, delivered.append) + assert counts == {"delivered_files": 1, "delivered_events": 1, "expired_files": 1, "rejected_files": 1} + assert [e["run_id"] for batch in delivered for e in batch] == ["p1"] + assert spool.pending_count() == 1 # the other folder's file is left for its owner + + +def test_flush_retries_later_when_storage_is_unavailable(dirs): + _auth, data = dirs + _state(data, "recording") + spool.issue_lease(data, lease_id="L1") + assert spool.spool_events(_events()) + + def busy(_events): + raise RuntimeError("database is locked") + + assert spool.flush_spool(data, busy)["delivered_files"] == 0 + assert spool.pending_count() == 1 + + +def test_spooled_events_go_through_normal_ingest_rules(dirs): + """Pause windows and the Observe switch still apply when spooled events arrive.""" + _auth, data = dirs + from server.db import init_db, rows + from shared.capture_control import initialize_run, set_state + + init_db() + initialize_run("2026-09-28T08:00:00+00:00") + runtime.tick() + before = datetime.now(timezone.utc) + assert spool.spool_events(_events("kept-run", "k1")) + set_state("pause") + paused = claude_hook_to_agent_events({ + "session_id": "s-delivery", "prompt_id": "paused-run", "hook_event_name": "PostToolUse", "tool_use_id": "x", "tool_name": "Read", + }, observed_at=datetime.now(timezone.utc).isoformat()) + # Written directly (as if spooled a moment before the pause took effect). + (spool.spool_dir() / "9-late.json").write_text(json.dumps({"data_dir": str(data.resolve()), "spooled_epoch": time.time(), "events": paused})) + set_state("resume") + runtime._STATE["last_flush"] = None + runtime.tick() + assert spool.pending_count() == 0 + stored = rows("SELECT metadata_json FROM events WHERE source = 'agent' AND observed_at >= ?", (before.isoformat(),)) + runs = {r["metadata"]["trace"]["run_id"] for r in stored} + assert "kept-run" in runs and "paused-run" not in runs + + +# --- the hook client ----------------------------------------------------------------------------- + +def test_client_spools_on_unreachable_openworkgraph_and_drops_without_lease(dirs, monkeypatch): + _auth, data = dirs + monkeypatch.setenv("WORKFLOW_OBSERVER_API", "http://127.0.0.1:9") # nothing listens + monkeypatch.setenv("OWG_AGENT_INGEST_TOKEN", "t") + with pytest.raises(OSError): + _agent_client.post_agent_events(_events()) + assert spool.pending_count() == 0 + _state(data, "recording") + spool.issue_lease(data, lease_id="L1") + assert _agent_client.post_agent_events(_events())["status"] == "spooled" + assert spool.pending_count() == 1 + + +def test_client_never_spools_a_rejection(dirs, monkeypatch): + from urllib.error import HTTPError + + _auth, data = dirs + _state(data, "recording") + spool.issue_lease(data, lease_id="L1") + + def rejected(request, timeout): + raise HTTPError(request.full_url, 401, "unauthorized", {}, None) + + monkeypatch.setattr(_agent_client, "urlopen", rejected) + monkeypatch.setenv("OWG_AGENT_INGEST_TOKEN", "t") + with pytest.raises(HTTPError): + _agent_client.post_agent_events(_events()) + assert spool.pending_count() == 0 + + +# --- diagnostics --------------------------------------------------------------------------------- + +@pytest.fixture() +def client(): + # Only the agent router: starting the full local app here would install its + # middleware on the shared app object that later tests use. + from fastapi import FastAPI + + from server.agent_routes import router + + app = FastAPI() + app.include_router(router) + diagnostics.reset_for_tests() + with TestClient(app) as c: + yield c + + +def test_diagnostics_locate_where_a_signal_stops(client): + from server.agent_auth import ensure_agent_ingest_token + from server.local_auth import ensure_api_token + + write = {"Authorization": "Bearer " + ensure_agent_ingest_token()} + read = {"Authorization": "Bearer " + ensure_api_token()} + hooks = {"events": _events("diag-run", "d1")} + assert client.post("/agent-ingest/v1/events", json=hooks, headers={"Authorization": "Bearer wrong", "X-OWG-Channel": "claude_code_hooks"}).status_code == 401 + assert client.post("/agent-ingest/v1/events", json=hooks, headers=write).status_code == 200 + assert client.post("/agent-ingest/v1/claude-otel", json={"resourceLogs": []}, headers=write).status_code == 202 + assert client.post("/agent-ingest/v1/claude-otel", content=b"[1]", headers={**write, "Content-Type": "application/json"}).status_code == 422 + + assert client.get("/v1/agent-telemetry/diagnostics", headers=write).status_code == 401 # write-only token cannot read + report = client.get("/v1/agent-telemetry/diagnostics", headers=read).json() + hooks_stats = report["channels"]["claude_code_hooks"] + assert hooks_stats["requests"] == 2 and hooks_stats["rejected"] == {"auth": 1} and hooks_stats["events_stored"] >= 1 + otel = report["channels"]["claude_code_otel_logs"] + assert otel["requests"] == 2 and otel["processed"] == 1 and otel["records_seen"] == 0 and otel["rejected"] == {"not_an_object": 1} + assert report["channels"]["codex_otel"]["requests"] == 0 and report["channels"]["codex_otel"]["last_received_at"] is None + assert "agent_spool" in report["configuration"] and "claude_code" in report["configuration"] + # Only counts, timestamps and reason codes: nothing from the payloads. + assert "diag-run" not in json.dumps(report) and "Read" not in json.dumps(report["channels"]) + + +# --- Claude Code hook refresh -------------------------------------------------------------------- + +def _settings(tmp_path, monkeypatch, hooks: dict) -> Path: + path = tmp_path / "claude" / "settings.json" + path.parent.mkdir(parents=True) + path.write_text(json.dumps({"hooks": hooks, "env": {"MY_VAR": "1"}}), encoding="utf-8") + monkeypatch.setenv("OWG_CLAUDE_SETTINGS_PATH", str(path)) + return path + + +def _owg_group(): + from adapters.claude_code_hook import settings_fragment + + return settings_fragment()["hooks"]["PostToolUse"] + + +OLD_EVENTS = ["SessionStart", "SessionEnd", "PostToolUse", "PostToolUseFailure", "PermissionRequest", + "PermissionDenied", "SubagentStart", "SubagentStop", "StopFailure"] + + +def test_outdated_owg_hooks_are_refreshed_and_user_hooks_kept(dirs, tmp_path, monkeypatch): + user_hook = {"hooks": [{"type": "command", "command": "echo mine"}]} + hooks = {event: _owg_group() for event in OLD_EVENTS} + hooks["Stop"] = [user_hook] + path = _settings(tmp_path, monkeypatch, hooks) + from server import agent_config_writer as writer + + assert writer.claude_missing_hook_events(["Stop", "UserPromptSubmit", "SessionStart"]) == ["Stop", "UserPromptSubmit"] + result = runtime.refresh_outdated_claude_hooks() + assert result and result["added"] == ["Stop", "UserPromptSubmit"] + data = json.loads(path.read_text()) + assert user_hook in data["hooks"]["Stop"] and len(data["hooks"]["Stop"]) == 2 + assert data["env"]["MY_VAR"] == "1" + assert list(path.parent.glob("settings.json.owg-backup-*")) + assert runtime.refresh_outdated_claude_hooks() is None # now current + + +def test_hooks_are_never_installed_or_refreshed_without_consent(dirs, tmp_path, monkeypatch): + path = _settings(tmp_path, monkeypatch, {"Stop": [{"hooks": [{"type": "command", "command": "echo mine"}]}]}) + before = path.read_text() + assert runtime.refresh_outdated_claude_hooks() is None # OpenWorkGraph's hooks were never installed + assert path.read_text() == before + + path.write_text(json.dumps({"hooks": {event: _owg_group() for event in OLD_EVENTS}})) + from server import connections + + monkeypatch.setattr(connections, "is_enabled", lambda client, kind: False) # Observe switched off + before = path.read_text() + assert runtime.refresh_outdated_claude_hooks() is None + assert path.read_text() == before + + +def test_capture_status_reports_blind_sensors_and_away(): + from server.collector_status import sensor_state + + heartbeat = {"session_id": "s", "activity": {"permissions": {"accessibility": True, "input_monitoring": False}, "away": True}} + assert sensor_state(heartbeat) == {"missing_permissions": ["input_monitoring"], "away": True} + assert sensor_state({}) == {"missing_permissions": [], "away": False} + assert sensor_state({"activity": {"permissions": {"accessibility": None}}})["missing_permissions"] == [] # unknown is not "missing" + source = (Path(__file__).resolve().parents[1] / "server" / "enterprise_app.py").read_text() + assert "value.update(sensor_state(main_module.COLLECTOR_STATUS))" in source diff --git a/tests/test_capture_correctness_v098.py b/tests/test_capture_correctness_v098.py new file mode 100644 index 0000000..3ae9aa2 --- /dev/null +++ b/tests/test_capture_correctness_v098.py @@ -0,0 +1,329 @@ +from __future__ import annotations + +"""Capture correctness: one collector, away time, document boundaries, honest health. + +The collector loop runs for real against a scripted clock, foreground window and +OS input-idle clock, so these tests check the evidence it actually emits. +""" + +import json +import subprocess +import sys +import textwrap +import time as real_time +from datetime import datetime as real_datetime, timedelta, timezone +from pathlib import Path + +import pytest + +import collector.main as cm +from collector.boundaries import document_key, is_material_change +from collector.instance_lock import EXIT_ALREADY_RUNNING, CollectorLock +from collector.platform import ActiveWindow +from connector.policy import prepare_event_for_gateway + +T0 = real_datetime(2026, 9, 28, 8, 0, tzinfo=timezone.utc) + + +class Clock: + def __init__(self, end: float) -> None: + self.t = 0.0 + self.end = end + + def monotonic(self) -> float: + return 1000.0 + self.t + + def sleep(self, seconds: float) -> None: + self.t += seconds + if self.t >= self.end: + cm.STOP = True + + +def drive(tmp_path, monkeypatch, world, *, end: float, permissions=None, config=None, keyboard=False): + """Run the real collector loop; ``world(t)`` returns (app, title, idle_seconds, locked).""" + clock = Clock(end) + + class FakeDatetime(real_datetime): + @classmethod + def now(cls, tz=None): + return T0 + timedelta(seconds=clock.t) + + fake_time = type("T", (), {"monotonic": staticmethod(clock.monotonic), "sleep": staticmethod(clock.sleep), "time": staticmethod(real_time.time)}) + events: list[dict] = [] + heartbeats: list[dict] = [] + monkeypatch.setattr(cm, "LOCAL_DIR", tmp_path) + monkeypatch.setattr(cm, "time", fake_time) + monkeypatch.setattr(cm, "datetime", FakeDatetime) + monkeypatch.setattr(cm, "persist_event", lambda event, _outbox: events.append(event)) + monkeypatch.setattr(cm, "post_heartbeat", lambda status, _url: heartbeats.append(status)) + monkeypatch.setattr(cm, "active_window", lambda: ActiveWindow(*world(clock.t)[:2])) + monkeypatch.setattr(cm, "system_idle_seconds", lambda: world(clock.t)[2]) + monkeypatch.setattr(cm, "screen_locked", lambda: world(clock.t)[3]) + monkeypatch.setattr(cm, "sensor_permissions", lambda: dict(permissions or {})) + if keyboard: + monkeypatch.setattr(cm.KeyboardActivitySensor, "start", lambda self: True) + monkeypatch.setattr(cm.KeyboardActivitySensor, "stop", lambda self: None) + cfg = { + "backend_url": "http://127.0.0.1:9", + "poll_seconds": 2, + "screen_interactions_enabled": False, + "keyboard_activity_enabled": keyboard, + "excluded_title_patterns": [], + **(config or {}), + } + tmp_path.mkdir(parents=True, exist_ok=True) + path = tmp_path / "config.json" + path.write_text(json.dumps(cfg), encoding="utf-8") + assert cm.run(path) == 0 + return events, heartbeats + + +def spans(events, kind=None): + return [e for e in events if kind is None or e["event_type"] == kind] + + +def seconds(event) -> float: + return (real_datetime.fromisoformat(event["observed_at"]) - T0).total_seconds() + + +def assert_no_overlap_and_no_hole(events, end): + timed = sorted((seconds(e), seconds(e) + e["duration_seconds"]) for e in events if e["event_type"] in {"focus_span", "away_span"}) + for (a1, b1), (a2, b2) in zip(timed, timed[1:]): + assert a2 == pytest.approx(b1), (a1, b1, a2, b2) + assert timed[0][0] == 0 and timed[-1][1] == pytest.approx(end) + + +# --- away ---------------------------------------------------------------------------------------- + +def test_idle_becomes_an_away_span_and_duration_keeps_wall_clock_meaning(tmp_path, monkeypatch): + last_input = {"t": 0.0} + + def world(t): + if t <= 150 or t >= 900: + last_input["t"] = t # typing + return ("Microsoft Word", "Report.docx", t - last_input["t"], False) + + events, heartbeats = drive(tmp_path, monkeypatch, world, end=1000) + away = spans(events, "away_span") + # Away starts once 300 s passed without input (the first 300 s stay with Word + # as reading/thinking time) and ends when input returns. + assert [(seconds(e), e["duration_seconds"]) for e in away][0] == (450, 120) # checkpointed + assert sum(e["duration_seconds"] for e in away) == pytest.approx(450) + assert {e["metadata"]["reason"] for e in away} == {"no_input"} + assert away[0]["metadata"]["idle_source"] == "os_input_clock" and away[0]["app"] == "Away" + focus = spans(events, "focus_span") + assert any(e["metadata"]["focus_boundary"] == "away" and seconds(e) + e["duration_seconds"] == 450 for e in focus) + assert seconds([e for e in focus if seconds(e) >= 900][0]) == 900 # Word resumes with input + assert_no_overlap_and_no_hole(events, 1000) + assert any(h["app"] == "Away" and h["activity"]["away"] for h in heartbeats) + + +def test_screen_lock_is_away_immediately(tmp_path, monkeypatch): + events, _ = drive(tmp_path, monkeypatch, lambda t: ("Mail", "Inbox", 0.0, 100 <= t < 200), end=300) + away = spans(events, "away_span") + assert [(seconds(e), e["duration_seconds"], e["metadata"]["reason"]) for e in away] == [(100, 100, "screen_locked")] + assert_no_overlap_and_no_hole(events, 300) + + +def test_without_any_idle_signal_it_never_guesses_away(tmp_path, monkeypatch): + # No OS idle clock (Linux today) and no working input sensors: long spans stay + # focus spans, as before, rather than inventing absences. + events, _ = drive(tmp_path, monkeypatch, lambda t: ("Terminal", "zsh", None, None), end=900) + assert not spans(events, "away_span") + assert sum(e["duration_seconds"] for e in spans(events, "focus_span")) == pytest.approx(900) + + +def test_away_spans_never_leave_the_computer(): + away = {"event_type": "away_span", "app": "Away", "metadata": {}, "event_id": "x"} + assert prepare_event_for_gateway(away, {"allowed_event_types": []}) is None + assert prepare_event_for_gateway(away, {"allowed_event_types": ["away_span"]}) is None + focus = {"event_type": "focus_span", "app": "Word", "metadata": {}, "event_id": "y"} + assert prepare_event_for_gateway(focus, {"allowed_event_types": []}) is not None + + +# --- same-app documents ------------------------------------------------------------------------ + +def test_document_switch_in_same_app_is_a_debounced_boundary(tmp_path, monkeypatch): + def world(t): + if t < 60: + title = "Report.docx" + elif t < 80: + title = "(2) Report.docx — Edited" # noise: unread badge + edited marker + elif t in (80.0, 82.0): + title = "Budget.docx" # a 1-poll flicker must not split... (confirmed on 2nd poll) + elif t < 100: + title = "Budget.docx" + elif t == 100: + title = "Save As" # a one-poll dialog: no boundary + else: + title = "Budget.docx" + return ("Microsoft Word", title, 0.0, False) + + events, _ = drive(tmp_path, monkeypatch, world, end=140) + focus = spans(events, "focus_span") + assert [(seconds(e), e["window_title"], e["metadata"]["focus_boundary"]) for e in focus] == [ + (0, "Report.docx", "document_change"), + (80, "Budget.docx", "shutdown"), # boundary backdated to when Budget.docx first appeared + ] + assert_no_overlap_and_no_hole(events, 140) + + +def test_old_default_value_gets_document_boundaries_and_app_only_is_opt_in(tmp_path, monkeypatch): + world = lambda t: ("Word", "A.docx" if t < 50 else "B.docx", 0.0, False) + # "application" was copied into every config.json as the old default. + events, _ = drive(tmp_path, monkeypatch, world, end=100, config={"change_detection": "application"}) + assert [e["window_title"] for e in spans(events, "focus_span")] == ["A.docx", "B.docx"] + events, _ = drive(tmp_path / "only", monkeypatch, world, end=100, config={"change_detection": "application_only"}) + assert [e["window_title"] for e in spans(events, "focus_span")] == ["A.docx"] + + +def test_example_config_ships_the_new_defaults(): + example = json.loads((Path(__file__).resolve().parents[1] / "config.example.json").read_text()) + assert example["change_detection"] == "application_and_document" + assert example["away_detection_enabled"] is True and example["away_after_seconds"] == 300 + + +def test_document_key_ignores_noise_but_not_real_changes(): + assert document_key("(3) Inbox - Outlook") == document_key("Inbox - Outlook") + assert document_key("● main.py — owg") == document_key("main.py — owg") + assert document_key("Budget.xlsx - Edited") == document_key("Budget.xlsx") + assert document_key("Mötesanteckningar — Redigerad") == document_key("Mötesanteckningar") + assert document_key("Q3 plan.docx (Not Responding)") == document_key("Q3 plan.docx") + assert document_key("Report.docx") != document_key("Budget.docx") + assert not is_material_change("report.docx", "") # a title going empty is not a new document + + +def test_title_privacy_modes_are_respected(tmp_path, monkeypatch): + world = lambda t: ("Word", "A.docx" if t < 50 else ("(3) B.docx" if t < 70 else "B.docx"), 0.0, False) + # "none": no title information is recorded, so none is used for boundaries. + events, _ = drive(tmp_path, monkeypatch, world, end=100, config={"window_title_mode": "none"}) + focus = spans(events, "focus_span") + assert len(focus) == 1 and focus[0]["window_title"] == "" + # "hash": only hashes are stored; the badge change is not a new document. + events, _ = drive(tmp_path / "h", monkeypatch, world, end=100, config={"window_title_mode": "hash"}) + focus = spans(events, "focus_span") + assert len(focus) == 2 and all(len(e["window_title"]) == 16 for e in focus) + assert "document_key" not in json.dumps(events) + + +# --- health -------------------------------------------------------------------------------------- + +def test_routine_checkpoints_do_not_produce_health_evidence(tmp_path, monkeypatch): + events, heartbeats = drive(tmp_path, monkeypatch, lambda t: ("Code", "main.py", 0.0, False), end=3600) + assert not spans(events, "capture_health") + assert len([e for e in spans(events, "focus_span") if e["metadata"]["focus_boundary"] == "periodic_checkpoint"]) == 29 + assert heartbeats[-1]["activity"]["diagnostics"]["focus_checkpoint_count"] == 29 + + +def test_real_degradation_is_reported_once_then_rate_limited(tmp_path, monkeypatch): + # The foreground window is unreadable for a stretch (e.g. automation denied). + events, _ = drive(tmp_path, monkeypatch, lambda t: ("Unknown" if 100 <= t < 1500 else "Code", "", 0.0, False), end=1600) + health = spans(events, "capture_health") + assert 1 <= len(health) <= 4 + first = health[0]["metadata"] + assert first["capture_health"]["active_window_unavailable"] >= 1 + assert "focus_checkpoint_count" not in first["capture_health"] and "focus_checkpoint_count" in first["diagnostics"] + + +def test_missing_macos_permission_is_stated_not_hidden(tmp_path, monkeypatch, capsys): + events, heartbeats = drive( + tmp_path, monkeypatch, lambda t: ("Code", "main.py", 0.0, False), end=20, + permissions={"accessibility": True, "input_monitoring": False}, keyboard=True, + ) + out = capsys.readouterr().out + assert "Keyboard activity: BLOCKED (macOS input monitoring permission not granted" in out + assert "Keyboard activity: ON" not in out + assert "Window titles and UI labels: ON" in out + assert heartbeats[0]["keyboard_sensor"] is False + assert heartbeats[0]["activity"]["permissions"] == {"accessibility": True, "input_monitoring": False} + health = spans(events, "capture_health") + assert health and health[0]["metadata"]["missing_permissions"] == ["input_monitoring"] + + +# --- one collector ----------------------------------------------------------------------------- + +def test_second_collector_for_the_same_data_folder_does_not_start(tmp_path, monkeypatch, capsys): + monkeypatch.setattr(cm, "LOCAL_DIR", tmp_path) + holder = CollectorLock(tmp_path) + assert holder.acquire() + try: + assert cm.run(tmp_path / "config.json") == EXIT_ALREADY_RUNNING + assert "already recording" in capsys.readouterr().out + other = CollectorLock(tmp_path / "demo") + assert other.acquire() # a different data folder (demo) is independent + other.release() + finally: + holder.release() + again = CollectorLock(tmp_path) + assert again.acquire() + again.release() + + +def test_lock_is_released_when_the_holder_crashes(tmp_path): + root = Path(__file__).resolve().parents[1] + script = textwrap.dedent(f""" + import sys, time + sys.path.insert(0, {str(root)!r}) + from collector.instance_lock import CollectorLock + held = CollectorLock({str(tmp_path)!r}) + assert held.acquire() + print("locked", flush=True) + time.sleep(60) + """) + proc = subprocess.Popen([sys.executable, "-c", script], stdout=subprocess.PIPE, text=True) + try: + assert proc.stdout.readline().strip() == "locked" + assert not CollectorLock(tmp_path).acquire() + finally: + proc.kill() + proc.wait(timeout=10) + lock = CollectorLock(tmp_path) + assert lock.acquire(), "a crashed holder must not leave a stale lock" + lock.release() + + +def test_supervisor_backs_off_instead_of_respawning_every_second(monkeypatch): + import collector.secure_main as sm + + spawned = [] + + class Worker: + def __init__(self, *args, **kwargs): + spawned.append(real_time.monotonic()) + + def poll(self): + return EXIT_ALREADY_RUNNING + + def terminate(self): + pass + + monkeypatch.setattr(sm, "initialize_run", lambda *_a: None) + monkeypatch.setattr(sm, "read_state", lambda: {"state": "recording", "generation": 1}) + monkeypatch.setattr(sm.subprocess, "Popen", Worker) + monkeypatch.setattr(sm.signal, "signal", lambda *_a: None) + stop_at = real_time.monotonic() + 3.0 + original_sleep = real_time.sleep + + def fake_sleep(seconds): + original_sleep(0.05) + if real_time.monotonic() >= stop_at: + sm._STOP = True + + monkeypatch.setattr(sm.time, "sleep", fake_sleep) + sm.run_supervisor(Path("/nonexistent/config.json")) + assert len(spawned) == 1 # without the back-off this is 3 + + +def test_untrusted_process_blocks_mouse_and_keyboard_too(tmp_path, monkeypatch, capsys): + """pynput needs Accessibility for both listeners on macOS (seen on a real Mac).""" + monkeypatch.setattr(cm.InteractionSensor, "start", lambda self: True) + monkeypatch.setattr(cm.InteractionSensor, "stop", lambda self: None) + _events, heartbeats = drive( + tmp_path, monkeypatch, lambda t: ("Code", "main.py", 0.0, False), end=10, + permissions={"accessibility": False, "input_monitoring": True}, keyboard=True, + config={"screen_interactions_enabled": True}, + ) + out = capsys.readouterr().out + assert "Screen interaction capture: BLOCKED (macOS accessibility permission" in out + assert "Keyboard activity: BLOCKED (macOS accessibility permission" in out + assert heartbeats[0]["keyboard_sensor"] is False diff --git a/tests/test_claude_code_adapter_v060.py b/tests/test_claude_code_adapter_v060.py index 852c946..f777cf6 100644 --- a/tests/test_claude_code_adapter_v060.py +++ b/tests/test_claude_code_adapter_v060.py @@ -45,15 +45,40 @@ def test_post_tool_use_keeps_structure_and_drops_native_content(): def test_content_heavy_hooks_are_ignored_instead_of_guessed(): for hook, extra in ( - ("UserPromptSubmit", {"prompt": "top secret prompt"}), ("PreToolUse", {"tool_input": {"command": "secret"}, "tool_name": "Bash"}), - ("Stop", {"last_assistant_message": "secret answer"}), ("MessageDisplay", {"message": "secret UI content"}), + # Turn hooks without a prompt_id have no turn identity; never fall back + # to the session, or a turn's Stop would look like the session ending. + ("UserPromptSubmit", {"prompt": "top secret prompt"}), + ("Stop", {"last_assistant_message": "secret answer"}), ): payload = {"session_id": "s1", "hook_event_name": hook, **extra} assert claude_hook_to_agent_events(payload, observed_at="2026-09-25T00:00:00Z") == [] +def test_turn_hooks_bound_one_prompt_without_copying_content(): + start = claude_hook_to_agent_events({ + "session_id": "s1", "prompt_id": "p1", "hook_event_name": "UserPromptSubmit", + "prompt": "top secret prompt", "transcript_path": "/Users/x/secret.jsonl", "custom_instructions": "private", + }, observed_at="2026-09-25T00:00:00Z") + stop = claude_hook_to_agent_events({ + "session_id": "s1", "prompt_id": "p1", "hook_event_name": "Stop", + "last_assistant_message": "secret answer", "stop_hook_active": False, + }, observed_at="2026-09-25T00:01:00Z") + tool = claude_hook_to_agent_events({ + "session_id": "s1", "prompt_id": "p1", "hook_event_name": "PostToolUse", "tool_use_id": "t", "tool_name": "Read", + }, observed_at="2026-09-25T00:00:30Z") + session_end = claude_hook_to_agent_events({"session_id": "s1", "hook_event_name": "SessionEnd"}, observed_at="2026-09-25T00:02:00Z") + assert [(e["operation"], e["status"], e["run_id"]) for e in start + stop] == [ + ("run_started", "running", "p1"), ("run_finished", "success", "p1"), + ] + # The turn groups with its tools; the session keeps its own boundary. + assert tool[0]["run_id"] == "p1" and session_end[0]["run_id"] == "s1" + serialized = json.dumps(start + stop) + for forbidden in ("top secret prompt", "secret answer", "secret.jsonl", "private"): + assert forbidden not in serialized + + def test_session_subagent_permission_and_failure_mappings_are_structural(): cases = [ ({"session_id": "s", "hook_event_name": "SessionStart"}, "run_started", "running"), @@ -138,7 +163,7 @@ def test_settings_fragment_registers_only_safe_supported_hooks(): fragment = claude_code_hook.settings_fragment("python") hooks = fragment["hooks"] assert set(hooks) == set(claude_code_hook.SUPPORTED_EVENTS) - assert "UserPromptSubmit" not in hooks + assert {"UserPromptSubmit", "Stop"} <= set(hooks) # turn boundaries; content never read assert "PreToolUse" not in hooks handler = hooks["PostToolUse"][0]["hooks"][0] assert handler["command"] == claude_code_hook.hook_command("python") From 5deff881b4568452d044f063d44f2272652646c7 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:18:41 +0200 Subject: [PATCH 02/13] Fix empty focus span when a document change coincides with a checkpoint; make pause-window test independent of Windows clock resolution --- collector/main.py | 6 +++++- tests/test_agent_delivery_v098.py | 19 +++++++++++++------ tests/test_capture_correctness_v098.py | 9 +++++++++ 3 files changed, 27 insertions(+), 7 deletions(-) diff --git a/collector/main.py b/collector/main.py index 1ed824e..1dcdc45 100644 --- a/collector/main.py +++ b/collector/main.py @@ -832,7 +832,11 @@ def close_span() -> None: boundary_mono = max(pending_mono, current_started_mono) boundary_wall = pending_wall if pending_mono >= current_started_mono else current_started_wall first_state = pending_state - emit_focus(boundary_mono, "document_change") + if boundary_mono > current_started_mono: + emit_focus(boundary_mono, "document_change") + # else: the new title appeared exactly when the current span + # began (a checkpoint), so relabel it rather than emit an + # empty span for the old document. diagnostics["document_boundary_count"] += 1 open_span(first_state, key, boundary_wall, boundary_mono) else: diff --git a/tests/test_agent_delivery_v098.py b/tests/test_agent_delivery_v098.py index 3bbd440..aea8448 100644 --- a/tests/test_agent_delivery_v098.py +++ b/tests/test_agent_delivery_v098.py @@ -4,7 +4,7 @@ import json import time -from datetime import datetime, timezone +from datetime import datetime, timedelta, timezone from pathlib import Path import pytest @@ -135,15 +135,22 @@ def test_spooled_events_go_through_normal_ingest_rules(dirs): init_db() initialize_run("2026-09-28T08:00:00+00:00") runtime.tick() - before = datetime.now(timezone.utc) - assert spool.spool_events(_events("kept-run", "k1")) - set_state("pause") + # Explicit times: on Windows the clock ticks ~15 ms, so "now" for the pause, + # the event and the resume can coincide and leave an empty pause window. + t0 = datetime.now(timezone.utc) + at = lambda seconds: (t0 + timedelta(seconds=seconds)).isoformat() + before = t0 - timedelta(seconds=60) + kept = claude_hook_to_agent_events({ + "session_id": "s-delivery", "prompt_id": "kept-run", "hook_event_name": "PostToolUse", "tool_use_id": "k1", "tool_name": "Read", + }, observed_at=at(-5)) + assert spool.spool_events(kept) + set_state("pause", at=at(1)) paused = claude_hook_to_agent_events({ "session_id": "s-delivery", "prompt_id": "paused-run", "hook_event_name": "PostToolUse", "tool_use_id": "x", "tool_name": "Read", - }, observed_at=datetime.now(timezone.utc).isoformat()) + }, observed_at=at(2)) # Written directly (as if spooled a moment before the pause took effect). (spool.spool_dir() / "9-late.json").write_text(json.dumps({"data_dir": str(data.resolve()), "spooled_epoch": time.time(), "events": paused})) - set_state("resume") + set_state("resume", at=at(3)) runtime._STATE["last_flush"] = None runtime.tick() assert spool.pending_count() == 0 diff --git a/tests/test_capture_correctness_v098.py b/tests/test_capture_correctness_v098.py index 3ae9aa2..67d6e7e 100644 --- a/tests/test_capture_correctness_v098.py +++ b/tests/test_capture_correctness_v098.py @@ -327,3 +327,12 @@ def test_untrusted_process_blocks_mouse_and_keyboard_too(tmp_path, monkeypatch, assert "Screen interaction capture: BLOCKED (macOS accessibility permission" in out assert "Keyboard activity: BLOCKED (macOS accessibility permission" in out assert heartbeats[0]["keyboard_sensor"] is False + + +def test_document_change_seen_at_a_checkpoint_never_emits_an_empty_span(tmp_path, monkeypatch): + # The new title first appears in the same poll as the 120 s checkpoint. + events, _ = drive(tmp_path, monkeypatch, lambda t: ("Word", "A.docx" if t < 120 else "B.docx", 0.0, False), end=200) + focus = spans(events, "focus_span") + assert all(e["duration_seconds"] > 0 for e in focus), [(seconds(e), e["duration_seconds"]) for e in focus] + assert [(seconds(e), e["window_title"]) for e in focus] == [(0, "A.docx"), (120, "B.docx")] + assert_no_overlap_and_no_hole(events, 200) From 37ef04e874a0b8e27168621d8977afc074919b4b Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:23:08 +0200 Subject: [PATCH 03/13] Ephemeral agent history: only the end of the session itself closes it, not a turn or subagent finishing --- server/agent_ingest.py | 7 ++++++- tests/test_agent_delivery_v098.py | 29 +++++++++++++++++++++++++++++ 2 files changed, 35 insertions(+), 1 deletion(-) diff --git a/server/agent_ingest.py b/server/agent_ingest.py index fa7f1bd..54a0b11 100644 --- a/server/agent_ingest.py +++ b/server/agent_ingest.py @@ -41,7 +41,12 @@ def _closed_agent_sessions(events: list[dict[str, Any]]) -> list[str]: if str(metadata.get("operation") or "") != "run_finished": continue session_id = str(event.get("session_id") or "").strip() - if session_id: + trace = metadata.get("trace") if isinstance(metadata.get("trace"), dict) else {} + run_id = str(trace.get("run_id") or "").strip() + # Only the end of the session itself closes it. A turn (Claude Stop, + # Cursor stop) or a subagent finishing is a run *inside* the session; + # purging then would delete the session mid-way and tombstone the rest. + if session_id and (not run_id or run_id == session_id): closed.append(session_id) return list(dict.fromkeys(closed)) diff --git a/tests/test_agent_delivery_v098.py b/tests/test_agent_delivery_v098.py index aea8448..3180ddb 100644 --- a/tests/test_agent_delivery_v098.py +++ b/tests/test_agent_delivery_v098.py @@ -293,3 +293,32 @@ def test_capture_status_reports_blind_sensors_and_away(): assert sensor_state({"activity": {"permissions": {"accessibility": None}}})["missing_permissions"] == [] # unknown is not "missing" source = (Path(__file__).resolve().parents[1] / "server" / "enterprise_app.py").read_text() assert "value.update(sensor_state(main_module.COLLECTOR_STATUS))" in source + + +def test_ephemeral_history_survives_turns_and_is_purged_at_session_end(dirs): + """A turn's Stop must not close the session in ephemeral mode (it used to purge + the whole session and tombstone every later turn).""" + from server.agent_ingest import ingest_agent_payloads + from server.db import init_db, rows + from shared.capture_control import initialize_run + from shared.history_policy import update_retention + + init_db() + initialize_run("2026-09-28T08:00:00+00:00") + update_retention(human_mode="ephemeral", human_days=None, agent_mode="ephemeral", agent_days=None) + + def hook(name, prompt=None, **extra): + payload = {"session_id": "eph-s", "hook_event_name": name, **extra} + if prompt: + payload["prompt_id"] = prompt + return claude_hook_to_agent_events(payload, observed_at=datetime.now(timezone.utc).isoformat()) + + def stored(): + return [r for r in rows("SELECT session_id FROM events WHERE source = 'agent'") if r["session_id"] == "eph-s"] + + ingest_agent_payloads(hook("SessionStart") + hook("UserPromptSubmit", "t1") + hook("Stop", "t1")) + ingest_agent_payloads(hook("UserPromptSubmit", "t2") + hook("PostToolUse", "t2", tool_use_id="x", tool_name="Read")) + ingest_agent_payloads(hook("SubagentStop", "t2", agent_id="a1", agent_type="explore")) + assert len(stored()) == 6, "turns and subagents must not purge the running session" + ingest_agent_payloads(hook("SessionEnd")) + assert stored() == [], "the end of the session itself still purges ephemeral history" From a4273bf472b0da8e0dddaeb6c5c76f56aed4f669 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:40:55 +0200 Subject: [PATCH 04/13] Revoke the agent spool lease in the app shutdown step (uvicorn re-raises SIGTERM before finally/atexit) and never re-issue it after stop --- server/agent_capture_runtime.py | 3 ++- server/enterprise_runner.py | 8 ++++++++ tests/test_agent_delivery_v098.py | 19 +++++++++++++++++++ 3 files changed, 29 insertions(+), 1 deletion(-) diff --git a/server/agent_capture_runtime.py b/server/agent_capture_runtime.py index 6f3a4bb..72389b2 100644 --- a/server/agent_capture_runtime.py +++ b/server/agent_capture_runtime.py @@ -66,7 +66,8 @@ def tick(now: float | None = None) -> None: from shared.capture_control import read_state current = _time.time() if now is None else now - recording = not _demo() and read_state().get("state") == "recording" + # Never re-issue a lease once stop() has run (the loop may be mid-tick). + recording = not _STOP.is_set() and not _demo() and read_state().get("state") == "recording" if recording: valid_until = _STATE.get("lease_valid_until") or 0.0 if not _STATE["lease_active"] or valid_until - current < LEASE_RENEW_BEFORE_SECONDS: diff --git a/server/enterprise_runner.py b/server/enterprise_runner.py index 572a72f..caf6b61 100644 --- a/server/enterprise_runner.py +++ b/server/enterprise_runner.py @@ -33,6 +33,14 @@ def main() -> None: org_join_routes.start_managed_setup_in_background() agent_capture_runtime.start() + # uvicorn re-raises SIGTERM/SIGINT after its graceful shutdown, which ends + # the process before a finally/atexit block can run. Revoke the recording + # lease in the app's own shutdown step so a normal quit stops agent spooling + # at once (a hard kill is still bounded by the lease's short TTL). + from shared.lifespan import extend_lifespan + from server.secure_app import app as secure_app + + extend_lifespan(secure_app, shutdown=agent_capture_runtime.stop) try: uvicorn.run(SECURE_APP, host=args.host, port=args.port) finally: diff --git a/tests/test_agent_delivery_v098.py b/tests/test_agent_delivery_v098.py index 3180ddb..257a358 100644 --- a/tests/test_agent_delivery_v098.py +++ b/tests/test_agent_delivery_v098.py @@ -27,6 +27,7 @@ def dirs(tmp_path, monkeypatch): monkeypatch.setenv("WORKFLOW_OBSERVER_DATA", str(data)) monkeypatch.setattr(runtime, "_demo", lambda: False) runtime._STATE.update(lease_active=False, lease_valid_until=None, last_flush=None) + runtime._STOP.clear() return auth, data @@ -322,3 +323,21 @@ def stored(): assert len(stored()) == 6, "turns and subagents must not purge the running session" ingest_agent_payloads(hook("SessionEnd")) assert stored() == [], "the end of the session itself still purges ephemeral history" + + +def test_no_lease_is_issued_after_stop(dirs): + from shared.capture_control import initialize_run + + initialize_run("2026-09-28T08:00:00+00:00") + runtime.tick() + assert spool.valid_lease() is not None + runtime.stop() + runtime.tick() # a tick already in flight when stop() ran must not bring it back + assert spool.valid_lease() is None + + +def test_real_launcher_revokes_the_lease_on_a_normal_quit(): + """uvicorn re-raises SIGTERM after shutdown, so finally/atexit never run; the + revoke must happen in the app's shutdown step.""" + source = (Path(__file__).resolve().parents[1] / "server" / "enterprise_runner.py").read_text() + assert "extend_lifespan(secure_app, shutdown=agent_capture_runtime.stop)" in source From b981199d448f074a6e6d94cb3b7f3ca455cf903d Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 17:49:01 +0200 Subject: [PATCH 05/13] Keep the Recording pill short; name missing macOS permissions in its hover text --- dashboard/gateway_panel.js | 3 ++- tests/js/capture_correctness.test.mjs | 2 +- 2 files changed, 3 insertions(+), 2 deletions(-) diff --git a/dashboard/gateway_panel.js b/dashboard/gateway_panel.js index 674f495..510e53c 100644 --- a/dashboard/gateway_panel.js +++ b/dashboard/gateway_panel.js @@ -122,7 +122,8 @@ if (state==='recording') { const names={accessibility:'Accessibility',input_monitoring:'Input Monitoring'}; const missing=(s.missing_permissions||[]).map(p=>names[p]||p); - label.textContent=`Recording · ${fmtDuration(s.run_elapsed_seconds)}`+(s.away?' · away':'')+(missing.length?` · macOS ${missing.join(' and ')} permission missing`:''); + // Short enough for the header; the hover text names each missing permission. + label.textContent=`Recording · ${fmtDuration(s.run_elapsed_seconds)}`+(s.away?' · away':'')+(missing.length?' · macOS permission missing':''); label.title=missing.length?`OpenWorkGraph cannot see everything: allow it under System Settings → Privacy & Security → ${missing.join(' / ')}, then restart OpenWorkGraph.`:''; dot.style.background='#d25a5a'; buttons.innerHTML=''; diff --git a/tests/js/capture_correctness.test.mjs b/tests/js/capture_correctness.test.mjs index e72021c..5f9534a 100644 --- a/tests/js/capture_correctness.test.mjs +++ b/tests/js/capture_correctness.test.mjs @@ -17,5 +17,5 @@ test('Connections rows say where each agent telemetry signal stands', () => { test('the Recording pill does not hide a blind sensor', () => { const js = read('dashboard/gateway_panel.js'); assert.doesNotThrow(() => new Function(js)); - assert.ok(js.includes('missing_permissions') && js.includes('permission missing') && js.includes("' · away'")); + assert.ok(js.includes('missing_permissions') && js.includes('macOS permission missing') && js.includes("' · away'")); }); From 1dbaced5d841d87666c9cdafaa0d235076c9ddea Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:37:48 +0200 Subject: [PATCH 06/13] Release capture correctness as v0.98.0 --- VERSION | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/VERSION b/VERSION index 05e39cb..95fce8c 100644 --- a/VERSION +++ b/VERSION @@ -1 +1 @@ -0.97.0 +0.98.0 From 3055fad8354cbed90f839c5b6836e10c980454b8 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:38:01 +0200 Subject: [PATCH 07/13] Align Python package version for v0.98.0 --- pyproject.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pyproject.toml b/pyproject.toml index 3963b60..791abbe 100644 --- a/pyproject.toml +++ b/pyproject.toml @@ -1,6 +1,6 @@ [project] name = "workflow-observer" -version = "0.97.0" +version = "0.98.0" description = "Local-first work evidence, self-hosted organizational context gateway, REST API, and MCP access." requires-python = ">=3.11" license = {file = "LICENSE"} From 823175638e0b0a4774eb74ffb3f257c4398d6f99 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:38:18 +0200 Subject: [PATCH 08/13] Align MCP bundle version for v0.98.0 --- mcpb/manifest.json | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) diff --git a/mcpb/manifest.json b/mcpb/manifest.json index 7da8394..9b32a63 100644 --- a/mcpb/manifest.json +++ b/mcpb/manifest.json @@ -2,9 +2,9 @@ "manifest_version": "0.3", "name": "openworkgraph-local", "display_name": "OpenWorkGraph", - "version": "0.97.0", + "version": "0.98.0", "description": "Connect Claude Desktop to the compact local OpenWorkGraph context surface.", - "long_description": "Uses the OpenWorkGraph installation already running on this computer. v0.97 adds named Gateway administrators, an employee roster, identity-bound personal invitations, optional company sign-in, and an employee /me view of organization-held evidence and recorded reads. v0.96 added an evidence-driven first-value dashboard layer that reconstructs the current session from existing privacy-hardened evidence without adding sensors, permissions, AI access, retention, or MCP capabilities. v0.95 added a framework-neutral custom-harness setup flow plus standalone Python and Node helpers for privacy-safe structural agent telemetry; arbitrary MCP-capable harnesses can separately read authorized OpenWorkGraph context. v0.94 added explicit local history retention, separate saved-history AI access, lightweight history navigation, and structural browser-agent lifecycle observation. Canonical workflow evidence remains primary; Context Pulse provides incremental factual updates. Retained history is separately user-controlled: list_history can navigate saved human and agent sessions only while a time-limited saved-history lease is active. The legacy 24-tool MCP entrypoint remains available for existing configurations while new connections use this compact surface. Prompts, model responses, tool arguments/results, typed text, clipboard contents, exception text, returned values, and hidden reasoning are not captured by the custom agent helpers.", + "long_description": "Uses the OpenWorkGraph installation already running on this computer. v0.98 hardens capture correctness with single-instance recording, local-only away spans, sparse degradation-only health evidence, truthful macOS permission status, debounced same-app document boundaries, Claude Code turn boundaries, and lease-gated durable structural agent delivery. v0.97 added named Gateway administrators, an employee roster, identity-bound personal invitations, optional company sign-in, and an employee /me view of organization-held evidence and recorded reads. v0.96 added an evidence-driven first-value dashboard layer that reconstructs the current session from existing privacy-hardened evidence without adding sensors, permissions, AI access, retention, or MCP capabilities. v0.95 added a framework-neutral custom-harness setup flow plus standalone Python and Node helpers for privacy-safe structural agent telemetry; arbitrary MCP-capable harnesses can separately read authorized OpenWorkGraph context. v0.94 added explicit local history retention, separate saved-history AI access, lightweight history navigation, and structural browser-agent lifecycle observation. Canonical workflow evidence remains primary; Context Pulse provides incremental factual updates. Retained history is separately user-controlled: list_history can navigate saved human and agent sessions only while a time-limited saved-history lease is active. The legacy 24-tool MCP entrypoint remains available for existing configurations while new connections use this compact surface. Prompts, model responses, tool arguments/results, typed text, clipboard contents, exception text, returned values, and hidden reasoning are not captured by the custom agent helpers.", "author": {"name": "Koyar Afrasyab / Kinvectum"}, "repository": {"type": "git", "url": "https://github.com/KAVentures/openworkgraph"}, "server": {"type": "node", "entry_point": "server/index.js", "mcp_config": {"command": "node", "args": ["${__dirname}/server/index.js"], "env": {}}}, From 5d22e09ee4578be9b30b917b28e6fe951ec34dc2 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:38:28 +0200 Subject: [PATCH 09/13] Align Python agent SDK version for v0.98.0 --- sdk/python/pyproject.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/python/pyproject.toml b/sdk/python/pyproject.toml index 9a013fd..6972ed5 100644 --- a/sdk/python/pyproject.toml +++ b/sdk/python/pyproject.toml @@ -4,7 +4,7 @@ build-backend = "setuptools.build_meta" [project] name = "openworkgraph-agent" -version = "0.97.0" +version = "0.98.0" description = "Dependency-free structural telemetry helper for custom OpenWorkGraph agent harnesses" requires-python = ">=3.10" license = {text = "Apache-2.0"} From d46a427f644b63c583829e1be5937f76636087a2 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:38:37 +0200 Subject: [PATCH 10/13] Align Node agent SDK version for v0.98.0 --- sdk/typescript/package.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/sdk/typescript/package.json b/sdk/typescript/package.json index 74da4fa..6bad59f 100644 --- a/sdk/typescript/package.json +++ b/sdk/typescript/package.json @@ -1,6 +1,6 @@ { "name": "@openworkgraph/agent", - "version": "0.97.0", + "version": "0.98.0", "description": "Dependency-free structural telemetry helper for custom OpenWorkGraph agent harnesses", "type": "module", "exports": { From cc06efeb2fd0e87a2ef2f1952e56e0138f04e908 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 18:38:55 +0200 Subject: [PATCH 11/13] Expect v0.98.0 release version --- tests/test_release_version_v087.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/tests/test_release_version_v087.py b/tests/test_release_version_v087.py index b0582a2..fed13a6 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.97.0" +EXPECTED_VERSION = "0.98.0" def test_release_version_sources_are_aligned(): From 0fbd9b15f8ea5d41892480020266177874ee2ed4 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:00:19 +0200 Subject: [PATCH 12/13] Gate agent spooling on Observe switch --- server/agent_spool.py | 36 +++++++++++++++++++++++++++++------- 1 file changed, 29 insertions(+), 7 deletions(-) diff --git a/server/agent_spool.py b/server/agent_spool.py index ba5f103..5597a7b 100644 --- a/server/agent_spool.py +++ b/server/agent_spool.py @@ -4,16 +4,16 @@ Native agent hooks (Claude Code) run outside OpenWorkGraph and POST each event once. If that POST fails, the event may be kept on disk and delivered later, -but only under a recording lease: +but only under a recording lease and while its framework's Observe switch is on: * Only a running, recording, non-demo OpenWorkGraph issues the lease. It is short (``LEASE_TTL_SECONDS``), renewed while recording, and revoked immediately on Pause, Stop and shutdown. A crashed OpenWorkGraph stops granting it within the TTL. -* A hook spools only while the lease is valid *and* the capture state file of - the data folder that issued it says "recording" (read directly, never - defaulted). So after the user presses Stop or quits OpenWorkGraph, agent - events are dropped, never quietly collected for later. +* A hook spools only while the lease is valid, the capture state file of the data + folder that issued it says "recording" (read directly, never defaulted), and + every switch-controlled framework in the batch still has Observe enabled. + So Pause, Stop, quit, or switching Observe off prevents new durable agent data. * Flushing goes through the normal ingest path, which applies the Observe switches, deletion tombstones, retention and pause/stop windows again. * The spool is bounded by file count, file size and age. @@ -133,9 +133,31 @@ def valid_lease(now: float | None = None) -> dict[str, Any] | None: return lease +def _observe_switch_allows(events: list[dict[str, Any]]) -> bool: + """Fail closed for durable buffering when a framework's Observe switch is off. + + This is intentionally checked before anything is written to disk. Normal + direct delivery is still rechecked by the ingest route, and flush rechecks it + again, so a switch-off is enforced at every persistence boundary. + """ + try: + from .connections import FRAMEWORK_CLIENTS, is_enabled + + for event in events: + framework = str(event.get("framework") or "") if isinstance(event, dict) else "" + client = FRAMEWORK_CLIENTS.get(framework) + if client and not is_enabled(client, "observe"): + return False + return True + except Exception: + # Spooling is optional reliability, not a reason to risk capturing after + # consent/state can no longer be verified. + return False + + def spool_events(events: list[dict[str, Any]], *, now: float | None = None) -> bool: - """Keep events for later delivery if, and only if, recording is leased.""" - if not events: + """Keep events only while recording is leased and Observe still allows it.""" + if not events or not _observe_switch_allows(events): return False lease = valid_lease(now) if lease is None: From 634e181bf5016dd2ab43e66c69978c2e251fd313 Mon Sep 17 00:00:00 2001 From: Kinvectum <134240819+KAVentures@users.noreply.github.com> Date: Mon, 28 Sep 2026 21:00:41 +0200 Subject: [PATCH 13/13] Test Observe-off agent spool boundary --- tests/test_agent_spool_observe_v098.py | 43 ++++++++++++++++++++++++++ 1 file changed, 43 insertions(+) create mode 100644 tests/test_agent_spool_observe_v098.py diff --git a/tests/test_agent_spool_observe_v098.py b/tests/test_agent_spool_observe_v098.py new file mode 100644 index 0000000..0b1617f --- /dev/null +++ b/tests/test_agent_spool_observe_v098.py @@ -0,0 +1,43 @@ +from __future__ import annotations + +import json +from datetime import datetime, timezone + +from server import agent_spool as spool +from server import connections +from shared.claude_code_adapter import claude_hook_to_agent_events + + +def _events() -> list[dict]: + return claude_hook_to_agent_events( + { + "session_id": "s-observe", + "prompt_id": "p-observe", + "hook_event_name": "PostToolUse", + "tool_use_id": "t-observe", + "tool_name": "Read", + }, + observed_at=datetime.now(timezone.utc).isoformat(), + ) + + +def test_observe_off_prevents_agent_event_from_being_spooled(tmp_path, monkeypatch): + auth = tmp_path / "auth" + data = tmp_path / "live" + auth.mkdir() + data.mkdir() + monkeypatch.setenv("WORKFLOW_OBSERVER_AUTH_DIR", str(auth)) + monkeypatch.setenv("WORKFLOW_OBSERVER_DATA", str(data)) + (data / "capture_control.json").write_text( + json.dumps({"state": "recording", "generation": 1, "skip_intervals": []}), + encoding="utf-8", + ) + spool.issue_lease(data, lease_id="observe-lease") + + connections._write_switch("claude_code", "observe", False) + assert spool.spool_events(_events()) is False + assert spool.pending_count() == 0 + + connections._write_switch("claude_code", "observe", True) + assert spool.spool_events(_events()) is True + assert spool.pending_count() == 1