From 67c4398f254b9473131c796ceb8e4533135633ee Mon Sep 17 00:00:00 2001 From: Eric Lee Date: Sun, 27 Sep 2026 21:45:44 -0700 Subject: [PATCH] fix(tui): /resume finds and resumes this workspace's saved sessions The TUI's Sessions picker (/resume, /sessions, /switch) said "0 resumable" in every workspace. Three defects stacked: - gatewayClient answered session.list locally with an empty list (a "Phase 2" stub since #572). - session.resume was stubbed to return the current session, so resuming by id blanked the screen and kept the old conversation. - The backend's list_sessions returned an arbitrary slice. Session files carry updated_at as a float or an ISO string, and sorting the mix raised inside a catch-all. Backend: - list_sessions sorts on a normalized epoch (NaN and unreadable values fall back to mtime), scopes to the caller's cwd, excludes the live session, clamps the limit, and runs off the loop. - resume can return transcript rows (include_messages): prompts, replies and tool calls. Injected reminders, meta messages, tool results and reasoning are left out. - resume validates the id and resolves /resume off the loop; a blank title matches nothing. - The resume reply carries the restored-tasks notice. - New delete_session control for the picker's `d d`. TUI: - session.list and session.resume ride those controls. Tool rows are summarized with the live toolContext, and notices are appended. - The model is refreshed after a resume. - The resume's session.stats fires after the picker's reset (the timer is armed after the last await). - gateway.ready fires once per spawned backend, so a resume that re-sends init can't forge a session over the repaint. - A failed resume with no live session falls back to a fresh one instead of leaving prompts queued. Co-Authored-By: Claude Opus 5.5 <noreply@anthropic.com> --- CHANGELOG.md | 11 + src/server/agent_server.py | 225 ++++++++++++++--- tests/server/test_session_resume_picker.py | 265 +++++++++++++++++++++ ui-tui/src/__tests__/gatewayClient.test.ts | 122 ++++++++++ ui-tui/src/app/useSessionLifecycle.ts | 11 +- ui-tui/src/gatewayClient.ts | 166 ++++++++++++- 6 files changed, 761 insertions(+), 39 deletions(-) create mode 100644 tests/server/test_session_resume_picker.py diff --git a/CHANGELOG.md b/CHANGELOG.md index ab2ec32e8..72db32cd8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -9,6 +9,17 @@ and this project adheres to [Semantic Versioning](https://semver.org/spec/v2.0.0 ### Fixed +- **`/resume` in the TUI finds and resumes saved sessions.** The Sessions + picker (`/resume`, `/sessions`, `/switch`) said "0 resumable" in every + workspace: the client answered its own list request with nothing, and + resuming by id silently kept the current conversation. It now lists this + workspace's saved sessions, newest first; picking one (or `/resume <id or + title>`) replays the conversation into the live session, repaints its + transcript, and restores its turn count. `d d` on a row deletes it. The + backend's session list is sorted again too: files store `updated_at` as a + number or, from the older writer, a date string, and sorting the mix failed + silently, so a list request returned an arbitrary slice instead of the + newest sessions. - **The TUI no longer slows to a crawl in long sessions.** The bundle shipped React's development build: the launcher runs `node dist/entry.js` with `NODE_ENV` unset, and React picks its build from that at load time. Its diff --git a/src/server/agent_server.py b/src/server/agent_server.py index 038ed8cef..93eeb3a47 100644 --- a/src/server/agent_server.py +++ b/src/server/agent_server.py @@ -58,6 +58,8 @@ import asyncio import json import logging +import math +import os import queue as _queue import re import tempfile @@ -67,6 +69,7 @@ from collections.abc import AsyncIterator, Mapping from dataclasses import dataclass, field from dataclasses import replace as _dc_replace +from datetime import datetime from pathlib import Path from typing import Any @@ -809,7 +812,31 @@ async def _handle_control_request(self, msg: dict) -> None: self._do_rewind(request_id, inner.get("turns", 1)) return if subtype == "list_sessions": - self._reply(request_id, {"sessions": _list_saved_sessions()}) + # Parses every saved session file, so it runs off the loop. + limit = inner.get("limit") + cwd = inner.get("cwd") + sessions = await asyncio.to_thread( + _list_saved_sessions, + min(max(limit, 1), 200) if isinstance(limit, int) and not isinstance(limit, bool) else 20, + cwd=cwd if isinstance(cwd, str) and cwd else None, + exclude=self.session_id, + ) + self._reply(request_id, {"sessions": sessions}) + return + if subtype == "delete_session": + # The resume picker's `d d`. The live session's file is the one + # being written, so it is not deletable from under it. + target = inner.get("session_id") + if not isinstance(target, str) or not target: + self._reply(request_id, {"ok": False, "error": "missing session_id"}) + elif target == self.session_id: + self._reply(request_id, {"ok": False, "error": "cannot delete the live session"}) + else: + from src.server.desktop_sessions import delete_session + + deleted = delete_session(_sessions_dir(), target) + self._reply(request_id, {"ok": True, "deleted": target} if deleted + else {"ok": False, "error": "session not found"}) return if subtype == "rename": name = inner.get("name") @@ -821,7 +848,19 @@ async def _handle_control_request(self, msg: dict) -> None: self._start_generate_title(request_id, inner.get("text")) return if subtype == "resume": - self._do_resume(request_id, inner.get("session_id")) + target = inner.get("session_id") + # `/resume <title>` names a session instead of its file; finding + # it reads every saved session, so off the loop. + if isinstance(target, str) and target and _saved_session_file(target) is None: + named = await asyncio.to_thread( + _saved_session_named, target, cwd=self.cwd, exclude=self.session_id, + ) + target = named or target + self._do_resume( + request_id, + target, + include_messages=bool(inner.get("include_messages")), + ) return if subtype == "get_activity": self._reply(request_id, self._activity_snapshot()) @@ -3354,9 +3393,14 @@ def _activity_snapshot(self) -> dict[str, Any]: "busy": turn_active or queued or goal_active or scheduled or background, } - def _do_resume(self, request_id: object, session_id: object) -> None: + def _do_resume( + self, request_id: object, session_id: object, *, include_messages: bool = False + ) -> None: """Load a saved conversation into this session (the original's /resume). - Idle-only — replacing the conversation mid-turn would race the worker.""" + Idle-only — replacing the conversation mid-turn would race the worker. + + ``include_messages`` adds the conversation as transcript rows, for a + client that repaints the resumed history (the Ink TUI's picker).""" with self._lock: active = self._current_abort is not None if active: @@ -3366,8 +3410,8 @@ def _do_resume(self, request_id: object, session_id: object) -> None: if not isinstance(session_id, str) or not session_id: self._reply(request_id, {"ok": False, "error": "missing session_id"}) return - f = _sessions_dir() / f"{session_id}.json" - if not f.exists(): + f = _saved_session_file(session_id) + if f is None: self._reply(request_id, {"ok": False, "error": "session not found"}) return from src.agent.conversation import Conversation @@ -3493,6 +3537,7 @@ def _do_resume(self, request_id: object, session_id: object) -> None: # only an ACTIVE goal carries over; turn count, timer, and the # token-spend baseline reset. Achieved/cleared goals stay gone. goal_notice = None + cron_notice = None saved_goal = data.get("goal") if isinstance(saved_goal, dict): try: @@ -3543,10 +3588,11 @@ def _do_resume(self, request_id: object, session_id: object) -> None: if isinstance(saved_sched, dict) else 0 ) if restored_n: - self._push_cron_state( + cron_notice = ( f"⏰ Restored {restored_n} scheduled task(s) " "from the saved session." ) + self._push_cron_state(cron_notice) else: self._push_cron_state() except Exception: # noqa: BLE001 — must not break resume @@ -3563,6 +3609,10 @@ def _do_resume(self, request_id: object, session_id: object) -> None: "cost": _cost_snapshot(), **({"mode_banner": mode_banner} if mode_banner else {}), **({"goal_notice": goal_notice} if goal_notice else {}), + # Also pushed as a cron_status line above; carried here too for + # a client that repaints the transcript once this reply lands. + **({"cron_notice": cron_notice} if cron_notice else {}), + **({"messages": _transcript_rows(conv.messages)} if include_messages else {}), }) except Exception as exc: # noqa: BLE001 logger.exception("[agent-server] resume failed") @@ -7025,6 +7075,80 @@ def _first_prompt_preview(msgs: list) -> str: return "" +# A saved session's id as a file stem: no separators, so it can't leave the +# sessions directory (the same set desktop_sessions._safe_id accepts). +_SESSION_ID_RE = re.compile(r"[A-Za-z0-9_-]+") + + +def _saved_session_file(session_id: str) -> Path | None: + """The saved file of a plain session id; ``None`` when there is none, or + when the id has path pieces that could reach outside the directory.""" + if not _SESSION_ID_RE.fullmatch(session_id): + return None + f = _sessions_dir() / f"{session_id}.json" + return f if f.exists() else None + + +def _saved_session_named(title: str, *, cwd: str, exclude: str | None) -> str | None: + """`/resume <title>`: the newest of this workspace's sessions renamed to + ``title`` (case-insensitive), as long as its id is a plain one.""" + wanted = title.strip().casefold() + if not wanted: # a blank title would match every unnamed session + return None + for row in _list_saved_sessions(200, cwd=cwd, exclude=exclude): + if row["name"].strip().casefold() == wanted and _SESSION_ID_RE.fullmatch(row["session_id"]): + return row["session_id"] + return None + +# The tool-input fields the TUI's trail summary reads (gatewayClient toolContext). +_TOOL_CONTEXT_FIELDS = ( + "pattern", "file_path", "path", "notebook_path", "command", "url", "query", + "description", "prompt", +) + + +def _transcript_rows(msgs: list) -> list[dict]: + """A saved conversation as the TUI's transcript rows: prompts and replies + as ``{role, text}``, each tool call as ``{role: "tool", name, input}``. + + What a live turn never painted stays out: meta and compact-summary + messages, injected ``<system-reminder>`` prompts (teammate mail, task + notices — the web client hides the same rows), tool results, and reasoning + items. A tool row keeps only the input fields its one-line summary reads, + so a Write doesn't ship the whole file it wrote. + """ + + def field(block: object, key: str) -> object: + return block.get(key) if isinstance(block, dict) else getattr(block, key, None) + + def injected(role: str, text: str) -> bool: + return role == "user" and text.lstrip().startswith("<system-reminder>") + + rows: list[dict] = [] + for m in msgs: + role = getattr(m, "role", None) + if role not in ("user", "assistant") or getattr(m, "isMeta", False) \ + or getattr(m, "isCompactSummary", False): + continue + content = getattr(m, "content", None) + blocks = [{"type": "text", "text": content}] if isinstance(content, str) else content + for block in blocks or []: + kind = field(block, "type") + if kind == "text": + text = field(block, "text") + if isinstance(text, str) and text.strip() and not injected(role, text): + rows.append({"role": role, "text": text}) + elif kind == "tool_use" and role == "assistant": + raw = field(block, "input") + summary = { + key: str(value)[:500] + for key, value in (raw.items() if isinstance(raw, dict) else ()) + if key in _TOOL_CONTEXT_FIELDS and isinstance(value, (str, int, float)) + } + rows.append({"role": "tool", "name": str(field(block, "name") or "tool"), "input": summary}) + return rows + + def _count_prompt_turns(msgs: list) -> int: """Real user prompts in a conversation — re-seeds ``_stats_turns`` after /resume and /rewind. A prompt is a user message that isn't an injected @@ -7045,30 +7169,71 @@ def _count_prompt_turns(msgs: list) -> int: return n -def _list_saved_sessions(limit: int = 20) -> list[dict]: - """Saved sessions, newest first (for /resume).""" - out: list[dict] = [] +def _updated_epoch(value: object, fallback: float) -> float: + """A session file's ``updated_at`` as epoch seconds. + + The agent server writes a float; the older session writer wrote an ISO + string. Sorting the two together raised, so this is the one sortable form + — anything unreadable (or non-finite, which would also break the JSON + reply) falls back to the file's mtime. + """ + epoch: float | None = None + if isinstance(value, (int, float)) and not isinstance(value, bool): + epoch = float(value) + elif isinstance(value, str): + try: + epoch = datetime.fromisoformat(value.replace("Z", "+00:00")).timestamp() + except (ValueError, OSError, OverflowError): # OSError: a pre-1970 local time on Windows + epoch = None + return epoch if epoch is not None and math.isfinite(epoch) else fallback + + +def _list_saved_sessions( + limit: int = 20, *, cwd: str | None = None, exclude: str | None = None +) -> list[dict]: + """Saved sessions, newest first (for /resume). + + ``cwd`` keeps the sessions that ran in that directory — a resume replays a + conversation into this session's workspace, so another project's rows only + crowd the picker. ``exclude`` drops one id: the live session is not + something to resume into itself. + """ try: - d = _sessions_dir() - if not d.exists(): - return [] - for f in d.glob("*.json"): - try: - data = json.loads(f.read_text(encoding="utf-8")) - except Exception: # noqa: BLE001 - continue - out.append({ - "session_id": data.get("session_id", f.stem), - "updated_at": data.get("updated_at", 0), - "preview": data.get("preview", ""), - "name": data.get("name") or "", - "message_count": data.get("message_count", 0), - "model": data.get("model", ""), - "cwd": data.get("cwd", ""), # for the TagTabs project filter - }) - out.sort(key=lambda s: s.get("updated_at", 0), reverse=True) - except Exception: # noqa: BLE001 - pass + files = list(_sessions_dir().glob("*.json")) + except OSError: + return [] + want = os.path.normpath(cwd) if cwd else None + out: list[dict] = [] + for f in files: + try: + data = json.loads(f.read_text(encoding="utf-8")) + mtime = f.stat().st_mtime + except Exception: # noqa: BLE001 — one bad file must not hide the rest + continue + if not isinstance(data, dict): + continue + session_id = data.get("session_id") + if not isinstance(session_id, str) or not session_id: + session_id = f.stem + name = data.get("name") + session_cwd = data.get("cwd") or "" + if session_id == exclude: + continue + if want is not None and ( + not isinstance(session_cwd, str) or not session_cwd + or os.path.normpath(session_cwd) != want + ): + continue + out.append({ + "session_id": session_id, + "updated_at": _updated_epoch(data.get("updated_at"), mtime), + "preview": data.get("preview", ""), + "name": name if isinstance(name, str) else "", + "message_count": data.get("message_count", 0), + "model": data.get("model", ""), + "cwd": session_cwd, # for the TagTabs project filter + }) + out.sort(key=lambda s: s["updated_at"], reverse=True) return out[:limit] diff --git a/tests/server/test_session_resume_picker.py b/tests/server/test_session_resume_picker.py new file mode 100644 index 000000000..1ba172987 --- /dev/null +++ b/tests/server/test_session_resume_picker.py @@ -0,0 +1,265 @@ +"""The TUI's resume picker is served by ``list_sessions`` and ``resume``. + +The picker showed "0 resumable" in every workspace: the client never asked, +and the listing it would have asked for came back unsorted — session files +carry ``updated_at`` as a float or, from the older writer, an ISO string, and +sorting the mix raised inside a catch-all. +""" +from __future__ import annotations + +import asyncio +import json +from pathlib import Path +from unittest import mock + +import pytest + +from src.agent.conversation import Conversation +from src.types.messages import AssistantMessage, UserMessage + + +def _session(cwd: str, session_id: str = "live"): + from src.server.agent_server import AgentServerConfig, _AgentSession + + emitted: list = [] + sess = _AgentSession( + session_id=session_id, + cwd=cwd, + config=AgentServerConfig(single_session=True), + loop=mock.MagicMock(), + out_queue=mock.MagicMock(), + ) + sess._emit = lambda env: emitted.append(env) + return sess, emitted + + +def _reply_of(emitted: list) -> dict: + return emitted[-1]["response"]["response"] + + +@pytest.fixture() +def sessions_dir(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> Path: + # Everything a resume touches (cost restore reads the config dir too) + # stays under tmp_path, never the developer's real ~/.clawcodex. + monkeypatch.setenv("CLAWCODEX_CONFIG_DIR", str(tmp_path / "config")) + directory = tmp_path / "sessions" + directory.mkdir() + monkeypatch.setattr("src.server.agent_server._sessions_dir", lambda: directory) + return directory + + +def _save(directory: Path, session_id: str, **fields) -> None: + data = {"session_id": session_id, "conversation": {"messages": []}, **fields} + (directory / f"{session_id}.json").write_text(json.dumps(data), encoding="utf-8") + + +def test_saved_sessions_list_newest_first_across_both_timestamp_formats(sessions_dir: Path) -> None: + from src.server.agent_server import _list_saved_sessions + + # Names sort opposite to their ages, so an unsorted (directory-order) + # listing cannot pass by accident. + _save(sessions_dir, "a-oldest", updated_at=1_700_000_000.0) + _save(sessions_dir, "b-iso", updated_at="2024-06-01T00:00:00") + _save(sessions_dir, "c-newer", updated_at=1_790_000_000.0) + _save(sessions_dir, "d-newest-iso", updated_at="2026-12-31T23:00:00Z") + (sessions_dir / "e-garbage.json").write_text("{not json", encoding="utf-8") + + rows = _list_saved_sessions(10) + + assert [row["session_id"] for row in rows] == ["d-newest-iso", "c-newer", "b-iso", "a-oldest"] + assert all(isinstance(row["updated_at"], float) for row in rows) + + +def test_a_non_finite_timestamp_cannot_break_the_listing_reply(sessions_dir: Path) -> None: + from src.server.agent_server import _list_saved_sessions + + (sessions_dir / "nan.json").write_text('{"session_id": "nan", "updated_at": NaN}', encoding="utf-8") + _save(sessions_dir, "fine", updated_at=1.0) + + rows = _list_saved_sessions(10) + + # NaN would sort unpredictably and make the NDJSON reply invalid JSON. + json.dumps(rows, allow_nan=False) + assert {row["session_id"] for row in rows} == {"nan", "fine"} + + +def test_list_sessions_clamps_its_limit(tmp_path: Path, sessions_dir: Path) -> None: + for n in range(3): + _save(sessions_dir, f"s{n}", cwd=str(tmp_path), updated_at=float(n)) + sess, emitted = _session(str(tmp_path)) + + asyncio.run(sess._handle_control_request( + {"request_id": "r1", "request": {"subtype": "list_sessions", "limit": 0}} + )) + + assert [row["session_id"] for row in _reply_of(emitted)["sessions"]] == ["s2"] + + +def test_list_sessions_keeps_this_workspace_and_leaves_out_the_live_session( + tmp_path: Path, sessions_dir: Path +) -> None: + here = str(tmp_path / "project") + _save(sessions_dir, "mine", cwd=here, updated_at=3.0, message_count=788, preview="prove it") + _save(sessions_dir, "live", cwd=here, updated_at=4.0) + _save(sessions_dir, "elsewhere", cwd=str(tmp_path / "other"), updated_at=5.0) + _save(sessions_dir, "no-cwd", updated_at=6.0) + sess, emitted = _session(here, session_id="live") + + asyncio.run(sess._handle_control_request( + {"request_id": "r1", "request": {"subtype": "list_sessions", "cwd": here + "/", "limit": 50}} + )) + + reply = _reply_of(emitted) + assert [row["session_id"] for row in reply["sessions"]] == ["mine"] + assert reply["sessions"][0]["message_count"] == 788 + assert reply["sessions"][0]["preview"] == "prove it" + + +def test_list_sessions_without_a_cwd_still_lists_every_workspace(tmp_path: Path, sessions_dir: Path) -> None: + _save(sessions_dir, "one", cwd=str(tmp_path / "a"), updated_at=1.0) + _save(sessions_dir, "two", cwd=str(tmp_path / "b"), updated_at=2.0) + sess, emitted = _session(str(tmp_path)) + + asyncio.run(sess._handle_control_request({"request_id": "r1", "request": {"subtype": "list_sessions"}})) + + assert [row["session_id"] for row in _reply_of(emitted)["sessions"]] == ["two", "one"] + + +def _saved_conversation(directory: Path) -> None: + conv = Conversation() + conv.messages.extend([ + UserMessage(content="prove the lemma"), + UserMessage(content="<system-reminder>Messages from your teammates follow.</system-reminder>"), + UserMessage(content="hook context the model saw but the user never typed", isMeta=True), + AssistantMessage(content=[ + {"type": "openai_responses_item", "item": {"type": "reasoning", "encrypted_content": "x"}}, + {"type": "text", "text": "Reading the notes first."}, + {"type": "tool_use", "id": "t1", "name": "Read", + "input": {"file_path": "/w/notes.md", "offset": 10}}, + ]), + UserMessage(content=[{"type": "tool_result", "tool_use_id": "t1", "content": "1\tbody"}]), + AssistantMessage(content=[ + {"type": "tool_use", "id": "t2", "name": "Write", + "input": {"file_path": "/w/proof.lean", "content": "x" * 50_000}}, + ]), + UserMessage(content=[{"type": "tool_result", "tool_use_id": "t2", "content": "ok"}]), + AssistantMessage(content=[{"type": "text", "text": "The lemma holds."}]), + UserMessage(content="summary of earlier turns", isCompactSummary=True), + UserMessage(content=[{"type": "text", "text": "[Image #1] now the multiple-of-four case"}]), + ]) + _save(directory, "old", conversation=conv.to_dict(), cwd=str(directory.parent), name="Liouville run") + + +def test_resume_can_return_the_conversation_as_transcript_rows(tmp_path: Path, sessions_dir: Path) -> None: + _saved_conversation(sessions_dir) + sess, emitted = _session(str(tmp_path)) + sess.session = mock.MagicMock() + + sess._do_resume("r", "old", include_messages=True) + + reply = _reply_of(emitted) + assert reply["ok"] is True + assert reply["messages"] == [ + {"role": "user", "text": "prove the lemma"}, + {"role": "assistant", "text": "Reading the notes first."}, + {"role": "tool", "name": "Read", "input": {"file_path": "/w/notes.md"}}, + {"role": "tool", "name": "Write", "input": {"file_path": "/w/proof.lean"}}, + {"role": "assistant", "text": "The lemma holds."}, + {"role": "user", "text": "[Image #1] now the multiple-of-four case"}, + ] + + +def test_resume_without_the_flag_keeps_its_old_reply(tmp_path: Path, sessions_dir: Path) -> None: + _saved_conversation(sessions_dir) + sess, emitted = _session(str(tmp_path)) + sess.session = mock.MagicMock() + + sess._do_resume("r", "old") + + reply = _reply_of(emitted) + assert reply["ok"] is True and reply["count"] == 10 + assert "messages" not in reply + + +def test_resume_finds_a_renamed_session_in_this_workspace_by_its_title( + tmp_path: Path, sessions_dir: Path +) -> None: + _saved_conversation(sessions_dir) + sess, emitted = _session(str(tmp_path)) + sess.session = mock.MagicMock() + + def resume(session, target: str) -> None: + asyncio.run(session._handle_control_request({"request_id": "r", "request": { + "subtype": "resume", "session_id": target, "include_messages": True}})) + + resume(sess, "liouville RUN") + + assert _reply_of(emitted)["messages"][0] == {"role": "user", "text": "prove the lemma"} + + elsewhere, emitted_elsewhere = _session(str(tmp_path / "other-project")) + elsewhere.session = mock.MagicMock() + resume(elsewhere, "Liouville run") + assert _reply_of(emitted_elsewhere) == {"ok": False, "error": "session not found"} + + +def test_a_blank_resume_target_matches_no_unnamed_session(tmp_path: Path, sessions_dir: Path) -> None: + _save(sessions_dir, "unnamed", cwd=str(tmp_path), updated_at=1.0) + sess, emitted = _session(str(tmp_path)) + sess.session = mock.MagicMock() + sess.session.conversation = "untouched" + + asyncio.run(sess._handle_control_request( + {"request_id": "r", "request": {"subtype": "resume", "session_id": " "}} + )) + + assert _reply_of(emitted) == {"ok": False, "error": "session not found"} + assert sess.session.conversation == "untouched" + + +def test_resume_never_builds_a_path_from_an_unsafe_id(tmp_path: Path, sessions_dir: Path) -> None: + (tmp_path / "outside.json").write_text(json.dumps({"conversation": {"messages": []}}), encoding="utf-8") + sess, emitted = _session(str(tmp_path)) + sess.session = mock.MagicMock() + sess.session.conversation = "untouched" + + sess._do_resume("r", "../outside") + + assert _reply_of(emitted) == {"ok": False, "error": "session not found"} + assert sess.session.conversation == "untouched" + + +def test_resume_carries_the_restored_scheduled_tasks_notice(tmp_path: Path, sessions_dir: Path) -> None: + from src.scheduled_tasks import SessionCronScheduler + + donor = SessionCronScheduler(jitter=False) + donor.create("*/5 * * * *", "check the build") + _save(sessions_dir, "old", scheduled_tasks=donor.snapshot()) + sess, emitted = _session(str(tmp_path)) + sess.session = mock.MagicMock() + sess.cron_scheduler = SessionCronScheduler(jitter=False) + + sess._do_resume("r", "old", include_messages=True) + + # Also pushed as its own cron_status line, which the TUI's repaint clears. + assert _reply_of(emitted)["cron_notice"] == "⏰ Restored 1 scheduled task(s) from the saved session." + + +def test_delete_session_removes_a_saved_session_but_never_the_live_one( + tmp_path: Path, sessions_dir: Path +) -> None: + _save(sessions_dir, "old") + _save(sessions_dir, "live") + sess, emitted = _session(str(tmp_path), session_id="live") + + def delete(target: str) -> dict: + asyncio.run(sess._handle_control_request( + {"request_id": "r", "request": {"subtype": "delete_session", "session_id": target}} + )) + return _reply_of(emitted) + + assert delete("old") == {"ok": True, "deleted": "old"} + assert not (sessions_dir / "old.json").exists() + assert delete("live") == {"ok": False, "error": "cannot delete the live session"} + assert (sessions_dir / "live.json").exists() + assert delete("../live") == {"ok": False, "error": "session not found"} + assert delete("old") == {"ok": False, "error": "session not found"} diff --git a/ui-tui/src/__tests__/gatewayClient.test.ts b/ui-tui/src/__tests__/gatewayClient.test.ts index 2d33534af..2658c9d1a 100644 --- a/ui-tui/src/__tests__/gatewayClient.test.ts +++ b/ui-tui/src/__tests__/gatewayClient.test.ts @@ -297,6 +297,128 @@ describe('GatewayClient NDJSON adapter', () => { proc.line({ response: { request_id: req.request_id, response }, type: 'control_response' }) } + // ── resume picker ────────────────────────────────────────────────────────── + // session.list and session.resume used to resolve locally, so /resume showed + // "0 resumable" in every workspace and a picked row left the old + // conversation in place. They now ride the list_sessions / resume controls. + + it('lists this workspace saved sessions for the resume picker', async () => { + proc.line(INIT) + const p = gw.request('session.list', { limit: 200 }) + await replyToControl('list_sessions', { + sessions: [{ message_count: 788, name: '', preview: '/math-team formal', session_id: 'ds_old', updated_at: 1_790_538_257 }] + }) + await expect(p).resolves.toEqual({ + sessions: [{ id: 'ds_old', message_count: 788, preview: '/math-team formal', started_at: 1_790_538_257, title: '' }] + }) + expect(seen.find(f => f.request?.subtype === 'list_sessions').request).toMatchObject({ cwd: '/ws', limit: 200 }) + }) + + it('keeps the switcher 1.5 s live-session poll off the backend', async () => { + proc.line(INIT) + await expect(gw.request('session.active_list', {})).resolves.toEqual({ sessions: [] }) + seen.push(...stdinFrames()) + expect(seen.some(f => f.request?.subtype === 'list_sessions')).toBe(false) + }) + + it('resumes a saved session through the backend and hands back its transcript', async () => { + proc.line(INIT) + const p = gw.request<any>('session.resume', { session_id: 'ds_old' }) + await replyToControl('resume', { + count: 4, + messages: [ + { role: 'user', text: 'prove the lemma' }, + { input: { file_path: '/ws/notes.md' }, name: 'Read', role: 'tool' }, + { role: 'assistant', text: 'It holds.' } + ], + mode_banner: 'Entered coordinator mode to match resumed session.', + ok: true, + session_turns: 5 + }) + await replyToControl('get_settings', { model: 'gpt-6-astra' }) + const r = await p + + expect(r).toMatchObject({ message_count: 4, session_id: 's1' }) + expect(r.messages).toEqual([ + { role: 'user', text: 'prove the lemma' }, + { context: 'notes.md', name: 'Read', role: 'tool' }, + { role: 'assistant', text: 'It holds.' }, + { role: 'system', text: 'Entered coordinator mode to match resumed session.' } + ]) + expect(r.info.model).toBe('gpt-6-astra') + expect(seen.find(f => f.request?.subtype === 'resume').request).toMatchObject({ + include_messages: true, + session_id: 'ds_old' + }) + // Stamped after the reply, not with it: the switcher zeroes the stats + // line when the reply lands and relies on this event to refill it. + expect(last('session.stats')).toBeUndefined() + await vi.waitFor(() => expect(last('session.stats')?.payload).toMatchObject({ session_turns: 5 })) + }) + + it('surfaces a refused resume instead of pretending it worked', async () => { + proc.line(INIT) + const p = gw.request('session.resume', { session_id: 'gone' }) + await replyToControl('resume', { error: 'session not found', ok: false }) + await expect(p).rejects.toThrow('session not found') + }) + + it('stamps the resumed stats after the caller reset, however late get_settings answers', async () => { + proc.line(INIT) + const order: string[] = [] + gw.on('event', (e: any) => e.type === 'session.stats' && order.push('stats')) + // resumeById resets the stats line synchronously when the reply lands. + const p = gw.request<any>('session.resume', { session_id: 'ds_old' }).then(r => (order.push('reset'), r)) + await replyToControl('resume', { count: 1, cron_notice: '⏰ Restored 1 scheduled task(s).', messages: [], ok: true, session_turns: 5 }) + await new Promise(resolve => setTimeout(resolve, 20)) + await replyToControl('get_settings', { model: 'm' }) + const r = await p + await vi.waitFor(() => expect(order).toEqual(['reset', 'stats'])) + expect(r.messages).toEqual([{ role: 'system', text: '⏰ Restored 1 scheduled task(s).' }]) + }) + + it('publishes gateway.ready once per backend, even when a resume re-sends init', async () => { + proc.line(INIT) + await vi.waitFor(() => expect(types()).toContain('gateway.ready')) + // A resume that flips coordinator mode re-emits system/init; a second + // ready would rerun startup and forge a session over the repaint. + proc.line({ ...INIT, tools: [{ name: 'Agent' }] }) + await vi.waitFor(() => expect(types().filter(t => t === 'session.info')).toHaveLength(2)) + expect(types().filter(t => t === 'gateway.ready')).toHaveLength(1) + expect(Object.values(last('session.info').payload.tools).flat()).toEqual(['Agent']) + }) + + it('publishes gateway.ready again for a respawned backend, so crash recovery still runs', async () => { + proc.line(INIT) + await vi.waitFor(() => expect(types().filter(t => t === 'gateway.ready')).toHaveLength(1)) + // useMainApp's exit handler calls start() again on the same client. + const respawned = new FakeProc() + harness.proc = respawned + gw.start() + respawned.line({ ...INIT, session_id: 's2' }) + await vi.waitFor(() => expect(types().filter(t => t === 'gateway.ready')).toHaveLength(2)) + }) + + it('reports a refused session list rather than an empty one', async () => { + proc.line(INIT) + const p = gw.request('session.list', { limit: 200 }) + await replyToControl('list_sessions', { error: 'sandbox refused to start', ok: false }) + await expect(p).rejects.toThrow('sandbox refused to start') + }) + + it('deletes a saved session through the backend, and reports a refusal', async () => { + proc.line(INIT) + const ok = gw.request('session.delete', { session_id: 'ds_old' }) + await replyToControl('delete_session', { deleted: 'ds_old', ok: true }) + await expect(ok).resolves.toEqual({ deleted: 'ds_old' }) + expect(seen.find(f => f.request?.subtype === 'delete_session').request).toMatchObject({ session_id: 'ds_old' }) + + seen = [] + const refused = gw.request('session.delete', { session_id: 's1' }) + await replyToControl('delete_session', { error: 'cannot delete the live session', ok: false }) + await expect(refused).rejects.toThrow('cannot delete the live session') + }) + // ── delegation control plane ─────────────────────────────────────────────── // These three RPCs had no case in request(), so they hit the `default:` arm // and resolved `{}` — the agents overlay's status readout, pause key and diff --git a/ui-tui/src/app/useSessionLifecycle.ts b/ui-tui/src/app/useSessionLifecycle.ts index 637abcef8..d91c7ab95 100644 --- a/ui-tui/src/app/useSessionLifecycle.ts +++ b/ui-tui/src/app/useSessionLifecycle.ts @@ -396,12 +396,21 @@ export function useSessionLifecycle(opts: UseSessionLifecycleOptions) { }) .catch((e: Error) => { + // Nothing to stay on (startup resume, crash recovery of a session + // that never saved): without a session every prompt only queues, + // so start a fresh one and say why. + if (!getUiState().sid) { + void newSession(`could not resume ${id}: ${e.message}`) + + return + } + sys(`error: ${e.message}`) patchUiState({ status: 'ready' }) }) }) }, - [closeSession, colsRef, gw, panel, resetSession, rpc, scrollRef, setHistoryItems, setSessionStartedAt, sys] + [closeSession, colsRef, gw, newSession, panel, resetSession, rpc, scrollRef, setHistoryItems, setSessionStartedAt, sys] ) const guardBusySessionSwitch = useCallback( diff --git a/ui-tui/src/gatewayClient.ts b/ui-tui/src/gatewayClient.ts index 2f3a31679..44c394277 100644 --- a/ui-tui/src/gatewayClient.ts +++ b/ui-tui/src/gatewayClient.ts @@ -26,8 +26,11 @@ import type { CostSnapshot, CronSnapshot, GatewayEvent, + GatewayTranscriptMessage, GoalSnapshot, PatchHunk, + SessionListItem, + SessionResumeResponse, StructuredDiffPayload } from './gatewayTypes.js' import { formatTotalCost, setLastCostSnapshot } from './lib/costSummary.js' @@ -51,6 +54,9 @@ const WORKTREE_RPC_TIMEOUT_MS = 600_000 * downsamples through Pillow. A timeout here would drop an image the user * watched themselves paste. */ const IMAGE_RPC_TIMEOUT_MS = 30_000 +/** Listing parses every saved session file, and a resume parses one that can + * run to megabytes; both outlast the 5 s default on a long history. */ +const SESSION_RPC_TIMEOUT_MS = 30_000 /** Fallback app version for the banner ("clawcodex v{version}"). * * The version normally comes from the backend's `system/init` frame, which @@ -146,6 +152,51 @@ function toolContext(input: any): string { return v == null ? '' : String(v) } +/** The backend's saved-session rows → the switcher's resumable history. + * `started_at` carries the last-updated time: the row's age should say when + * the session was last used, not when it began. */ +function toSessionListItems(rows: unknown): SessionListItem[] { + if (!Array.isArray(rows)) {return []} + + return rows.flatMap(row => { + const s = (row ?? {}) as Record<string, unknown> + const id = typeof s.session_id === 'string' ? s.session_id : '' + + return id + ? [ + { + id, + message_count: typeof s.message_count === 'number' ? s.message_count : 0, + preview: typeof s.preview === 'string' ? s.preview : '', + started_at: typeof s.updated_at === 'number' ? s.updated_at : 0, + title: typeof s.name === 'string' ? s.name : '' + } + ] + : [] + }) +} + +/** The backend's resume rows → transcript rows. A tool row's one-line summary + * is built here with the same toolContext a live tool.start uses, so a + * resumed trail reads like the one the session showed while it ran. */ +function toResumedTranscript(rows: unknown): GatewayTranscriptMessage[] { + if (!Array.isArray(rows)) {return []} + + return rows.flatMap((row): GatewayTranscriptMessage[] => { + const r = (row ?? {}) as Record<string, unknown> + + if (r.role === 'tool') { + return [{ context: toolContext(r.input), name: String(r.name ?? 'tool'), role: 'tool' }] + } + + if ((r.role === 'user' || r.role === 'assistant' || r.role === 'system') && typeof r.text === 'string') { + return [{ role: r.role, text: r.text }] + } + + return [] + }) +} + /** Shorten an absolute path to a workspace-relative path (or basename). */ function relativizePath(p: string): string { const ws = (process.env.CLAWCODEX_WORKSPACE || process.env.CLAWCODEX_CWD || process.cwd()).replace(/\/+$/, '') @@ -899,6 +950,11 @@ export class GatewayClient extends EventEmitter { private proc: ChildProcess | null = null private readyPromise: Promise<void> private readyResolve: (() => void) | null = null + // gateway.ready runs the app's startup (forge or resume a session), so it + // fires once per spawned backend. A resume that flips coordinator mode + // re-sends system/init; that one only refreshes session.info — as a second + // ready it forged a new session over the transcript just repainted. + private readyPublished = false private readyTimer: null | ReturnType<typeof setTimeout> = null private reqId = 0 private sessionId = '' @@ -926,6 +982,7 @@ export class GatewayClient extends EventEmitter { // ── lifecycle ──────────────────────────────────────────────────────────── start(): void { + this.readyPublished = false const cmd = resolveAgentCmd() const cwd = process.env.CLAWCODEX_WORKSPACE || process.env.CLAWCODEX_CWD || process.cwd() const env = { ...process.env, PYTHONUNBUFFERED: '1' } @@ -1152,12 +1209,19 @@ export class GatewayClient extends EventEmitter { // the live mode, so the client can't step into bypassPermissions // unconditionally or desync a cursor after /mode. return this.controlQuery('cycle_permission_mode', {}).then(r => (r ?? {}) as T) + case 'session.resume': { + const wanted = String((params as any)?.session_id ?? '') + + if (wanted) { + return this.readyPromise.then(() => this.resumeSaved(wanted) as Promise<T>) + } + + return this.readyPromise.then(() => ({ info: this.sessionInfo ?? undefined, session_id: this.sessionId }) as T) + } case 'session.activate': case 'session.create': - - case 'session.resume': // clawcodex runs a single agent-server session; hand back its id once // system/init has set it. The app then enables the composer. return this.readyPromise.then(() => ({ info: this.sessionInfo ?? undefined, session_id: this.sessionId }) as T) @@ -1286,12 +1350,48 @@ export class GatewayClient extends EventEmitter { } case 'session.active_list': - - case 'session.list': - // Single agent-server session in the basic port; the switcher/resume - // list is Phase 2. Resolve locally so the 1.5s poll doesn't spam the - // backend with list_sessions. + // One agent-server session, so no other live sessions to switch to. + // Resolved locally: the switcher polls this every 1.5 s. return Promise.resolve({ sessions: [] } as T) + case 'session.list': { + // The resumable history: this workspace's saved sessions, newest + // first, without the live one (the backend drops its own id). + const limit = typeof (params as any)?.limit === 'number' ? (params as any).limit : undefined + + return this.readyPromise + .then(() => + this.controlQuery('list_sessions', { cwd: this.sessionInfo?.cwd, limit }, SESSION_RPC_TIMEOUT_MS) + ) + .then(r => { + if (!r) { + throw new Error('session list timed out') + } + + // A refusal (a session that failed to start) is not "0 resumable". + if ((r as any).ok === false) { + throw new Error(String((r as any).error ?? 'could not list sessions')) + } + + return { sessions: toSessionListItems((r as any).sessions) } as T + }) + } + + case 'session.delete': { + // The switcher's `d d` on a resumable row: remove its saved file. + const target = String((params as any)?.session_id ?? '') + + return this.controlQuery('delete_session', { session_id: target }).then(r => { + if (!r) { + throw new Error('delete timed out') + } + + if ((r as any).ok === false) { + throw new Error(String((r as any).error ?? 'delete failed')) + } + + return { deleted: String((r as any).deleted ?? target) } as T + }) + } case 'session.interrupt': this.sendControl('interrupt', {}) @@ -2437,7 +2537,12 @@ export class GatewayClient extends EventEmitter { this.sessionId = String(msg.session_id ?? '') this.sessionInfo = this.toSessionInfo(msg) this.readyResolve?.() - this.publish({ payload: {}, session_id: this.sessionId, type: 'gateway.ready' }) + + if (!this.readyPublished) { + this.readyPublished = true + this.publish({ payload: {}, session_id: this.sessionId, type: 'gateway.ready' }) + } + this.publish({ payload: this.sessionInfo, session_id: this.sessionId, type: 'session.info' }) } else if (msg.subtype === 'status') { if (typeof msg.permission_mode === 'string' && msg.permission_mode) { @@ -2851,6 +2956,51 @@ export class GatewayClient extends EventEmitter { }) } + /** Replay a saved session into the live agent-server session (the + * backend's `resume` control) and hand the switcher its transcript. The + * session keeps its live id: the resumed conversation continues there. */ + private async resumeSaved(sessionId: string): Promise<SessionResumeResponse> { + const r = (await this.controlQuery( + 'resume', + { include_messages: true, session_id: sessionId }, + SESSION_RPC_TIMEOUT_MS + )) as any + + if (!r) { + throw new Error('resume timed out') + } + + if (r.ok === false) { + throw new Error(String(r.error ?? 'resume failed')) + } + + // A resume puts the saved session's model back when it ran on this + // provider, so re-read it rather than keep showing the launch model. + const settings = (await this.controlQuery('get_settings', {})) as any + const model = settings?.fusion || settings?.model + + if (this.sessionInfo && typeof model === 'string' && model) { + this.sessionInfo = { ...this.sessionInfo, model } + } + + const notices = [r.mode_banner, r.goal_notice, r.cron_notice].filter( + (text): text is string => typeof text === 'string' && text.trim() !== '' + ) + + // The switcher resets the session (stats line included) when this reply + // lands, and counts on the resume's session.stats to re-stamp it after. + // Armed after the last await: the timer then fires once the reply's + // promise chain, and so the reset, has run. + setTimeout(() => this.publishSessionStats(r), 0) + + return { + info: this.sessionInfo ?? undefined, + message_count: typeof r.count === 'number' ? r.count : undefined, + messages: [...toResumedTranscript(r.messages), ...notices.map(text => ({ role: 'system' as const, text }))], + session_id: this.sessionId + } + } + /** Stats-line refresh from a clear/resume reply's rider (session_turns + * cost snapshot). Silently a no-op for replies without the fields. */ private publishSessionStats(r: unknown): void {