diff --git a/src/context_system/builder.py b/src/context_system/builder.py index 1eaa55d6..7501df1c 100644 --- a/src/context_system/builder.py +++ b/src/context_system/builder.py @@ -101,8 +101,8 @@ def _build_workspace_section(root: Path, current: Path) -> str: f"- Today's date: {date.today().isoformat()}", f"- Workspace root: {workspace.workspace_root}", f"- Current directory: {workspace.current_directory}", - f"- Python files: {workspace.python_file_count}", - f"- Test files: {workspace.test_file_count}", + f"- Python files: {_count_text(workspace.python_file_count, workspace.counts_partial)}", + f"- Test files: {_count_text(workspace.test_file_count, workspace.counts_partial)}", ] if workspace.key_files: lines.append(f"- Key files: {', '.join(workspace.key_files)}") @@ -111,6 +111,11 @@ def _build_workspace_section(root: Path, current: Path) -> str: return "\n".join(lines) +def _count_text(count: int, partial: bool) -> str: + """``1234``, or ``1234+ (partial scan)`` when the walk hit its budget.""" + return f"{count}+ (partial scan)" if partial else str(count) + + def _build_git_section(cwd: str) -> str: from .git_context import collect_git_context, format_git_status try: diff --git a/src/context_system/models.py b/src/context_system/models.py index 27b70f33..79d0e8a9 100644 --- a/src/context_system/models.py +++ b/src/context_system/models.py @@ -105,3 +105,6 @@ class WorkspaceSnapshot: key_files: tuple[str, ...] python_file_count: int test_file_count: int + #: True when the file walk hit its directory budget, so the counts are + #: lower bounds — rendered as ``N+ (partial scan)``. + counts_partial: bool = False diff --git a/src/context_system/workspace_snapshot.py b/src/context_system/workspace_snapshot.py index e3ff6caf..e40adf15 100644 --- a/src/context_system/workspace_snapshot.py +++ b/src/context_system/workspace_snapshot.py @@ -1,18 +1,61 @@ +"""The ``## Runtime Context`` facts about a workspace: shape, key files, counts. + +The Python/test file counts come from ONE bounded ``os.scandir`` walk that +never descends into vendored or generated directories (``node_modules``, +``.git``, virtualenvs, caches, ``.clawcodex`` worktree checkouts …) and stops +after :data:`MAX_SCANNED_DIRS` directories. The previous implementation used +``Path.rglob`` twice, which walks *everything* and filters afterwards: on a +workspace with a few front-end packages (hundreds of thousands of +``node_modules`` directories) each walk took ~20 s, and the system prompt is +built at spawn and again on every resume/clear — so opening a saved session +from the web sidebar sat on this for ~45 s. The counts are a hint for the +model, not an inventory: a bounded scan that says "N+ (partial)" is the right +trade, and it keeps the cost proportional to the project rather than to +whatever a package manager left on disk. +""" + from __future__ import annotations +import os +from collections import deque from pathlib import Path from .models import WorkspaceSnapshot -_IGNORED_NAMES = { +#: Directory names the walk never enters. Vendored trees, VCS internals, +#: virtualenvs, tool caches, build output and the per-repo ``.clawcodex`` +#: directory (its ``worktrees/`` holds whole extra checkouts of the repo, +#: which would count every file again per worktree). +_IGNORED_NAMES = frozenset({ ".git", + ".hg", + ".svn", ".venv", + "venv", + ".tox", + ".nox", + ".eggs", + ".cache", ".pytest_cache", ".mypy_cache", ".ruff_cache", "__pycache__", "node_modules", -} + ".clawcodex", + ".claude", + "site-packages", +}) + +#: Upper bound on directories one snapshot visits (after pruning). Roughly a +#: few tens of milliseconds on a warm disk; past it the counts are reported as +#: partial rather than the walk running on. +MAX_SCANNED_DIRS = 2500 + +#: Upper bound on directory entries one snapshot reads, so a single flat +#: directory of a million files (a dataset dump next to the code) is bounded +#: the same way a deep tree is. +MAX_SCANNED_ENTRIES = 200_000 + _KEY_FILE_CANDIDATES = ( "README.md", "CLAWCODEX.md", @@ -29,6 +72,7 @@ def build_workspace_snapshot( *, cwd: str | Path | None = None, top_level_limit: int = 12, + max_dirs: int = MAX_SCANNED_DIRS, ) -> WorkspaceSnapshot: root = Path(workspace_root).expanduser().resolve() current = Path(cwd).expanduser().resolve() if cwd is not None else root @@ -49,8 +93,7 @@ def build_workspace_snapshot( break key_files = tuple(name for name in _KEY_FILE_CANDIDATES if (root / name).exists()) - python_file_count = sum(1 for path in root.rglob("*.py") if _is_countable(path)) - test_file_count = sum(1 for path in root.rglob("test_*.py") if _is_countable(path)) + python_file_count, test_file_count, partial = count_python_files(root, max_dirs=max_dirs) return WorkspaceSnapshot( workspace_root=root, @@ -59,9 +102,56 @@ def build_workspace_snapshot( key_files=key_files, python_file_count=python_file_count, test_file_count=test_file_count, + counts_partial=partial, ) +def count_python_files(root: Path, *, max_dirs: int = MAX_SCANNED_DIRS) -> tuple[int, int, bool]: + """``(python_files, test_files, partial)`` under ``root``, pruned and bounded. + + One breadth-first ``scandir`` walk: ignored directory names are skipped + *before* they are entered (this is what makes the walk cheap — filtering + ``rglob`` output afterwards still pays for every directory it crawled), + symlinks are never followed (nor counted: a symlinked ``.py`` is not a + file of this workspace), and once ``max_dirs`` directories or + :data:`MAX_SCANNED_ENTRIES` entries have been scanned the walk stops and + ``partial`` is True. Breadth-first so a partial + scan still covers the project's shallow structure rather than one deep + corner of it. Entries are visited in name order, so a partial count is + deterministic for a given tree. + """ + python_files = 0 + test_files = 0 + scanned = 0 + entries_seen = 0 + pending: deque[str] = deque([str(root)]) + while pending: + if scanned >= max_dirs or entries_seen >= MAX_SCANNED_ENTRIES: + return python_files, test_files, True + directory = pending.popleft() + scanned += 1 + try: + with os.scandir(directory) as scan: + items = sorted(scan, key=lambda entry: entry.name) + except OSError: + continue + entries_seen += len(items) + for entry in items: + name = entry.name + try: + if entry.is_dir(follow_symlinks=False): + if name not in _IGNORED_NAMES: + pending.append(entry.path) + continue + if name.endswith(".py") and entry.is_file(follow_symlinks=False): + python_files += 1 + if name.startswith("test_"): + test_files += 1 + except OSError: + continue + return python_files, test_files, False + + def _is_within(child: Path, parent: Path) -> bool: try: child.relative_to(parent) @@ -70,5 +160,4 @@ def _is_within(child: Path, parent: Path) -> bool: return False -def _is_countable(path: Path) -> bool: - return path.is_file() and not any(part in _IGNORED_NAMES for part in path.parts) +__all__ = ["MAX_SCANNED_DIRS", "MAX_SCANNED_ENTRIES", "build_workspace_snapshot", "count_python_files"] diff --git a/src/server/agent_server.py b/src/server/agent_server.py index b6a703e5..c7c313ba 100644 --- a/src/server/agent_server.py +++ b/src/server/agent_server.py @@ -404,7 +404,10 @@ async def _handle_control_request(self, msg: dict) -> None: # refuse (critic C8). ``interrupt`` is exempt (a benign abort of a # non-existent turn). Mirrors the ``session not ready`` pattern the # permission-mode handlers already use. - if self.init_error is not None and subtype != "interrupt": + # ``get_activity`` is exempt too: a runtime that refused to start + # has no work, and answering "refused" would read as busy — leaving + # it unreleasable by a conditional close. + if self.init_error is not None and subtype not in ("interrupt", "get_activity"): self._reply(request_id, {"ok": False, "error": self.init_error}) return if subtype == "interrupt": @@ -800,6 +803,9 @@ async def _handle_control_request(self, msg: dict) -> None: if subtype == "resume": self._do_resume(request_id, inner.get("session_id")) return + if subtype == "get_activity": + self._reply(request_id, self._activity_snapshot()) + return if subtype == "branch": self._do_branch(request_id) return @@ -3096,6 +3102,41 @@ def _do_branch(self, request_id: object) -> None: logger.exception("[agent-server] branch failed") self._reply(request_id, {"ok": False, "error": str(exc)}) + def _activity_snapshot(self) -> dict[str, Any]: + """What this session is doing that a gateway cannot see from outside. + + A gateway knows the turns IT submitted. It does not know about a + /goal continuation the agent queued for itself, a /loop or cron job + waiting to fire, a prompt still in the inbox, or a background shell + the agent started — all of which an idle-looking session may carry, + and all of which ``session.close`` with ``if_idle`` must not kill. + """ + with self._lock: + turn_active = self._current_abort is not None + goal_active = self._goal_mgr is not None and self._goal_mgr.is_active() + queued = not self._inbox.empty() + scheduled = False + try: + scheduled = bool(self.cron_scheduler.list_jobs()) or self.cron_scheduler.wakeup_info() is not None + except Exception: # noqa: BLE001 — the scheduler is optional state + logger.debug("[agent-server] scheduler probe failed", exc_info=True) + background = False + try: + background = self._bgtasks is not None and any( + getattr(task, "status", "") == "running" for task in self._bgtasks.list() + ) + except Exception: # noqa: BLE001 + logger.debug("[agent-server] background task probe failed", exc_info=True) + return { + "ok": True, + "turn_active": turn_active, + "queued": queued, + "goal_active": goal_active, + "scheduled": scheduled, + "background": background, + "busy": turn_active or queued or goal_active or scheduled or background, + } + def _do_resume(self, request_id: object, session_id: object) -> None: """Load a saved conversation into this session (the original's /resume). Idle-only — replacing the conversation mid-turn would race the worker.""" diff --git a/src/server/desktop_gateway.py b/src/server/desktop_gateway.py index 91fe0c43..e2ed0f1a 100644 --- a/src/server/desktop_gateway.py +++ b/src/server/desktop_gateway.py @@ -51,6 +51,16 @@ async def send_event( await _send_json(websocket, {"method": "event", "params": params}) +#: Pure reads that answer beside whatever the socket is doing, instead of +#: queueing behind it. Everything else runs one call at a time per socket, +#: in order — a prompt after an interrupt, a close after a submit — which is +#: what a client relying on the answer order needs. These two are file reads +#: off the event loop: the sidebar click's transcript must not sit behind +#: the runtime attach the previous click started, or the teardown of the +#: session being left. +CONCURRENT_METHODS = frozenset({"session.history", "projects.tree"}) + + async def handle_gateway_socket(websocket: WebSocket, state: DesktopServeState) -> None: """Accept one gateway socket and pump it until disconnect.""" await websocket.accept() @@ -59,6 +69,7 @@ async def handle_gateway_socket(websocket: WebSocket, state: DesktopServeState) conn = GatewayConnection(websocket=websocket, state=state) await conn.on_open() + concurrent: set[asyncio.Task] = set() try: while True: try: @@ -69,8 +80,15 @@ async def handle_gateway_socket(websocket: WebSocket, state: DesktopServeState) continue if not isinstance(frame, dict): continue + if frame.get("method") in CONCURRENT_METHODS: + task = asyncio.create_task(_dispatch(conn, frame)) + concurrent.add(task) + task.add_done_callback(concurrent.discard) + continue await _dispatch(conn, frame) finally: + for task in concurrent: + task.cancel() await conn.on_close() @@ -92,6 +110,8 @@ async def _dispatch(conn: Any, frame: dict[str, Any]) -> None: try: result = await handler(params) + except asyncio.CancelledError: + raise except Exception as exc: # noqa: BLE001 — one bad call must not drop the socket logger.warning("gateway: %s failed", method, exc_info=True) if request_id is not None: diff --git a/src/server/desktop_gateway_methods.py b/src/server/desktop_gateway_methods.py index d11c703b..3929ea80 100644 --- a/src/server/desktop_gateway_methods.py +++ b/src/server/desktop_gateway_methods.py @@ -142,6 +142,23 @@ def __init__(self, session_id: str, state: DesktopServeState) -> None: self.pump_task: asyncio.Task | None = None self.init_info: dict[str, Any] = {} self.init_seen = asyncio.Event() + # The saved session this runtime replays (``session.resume``), or None + # for one created fresh. ``state.sessions`` is keyed by RUNTIME id and + # a resumed row gets a fresh runtime, so this is how a second click on + # the same row finds the runtime it already has instead of spawning + # another. + self.stored_id: str | None = None + # Set once ``_create`` has finished attaching this runtime (spawned, + # capability-negotiated, stored conversation loaded) — or given up. + # A concurrent resume of the same row waits on it rather than racing + # a second spawn. + self.attached = asyncio.Event() + # A prompt is being answered: set on submit, cleared by the turn's + # ``result`` frame. What ``session.close`` with ``if_idle`` refuses on. + self.turn_active = False + # The agent's frame stream ended on an error: nothing runs, nothing + # can answer a control query, and a conditional close need not ask. + self.dead = False # My queries INTO the agent (control_request → control_response). self._pending_control: dict[str, asyncio.Future] = {} # The agent's asks OF the user (can_use_tool …), keyed by request_id; @@ -164,6 +181,10 @@ def __init__(self, session_id: str, state: DesktopServeState) -> None: # in between outranks it. self.user_titled = False self.sockets: set[WebSocket] = set() + # The sockets that asked for this runtime (created or resumed it) and + # have not let go: a conditional close from one of them is refused + # while another still holds it — a desktop tile, another window. + self.holders: set[WebSocket] = set() # Scheduled session.info refreshes; held so they aren't GC'd mid-flight. self._background: set[asyncio.Task] = set() # tool_use_id -> tool name, so a tool_result can label its row. @@ -201,14 +222,28 @@ async def shutdown(self) -> None: self._background.clear() if self.pump_task is not None: self.pump_task.cancel() + # Answer every waiting control query with "no reply" rather than + # cancelling it: a cancelled future raises CancelledError inside the + # RPC handler awaiting it, which unwinds the gateway's socket loop and + # drops that window's socket — the very window a retirement or + # another window's close should merely inform. None is the answer + # every caller already degrades on. Before the agent's own shutdown, + # so a handler on another socket unblocks now, not after the worker + # join. + for fut in self._pending_control.values(): + if not fut.done(): + fut.set_result(None) + self._pending_control.clear() + # Dead from here on: a query issued during or after the agent's own + # shutdown (a scheduled info refresh, a resume's tail on a runtime + # another window closed) answers None at once instead of waiting the + # control timeout on a pump that is already cancelled. + self.dead = True if self.agent is not None: try: await self.agent.shutdown() except Exception: # noqa: BLE001 pass - for fut in self._pending_control.values(): - if not fut.done(): - fut.cancel() # ── broadcast ──────────────────────────────────────────────────────────── @@ -229,7 +264,9 @@ async def publish_session_info(self, **extra: Any) -> None: deadlock until the timeout. Schedule it with ``refresh_session_info``. """ settings = await self.control_query("get_settings", {}) - payload: dict[str, Any] = {"running": False, "desktop_contract": DESKTOP_CONTRACT} + # Not a hardcoded False: this republish follows a resume reply that + # may have said running, and the desktop reads it as the turn state. + payload: dict[str, Any] = {"running": self.turn_active, "desktop_contract": DESKTOP_CONTRACT} if isinstance(settings, dict): model = settings.get("fusion") or settings.get("model") if model: @@ -277,10 +314,50 @@ async def _pump(self) -> None: raise except Exception: # noqa: BLE001 logger.exception("desktop session %s pump died", self.session_id) + # The turn's end first, then the retirement: the teardown it + # schedules cancels THIS task, so every send happens before it. await self._broadcast( "message.complete", {"text": "", "status": "error", "error": "backend session ended unexpectedly"}, ) + await self._retire() + return + # The agent closed its stream on its own: the runtime is over. + await self._retire() + + async def _retire(self) -> None: + """The agent's stream is gone: unregister this runtime and tear it down. + + Left registered, a dead runtime would be handed back by the reuse + lookup on the next click and every prompt to it would fail; left + "busy", it could never be released. So it leaves the registry here, + every window learns (``session.closed``), and the teardown runs as + its own task — this is the pump's task, which the teardown cancels. + """ + self.dead = True + self.turn_active = False + self._pending_asks.clear() + self._pending_question = None + # An attach waiting for system/init must not wait the full control + # timeout for a stream that will never send it. + self.init_seen.set() + if self.state.sessions.get(self.session_id) is not self: + return + self.state.sessions.pop(self.session_id, None) + await self._broadcast("session.closed", {}) + task = asyncio.create_task(self._teardown()) + self.state.teardowns.add(task) + task.add_done_callback(self.state.teardowns.discard) + + async def _teardown(self) -> None: + try: + await self.shutdown() + except Exception: # noqa: BLE001 — a failed teardown must not raise into the loop + logger.debug("desktop session %s: shutdown failed", self.session_id, exc_info=True) + try: + await self.state.manager.stop_session(self.session_id) + except Exception: # noqa: BLE001 — index upkeep is best-effort + pass async def _route(self, frame: dict[str, Any]) -> None: kind = frame.get("type") @@ -305,9 +382,15 @@ async def _route(self, frame: dict[str, Any]) -> None: await self._broadcast("session.info", _init_session_info(frame)) return + # A turn the agent started for itself (a /goal continuation, a loop + # firing) has no submit_prompt; its frames are how the gateway learns + # it is running. + if kind in ("assistant", "stream_event"): + self.turn_active = True for type_, payload in translate_frame(frame, self._tool_names): await self._broadcast(type_, payload) if kind == "result": + self.turn_active = False # The renderer clears busy from message.complete; refresh the info # line too (permission mode may have flipped server-side). mode = frame.get("permission_mode") @@ -437,9 +520,17 @@ async def _upgrade_title(self, text: str, fallback: str) -> None: return await self._apply_title(name) + @property + def idle(self) -> bool: + """No turn running and nothing waiting on the user (an approval, a question).""" + if self.dead: + return True + return not self.turn_active and not self._pending_asks and self._pending_question is None + async def submit_prompt(self, text: str) -> None: if not self.ready: raise ValueError(f"session {self.session_id} is still starting") + self.turn_active = True await self._broadcast("message.start", {}) await self.agent.send_to_agent( {"type": "user", "message": {"role": "user", "content": text}} @@ -460,11 +551,12 @@ async def interrupt(self) -> None: async def control_query(self, subtype: str, params: dict[str, Any], timeout: float = CONTROL_TIMEOUT_S) -> Any: - # No agent to ask yet: answer the way a timeout does. Every caller - # already handles "no reply" (`if not isinstance(result, dict)`) by - # degrading to its own fallback, which is the honest outcome — the - # alternative was an AttributeError that failed the whole RPC. - if not self.ready: + # No agent to ask yet — or no agent left to answer: answer the way a + # timeout does. Every caller already handles "no reply" (`if not + # isinstance(result, dict)`) by degrading to its own fallback, which + # is the honest outcome — the alternative was an AttributeError that + # failed the whole RPC, or a 30 s wait on a stream that has ended. + if not self.ready or self.dead: return None rid = f"srv-{uuid.uuid4().hex[:12]}" fut: asyncio.Future = asyncio.get_running_loop().create_future() @@ -955,6 +1047,49 @@ def _clean(value: Any) -> str | None: return None +def _prepare_workspace(path: str, create: bool) -> str: + """An absolute directory a session can run in, created when asked. + + Raises ``ValueError`` — the gateway's "this is the caller's mistake" error + — for a relative path, a path that is a file, or a folder that does not + exist when ``create`` is off. Runs off the event loop (``makedirs`` and + the stats can block on a slow volume). + """ + import os + + expanded = os.path.expanduser(path.strip()) + if not os.path.isabs(expanded): + raise ValueError(f"workspace path must be absolute: {path}") + normalized = os.path.normpath(expanded) + if os.path.isdir(normalized): + return normalized + if os.path.exists(normalized): + raise ValueError(f"not a directory: {normalized}") + if not create: + raise ValueError(f"no such directory: {normalized}") + try: + os.makedirs(normalized, exist_ok=True) + except OSError as exc: + raise ValueError(f"cannot create {normalized}: {exc.strerror or exc}") from exc + return normalized + + +def _create_session_worktree(cwd: str) -> dict[str, Any]: + """A fresh worktree of the repo at ``cwd`` (the CLI's bare ``--worktree``).""" + from src.utils.worktree_session import WorktreeError, create_worktree_for_session + + try: + session = create_worktree_for_session(None, cwd=cwd) + except WorktreeError as exc: + raise ValueError(str(exc)) from exc + return { + "name": session.worktree_name, + "path": session.worktree_path, + "branch": session.worktree_branch, + "repo_root": session.repo_root, + } + + def _positive_int(value: Any, fallback: int) -> int: """A positive integer from a JSON field, or ``fallback``. @@ -986,6 +1121,7 @@ def __init__(self, websocket: WebSocket, state: DesktopServeState) -> None: self.method_handlers = { "session.create": self.session_create, "session.resume": self.session_resume, + "session.history": self.session_history, "session.activate": self.session_activate, "session.close": self.session_close, "session.active_list": self.session_active_list, @@ -1043,6 +1179,7 @@ async def on_open(self) -> None: async def on_close(self) -> None: for session in self.state.sessions.values(): session.sockets.discard(self.websocket) + session.holders.discard(self.websocket) # ── helpers ────────────────────────────────────────────────────────────── @@ -1062,18 +1199,48 @@ def _session(self, params: dict[str, Any]) -> DesktopSession: raise ValueError(f"session {session_id} is still starting") return session + def _live_session_for(self, stored_id: str) -> DesktopSession | None: + """The runtime this server already has for ``stored_id``, if any. + + Either the runtime itself (a session created here saves under its own + id) or the fresh runtime a resume of that row spawned (``stored_id``). + Without the second match every click on a row the server had already + replayed spawned yet another runtime and left the previous one alive. + """ + direct = self.state.sessions.get(stored_id) + if direct is not None and not direct.dead: + return direct + for session in self.state.sessions.values(): + if session.stored_id == stored_id and not session.dead: + return session + return None + async def _create(self, cwd: str | None, resume: str | None, params: dict[str, Any] | None = None) -> DesktopSession: manager = self.state.manager workspace = cwd or self.state.workspace - if resume and resume in self.state.sessions: - return self.state.sessions[resume] + if resume: + existing = self._live_session_for(resume) + if existing is not None: + # A second resume while the first is still attaching waits + # for it — two rapid clicks must not become two runtimes. + # Uncapped: the event is set on every exit of _create. + await existing.attached.wait() + # Still registered: the attach succeeded. This socket may + # not be the one that spawned it — a second window opening + # the same row — and it needs the turn events too. + if existing.session_id in self.state.sessions: + existing.sockets.add(self.websocket) + existing.holders.add(self.websocket) + return existing # A resumed stored session still gets a fresh runtime session: spawn, # then load the stored conversation via the `resume` control below. info = manager.create_session(cwd=workspace) session_id = info.id session = DesktopSession(session_id, self.state) + session.stored_id = resume session.sockets.add(self.websocket) + session.holders.add(self.websocket) self.state.sessions[session_id] = session # Honor the composer's provider/model/effort selection at spawn time, so # a session can use a working provider even when the config default is @@ -1109,7 +1276,24 @@ async def _create(self, cwd: str | None, resume: str | None, # sessionless call (model.options, commands.catalog …) from any # window picks it up. Re-raised untouched; this only cleans up. self.state.sessions.pop(session_id, None) + session.attached.set() raise + try: + await self._attach(session, resume, params) + finally: + session.attached.set() + # The agent's stream ended while the runtime was being attached: it + # has already retired itself, and handing its id out would only make + # the next prompt fail. Say so — the caller's socket stays up. + if session.dead: + raise ValueError(f"session {session_id} ended while starting") + return session + + async def _attach(self, session: DesktopSession, resume: str | None, + params: dict[str, Any]) -> None: + """Finish a freshly spawned runtime: init, capabilities, stored replay.""" + manager = self.state.manager + session_id = session.session_id try: manager.mark_running(session_id) except Exception: # noqa: BLE001 — index upkeep is best-effort @@ -1124,7 +1308,7 @@ async def _create(self, cwd: str | None, resume: str | None, # a client that ignored the question would park the session's worker # thread until the ask timeout. Clients that render questions ask for # the real thing here; the ones that don't are left exactly as before. - if _wants_questions(params): + if _wants_questions(params) and not session.dead: reply = await session.control_query("set_ask_user_interactive", {"enabled": True}) if isinstance(reply, dict) and reply.get("ok") is not False: session.asks_questions = True @@ -1132,7 +1316,7 @@ async def _create(self, cwd: str | None, resume: str | None, logger.warning( "session %s: interactive questions refused: %r", session_id, reply ) - if resume: + if resume and not session.dead: # A stored session brings its own name; auto-titling would rename # someone's saved conversation after whatever they type next. session.titled = True @@ -1140,7 +1324,9 @@ async def _create(self, cwd: str | None, resume: str | None, if not isinstance(reply, dict) or reply.get("ok") is False: logger.warning("session %s: resume of %s refused: %r", session_id, resume, reply) - return session + # Not a replay of that row after all: the next click must + # try again rather than reuse an empty conversation. + session.stored_id = None # ── methods ────────────────────────────────────────────────────────────── @@ -1153,6 +1339,9 @@ async def projects_tree(self, params: dict[str, Any]) -> dict[str, Any]: return await _asyncio.to_thread(self._build_projects_tree, preview_limit) def _build_projects_tree(self, preview_limit: int) -> dict[str, Any]: + import os + from concurrent.futures import ThreadPoolExecutor + from src.server.desktop_projects import build_project_tree, canonical_workspace_path from src.server.desktop_sessions import list_session_rows from src.utils.git import get_repo_root, list_worktrees @@ -1163,7 +1352,9 @@ def _build_projects_tree(self, preview_limit: int) -> dict[str, Any]: # just-created session still places into its repo immediately. seen = {r.get("id") for r in rows} active_cwd: str | None = None - for session in self.state.sessions.values(): + # A copy: this runs on a worker thread while the loop adds and + # removes sessions (projects.tree is served beside a resume). + for session in list(self.state.sessions.values()): info = getattr(session, "init_info", None) or {} cwd = info.get("cwd") if isinstance(info, dict) else None if cwd: @@ -1177,19 +1368,38 @@ def _build_projects_tree(self, preview_limit: int) -> dict[str, Any]: }) seen.add(sid) - # Per-cwd / per-repo memoized git probes: a tree can hold many sessions - # in the same repo, so probe each distinct path once. - repo_cache: dict[str, str | None] = {} - wt_cache: dict[str, list[str]] = {} + # Git probes are memoized ACROSS rebuilds on the serve state: the + # tree is rebuilt after every turn end and every session switch, and + # a sessions dir accumulates thousands of distinct cwds, so probing + # each one every time cost seconds per rebuild. + cache = self.state.probe_cache workspace_cache: dict[str, str | None] = {} def worktrees_of(repo_root: str) -> list[str]: # ``git worktree list`` is repo-global and main-first from ANY # worktree in the repo, so this is correct whether keyed by the # main root or a linked-worktree path. - if repo_root not in wt_cache: - wt_cache[repo_root] = [w.path for w in list_worktrees(repo_root) if w.path] - return wt_cache[repo_root] + cached = cache.worktrees(repo_root) + if cached is None: + cached = [w.path for w in list_worktrees(repo_root) if w.path] + cache.set_worktrees(repo_root, cached) + return cached + + def probe_toplevel(cwd: str) -> str | None: + # A cwd that is gone — a deleted temp dir, an unmounted volume — + # is not a repo, and needs no git call to say so. + if not os.path.isdir(cwd): + return None + return get_repo_root(cwd) or None + + # Every cwd the cache cannot answer, probed in parallel: each is one + # subprocess, and the first tree after boot has all of them to do. + wanted = {c for c in (str(r.get("cwd") or "").strip() for r in rows) if c} + missing = [c for c in wanted if not cache.has_repo_root(c)] + if missing: + with ThreadPoolExecutor(max_workers=min(8, len(missing))) as pool: + for cwd, top in zip(missing, pool.map(probe_toplevel, missing)): + cache.set_repo_root(cwd, top) def repo_root_of(cwd: str) -> str | None: # ``rev-parse --show-toplevel`` inside a LINKED worktree returns the @@ -1197,14 +1407,13 @@ def repo_root_of(cwd: str) -> str | None: # worktree into its own project. Resolve the MAIN worktree root # (the first ``git worktree list`` entry) so linked worktrees group # as lanes under their repo, matching the renderer's tree. - if cwd not in repo_cache: - top = get_repo_root(cwd) - if top: - worktrees = worktrees_of(top) - repo_cache[cwd] = worktrees[0] if worktrees else top - else: - repo_cache[cwd] = None - return repo_cache[cwd] + if not cache.has_repo_root(cwd): + cache.set_repo_root(cwd, probe_toplevel(cwd)) + top = cache.repo_root(cwd) + if not top: + return None + worktrees = worktrees_of(top) + return worktrees[0] if worktrees else top def workspace_path_of(cwd: str) -> str | None: if cwd not in workspace_cache: @@ -1221,12 +1430,109 @@ def workspace_path_of(cwd: str) -> str | None: ) async def session_create(self, params: dict[str, Any]) -> dict[str, Any]: - session = await self._create(params.get("cwd"), None, params) - return { + """Spawn a fresh session. + + ``cwd`` names the workspace; with ``create_dir`` a folder that does + not exist yet is created (the "new workspace" flow, mkdir -p). With + ``worktree`` the session runs in a fresh git worktree of that repo + (``.clawcodex/worktrees/``, the CLI's ``--worktree``), which is + left in place when the session ends — a browser tab has no exit + dialog to offer keep-or-remove. Both are validated here so a bad path + is an error with a name rather than a runtime that fails to start. A + plain ``cwd`` passes through as it always has: the runtime, not the + gateway, is the judge of a workspace it is merely pointed at. + """ + cwd = _clean(params.get("cwd")) + wants_worktree = params.get("worktree") is True + create_dir = params.get("create_dir") is True + if cwd is not None and (create_dir or wants_worktree): + import os + + if wants_worktree and not os.path.isdir(os.path.expanduser(cwd)): + # Checked before anything is created: a folder made here + # would never be a repository, and a worktree refusal after + # the mkdir would leave an empty folder behind. + raise ValueError( + "A worktree needs an existing git repository; a new folder is " + "not one. Turn Worktree off, or pick a repository." + ) + cwd = await asyncio.to_thread(_prepare_workspace, cwd, create_dir) + worktree: dict[str, Any] | None = None + if wants_worktree: + worktree = await asyncio.to_thread( + _create_session_worktree, cwd or self.state.workspace + ) + cwd = worktree["path"] + # The next sidebar tree must show the new lane. Every repo's list, + # not just this one's: the tree keys worktree lists by whatever + # ``git rev-parse`` printed for a cwd, which need not spell the + # repo root the way the worktree helper does (symlinked /tmp …). + self.state.probe_cache.forget_worktrees() + session = await self._create(cwd, None, params) + reply: dict[str, Any] = { "session_id": session.session_id, "stored_session_id": session.session_id, "info": _init_session_info(session.init_info), } + if worktree is not None: + reply["worktree"] = worktree + return reply + + async def session_history(self, params: dict[str, Any]) -> dict[str, Any]: + """A saved session's transcript, cold: no runtime is spawned or touched. + + What the web client renders the moment a sidebar row is clicked; the + runtime attaches afterwards (``session.resume``) without the reader + waiting on it. Same message shape as ``session.resume`` returns. + + When this server already has a runtime replaying the row, its OWN + record is preferred — it holds the turns run since the row was + resumed, which the row's file never learns about. + """ + wanted = str(params.get("session_id") or "") + if not wanted: + raise ValueError("session_id required") + from src.server.desktop_sessions import load_session_messages + + sessions_dir = self.state.saved_sessions_dir() + live = self._live_session_for(wanted) + stored = None + if live is not None and live.session_id != wanted: + stored = await asyncio.to_thread(load_session_messages, sessions_dir, live.session_id) + if stored is None: + stored = await asyncio.to_thread(load_session_messages, sessions_dir, wanted) + if stored is None: + if live is None: + raise ValueError(f"unknown session: {wanted}") + # A runtime that never saved (nothing typed since it was made): + # there is nothing to show, and that is not an error. + info = _init_session_info(live.init_info) + info["running"] = live.turn_active + return { + "session_id": wanted, + "stored_session_id": wanted, + "live_session_id": live.session_id, + "found": False, + "messages": [], + "message_count": 0, + "info": info, + } + info: dict[str, Any] = {key: stored[key] for key in ("cwd", "model", "provider") if stored.get(key)} + if live is not None: + info["running"] = live.turn_active + reply: dict[str, Any] = { + "session_id": wanted, + "stored_session_id": wanted, + "found": True, + "messages": stored["messages"], + "message_count": stored["message_count"], + "info": info, + } + if stored.get("title"): + reply["title"] = stored["title"] + if live is not None: + reply["live_session_id"] = live.session_id + return reply async def session_resume(self, params: dict[str, Any]) -> dict[str, Any]: wanted = str(params.get("session_id") or "") or None @@ -1238,6 +1544,11 @@ async def session_resume(self, params: dict[str, Any]) -> dict[str, Any]: # this, so the user reads a switch that held as a revert. Re-read the # live settings and tell the truth, here where both paths converge. settings = await session.control_query("get_settings", {}) + # The stream can end during that round trip, or another window can + # close the runtime outright; a dead or unregistered id handed out + # here would only fail the next prompt. + if session.dead or self.state.sessions.get(session.session_id) is not session: + raise ValueError(f"session {session.session_id} ended while starting") if isinstance(settings, dict): live_model = settings.get("fusion") or settings.get("model") if live_model: @@ -1245,19 +1556,36 @@ async def session_resume(self, params: dict[str, Any]) -> dict[str, Any]: if settings.get("provider"): session.init_info["provider"] = str(settings["provider"]) session.refresh_session_info() + info = _init_session_info(session.init_info) + # Reattaching to a runtime mid-turn is a normal path now (a busy + # runtime is kept and handed back); the client must adopt it as + # running, not draw the turn's remaining deltas under its next prompt. + info["running"] = session.turn_active + # The row this runtime replays — which is NOT the row asked for when + # the agent refused the replay (stored_id was cleared): echoing the + # asked-for row would make the client keep an empty runtime as the + # conversation on screen. + stored_row = (session.stored_id or session.session_id) if wanted else session.session_id response: dict[str, Any] = { "session_id": session.session_id, - "stored_session_id": wanted or session.session_id, + "stored_session_id": stored_row, "resumed": wanted or session.session_id, "message_count": 0, "messages": [], - "info": _init_session_info(session.init_info), + "info": info, } omit = bool(params.get("omit_messages") or params.get("lazy")) if wanted and not omit: from src.server.desktop_sessions import load_session_messages - stored = load_session_messages(self.state.saved_sessions_dir(), wanted) + sessions_dir = self.state.saved_sessions_dir() + stored = None + # A runtime already replaying this row has the complete record + # (see session_history); the row's own file is the fallback. + if session.session_id != wanted: + stored = await asyncio.to_thread(load_session_messages, sessions_dir, session.session_id) + if stored is None: + stored = await asyncio.to_thread(load_session_messages, sessions_dir, wanted) if stored is not None: response["messages"] = stored["messages"] response["message_count"] = stored["message_count"] @@ -1273,15 +1601,68 @@ async def session_activate(self, params: dict[str, Any]) -> dict[str, Any]: return await self.session_create(params) async def session_close(self, params: dict[str, Any]) -> dict[str, Any]: + """Shut a runtime down. + + ``if_idle`` closes only a runtime with no turn running and nothing + waiting on the user — the web client's "let go of the session I am + leaving" — and answers ``closed: False`` otherwise. Checked here, on + the server's knowledge, because a client that just attached to a + runtime another window is driving does not know it is busy. + """ session_id = str(params.get("session_id") or "") + if params.get("if_idle") is True: + session = self.state.sessions.get(session_id) + if session is None: + return {"ok": True, "closed": False, "reason": "gone"} + # This socket lets go; anyone else still holding it keeps it — + # a desktop tile, another window that opened the same row. + session.holders.discard(self.websocket) + if session.holders: + return {"ok": True, "closed": False, "reason": "held"} + if not await self._session_is_idle(session): + return {"ok": True, "closed": False, "reason": "busy"} + # Re-checked after the await: a prompt from another socket can + # land while the agent answers, a replaced entry is not the + # runtime that was asked about, and a holder may have arrived. + if self.state.sessions.get(session_id) is not session: + return {"ok": True, "closed": False, "reason": "gone"} + if session.holders: + return {"ok": True, "closed": False, "reason": "held"} + if not session.idle: + return {"ok": True, "closed": False, "reason": "busy"} session = self.state.sessions.pop(session_id, None) if session is not None: - await session.shutdown() - try: - await self.state.manager.stop_session(session_id) - except Exception: # noqa: BLE001 — index upkeep is best-effort - pass - return {"ok": True} + # Every window on this runtime learns it is gone, so the next + # prompt reconnects instead of failing with "unknown session". + await session._broadcast("session.closed", {}) + # The teardown (SessionEnd hooks, a bounded worker join) runs + # behind the reply: calls on this socket are served in order, and + # the call after a close is the open of the session being moved + # to — it must not wait seconds on the one being left. + task = asyncio.create_task(session._teardown()) + self.state.teardowns.add(task) + task.add_done_callback(self.state.teardowns.discard) + return {"ok": True, "closed": session is not None} + + async def _session_is_idle(self, session: DesktopSession) -> bool: + """Idle by the gateway's own knowledge AND the agent's. + + The gateway sees the turns it submitted and the asks it relayed. The + agent also knows about work it started for itself — a /goal + continuation, a /loop or cron job waiting to fire, a queued prompt, a + background shell — which ``get_activity`` reports. No reply (an agent + too old for the control, or too busy to answer) reads as busy: the + cost of a wrong "idle" is killing someone's loop. + """ + if not session.idle: + return False + if session.dead: + return True + activity = await session.control_query("get_activity", {}, timeout=5.0) + if not isinstance(activity, dict) or activity.get("ok") is False: + return False + return activity.get("busy") is not True + async def session_active_list(self, _: dict[str, Any]) -> dict[str, Any]: sessions = [] diff --git a/src/server/desktop_projects.py b/src/server/desktop_projects.py index 0ead2e0e..6d12b797 100644 --- a/src/server/desktop_projects.py +++ b/src/server/desktop_projects.py @@ -31,11 +31,98 @@ from __future__ import annotations import os +import threading +import time from typing import Any, Callable NO_PROJECT_ID = "__no_project__" +class ProbeCache: + """Memoized git probes for the sidebar tree, shared across rebuilds. + + ``projects.tree`` shells out once per distinct session cwd (``git + rev-parse --show-toplevel``) and once per repo (``git worktree list``). A + sessions directory accumulates thousands of distinct cwds over time — + temp dirs, worktrees, test fixtures — so probing them all on every rebuild + cost seconds, and the tree is rebuilt after every turn end and every + session switch. A repo root answer is kept for ``ttl_s`` (a directory's + repo does not move); the worktree list for ``worktree_ttl_s``, and + :meth:`forget_worktrees` drops it early when this server adds one. + Thread-safe: the tree is built off the event loop and probes run in a + pool. + """ + + def __init__( + self, + *, + ttl_s: float = 300.0, + negative_ttl_s: float = 30.0, + worktree_ttl_s: float = 30.0, + clock: Callable[[], float] = time.monotonic, + ) -> None: + self._ttl_s = ttl_s + # "Not a repo" is kept for less time: a folder just created for a + # new workspace gets its ``git init`` moments later. + self._negative_ttl_s = negative_ttl_s + self._worktree_ttl_s = worktree_ttl_s + self._clock = clock + self._lock = threading.Lock() + # cwd → (expires_at, toplevel or None) + self._repo_root: dict[str, tuple[float, str | None]] = {} + # repo root → (expires_at, worktree paths, main first) + self._worktrees: dict[str, tuple[float, list[str]]] = {} + + def lookup(self, cwd: str) -> tuple[bool, str | None]: + """``(cached, toplevel)`` for ``cwd``, read against one clock sample. + + One sample, so an entry cannot be "there" for the found check and + "expired" for the value read a moment later. + """ + now = self._clock() + with self._lock: + entry = self._repo_root.get(cwd) + if entry is None or entry[0] <= now: + return False, None + return True, entry[1] + + def has_repo_root(self, cwd: str) -> bool: + return self.lookup(cwd)[0] + + def repo_root(self, cwd: str) -> str | None: + """The cached toplevel for ``cwd`` (None: not a repo, or not cached).""" + return self.lookup(cwd)[1] + + def set_repo_root(self, cwd: str, toplevel: str | None) -> None: + ttl = self._ttl_s if toplevel else self._negative_ttl_s + with self._lock: + self._repo_root[cwd] = (self._clock() + ttl, toplevel) + + def worktrees(self, repo_root: str) -> list[str] | None: + with self._lock: + entry = self._worktrees.get(repo_root) + if entry is None or entry[0] <= self._clock(): + return None + return list(entry[1]) + + def set_worktrees(self, repo_root: str, paths: list[str]) -> None: + with self._lock: + self._worktrees[repo_root] = (self._clock() + self._worktree_ttl_s, list(paths)) + + def forget_worktrees(self, repo_root: str | None = None) -> None: + """Drop the worktree list for one repo, or every repo's.""" + with self._lock: + if repo_root is None: + self._worktrees.clear() + else: + self._worktrees.pop(repo_root, None) + + def clear(self) -> None: + with self._lock: + self._repo_root.clear() + self._worktrees.clear() + + def canonical_workspace_path(path: str) -> str | None: """Return the real path for an existing directory, else ``None``. @@ -262,4 +349,4 @@ def _lane_for(repo_root: str, cwd: str) -> tuple[str, str, bool, str]: } -__all__ = ["build_project_tree", "canonical_workspace_path", "NO_PROJECT_ID"] +__all__ = ["ProbeCache", "build_project_tree", "canonical_workspace_path", "NO_PROJECT_ID"] diff --git a/src/server/desktop_serve.py b/src/server/desktop_serve.py index 10136342..fff62a4a 100644 --- a/src/server/desktop_serve.py +++ b/src/server/desktop_serve.py @@ -65,6 +65,11 @@ class DesktopServeState: sessions: dict[str, Any] = field(default_factory=dict) # Saved-transcript dir override (tests); default resolves per request. sessions_dir: Path | None = None + # Git probe answers the sidebar tree reuses across rebuilds. + probe_cache: Any = field(default_factory=lambda: _new_probe_cache()) + # Session teardowns in flight: session.close answers before they finish, + # and they are held here so they are not garbage-collected mid-shutdown. + teardowns: set[Any] = field(default_factory=set) def spawn_for(self, provider: str | None, model: str | None, effort: str | None) -> Callable[..., Awaitable[Any]]: @@ -103,6 +108,20 @@ async def shutdown(self) -> None: except Exception: # noqa: BLE001 — teardown must not raise pass self.sessions.clear() + # Teardowns a session.close left running: let them finish their + # SessionEnd hooks before the process goes. + pending = [task for task in self.teardowns if not task.done()] + if pending: + import asyncio + + await asyncio.gather(*pending, return_exceptions=True) + self.teardowns.clear() + + +def _new_probe_cache() -> Any: + from src.server.desktop_projects import ProbeCache + + return ProbeCache() def _token_ok(state: DesktopServeState, presented: str | None) -> bool: diff --git a/src/server/desktop_sessions.py b/src/server/desktop_sessions.py index ab415e2e..7882b1e9 100644 --- a/src/server/desktop_sessions.py +++ b/src/server/desktop_sessions.py @@ -14,6 +14,7 @@ import json import logging +import threading from pathlib import Path from typing import Any @@ -60,6 +61,23 @@ def _row_from_file(path: Path, data: dict[str, Any]) -> dict[str, Any]: } +# Sidebar rows by file path, keyed on the (mtime, size, inode) they were read +# at. A session file is a whole conversation (up to a few MB), and the sidebar +# tree is rebuilt after every turn end and every session switch: re-parsing +# thousands of unchanged files each time cost ~1 s per rebuild. The stamp is +# re-checked with one ``stat`` per file, so an edited, replaced or deleted +# file is never served stale — the inode catches an atomic replace of equal +# size within one mtime tick on a coarse filesystem. +_ROW_CACHE: dict[str, tuple[tuple[int, int, int], dict[str, Any]], ] = {} +_ROW_CACHE_LOCK = threading.Lock() + + +def clear_session_row_cache() -> None: + """Forget every cached sidebar row (tests, or a sessions-dir switch).""" + with _ROW_CACHE_LOCK: + _ROW_CACHE.clear() + + def list_session_rows( sessions_dir: Path, *, @@ -67,25 +85,47 @@ def list_session_rows( offset: int = 0, min_messages: int = 0, ) -> dict[str, Any]: - """Paginated sidebar listing, newest-first by file mtime.""" + """Paginated sidebar listing, newest-first by file mtime. + + Unchanged files come from :data:`_ROW_CACHE`; only a file whose + ``(mtime, size)`` moved since it was last read is parsed again. + """ + stamped: list[tuple[Path, tuple[int, int, int], float]] = [] try: - files = sorted( - sessions_dir.glob("*.json"), - key=lambda p: p.stat().st_mtime, - reverse=True, - ) + for path in sessions_dir.glob("*.json"): + try: + stat = path.stat() + except OSError: + continue # deleted between the listing and the stat + stamped.append((path, (stat.st_mtime_ns, stat.st_size, stat.st_ino), stat.st_mtime)) except OSError: - files = [] + stamped = [] + stamped.sort(key=lambda item: item[2], reverse=True) rows: list[dict[str, Any]] = [] - for path in files: - data = _read_session_file(path) - if data is None: - continue - row = _row_from_file(path, data) + seen: set[str] = set() + for path, stamp, _mtime in stamped: + key = str(path) + seen.add(key) + with _ROW_CACHE_LOCK: + cached = _ROW_CACHE.get(key) + if cached is not None and cached[0] == stamp: + row = cached[1] + else: + data = _read_session_file(path) + if data is None: + continue + row = _row_from_file(path, data) + with _ROW_CACHE_LOCK: + _ROW_CACHE[key] = (stamp, row) if row["message_count"] < min_messages: continue - rows.append(row) + # A copy per call: callers annotate rows (``is_active`` …) and the + # cached one must stay as the file said. + rows.append(dict(row)) + with _ROW_CACHE_LOCK: + for key in [k for k in _ROW_CACHE if k not in seen]: + _ROW_CACHE.pop(key, None) window = rows[offset : offset + limit] if limit > 0 else rows[offset:] return { @@ -224,12 +264,21 @@ def load_session_messages(sessions_dir: Path, session_id: str) -> dict[str, Any] # the sidebar reads it from this same file, and a header that disagreed # with the row the user just clicked is its own small confusion. name = data.get("name") - return { + result: dict[str, Any] = { "messages": messages, "message_count": len(messages), "session_id": str(data.get("session_id") or safe), "title": str(name) if isinstance(name, str) and name.strip() else "", } + # The facts a client needs to SHOW a stored session before any runtime + # exists for it: where it ran and what it ran on. Same keys as the + # ``session.info`` payload, so the header and the model chip read a cold + # transcript and a live one alike. + for key in ("cwd", "model", "provider"): + value = data.get(key) + if isinstance(value, str) and value: + result[key] = value + return result def _write_session_file(path: Path, data: dict[str, Any]) -> bool: diff --git a/tests/server/test_desktop_gateway.py b/tests/server/test_desktop_gateway.py index 171a7725..1f68c51b 100644 --- a/tests/server/test_desktop_gateway.py +++ b/tests/server/test_desktop_gateway.py @@ -16,6 +16,7 @@ import asyncio import contextlib import json +import time from pathlib import Path from types import SimpleNamespace from unittest.mock import patch @@ -25,6 +26,7 @@ from src.providers.base import ChatResponse from src.server.agent_server import AgentServerConfig, make_spawn_agent +from src.server.desktop_gateway_methods import CONTROL_TIMEOUT_S from src.server.desktop_serve import DesktopServeState, build_app from src.server.session_manager import SessionManager @@ -96,6 +98,20 @@ def __init__(self) -> None: # saved as the default for new sessions. Tests flip it to model a # session that may not write the host's settings. self.persist_preferences = True + # Leave a user message unanswered (the turn stays open until the test + # pushes its own ``result`` frame). + self.hold_turns = False + # What ``get_activity`` reports beyond the gateway's own view: a goal + # continuation, a loop, a background shell. + self.busy = False + # A shutdown that takes a while (SessionEnd hooks, a worker join). + self.shutdown_delay_s = 0.0 + # A slow ``get_activity`` answer, for the close's re-check. + self.activity_delay_s = 0.0 + # A slow ``get_settings`` answer, for a death during the resume's tail. + self.settings_delay_s = 0.0 + # Die before ever announcing system/init. + self.die_before_init = False async def send_to_agent(self, frame: dict) -> None: self.inbound.append(frame) @@ -105,6 +121,10 @@ async def send_to_agent(self, frame: dict) -> None: reply: dict | None = None if subtype == "resume": reply = {"ok": True} + elif subtype == "get_activity": + if self.activity_delay_s: + await asyncio.sleep(self.activity_delay_s) + reply = {"ok": True, "busy": self.busy} elif subtype == "set_model": # Record the switch so get_settings reports the new state. self.model = request.get("model") or self.model @@ -154,6 +174,8 @@ async def send_to_agent(self, frame: dict) -> None: else: reply = {"ok": False, "error": "usage: /recap [on|off|status]"} elif subtype == "get_settings": + if self.settings_delay_s: + await asyncio.sleep(self.settings_delay_s) reply = { "model": self.model, "provider": self.provider, @@ -174,6 +196,8 @@ async def send_to_agent(self, frame: dict) -> None: ) return if frame.get("type") == "user": + if self.hold_turns: + return # One scripted streamed turn per user message. await self.queue.put( { @@ -205,6 +229,8 @@ async def send_to_agent(self, frame: dict) -> None: ) async def messages_from_agent(self): + if self.die_before_init: + raise RuntimeError("stream died before init") yield { "type": "system", "subtype": "init", @@ -213,21 +239,29 @@ async def messages_from_agent(self): "model": "fake", } while True: - yield await self.queue.get() + item = await self.queue.get() + # An exception queued by a test stands for the agent stream dying. + if isinstance(item, BaseException): + raise item + yield item async def shutdown(self) -> None: + if self.shutdown_delay_s: + await asyncio.sleep(self.shutdown_delay_s) self.shutdown_called = True class FakeManager: def __init__(self) -> None: self.created: list[str] = [] + self.cwds: list[str] = [] self._n = 0 def create_session(self, cwd: str): self._n += 1 session_id = f"fake-{self._n}" self.created.append(session_id) + self.cwds.append(cwd) return SimpleNamespace(id=session_id, cwd=cwd) def mark_running(self, session_id: str) -> None: @@ -1047,3 +1081,683 @@ def test_effort_change_round_trips_and_reports_persisted(tmp_path: Path) -> None _rpc(ws, 5, "config.set", {"session_id": sid, "key": "effort", "value": "bogus"}) result = _drain_for_response(ws, 5, events)["result"] assert result["ok"] is False and "invalid effort" in result["error"] + + +# ─── opening a saved session: cold history, runtime reuse ──────────────────── + + +def _write_saved(sessions_dir: Path, session_id: str, messages: list, **extra) -> None: + sessions_dir.mkdir(exist_ok=True) + payload = { + "session_id": session_id, + "preview": "hello?", + "message_count": len(messages), + "cwd": "/tmp/where", + "model": "m-stored", + "provider": "p-stored", + "conversation": {"messages": messages}, + **extra, + } + (sessions_dir / f"{session_id}.json").write_text(json.dumps(payload), encoding="utf-8") + + +def test_session_history_reads_the_saved_transcript_without_a_runtime(tmp_path: Path) -> None: + """The sidebar click renders from this — no spawn, no control round-trip.""" + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved( + state.sessions_dir, "old-chat", + [{"role": "user", "content": "hello?"}, + {"role": "assistant", "content": [{"type": "text", "text": "hi back"}]}], + name="My chat", + ) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.history", {"session_id": "old-chat"}) + result = _drain_for_response(ws, 1, events)["result"] + + assert result["found"] is True + assert result["stored_session_id"] == "old-chat" + assert result["title"] == "My chat" + assert [m["role"] for m in result["messages"]] == ["user", "assistant"] + assert result["info"] == {"cwd": "/tmp/where", "model": "m-stored", "provider": "p-stored"} + assert agents == [] and state.manager.created == [] + + +def test_session_history_of_an_unknown_row_is_an_error(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + state.sessions_dir.mkdir() + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + _rpc(ws, 1, "session.history", {"session_id": "nope"}) + reply = _drain_for_response(ws, 1, []) + + assert "unknown session" in reply["error"]["message"] + + +def test_resuming_the_same_row_twice_reuses_its_runtime(tmp_path: Path) -> None: + """Every click used to spawn a fresh runtime and leave the last one alive. + + ``state.sessions`` is keyed by runtime id, and a resumed row's runtime has + a different id from the row, so the "already live" check never matched a + row. The second resume must come back with the same runtime, spawn-free. + """ + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "old-chat", [{"role": "user", "content": "hello?"}]) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.resume", {"session_id": "old-chat"}) + first = _drain_for_response(ws, 1, events)["result"] + _rpc(ws, 2, "session.resume", {"session_id": "old-chat", "omit_messages": True}) + second = _drain_for_response(ws, 2, events)["result"] + + assert first["session_id"] == "fake-1" + assert second["session_id"] == "fake-1" + assert second["stored_session_id"] == "old-chat" + assert second.get("messages_omitted") is True + assert len(agents) == 1 and state.manager.created == ["fake-1"] + assert state.sessions["fake-1"].stored_id == "old-chat" + + +def test_a_live_replay_answers_history_from_its_own_record(tmp_path: Path) -> None: + """Turns run after a resume are saved under the RUNTIME's id; the row the + user clicks still names the original file. Opening the row again must + show the conversation as it is now, not as the row's file left it.""" + state, _agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "old-chat", [{"role": "user", "content": "first"}]) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.resume", {"session_id": "old-chat", "omit_messages": True}) + runtime = _drain_for_response(ws, 1, events)["result"]["session_id"] + # The runtime saved its own, longer record after a turn. + _write_saved( + state.sessions_dir, runtime, + [{"role": "user", "content": "first"}, {"role": "user", "content": "second"}], + ) + _rpc(ws, 2, "session.history", {"session_id": "old-chat"}) + history = _drain_for_response(ws, 2, events)["result"] + _rpc(ws, 3, "session.resume", {"session_id": "old-chat"}) + resumed = _drain_for_response(ws, 3, events)["result"] + + assert history["live_session_id"] == runtime + assert [m["content"] for m in history["messages"]] == ["first", "second"] + assert [m["content"] for m in resumed["messages"]] == ["first", "second"] + + +def test_two_concurrent_resumes_of_one_row_share_a_spawn(tmp_path: Path) -> None: + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "old-chat", [{"role": "user", "content": "hello?"}]) + gate = asyncio.Event() + base_spawn = state.spawn_agent + + async def slow_spawn(session_id, cwd, resume): + await gate.wait() + return await base_spawn(session_id, cwd, resume) + + state.spawn_agent = slow_spawn + + async def run() -> tuple[dict, dict]: + from src.server.desktop_gateway_methods import GatewayConnection + + class _Socket: + async def send_json(self, obj): # pragma: no cover - no pushes read + pass + + conn = GatewayConnection(websocket=_Socket(), state=state) # type: ignore[arg-type] + first = asyncio.create_task(conn.session_resume({"session_id": "old-chat"})) + await asyncio.sleep(0) + second = asyncio.create_task(conn.session_resume({"session_id": "old-chat"})) + await asyncio.sleep(0) + gate.set() + return await first, await second + + a, b = asyncio.run(run()) + assert a["session_id"] == b["session_id"] == "fake-1" + assert len(agents) == 1 + + +# ─── session.create: a new folder, a worktree ──────────────────────────────── + + +def test_session_create_can_make_the_workspace_folder(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + target = tmp_path / "fresh" / "project" + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + # A plain cwd passes through untouched, as it always has (the desktop + # client sends paths the gateway is not the judge of): nothing made. + _rpc(ws, 1, "session.create", {"cwd": str(target)}) + plain = _drain_for_response(ws, 1, [])["result"] + assert not target.exists() + _rpc(ws, 2, "session.create", {"cwd": str(target), "create_dir": True}) + created = _drain_for_response(ws, 2, [])["result"] + _rpc(ws, 3, "session.create", {"cwd": "relative/path", "create_dir": True}) + relative = _drain_for_response(ws, 3, []) + + assert plain["session_id"] == "fake-1" + assert target.is_dir() + assert created["session_id"] == "fake-2" + assert state.manager.cwds == [str(target), str(target)] + assert "must be absolute" in relative["error"]["message"] + + +def test_session_create_can_isolate_the_session_in_a_worktree(tmp_path: Path) -> None: + import subprocess + + import os + + repo = tmp_path / "repo" + repo.mkdir() + # The real environment plus an identity: git on Windows needs SYSTEMROOT, + # TEMP and friends to start at all. + env = {**os.environ, "GIT_AUTHOR_NAME": "t", "GIT_AUTHOR_EMAIL": "t@t", + "GIT_COMMITTER_NAME": "t", "GIT_COMMITTER_EMAIL": "t@t", + "HOME": str(tmp_path), "USERPROFILE": str(tmp_path)} + for args in (["init", "-q", "-b", "main"], ["commit", "-q", "--allow-empty", "-m", "root"]): + subprocess.run(["git", *args], cwd=repo, check=True, env=env) + state, _agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + _rpc(ws, 1, "session.create", {"cwd": str(repo), "worktree": True}) + created = _drain_for_response(ws, 1, [])["result"] + _rpc(ws, 2, "session.create", {"cwd": str(tmp_path), "worktree": True}) + refused = _drain_for_response(ws, 2, []) + + worktree = created["worktree"] + assert Path(worktree["path"]).is_dir() + assert Path(worktree["path"]).resolve().parent == (repo / ".clawcodex" / "worktrees").resolve() + assert Path(worktree["repo_root"]).resolve() == repo.resolve() + assert state.manager.cwds == [worktree["path"]] + assert "git repository" in refused["error"]["message"] + + +def test_session_close_if_idle_refuses_a_busy_runtime(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.create", {}) + sid = _drain_for_response(ws, 1, events)["result"]["session_id"] + _agents[0].hold_turns = True + _rpc(ws, 2, "prompt.submit", {"session_id": sid, "text": "hi"}) + _drain_for_response(ws, 2, events) + # Mid-turn: the conditional close is refused, the runtime stays. + _rpc(ws, 3, "session.close", {"session_id": sid, "if_idle": True}) + refused = _drain_for_response(ws, 3, events)["result"] + assert refused == {"ok": True, "closed": False, "reason": "busy"} + assert sid in state.sessions + # The turn ends: now it is idle and goes. + _agents[0].queue.put_nowait({"type": "result", "subtype": "success", + "result": "done", "permission_mode": "default"}) + _drain_for_event(ws, "message.complete", events) + _rpc(ws, 4, "session.close", {"session_id": sid, "if_idle": True}) + closed = _drain_for_response(ws, 4, events)["result"] + assert closed == {"ok": True, "closed": True} + assert sid not in state.sessions + + +def test_session_close_if_idle_trusts_the_agent_about_its_own_work(tmp_path: Path) -> None: + """A /goal continuation or a /loop is invisible to the gateway: the agent + says busy through ``get_activity``, and the conditional close is refused.""" + state, agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.create", {}) + sid = _drain_for_response(ws, 1, events)["result"]["session_id"] + agents[0].busy = True + _rpc(ws, 2, "session.close", {"session_id": sid, "if_idle": True}) + assert _drain_for_response(ws, 2, events)["result"] == {"ok": True, "closed": False, "reason": "busy"} + agents[0].busy = False + _rpc(ws, 3, "session.close", {"session_id": sid, "if_idle": True}) + assert _drain_for_response(ws, 3, events)["result"] == {"ok": True, "closed": True} + + +def test_session_close_if_idle_refuses_while_an_approval_is_pending(tmp_path: Path) -> None: + state, agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.create", {}) + sid = _drain_for_response(ws, 1, events)["result"]["session_id"] + agents[0].queue.put_nowait({ + "type": "control_request", "request_id": "ask-1", + "request": {"subtype": "can_use_tool", "tool_name": "Bash", "input": {"command": "ls"}}, + }) + _drain_for_event(ws, "approval.request", events) + _rpc(ws, 2, "session.close", {"session_id": sid, "if_idle": True}) + assert _drain_for_response(ws, 2, events)["result"]["closed"] is False + _rpc(ws, 3, "approval.respond", {"session_id": sid, "choice": "once"}) + _drain_for_response(ws, 3, events) + _rpc(ws, 4, "session.close", {"session_id": sid, "if_idle": True}) + assert _drain_for_response(ws, 4, events)["result"]["closed"] is True + + +def test_session_close_answers_before_a_slow_teardown_and_tells_every_window(tmp_path: Path) -> None: + """The reply — and the next call on the socket, the open of the session + being moved to — must not wait on SessionEnd hooks and the worker join.""" + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "next-row", [{"role": "user", "content": "hi"}]) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.create", {}) + sid = _drain_for_response(ws, 1, events)["result"]["session_id"] + agents[0].shutdown_delay_s = 5.0 + _rpc(ws, 2, "session.close", {"session_id": sid}) + _rpc(ws, 3, "session.history", {"session_id": "next-row"}) + closed = _drain_for_response(ws, 2, events)["result"] + history = _drain_for_response(ws, 3, events)["result"] + # Both answered while the teardown is still running: ordering, not + # a wall-clock bound, is the claim. + assert agents[0].shutdown_called is False + + assert closed == {"ok": True, "closed": True} + assert history["message_count"] == 1 + assert any(e.get("type") == "session.closed" and e.get("session_id") == sid for e in events) + assert sid not in state.sessions + + +def test_a_second_window_resuming_the_same_row_gets_the_turn_events(tmp_path: Path) -> None: + """The reuse path must subscribe the resuming socket: before, a second + window handed the first window's runtime saw none of its own turn.""" + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "shared-row", [{"role": "user", "content": "hi"}]) + + with TestClient(build_app(state)) as client, _connect(client) as first, _connect(client) as second: + first.receive_json() + second.receive_json() + events_a: list[dict] = [] + events_b: list[dict] = [] + _rpc(first, 1, "session.resume", {"session_id": "shared-row", "omit_messages": True}) + runtime = _drain_for_response(first, 1, events_a)["result"]["session_id"] + _rpc(second, 1, "session.resume", {"session_id": "shared-row", "omit_messages": True}) + assert _drain_for_response(second, 1, events_b)["result"]["session_id"] == runtime + assert len(agents) == 1 + _rpc(second, 2, "prompt.submit", {"session_id": runtime, "text": "from b"}) + _drain_for_response(second, 2, events_b) + complete = _drain_for_event(second, "message.complete", events_b) + assert complete["session_id"] == runtime + + +def test_resume_reports_a_runtime_mid_turn_as_running(tmp_path: Path) -> None: + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "busy-row", [{"role": "user", "content": "hi"}]) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.resume", {"session_id": "busy-row", "omit_messages": True}) + runtime = _drain_for_response(ws, 1, events)["result"]["session_id"] + agents[0].hold_turns = True + _rpc(ws, 2, "prompt.submit", {"session_id": runtime, "text": "go"}) + _drain_for_response(ws, 2, events) + _rpc(ws, 3, "session.resume", {"session_id": "busy-row", "omit_messages": True}) + again = _drain_for_response(ws, 3, events)["result"] + _rpc(ws, 4, "session.history", {"session_id": "busy-row"}) + history = _drain_for_response(ws, 4, events)["result"] + + assert again["session_id"] == runtime + assert again["info"]["running"] is True + assert history["info"]["running"] is True + + +def test_session_history_of_a_live_runtime_that_never_saved_is_empty_not_an_error(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + state.sessions_dir.mkdir() + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.create", {}) + sid = _drain_for_response(ws, 1, events)["result"]["session_id"] + _rpc(ws, 2, "session.history", {"session_id": sid}) + history = _drain_for_response(ws, 2, events)["result"] + + assert history["found"] is False + assert history["messages"] == [] and history["live_session_id"] == sid + assert history["info"]["running"] is False + + +def test_a_refused_replay_is_not_reused(tmp_path: Path) -> None: + """A runtime whose ``resume`` control the agent refused holds no + conversation; the next click must spawn again rather than adopt it.""" + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "row", [{"role": "user", "content": "hi"}]) + + class Refusing(FakeAgent): + async def send_to_agent(self, frame: dict) -> None: + request = frame.get("request") or {} + if frame.get("type") == "control_request" and request.get("subtype") == "resume": + self.inbound.append(frame) + await self.queue.put({ + "type": "control_response", + "response": {"request_id": frame["request_id"], "response": {"ok": False, "error": "nope"}}, + }) + return + await super().send_to_agent(frame) + + async def spawn(session_id, cwd, resume): + agent = Refusing() + agents.append(agent) + return agent + + state.spawn_agent = spawn + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.resume", {"session_id": "row", "omit_messages": True}) + reply = _drain_for_response(ws, 1, events)["result"] + first = reply["session_id"] + _rpc(ws, 2, "session.resume", {"session_id": "row", "omit_messages": True}) + second = _drain_for_response(ws, 2, events)["result"]["session_id"] + + assert first != second and len(agents) == 2 + # …and the reply does not claim the row it failed to replay, so a client + # never keeps the empty runtime as that conversation. + assert reply["stored_session_id"] == first + + +def test_prepare_workspace_expands_home_and_refuses_a_file(tmp_path: Path, monkeypatch) -> None: + from src.server.desktop_gateway_methods import _prepare_workspace + + # HOME for POSIX, USERPROFILE for Windows (ntpath ignores HOME). + monkeypatch.setenv("HOME", str(tmp_path)) + monkeypatch.setenv("USERPROFILE", str(tmp_path)) + (tmp_path / "proj").mkdir() + (tmp_path / "notes.txt").write_text("x", encoding="utf-8") + + assert _prepare_workspace("~/proj", False) == str(tmp_path / "proj") + assert _prepare_workspace("~/fresh/deep", True) == str(tmp_path / "fresh" / "deep") + assert (tmp_path / "fresh" / "deep").is_dir() + with pytest.raises(ValueError, match="not a directory"): + _prepare_workspace(str(tmp_path / "notes.txt"), True) + + +def test_worktree_on_a_new_folder_is_refused_before_anything_is_created(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + target = tmp_path / "brand-new" + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + _rpc(ws, 1, "session.create", {"cwd": str(target), "create_dir": True, "worktree": True}) + refused = _drain_for_response(ws, 1, []) + + assert "existing git repository" in refused["error"]["message"] + assert not target.exists() + + +def test_session_close_if_idle_rechecks_after_asking_the_agent(tmp_path: Path) -> None: + """A prompt that lands while the agent is answering ``get_activity`` must + keep the runtime: the idle verdict is re-taken after the await.""" + state, agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as first, _connect(client) as second: + first.receive_json() + second.receive_json() + events_a: list[dict] = [] + events_b: list[dict] = [] + _rpc(first, 1, "session.create", {}) + sid = _drain_for_response(first, 1, events_a)["result"]["session_id"] + agent = agents[0] + agent.hold_turns = True + agent.activity_delay_s = 0.5 + # The close starts asking the agent; meanwhile a prompt arrives on + # another socket. + _rpc(first, 2, "session.close", {"session_id": sid, "if_idle": True}) + _rpc(second, 1, "prompt.submit", {"session_id": sid, "text": "late"}) + _drain_for_response(second, 1, events_b) + closed = _drain_for_response(first, 2, events_a)["result"] + + assert closed == {"ok": True, "closed": False, "reason": "busy"} + assert sid in state.sessions + + +def test_a_dead_runtime_is_not_busy(tmp_path: Path) -> None: + """The pump dying mid-turn clears the busy markers, so the runtime can be + released with ``if_idle`` instead of refusing forever.""" + state, agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + _rpc(ws, 1, "session.create", {}) + sid = _drain_for_response(ws, 1, events)["result"]["session_id"] + agents[0].hold_turns = True + _rpc(ws, 2, "prompt.submit", {"session_id": sid, "text": "hi"}) + _drain_for_response(ws, 2, events) + agents[0].queue.put_nowait(RuntimeError("stream died")) + complete = _drain_for_event(ws, "message.complete", events) + assert complete["payload"]["status"] == "error" + # Retired on the spot: the turn's end first, then unregistered and + # every window told, so the next click spawns afresh instead of + # being handed the corpse. + closed = _drain_for_event(ws, "session.closed", events) + assert closed["session_id"] == sid + assert sid not in state.sessions + _rpc(ws, 3, "prompt.submit", {"session_id": sid, "text": "anyone?"}) + assert "unknown session" in _drain_for_response(ws, 3, events)["error"]["message"] + _rpc(ws, 4, "session.close", {"session_id": sid, "if_idle": True}) + assert _drain_for_response(ws, 4, events)["result"] == {"ok": True, "closed": False, "reason": "gone"} + + +def test_get_activity_is_answered_by_a_runtime_that_refused_to_start() -> None: + from tests.server.test_cron_control import _control, _last_reply, _make_session + + sess, emitted, _clock = _make_session() + sess.init_error = "sandbox unavailable" + _control(sess, "get_activity") + reply = _last_reply(emitted) + + assert reply["ok"] is True and reply["busy"] is False + # …while a control that does work is still refused. + _control(sess, "rename", name="x") + assert _last_reply(emitted)["ok"] is False + + +def test_a_conditional_close_keeps_a_runtime_another_window_holds(tmp_path: Path) -> None: + """A window browsing rows must not close the runtime a desktop tile or + another window opened and sits on; only the last holder's release closes.""" + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "shared", [{"role": "user", "content": "hi"}]) + + with TestClient(build_app(state)) as client, _connect(client) as first, _connect(client) as second: + first.receive_json() + second.receive_json() + _rpc(first, 1, "session.resume", {"session_id": "shared", "omit_messages": True}) + runtime = _drain_for_response(first, 1, [])["result"]["session_id"] + _rpc(second, 1, "session.resume", {"session_id": "shared", "omit_messages": True}) + assert _drain_for_response(second, 1, [])["result"]["session_id"] == runtime + _rpc(second, 2, "session.close", {"session_id": runtime, "if_idle": True}) + assert _drain_for_response(second, 2, [])["result"] == {"ok": True, "closed": False, "reason": "held"} + assert runtime in state.sessions + _rpc(first, 2, "session.close", {"session_id": runtime, "if_idle": True}) + assert _drain_for_response(first, 2, [])["result"] == {"ok": True, "closed": True} + assert len(agents) == 1 + + +def test_a_disconnected_window_no_longer_holds_its_runtime(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "shared", [{"role": "user", "content": "hi"}]) + + with TestClient(build_app(state)) as client, _connect(client) as first: + first.receive_json() + _rpc(first, 1, "session.resume", {"session_id": "shared", "omit_messages": True}) + runtime = _drain_for_response(first, 1, [])["result"]["session_id"] + with _connect(client) as second: + second.receive_json() + _rpc(second, 1, "session.resume", {"session_id": "shared", "omit_messages": True}) + _drain_for_response(second, 1, []) + # The second window went away: its hold went with it. + for _ in range(50): + if len(state.sessions[runtime].holders) == 1: + break + time.sleep(0.02) + _rpc(first, 2, "session.close", {"session_id": runtime, "if_idle": True}) + assert _drain_for_response(first, 2, [])["result"] == {"ok": True, "closed": True} + + +def test_a_dying_runtime_answers_its_waiting_queries_instead_of_dropping_sockets(tmp_path: Path) -> None: + """A control query awaiting a runtime that dies gets "no reply", not a + CancelledError that unwinds the handler and closes the window's socket.""" + state, agents = _fake_state(tmp_path) + + with TestClient(build_app(state)) as client, _connect(client) as first, _connect(client) as second: + first.receive_json() + second.receive_json() + events_a: list[dict] = [] + events_b: list[dict] = [] + _rpc(first, 1, "session.create", {}) + sid = _drain_for_response(first, 1, events_a)["result"]["session_id"] + # The fake never answers get_context_usage: the second window waits. + _rpc(second, 1, "session.usage", {"session_id": sid}) + time.sleep(0.1) + agents[0].queue.put_nowait(RuntimeError("stream died")) + usage = _drain_for_response(second, 1, events_b) + assert usage["result"] == {} + # The socket is still served. + _rpc(second, 2, "setup.status", {}) + assert _drain_for_response(second, 2, events_b)["result"] == {"provider_configured": True} + _rpc(first, 2, "setup.status", {}) + assert _drain_for_response(first, 2, events_a)["result"] == {"provider_configured": True} + assert sid not in state.sessions + + +def test_the_opener_survives_a_runtime_dying_mid_attach(tmp_path: Path) -> None: + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "row", [{"role": "user", "content": "hi"}]) + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + events: list[dict] = [] + # With the questions capability the attach waits on + # set_ask_user_interactive, which the fake never answers. + _rpc(ws, 1, "session.resume", { + "session_id": "row", "omit_messages": True, + "capabilities": {"ask_user_question": True}, + }) + for _ in range(100): + if agents and any( + (f.get("request") or {}).get("subtype") == "set_ask_user_interactive" for f in agents[0].inbound + ): + break + time.sleep(0.02) + agents[0].queue.put_nowait(RuntimeError("stream died")) + reply = _drain_for_response(ws, 1, events) + assert "ended while starting" in reply["error"]["message"] + _rpc(ws, 2, "setup.status", {}) + assert _drain_for_response(ws, 2, events)["result"] == {"provider_configured": True} + assert state.sessions == {} + + +def test_a_stream_that_dies_before_init_fails_the_attach_at_once(tmp_path: Path) -> None: + state, _agents = _fake_state(tmp_path) + base_spawn = state.spawn_agent + + async def spawn(session_id, cwd, resume): + agent = await base_spawn(session_id, cwd, resume) + agent.die_before_init = True + return agent + + state.spawn_agent = spawn + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + started = time.monotonic() + _rpc(ws, 1, "session.create", {}) + reply = _drain_for_response(ws, 1, []) + elapsed = time.monotonic() - started + + assert "ended while starting" in reply["error"]["message"] + # What this discriminates is the control timeout the attach used to wait. + assert elapsed < CONTROL_TIMEOUT_S / 2, f"waited {elapsed:.0f}s on a stream that never sent init" + assert state.sessions == {} + + +def test_a_runtime_dying_during_the_resume_tail_is_not_handed_out(tmp_path: Path) -> None: + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "row", [{"role": "user", "content": "hi"}]) + base_spawn = state.spawn_agent + + async def spawn(session_id, cwd, resume): + agent = await base_spawn(session_id, cwd, resume) + agent.settings_delay_s = 0.5 + return agent + + state.spawn_agent = spawn + + with TestClient(build_app(state)) as client, _connect(client) as ws: + ws.receive_json() + _rpc(ws, 1, "session.resume", {"session_id": "row", "omit_messages": True}) + for _ in range(100): + if agents and any( + (f.get("request") or {}).get("subtype") == "get_settings" for f in agents[0].inbound + ): + break + time.sleep(0.02) + agents[0].queue.put_nowait(RuntimeError("stream died")) + reply = _drain_for_response(ws, 1, []) + + assert "ended while starting" in reply["error"]["message"] + assert state.sessions == {} + + +def test_a_resume_whose_runtime_was_closed_outright_meanwhile_is_refused(tmp_path: Path) -> None: + """A desktop tile's unconditional close during the resume's tail must not + leave the resuming window with an id the registry no longer has.""" + state, agents = _fake_state(tmp_path) + state.sessions_dir = tmp_path / "saved" + _write_saved(state.sessions_dir, "row", [{"role": "user", "content": "hi"}]) + base_spawn = state.spawn_agent + + async def spawn(session_id, cwd, resume): + agent = await base_spawn(session_id, cwd, resume) + agent.settings_delay_s = 0.5 + return agent + + state.spawn_agent = spawn + + with TestClient(build_app(state)) as client, _connect(client) as first, _connect(client) as second: + first.receive_json() + second.receive_json() + _rpc(first, 1, "session.resume", {"session_id": "row", "omit_messages": True}) + for _ in range(100): + if agents and any( + (f.get("request") or {}).get("subtype") == "get_settings" for f in agents[0].inbound + ): + break + time.sleep(0.02) + runtime = next(iter(state.sessions)) + _rpc(second, 1, "session.close", {"session_id": runtime}) + assert _drain_for_response(second, 1, [])["result"]["closed"] is True + reply = _drain_for_response(first, 1, []) + + assert "ended while starting" in reply["error"]["message"] + assert state.sessions == {} diff --git a/tests/server/test_desktop_projects.py b/tests/server/test_desktop_projects.py index 142c0d64..4182aac4 100644 --- a/tests/server/test_desktop_projects.py +++ b/tests/server/test_desktop_projects.py @@ -183,3 +183,39 @@ def test_unresolved_cwd_falls_into_home_bucket(tmp_path): ) assert [project["id"] for project in tree["projects"]] == [NO_PROJECT_ID] + + +# ─── ProbeCache ────────────────────────────────────────────────────────────── + + +def test_probe_cache_answers_within_its_ttl_and_forgets_after(): + from src.server.desktop_projects import ProbeCache + + now = [100.0] + cache = ProbeCache(ttl_s=10.0, worktree_ttl_s=2.0, clock=lambda: now[0]) + + assert cache.has_repo_root("/a") is False + cache.set_repo_root("/a", "/repo") + cache.set_repo_root("/b", None) # "not a repo" is an answer too + cache.set_worktrees("/repo", ["/repo", "/repo/.wt/x"]) + + assert cache.has_repo_root("/a") and cache.repo_root("/a") == "/repo" + assert cache.has_repo_root("/b") and cache.repo_root("/b") is None + assert cache.worktrees("/repo") == ["/repo", "/repo/.wt/x"] + + now[0] += 3.0 + assert cache.worktrees("/repo") is None # worktree lists expire sooner + assert cache.has_repo_root("/a") + now[0] += 8.0 + assert cache.has_repo_root("/a") is False + + +def test_probe_cache_forget_worktrees_drops_one_repo(): + from src.server.desktop_projects import ProbeCache + + cache = ProbeCache() + cache.set_worktrees("/r1", ["/r1"]) + cache.set_worktrees("/r2", ["/r2"]) + cache.forget_worktrees("/r1") + assert cache.worktrees("/r1") is None + assert cache.worktrees("/r2") == ["/r2"] diff --git a/tests/server/test_desktop_sessions.py b/tests/server/test_desktop_sessions.py index 416d2df6..92f75829 100644 --- a/tests/server/test_desktop_sessions.py +++ b/tests/server/test_desktop_sessions.py @@ -522,3 +522,61 @@ def test_conversation_keeps_a_step_usage_and_model_on_disk() -> None: assert stored[1]["usage"] == {"input_tokens": 3, "output_tokens": 2} assert stored[1]["model"] == "deepseek-v4-flash" assert "usage" not in stored[0] + + +# ─── the sidebar row cache ─────────────────────────────────────────────────── + + +def test_list_session_rows_rereads_only_changed_files(tmp_path: Path, monkeypatch) -> None: + from src.server import desktop_sessions + + desktop_sessions.clear_session_row_cache() + d = tmp_path / "sessions" + d.mkdir() + _write_session(d, "a", preview="A", count=1, age_s=20) + _write_session(d, "b", preview="B", count=1, age_s=10) + reads: list[str] = [] + real = desktop_sessions._read_session_file + + def counting(path): + reads.append(path.stem) + return real(path) + + monkeypatch.setattr(desktop_sessions, "_read_session_file", counting) + + first = desktop_sessions.list_session_rows(d, limit=0) + assert [r["id"] for r in first["sessions"]] == ["b", "a"] + assert sorted(reads) == ["a", "b"] + + # Nothing changed: no file is parsed again, rows are equal but not shared. + second = desktop_sessions.list_session_rows(d, limit=0) + assert sorted(reads) == ["a", "b"] + assert second["sessions"] == first["sessions"] + second["sessions"][0]["is_active"] = True + assert desktop_sessions.list_session_rows(d, limit=0)["sessions"][0]["is_active"] is False + + # A rewritten file is parsed again and its new content served. + _write_session(d, "a", preview="A2", count=3) + third = desktop_sessions.list_session_rows(d, limit=0) + assert reads.count("a") == 2 + assert [r["id"] for r in third["sessions"]] == ["a", "b"] + assert third["sessions"][0]["preview"] == "A2" + + # A deleted file leaves the cache too. + (d / "b.json").unlink() + assert [r["id"] for r in desktop_sessions.list_session_rows(d, limit=0)["sessions"]] == ["a"] + assert "b.json" not in " ".join(desktop_sessions._ROW_CACHE) + desktop_sessions.clear_session_row_cache() + + +def test_load_session_messages_carries_the_stored_session_facts(tmp_path: Path) -> None: + from src.server.desktop_sessions import load_session_messages + + d = tmp_path / "sessions" + d.mkdir() + _write_session(d, "s", preview="P", count=1, messages=[{"role": "user", "content": "hi"}]) + + stored = load_session_messages(d, "s") + + assert stored is not None + assert (stored["cwd"], stored["model"], stored["provider"]) == ("/tmp/w", "m1", "p1") diff --git a/tests/server/test_gateway_questions.py b/tests/server/test_gateway_questions.py index 990c1b58..ba23dca5 100644 --- a/tests/server/test_gateway_questions.py +++ b/tests/server/test_gateway_questions.py @@ -1215,6 +1215,12 @@ class _Session: session_id = "rt1" init_info = {"model": "launch-default", "provider": "deepseek", "cwd": "/w"} titled = True + # What the reply's ``running`` reads: a stub runtime between turns. + turn_active = False + # The row this runtime replays (the reply echoes it). + stored_id = "stored-1" + # Its stream is up. + dead = False def __init__(self) -> None: self.refreshed = False @@ -1235,6 +1241,9 @@ async def _create(cwd: Any, resume: Any, params: Any) -> Any: connection._create = _create # type: ignore[method-assign] class _State: + # The registry the reply's liveness check consults. + sessions = {"rt1": session} + def saved_sessions_dir(self): # pragma: no cover - not reached raise AssertionError diff --git a/tests/test_workspace_snapshot.py b/tests/test_workspace_snapshot.py new file mode 100644 index 00000000..f00da5d8 --- /dev/null +++ b/tests/test_workspace_snapshot.py @@ -0,0 +1,118 @@ +"""The workspace snapshot's file walk: pruned before entering, bounded, honest. + +The system prompt's ``## Runtime Context`` counts used to come from two +``Path.rglob`` passes that crawled every directory under the workspace — +``node_modules`` included — and filtered afterwards, ~20 s per system-prompt +build on a repo with a few front-end packages (built at spawn and again on +every resume, so ~45 s to open a saved session from the web sidebar). +""" + +from __future__ import annotations + +from pathlib import Path + +import pytest + +from src.context_system.builder import _build_workspace_section +from src.context_system.workspace_snapshot import ( + MAX_SCANNED_DIRS, + build_workspace_snapshot, + count_python_files, +) + + +def _touch(path: Path) -> None: + path.parent.mkdir(parents=True, exist_ok=True) + path.write_text("", encoding="utf-8") + + +def test_counts_python_and_test_files_outside_ignored_trees(tmp_path: Path) -> None: + _touch(tmp_path / "src" / "app.py") + _touch(tmp_path / "src" / "pkg" / "util.py") + _touch(tmp_path / "tests" / "test_app.py") + _touch(tmp_path / "README.md") + # Vendored and generated trees must neither count nor be entered. + _touch(tmp_path / "node_modules" / "left-pad" / "setup.py") + _touch(tmp_path / ".venv" / "lib" / "site.py") + _touch(tmp_path / ".git" / "hooks" / "test_hook.py") + _touch(tmp_path / ".clawcodex" / "worktrees" / "wt-1" / "src" / "app.py") + + snapshot = build_workspace_snapshot(tmp_path) + + assert snapshot.python_file_count == 3 + assert snapshot.test_file_count == 1 + assert snapshot.counts_partial is False + assert snapshot.key_files == ("README.md",) + assert "node_modules/" not in snapshot.top_level_entries + assert "src/" in snapshot.top_level_entries + + +def test_ignored_directories_are_pruned_not_filtered(tmp_path: Path) -> None: + """A huge ignored subtree costs nothing: the walk never steps into it. + + With a budget of three directories (root, ``src``, ``tests``) the walk + only stays within budget if ``node_modules`` — deeper than the budget on + its own — was skipped before being entered rather than crawled and then + filtered out. + """ + _touch(tmp_path / "src" / "app.py") + _touch(tmp_path / "tests" / "test_app.py") + deep = tmp_path / "node_modules" + for index in range(20): + deep = deep / f"dep-{index}" + _touch(deep / "vendored.py") + + python_files, test_files, partial = count_python_files(tmp_path, max_dirs=3) + + assert (python_files, test_files, partial) == (2, 1, False) + + +def test_the_walk_stops_at_its_budget_and_says_so(tmp_path: Path) -> None: + for index in range(10): + _touch(tmp_path / f"pkg-{index:02d}" / "mod.py") + + python_files, _tests, partial = count_python_files(tmp_path, max_dirs=4) + + # Root plus the first three packages in name order — deterministic. + assert partial is True + assert python_files == 3 + + snapshot = build_workspace_snapshot(tmp_path, max_dirs=4) + assert snapshot.counts_partial is True + + section = _build_workspace_section(tmp_path, tmp_path) + assert "- Python files: 10" in section # the default budget covers it all + assert "partial" not in section + + +def test_partial_counts_are_rendered_as_lower_bounds(tmp_path: Path, monkeypatch: pytest.MonkeyPatch) -> None: + for index in range(6): + _touch(tmp_path / f"pkg-{index}" / "test_mod.py") + + from src.context_system import workspace_snapshot + + monkeypatch.setattr(workspace_snapshot, "MAX_SCANNED_DIRS", 2) + # The builder calls through the module-level default, so patch the + # function's default the way a smaller budget would arrive in production. + monkeypatch.setattr( + workspace_snapshot, + "build_workspace_snapshot", + lambda root, cwd=None: build_workspace_snapshot(root, cwd=cwd, max_dirs=2), + ) + + section = _build_workspace_section(tmp_path, tmp_path) + + assert "- Python files: 1+ (partial scan)" in section + assert "- Test files: 1+ (partial scan)" in section + + +def test_symlink_loops_do_not_hang_the_walk(tmp_path: Path) -> None: + _touch(tmp_path / "src" / "app.py") + try: + (tmp_path / "src" / "loop").symlink_to(tmp_path, target_is_directory=True) + except (OSError, NotImplementedError): # pragma: no cover - platform + pytest.skip("symlinks unavailable") + + python_files, _tests, partial = count_python_files(tmp_path, max_dirs=MAX_SCANNED_DIRS) + + assert (python_files, partial) == (1, False) diff --git a/ui-web/README.md b/ui-web/README.md index f283dc4a..775e8e07 100644 --- a/ui-web/README.md +++ b/ui-web/README.md @@ -237,6 +237,64 @@ Harness, selecting a session opens its folder unless it was manually collapsed; manual toggles last for the page's lifetime. Filtering temporarily expands matching folders and restores their previous state when cleared. +## Opening a saved session + +Clicking a row is two round-trips, the way the reference opens a session. +`session.history` reads the stored transcript cold — a file read, tens of +milliseconds even for a multi-megabyte conversation — and the client renders +it at once: title, workspace, model chip, nodes, trajectory timings. Then +`session.resume` attaches the runtime that will answer the next prompt +(provider, tool registry, system prompt, the stored conversation loaded into +it) behind the transcript; the composer says *Connecting the agent…* +meanwhile, and a prompt sent during that window waits for the attach rather +than starting a session of its own. Before this the whole click sat on the +attach, and the attach sat on a system-prompt walk of the workspace that +took ~20 s per build on a repo with a few `node_modules` trees (built at +spawn and again on resume — ~45 s to open a session). + +A row the backend has already replayed comes back with the same runtime +(subscribed to this socket too, so a second window sees the turns): the +backend keys live sessions by runtime id, matches a resume on the stored id +the runtime replays, and a second click while the first is still attaching +waits on it instead of spawning twice. A runtime handed back mid-turn is +adopted as running. `session.history` and `projects.tree` are pure reads +the gateway serves beside whatever else the socket is doing, so a click's +transcript never queues behind the previous click's attach or the teardown +of the session being left. + +Leaving a session that was only *looked at* — nothing sent to it, no +approval or question answered — releases its runtime: `session.close` with +`if_idle`, which the backend refuses while another window or a desktop tile +still holds the runtime (each socket that opened it is a holder until it +lets go or disconnects), while a turn runs or an ask is pending, and again +on the agent's own word (`get_activity`: a `/goal` continuation, a `/loop` +or cron job waiting to fire, a queued prompt, a background shell). So +browsing through saved sessions does not leave a trail of idle agents, +while a session that was used stays up; "used" survives a reload with the +remembered session. A runtime that goes away — closed by its last holder, or +its agent stream ended — tells every window (`session.closed`) and leaves +the registry, so the next prompt or click reconnects the conversation (a +prompt to a runtime the backend no longer has is sent again after the +reconnect); the row replays from the record the runtime saved under its own +id. A backend without `session.history` gets the one-call resume, +transcript included, as before. + +## New session + +**New session** (the sidebar button, the brand mark, `⌘⇧N`) opens a dialog +rather than starting a session on the spot. It offers every workspace the +sidebar knows plus **Create new workspace…**, which takes an absolute folder +path and creates the folder if it is not there yet (`session.create` with +`create_dir`). The **Worktree** switch runs the session in a fresh git +worktree of that repo — the CLI's `--worktree`, under +`.clawcodex/worktrees/` — so parallel sessions in one repo cannot step +on each other's files; the worktree is left in place when the session ends, +since a browser tab has no exit dialog to offer keep-or-remove. The folder +is created as the backend's own user, anywhere it can write (`~` expands). +A refusal (a relative path, a folder that is not a git repository, Worktree +on a folder that does not exist yet) stays in the dialog for correcting, +before anything is created. + ## The session across a reload A reload lands back on the session the window was on. The client remembers @@ -248,9 +306,10 @@ the backend still has it, the reply is the very same session, its running turn included; once it is gone, the runtime's own record — the complete one — is replayed into a new runtime. A runtime that never saved, because nothing was typed after resuming a row, has no record, so the row it came from is -replayed instead and the blank runtime the first attempt spawned is closed. -A session the backend no longer knows is forgotten without a notice: the -hero is the honest place to land. +replayed instead (which releases the blank runtime the first attempt landed +on, as any navigation away from an idle runtime does). A session the backend +no longer knows is forgotten without a notice: the hero is the honest place +to land. ## Trajectory diff --git a/ui-web/src/App.tsx b/ui-web/src/App.tsx index 07dbd214..79061410 100644 --- a/ui-web/src/App.tsx +++ b/ui-web/src/App.tsx @@ -6,8 +6,9 @@ import { SettingsOverlay } from './settings/SettingsOverlay.tsx' import { SidebarRight } from './sidebar-right/SidebarRight.tsx' import { closeSidebar, resetSidebar } from './sidebar-right/store.ts' import { AppFrame } from './layout/AppFrame.tsx' +import { NewSessionDialog, openNewSessionDialog } from './sidebar/NewSessionDialog.tsx' import { Sidebar } from './sidebar/Sidebar.tsx' -import { createSession, start } from './state/actions.ts' +import { start } from './state/actions.ts' import { $detailsOpen, $detailsWidth, openDetails, toggleSidebar } from './state/layout.ts' import { $bootError, $bootPhase, $sessionId, $workspace } from './state/store.ts' import { installTheme } from './state/theme.ts' @@ -97,7 +98,7 @@ export function App() { // stealing it would surprise the user in their own browser. if (event.key === 'n' && event.shiftKey) { event.preventDefault() - void createSession({ cwd: $workspace.get() }) + openNewSessionDialog() } } @@ -121,9 +122,10 @@ export function App() { details={detailsOpen ? : null} sidebar={state => } /> - {/* Outside the frame: it covers the whole app, including the sidebar - it is opened from. */} + {/* Outside the frame: they cover the whole app, including the sidebar + they are opened from. */} + ) } diff --git a/ui-web/src/conversation/InputBar.tsx b/ui-web/src/conversation/InputBar.tsx index db3bf968..28277e23 100644 --- a/ui-web/src/conversation/InputBar.tsx +++ b/ui-web/src/conversation/InputBar.tsx @@ -17,7 +17,7 @@ import type { ModelOptionsResult, } from '../gateway/protocol.ts' import { attachImage, searchFiles } from '../state/actions.ts' -import { $commands, $notice } from '../state/store.ts' +import { $commands, $notice, $sessionAttaching } from '../state/store.ts' import { ArrowUpIcon, PlusIcon, SlashSquareIcon, StopIcon, XIcon } from '../ui/icons.tsx' import { ContextMeter } from './ContextMeter.tsx' import { @@ -112,6 +112,7 @@ export function InputBar({ }: InputBarProps) { const commands = useStore($commands) const notice = useStore($notice) + const attaching = useStore($sessionAttaching) const textarea = useRef(null) const card = useRef(null) const [highlight, setHighlight] = useState(0) @@ -484,7 +485,7 @@ export function InputBar({ return (
- {notice.text !== '' && ( + {notice.text !== '' ? (
{notice.text}
+ ) : ( + // The transcript is up before its runtime is: say so, since a prompt + // sent now waits for the agent rather than going out at once. + attaching && ( +
+ Connecting the agent… +
+ ) )}
{mention !== null && files.length > 0 && ( diff --git a/ui-web/src/gateway/protocol.ts b/ui-web/src/gateway/protocol.ts index 1c6d4689..b061e4de 100644 --- a/ui-web/src/gateway/protocol.ts +++ b/ui-web/src/gateway/protocol.ts @@ -243,6 +243,36 @@ export interface SessionCreateResult { info?: SessionInfoPayload session_id: string stored_session_id?: string + /** Present when the session was created with `worktree: true`. */ + worktree?: SessionWorktree +} + +/** The git worktree a session was isolated in (`session.create` → `worktree`). */ +export interface SessionWorktree { + branch?: string + name?: string + path: string + repo_root?: string +} + +/** + * `session.history` — a saved session's transcript, read cold: no runtime is + * spawned or asked. What the sidebar click renders at once; the runtime + * attaches afterwards through `session.resume`. + * + * `found` is false for a live runtime that never saved (nothing typed into it + * yet): no messages, and not an error. `live_session_id` names the runtime + * already replaying this row when the backend has one. + */ +export interface SessionHistoryResult { + found?: boolean + info?: SessionInfoPayload + live_session_id?: string + message_count?: number + messages?: StoredMessage[] + session_id?: string + stored_session_id?: string + title?: string } export interface StoredMessage { diff --git a/ui-web/src/sidebar/NewSessionDialog.module.css b/ui-web/src/sidebar/NewSessionDialog.module.css new file mode 100644 index 00000000..60de25ea --- /dev/null +++ b/ui-web/src/sidebar/NewSessionDialog.module.css @@ -0,0 +1,171 @@ +/* The New session dialog. Same chrome as the Full-access confirmation: a + centred card on the mask, above the composer's menus and the settings + takeover. */ + +.scrim { + position: fixed; + inset: 0; + z-index: 300; + display: grid; + place-items: center; + padding: 24px; + background: var(--cc-alias-bg-mask-1); +} + +.dialog { + box-sizing: border-box; + display: flex; + flex-direction: column; + gap: 14px; + width: min(440px, 100%); + margin: 0; + padding: 18px; + border: 1px solid var(--cc-alias-border-inverted); + border-radius: 14px; + background: var(--cc-specific-menu); + box-shadow: var(--cc-shadow-lv3); +} + +.head { + display: flex; + align-items: center; + justify-content: space-between; +} + +.title { + color: var(--cc-alias-label-primary); + font-size: 15px; + font-weight: 600; + line-height: 22px; +} + +.close { + display: inline-flex; + align-items: center; + justify-content: center; + width: 24px; + height: 24px; + padding: 0; + border: none; + border-radius: 50%; + background: transparent; + color: var(--cc-alias-label-secondary); + cursor: pointer; +} + +.close:hover { + background: var(--cc-alias-interactive-bg-hover); +} + +.field { + display: flex; + flex-direction: column; + gap: 6px; +} + +.label { + color: var(--cc-alias-label-primary); + font-size: 13px; + font-weight: 500; + line-height: 18px; +} + +.hint { + color: var(--cc-alias-label-tertiary); + font-size: 12px; + line-height: 16px; +} + +.select, +.input { + box-sizing: border-box; + width: 100%; + height: 34px; + padding: 0 10px; + border: 1px solid var(--cc-alias-border-l3); + border-radius: 8px; + background: var(--cc-alias-bg-layer-1); + color: var(--cc-alias-label-primary); + font: inherit; + font-size: 13px; +} + +.select:focus-visible, +.input:focus-visible { + outline: 2px solid var(--cc-alias-button-info-fill); + outline-offset: -1px; +} + +.input::placeholder { + color: var(--cc-alias-label-caption); +} + +.switchRow { + display: flex; + align-items: center; + justify-content: space-between; + gap: 12px; + cursor: pointer; +} + +.switchText { + display: flex; + flex-direction: column; + gap: 2px; +} + +/* A switch drawn from the checkbox itself, so the state stays in the form. */ +.switch { + appearance: none; + flex: none; + position: relative; + width: 36px; + height: 20px; + margin: 0; + border-radius: 10px; + background: var(--cc-alias-border-l4); + cursor: pointer; + transition: background 120ms ease; +} + +.switch::after { + content: ''; + position: absolute; + top: 2px; + left: 2px; + width: 16px; + height: 16px; + border-radius: 50%; + background: var(--cc-static-neutral-00); + box-shadow: 0 1px 2px rgba(0, 0, 0, 0.2); + transition: transform 120ms ease; +} + +.switch:checked { + background: var(--cc-alias-button-info-fill); +} + +.switch:checked::after { + transform: translateX(16px); +} + +.switch:focus-visible { + outline: 2px solid var(--cc-alias-button-info-fill); + outline-offset: 2px; +} + +.error { + padding: 8px 10px; + border-radius: 8px; + background: var(--cc-alias-interactive-bg-hover-danger); + color: var(--cc-static-red-600); + font-size: 12px; + line-height: 18px; + overflow-wrap: anywhere; +} + +.actions { + display: flex; + justify-content: flex-end; + gap: 8px; +} diff --git a/ui-web/src/sidebar/NewSessionDialog.test.tsx b/ui-web/src/sidebar/NewSessionDialog.test.tsx new file mode 100644 index 00000000..265107c6 --- /dev/null +++ b/ui-web/src/sidebar/NewSessionDialog.test.tsx @@ -0,0 +1,148 @@ +import { act, cleanup, fireEvent, render, screen } from '@testing-library/react' +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest' + +import type { ProjectNode } from '../gateway/protocol.ts' +import { $newSessionDialog, $projects, $workspace } from '../state/store.ts' + +const createSession = vi.fn<(options?: Record) => Promise>() + +vi.mock('../state/actions.ts', () => ({ + createSession: (options?: Record) => createSession(options), +})) + +const { NEW_WORKSPACE, NewSessionDialog, knownWorkspaces, openNewSessionDialog } = await import('./NewSessionDialog.tsx') + +function project(path: string | null): ProjectNode { + return { id: path ?? 'home', label: path ?? 'Home', path, repos: [] } +} + +beforeEach(() => { + createSession.mockReset() + createSession.mockResolvedValue(null) + $workspace.set('/work/current') + $projects.set([project('/work/alpha'), project('/work/current'), project(null)]) + $newSessionDialog.set(false) +}) + +afterEach(() => { + cleanup() + $newSessionDialog.set(false) + $projects.set([]) + $workspace.set('') +}) + +describe('knownWorkspaces', () => { + it('lists the current workspace first, then the sidebar folders, once each, never Home', () => { + expect(knownWorkspaces('/work/current', [project('/work/alpha'), project('/work/current'), project(null)])) + .toEqual(['/work/current', '/work/alpha']) + expect(knownWorkspaces('', [project('/a')])).toEqual(['/a']) + }) + + it('offers a repo’s worktree lanes as well as the repo', () => { + const repo: ProjectNode = { + id: '/repo', + label: 'repo', + path: '/repo', + repos: [{ + groups: [ + { id: 'main', label: 'main', path: '/repo', sessions: [] }, + { id: 'wt', label: 'feature', path: '/repo/.clawcodex/worktrees/feature', sessions: [] }, + ], + id: '/repo', + label: 'repo', + path: '/repo', + }], + } + + expect(knownWorkspaces('', [repo])).toEqual(['/repo', '/repo/.clawcodex/worktrees/feature']) + }) +}) + +describe('NewSessionDialog', () => { + it('is closed until opened, and starts a session in the chosen workspace', async () => { + render() + + expect(screen.queryByRole('dialog')).toBeNull() + + act(() => { + openNewSessionDialog() + }) + + const select = screen.getByRole('combobox') as HTMLSelectElement + expect(select.value).toBe('/work/current') + expect(screen.queryByPlaceholderText('/absolute/path/to/project')).toBeNull() + + fireEvent.change(select, { target: { value: '/work/alpha' } }) + fireEvent.click(screen.getByRole('button', { name: 'Create' })) + await act(async () => { + await Promise.resolve() + }) + + expect(createSession).toHaveBeenCalledWith({ cwd: '/work/alpha' }) + expect($newSessionDialog.get()).toBe(false) + }) + + it('creates a new workspace from an absolute path, in a worktree when asked', async () => { + $newSessionDialog.set(true) + render() + + const select = screen.getByRole('combobox') as HTMLSelectElement + fireEvent.change(select, { target: { value: NEW_WORKSPACE } }) + + const create = screen.getByRole('button', { name: 'Create' }) as HTMLButtonElement + expect(create.disabled).toBe(true) + + fireEvent.change(screen.getByPlaceholderText('/absolute/path/to/project'), { + target: { value: '/work/fresh ' }, + }) + fireEvent.click(screen.getByRole('switch')) + expect(create.disabled).toBe(false) + + fireEvent.click(create) + await act(async () => { + await Promise.resolve() + }) + + expect(createSession).toHaveBeenCalledWith({ createDir: true, cwd: '/work/fresh', worktree: true }) + expect($newSessionDialog.get()).toBe(false) + }) + + it('keeps the dialog open with the reason when the backend refuses', async () => { + createSession.mockResolvedValue('no such directory: /nope') + $newSessionDialog.set(true) + render() + + fireEvent.click(screen.getByRole('button', { name: 'Create' })) + await act(async () => { + await Promise.resolve() + }) + + expect(screen.getByRole('alert').textContent).toBe('no such directory: /nope') + expect($newSessionDialog.get()).toBe(true) + }) + + it('closes on Cancel and on Escape', () => { + $newSessionDialog.set(true) + render() + + fireEvent.click(screen.getByRole('button', { name: 'Cancel' })) + expect($newSessionDialog.get()).toBe(false) + + act(() => { + openNewSessionDialog() + }) + fireEvent.keyDown(document, { key: 'Escape' }) + expect($newSessionDialog.get()).toBe(false) + expect(createSession).not.toHaveBeenCalled() + }) + + it('offers only the new-workspace path when no workspace is known', () => { + $workspace.set('') + $projects.set([]) + $newSessionDialog.set(true) + render() + + expect((screen.getByRole('combobox') as HTMLSelectElement).value).toBe(NEW_WORKSPACE) + expect(screen.getByPlaceholderText('/absolute/path/to/project')).toBeTruthy() + }) +}) diff --git a/ui-web/src/sidebar/NewSessionDialog.tsx b/ui-web/src/sidebar/NewSessionDialog.tsx new file mode 100644 index 00000000..43d800f1 --- /dev/null +++ b/ui-web/src/sidebar/NewSessionDialog.tsx @@ -0,0 +1,242 @@ +import { useStore } from '@nanostores/react' +import { useEffect, useMemo, useRef, useState, type FormEvent } from 'react' + +import { createSession } from '../state/actions.ts' +import { $newSessionDialog, $projects, $workspace } from '../state/store.ts' +import { Button } from '../ui/primitives/Button.tsx' +import { XIcon } from '../ui/icons.tsx' +import css from './NewSessionDialog.module.css' + +/** The select value that reveals the folder-path field. */ +export const NEW_WORKSPACE = '__new_workspace__' + +/** The last path segment, for a label; the whole path when it has none. */ +function baseName(path: string): string { + const segments = path.split(/[/\\]/).filter(Boolean) + + return segments[segments.length - 1] ?? path +} + +/** The shape of a sidebar project this dialog reads: its path and its lanes'. */ +export interface WorkspaceSource { + path?: string | null + repos?: readonly { groups?: readonly { path?: string | null }[] }[] +} + +/** + * The workspaces the dialog offers: the current one first, then every folder + * the sidebar knows a session in — a repo and each of its worktree lanes, + * which are the natural places to start another session — without repeats. + * "Home" (sessions with no folder) has no path to start a session in, so it + * is not a choice. + */ +export function knownWorkspaces(current: string, projects: readonly WorkspaceSource[]): string[] { + const seen = new Set() + const paths: string[] = [] + const candidates = [ + current, + ...projects.flatMap(project => [ + project.path ?? '', + ...(project.repos ?? []).flatMap(repo => (repo.groups ?? []).map(lane => lane.path ?? '')), + ]), + ] + + for (const path of candidates) { + if (path === '' || seen.has(path)) continue + + seen.add(path) + paths.push(path) + } + + return paths +} + +export function openNewSessionDialog(): void { + $newSessionDialog.set(true) +} + +export function closeNewSessionDialog(): void { + $newSessionDialog.set(false) +} + +/** + * The New session dialog: which workspace, or a new one, and whether to + * isolate the session in a git worktree. + * + * A session runs somewhere, and until now the only somewhere was the current + * workspace: starting work in another project meant browsing to it first. + * The dialog puts the choice where the intent is. "Create new workspace…" + * takes an absolute path and makes the folder if it is not there yet; the + * worktree switch runs the session in a fresh checkout of the repo, the + * CLI's `--worktree`, so parallel sessions cannot step on each other's files. + * Errors stay in the dialog: a path the backend refuses is corrected here, + * not read off a status line behind a closed dialog. + */ +export function NewSessionDialog() { + const open = useStore($newSessionDialog) + + if (!open) return null + + return +} + +function NewSessionForm() { + const workspace = useStore($workspace) + const projects = useStore($projects) + const workspaces = useMemo(() => knownWorkspaces(workspace, projects), [projects, workspace]) + const [choice, setChoice] = useState(() => workspaces[0] ?? NEW_WORKSPACE) + const [path, setPath] = useState('') + const [worktree, setWorktree] = useState(false) + const [error, setError] = useState('') + const [creating, setCreating] = useState(false) + + const close = closeNewSessionDialog + + useEffect(() => { + const onKeyDown = (event: KeyboardEvent) => { + if (event.key === 'Escape') { + event.stopPropagation() + closeNewSessionDialog() + } + } + + document.addEventListener('keydown', onKeyDown) + + return () => { + document.removeEventListener('keydown', onKeyDown) + } + }, []) + + // The scrim closes on a click that STARTED on it: a drag that begins in + // the path field and ends outside must not throw the form away. + const pressedOnScrim = useRef(false) + + const creatingNew = choice === NEW_WORKSPACE + const target = creatingNew ? path.trim() : choice + const canCreate = target !== '' && !creating + + const onSubmit = async (event: FormEvent) => { + event.preventDefault() + + if (!canCreate) return + + setCreating(true) + setError('') + + const failure = await createSession({ + cwd: target, + ...(creatingNew && { createDir: true }), + ...(worktree && { worktree: true }), + }) + + setCreating(false) + + if (failure === null) close() + else setError(failure) + } + + return ( +
{ + if (pressedOnScrim.current && event.target === event.currentTarget) close() + + pressedOnScrim.current = false + }} + onMouseDown={event => { + pressedOnScrim.current = event.target === event.currentTarget + }} + > +
{ + event.stopPropagation() + }} + onSubmit={event => { + void onSubmit(event) + }} + role="dialog" + > +
+ + New session + + +
+ + + + {creatingNew && ( + + )} + + + + {error !== '' && ( +
+ {error} +
+ )} + +
+ + +
+
+
+ ) +} diff --git a/ui-web/src/sidebar/Sidebar.test.tsx b/ui-web/src/sidebar/Sidebar.test.tsx index e2afd1d9..303a6f1d 100644 --- a/ui-web/src/sidebar/Sidebar.test.tsx +++ b/ui-web/src/sidebar/Sidebar.test.tsx @@ -144,3 +144,16 @@ describe('workspace folder expansion', () => { expect(screen.getByText('gamma conversation')).toBeTruthy() }) }) + +describe('New session', () => { + it('opens the dialog instead of starting a session on the spot', async () => { + const { $newSessionDialog } = await import('../state/store.ts') + $newSessionDialog.set(false) + + render() + fireEvent.click(screen.getByRole('button', { name: 'New session' })) + + expect($newSessionDialog.get()).toBe(true) + $newSessionDialog.set(false) + }) +}) diff --git a/ui-web/src/sidebar/Sidebar.tsx b/ui-web/src/sidebar/Sidebar.tsx index 6c37f3b2..d3a34f65 100644 --- a/ui-web/src/sidebar/Sidebar.tsx +++ b/ui-web/src/sidebar/Sidebar.tsx @@ -2,7 +2,7 @@ import { useStore } from '@nanostores/react' import { useEffect, useMemo, useState } from 'react' import type { ProjectNode, SessionRow } from '../gateway/protocol.ts' -import { createSession, resumeSession } from '../state/actions.ts' +import { resumeSession } from '../state/actions.ts' import { toggleSidebar } from '../state/layout.ts' import { BrandMark } from '../ui/BrandMark.tsx' import { @@ -28,6 +28,7 @@ import { XIcon, } from '../ui/icons.tsx' import { filterProjects, isBlankSession, visibleSessions } from './filter.ts' +import { openNewSessionDialog } from './NewSessionDialog.tsx' import { absoluteTime, relativeTime } from './recency.ts' import css from './Sidebar.module.css' @@ -218,9 +219,7 @@ export function Sidebar({ collapsed }: SidebarProps) { {!collapsed && (