Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 1 addition & 1 deletion VERSION
Original file line number Diff line number Diff line change
@@ -1 +1 @@
0.97.0
0.98.0
30 changes: 27 additions & 3 deletions adapters/_agent_client.py
Original file line number Diff line number Diff line change
Expand Up @@ -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
Expand Down Expand Up @@ -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")
Expand All @@ -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",
)
Expand All @@ -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
2 changes: 2 additions & 0 deletions adapters/claude_code_hook.py
Original file line number Diff line number Diff line change
Expand Up @@ -20,6 +20,8 @@
SUPPORTED_EVENTS = [
"SessionStart",
"SessionEnd",
"UserPromptSubmit",
"Stop",
"PostToolUse",
"PostToolUseFailure",
"PermissionRequest",
Expand Down
41 changes: 41 additions & 0 deletions collector/boundaries.py
Original file line number Diff line number Diff line change
@@ -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
73 changes: 73 additions & 0 deletions collector/instance_lock.py
Original file line number Diff line number Diff line change
@@ -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()
4 changes: 4 additions & 0 deletions collector/interactions.py
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
Loading
Loading