diff --git a/tests/test_coordinator_adapter.py b/tests/test_coordinator_adapter.py index ff71e188..4c2e2fcd 100644 --- a/tests/test_coordinator_adapter.py +++ b/tests/test_coordinator_adapter.py @@ -85,6 +85,44 @@ def test_emit_created_calls_collector_with_coord_fields() -> None: ) +def test_emit_created_seeds_resolved_display_name(tmp_path: Any) -> None: + """The collector seed uses the resolved display name (alias > title > + name), not the synthetic ``ws.name``. A coordinator carrying a + persisted LLM auto-title (written by ``update_workstream_title``) then + shows that title in the live cluster tree instead of reverting to + ``ws-xxxx``. Regression guard for the adapter half of the + coordinator-title-persistence fix — the server-side ``_coordinator_rows`` + half is pinned in test_coordinator_endpoints.py.""" + from turnstone.core.storage import init_storage, reset_storage + + reset_storage() + backend = init_storage("sqlite", path=str(tmp_path / "adapter.db"), run_migrations=False) + try: + # Titled coordinator → the title surfaces over the placeholder name. + backend.register_workstream( + "coord-1", + node_id="console", + user_id="u1", + name="ws-c0c0", + kind=WorkstreamKind.COORDINATOR, + ) + backend.update_workstream_title("coord-1", "Investigate the title bug") + adapter, collector = _make_adapter() + adapter.emit_created(_make_ws(name="ws-c0c0")) + assert ( + collector.emit_console_ws_created.call_args.kwargs["name"] + == "Investigate the title bug" + ) + + # A user alias outranks the auto-title (alias > title > name). + assert backend.set_workstream_alias("coord-1", "Pinned name") + collector.emit_console_ws_created.reset_mock() + adapter._fanout_console_ws_created(_make_ws(name="ws-c0c0")) + assert collector.emit_console_ws_created.call_args.kwargs["name"] == "Pinned name" + finally: + reset_storage() + + def test_emit_state_calls_collector_state() -> None: """Post-rich-payload, emit_state passes tokens / context_ratio / activity / activity_state / content kwargs read from ws.ui's diff --git a/tests/test_coordinator_endpoints.py b/tests/test_coordinator_endpoints.py index 61362198..2439eb99 100644 --- a/tests/test_coordinator_endpoints.py +++ b/tests/test_coordinator_endpoints.py @@ -2433,6 +2433,54 @@ def test_coordinator_rows_persisted_cluster_wide(storage): assert {r["name"] for r in rows} == {"alice-closed", "bob-closed", "orphan-closed"} +def test_coordinator_rows_surface_persisted_title(storage): + """Regression for the coordinator-title-persistence bug. + + The LLM auto-title (``update_workstream_title``) and the user alias + (``set_workstream_alias``) live only in ``workstreams.title`` / + ``workstreams.alias``. ``_coordinator_rows`` must resolve the + display name ``alias > title > name`` from the persisted row for BOTH + lanes — the in-memory ``ws.name`` is the synthetic ``ws-xxxx`` + placeholder. Before the fix the read path hardcoded ``title=""`` and + used ``ws.name`` / the ``name`` column, so a generated title was + written but never read back: it reverted to ``ws-xxxx`` on every + dashboard refresh.""" + from turnstone.console.server import _coordinator_rows + from turnstone.core.workstream import WorkstreamKind + + mgr = _build_mgr(storage) + + # In-memory lane: a LIVE coordinator titled after creation. The + # manager assigned the placeholder ``ws.name``; the title is in the DB. + live = mgr.create(user_id="alice", name="ws-abcd") + storage.update_workstream_title(live.id, "Refactor the auth layer") + + # Persisted lane: a closed coordinator (evicted from the manager) + # carrying BOTH a title and a user alias — the alias must win. + storage.register_workstream( + "f" * 32, + node_id="console", + user_id="bob", + name="ws-f0f0", + state="closed", + kind=WorkstreamKind.COORDINATOR, + parent_ws_id=None, + ) + storage.update_workstream_title("f" * 32, "auto-generated title") + assert storage.set_workstream_alias("f" * 32, "Bob's pinned name") + + request = _persisted_rows_request(storage, mgr, "alice", frozenset({"read"})) + rows = {r["id"]: r for r in _coordinator_rows(request)} + + live_row = rows[live.id] + assert live_row["name"] == "Refactor the auth layer" + assert live_row["title"] == "Refactor the auth layer" + + closed_row = rows["f" * 32] + assert closed_row["name"] == "Bob's pinned name" # alias > title > name + assert closed_row["title"] == "auto-generated title" + + # --------------------------------------------------------------------------- # Stage 2 P1.5 — coord attachment surface parity with interactive # --------------------------------------------------------------------------- diff --git a/tests/test_session.py b/tests/test_session.py index 626f3bdf..5d3f2c42 100644 --- a/tests/test_session.py +++ b/tests/test_session.py @@ -150,6 +150,27 @@ def _send_with_mocks(session, responses, mock_execute, **extra_patches): yield save_msg +def _capturing_thread_cls(): + """Return a no-op ``threading.Thread`` stand-in plus the list it records + each constructed thread's ``target`` into. + + Patched over ``session.threading.Thread`` so a test can assert WHICH + callable was scheduled (e.g. ``_generate_title``) without the thread + actually running — ``start()`` is a no-op, so no background LLM call + fires. + """ + started: list = [] + + class _CaptureThread: + def __init__(self, *a, target=None, **kw): + started.append(target) + + def start(self): + pass + + return _CaptureThread, started + + def _user_pending(session) -> list[tuple[str, str]]: """Return user-channel queued nudges as ``(type, text)`` tuples. @@ -1010,6 +1031,66 @@ class TestTitleRetry: # Restore for cleanup session._ws_id = original_ws_id + def test_title_fires_after_send_not_after_tool_free_turn(self, tmp_db): + """Auto-title fires right after the user turn is recorded, BEFORE + tools run — it no longer waits for a tool-call-free assistant + turn. Coordinators spend nearly every turn in tool calls and may + never reach that terminal text turn, so the old end-of-turn + trigger almost never fired for them (the timing half of the + coordinator-title bug).""" + session = _make_session() + assert session._title_generated is False + # The assistant's opening turn is ALL tool calls — under the old + # trigger no title would generate until a later text-only turn. + responses = [ + { + "role": "assistant", + "content": "working", + "tool_calls": [ + { + "id": "c1", + "type": "function", + "function": {"name": "echo", "arguments": "{}"}, + } + ], + }, + {"role": "assistant", "content": "done"}, + ] + capture_cls, started = _capturing_thread_cls() + + def mock_execute(_tool_calls): + # The title must already be scheduled by the time tools run. + assert session._title_generated is True + return [("c1", "ok")], None + + with ( + _send_with_mocks(session, responses, mock_execute), + patch("turnstone.core.session.threading.Thread", capture_cls), + ): + session.send("refactor the auth layer") + + assert session._title_generated is True + assert session._generate_title in started + + def test_title_not_generated_for_blank_or_wake_send(self, tmp_db): + """Blank input and synthetic wake sends don't burn the one-shot + auto-title — ``_generate_title`` needs first-user-message text, + and a wake carries none.""" + capture_cls, started = _capturing_thread_cls() + + def mock_execute(_tool_calls): + return [], None + + for user_input, kwargs in ((" ", {}), ("a real message", {"from_wake": True})): + session = _make_session() + with ( + _send_with_mocks(session, [{"role": "assistant", "content": "ok"}], mock_execute), + patch("turnstone.core.session.threading.Thread", capture_cls), + ): + session.send(user_input, **kwargs) + assert session._generate_title not in started + assert session._title_generated is False + class TestLiveConfigUpdate: """ConfigStore-backed sessions pick up settings changes at point-of-use.""" diff --git a/turnstone/console/coordinator_adapter.py b/turnstone/console/coordinator_adapter.py index 7ebc492d..52e6b310 100644 --- a/turnstone/console/coordinator_adapter.py +++ b/turnstone/console/coordinator_adapter.py @@ -23,6 +23,7 @@ from turnstone.core.child_event_bus import ChildEventBus from turnstone.core.child_source import ClusterChildSource from turnstone.core.children_registry import ChildrenRegistry from turnstone.core.log import get_logger +from turnstone.core.memory import get_workstream_display_name, get_workstream_display_names from turnstone.core.workstream import Workstream, WorkstreamKind, WorkstreamState if TYPE_CHECKING: @@ -37,6 +38,19 @@ if TYPE_CHECKING: log = get_logger(__name__) +def _coord_display_name(ws: Workstream) -> str: + """Resolve a coordinator's display name (``alias > title > name``). + + ``ws.name`` is the synthetic ``ws-xxxx`` placeholder; the persisted + auto-title (``update_workstream_title``) and user alias live only in + the DB. Seeding the collector with the resolved name means a + rehydrated coordinator shows its title in the live cluster tree + immediately, rather than reverting to ``ws-xxxx`` until a (for + coordinators, rarely-firing) ``on_rename`` event arrives. + """ + return get_workstream_display_name(ws.id) or ws.name + + class CoordinatorAdapter: """Bridges SessionManager to the console's coordinator transport.""" @@ -132,7 +146,7 @@ class CoordinatorAdapter: try: self._collector.emit_console_ws_created( ws.id, - name=ws.name, + name=_coord_display_name(ws), user_id=ws.user_id, kind=ws.kind.value, state=ws.state.value, @@ -466,11 +480,16 @@ class CoordinatorAdapter: # creates happened before the collector was wired up and their # rows never showed on the snapshot. (Coord-specific — interactive # has no analogous pseudo-node.) - for ws in mgr.list_all(): + coords = mgr.list_all() + # One round-trip for every coordinator's display name instead of a + # per-``ws`` ``_coord_display_name`` lookup (N+1); cold path, but + # the bulk helper is right there. + seed_names = get_workstream_display_names([ws.id for ws in coords]) + for ws in coords: try: collector.emit_console_ws_created( ws.id, - name=ws.name, + name=seed_names.get(ws.id) or ws.name, user_id=ws.user_id or "", kind=WorkstreamKind.COORDINATOR.value, state=ws.state.value, diff --git a/turnstone/console/server.py b/turnstone/console/server.py index d86ff09d..0c4341a2 100644 --- a/turnstone/console/server.py +++ b/turnstone/console/server.py @@ -56,6 +56,7 @@ from turnstone.core.auth import ( require_permission, ) from turnstone.core.deadline import DeadlineExceededError, run_with_deadline +from turnstone.core.memory import get_workstream_display_names from turnstone.core.rendezvous import NoAvailableNodeError from turnstone.core.session_replay import session_replay_preamble from turnstone.core.session_routes import ( @@ -853,6 +854,14 @@ def _coordinator_rows(request: Request) -> list[dict[str, Any]]: In-memory wins on ws_id conflict so live state stays authoritative for active sessions. + Display name resolves ``alias > title > name`` from the persisted + row for BOTH lanes. ``ws.name`` on the in-memory Workstream is the + synthetic ``ws-xxxx`` placeholder; the LLM auto-title + (``update_workstream_title``) and the user alias + (``set_workstream_alias``) live only in the DB, so without the + persisted lookup the live lane would show ``ws-xxxx`` and the + auto-title would never survive a dashboard refresh. + Trusted-team visibility (post-#400): the cluster dashboard shows every coordinator regardless of caller identity; ``user_id`` is surfaced on each row as display metadata. @@ -870,6 +879,54 @@ def _coordinator_rows(request: Request) -> list[dict[str, Any]]: val = getattr(sess, name, "") if sess else "" return val if isinstance(val, str) else "" + # Persisted coordinator rows serve two purposes: (1) surface + # closed / error / deleted coordinators the manager has already + # evicted from ``self._workstreams``, and (2) supply the persisted + # display name (``alias > title > name``) for the LIVE coordinators + # too — ``ws.name`` is the synthetic placeholder. Cluster-wide + # (trusted-team visibility). Indexed by ws_id so both lanes resolve + # the same way. + storage = getattr(request.app.state, "auth_storage", None) + persisted: list[Any] = [] + if storage is not None: + try: + persisted = storage.list_workstreams( + kind=WorkstreamKind.COORDINATOR, + user_id=None, + limit=200, + ) + except Exception: + log.debug("cluster_workstreams.coord_persisted_failed", exc_info=True) + persisted = [] + # SQLAlchemy Row — access via _mapping so future SELECT reorders / + # new columns don't silently corrupt the projection (per the + # storage-protocol guidance on list_workstreams). Test doubles must + # expose the same ._mapping attribute. + meta: dict[str, Any] = {} + for row in persisted: + m = row._mapping + rid = m.get("ws_id") or "" + if rid: + meta[rid] = m + + # Live coordinators resolve their display name through the bulk + # helper keyed on their EXACT ids (one round-trip, no row cap) rather + # than the ``limit=200`` ``meta`` map: a live coord that has dropped + # below the 200-row ``updated DESC`` window would otherwise revert to + # its synthetic ``ws.name``. Closed/evicted rows (the persisted lane + # below) already carry alias/title in their own ``_mapping``. + live_display = get_workstream_display_names([ws.id for ws in wss]) if wss else {} + + def _display_name(ws_id: str, fallback: str) -> str: + m = meta.get(ws_id) + if m is None: + return fallback + return m.get("alias") or m.get("title") or m.get("name") or fallback + + def _title(ws_id: str) -> str: + m = meta.get(ws_id) + return str(m.get("title") or "") if m is not None else "" + rows: list[dict[str, Any]] = [] seen: set[str] = set() for ws in wss: @@ -877,9 +934,9 @@ def _coordinator_rows(request: Request) -> list[dict[str, Any]]: rows.append( { "id": ws.id, - "name": ws.name, + "name": live_display.get(ws.id) or ws.name, "state": ws.state.value, - "title": "", + "title": _title(ws.id), "node": "console", "server_url": "", "model": _str_sess_attr(sess, "model"), @@ -896,30 +953,7 @@ def _coordinator_rows(request: Request) -> list[dict[str, Any]]: ) seen.add(ws.id) - # Second lane — persisted coordinator rows, used to surface - # closed / error / deleted coordinators the manager has already - # evicted from ``self._workstreams``. Cluster-wide (trusted-team - # visibility). - storage = getattr(request.app.state, "auth_storage", None) - if storage is None: - return rows - try: - persisted = storage.list_workstreams( - kind=WorkstreamKind.COORDINATOR, - user_id=None, - limit=200, - ) - except Exception: - log.debug("cluster_workstreams.coord_persisted_failed", exc_info=True) - return rows - for row in persisted: - # SQLAlchemy Row — access via _mapping so future SELECT reorders - # / new columns don't silently corrupt the projection (per the - # storage-protocol guidance on list_workstreams). Test doubles - # must expose the same ._mapping attribute; positional indexing - # was removed because it hard-coded column offsets that drift - # with migrations. m = row._mapping row_id = m.get("ws_id") or "" if not row_id or row_id in seen: @@ -928,9 +962,9 @@ def _coordinator_rows(request: Request) -> list[dict[str, Any]]: rows.append( { "id": row_id, - "name": m.get("name") or f"coord-{row_id[:4]}", + "name": _display_name(row_id, f"coord-{row_id[:4]}"), "state": str(m.get("state") or "idle"), - "title": "", + "title": _title(row_id), "node": "console", "server_url": "", "model": "", diff --git a/turnstone/core/session.py b/turnstone/core/session.py index 78f211c7..467599e6 100644 --- a/turnstone/core/session.py +++ b/turnstone/core/session.py @@ -2425,10 +2425,15 @@ class ChatSession: ws_id = self._ws_id # Capture before async work log.info("ws.title.gen_start", ws_id=ws_id[:8]) try: - # Gather first user message and first assistant reply + # Gather first user message and first assistant reply. + # Snapshot ``self.messages`` (C-level atomic copy under the + # GIL): this runs in a background thread that may now fire + # while the main ``send`` loop is still streaming and + # appending turns, so iterating the live list directly could + # raise "list changed size during iteration". user_msg = "" asst_msg = "" - for m in self.messages: + for m in list(self.messages): content = m.text # joins text blocks; multipart attachments contribute none if m.role is Role.USER and not user_msg: user_msg = content[:300] @@ -4150,6 +4155,20 @@ class ChatSession: # legacy per-message ``_reminders`` side-channel splice. self._emit_pending_user_nudges() + # Auto-title from the opening user message — fire NOW rather than + # waiting for the assistant's final tool-call-free turn. The old + # trigger sat in the ``not tool_calls`` branch of the loop below; + # coordinators spend nearly every turn in tool calls and may never + # reach that terminal text turn, so the title almost never + # generated for them. Gate on a real user message: synthetic wake + # sends carry no content and ``_generate_title`` would no-op on the + # empty/attachment-only case anyway (it needs first-user-message + # text). The background thread snapshots ``self.messages`` so it is + # safe to run concurrently with the streaming turn started below. + if not self._title_generated and user_input.strip() and not from_wake: + self._title_generated = True + threading.Thread(target=self._generate_title, daemon=True).start() + # A fresh session composed its system prefix at __init__ with an empty # history, so memory selection fell back to recency (no query, no rerank). # Recompose once the first real user message exists so the opening turn @@ -4292,10 +4311,6 @@ class ChatSession: self._compact_messages(auto=True) # Update status bar with post-compaction token counts self._print_status_line() - # Auto-title session after first exchange - if not self._title_generated: - self._title_generated = True - threading.Thread(target=self._generate_title, daemon=True).start() # Flush any queued messages that weren't injected # (no tool calls → no advisory seam to inject at). # If anything drained, the model hasn't seen those diff --git a/turnstone/core/storage/_postgresql.py b/turnstone/core/storage/_postgresql.py index 5e1ddc28..22dfcf75 100644 --- a/turnstone/core/storage/_postgresql.py +++ b/turnstone/core/storage/_postgresql.py @@ -1052,6 +1052,12 @@ class PostgreSQLBackend: workstreams.c.skill_id, workstreams.c.skill_version, workstreams.c.user_id, + # Appended after ``user_id`` so positional fallbacks in + # consumers (``_coord_children_row`` et al.) that index + # up to row[9] stay valid; ``_coordinator_rows`` reads + # these by name to surface the persisted display title. + workstreams.c.title, + workstreams.c.alias, ) .order_by(workstreams.c.updated.desc()) .limit(limit) diff --git a/turnstone/core/storage/_protocol.py b/turnstone/core/storage/_protocol.py index 132881fe..ba6e409f 100644 --- a/turnstone/core/storage/_protocol.py +++ b/turnstone/core/storage/_protocol.py @@ -679,7 +679,9 @@ class StorageBackend(Protocol): Returns a list of SQLAlchemy ``Row`` objects. **Prefer dict access via ``row._mapping[]``**; positional indexing is brittle against future SELECT reorders and against new columns appearing in the - tail (the select currently ends with ``user_id``). + tail (the select currently ends with ``user_id, title, alias`` — + ``title``/``alias`` were appended after ``user_id`` so existing + positional fallbacks that index up to row[9] stay valid). """ ... diff --git a/turnstone/core/storage/_sqlite.py b/turnstone/core/storage/_sqlite.py index a60540df..b89d1247 100644 --- a/turnstone/core/storage/_sqlite.py +++ b/turnstone/core/storage/_sqlite.py @@ -1209,6 +1209,12 @@ class SQLiteBackend: workstreams.c.skill_id, workstreams.c.skill_version, workstreams.c.user_id, + # Appended after ``user_id`` so positional fallbacks in + # consumers (``_coord_children_row`` et al.) that index + # up to row[9] stay valid; ``_coordinator_rows`` reads + # these by name to surface the persisted display title. + workstreams.c.title, + workstreams.c.alias, ) .order_by(workstreams.c.updated.desc()) .limit(limit)