diff --git a/docs/CONTEXT_PULSE.md b/docs/CONTEXT_PULSE.md new file mode 100644 index 00000000..1a7bd697 --- /dev/null +++ b/docs/CONTEXT_PULSE.md @@ -0,0 +1,76 @@ +# Context Pulse + +Context Pulse is the incremental memory surface of OpenWorkGraph's compact Context MCP. +It complements, rather than replaces, `get_workflow_trace`. + +## Purpose + +A connected AI often does not need the entire retained workflow history on every turn. +`get_context_pulse` answers two narrower factual questions: + +1. What canonical evidence arrived since this caller last checked? +2. Which evidence-backed long-horizon findings are new or materially changed? + +The AI remains responsible for interpretation and suggestions. OpenWorkGraph does not +turn a repeated pattern into a recommendation, policy, permission, or productivity score. + +## Cursor contract + +The first call establishes a baseline and returns `next_cursor`. The caller should pass +that cursor back unchanged on its next check. The cursor is caller-owned and is not a +server-side subscription or timer. + +Recent evidence is bookmarked by the monotonic local database arrival ID rather than by +`observed_at`. This means an event that arrives late with an older observation timestamp +is still returned on the next pulse. + +If one pulse needs multiple pages, OpenWorkGraph freezes the event watermark and finding +snapshot until both are drained. Events that arrive during paging wait for the next +completed pulse, so they are neither mixed into the current snapshot nor skipped. + +The encoded cursor contains no captured titles, names, surfaces, or other observed text. +It contains watermarks plus opaque finding IDs/versions needed to know what has already +been delivered. + +## Findings v1 + +The initial factual finding set is intentionally conservative and uses the factual context +timeline, not inferred task labels: + +- **Surface engagement**: engaged/foreground seconds, span count and active-day count for + a surface over the rolling lookback. A change becomes material after another active day + or another 15 minutes of engaged time. +- **Repeated surface transition**: a directional transition between two different work + surfaces observed at least three times within sessions. Each additional occurrence is + material. + +Every emitted finding carries evidence event IDs and explicitly reports that it is a +factual aggregate, that task inference was not used, and that it is not advice. + +The first call labels findings `baseline`. Later calls emit only `new` or `changed` +findings. If `finding_limit` is smaller than the number of changed findings, unseen +findings are not marked delivered; subsequent calls continue the same frozen snapshot. + +## Privacy and disclosure + +Context Pulse does not add new capture. It uses evidence OpenWorkGraph already stores. +Typed text and clipboard contents remain uncaptured. + +The existing Context disclosure boundary applies to the entire Pulse response: + +- **Redacted** remains the default AI representation. +- **Full** is available only when the user selected it and organization policy permits it. +- Organization policy can restrict disclosure. + +Canonical local evidence is never rewritten by the Pulse or by AI redaction. + +## Pull, not push + +Context Pulse is a protocol capability, not an autonomous background sender. An MCP client +calls it when that client chooses to refresh context. Clients that support schedules, +loops, or long-running agent logic can call it periodically; ordinary conversational +clients can call it at session start or whenever fresh workflow context is relevant. + +This keeps the core primitive portable across ChatGPT, Claude, Codex, other MCP clients, +and future local or Gateway-based integrations without requiring any one client's +scheduling model. diff --git a/mcp_server/compact_hardening.py b/mcp_server/compact_hardening.py index eca9d64d..ec717c55 100644 --- a/mcp_server/compact_hardening.py +++ b/mcp_server/compact_hardening.py @@ -58,6 +58,39 @@ def apply_compact_hardening(compact_module: ModuleType) -> None: return compact_module._is_readable_step_input = _readable_step_input compact_module.secure_runtime = _CompactRuntimeProxy(compact_module.secure_runtime) + + @compact_module.mcp.tool() + def get_context_pulse( + cursor: str | None = None, + recent_limit: int = 100, + finding_limit: int = 20, + lookback_days: int = 30, + ) -> dict[str, Any]: + """Return what changed since the last check plus changed factual findings. + + Pass ``next_cursor`` back unchanged on the next call. The recent section is + canonical evidence that arrived since this caller's bookmark. Findings are + deterministic evidence-backed counts/aggregates over the rolling lookback, + never recommendations or task-label guesses. On the first call they are a + baseline; later calls omit findings that have not materially changed. + """ + name = "get_context_pulse" + compact_module.core._begin(name) + params: dict[str, Any] = { + "recent_limit": min(max(1, int(recent_limit)), 500), + "finding_limit": min(max(0, int(finding_limit)), 50), + "lookback_days": min(max(7, int(lookback_days)), 90), + } + if cursor not in (None, ""): + params["cursor"] = cursor + return compact_module.core._finish( + name, + compact_module.secure_runtime.secure_get("/v1/context-pulse", params), + ) + + # Keep the function addressable for direct unit tests/debugging in addition to + # registering it with both compact stdio and compact HTTP MCP transports. + compact_module.get_context_pulse = get_context_pulse compact_module._V0873_HARDENING_APPLIED = True diff --git a/mcpb/manifest.json b/mcpb/manifest.json index 7c4539ad..7ce90500 100644 --- a/mcpb/manifest.json +++ b/mcpb/manifest.json @@ -4,7 +4,7 @@ "display_name": "OpenWorkGraph", "version": "0.91.0", "description": "Connect Claude Desktop to the compact local OpenWorkGraph context surface.", - "long_description": "Uses the OpenWorkGraph installation already running on this computer. New connections expose canonical workflow evidence first through the compact read-oriented Context MCP surface, while derived task and pattern views remain optional and non-authoritative. The legacy 24-tool stdio entrypoint remains available for existing configurations. Workflow evidence remains in the local OpenWorkGraph store until Claude requests it through MCP. AI access can be turned off instantly from the OpenWorkGraph dashboard.", + "long_description": "Uses the OpenWorkGraph installation already running on this computer. New connections expose canonical workflow evidence first through the compact read-oriented Context MCP surface, while Context Pulse provides incremental factual updates and changed long-horizon findings. Derived task and pattern views remain optional and non-authoritative. The legacy 24-tool stdio entrypoint remains available for existing configurations. Workflow evidence remains in the local OpenWorkGraph store until Claude requests it through MCP. AI access can be turned off instantly from the OpenWorkGraph dashboard.", "author": { "name": "Koyar Afrasyab / Kinvectum" }, @@ -23,6 +23,7 @@ }, "tools": [ {"name": "get_current_work_context", "description": "Read an optional compact derived overview with pointers to canonical evidence."}, + {"name": "get_context_pulse", "description": "Read what changed since the last check plus new or materially changed factual findings; pass the returned cursor back next time."}, {"name": "search_work", "description": "Search prior work through evidence or semantic layers."}, {"name": "get_workflow_trace", "description": "Read the primary canonical paginated workflow evidence, optionally for one session."}, {"name": "get_work_profile", "description": "Read derived workflow signals without productivity scoring."}, diff --git a/server/context_pulse.py b/server/context_pulse.py new file mode 100644 index 00000000..a2e37c1b --- /dev/null +++ b/server/context_pulse.py @@ -0,0 +1,397 @@ +from __future__ import annotations + +"""Incremental, evidence-backed context updates for AI clients. + +Context Pulse answers two factual questions without becoming an advice engine: + +1. What canonical evidence arrived since this caller last checked? +2. Which long-horizon factual aggregates are new or materially changed? + +The cursor intentionally contains no titles, names, surfaces or other captured text. +It keeps only monotonic database watermarks and hashes/versions of findings. +""" + +import base64 +from collections import defaultdict +from datetime import datetime, timedelta, timezone +import hashlib +import json +from typing import Any + +from shared.evidence import RAW_RICH_EVIDENCE_CONTRACT, rich_evidence_row +from .context_layers import factual_context_timeline +from .db import connect + +_CURSOR_VERSION = 1 +_MAX_CURSOR_BYTES = 32_768 +_MAX_TRACKED_FINDINGS = 50 +_MAX_FINDING_SOURCE_EVENTS = 100_000 + + +def _encode_cursor(value: dict[str, Any]) -> str: + raw = json.dumps(value, ensure_ascii=False, separators=(",", ":"), sort_keys=True).encode("utf-8") + return base64.urlsafe_b64encode(raw).decode("ascii").rstrip("=") + + +def _decode_cursor(value: str | None) -> dict[str, Any]: + if not value: + return {} + try: + padded = str(value) + "=" * (-len(str(value)) % 4) + raw = base64.urlsafe_b64decode(padded.encode("ascii")) + if len(raw) > _MAX_CURSOR_BYTES: + raise ValueError("Context Pulse cursor is too large") + data = json.loads(raw.decode("utf-8")) + except ValueError: + raise + except Exception as exc: + raise ValueError("Invalid Context Pulse cursor") from exc + if not isinstance(data, dict) or int(data.get("v") or 0) != _CURSOR_VERSION: + raise ValueError("Invalid Context Pulse cursor version") + return data + + +def _event_dict(row: Any) -> dict[str, Any]: + value = dict(row) + raw = value.get("metadata_json") + if isinstance(raw, str): + try: + parsed = json.loads(raw or "{}") + value["metadata"] = parsed if isinstance(parsed, dict) else {} + except Exception: + value["metadata"] = {} + return value + + +def _finding_id(kind: str, *parts: str) -> str: + material = "\x1f".join([kind, *(str(part).strip().casefold() for part in parts)]) + return f"finding:{kind}:" + hashlib.sha256(material.encode("utf-8")).hexdigest()[:16] + + +def _parse_ts(value: Any) -> float | None: + try: + return datetime.fromisoformat(str(value or "").replace("Z", "+00:00")).timestamp() + except Exception: + return None + + +def _parse_datetime(value: Any) -> datetime: + try: + parsed = datetime.fromisoformat(str(value or "").replace("Z", "+00:00")) + return parsed if parsed.tzinfo is not None else parsed.replace(tzinfo=timezone.utc) + except Exception as exc: + raise ValueError("Invalid Context Pulse snapshot time") from exc + + +def _source_rows( + *, + snapshot_max_id: int, + snapshot_at: str, + lookback_days: int, +) -> tuple[list[dict[str, Any]], bool]: + since = (_parse_datetime(snapshot_at).astimezone(timezone.utc) - timedelta(days=lookback_days)).isoformat() + with connect() as conn: + rows = conn.execute( + "SELECT * FROM events WHERE id <= ? AND observed_at >= ? " + "ORDER BY observed_at ASC, id ASC LIMIT ?", + (snapshot_max_id, since, _MAX_FINDING_SOURCE_EVENTS + 1), + ).fetchall() + truncated = len(rows) > _MAX_FINDING_SOURCE_EVENTS + return [_event_dict(row) for row in rows[:_MAX_FINDING_SOURCE_EVENTS]], truncated + + +def _surface_findings(timeline: list[dict[str, Any]], lookback_days: int) -> list[dict[str, Any]]: + grouped: dict[str, dict[str, Any]] = {} + for row in timeline: + surface = str(row.get("work_surface") or row.get("container_app") or "").strip() + if not surface: + continue + bucket = grouped.setdefault(surface, { + "engaged_seconds": 0.0, + "foreground_seconds": 0.0, + "span_count": 0, + "days": set(), + "first_observed_at": None, + "last_observed_at": None, + "evidence_event_ids": [], + }) + bucket["engaged_seconds"] += float(row.get("engaged_seconds") or 0.0) + bucket["foreground_seconds"] += float(row.get("foreground_seconds") or 0.0) + bucket["span_count"] += 1 + started = str(row.get("started_at") or "") + if started: + bucket["days"].add(started[:10]) + if bucket["first_observed_at"] is None or started < bucket["first_observed_at"]: + bucket["first_observed_at"] = started + if bucket["last_observed_at"] is None or started > bucket["last_observed_at"]: + bucket["last_observed_at"] = started + bucket["evidence_event_ids"].extend(str(x) for x in (row.get("evidence_event_ids") or []) if x) + + findings: list[dict[str, Any]] = [] + for surface, value in grouped.items(): + engaged = round(float(value["engaged_seconds"]), 3) + active_days = len(value["days"]) + # Foreground presence alone is not engagement. In particular, a window + # left open across multiple days must not become a long-horizon finding. + # Multi-day surfaces need at least five engaged minutes total; a one-day + # surface needs at least thirty engaged minutes to be worth surfacing. + if engaged < 300 or (active_days < 2 and engaged < 1800): + continue + findings.append({ + "finding_id": _finding_id("surface_engagement", surface), + "finding_kind": "surface_engagement", + "surface": surface, + "engaged_seconds": engaged, + "foreground_seconds": round(float(value["foreground_seconds"]), 3), + "span_count": int(value["span_count"]), + "active_days": active_days, + "first_observed_at": value["first_observed_at"], + "last_observed_at": value["last_observed_at"], + "evidence_event_ids": list(dict.fromkeys(value["evidence_event_ids"]))[-12:], + "lookback_days": lookback_days, + # Engagement changes are material every 15 engaged minutes, or when a + # new active day appears. Tiny focus/heartbeat changes do not spam AI. + "material_version": f"d{active_days}:q{int(engaged // 900)}", + "factual_aggregate": True, + "task_inference_used": False, + "advice": False, + }) + return findings + + +def _transition_findings(timeline: list[dict[str, Any]], lookback_days: int) -> list[dict[str, Any]]: + by_session: dict[str, list[dict[str, Any]]] = defaultdict(list) + for row in timeline: + session_id = str(row.get("session_id") or "").strip() + if not session_id: + # Missing session identity means we cannot safely assert adjacency; + # grouping all such rows together would manufacture transitions. + continue + by_session[session_id].append(row) + + grouped: dict[tuple[str, str], dict[str, Any]] = {} + for rows in by_session.values(): + rows.sort(key=lambda row: str(row.get("started_at") or "")) + for left, right in zip(rows, rows[1:]): + source = str(left.get("work_surface") or left.get("container_app") or "").strip() + destination = str(right.get("work_surface") or right.get("container_app") or "").strip() + if not source or not destination or source == destination: + continue + left_end = _parse_ts(left.get("ended_at")) + right_start = _parse_ts(right.get("started_at")) + if left_end is None or right_start is None: + continue + gap = right_start - left_end + if gap < -1.0 or gap > 300.0: + continue + key = (source, destination) + value = grouped.setdefault(key, { + "count": 0, + "first_observed_at": None, + "last_observed_at": None, + "evidence_event_ids": [], + }) + value["count"] += 1 + observed = str(right.get("started_at") or "") + if value["first_observed_at"] is None or observed < value["first_observed_at"]: + value["first_observed_at"] = observed + if value["last_observed_at"] is None or observed > value["last_observed_at"]: + value["last_observed_at"] = observed + value["evidence_event_ids"].extend(str(x) for x in (left.get("evidence_event_ids") or []) if x) + value["evidence_event_ids"].extend(str(x) for x in (right.get("evidence_event_ids") or []) if x) + + findings: list[dict[str, Any]] = [] + for (source, destination), value in grouped.items(): + count = int(value["count"]) + if count < 3: + continue + findings.append({ + "finding_id": _finding_id("surface_transition", source, destination), + "finding_kind": "repeated_surface_transition", + "source_surface": source, + "destination_surface": destination, + "occurrence_count": count, + "first_observed_at": value["first_observed_at"], + "last_observed_at": value["last_observed_at"], + "evidence_event_ids": list(dict.fromkeys(value["evidence_event_ids"]))[-12:], + "lookback_days": lookback_days, + "material_version": f"n{count}", + "factual_aggregate": True, + "task_inference_used": False, + "advice": False, + }) + return findings + + +def _findings( + *, + snapshot_max_id: int, + snapshot_at: str, + lookback_days: int, +) -> tuple[list[dict[str, Any]], bool]: + raw_rows, truncated = _source_rows( + snapshot_max_id=snapshot_max_id, + snapshot_at=snapshot_at, + lookback_days=lookback_days, + ) + timeline = factual_context_timeline(_raw_events=raw_rows) + findings = [ + *_surface_findings(timeline, lookback_days), + *_transition_findings(timeline, lookback_days), + ] + findings.sort( + key=lambda item: ( + 0 if item.get("finding_kind") == "repeated_surface_transition" else 1, + -int(item.get("occurrence_count") or item.get("active_days") or 0), + str(item.get("finding_id") or ""), + ) + ) + return findings[:_MAX_TRACKED_FINDINGS], truncated or len(findings) > _MAX_TRACKED_FINDINGS + + +def context_pulse( + *, + cursor: str | None = None, + recent_limit: int = 100, + finding_limit: int = 20, + lookback_days: int = 30, +) -> dict[str, Any]: + """Return one incremental factual update and an opaque bookmark for the next check.""" + state = _decode_cursor(cursor) + bootstrap = not bool(state) + page_limit = max(1, min(int(recent_limit), 500)) + output_finding_limit = max(0, min(int(finding_limit), 50)) + + # Once a cursor exists it owns the lookback contract. A caller can start a new + # pulse (no cursor) to choose another lookback; it need not repeat the original + # value on every subsequent tool call. + if state: + lookback_days = max(7, min(int(state.get("lookback_days") or 30), 90)) + else: + lookback_days = max(7, min(int(lookback_days), 90)) + + last_event_id = max(0, int(state.get("last_event_id") or 0)) if state else 0 + pending_snapshot = max(0, int(state.get("pending_snapshot_max_id") or 0)) if state else 0 + pending_snapshot_at = str(state.get("pending_snapshot_at") or "") if state else "" + previous_versions = state.get("finding_versions") if state else {} + if not isinstance(previous_versions, dict): + raise ValueError("Invalid Context Pulse finding state") + previous_versions = { + str(key): str(value) + for key, value in list(previous_versions.items())[:_MAX_TRACKED_FINDINGS] + } + baseline_pending = bool(state.get("baseline_pending")) if state else True + + with connect() as conn: + current_max_id = int(conn.execute("SELECT COALESCE(MAX(id), 0) FROM events").fetchone()[0]) + snapshot_max_id = pending_snapshot or current_max_id + snapshot_at = pending_snapshot_at or datetime.now(timezone.utc).isoformat() + + if bootstrap: + with connect() as conn: + db_rows = conn.execute( + "SELECT * FROM events WHERE id <= ? ORDER BY id DESC LIMIT ?", + (snapshot_max_id, page_limit), + ).fetchall() + visible = list(reversed(db_rows)) + has_more = False + recent_mode = "bootstrap_tail" + else: + with connect() as conn: + db_rows = conn.execute( + "SELECT * FROM events WHERE id > ? AND id <= ? ORDER BY id ASC LIMIT ?", + (last_event_id, snapshot_max_id, page_limit + 1), + ).fetchall() + has_more = len(db_rows) > page_limit + visible = db_rows[:page_limit] + recent_mode = "since_last_pulse" + + recent_rows = [rich_evidence_row(dict(row), include_identity=False) for row in visible] + findings, findings_truncated = _findings( + snapshot_max_id=snapshot_max_id, + snapshot_at=snapshot_at, + lookback_days=lookback_days, + ) + current_versions = { + str(item["finding_id"]): str(item["material_version"]) + for item in findings + } + + # Keep only still-current delivered state. Crucially, a finding does not enter + # the cursor until it has actually been returned to the caller; a low + # finding_limit therefore cannot silently mark unseen findings as delivered. + next_versions = { + finding_id: version + for finding_id, version in previous_versions.items() + if finding_id in current_versions + } + candidates: list[tuple[dict[str, Any], str]] = [] + if output_finding_limit > 0: + for item in findings: + finding_id = str(item["finding_id"]) + version = str(item["material_version"]) + previous = previous_versions.get(finding_id) + if previous == version: + continue + status = "baseline" if baseline_pending and previous is None else ("new" if previous is None else "changed") + candidates.append((item, status)) + + selected = candidates[:output_finding_limit] if output_finding_limit > 0 else [] + changed: list[dict[str, Any]] = [] + for item, status in selected: + out = {key: value for key, value in item.items() if key != "material_version"} + out["status"] = status + changed.append(out) + next_versions[str(item["finding_id"])] = str(item["material_version"]) + + findings_has_more = output_finding_limit > 0 and len(candidates) > len(selected) + next_baseline_pending = bool(baseline_pending and findings_has_more) + + if bootstrap: + next_last_event_id = snapshot_max_id + elif has_more and visible: + next_last_event_id = int(visible[-1]["id"]) + else: + next_last_event_id = snapshot_max_id + + # Freeze the same snapshot while either the evidence page or finding page has + # more to deliver. New arrivals wait for the next completed pulse. + keep_snapshot = bool(has_more or findings_has_more) + next_cursor = _encode_cursor({ + "v": _CURSOR_VERSION, + "last_event_id": next_last_event_id, + "pending_snapshot_max_id": snapshot_max_id if keep_snapshot else 0, + "pending_snapshot_at": snapshot_at if keep_snapshot else "", + "lookback_days": lookback_days, + "baseline_pending": next_baseline_pending, + "finding_versions": dict(list(next_versions.items())[:_MAX_TRACKED_FINDINGS]), + }) + + return { + "bootstrap": bootstrap, + "recent_mode": recent_mode, + "recent_evidence": recent_rows, + "recent_returned": len(recent_rows), + "recent_has_more": has_more, + "snapshot_max_event_watermark": snapshot_max_id, + "snapshot_at": snapshot_at, + "findings": changed, + "findings_returned": len(changed), + "findings_has_more": findings_has_more, + "findings_lookback_days": lookback_days, + "findings_source_truncated": findings_truncated, + "next_cursor": next_cursor, + "cursor_contract": "opaque_bookmark; pass next_cursor unchanged to the next get_context_pulse call", + "canonical_evidence_tool": "get_workflow_trace", + "recent_evidence_contract": dict(RAW_RICH_EVIDENCE_CONTRACT), + "finding_contract": { + "factual_aggregates_only": True, + "task_inference_used": False, + "recommendations_or_advice": False, + "unchanged_findings_omitted_after_bootstrap": True, + "material_change": "new support, new active day, or another 15 minutes of engaged time depending on finding kind", + }, + } + + +__all__ = ["context_pulse"] diff --git a/server/context_pulse_routes.py b/server/context_pulse_routes.py new file mode 100644 index 00000000..eb759f57 --- /dev/null +++ b/server/context_pulse_routes.py @@ -0,0 +1,35 @@ +from __future__ import annotations + +from typing import Any + +from fastapi import HTTPException + +from .context_pulse import context_pulse +from .secure_app import app + + +@app.get("/v1/context-pulse") +def get_context_pulse( + cursor: str | None = None, + recent_limit: int = 100, + finding_limit: int = 20, + lookback_days: int = 30, +) -> dict[str, Any]: + """Return an incremental factual update for a connected AI client. + + The endpoint is read-only. The caller owns the opaque cursor and passes the + returned ``next_cursor`` back on the next check. AI disclosure/redaction is + still enforced by the established Context boundary middleware. + """ + try: + return context_pulse( + cursor=cursor, + recent_limit=recent_limit, + finding_limit=finding_limit, + lookback_days=lookback_days, + ) + except ValueError as exc: + raise HTTPException(status_code=400, detail=str(exc)) from exc + + +__all__ = ["app", "get_context_pulse"] diff --git a/server/enterprise_runner.py b/server/enterprise_runner.py index fb278727..e6fb40ed 100644 --- a/server/enterprise_runner.py +++ b/server/enterprise_runner.py @@ -26,6 +26,7 @@ def main() -> None: import server.v0571_polish # noqa: F401 import server.evidence_paging # noqa: F401 import server.work_profile_routes # noqa: F401 + import server.context_pulse_routes # noqa: F401 import server.browser_signal_routes # noqa: F401 import server.agent_dashboard_control_plane # noqa: F401 # Import last so its HTML middleware injects the privacy-safe dashboard loader diff --git a/tests/test_ai_context_mcp_v092.py b/tests/test_ai_context_mcp_v092.py index 02ca5a75..2ae8c102 100644 --- a/tests/test_ai_context_mcp_v092.py +++ b/tests/test_ai_context_mcp_v092.py @@ -28,6 +28,7 @@ # Context tools of the compact MCP server (experimental governance tools are off). CONTEXT_TOOL_ARGS = { "get_current_work_context": {}, + "get_context_pulse": {}, "search_work": {"query": "Contract"}, "get_workflow_trace": {"limit": 50}, "get_work_profile": {}, diff --git a/tests/test_compact_mcp_v087.py b/tests/test_compact_mcp_v087.py index 5d737fab..09d36ad3 100644 --- a/tests/test_compact_mcp_v087.py +++ b/tests/test_compact_mcp_v087.py @@ -12,6 +12,7 @@ DEFAULT_TOOLS = { "get_current_work_context", + "get_context_pulse", "search_work", "get_workflow_trace", "get_work_profile", diff --git a/tests/test_context_pulse_v094.py b/tests/test_context_pulse_v094.py new file mode 100644 index 00000000..8089cee2 --- /dev/null +++ b/tests/test_context_pulse_v094.py @@ -0,0 +1,214 @@ +from __future__ import annotations + +import base64 +import json +import os +import subprocess +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] + + +def _run(code: str, tmp_path: Path, timeout: int = 120) -> str: + env = os.environ.copy() + env.update({ + "WORKFLOW_OBSERVER_DATA": str(tmp_path / "data"), + "WORKFLOW_OBSERVER_AUTH_DIR": str(tmp_path / "auth"), + "WORKFLOW_OBSERVER_CONFIG": str(tmp_path / "config.json"), + "PYTHONPATH": str(ROOT), + }) + result = subprocess.run( + [sys.executable, "-c", code], + cwd=ROOT, + env=env, + text=True, + capture_output=True, + timeout=timeout, + ) + assert result.returncode == 0, f"stdout={result.stdout}\nstderr={result.stderr}" + return result.stdout + + +def test_pulse_bootstrap_findings_cursor_and_late_backdated_event(tmp_path): + _run(r''' +import base64, json +from datetime import datetime, timedelta, timezone +from server.db import init_db, insert_events +from server.context_pulse import context_pulse + +init_db() +now=datetime.now(timezone.utc) +events=[] +for day in range(3): + base=now-timedelta(days=day, hours=1) + for index,(app,offset) in enumerate([('Gmail',0),('ChatGPT',120),('Gmail',240)]): + events.append({ + 'event_id':f'base-{day}-{index}', + 'observed_at':(base+timedelta(seconds=offset)).isoformat(), + 'device_id':'d','session_id':f's{day}','app':app, + 'window_title':f'{app} work','event_type':'focus_span','duration_seconds':60, + 'metadata':{'activity':{'foreground_seconds':60,'engaged_seconds':50,'active_input_seconds':20}}, + }) +assert insert_events(events)==9 + +first=context_pulse(recent_limit=2, finding_limit=20, lookback_days=30) +assert first['bootstrap'] is True +assert first['recent_mode']=='bootstrap_tail' +assert first['recent_returned']==2 +assert first['recent_has_more'] is False +assert any(x['finding_kind']=='repeated_surface_transition' and x['status']=='baseline' for x in first['findings']), first +assert any(x['finding_kind']=='surface_engagement' and x['status']=='baseline' for x in first['findings']), first +assert all(x['factual_aggregate'] and not x['task_inference_used'] and not x['advice'] for x in first['findings']) + +# Cursor carries only watermarks and opaque finding IDs/versions, never captured labels. +raw=base64.urlsafe_b64decode(first['next_cursor']+'='*(-len(first['next_cursor'])%4)).decode() +assert 'Gmail' not in raw and 'ChatGPT' not in raw and 'work' not in raw, raw + +# Arrives later in storage but has an old observed_at: an arrival-id watermark must still return it. +late={ + 'event_id':'late-backdated','observed_at':(now-timedelta(days=10)).isoformat(), + 'device_id':'d','session_id':'late','app':'Notes','window_title':'Late note', + 'event_type':'focus_span','duration_seconds':5,'metadata':{'activity':{'engaged_seconds':5}}, +} +assert insert_events([late])==1 +second=context_pulse(cursor=first['next_cursor'], recent_limit=10, lookback_days=30) +assert second['bootstrap'] is False +assert [x['event_id'] for x in second['recent_evidence']]==['late-backdated'], second + +third=context_pulse(cursor=second['next_cursor'], recent_limit=10, lookback_days=30) +assert third['recent_evidence']==[] +assert third['findings']==[], third +''', tmp_path) + + +def test_pulse_freezes_multi_page_snapshot_and_does_not_skip_new_arrivals(tmp_path): + _run(r''' +from datetime import datetime, timedelta, timezone +from server.db import init_db, insert_events +from server.context_pulse import context_pulse + +init_db() +now=datetime.now(timezone.utc) +def event(eid, seconds): + return { + 'event_id':eid,'observed_at':(now+timedelta(seconds=seconds)).isoformat(), + 'device_id':'d','session_id':'s','app':'Editor','window_title':'Editor', + 'event_type':'ui_click','duration_seconds':0,'metadata':{'action':'click'}, + } + +assert insert_events([event('seed',0)])==1 +baseline=context_pulse(recent_limit=10) +assert insert_events([event('a',1),event('b',2),event('c',3)])==3 +page1=context_pulse(cursor=baseline['next_cursor'],recent_limit=1) +assert [x['event_id'] for x in page1['recent_evidence']]==['a'] +assert page1['recent_has_more'] is True +frozen=page1['snapshot_max_event_watermark'] + +# This arrival happens while the frozen snapshot is being paged. +assert insert_events([event('after-snapshot',4)])==1 +page2=context_pulse(cursor=page1['next_cursor'],recent_limit=1) +assert page2['snapshot_max_event_watermark']==frozen +assert [x['event_id'] for x in page2['recent_evidence']]==['b'] +page3=context_pulse(cursor=page2['next_cursor'],recent_limit=1) +assert [x['event_id'] for x in page3['recent_evidence']]==['c'] +assert page3['recent_has_more'] is False + +next_pulse=context_pulse(cursor=page3['next_cursor'],recent_limit=10) +assert [x['event_id'] for x in next_pulse['recent_evidence']]==['after-snapshot'] +''', tmp_path) + + +def test_finding_limit_drains_same_snapshot_without_marking_unseen_findings_seen(tmp_path): + _run(r''' +from datetime import datetime, timedelta, timezone +from server.db import init_db, insert_events +from server.context_pulse import context_pulse + +init_db() +now=datetime.now(timezone.utc) +events=[] +for day in range(3): + base=now-timedelta(days=day, hours=1) + for index,(app,offset) in enumerate([('Gmail',0),('ChatGPT',120),('Gmail',240)]): + events.append({ + 'event_id':f'limit-{day}-{index}', + 'observed_at':(base+timedelta(seconds=offset)).isoformat(), + 'device_id':'d','session_id':f'limit-s{day}','app':app, + 'window_title':app,'event_type':'focus_span','duration_seconds':120, + 'metadata':{'activity':{'foreground_seconds':120,'engaged_seconds':100,'active_input_seconds':40}}, + }) +assert insert_events(events)==9 + +page=context_pulse(recent_limit=20,finding_limit=1) +assert page['findings_returned']==1 and page['findings_has_more'] is True, page +frozen=page['snapshot_max_event_watermark'] +seen=set() + +while True: + assert page['snapshot_max_event_watermark']==frozen + assert page['findings_returned']==1, page + item=page['findings'][0] + assert item['status']=='baseline' + assert item['finding_id'] not in seen + seen.add(item['finding_id']) + if not page['findings_has_more']: + break + page=context_pulse(cursor=page['next_cursor'],recent_limit=20,finding_limit=1) + +assert len(seen)>=3, seen +settled=context_pulse(cursor=page['next_cursor'],recent_limit=20,finding_limit=1) +assert settled['findings']==[], settled +''', tmp_path) + + +def test_findings_do_not_turn_idle_foreground_or_missing_sessions_into_patterns(tmp_path): + _run(r''' +from server.context_pulse import _surface_findings, _transition_findings + +idle=[] +for day in range(3): + idle.append({ + 'session_id':f'idle-{day}', + 'started_at':f'2026-09-{20+day:02d}T10:00:00+00:00', + 'ended_at':f'2026-09-{20+day:02d}T11:00:00+00:00', + 'work_surface':'Claude', + 'container_app':'Claude', + 'foreground_seconds':3600, + 'engaged_seconds':0, + 'evidence_event_ids':[f'idle-{day}'], + }) +assert _surface_findings(idle,30)==[], _surface_findings(idle,30) + +# Without a session ID, chronology across rows is insufficient evidence that two +# surfaces were adjacent parts of one workflow. Do not manufacture a transition. +missing=[] +for i in range(8): + surface='A' if i%2==0 else 'B' + missing.append({ + 'session_id':'', + 'started_at':f'2026-09-25T10:{i:02d}:00+00:00', + 'ended_at':f'2026-09-25T10:{i:02d}:30+00:00', + 'work_surface':surface, + 'container_app':surface, + 'foreground_seconds':30, + 'engaged_seconds':20, + 'evidence_event_ids':[f'missing-{i}'], + }) +assert _transition_findings(missing,30)==[], _transition_findings(missing,30) +''', tmp_path) + + +def test_context_pulse_architecture_stays_observe_independent_and_factual(): + service=(ROOT/'server'/'context_pulse.py').read_text(encoding='utf-8') + assert 'hidden_frameworks' not in service + assert 'observation_active' not in service + assert 'candidate_tasks' not in service + assert 'factual_context_timeline' in service + assert 'task_inference_used' in service + assert 'recommendations_or_advice' in service + + +def test_context_pulse_route_is_registered_in_production_runner(): + runner=(ROOT/'server'/'enterprise_runner.py').read_text(encoding='utf-8') + assert 'import server.context_pulse_routes' in runner diff --git a/tests/test_mcpb_manifest_v087.py b/tests/test_mcpb_manifest_v087.py index cc44c489..07acfa26 100644 --- a/tests/test_mcpb_manifest_v087.py +++ b/tests/test_mcpb_manifest_v087.py @@ -7,6 +7,7 @@ ROOT = Path(__file__).resolve().parents[1] EXPECTED_COMPACT_TOOLS = { "get_current_work_context", + "get_context_pulse", "search_work", "get_workflow_trace", "get_work_profile",