"""Workflow-event and conversation timeline model (#637, Phase 1). Operators cannot browse a unified timeline of workflow events, decisions, tool calls, and handoffs: the evidence is scattered across control-plane events, Gitea canonical handoff comments, and local logs. This module defines one durable, versioned event schema and per-source adapters that normalise those scattered records into a single ``WorkflowEvent`` stream, plus a read-only query layer (filter by issue / PR / session, stable ordering, pagination) that the ``/api/v1/timeline`` route serves. Design rules honoured here: - **Read-only.** Sources are read; nothing is mutated. The control-plane database is opened through a ``mode=ro`` URI so a missing or unwritable DB degrades to a reason instead of creating directories or running migrations. - **Fail-soft per source.** An unavailable source degrades to a status with a reason rather than raising, and a source that could not run is never rendered as an empty-and-healthy timeline. - **Answerable filters only.** Each source declares which filter dimensions it can actually answer. A filter dimension no source that ran can carry is refused with an explicit reason rather than silently matching nothing: an empty page from an unanswerable filter reads to an operator as "no such activity", which is a different — and false — statement. - **Redaction at the boundary, fail closed.** Every free-text field (event messages, redacted tool arguments, decision/proof text) is run through the console redaction policy before it leaves this module, and *before* any structured value is derived from it — evidence references are extracted from redacted text, then independently revalidated before serialization. An unredactable value becomes the placeholder, and a value that cannot be proven safe is dropped — an unredacted payload is never emitted, and a generation error never drops raw data to a caller or a log. - **Stable ordering.** Events sort by ``(timestamp, source_rank, event_key)`` with a deterministic tiebreak, so pagination is stable across calls and events with equal or missing timestamps keep a fixed order. Non-goals (from the issue): no full chat replay, no mutation of historical events, no unredacted tool-argument storage. """ from __future__ import annotations import re import sqlite3 from dataclasses import dataclass, replace from datetime import datetime, timezone from typing import Any, Callable, Iterable import control_plane_db from webui import console_redaction # The schema is versioned so consumers can branch on shape. Bump on any # breaking change to WorkflowEvent's serialized form. TIMELINE_SCHEMA_VERSION = 1 # Known event sources and their deterministic ordering rank. When two events # carry the same timestamp, the source rank breaks the tie before the # per-source event key, so a control-plane event and a handoff comment minted # in the same second always sort in a fixed order. SOURCE_CONTROL_PLANE = "control_plane" SOURCE_GITEA_HANDOFF = "gitea_handoff" _SOURCE_RANK = { SOURCE_CONTROL_PLANE: 0, SOURCE_GITEA_HANDOFF: 1, } # The filter dimensions the query layer accepts. FILTER_ISSUE = "issue" FILTER_PR = "pr" FILTER_SESSION = "session" # Which dimensions each source can actually answer. This is a property of the # underlying records, not of the query code: the control-plane ``events`` table # is (event_id, work_item_id, event_type, message, created_at) and carries no # session identity at all, so no control-plane event can ever match a session # filter. A CTH handoff comment can declare its session as a field, so the # handoff source answers all three. Filtering on a dimension the surviving # sources cannot carry is refused in ``load_timeline`` rather than answered # with an empty page. _SOURCE_FILTER_SUPPORT: dict[str, tuple[str, ...]] = { SOURCE_CONTROL_PLANE: (FILTER_ISSUE, FILTER_PR), SOURCE_GITEA_HANDOFF: (FILTER_ISSUE, FILTER_PR, FILTER_SESSION), } # Why a source cannot answer a dimension, for the refusal reason an operator reads. _SOURCE_FILTER_LIMITS: dict[tuple[str, str], str] = { (SOURCE_CONTROL_PLANE, FILTER_SESSION): ( "control-plane events carry no session identity " "(the events table has no session column)" ), } # A timestamp far in the future so events with no parseable timestamp sort # last (after everything real) instead of first, without raising. _MISSING_TS_SORT = "9999-12-31T23:59:59Z" def _parse_ts(value: str | None) -> str | None: """Normalise a timestamp to ``...Z`` UTC, or None when unparseable.""" if not value: return None text = str(value).strip() if not text: return None candidate = text[:-1] + "+00:00" if text.endswith("Z") else text try: parsed = datetime.fromisoformat(candidate) except ValueError: return None if parsed.tzinfo is None: parsed = parsed.replace(tzinfo=timezone.utc) return parsed.astimezone(timezone.utc).replace(microsecond=0).isoformat().replace("+00:00", "Z") def _redact(value: Any) -> Any: """Redact a single free-text field, failing closed to the placeholder.""" if value is None: return None return console_redaction.redact_text(str(value)) @dataclass(frozen=True) class WorkflowEvent: """One normalised timeline event. Every field is optional except ``source``/``event_type``/``event_key`` because sources carry different subsets. The class is frozen so an adapted event is an immutable record; a consumer that needs a variant builds a new one rather than mutating history. """ source: str event_type: str event_key: str timestamp: str | None = None actor: str | None = None role: str | None = None issue_number: int | None = None pr_number: int | None = None session_id: str | None = None tool_name: str | None = None decision: str | None = None message: str | None = None correlation_id: str | None = None evidence_refs: tuple[str, ...] = () sensitive: bool = False def sort_key(self) -> tuple[str, int, str]: return ( self.timestamp or _MISSING_TS_SORT, _SOURCE_RANK.get(self.source, 99), self.event_key, ) def to_dict(self) -> dict[str, Any]: return { "source": self.source, "event_type": self.event_type, "event_key": self.event_key, "timestamp": self.timestamp, "actor": self.actor, "role": self.role, "issue_number": self.issue_number, "pr_number": self.pr_number, "session_id": self.session_id, "tool_name": self.tool_name, "decision": self.decision, "message": self.message, "correlation_id": self.correlation_id, "evidence_refs": list(self.evidence_refs), "sensitive": self.sensitive, } # --------------------------------------------------------------------------- # # Adapters — pure functions from a source's raw records to WorkflowEvents. # # Each is total: a malformed record is skipped, never raised on. # # --------------------------------------------------------------------------- # # Event types whose payload is treated as sensitive and always redaction-hard # (they can carry lease/session provenance or tool arguments). _SENSITIVE_EVENT_HINTS = ("lease", "capability", "token", "auth", "secret") # Reference tokens (issue/PR/comment ids) and SHAs parsed out of proof text. _EVIDENCE_REF_RE = re.compile(r"(?:#|PR\s*#?|issue\s*#?|comment\s*#?)(\d+)", re.IGNORECASE) # A commit reference is only recognised when the text *declares* it as one. # A bare lowercase hex run is not evidence of anything: at 40 characters it is # exactly the shape of a Gitea personal access token, and at 7 it also matches # ordinary words such as "defaced". Requiring an anchoring keyword keeps real # references ("commit abc1234", "at head a209756...", "base caaae9b6") usable # while refusing to lift an undeclared secret-shaped run out of free text. _SHA_RE = re.compile( r"(?i:\b(?:commit|sha|head|base|parent|revision|rev|merge[- ]base)\b[\s:=@#]*)" r"([0-9a-f]{7,40})\b" ) # Shapes a serialized evidence reference is allowed to take. Anything else is # dropped rather than emitted. _REF_ISSUE_SHAPE = re.compile(r"^#[0-9]{1,9}$") _REF_SHA_SHAPE = re.compile(r"^[0-9a-f]{7,40}$") # A long undelimited hex run with no declaring context is treated as credential # material wherever it appears, never as an identifier. _BARE_SECRET_SHAPE = re.compile(r"^[0-9a-f]{32,}$") def _kind_to_numbers(kind: str | None, number: int | None) -> tuple[int | None, int | None]: """Map a control-plane work-item (kind, number) to (issue_no, pr_no).""" if number is None: return (None, None) if kind == "pr": return (None, int(number)) if kind == "issue": return (int(number), None) return (None, None) def _correlation_for(kind: str | None, number: int | None) -> str | None: if number is None or kind not in ("issue", "pr"): return None return f"{kind}#{number}" def _extract_evidence_refs(*texts: str | None) -> tuple[str, ...]: """Extract issue/PR and declared-commit references from **redacted** text. Callers must pass text that has already been through :func:`_redact`; this function derives a structured field from its input, so extracting ahead of redaction would republish whatever redaction was about to remove. Every reference is revalidated by :func:`_validated_evidence_refs` before it is serialized. """ refs: list[str] = [] for text in texts: if not text: continue for match in _EVIDENCE_REF_RE.finditer(text): token = f"#{match.group(1)}" if token not in refs: refs.append(token) for match in _SHA_RE.finditer(text): token = match.group(1) if token not in refs: refs.append(token) return tuple(refs) def _validated_evidence_refs(refs: Iterable[str]) -> tuple[tuple[str, ...], bool]: """Independently revalidate references immediately before serialization. Extraction is not trusted on its own. A reference survives only when it has a known reference shape and is unchanged by a second redaction pass — a value the redaction policy would alter is credential material that must not be emitted as a structured field. A full 40-character SHA stays usable because extraction only accepts a hex run the source text explicitly declared as a commit. Returns ``(safe_refs, dropped_any)``; ``dropped_any`` marks the event sensitive so the drop is visible rather than silent. """ safe: list[str] = [] dropped = False for ref in refs or (): try: token = str(ref).strip() if not token: continue recognised = bool(_REF_ISSUE_SHAPE.match(token) or _REF_SHA_SHAPE.match(token)) if not recognised: dropped = True continue if _redact(token) != token: dropped = True continue if token not in safe: safe.append(token) except Exception: # Fail closed: a reference that cannot be proven safe is dropped. dropped = True continue return (tuple(safe), dropped) def _safe_session_id(value: Any) -> str | None: """Return a session identifier only when it is safe to emit. The value is authoritative source data — a session the record names for itself — but it is still free text. It is dropped when redaction alters it or when it is a bare secret-shaped hex run, so a credential parked in a session field can never reach the payload or be echoed back by a filter. """ if value is None: return None text = str(value).strip() if not text: return None if _BARE_SECRET_SHAPE.match(text): return None return text if _redact(text) == text else None def adapt_cp_events(rows: Iterable[dict[str, Any]]) -> list[WorkflowEvent]: """Adapt control-plane ``events`` rows (joined to work_items) into events. Each row is expected to carry ``event_id``, ``event_type``, ``message``, ``created_at`` and the joined work-item ``kind``/``number``. Rows missing an id or type are skipped so a partially written table never raises. """ events: list[WorkflowEvent] = [] for row in rows or []: try: event_id = row.get("event_id") event_type = (row.get("event_type") or "").strip() if event_id is None or not event_type: continue kind = row.get("kind") number = row.get("number") issue_no, pr_no = _kind_to_numbers(kind, number) sensitive = any(hint in event_type.lower() for hint in _SENSITIVE_EVENT_HINTS) events.append( WorkflowEvent( source=SOURCE_CONTROL_PLANE, event_type=event_type, event_key=f"cp:{event_id}", timestamp=_parse_ts(row.get("created_at")), issue_number=issue_no, pr_number=pr_no, # No session_id: the control-plane events table is # (event_id, work_item_id, event_type, message, created_at) # and records no session. Inventing one from the work item # or the message text would be a guess, so this source # declares the session dimension unsupported instead # (_SOURCE_FILTER_SUPPORT) and the query layer refuses a # session filter it cannot honestly answer. message=_redact(row.get("message")), correlation_id=_correlation_for(kind, number), sensitive=sensitive, ) ) except Exception: # A single malformed row must not sink the whole adaptation. continue return events def adapt_cth_comments( comments: Iterable[dict[str, Any]], *, kind: str, number: int, ) -> list[WorkflowEvent]: """Adapt Gitea Canonical Thread Handoff (CTH) comments into events. Only comments that parse as a CTH (``canonical_thread_handoff.parse_cth_comment``) become events; ordinary comments are ignored. ``kind``/``number`` scope the events to the issue or PR the comments belong to. """ # Imported lazily so this module has no import-time dependency on the # handoff parser when only the control-plane adapter is used. from canonical_thread_handoff import parse_cth_comment issue_no, pr_no = _kind_to_numbers(kind, number) correlation = _correlation_for(kind, number) events: list[WorkflowEvent] = [] for comment in comments or []: try: body = comment.get("body") or "" parsed = parse_cth_comment(body) if not parsed: continue fields = parsed.get("fields") or {} cth_type = parsed.get("cth_type") or "handoff" comment_id = comment.get("id") # Redaction runs first, and every derived value is taken from the # redacted text — deriving evidence refs from the raw proof would # re-emit exactly what redaction was about to remove. decision = _redact(fields.get("decision")) proof = _redact(fields.get("proof")) next_action = _redact(fields.get("next action")) refs, refs_dropped = _validated_evidence_refs( _extract_evidence_refs(proof, decision) ) events.append( WorkflowEvent( source=SOURCE_GITEA_HANDOFF, event_type=f"handoff:{cth_type}", event_key=f"cth:{kind}:{number}:{comment_id}", timestamp=_parse_ts(comment.get("created_at")), actor=_redact((comment.get("user") or {}).get("login")), role=_redact(fields.get("next owner")), issue_number=issue_no, pr_number=pr_no, # A CTH names its own session when the producer records one; # it is read from that declared field, never inferred from # unrelated text. session_id=_safe_session_id(fields.get("session")), decision=decision, message=next_action or _redact(fields.get("status")), correlation_id=correlation, evidence_refs=refs, sensitive=refs_dropped, ) ) except Exception: continue return events # --------------------------------------------------------------------------- # # Read-only control-plane event source. # # --------------------------------------------------------------------------- # _CP_EVENTS_QUERY = """ SELECT e.event_id AS event_id, e.event_type AS event_type, e.message AS message, e.created_at AS created_at, w.kind AS kind, w.number AS number FROM events e JOIN work_items w ON e.work_item_id = w.work_item_id WHERE w.remote = ? AND w.org = ? AND w.repo = ? """ @dataclass(frozen=True) class SourceStatus: """Fail-soft status for one timeline source. ``supported_filters`` states which filter dimensions this source's records can carry; ``unsupported_filters`` names the requested dimensions it cannot, so an operator can see *why* a source contributed nothing rather than being left to read an empty list as an absence of activity. """ name: str ok: bool reason: str | None = None count: int = 0 supported_filters: tuple[str, ...] = () unsupported_filters: tuple[str, ...] = () def to_dict(self) -> dict[str, Any]: return { "name": self.name, "ok": self.ok, "reason": self.reason, "count": self.count, "supported_filters": list(self.supported_filters), "unsupported_filters": list(self.unsupported_filters), } def _cp_status(*, ok: bool, reason: str | None = None, count: int = 0) -> SourceStatus: return SourceStatus( SOURCE_CONTROL_PLANE, ok=ok, reason=reason, count=count, supported_filters=_SOURCE_FILTER_SUPPORT[SOURCE_CONTROL_PLANE], ) def _handoff_status(*, ok: bool, reason: str | None = None, count: int = 0) -> SourceStatus: return SourceStatus( SOURCE_GITEA_HANDOFF, ok=ok, reason=reason, count=count, supported_filters=_SOURCE_FILTER_SUPPORT[SOURCE_GITEA_HANDOFF], ) def read_cp_events( *, remote: str, org: str, repo: str, db_path: str | None = None, ) -> tuple[list[WorkflowEvent], SourceStatus]: """Read scoped control-plane events read-only. Never creates the DB. Opens the SQLite file through a ``mode=ro`` URI: a health/timeline read must never create directories or run the schema migration that ``ControlPlaneDB()`` performs on construction. A missing or unreadable DB degrades to a status with a reason. """ path = (db_path or control_plane_db.default_db_path()).strip() conn: sqlite3.Connection | None = None try: conn = sqlite3.connect(f"file:{path}?mode=ro", uri=True) conn.row_factory = sqlite3.Row cursor = conn.execute(_CP_EVENTS_QUERY, (remote, org, repo)) rows = [dict(r) for r in cursor.fetchall()] except sqlite3.OperationalError as exc: return ([], _cp_status(ok=False, reason=f"control-plane DB unavailable: {exc}")) except sqlite3.Error as exc: return ([], _cp_status(ok=False, reason=f"control-plane read failed: {exc}")) finally: if conn is not None: conn.close() events = adapt_cp_events(rows) return (events, _cp_status(ok=True, count=len(events))) # --------------------------------------------------------------------------- # # Filter, sort, paginate. # # --------------------------------------------------------------------------- # def filter_events( events: Iterable[WorkflowEvent], *, issue: int | None = None, pr: int | None = None, session: str | None = None, ) -> list[WorkflowEvent]: """Filter events by issue number, PR number, and/or session id. Filters are conjunctive. A filter that names a dimension an event does not carry excludes that event (an issue filter excludes PR-only events). """ out: list[WorkflowEvent] = [] for ev in events: if issue is not None and ev.issue_number != issue: continue if pr is not None and ev.pr_number != pr: continue if session is not None and ev.session_id != session: continue out.append(ev) return out def sort_events(events: Iterable[WorkflowEvent]) -> list[WorkflowEvent]: """Return events in stable timeline order (ascending).""" return sorted(events, key=lambda ev: ev.sort_key()) @dataclass(frozen=True) class TimelinePage: """One page of the sorted, filtered timeline.""" events: tuple[WorkflowEvent, ...] total: int limit: int offset: int @property def next_offset(self) -> int | None: nxt = self.offset + len(self.events) return nxt if nxt < self.total else None def to_dict(self) -> dict[str, Any]: return { "events": [ev.to_dict() for ev in self.events], "pagination": { "total": self.total, "limit": self.limit, "offset": self.offset, "returned": len(self.events), "next_offset": self.next_offset, "has_more": self.next_offset is not None, }, } _MAX_LIMIT = 500 _DEFAULT_LIMIT = 50 def _coerce_bounds(limit: int | None, offset: int | None) -> tuple[int, int]: try: lim = int(limit) if limit is not None else _DEFAULT_LIMIT except (TypeError, ValueError): lim = _DEFAULT_LIMIT try: off = int(offset) if offset is not None else 0 except (TypeError, ValueError): off = 0 lim = max(1, min(lim, _MAX_LIMIT)) off = max(0, off) return (lim, off) def paginate(events: list[WorkflowEvent], *, limit: int | None, offset: int | None) -> TimelinePage: lim, off = _coerce_bounds(limit, offset) window = events[off : off + lim] return TimelinePage(events=tuple(window), total=len(events), limit=lim, offset=off) # --------------------------------------------------------------------------- # # Composition — load_timeline aggregates all sources, fail-soft. # # --------------------------------------------------------------------------- # # A comment source is a callable that, given (kind, number), returns the raw # Gitea comment list for that issue/PR. The route supplies a live fail-soft # fetcher; tests supply a fixture. When None, the handoff source is reported as # not-run (never silently empty-and-healthy). CommentSource = Callable[[str, int], list[dict[str, Any]]] @dataclass(frozen=True) class TimelineSnapshot: """One answered timeline query. ``ok`` is False when the query could not be answered as asked — currently when a requested filter dimension no surviving source can carry was supplied. The page is then empty *and* the snapshot says so, because an ``ok`` empty page is a claim that no such activity exists. """ schema_version: int remote: str org: str repo: str filters: dict[str, Any] page: TimelinePage sources: tuple[SourceStatus, ...] ok: bool = True error: dict[str, Any] | None = None def to_dict(self) -> dict[str, Any]: return { "ok": self.ok, "error": self.error, "schema_version": self.schema_version, "scope": {"remote": self.remote, "org": self.org, "repo": self.repo}, "filters": self.filters, "sources": [s.to_dict() for s in self.sources], **self.page.to_dict(), } def _unanswerable_reasons( statuses: Iterable[SourceStatus], unanswerable: Iterable[str] ) -> list[dict[str, str]]: """Explain, per source, why each unanswerable dimension went unanswered.""" out: list[dict[str, str]] = [] for status in statuses: for dim in unanswerable: if dim not in status.supported_filters: reason = _SOURCE_FILTER_LIMITS.get( (status.name, dim), f"this source's records carry no {dim} identity" ) elif not status.ok: reason = ( f"this source can carry {dim} but did not run: " f"{status.reason or 'unavailable'}" ) else: continue out.append({"source": status.name, "filter": dim, "reason": reason}) return out def load_timeline( *, remote: str, org: str, repo: str, issue: int | None = None, pr: int | None = None, session: str | None = None, limit: int | None = None, offset: int | None = None, db_path: str | None = None, comment_source: CommentSource | None = None, ) -> TimelineSnapshot: """Aggregate every timeline source into one filtered, paginated snapshot. Sources are read independently and fail soft: an unavailable source contributes a ``SourceStatus`` with ``ok=False`` and a reason, and never collapses the whole timeline. The handoff source only runs when a specific issue or PR is requested (a handoff comment belongs to one thread) and a ``comment_source`` is available; otherwise it is reported as ``not run`` rather than as an empty-and-healthy source. A filter dimension that no surviving source can carry — a ``session`` filter when the only source that ran is the control plane, whose events record no session — is refused with ``ok=False`` and a structured error instead of being answered with an empty page. """ all_events: list[WorkflowEvent] = [] statuses: list[SourceStatus] = [] cp_events, cp_status = read_cp_events(remote=remote, org=org, repo=repo, db_path=db_path) all_events.extend(cp_events) statuses.append(cp_status) # Gitea handoff comments are thread-scoped: only fetch when the caller # narrowed to one issue or PR, and only when a source was provided. handoff_target: tuple[str, int] | None = None if pr is not None: handoff_target = ("pr", pr) elif issue is not None: handoff_target = ("issue", issue) if handoff_target is None: statuses.append( _handoff_status( ok=False, reason="not run: handoff comments are thread-scoped; filter by issue or pr to include them", ) ) elif comment_source is None: statuses.append( _handoff_status( ok=False, reason="not run: no comment source configured for this timeline read", ) ) else: kind, number = handoff_target try: comments = comment_source(kind, number) or [] handoff_events = adapt_cth_comments(comments, kind=kind, number=number) all_events.extend(handoff_events) statuses.append(_handoff_status(ok=True, count=len(handoff_events))) except Exception as exc: # fail soft: a fetch/parse error degrades this source only statuses.append(_handoff_status(ok=False, reason=f"handoff source failed: {exc}")) requested = tuple( name for name, value in ((FILTER_ISSUE, issue), (FILTER_PR, pr), (FILTER_SESSION, session)) if value is not None ) statuses = [ replace( status, unsupported_filters=tuple( dim for dim in requested if dim not in status.supported_filters ), ) for status in statuses ] filters = {"issue": issue, "pr": pr, "session": session} # A dimension is answerable only if a source that actually ran can carry it. # If none can, refuse: an empty page would assert "no such activity", which # is a claim this timeline is not in a position to make. answerable: set[str] = set() for status in statuses: if status.ok: answerable.update(status.supported_filters) unanswerable = tuple(dim for dim in requested if dim not in answerable) if unanswerable: return TimelineSnapshot( schema_version=TIMELINE_SCHEMA_VERSION, remote=remote, org=org, repo=repo, filters=filters, page=paginate([], limit=limit, offset=offset), sources=tuple(statuses), ok=False, error={ "code": "filter_not_supported", "unsupported_filters": list(unanswerable), "detail": ( "no timeline source that ran can answer " + ", ".join(f"'{dim}'" for dim in unanswerable) + "; the result is refused rather than returned empty" ), "sources": _unanswerable_reasons(statuses, unanswerable), }, ) filtered = filter_events(all_events, issue=issue, pr=pr, session=session) ordered = sort_events(filtered) page = paginate(ordered, limit=limit, offset=offset) return TimelineSnapshot( schema_version=TIMELINE_SCHEMA_VERSION, remote=remote, org=org, repo=repo, filters=filters, page=page, sources=tuple(statuses), ) def snapshot_to_dict(snapshot: TimelineSnapshot) -> dict[str, Any]: return snapshot.to_dict()