"You" was the else-branch of session_label: anything the board could not attribute to an agent it attributed to the person. Every session read back from disk was one of those, because the agent id lived only in the session registry and never reached the persisted events — so past agent runs came back wearing the human's label, carrying their own closing reports underneath it. Identity is now a small whole file beside each event log (state/sessions/<id>.who.json): agent id, the agent's name, the model it rode, and the task. A file rather than a key on the events, because the logs are append-only JSONL whose first line every reader takes for an event — and because the name and the model are nowhere in the stream, so this is the only thing a restart can read them back from. It is rewritten only when what the board knows changes, which also covers an agent id that arrives on a later event. load_disk_sessions() reads it back, and the label now has three registers instead of two: the agent's name (persisted, so a restart no longer costs it), "You" only for a session positively recorded as carrying no agent, and a neutral "Session · <id>" for a log written before any of this was recorded. Old logs are not retro-attributed in either direction. agentFor() in board.html falls back to the persisted identity when this board no longer holds the live record, so a replayed agent session wears its model chip from what was written rather than from what happens to be in memory. What depends on liveness (Hold, the worktree branch) finds nothing there and stays silent, as before. tests/test_session_identity.py drives the real ingest → persist → reload path and the page's own chip functions in node. Co-Authored-By: Claude Opus 5 <noreply@anthropic.com>
138 lines
4.9 KiB
Python
138 lines
4.9 KiB
Python
"""Shared in-memory state, persistence of event logs, and the SSE fan-out.
|
|
|
|
All cross-thread registries live here, guarded by LOCK where they are
|
|
mutated from several threads. Modules communicate through this state rather
|
|
than importing each other's internals.
|
|
"""
|
|
|
|
from __future__ import annotations
|
|
|
|
import json
|
|
import queue
|
|
import threading
|
|
import time
|
|
from pathlib import Path
|
|
|
|
import config
|
|
|
|
LOCK = threading.Lock()
|
|
CLIENTS: set[queue.Queue] = set() # one queue per open SSE connection
|
|
SESSIONS: dict[str, dict] = {} # session_id -> meta
|
|
EVENTS: dict[str, list[dict]] = {} # session_id -> slim events
|
|
BOARD_EVENTS: list[dict] = [] # moves + agent lifecycle
|
|
AGENTS: dict[str, dict] = {} # agent_id -> launch record
|
|
EXPECTED_MOVES: dict[tuple[str, str], tuple[str, float]] = {} # (file, to) -> (actor, ts)
|
|
COMMIT_HOOKS: list = [] # run after a board-made task commit
|
|
|
|
# The port actually being served; board.py sets it from --port at startup so
|
|
# launched agents know where to report events.
|
|
serve_port = config.PORT
|
|
|
|
# The last card archived through this board — the scope of the ⌘Z undo.
|
|
LAST_ARCHIVED: dict | None = None
|
|
|
|
# A session's identity sidecar, beside its <sid>.jsonl event log. Not a
|
|
# `.jsonl` itself, so no reader that globs the event logs picks it up.
|
|
IDENTITY_SUFFIX = ".who.json"
|
|
|
|
|
|
def broadcast(payload: dict) -> None:
|
|
msg = json.dumps(payload)
|
|
with LOCK:
|
|
clients = list(CLIENTS)
|
|
for q in clients:
|
|
try:
|
|
q.put_nowait(msg)
|
|
except queue.Full:
|
|
pass
|
|
|
|
|
|
def persist(name: str, record: dict) -> None:
|
|
try:
|
|
config.SESSIONS_DIR.mkdir(parents=True, exist_ok=True)
|
|
with (config.SESSIONS_DIR / name).open("a", encoding="utf-8") as fh:
|
|
fh.write(json.dumps(record) + "\n")
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
def _session_file(name: str) -> Path | None:
|
|
if "/" in name or ".." in name:
|
|
return None
|
|
return config.SESSIONS_DIR / name
|
|
|
|
|
|
def persist_identity(sid: str, identity: dict) -> None:
|
|
"""Who a session belonged to, written beside its event log.
|
|
|
|
A whole small file of its own rather than a key on the events: the logs
|
|
are append-only JSONL whose first line every reader takes for an event,
|
|
and identity is a property of the session, not of anything that happened
|
|
inside it. The agent's *name* and *model* live only in board memory, so
|
|
this file is the only thing a restart can read them back from.
|
|
|
|
Rewritten whenever what we know changes — an agent id that arrives on a
|
|
later event, a name that was not registered yet when the first event
|
|
landed. The file is written whole, so the last write is simply the truth.
|
|
"""
|
|
path = _session_file(f"{sid}{IDENTITY_SUFFIX}")
|
|
if path is None:
|
|
return
|
|
try:
|
|
config.SESSIONS_DIR.mkdir(parents=True, exist_ok=True)
|
|
path.write_text(json.dumps(identity), encoding="utf-8")
|
|
except OSError:
|
|
pass
|
|
|
|
|
|
def read_identity(sid: str) -> dict | None:
|
|
"""The persisted identity, or None when nothing was ever recorded — the
|
|
difference between "this session was the person" and "we do not know",
|
|
which is exactly what the label must not blur."""
|
|
path = _session_file(f"{sid}{IDENTITY_SUFFIX}")
|
|
if path is None or not path.is_file():
|
|
return None
|
|
try:
|
|
data = json.loads(path.read_text(encoding="utf-8"))
|
|
except (OSError, json.JSONDecodeError, ValueError):
|
|
return None
|
|
return data if isinstance(data, dict) else None
|
|
|
|
|
|
def record_board_event(event: dict) -> None:
|
|
event["ts"] = time.time()
|
|
with LOCK:
|
|
BOARD_EVENTS.append(event)
|
|
del BOARD_EVENTS[:-config.BOARD_EVENTS_CAP]
|
|
persist("board.jsonl", event)
|
|
broadcast({"type": "board_event", "event": event})
|
|
|
|
|
|
def task_committed(filename: str) -> None:
|
|
"""A board-made move committed itself. Registered hooks turn that into
|
|
whatever else should follow — sync.py's push, when the gate is on. The
|
|
hook is a registry rather than an import so taskfiles stays to the left
|
|
of everything that reacts to it; a hook that raises must never break a
|
|
move that has already happened on disk."""
|
|
for hook in list(COMMIT_HOOKS):
|
|
try:
|
|
hook(filename)
|
|
except Exception: # noqa: BLE001 — the move is done; nothing may undo it
|
|
pass
|
|
|
|
|
|
def expect_move(filename: str, target: str, actor: str) -> None:
|
|
"""Tell the watcher who is about to move a file so it can attribute it."""
|
|
with LOCK:
|
|
EXPECTED_MOVES[(filename, target)] = (actor, time.time())
|
|
|
|
|
|
def claim_expected(filename: str, target: str) -> str:
|
|
with LOCK:
|
|
actor_ts = EXPECTED_MOVES.pop((filename, target), None)
|
|
# forget stale expectations while we're here
|
|
cutoff = time.time() - 30
|
|
for key in [k for k, (_, ts) in EXPECTED_MOVES.items() if ts < cutoff]:
|
|
EXPECTED_MOVES.pop(key, None)
|
|
return actor_ts[0] if actor_ts else "disk"
|