From 3d9ce6b94ae0f93b87c53ea0a4696ebf81f155ec Mon Sep 17 00:00:00 2001 From: arelchan <204152633+arelchan@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:31:57 +0800 Subject: [PATCH 1/5] fix(*): show the diff of a file a shell command rewrote A file an exec command wrote reached the desk diff as a bare row: the runtime listed the working directory on either side of the command and knew a file had changed, never what it changed from. A rewrite showed a bare M with no counts, and opening any such row fell back to the file viewer instead of the diff pane. The checkpoint's shadow git repo now supplies the missing half. Just before each command the tree is staged into a per-process index of its own (never the index the turn commit reads), and afterwards the old contents of every rewritten or removed file are read back from that tree. Each file_written entry gains added/removed counts and a unified diff (created files are measured against nothing), and a listed removal carries the text it held. The page draws those rows with their hunks, so they open in the diff pane like a file tool's. - The first staging of a directory the shadow repo has never indexed hashes every file. It is started in the background on the directory's first turn in the process (warm returns at once; repo setup runs on the staging thread too), so it overlaps the model's first reply. - A command waits at most 6s for its staging. Past that it is not run: the call fails with a reply telling the model to run it again, and the staging carries on. A command whose tree cannot be staged at all (checkpoint off, outside the repo, git error) runs without a diff. - Stagings run on daemon threads with blocking git calls: an asyncio subprocess still starting when its loop closes hangs the close on CPython 3.12 macOS. add and write-tree run under one lock per index, and the add has its own 600s ceiling instead of the 30s one, which killed a cold staging of a large tree over and over. - The diffs share the event's 512 KiB budget; past it the counts still go and the diff is dropped whole. Co-authored-by: Claude (claude-opus-5-5) --- CONTEXT.md | 7 +- raven/agent/loop/_shared.py | 97 +++-- raven/agent/loop/checkpoint.py | 305 ++++++++++++++- raven/agent/loop/turn_path.py | 54 ++- raven/rpc/models.py | 25 +- rpc-schema/openrpc.json | 23 +- tests/test_agent_loop_session_stamps.py | 335 ++++++++++++++++- tests/test_runtime_checkpoint.py | 368 +++++++++++++++++++ tests/test_runtime_checkpoint_deep.py | 6 + ui-tui/src/rpc/generated.ts | 14 +- ui-web/src/features/desk/store.ts | 6 +- ui-web/src/features/workspace/record.test.ts | 98 +++++ ui-web/src/features/workspace/record.ts | 44 ++- ui-web/src/features/workspace/types.ts | 4 + ui-web/src/rpc/fixtures/turn.test.ts | 5 +- ui-web/src/rpc/fixtures/turn.ts | 11 +- ui-web/src/rpc/generated.ts | 14 +- 17 files changed, 1351 insertions(+), 65 deletions(-) diff --git a/CONTEXT.md b/CONTEXT.md index ed55a7c61..479346522 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -566,7 +566,12 @@ _Avoid_: a fourth door -- a tool reaching the table any other way skips admissio A once-per-turn commit of the session workspace into a shadow git repo (separate from the user's `.git`), so an interrupted or failed turn can be rolled back. One `CheckpointService` per working directory, cached by `AgentLoop._turn_checkpoint()` and keyed on the directory -the running turn is bound to. +the running turn is bound to. The same repo also answers what a file an `exec` command +rewrote or removed held before it: `CheckpointService.warm` starts staging the tree when +a directory's first turn in this process starts, `stage_tree` stages it again just before +each command into a per-process index of its own (never the index the turn commit reads), +and `read_blobs` reads the old contents back for the command's `file_written` diff and +`file_removed` body. _Avoid_: "shadow git" as the term — Checkpoint is the per-turn snapshot it produces. **Empty-Response Recovery** (`agent/loop/recovery.py`): diff --git a/raven/agent/loop/_shared.py b/raven/agent/loop/_shared.py index cc86794cc..8d9dd54ff 100644 --- a/raven/agent/loop/_shared.py +++ b/raven/agent/loop/_shared.py @@ -8,6 +8,7 @@ from __future__ import annotations import asyncio +import difflib import json import os import time @@ -15,7 +16,7 @@ from dataclasses import dataclass, field, replace from datetime import datetime from pathlib import Path -from typing import TYPE_CHECKING, Any, Awaitable, Callable, Collection +from typing import TYPE_CHECKING, Any, Awaitable, Callable, Collection, Mapping from uuid import uuid4 from loguru import logger @@ -578,6 +579,16 @@ def _file_removed_payload(removals: Any) -> list[dict[str, Any]] | None: return out or None +#: What the model reads for a command held back because the tree its file +#: changes are measured against is still being staged. Not run rather than run +#: unmeasured: the staging carries on, and the same command a moment later is. +_EXEC_NOT_STAGED_REPLY = ( + "Error: the command was not run. Raven is still snapshotting the working directory, " + "which it needs to record what this command changes (the first snapshot of a large " + "directory takes a few seconds). Run the same command again." +) + + #: A file the listing found is counted in lines only when it is text this size #: or under. Past it the count is unknown rather than wrong: reading a gigabyte #: to number it would cost the turn more than the row it draws is worth. @@ -590,56 +601,100 @@ def _file_written_payload( after: dict[str, tuple[int, int]] | None, *, already: Collection[str] = (), + before: Mapping[str, bytes] | None = None, ) -> list[dict[str, Any]] | None: """The files a command left behind, as plain mappings, or ``None`` for none. The other half of ``_file_change_payload``: a file tool reports what it wrote, a command reports its output and nothing else, so this is read off - two listings of the working directory instead of off a result. Sizes and a - line count rather than contents -- one command can write a hundred files, - and what a row draws is that they were written and how big they are. + two listings of the working directory instead of off a result. ``lines`` belongs to a created file alone, and ``None`` there means unknown: - too large to read, or not text. A rewritten file has no count at all, since - the listing never held the old content and a number against nothing would - read as a change nobody measured. + too large to read, or not text. + + ``before`` is what a rewritten file held when the command started, when the + shadow repo could say (see ``CheckpointService.stage_tree``). With it, and + for every created file, the entry also carries the change itself: ``added`` + and ``removed`` line counts and a unified ``diff``. Without it a rewrite has + neither -- a number against contents nobody held would read as a change + somebody measured. The diffs share one budget per event, the way removals' + bodies do; past it the counts still go and the diff is dropped whole. ``already`` are the paths this same call accounted for by name. The listing sees those too, and reporting one again would draw a single write twice. """ accounted = {os.path.realpath(path) for path in already if isinstance(path, str) and path} + held = before or {} + budget = _FILE_CHANGE_MAX_CHARS out: list[dict[str, Any]] = [] - for path in created: + for path, was_created in [*((path, True) for path in created), *((path, False) for path in modified)]: if os.path.realpath(path) in accounted: continue size = (after or {}).get(path, (0, 0))[0] - out.append({"path": path, "created": True, "size": size, "lines": _text_line_count(path, size)}) - for path in modified: - if os.path.realpath(path) in accounted: - continue - out.append({"path": path, "created": False, "size": (after or {}).get(path, (0, 0))[0], "lines": None}) + old = "" if was_created else _decoded(held.get(path)) + text = None if old is None else _small_text(path, size) + entry: dict[str, Any] = { + "path": path, + "created": was_created, + "size": size, + "lines": len(text.splitlines()) if was_created and text is not None else None, + } + if text is not None and old is not None: + diff, added, removed = _line_diff(old, text, os.path.basename(path)) + entry["added"], entry["removed"] = added, removed + if diff and len(diff) <= budget: + entry["diff"] = diff + budget -= len(diff) + out.append(entry) return out or None -def _text_line_count(path: str, size: int) -> int | None: - """Lines in a file the listing found, or ``None`` when it cannot be counted.""" +def _small_text(path: str, size: int) -> str | None: + """A file the listing found, as text, or ``None`` when it is not worth reading.""" if size > _FILE_WRITTEN_TEXT_MAX_BYTES: return None try: - return len(Path(path).read_text(encoding="utf-8").splitlines()) + return Path(path).read_text(encoding="utf-8") except (OSError, UnicodeDecodeError): return None -def _listing_removals(deleted: Collection[str], *, already: Collection[str] = ()) -> list[FileRemoval]: +def _decoded(raw: bytes | None) -> str | None: + if raw is None or len(raw) > _FILE_WRITTEN_TEXT_MAX_BYTES: + return None + try: + return raw.decode("utf-8") + except UnicodeDecodeError: + return None + + +def _line_diff(old: str, new: str, name: str) -> tuple[str | None, int, int]: + """A unified diff of ``old`` to ``new`` with its added and removed line counts.""" + rows = list(difflib.unified_diff(old.splitlines(), new.splitlines(), fromfile=name, tofile=name, lineterm="")) + body = rows[2:] + added = sum(1 for row in body if row.startswith("+")) + removed = sum(1 for row in body if row.startswith("-")) + return ("\n".join(rows) if rows else None), added, removed + + +def _listing_removals( + deleted: Collection[str], *, already: Collection[str] = (), before: Mapping[str, bytes] | None = None +) -> list[FileRemoval]: """Files a listing says went, for the deletions no tool reported itself. - Without a body: the file was gone before anything read it, and the turn only - knows it was there when the command started. ``already`` are the removals - the call reported by name, which the listing sees as well. + With a body only when ``before`` holds one -- the shadow repo's copy from + just before the command. Otherwise the file was gone before anything read + it, and the turn only knows it was there when the command started. + ``already`` are the removals the call reported by name, which the listing + sees as well. """ accounted = {os.path.realpath(path) for path in already if isinstance(path, str) and path} - return [FileRemoval(path=path) for path in deleted if os.path.realpath(path) not in accounted] + held = before or {} + return [ + FileRemoval(path=path, before=_decoded(held.get(path))) + for path in deleted + if os.path.realpath(path) not in accounted + ] def monotonic() -> float: diff --git a/raven/agent/loop/checkpoint.py b/raven/agent/loop/checkpoint.py index 330b715a0..aed623133 100644 --- a/raven/agent/loop/checkpoint.py +++ b/raven/agent/loop/checkpoint.py @@ -30,7 +30,15 @@ from __future__ import annotations import asyncio +import concurrent.futures +import contextlib +import os +import shutil +import subprocess +import threading +import time from pathlib import Path +from typing import Collection from loguru import logger @@ -136,6 +144,50 @@ # fits comfortably; tight enough to detect a real hang within one turn. _GIT_TIMEOUT_SECONDS = 30.0 +# How long a command waits for the tree staged in front of it. Staging a tree +# git has seen before is a stat walk (tens of milliseconds on a repo of a few +# thousand files); the first one in a directory hashes and writes every file and +# can take seconds. Past this the command is held back and the staging carries +# on behind it, so the retry finds the index warm. +_STAGE_WAIT_SECONDS = 6.0 + +# The ceiling on a staging's ``git add``. Not ``_GIT_TIMEOUT_SECONDS``: that one +# bounds a call something waits behind, and a staging runs on a thread nothing +# waits behind past ``_STAGE_WAIT_SECONDS``. Killed at 30s, the first staging of +# a large directory (tens of seconds while a turn's commit hashes the same tree) +# was started over and killed again, and every command was held back for good. +_STAGE_ADD_TIMEOUT_SECONDS = 600.0 + +# A command's staging index left behind by a process that is gone. Age rather +# than a liveness probe: asking whether a pid is alive terminates it on Windows. +_STAGE_INDEX_STALE_SECONDS = 7 * 24 * 3600 + + +# The staging running on each index, across every service in this process: two +# services for one directory share the index file, so they share its one run. +_STAGING: dict[Path, "concurrent.futures.Future[str | None]"] = {} +_STAGE_LOCKS: dict[Path, threading.Lock] = {} + + +class StagingTimeoutError(Exception): + """The tree a command is to be measured against is still being staged.""" + + +async def _within(staging: "concurrent.futures.Future[str | None]", deadline: float) -> bool: + """Whether ``staging`` finished by ``deadline``, waited for without owning it. + + Detached rather than awaited: a staging still running when the loop closes + must not hold the close up, and a cancelled waiter is what lets its late + result be dropped once the loop is gone. + """ + waiter = asyncio.wrap_future(staging) + try: + done, _ = await asyncio.wait({waiter}, timeout=max(0.0, deadline - time.monotonic())) + finally: + if not waiter.done(): + waiter.cancel() + return bool(done) + class CheckpointService: """Shadow-git working-tree snapshots, one commit per turn.""" @@ -174,6 +226,15 @@ def __init__(self, workspace: Path, shadow_dir: str = ".raven/shadow.git") -> No self._shadow_rel = shadow_dir self._ready = False self._commit_count = 0 + self._stage_pruned = False + self._warmed = False + + def covers(self, path: Path | str) -> bool: + """Whether ``path`` lies in the work-tree this repo snapshots.""" + try: + return Path(path).expanduser().resolve().is_relative_to(self._workspace) + except OSError: + return False async def _git(self, *args: str) -> tuple[int, str, str]: """Run a git command against the shadow repo. Returns (rc, out, err). @@ -183,6 +244,10 @@ async def _git(self, *args: str) -> tuple[int, str, str]: without this, ``edited_files`` would land in the recovery prompt as ``"\\346\\265\\213"`` gibberish. """ + rc, out, err = await self._run(args) + return rc, out.decode(errors="replace"), err.decode(errors="replace") + + def _command(self, args: tuple[str, ...], index: Path | None) -> tuple[tuple[str, ...], dict[str, str] | None]: cmd = ( "git", f"--git-dir={self._git_dir}", @@ -191,16 +256,32 @@ async def _git(self, *args: str) -> tuple[int, str, str]: "core.quotePath=false", *args, ) + return cmd, None if index is None else {**os.environ, "GIT_INDEX_FILE": str(index)} + + async def _run( + self, + args: tuple[str, ...], + *, + index: Path | None = None, + stdin: bytes | None = None, + timeout: float | None = None, + ) -> tuple[int, bytes, bytes]: + """``_git`` with the raw bytes, and optionally against another index.""" + if timeout is None: + timeout = _GIT_TIMEOUT_SECONDS + cmd, env = self._command(args, index) proc = await asyncio.create_subprocess_exec( *cmd, + stdin=asyncio.subprocess.PIPE if stdin is not None else asyncio.subprocess.DEVNULL, stdout=asyncio.subprocess.PIPE, stderr=asyncio.subprocess.PIPE, cwd=str(self._workspace), + env=env, ) try: out, err = await asyncio.wait_for( - proc.communicate(), - timeout=_GIT_TIMEOUT_SECONDS, + proc.communicate() if stdin is None else proc.communicate(stdin), + timeout=timeout, ) except asyncio.TimeoutError: # NFS / index-lock / disk-full pathology: don't leak a zombie, @@ -213,11 +294,11 @@ async def _git(self, *args: str) -> tuple[int, str, str]: pass logger.debug( "checkpoint git timed out after {}s: {}", - _GIT_TIMEOUT_SECONDS, + timeout, " ".join(args[:2]), ) - return -1, "", "timeout" - return proc.returncode or 0, out.decode(errors="replace"), err.decode(errors="replace") + return -1, b"", b"timeout" + return proc.returncode or 0, out, err async def _ensure_init(self) -> bool: """Lazily initialize the shadow repo. Idempotent; returns readiness.""" @@ -305,6 +386,218 @@ async def commit_turn(self, label: str) -> tuple[str | None, list[str]]: logger.debug("checkpoint commit error: {}", exc) return None, [] + async def stage_tree(self) -> str | None: + """The work-tree as it stands, as a tree id in the shadow repo. + + Taken just before a command runs, so what the command changed can be + diffed afterwards against the contents it replaced -- a command names no + file it writes, and by the time it returns the old text is gone. Staged + into an index of its own: the per-turn commit reads the shared index as + "what the last turn left", and moving it here would make that commit miss + everything this turn did before the command. Nothing is committed; the + tree is only read back within the same call. + + ``None`` when the tree cannot be staged at all (git failed), which costs + the call its diff and nothing else. :class:`StagingTimeoutError` when it is + still being staged after :data:`_STAGE_WAIT_SECONDS`: the caller is to + hold the command back rather than run it unmeasured, because a retry a + moment later finds it done. A staging already running -- the warm-up + :meth:`warm` started, or an earlier call's -- is waited for inside the + same budget and then followed by a fresh one, which is then a stat walk. + One that runs out of budget is left running rather than killed, because + what it has hashed is what makes the next one fast. + """ + index = self._stage_path() + deadline = time.monotonic() + _STAGE_WAIT_SECONDS + # Started before this call, so its tree may predate what the calls in + # between wrote: waited out, never used. Waited for before the repo + # setup below, because a warm-up runs that same setup on its thread and + # two ``git config`` writes at once fail on the config lock. + earlier = _STAGING.get(index) + if earlier is not None and not earlier.done() and not await _within(earlier, deadline): + raise StagingTimeoutError + if not await self._ensure_init(): + return None + self._prepare_index(index) + # One another caller started while this one waited is as fresh as a new + # one would be, and a second ``git add`` on the same index would only + # fail on its lock. + staging = _STAGING.get(index) + if staging is None or staging is earlier: + staging = self._start_stage(index) + if not await _within(staging, deadline): + raise StagingTimeoutError + return staging.result() + + async def warm(self) -> None: + """Start a staging in the background, so the first command finds the index warm. + + The first staging in a directory the shadow repo has never indexed hashes + every file in it, which on a large tree is seconds. Started when the turn + starts, it runs while the model writes its first reply instead of in front + of the command. Returns at once: the repo's own setup (``git init``, its + config, seeding the index) runs on the staging's thread too, so the turn + waits for none of it. Once per service, which is once per directory per + process: after that every command's own staging keeps the index warm, and + a warm-up would only repeat the stat walk the next command does. + """ + if self._warmed: + return + self._warmed = True + index = self._stage_path() + running = _STAGING.get(index) + if running is not None and not running.done(): + return + staging: concurrent.futures.Future[str | None] = concurrent.futures.Future() + _STAGING[index] = staging + threading.Thread(target=self._warm_up, args=(index, staging), name="raven-stage", daemon=True).start() + + def _warm_up(self, index: Path, staging: "concurrent.futures.Future[str | None]") -> None: + # A loop of the thread's own for the repo setup: the turn's loop is not + # to wait on it, and one that closes while the setup is still starting a + # git process would hold its close up. + try: + ready = asyncio.run(self._ensure_init()) + except Exception as exc: # noqa: BLE001 -- a warm-up never breaks anything + logger.debug("checkpoint warm-up init error: {}", exc) + ready = False + if not ready: + staging.set_running_or_notify_cancel() + staging.set_result(None) + return + self._prepare_index(index) + self._stage(index, staging) + + def _start_stage(self, index: Path) -> "concurrent.futures.Future[str | None]": + staging: concurrent.futures.Future[str | None] = concurrent.futures.Future() + _STAGING[index] = staging + threading.Thread(target=self._stage, args=(index, staging), name="raven-stage", daemon=True).start() + return staging + + def _stage(self, index: Path, result: "concurrent.futures.Future[str | None]") -> None: + """``git add -A`` then ``git write-tree`` on the staging index, on a thread of its own. + + A thread and blocking runs rather than asyncio subprocesses: a staging + outlives the call that started it whenever it is slow, and an asyncio + subprocess still starting when its loop closes holds the close up for + good (CPython 3.12, macOS). Both steps under one lock per index, because + ``write-tree`` writes the index back too and a second staging's ``add`` + beside it fails on the index lock. + """ + result.set_running_or_notify_cancel() + with _STAGE_LOCKS.setdefault(index, threading.Lock()): + if self._stage_step(("add", "-A"), index, timeout=_STAGE_ADD_TIMEOUT_SECONDS) is None: + result.set_result(None) + return + tree = self._stage_step(("write-tree",), index) + result.set_result(tree or None) + + def _stage_step(self, args: tuple[str, ...], index: Path, *, timeout: float | None = None) -> str | None: + """One git step of a staging: its output, or ``None`` when it failed.""" + if timeout is None: + timeout = _GIT_TIMEOUT_SECONDS + cmd, env = self._command(args, index) + try: + done = subprocess.run( + cmd, + cwd=str(self._workspace), + env=env, + stdin=subprocess.DEVNULL, + capture_output=True, + timeout=timeout, + ) + except subprocess.TimeoutExpired: + logger.debug("checkpoint stage timed out after {}s: {}", timeout, " ".join(args)) + # A killed git leaves the index lock behind, and every later stage + # would then fail on it. The index is this process's own and the + # stage lock is held, so no live git owns the lock. + with contextlib.suppress(OSError): + Path(f"{index}.lock").unlink() + return None + except OSError as exc: + logger.debug("checkpoint stage error: {}", exc) + return None + if done.returncode != 0: + logger.debug("checkpoint stage failed: {}", done.stderr.decode(errors="replace").strip()) + return None + return done.stdout.decode(errors="replace").strip() + + async def read_blobs(self, tree: str, paths: Collection[str], *, max_bytes: int) -> dict[str, bytes]: + """What each of ``paths`` held in ``tree``, keyed by the path as given. + + A path is left out when the tree never had it (new, ignored, excluded), + when it lies outside the work-tree, or when it held more than + ``max_bytes`` -- a caller reads a missing key as "not known", never as + "empty". + """ + rel_of: dict[str, str] = {} + for path in paths: + try: + rel = Path(path).resolve().relative_to(self._workspace).as_posix() + except (OSError, ValueError): + continue + if "\n" not in rel: + rel_of[path] = rel + if not rel_of: + return {} + wanted = list(rel_of.items()) + query = "".join(f"{tree}:{rel}\n" for _, rel in wanted).encode() + rc, out, _ = await self._run(("cat-file", "--batch-check"), stdin=query) + lines = out.decode(errors="replace").splitlines() + if rc != 0 or len(lines) != len(wanted): + return {} + small: list[tuple[str, str]] = [] + for (path, _), line in zip(wanted, lines): + parts = line.split() + if len(parts) == 3 and parts[1] == "blob" and parts[2].isdigit() and int(parts[2]) <= max_bytes: + small.append((path, parts[0])) + if not small: + return {} + rc, out, _ = await self._run(("cat-file", "--batch"), stdin="".join(f"{sha}\n" for _, sha in small).encode()) + if rc != 0: + return {} + found: dict[str, bytes] = {} + at = 0 + for path, _ in small: + end = out.find(b"\n", at) + if end < 0: + break + header = out[at:end].split() + if len(header) != 3 or not header[2].isdigit(): + break + size = int(header[2]) + found[path] = out[end + 1 : end + 1 + size] + at = end + 1 + size + 1 + return found + + def _stage_index(self) -> Path: + """This process's staging index, seeded from the shared one. + + Per process because two processes can run commands in one directory and + an index is a single file. Seeded so a first staging finds git's stat + cache already warm from the last turn's commit instead of hashing the + whole tree again. + """ + index = self._stage_path() + self._prepare_index(index) + return index + + def _stage_path(self) -> Path: + return self._git_dir / f"exec-{os.getpid()}.index" + + def _prepare_index(self, index: Path) -> None: + if not self._stage_pruned: + self._stage_pruned = True + now = time.time() + for other in self._git_dir.glob("exec-*.index"): + with contextlib.suppress(OSError): + if now - other.stat().st_mtime > _STAGE_INDEX_STALE_SECONDS: + other.unlink() + shared = self._git_dir / "index" + if not index.exists() and shared.exists(): + with contextlib.suppress(OSError): + shutil.copyfile(shared, index) + async def _maybe_gc(self) -> None: """Periodic ``git gc --auto`` so long-lived sessions don't accumulate loose objects forever. ``--auto`` is a no-op below ``gc.auto`` (256 @@ -326,4 +619,4 @@ async def _maybe_gc(self) -> None: logger.debug("checkpoint gc failed: {}", err.strip()) -__all__ = ["CheckpointService"] +__all__ = ["CheckpointService", "StagingTimeoutError"] diff --git a/raven/agent/loop/turn_path.py b/raven/agent/loop/turn_path.py index ee806f243..a8a4c8456 100644 --- a/raven/agent/loop/turn_path.py +++ b/raven/agent/loop/turn_path.py @@ -10,6 +10,8 @@ from raven.agent.loop._shared import ( _ABORTED_ACTION_REPLY, _DELEGATED_KEY, + _EXEC_NOT_STAGED_REPLY, + _FILE_WRITTEN_TEXT_MAX_BYTES, _HOOK_INJECTED_KEY, _MAX_ITER_STATIC_FALLBACK, _MAX_ITER_SYNTHESIS_PROMPT, @@ -82,6 +84,7 @@ uuid4, workdir, ) +from raven.agent.loop.checkpoint import StagingTimeoutError from raven.agent.loop.dead_end import NO_RESPONSE_FALLBACK, dead_reasons from raven.agent.loop.first_call import FirstCallGuard from raven.agent.loop.recovery import ContinuationGate, DraftGate, cut_reasoning_head, lower_reasoning_effort @@ -103,6 +106,8 @@ from raven.token_wise.turn_spend import TurnSpend if TYPE_CHECKING: + from pathlib import Path + from raven.agent.loop.checkpoint import CheckpointService from raven.contracts.token_strategy import UsageSnapshot from raven.providers.base import ErrorClassification @@ -306,6 +311,19 @@ def _turn_checkpoint(self) -> "CheckpointService | None": self._checkpoints[target] = None return self._checkpoints[target] + async def _exec_baseline(self, root: Path) -> "tuple[CheckpointService, str] | None": + """The tree a command's changes are read against, staged just before it runs. + + Only where the turn's shadow repo covers the directory the command runs + in: a ``working_dir`` outside it has no copy of what its files held, and + a rewrite there is reported without a diff, as it always was. + """ + repo = self._turn_checkpoint() + if repo is None or not repo.covers(root): + return None + tree = await repo.stage_tree() + return None if tree is None else (repo, tree) + def _stash_recovery(self, session_key: str, outcome: "LoopOutcome") -> None: """Remember an interrupted turn's snapshot so the next turn in this session gets a recovery prompt. No-op unless checkpoint is enabled @@ -725,6 +743,10 @@ async def _run_agent_loop( # noqa: C901 (cc 100: pre-existing, above the ceilin # files is seen. Per turn for the reason the counters above are: the loop # is a singleton and another session's turn is running beside this one. removal_watch = RemovalWatch() + # Behind the model's first reply rather than in front of the first + # command, and only on a directory's first turn: see CheckpointService.warm. + if self.tools.get("exec") is not None and (repo := self._turn_checkpoint()) is not None: + await repo.warm() # Empty-response recovery state, local to the turn — the AgentLoop is a # long-lived singleton shared across sessions, so per-instance counters # would leak across turns; resetting here gives clean per-turn budgets. @@ -1301,6 +1323,7 @@ def _hook_rollback(decision) -> bool: # other call, which is what the empty diff of two Nones means. exec_before: workdir_snapshot.Snapshot | None = None exec_after: workdir_snapshot.Snapshot | None = None + exec_baseline: tuple[CheckpointService, str] | None = None tracker = self.strategies.get("usage_tracker") if tracker is not None: await tracker.record_tool_call(tool_call.name, tool_call.id) @@ -1357,9 +1380,19 @@ def _hook_rollback(decision) -> bool: ) if exec_root is not None: exec_before = await asyncio.to_thread(workdir_snapshot.take, exec_root) - result = await self.tools.execute( - tool_call.name, tool_call.arguments, run_meta=tool_call.run_meta - ) + held_back = False + if exec_before is not None: + try: + exec_baseline = await self._exec_baseline(exec_root) + except StagingTimeoutError: + held_back = True + exec_before = None + if held_back: + result = _EXEC_NOT_STAGED_REPLY + else: + result = await self.tools.execute( + tool_call.name, tool_call.arguments, run_meta=tool_call.run_meta + ) duration_ms = int((time.monotonic() - tool_t0) * 1000) if exec_before is not None: exec_after = await asyncio.to_thread(workdir_snapshot.take, exec_root) @@ -1411,15 +1444,26 @@ def _hook_rollback(decision) -> bool: if isinstance(getattr(tool_change, "path", None), str): accounted.append(tool_change.path) created, modified, deleted = workdir_snapshot.diff(exec_before, exec_after) + # What the rewritten and removed files held when the command + # started, from the tree staged in front of it: the listing + # knows they changed, and this is what they changed from. + exec_held: dict[str, bytes] = {} + if exec_baseline is not None and (modified or deleted): + baseline_repo, baseline_tree = exec_baseline + exec_held = await baseline_repo.read_blobs( + baseline_tree, [*modified, *deleted], max_bytes=_FILE_WRITTEN_TEXT_MAX_BYTES + ) # Off the loop for the reason the walks above are: this reads # every created file to number its lines, and one command can # create hundreds. The removals stay here -- they read nothing. tool_written = ( - await asyncio.to_thread(_file_written_payload, created, modified, exec_after, already=accounted) + await asyncio.to_thread( + _file_written_payload, created, modified, exec_after, already=accounted, before=exec_held + ) if created or modified else None ) - tool_removed.extend(_listing_removals(deleted, already=accounted)) + tool_removed.extend(_listing_removals(deleted, already=accounted, before=exec_held)) if emit_tool_event: await on_tool_event( "complete", diff --git a/raven/rpc/models.py b/raven/rpc/models.py index aaffccde5..a1bdefc95 100644 --- a/raven/rpc/models.py +++ b/raven/rpc/models.py @@ -656,8 +656,10 @@ class FileWritten(_Strict): Neither a :class:`FileChange` nor a :class:`FileRemoval`: a command reports its output and nothing else, so what is known of the file is what two listings of the directory said about it -- that it is there, how big it is, - and whether it was there before. No contents either way, because one command - can write a hundred files and a row draws none of their text. + and whether it was there before. What it changed from is known only when the + working directory's shadow repo held a copy from just before the command; + then, and for any created text file, the change itself rides along as + counts and a unified diff. """ path: str = Field(description="Absolute path of the file the command wrote.") @@ -677,6 +679,25 @@ class FileWritten(_Strict): "change therefore has no number." ), ) + added: int | None = Field( + default=None, + description=( + "Lines the command added to the file. Absent when the change could not be measured: " + "not text, too large, or a rewrite whose previous contents were never captured." + ), + ) + removed: int | None = Field( + default=None, + description="Lines the command removed from the file. Absent exactly when added is.", + ) + diff: str | None = Field( + default=None, + description=( + "Unified diff of the change, when it was measured and small enough to carry. " + "Absent past the event's budget even when the counts are present: a partial " + "diff reads as a smaller change than the one that happened." + ), + ) class ToolCompletePayload(_Strict): diff --git a/rpc-schema/openrpc.json b/rpc-schema/openrpc.json index ab746a5c6..447459985 100644 --- a/rpc-schema/openrpc.json +++ b/rpc-schema/openrpc.json @@ -12317,7 +12317,7 @@ } }, "FileWritten": { - "description": "One file a command left behind, found by listing its working directory. Neither a FileChange nor a FileRemoval: a command reports its output and nothing else, so what is known of the file is that it is there, how big it is, and whether it was there before.", + "description": "One file a command left behind, found by listing its working directory. Neither a FileChange nor a FileRemoval: a command reports its output and nothing else, so what is known of the file is that it is there, how big it is, and whether it was there before. What it changed from is known only when the working directory's shadow repo held a copy from just before the command; then, and for any created text file, the change rides along as counts and a unified diff.", "type": "object", "additionalProperties": false, "required": [ @@ -12344,6 +12344,27 @@ "null" ], "description": "Lines in a created file, when it could be counted. Null, not absent: the key is always sent, and null says the count is unknown. Too large to read, not text, or a file that already existed, whose change therefore has no number." + }, + "added": { + "type": [ + "integer", + "null" + ], + "description": "Lines the command added to the file. Absent when the change could not be measured: not text, too large, or a rewrite whose previous contents were never captured." + }, + "removed": { + "type": [ + "integer", + "null" + ], + "description": "Lines the command removed from the file. Absent exactly when added is." + }, + "diff": { + "type": [ + "string", + "null" + ], + "description": "Unified diff of the change, when it was measured and small enough to carry. Absent past the event's budget even when the counts are present: a partial diff reads as a smaller change than the one that happened." } } }, diff --git a/tests/test_agent_loop_session_stamps.py b/tests/test_agent_loop_session_stamps.py index 0b3a40c06..d7386aebd 100644 --- a/tests/test_agent_loop_session_stamps.py +++ b/tests/test_agent_loop_session_stamps.py @@ -10,9 +10,11 @@ from __future__ import annotations +import asyncio import json import tempfile import threading +import time from pathlib import Path from typing import Any @@ -21,7 +23,8 @@ from raven.agent import workdir from raven.agent.loop import AgentLoop from raven.agent.loop._shared import _FILE_WRITTEN_TEXT_MAX_BYTES -from raven.agent.loop.bundles import ToolWiring, TurnPolicy +from raven.agent.loop.bundles import EngineWiring, ToolWiring, TurnPolicy +from raven.config.raven import CheckpointConfig, RuntimeConfig from raven.contracts.tool import FileChange, FileRemoval, Tool, ToolResult from raven.providers.base import LLMProvider, LLMResponse from raven.spine.events import ToolEvent, ToolPhase @@ -56,6 +59,11 @@ def get_default_model(self) -> str: def workspace(): with tempfile.TemporaryDirectory() as td: yield Path(td) + # A turn's warm-up stages the tree on a thread of its own; removing the + # directory under a git still writing into it fails the cleanup. + for thread in threading.enumerate(): + if thread.name == "raven-stage": + thread.join(30) def _make_agent(workspace: Path) -> AgentLoop: @@ -574,20 +582,30 @@ async def execute(self, command: str = "", **kwargs: Any) -> Any: return "ran" -async def _run_command_turn(workspace: Path, work: Path, script: list[LLMResponse], *extra_tools: Tool): +async def _run_command_turn( + workspace: Path, + work: Path, + script: list[LLMResponse], + *extra_tools: Tool, + checkpoint: bool = False, + provider: LLMProvider | None = None, +): """One real turn whose working directory is ``work``, as a served turn has. Bound rather than defaulted so the listing covers the directory the command - ran in and not the session store beside it. The checkpoint is off because - its shadow repo is a second tree inside that same directory, built for a - recovery nothing here tests. + ran in and not the session store beside it. The checkpoint is off unless a + test asks for it: its shadow repo is what a command's diff is read against, + and without it a rewrite is reported with no measure of what changed. """ agent = AgentLoop( - provider=ScriptedProvider(script), + provider=provider or ScriptedProvider(script), workspace=workspace, model="stub", policy=TurnPolicy(max_iterations=4, interactive=False), tools=ToolWiring(restrict_to_workspace=True), + engine=EngineWiring( + runtime_config=RuntimeConfig(checkpoint=CheckpointConfig(policy="always" if checkpoint else "never")) + ), ) for tool in extra_tools: agent.tools.register(tool) @@ -638,9 +656,10 @@ async def test_a_file_a_command_created_reaches_the_event_and_the_stored_entry(w @pytest.mark.asyncio async def test_a_file_a_command_rewrote_is_not_reported_as_a_new_one(workspace): - """A rewrite carries no count. The listing holds sizes, never contents, so - the old text was never known and a number against it would be invented -- - and a client that drew this as a creation would claim the whole file is new.""" + """Without a shadow repo a rewrite carries no count. The listing holds sizes, + never contents, so the old text was never known and a number against it + would be invented -- and a client that drew this as a creation would claim + the whole file is new.""" work = workspace / "work" work.mkdir() kept = work / "kept.txt" @@ -659,6 +678,284 @@ async def test_a_file_a_command_rewrote_is_not_reported_as_a_new_one(workspace): assert written[0]["created"] is False assert written[0]["lines"] is None assert written[0]["size"] == len("three\nfour\nfive\n") + assert "added" not in written[0] and "removed" not in written[0] and "diff" not in written[0] + + +@pytest.mark.asyncio +async def test_a_file_a_command_rewrote_carries_its_diff_when_the_shadow_repo_held_it(workspace): + """The tree staged in front of the command holds what the file said, so the + change is measured the way a file tool's is: counts and a unified diff of + only what the command did, not the whole file again. And it is stored with + the entry, so a reload draws what the live page drew.""" + work = workspace / "work" + work.mkdir() + kept = work / "kept.txt" + kept.write_text("one\ntwo\nthree\n", encoding="utf-8") + + completes = await _run_command_turn( + workspace, + work, + _command_script(), + _CommandTool(lambda: kept.write_text("one\n2\nthree\nfour\n", encoding="utf-8")), + checkpoint=True, + ) + + written = completes[0]["file_written"] + assert len(written) == 1, written + assert written[0]["created"] is False + assert (written[0]["added"], written[0]["removed"]) == (2, 1) + body = written[0]["diff"].splitlines() + assert "-two" in body and "+2" in body and "+four" in body + assert " one" in body, "unchanged lines are context, not a rewrite" + tool_entry = next(m for m in _persisted_messages(workspace) if m.get("role") == "tool") + assert tool_entry["file_written"] == written + + +@pytest.mark.asyncio +async def test_a_file_a_command_created_carries_its_diff(workspace): + """A new file needs no earlier copy: everything in it was added. The same + shape as a rewrite's, so a client draws both the one way.""" + work = workspace / "work" + work.mkdir() + made = work / "made.txt" + + completes = await _run_command_turn( + workspace, work, _command_script(), _CommandTool(lambda: made.write_text("one\ntwo\n", encoding="utf-8")) + ) + + written = completes[0]["file_written"] + assert (written[0]["added"], written[0]["removed"], written[0]["lines"]) == (2, 0, 2) + assert written[0]["diff"].splitlines()[2:] == ["@@ -0,0 +1,2 @@", "+one", "+two"] + + +@pytest.mark.asyncio +async def test_a_command_that_ran_before_this_one_is_not_part_of_its_diff(workspace): + """The tree is staged in front of each command, not once per turn: a second + command's diff is what the second command did, against the file as the + first one left it.""" + work = workspace / "work" + work.mkdir() + kept = work / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + steps = iter(["one\ntwo\n", "one\ntwo\nthree\n"]) + script = [ + _tool_call("c1", "exec", {"command": "first"}), + _tool_call("c2", "exec", {"command": "second"}), + LLMResponse(content="done", finish_reason="stop"), + ] + + completes = await _run_command_turn( + workspace, + work, + script, + _CommandTool(lambda: kept.write_text(next(steps), encoding="utf-8")), + checkpoint=True, + ) + + assert [(c["file_written"][0]["added"], c["file_written"][0]["removed"]) for c in completes] == [(1, 0), (1, 0)] + assert "+three" in completes[1]["file_written"][0]["diff"].splitlines() + assert "+two" not in completes[1]["file_written"][0]["diff"].splitlines() + + +@pytest.mark.asyncio +async def test_a_rewrite_the_shadow_repo_does_not_hold_is_reported_without_a_diff(workspace): + """A file the shadow repo excludes (here a ``.env``, kept out as a likely + credential) has no earlier copy, so its rewrite is reported bare rather + than measured against nothing.""" + work = workspace / "work" + work.mkdir() + secret = work / ".env" + secret.write_text("A=1\n", encoding="utf-8") + + completes = await _run_command_turn( + workspace, + work, + _command_script(), + _CommandTool(lambda: secret.write_text("A=2\n", encoding="utf-8")), + checkpoint=True, + ) + + written = completes[0]["file_written"] + assert len(written) == 1, written + assert "added" not in written[0] and "diff" not in written[0] + + +@pytest.mark.asyncio +async def test_a_file_a_command_removed_carries_what_it_held_when_the_shadow_repo_had_it(workspace): + """The listing sees the file go after it is gone; the staged tree still has + it, which is the body a deletion row draws.""" + work = workspace / "work" + work.mkdir() + doomed = work / "doomed.txt" + doomed.write_text("one\ntwo\n", encoding="utf-8") + + completes = await _run_command_turn( + workspace, work, _command_script("rm doomed.txt"), _CommandTool(doomed.unlink), checkpoint=True + ) + + removed = completes[0]["file_removed"] + assert len(removed) == 1, removed + assert removed[0]["before"] == "one\ntwo\n" + tool_entry = next(m for m in _persisted_messages(workspace) if m.get("role") == "tool") + assert tool_entry["file_removed"] == [{"path": removed[0]["path"], "del": 2}] + + +@pytest.mark.asyncio +async def test_the_real_command_tool_rewrite_is_measured_against_the_staged_tree(workspace): + """The stub above changes files in Python. A real shell writes through its + own cwd, which is the tree the stage covers -- the one agreement between the + shell and the shadow repo only the real tool can show.""" + work = workspace / "work" + work.mkdir() + (work / "keep.md").write_text("one\n", encoding="utf-8") + + completes = await _run_command_turn(workspace, work, _command_script("echo two >> keep.md"), checkpoint=True) + + written = completes[0]["file_written"] + assert [Path(w["path"]).resolve() for w in written] == [(work / "keep.md").resolve()] + assert (written[0]["added"], written[0]["removed"]) == (1, 0) + assert "+two" in written[0]["diff"].splitlines() + + +class _SlowFirstReply(ScriptedProvider): + """A model that takes its time over the first reply, the way a real one does.""" + + def __init__(self, script, delay: float) -> None: + super().__init__(script) + self._delay = delay + + async def chat(self, *args: Any, **kwargs: Any) -> Any: + if self._delay: + await asyncio.sleep(self._delay) + self._delay = 0 + return await super().chat(*args, **kwargs) + + +@pytest.mark.asyncio +async def test_the_first_command_in_a_cold_directory_is_measured_when_the_model_took_its_time(workspace, monkeypatch): + """The first staging in a directory the shadow repo has never indexed hashes + the whole tree -- seconds on a large one, far past what a command waits. It + is started when the turn starts, so it runs while the model writes its first + reply, and the command that follows finds the index warm.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + work = workspace / "work" + work.mkdir() + kept = work / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + real = subprocess.run + cold = [True] + + def _cold_first(cmd, **kwargs): + if "add" in cmd and cold[0]: + cold[0] = False + time.sleep(1.0) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 0.6) + monkeypatch.setattr(cp_module.subprocess, "run", _cold_first) + script = _command_script() + + completes = await _run_command_turn( + workspace, + work, + script, + _CommandTool(lambda: kept.write_text("one\ntwo\n", encoding="utf-8")), + checkpoint=True, + provider=_SlowFirstReply(script, delay=2.5), + ) + + written = completes[0]["file_written"] + assert (written[0]["added"], written[0]["removed"]) == (1, 0), written + + +@pytest.mark.asyncio +async def test_a_command_is_held_back_until_the_snapshot_is_ready(workspace, monkeypatch): + """Run unmeasured, a command in a directory still being snapshotted leaves a + change nobody can show. It is not run instead: the call fails with a reply + that says why and to run it again, the snapshot carries on, and the retry + runs and is measured.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + from raven.agent.loop._shared import _EXEC_NOT_STAGED_REPLY + + work = workspace / "work" + work.mkdir() + kept = work / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + real = subprocess.run + cold = [True] + + def _cold_first(cmd, **kwargs): + if "add" in cmd and cold[0]: + cold[0] = False + time.sleep(0.5) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 0.1) + monkeypatch.setattr(cp_module.subprocess, "run", _cold_first) + runs: list[int] = [] + + def _append() -> None: + runs.append(1) + kept.write_text("one\ntwo\n", encoding="utf-8") + + script = [ + _tool_call("c1", "exec", {"command": "echo two >> kept.txt"}), + _tool_call("c2", "exec", {"command": "echo two >> kept.txt"}), + LLMResponse(content="done", finish_reason="stop"), + ] + + class _Retrying(ScriptedProvider): + async def chat(self, *args: Any, **kwargs: Any) -> Any: + if len(self._script) == 2: + # The retry comes after the model has read the refusal; by then + # the budget only has to cover a warm staging, however loaded + # the machine running this is. + await asyncio.sleep(0.6) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 30.0) + return await super().chat(*args, **kwargs) + + completes = await _run_command_turn( + workspace, work, script, _CommandTool(_append), checkpoint=True, provider=_Retrying(script) + ) + + held, retried = completes + assert held["ok"] is False + assert held["result_preview"] == _EXEC_NOT_STAGED_REPLY + assert held["file_written"] is None + assert retried["ok"] is True + assert (retried["file_written"][0]["added"], retried["file_written"][0]["removed"]) == (1, 0) + assert runs == [1], "the held-back call must not have run the command" + + +@pytest.mark.asyncio +async def test_a_command_outside_the_shadow_repo_is_reported_without_a_diff(workspace): + """``working_dir`` can point anywhere, and the turn's shadow repo only holds + its own directory. Out there a rewrite is reported bare, as it always was, + rather than read against a tree that never held the file.""" + work = workspace / "work" + work.mkdir() + other = workspace / "other" + other.mkdir() + kept = other / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + + completes = await _run_command_turn( + workspace, + work, + _command_script("printf 'two\\n' >> kept.txt", working_dir=str(other)), + checkpoint=True, + ) + + written = completes[0]["file_written"] + assert [Path(w["path"]).resolve() for w in written] == [kept.resolve()] + assert "added" not in written[0] and "diff" not in written[0] @pytest.mark.asyncio @@ -947,3 +1244,23 @@ def watched(root: Any) -> Any: assert roots == [] assert completes[0]["file_written"] is None + + +def test_past_the_events_diff_budget_the_counts_still_go_and_the_diff_does_not(tmp_path, monkeypatch): + """One command can rewrite a hundred files. The diffs share the event's + budget, and one past it is dropped whole -- half a diff reads as a smaller + change -- while its counts, which cost nothing, still say how big it was.""" + from raven.agent.loop import _shared + + monkeypatch.setattr(_shared, "_FILE_CHANGE_MAX_CHARS", 60) + first = tmp_path / "a.txt" + second = tmp_path / "b.txt" + first.write_text("one\ntwo\n", encoding="utf-8") + second.write_text("three\nfour\n", encoding="utf-8") + after = {str(p): (p.stat().st_size, 0) for p in (first, second)} + + out = _shared._file_written_payload([str(first), str(second)], [], after) + + assert out is not None + assert "diff" in out[0] and "diff" not in out[1] + assert (out[1]["added"], out[1]["removed"]) == (2, 0) diff --git a/tests/test_runtime_checkpoint.py b/tests/test_runtime_checkpoint.py index 2f7276a2e..712c8ff33 100644 --- a/tests/test_runtime_checkpoint.py +++ b/tests/test_runtime_checkpoint.py @@ -11,8 +11,11 @@ from __future__ import annotations +import asyncio import subprocess import tempfile +import threading +import time from pathlib import Path import pytest @@ -29,6 +32,11 @@ def workspace(): with tempfile.TemporaryDirectory() as td: yield Path(td) + # A turn's warm-up stages the tree on a thread of its own; removing the + # directory under a git still writing into it fails the cleanup. + for thread in threading.enumerate(): + if thread.name == "raven-stage": + thread.join(30) async def _run_turn_body(agent: AgentLoop, workspace: Path): @@ -80,6 +88,366 @@ async def test_checkpoint_does_not_touch_user_git(workspace): assert _git_count(workspace) == count_before, "no commits added to user repo" +async def test_a_staged_tree_holds_what_a_file_said_before_it_changed(workspace): + """What a command's diff is read against: the file as it stood when the tree + was staged, after the file has been rewritten on disk.""" + svc = CheckpointService(workspace) + kept = workspace / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + + tree = await svc.stage_tree() + assert tree is not None + kept.write_text("two\n", encoding="utf-8") + (workspace / "new.txt").write_text("new\n", encoding="utf-8") + + held = await svc.read_blobs(tree, [str(kept), str(workspace / "new.txt")], max_bytes=1024) + assert held == {str(kept): b"one\n"}, "a file the tree never had is unknown, not empty" + + +async def test_a_blob_is_left_out_past_the_cap_or_outside_the_work_tree(workspace, tmp_path_factory): + svc = CheckpointService(workspace) + big = workspace / "big.txt" + big.write_text("x" * 64, encoding="utf-8") + outside = tmp_path_factory.mktemp("outside") / "o.txt" + outside.write_text("o\n", encoding="utf-8") + + tree = await svc.stage_tree() + assert tree is not None + + assert await svc.read_blobs(tree, [str(big), str(outside)], max_bytes=63) == {} + assert await svc.read_blobs(tree, [str(big)], max_bytes=64) == {str(big): b"x" * 64} + + +async def test_staging_a_tree_leaves_the_turn_commit_its_own_changes(workspace): + """The stage has an index of its own. Staged into the shared one, the turn's + commit would find this file already staged and report the turn as having + changed nothing -- the recovery prompt would lose it.""" + svc = CheckpointService(workspace) + (workspace / "a.py").write_text("print(1)\n", encoding="utf-8") + + assert await svc.stage_tree() is not None + cid, changed = await svc.commit_turn("turn 1") + + assert cid is not None + assert changed == ["a.py"] + + +async def test_a_stale_staging_index_is_pruned_and_a_fresh_one_kept(workspace): + """A process that died leaves its staging index behind. Pruned by age, never + by asking whether its pid is alive -- on Windows that question kills it.""" + import os + import time + + svc = CheckpointService(workspace) + assert await svc.stage_tree() is not None + git_dir = workspace / ".raven" / "shadow.git" + stale = git_dir / "exec-1.index" + fresh = git_dir / "exec-2.index" + stale.write_bytes(b"") + fresh.write_bytes(b"") + old = time.time() - 8 * 24 * 3600 + os.utime(stale, (old, old)) + + again = CheckpointService(workspace) + assert await again.stage_tree() is not None + + assert not stale.exists() + assert fresh.exists() + assert (git_dir / f"exec-{os.getpid()}.index").exists() + + +async def test_a_stage_that_timed_out_leaves_no_lock_behind(workspace, monkeypatch): + """A killed ``git add`` leaves its ``index.lock``, and every later stage would + fail on it for the rest of the process -- one slow command would cost every + command after it its diff.""" + import os + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + svc = CheckpointService(workspace) + assert await svc.stage_tree() is not None + lock = workspace / ".raven" / "shadow.git" / f"exec-{os.getpid()}.index.lock" + real = subprocess.run + + def _killed(cmd, **kwargs): + lock.write_bytes(b"") + raise subprocess.TimeoutExpired(cmd, 0.05) + + monkeypatch.setattr(cp_module.subprocess, "run", _killed) + assert await svc.stage_tree() is None + assert not lock.exists() + + monkeypatch.setattr(cp_module.subprocess, "run", real) + assert await svc.stage_tree() is not None + + +async def test_a_slow_first_stage_is_left_to_warm_the_index_behind_the_command(workspace, monkeypatch): + """The first stage in a directory hashes every file and can take seconds. + Past the budget the call is told so, the stage keeps running, a command + that arrives meanwhile does not start a second one, and the one after it + finds the index warm.""" + import asyncio + import subprocess + import threading + + import raven.agent.loop.checkpoint as cp_module + + svc = CheckpointService(workspace) + (workspace / "a.txt").write_text("a\n", encoding="utf-8") + real = subprocess.run + release = threading.Event() + adds: list[object] = [] + + def _slow(cmd, **kwargs): + if "add" in cmd: + adds.append(cmd) + release.wait(10) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 0.05) + monkeypatch.setattr(cp_module.subprocess, "run", _slow) + + with pytest.raises(cp_module.StagingTimeoutError): + await svc.stage_tree() + with pytest.raises(cp_module.StagingTimeoutError): + await svc.stage_tree() + assert len(adds) == 1, "a stage still running is not started again" + + release.set() + assert await asyncio.wrap_future(cp_module._STAGING[svc._stage_path()]) is not None + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 30.0) + assert await svc.stage_tree() is not None + + +async def test_a_stage_waits_out_the_warm_up_and_stages_again(workspace, monkeypatch): + """The warm-up started with the turn may still be running when the first + command arrives. Its tree was taken before whatever the turn has done since, + so it is waited for and never used; the staging after it is the one read.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + svc = CheckpointService(workspace) + (workspace / "a.txt").write_text("a\n", encoding="utf-8") + real = subprocess.run + adds: list[object] = [] + + def _count(cmd, **kwargs): + if "add" in cmd: + adds.append(cmd) + if len(adds) == 1: + time.sleep(0.3) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module.subprocess, "run", _count) + + await svc.warm() + (workspace / "a.txt").write_text("changed while warming\n", encoding="utf-8") + tree = await svc.stage_tree() + + assert tree is not None + assert len(adds) == 2 + held = await svc.read_blobs(tree, [str(workspace / "a.txt")], max_bytes=1024) + assert held == {str(workspace / "a.txt"): b"changed while warming\n"} + + +async def test_a_warm_up_already_running_is_not_started_again(workspace, monkeypatch): + """Two turns starting in one directory warm it once: the second warm-up + would only queue a second hash of the same tree behind the first.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + adds: list[object] = [] + + def _count(cmd, **kwargs): + if "add" in cmd: + adds.append(cmd) + time.sleep(0.2) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module.subprocess, "run", _count) + svc = CheckpointService(workspace) + + await svc.warm() + await CheckpointService(workspace).warm() + await asyncio.wrap_future(cp_module._STAGING[svc._stage_path()]) + + assert len(adds) == 1 + + +async def test_a_warm_up_returns_before_the_repo_is_even_set_up(workspace, monkeypatch): + """The turn awaits ``warm`` in front of its first model call, so everything + the warm-up does -- the repo's own ``git init`` and config as much as the + staging -- belongs on its thread. Awaited, a slow setup would be a first + reply that waits for it.""" + import raven.agent.loop.checkpoint as cp_module + + real = CheckpointService._ensure_init + + async def _slow_init(self) -> bool: + await asyncio.sleep(1.0) + return await real(self) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(CheckpointService, "_ensure_init", _slow_init) + svc = CheckpointService(workspace) + + started = time.monotonic() + await svc.warm() + assert time.monotonic() - started < 0.2 + + assert await asyncio.wrap_future(cp_module._STAGING[svc._stage_path()]) is not None + + +async def test_a_directory_is_warmed_once_and_not_every_turn(workspace, monkeypatch): + """Every turn starts with a warm-up call, and only the first does any work: + after it each command's own staging keeps the index warm, so another would + be one more stat walk of the whole tree per message for nothing.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + adds: list[object] = [] + + def _count(cmd, **kwargs): + if "add" in cmd: + adds.append(cmd) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module.subprocess, "run", _count) + svc = CheckpointService(workspace) + + await svc.warm() + await asyncio.wrap_future(cp_module._STAGING[svc._stage_path()]) + await svc.warm() + await svc.warm() + + assert len(adds) == 1 + + +async def test_two_commands_waiting_on_one_warm_up_share_the_staging_after_it(workspace, monkeypatch): + """Two sessions in one directory both find the warm-up running. The first to + wake starts the next staging; the second takes that one rather than starting + a third ``git add`` on the same index, which would fail on its lock.""" + import asyncio + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + svc = CheckpointService(workspace) + (workspace / "a.txt").write_text("a\n", encoding="utf-8") + real = subprocess.run + adds: list[object] = [] + + def _count(cmd, **kwargs): + if "add" in cmd: + adds.append(cmd) + if len(adds) == 1: + time.sleep(0.3) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module.subprocess, "run", _count) + + await svc.warm() + first, second = await asyncio.gather(svc.stage_tree(), CheckpointService(workspace).stage_tree()) + + assert first is not None and second is not None + assert len(adds) == 2 + + +async def test_a_slow_first_staging_is_not_cut_off_at_the_git_call_ceiling(workspace, monkeypatch): + """Every command is held back until the staging finishes. Killed at the + ceiling a turn's own git calls use, a staging that needs longer is started + over, killed again, and never finishes -- so no command would ever run.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + ceilings: list[float] = [] + + def _spy(cmd, **kwargs): + if "add" in cmd: + ceilings.append(kwargs["timeout"]) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module.subprocess, "run", _spy) + assert await CheckpointService(workspace).stage_tree() is not None + + assert ceilings == [cp_module._STAGE_ADD_TIMEOUT_SECONDS] + assert ceilings[0] > cp_module._GIT_TIMEOUT_SECONDS * 10 + + +async def test_a_first_stage_starts_from_the_last_turns_index(workspace, monkeypatch): + """The turn's commit has already hashed the tree into the shared index, and a + staging index copied from it only has to stat what changed since. Started + empty, the first command in every process would pay for hashing the whole + tree again.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + (workspace / "a.txt").write_text("a\n", encoding="utf-8") + assert (await CheckpointService(workspace).commit_turn("turn 1"))[0] is not None + shared = (workspace / ".raven" / "shadow.git" / "index").read_bytes() + real = subprocess.run + started_from: list[bytes] = [] + + def _spy(cmd, **kwargs): + if "add" in cmd: + started_from.append(Path(kwargs["env"]["GIT_INDEX_FILE"]).read_bytes()) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module.subprocess, "run", _spy) + assert await CheckpointService(workspace).stage_tree() is not None + assert started_from == [shared] + + +def test_a_staging_still_running_does_not_hold_the_loop_open(workspace, monkeypatch): + """Closing a loop cancels what is left on it. An asyncio subprocess still + starting at that moment never finishes cancelling (CPython 3.12, macOS) and + the close hangs -- which, for the gateway, is a shutdown that never ends.""" + import asyncio + import subprocess + import threading + import time + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + release = threading.Event() + + def _slow(cmd, **kwargs): + release.wait(10) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 0.05) + monkeypatch.setattr(cp_module.subprocess, "run", _slow) + + async def one_call() -> bool: + try: + await CheckpointService(workspace).stage_tree() + except cp_module.StagingTimeoutError: + return True + return False + + started = time.monotonic() + assert asyncio.run(one_call()) is True + assert time.monotonic() - started < 5 + release.set() + + def _os_environ(): import os diff --git a/tests/test_runtime_checkpoint_deep.py b/tests/test_runtime_checkpoint_deep.py index c1150a9aa..decea32e3 100644 --- a/tests/test_runtime_checkpoint_deep.py +++ b/tests/test_runtime_checkpoint_deep.py @@ -17,6 +17,7 @@ import os import shutil import tempfile +import threading from pathlib import Path import pytest @@ -34,6 +35,11 @@ def workspace(): with tempfile.TemporaryDirectory() as td: yield Path(td) + # A turn's warm-up stages the tree on a thread of its own; removing the + # directory under a git still writing into it fails the cleanup. + for thread in threading.enumerate(): + if thread.name == "raven-stage": + thread.join(30) async def _run_turn_body(agent: AgentLoop, workspace: Path): diff --git a/ui-tui/src/rpc/generated.ts b/ui-tui/src/rpc/generated.ts index 5072b18da..beb8d74ad 100644 --- a/ui-tui/src/rpc/generated.ts +++ b/ui-tui/src/rpc/generated.ts @@ -298,7 +298,7 @@ export interface TranscriptFileRemoval { del: number; } /** - * One file a command left behind, found by listing its working directory. Neither a FileChange nor a FileRemoval: a command reports its output and nothing else, so what is known of the file is that it is there, how big it is, and whether it was there before. + * One file a command left behind, found by listing its working directory. Neither a FileChange nor a FileRemoval: a command reports its output and nothing else, so what is known of the file is that it is there, how big it is, and whether it was there before. What it changed from is known only when the working directory's shadow repo held a copy from just before the command; then, and for any created text file, the change rides along as counts and a unified diff. * * This interface was referenced by `RavenRpcRoot`'s JSON-Schema * via the `definition` "FileWritten". @@ -320,6 +320,18 @@ export interface FileWritten { * Lines in a created file, when it could be counted. Null, not absent: the key is always sent, and null says the count is unknown. Too large to read, not text, or a file that already existed, whose change therefore has no number. */ lines?: number | null; + /** + * Lines the command added to the file. Absent when the change could not be measured: not text, too large, or a rewrite whose previous contents were never captured. + */ + added?: number | null; + /** + * Lines the command removed from the file. Absent exactly when added is. + */ + removed?: number | null; + /** + * Unified diff of the change, when it was measured and small enough to carry. Absent past the event's budget even when the counts are present: a partial diff reads as a smaller change than the one that happened. + */ + diff?: string | null; } /** * Why a turn's transcript stops where it does. diff --git a/ui-web/src/features/desk/store.ts b/ui-web/src/features/desk/store.ts index eccd0d284..a117bfd47 100644 --- a/ui-web/src/features/desk/store.ts +++ b/ui-web/src/features/desk/store.ts @@ -491,9 +491,9 @@ export function openDeskFile(path: string): void { export function openDeskDiff(change: WsChange): void { readItem('diff', `${change.key}:${change.turn}`) - /* A row with no hunks has no patch to draw: a command reports the files it - left behind and never how it changed them, so the listing that made the row - knows a count and nothing else. The file as it stands is the nearest thing + /* A row with no hunks has no patch to draw: a command whose change the + runtime could not measure (no earlier copy of the file, not text, or too + large) leaves a row that knows a count at most. The file as it stands is the nearest thing to the change and is what the reader clicked for -- an empty patch pane is not. A removal keeps its pane: there the missing hunk IS the answer, and there is no file left to open. */ diff --git a/ui-web/src/features/workspace/record.test.ts b/ui-web/src/features/workspace/record.test.ts index 08a878b5a..1e6fff0c2 100644 --- a/ui-web/src/features/workspace/record.test.ts +++ b/ui-web/src/features/workspace/record.test.ts @@ -488,6 +488,104 @@ describe('recording the files a command left behind', () => { expect(store.shared().changes.map((c) => c.kind).sort()).toEqual(['delete', 'write']) }) + /* The runtime read the file against the tree it staged in front of the + command, so the change comes measured, and the row is drawn like any file + tool's -- which is what sends a reader who opens it to the patch. */ + it('draws a measured rewrite with its hunk and counts', () => { + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ + path: '/w/log.json', created: false, size: 12, lines: null, added: 1, removed: 1, + diff: '--- log.json\n+++ log.json\n@@ -1,2 +1,2 @@\n a\n-b\n+B', + }]) + + const row = rowFor('/w/log.json') + expect(row?.kind).toBe('write') + expect(row?.hunks).toHaveLength(1) + expect([row?.add, row?.del]).toEqual([1, 1]) + expect(row?.listed).toBe(true) + }) + + /* Past the event's budget the diff is dropped and the counts still arrive: + the row says how big the change was even with no patch to open. */ + it('carries the counts of a measured change that came without its diff', () => { + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, + [{ path: '/w/log.json', created: false, size: 12, lines: null, added: 40, removed: 7 }]) + + const row = rowFor('/w/log.json') + expect(row?.hunks).toHaveLength(0) + expect([row?.add, row?.del]).toEqual([40, 7]) + }) + + /* A command that changes a file a tool already wrote this turn: its diff is + against the file as the tool left it, so it is one more hunk on the row, + not a second row and not a replacement for the first. */ + it('adds a measured change to the row a file tool already made', () => { + const args = { path: '/w/notes.md', content: 'a\n' } + wsOnTool('write_file', args) + wsOnToolDone('write_file', args, true, '', null, + '--- a/w/notes.md\n+++ b/w/notes.md\n@@ -0,0 +1,1 @@\n+a', { path: '/w/notes.md', after: 'a\n' }) + + wsOnToolDone('exec', { command: 'echo b >> notes.md' }, true, '', null, undefined, undefined, undefined, [{ + path: '/w/notes.md', created: false, size: 4, lines: null, added: 1, removed: 0, + diff: '--- notes.md\n+++ notes.md\n@@ -1,1 +1,2 @@\n a\n+b', + }]) + + expect(store.shared().changes).toHaveLength(1) + const row = rowFor('/w/notes.md') + expect(row?.kind).toBe('add') + expect(row?.hunks).toHaveLength(2) + expect(row?.add).toBe(2) + }) + + /* A listing row that carries a hunk is still the listing's: a file tool that + follows under the path the model typed takes it over rather than opening a + second row for the same file. */ + it('lets a file tool take over a command\'s measured row', () => { + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ + path: '/w/notes.md', created: true, size: 2, lines: 1, added: 1, removed: 0, + diff: '--- notes.md\n+++ notes.md\n@@ -0,0 +1,1 @@\n+b', + }]) + + const args = { path: 'notes.md', old_text: 'b', new_text: 'B' } + wsOnTool('edit_file', args) + + expect(store.shared().changes).toHaveLength(1) + const row = rowFor('notes.md') + expect(row?.kind).toBe('add') + expect(row?.hunks).toHaveLength(2) + }) + + /* The row's change so far was never measured, so a measured second change + would put a partial count on it and a patch that shows only the half the + reader did not ask about. */ + it('leaves an unmeasured row bare when a later command is measured', () => { + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, + [{ path: '/w/.env', created: false, size: 4, lines: null }]) + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ + path: '/w/.env', created: false, size: 4, lines: null, added: 1, removed: 1, + diff: '--- .env\n+++ .env\n@@ -1,1 +1,1 @@\n-A=1\n+A=2', + }]) + + const row = rowFor('/w/.env') + expect(row?.hunks).toHaveLength(0) + expect([row?.add, row?.del]).toEqual([0, 0]) + }) + + it('replays a stored measured change as the same row', () => { + const written = { + path: '/w/log.json', created: false, size: 12, lines: null, added: 1, removed: 1, + diff: '--- log.json\n+++ log.json\n@@ -1,2 +1,2 @@\n a\n-b\n+B', + } + wsOnHistory([ + { role: 'user', text: 'regenerate it' }, + { role: 'assistant', tool_calls: [{ id: 'c1', name: 'exec', arguments: JSON.stringify({ command: 'make' }) }] }, + { role: 'tool', tool_call_id: 'c1', file_written: [written] }, + ]) + + const row = rowFor('/w/log.json') + expect(row?.hunks).toHaveLength(1) + expect([row?.add, row?.del]).toEqual([1, 1]) + }) + /* A reload reads the same shape back: unlike a removal there is nothing to reduce, so live and replayed rows are identical. */ it('replays the stored listing as the same rows', () => { diff --git a/ui-web/src/features/workspace/record.ts b/ui-web/src/features/workspace/record.ts index 25cf247e6..00ca53870 100644 --- a/ui-web/src/features/workspace/record.ts +++ b/ui-web/src/features/workspace/record.ts @@ -132,33 +132,57 @@ function wsRecordRemoval(path: string, before?: string, lines?: number | null): desk would list one file twice. The tool's account is the one that can say what changed, so the listing's row carries on under the tool's spelling instead -- keeping the verdict the listing is better placed to know, that - the file was new. Only a row the listing made is taken this way, which is - what carrying no hunk means; a removal's bare row is its own answer. */ + the file was new. Only a row the listing made is taken this way; a + removal's bare row is its own answer. */ function adoptListing(key: string): void { const WS = record() const row = WS.changes.find((x) => x.turn === WS.turn && x.key !== key - && !x.hunks.length && x.kind !== 'delete' && sameFile(x.key, key)) + && x.listed && x.kind !== 'delete' && sameFile(x.key, key)) if (!row) return const { dir, name } = labelFor(key) row.key = key row.dir = dir row.name = name + row.listed = false } /* What a command left behind, which no tool result names: the runtime lists the directory the turn's tools run in before and after an `exec` and reports the - difference. A listing knows a file is there, how big it is and whether it was - there before -- never how it changed -- so the row carries a count and no - hunk, and the desk sends a reader who opens it to the file itself. + difference. Where it also held the file's previous contents (the working + directory's shadow repo, staged just before the command) or the file is new, + it measured the change and sends the diff, and the row is drawn like any + other. Otherwise the row carries what it can -- the counts, or for a created + file its lines -- and no hunk, and the desk sends a reader who opens it to + the file itself. - A row this turn already holds for the path is left alone: it came from a file - tool, whose arguments say everything a listing cannot. */ + A row this turn already holds for the path takes a measured change as one + more hunk, the way a second edit does: the diff was read against the file as + the earlier calls left it, so it is only what the command did. An unmeasured + one is dropped there -- it says nothing the row does not, and a row whose + change so far is unmeasured would read a partial count as the whole. */ function wsRecordWritten(w: FileWritten): void { const WS = record() const key = String(w.path) - if (WS.changes.some((x) => x.turn === WS.turn && sameFile(x.key, key))) return + const hunk = w.diff ? hunks.fromUnified(w.diff) : null + const had = WS.changes.find((x) => x.turn === WS.turn && sameFile(x.key, key)) + if (had) { + if (hunk && had.hunks.length && had.kind !== 'delete') { + had.hunks.push(hunk) + had.add += hunk.add + had.del += hunk.del + } + return + } const c = rowFor(key, w.created ? 'add' : 'write') - if (w.created) c.add = w.lines == null ? 0 : w.lines + c.listed = true + if (hunk) { + c.hunks.push(hunk) + c.add = hunk.add + c.del = hunk.del + } else if (w.added != null) { + c.add = w.added + c.del = w.removed == null ? 0 : w.removed + } else if (w.created) c.add = w.lines == null ? 0 : w.lines } /* ── tool-event hooks ────────────────────────────────────────────────── diff --git a/ui-web/src/features/workspace/types.ts b/ui-web/src/features/workspace/types.ts index 87774731b..8a6cb91f2 100644 --- a/ui-web/src/features/workspace/types.ts +++ b/ui-web/src/features/workspace/types.ts @@ -22,6 +22,10 @@ export interface WsChange { edited -- either way nothing here can say what the file holds. Read only when the file is removed and the runtime caught none of its contents. */ body?: string | null + /* Made by a command's listing rather than by a file tool. Said outright + because the hunk no longer tells them apart: a listing that could measure + the change carries one too. */ + listed?: boolean turn: number open: boolean auto?: boolean diff --git a/ui-web/src/rpc/fixtures/turn.test.ts b/ui-web/src/rpc/fixtures/turn.test.ts index 8d4991c64..129c64bfe 100644 --- a/ui-web/src/rpc/fixtures/turn.test.ts +++ b/ui-web/src/rpc/fixtures/turn.test.ts @@ -41,7 +41,10 @@ describe('the frames a scripted turn pushes', () => { .filter((e) => e?.type === 'tool.complete' && e.payload?.file_written) .map((e) => (e!.payload as unknown as ToolCompleteEvent['payload']).file_written) expect(written).toEqual([[ - { path: '~/work/raven/research/tally.txt', created: true, size: 96, lines: 4 }, + { + path: '~/work/raven/research/tally.txt', created: true, size: 89, lines: 4, added: 4, removed: 0, + diff: '--- tally.txt\n+++ tally.txt\n@@ -0,0 +1,4 @@\n+vendor leads replies\n+Clay 1240 88\n+Apollo 980 61\n+Unify 410 37', + }, { path: '~/work/raven/research/run.log', created: false, size: 412, lines: null }, ]]) }) diff --git a/ui-web/src/rpc/fixtures/turn.ts b/ui-web/src/rpc/fixtures/turn.ts index 01b382fdf..f75bc6932 100644 --- a/ui-web/src/rpc/fixtures/turn.ts +++ b/ui-web/src/rpc/fixtures/turn.ts @@ -60,7 +60,8 @@ export type ScriptEvent = { d?: number } & ( | { t: 't+'; id: number; n: string; a?: string | ToolArgs } | { t: 't-'; id: number; r: string; ok?: boolean; ms?: number; diff?: string[]; meta?: DeliveryMeta; removed?: Array<{ path: string; before: string }>; - wrote?: Array<{ path: string; created: boolean; size: number; lines: number | null }> } + wrote?: Array<{ path: string; created: boolean; size: number; lines: number | null; + added?: number; removed?: number; diff?: string }> } | DagEntry | { t: 'end' } ) @@ -153,11 +154,13 @@ const GTM_FILE_EVENTS: ScriptEvent[] = [ { t:'t-', d:110, id:11, ok:true, r:'', ms:110, removed:[{ path:'~/work/raven/research/gtm-notes.md', before: GTM_SUPERSEDED }] }, /* A command the runtime has no result to read: what it left on disk is known - only from listing the directory before and after it, which is where a - created file's line count comes from and why the rewritten one has none. */ + only from listing the directory before and after it. The created file is + measured against nothing and carries its diff; the log is one the shadow + repo excludes, so its rewrite has no earlier copy and no measure. */ { t:'t+', d:150, id:12, n:'exec', a:'python3 scripts/tally.py > research/tally.txt && date >> research/run.log' }, { t:'t-', d:260, id:12, ok:true, r:'', ms:260, - wrote:[{ path:'~/work/raven/research/tally.txt', created:true, size:96, lines:4 }, + wrote:[{ path:'~/work/raven/research/tally.txt', created:true, size:89, lines:4, added:4, removed:0, + diff:'--- tally.txt\n+++ tally.txt\n@@ -0,0 +1,4 @@\n+vendor leads replies\n+Clay 1240 88\n+Apollo 980 61\n+Unify 410 37' }, { path:'~/work/raven/research/run.log', created:false, size:412, lines:null }] }, ]; diff --git a/ui-web/src/rpc/generated.ts b/ui-web/src/rpc/generated.ts index 389c8c4d3..9d25cad28 100644 --- a/ui-web/src/rpc/generated.ts +++ b/ui-web/src/rpc/generated.ts @@ -246,7 +246,7 @@ export interface TranscriptFileRemoval { del: number; } /** - * One file a command left behind, found by listing its working directory. Neither a FileChange nor a FileRemoval: a command reports its output and nothing else, so what is known of the file is that it is there, how big it is, and whether it was there before. + * One file a command left behind, found by listing its working directory. Neither a FileChange nor a FileRemoval: a command reports its output and nothing else, so what is known of the file is that it is there, how big it is, and whether it was there before. What it changed from is known only when the working directory's shadow repo held a copy from just before the command; then, and for any created text file, the change rides along as counts and a unified diff. */ export interface FileWritten { /** @@ -265,6 +265,18 @@ export interface FileWritten { * Lines in a created file, when it could be counted. Null, not absent: the key is always sent, and null says the count is unknown. Too large to read, not text, or a file that already existed, whose change therefore has no number. */ lines?: number | null; + /** + * Lines the command added to the file. Absent when the change could not be measured: not text, too large, or a rewrite whose previous contents were never captured. + */ + added?: number | null; + /** + * Lines the command removed from the file. Absent exactly when added is. + */ + removed?: number | null; + /** + * Unified diff of the change, when it was measured and small enough to carry. Absent past the event's budget even when the counts are present: a partial diff reads as a smaller change than the one that happened. + */ + diff?: string | null; } /** * Why a turn's transcript stops where it does. From 29b4e642fd62d480cb5f760e75c59d238e39b6f4 Mon Sep 17 00:00:00 2001 From: arelchan <204152633+arelchan@users.noreply.github.com> Date: Tue, 29 Sep 2026 10:51:56 +0800 Subject: [PATCH 2/5] fix(agent): keep the turn commit from racing the checkpoint warm-up The warm-up runs the shadow repo's setup on its own thread, and a turn that ended before that setup did ran the same setup beside it. Two setups at once fail on the repo's config lock, and the one that lost was the turn's commit. The setup the warm-up started is now waited for, and repeated only if it did not succeed. The tests' wait for the staging thread moves from three workspace fixtures into one teardown hook in tests/conftest.py, ahead of every fixture finalizer: any test whose turn starts a warm-up can hit the same cleanup race. Co-authored-by: Claude (claude-opus-5-5) --- raven/agent/loop/checkpoint.py | 31 +++++++++++++++++-- tests/conftest.py | 14 +++++++++ tests/test_agent_loop_session_stamps.py | 5 ---- tests/test_runtime_checkpoint.py | 40 +++++++++++++++++++++---- tests/test_runtime_checkpoint_deep.py | 6 ---- 5 files changed, 77 insertions(+), 19 deletions(-) diff --git a/raven/agent/loop/checkpoint.py b/raven/agent/loop/checkpoint.py index aed623133..13319fbde 100644 --- a/raven/agent/loop/checkpoint.py +++ b/raven/agent/loop/checkpoint.py @@ -228,6 +228,7 @@ def __init__(self, workspace: Path, shadow_dir: str = ".raven/shadow.git") -> No self._commit_count = 0 self._stage_pruned = False self._warmed = False + self._initializing: concurrent.futures.Future[bool] | None = None def covers(self, path: Path | str) -> bool: """Whether ``path`` lies in the work-tree this repo snapshots.""" @@ -301,7 +302,28 @@ async def _run( return proc.returncode or 0, out, err async def _ensure_init(self) -> bool: - """Lazily initialize the shadow repo. Idempotent; returns readiness.""" + """Lazily initialize the shadow repo. Idempotent; returns readiness. + + A warm-up runs this same setup on its own thread (:meth:`warm`), and two + at once fail on the repo's config lock -- which costs the turn its + commit when the turn ends before the warm-up's setup does. So a setup + already running is waited for, and repeated only if it did not succeed. + """ + if self._ready: + return True + running = self._initializing + if running is not None and not running.done(): + waiter = asyncio.wrap_future(running) + try: + await waiter + finally: + if not waiter.done(): + waiter.cancel() + if self._ready: + return True + return await self._init_repo() + + async def _init_repo(self) -> bool: if self._ready: return True try: @@ -450,17 +472,22 @@ async def warm(self) -> None: return staging: concurrent.futures.Future[str | None] = concurrent.futures.Future() _STAGING[index] = staging + if not self._ready: + self._initializing = concurrent.futures.Future() threading.Thread(target=self._warm_up, args=(index, staging), name="raven-stage", daemon=True).start() def _warm_up(self, index: Path, staging: "concurrent.futures.Future[str | None]") -> None: # A loop of the thread's own for the repo setup: the turn's loop is not # to wait on it, and one that closes while the setup is still starting a # git process would hold its close up. + initializing = self._initializing try: - ready = asyncio.run(self._ensure_init()) + ready = asyncio.run(self._init_repo()) except Exception as exc: # noqa: BLE001 -- a warm-up never breaks anything logger.debug("checkpoint warm-up init error: {}", exc) ready = False + if initializing is not None: + initializing.set_result(ready) if not ready: staging.set_running_or_notify_cancel() staging.set_result(None) diff --git a/tests/conftest.py b/tests/conftest.py index 20040c91a..2e9e75b02 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -323,6 +323,20 @@ def _spent(started: _Clocks, now: _Clocks) -> tuple[float, float, float]: return wall, wall - cpu - queued_s, queued_s +@pytest.hookimpl(tryfirst=True) +def pytest_runtest_teardown(item: pytest.Item, nextitem: pytest.Item | None) -> None: + """Let a checkpoint warm-up the test started finish before its fixtures go. + + A turn stages its working directory into the shadow repo on a thread of its + own (``CheckpointService.warm``), and a fixture removing that directory + under a git still writing into it fails the cleanup. Ahead of the fixture + finalizers, which run in the default teardown after this one. + """ + for thread in threading.enumerate(): + if thread.name == "raven-stage": + thread.join(30) + + @pytest.hookimpl(hookwrapper=True) def pytest_runtest_protocol(item: pytest.Item, nextitem: pytest.Item | None): item.stash[_CLOCKS] = _clocks() diff --git a/tests/test_agent_loop_session_stamps.py b/tests/test_agent_loop_session_stamps.py index d7386aebd..7c557e3c0 100644 --- a/tests/test_agent_loop_session_stamps.py +++ b/tests/test_agent_loop_session_stamps.py @@ -59,11 +59,6 @@ def get_default_model(self) -> str: def workspace(): with tempfile.TemporaryDirectory() as td: yield Path(td) - # A turn's warm-up stages the tree on a thread of its own; removing the - # directory under a git still writing into it fails the cleanup. - for thread in threading.enumerate(): - if thread.name == "raven-stage": - thread.join(30) def _make_agent(workspace: Path) -> AgentLoop: diff --git a/tests/test_runtime_checkpoint.py b/tests/test_runtime_checkpoint.py index 712c8ff33..989177d33 100644 --- a/tests/test_runtime_checkpoint.py +++ b/tests/test_runtime_checkpoint.py @@ -14,7 +14,6 @@ import asyncio import subprocess import tempfile -import threading import time from pathlib import Path @@ -32,11 +31,6 @@ def workspace(): with tempfile.TemporaryDirectory() as td: yield Path(td) - # A turn's warm-up stages the tree on a thread of its own; removing the - # directory under a git still writing into it fails the cleanup. - for thread in threading.enumerate(): - if thread.name == "raven-stage": - thread.join(30) async def _run_turn_body(agent: AgentLoop, workspace: Path): @@ -305,6 +299,40 @@ async def _slow_init(self) -> bool: assert await asyncio.wrap_future(cp_module._STAGING[svc._stage_path()]) is not None +async def test_a_turn_that_ends_during_the_warm_up_still_commits(workspace, monkeypatch): + """A turn can end before its warm-up has set the repo up. Two setups at once + fail on the config lock, and the one that lost was the turn's commit -- so + the commit waits for the warm-up's setup instead of racing it.""" + import threading + + import raven.agent.loop.checkpoint as cp_module + + real = CheckpointService._init_repo + active = [0] + overlap = [0] + + async def _tracked(self) -> bool: + active[0] += 1 + overlap[0] = max(overlap[0], active[0]) + try: + if threading.current_thread().name == "raven-stage": + await asyncio.sleep(0.3) + return await real(self) + finally: + active[0] -= 1 + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(CheckpointService, "_init_repo", _tracked) + (workspace / "a.py").write_text("print(1)\n", encoding="utf-8") + svc = CheckpointService(workspace) + + await svc.warm() + cid, changed = await svc.commit_turn("turn 1") + + assert cid is not None and changed == ["a.py"] + assert overlap[0] == 1, "the commit's setup ran beside the warm-up's" + + async def test_a_directory_is_warmed_once_and_not_every_turn(workspace, monkeypatch): """Every turn starts with a warm-up call, and only the first does any work: after it each command's own staging keeps the index warm, so another would diff --git a/tests/test_runtime_checkpoint_deep.py b/tests/test_runtime_checkpoint_deep.py index decea32e3..c1150a9aa 100644 --- a/tests/test_runtime_checkpoint_deep.py +++ b/tests/test_runtime_checkpoint_deep.py @@ -17,7 +17,6 @@ import os import shutil import tempfile -import threading from pathlib import Path import pytest @@ -35,11 +34,6 @@ def workspace(): with tempfile.TemporaryDirectory() as td: yield Path(td) - # A turn's warm-up stages the tree on a thread of its own; removing the - # directory under a git still writing into it fails the cleanup. - for thread in threading.enumerate(): - if thread.name == "raven-stage": - thread.join(30) async def _run_turn_body(agent: AgentLoop, workspace: Path): From 5150f0a23f8cc7e895666e9db40e81ecda312cd9 Mon Sep 17 00:00:00 2001 From: arelchan <204152633+arelchan@users.noreply.github.com> Date: Tue, 29 Sep 2026 15:11:14 +0800 Subject: [PATCH 3/5] fix(*): let a held-back command's retry reuse the staging it waited on A command waited out any staging already running, then always started a fresh one inside the same budget. A staging near the budget refused most commands, and one past it refused every command for good: each retry threw away the staging the last one had waited for and started another equally slow one. The latest staging is now reused whenever it started after the last write into the directory. note_write marks one after every tool call and at the start of every turn, so a file the user saved between two messages is never shown as a command's change. A retry waits for the staging the refused call left running, and a warm-up is the first command's baseline, so each command costs one staging and a staging of any length is waited out by retries. A held-back command is now logged. A staging's future is marked running as it is created. The warm-up sets the repo up before it stages, and a waiter that timed out in that window cancelled it, which raised InvalidStateError on its thread and handed every later reuser a CancelledError. Review follow-ups: the page takes the runtime's added/removed counts over ones re-read from the patch, and fromUnified skips only a real ---/+++ header pair, not every line that starts that way (a removed "-- note" line was dropped, row and count both). The dead _stage_index is gone, with its rationale moved to _stage_path, and the _GIT_TIMEOUT_SECONDS comment no longer claims every git call. Co-authored-by: Claude (claude-opus-5-5) --- raven/agent/loop/_shared.py | 4 +- raven/agent/loop/checkpoint.py | 114 ++++++++------ raven/agent/loop/turn_path.py | 12 +- tests/test_agent_loop_session_stamps.py | 72 +++++++++ tests/test_runtime_checkpoint.py | 149 ++++++++++++++++++- ui-web/src/features/workspace/record.test.ts | 13 ++ ui-web/src/features/workspace/record.ts | 19 +-- ui-web/src/lib/hunks.test.ts | 13 ++ ui-web/src/lib/hunks.ts | 14 +- 9 files changed, 345 insertions(+), 65 deletions(-) diff --git a/raven/agent/loop/_shared.py b/raven/agent/loop/_shared.py index 8d9dd54ff..18721ea47 100644 --- a/raven/agent/loop/_shared.py +++ b/raven/agent/loop/_shared.py @@ -584,8 +584,8 @@ def _file_removed_payload(removals: Any) -> list[dict[str, Any]] | None: #: unmeasured: the staging carries on, and the same command a moment later is. _EXEC_NOT_STAGED_REPLY = ( "Error: the command was not run. Raven is still snapshotting the working directory, " - "which it needs to record what this command changes (the first snapshot of a large " - "directory takes a few seconds). Run the same command again." + "which it needs to record what this command changes. Run the same command again: " + "it waits for the snapshot already under way rather than starting another." ) diff --git a/raven/agent/loop/checkpoint.py b/raven/agent/loop/checkpoint.py index 13319fbde..1a17fa793 100644 --- a/raven/agent/loop/checkpoint.py +++ b/raven/agent/loop/checkpoint.py @@ -137,11 +137,12 @@ _GC_EVERY_N_COMMITS = 50 -# Upper bound on any single git subprocess. Without this, an NFS lock, a -# held ``.git/index.lock``, or a full disk could hang ``communicate()`` -# indefinitely and brick the agent loop — violating this service's -# "never break a turn" contract. Generous enough that normal cold-init -# fits comfortably; tight enough to detect a real hang within one turn. +# Upper bound on a git subprocess the turn waits behind (a staging's ``git add`` +# runs on a thread nothing waits behind and has its own, below). Without this, +# an NFS lock, a held ``.git/index.lock``, or a full disk could hang +# ``communicate()`` indefinitely and brick the agent loop — violating this +# service's "never break a turn" contract. Generous enough that normal +# cold-init fits comfortably; tight enough to detect a real hang within one turn. _GIT_TIMEOUT_SECONDS = 30.0 # How long a command waits for the tree staged in front of it. Staging a tree @@ -166,6 +167,12 @@ # The staging running on each index, across every service in this process: two # services for one directory share the index file, so they share its one run. _STAGING: dict[Path, "concurrent.futures.Future[str | None]"] = {} +# When that staging started, and when anything was last known to write into the +# directory an index covers (``note_write``), on the same monotonic clock. A +# staging that started after the last write holds what a command about to run +# would change, whether it is still running or long done. +_STAGE_STARTED: dict[Path, float] = {} +_WRITTEN_AT: dict[Path, float] = {} _STAGE_LOCKS: dict[Path, threading.Lock] = {} @@ -420,37 +427,53 @@ async def stage_tree(self) -> str | None: tree is only read back within the same call. ``None`` when the tree cannot be staged at all (git failed), which costs - the call its diff and nothing else. :class:`StagingTimeoutError` when it is - still being staged after :data:`_STAGE_WAIT_SECONDS`: the caller is to - hold the command back rather than run it unmeasured, because a retry a - moment later finds it done. A staging already running -- the warm-up - :meth:`warm` started, or an earlier call's -- is waited for inside the - same budget and then followed by a fresh one, which is then a stat walk. - One that runs out of budget is left running rather than killed, because - what it has hashed is what makes the next one fast. + the call its diff and nothing else. :class:`StagingTimeoutError` when it + is still being staged after :data:`_STAGE_WAIT_SECONDS`: the caller is + to hold the command back rather than run it unmeasured. The staging is + left running rather than killed, and the retry waits for that same one + -- nothing has written since it began -- so a staging of any length is + waited out by retries instead of being started over by each of them. + + The latest staging is reused whenever it started after the last write + (:meth:`note_write`): the warm-up for the turn's first command, the one + a held-back command left running for its retry. Only a staging that may + predate a write is replaced, and then the one running is waited out + first, because a warm-up runs the repo setup on its thread and two + ``git config`` writes at once fail on the config lock. """ index = self._stage_path() deadline = time.monotonic() + _STAGE_WAIT_SECONDS - # Started before this call, so its tree may predate what the calls in - # between wrote: waited out, never used. Waited for before the repo - # setup below, because a warm-up runs that same setup on its thread and - # two ``git config`` writes at once fail on the config lock. - earlier = _STAGING.get(index) - if earlier is not None and not earlier.done() and not await _within(earlier, deadline): - raise StagingTimeoutError - if not await self._ensure_init(): - return None - self._prepare_index(index) - # One another caller started while this one waited is as fresh as a new - # one would be, and a second ``git add`` on the same index would only - # fail on its lock. - staging = _STAGING.get(index) - if staging is None or staging is earlier: - staging = self._start_stage(index) + staging = self._current_staging(index) + if staging is None: + earlier = _STAGING.get(index) + if earlier is not None and not earlier.done() and not await _within(earlier, deadline): + raise StagingTimeoutError + if not await self._ensure_init(): + return None + self._prepare_index(index) + # One another caller started while this one waited is as fresh as + # a new one would be, and a second ``git add`` on the same index + # would only queue behind it. + staging = self._current_staging(index) or self._start_stage(index) if not await _within(staging, deadline): raise StagingTimeoutError return staging.result() + def note_write(self) -> None: + """Say that something may just have written into the work-tree. + + Called after every tool call and at the start of every turn -- a file + the user saved between two messages is as much a write as a command's. + A staging older than the last write is no command's baseline. + """ + _WRITTEN_AT[self._stage_path()] = time.monotonic() + + def _current_staging(self, index: Path) -> "concurrent.futures.Future[str | None] | None": + staging = _STAGING.get(index) + if staging is None or _STAGE_STARTED.get(index, 0.0) < _WRITTEN_AT.get(index, 0.0): + return None + return staging + async def warm(self) -> None: """Start a staging in the background, so the first command finds the index warm. @@ -470,8 +493,7 @@ async def warm(self) -> None: running = _STAGING.get(index) if running is not None and not running.done(): return - staging: concurrent.futures.Future[str | None] = concurrent.futures.Future() - _STAGING[index] = staging + staging = self._register_stage(index) if not self._ready: self._initializing = concurrent.futures.Future() threading.Thread(target=self._warm_up, args=(index, staging), name="raven-stage", daemon=True).start() @@ -489,16 +511,26 @@ def _warm_up(self, index: Path, staging: "concurrent.futures.Future[str | None]" if initializing is not None: initializing.set_result(ready) if not ready: - staging.set_running_or_notify_cancel() staging.set_result(None) return self._prepare_index(index) self._stage(index, staging) def _start_stage(self, index: Path) -> "concurrent.futures.Future[str | None]": + staging = self._register_stage(index) + threading.Thread(target=self._stage, args=(index, staging), name="raven-stage", daemon=True).start() + return staging + + @staticmethod + def _register_stage(index: Path) -> "concurrent.futures.Future[str | None]": staging: concurrent.futures.Future[str | None] = concurrent.futures.Future() + # Running from the start, so a waiter that gives up and cancels its + # wrapper cannot cancel the staging itself: the warm-up sets its repo up + # before it stages, and a staging cancelled in that window lost its + # result and handed every later reuser a CancelledError. + staging.set_running_or_notify_cancel() _STAGING[index] = staging - threading.Thread(target=self._stage, args=(index, staging), name="raven-stage", daemon=True).start() + _STAGE_STARTED[index] = time.monotonic() return staging def _stage(self, index: Path, result: "concurrent.futures.Future[str | None]") -> None: @@ -511,7 +543,6 @@ def _stage(self, index: Path, result: "concurrent.futures.Future[str | None]") - ``write-tree`` writes the index back too and a second staging's ``add`` beside it fails on the index lock. """ - result.set_running_or_notify_cancel() with _STAGE_LOCKS.setdefault(index, threading.Lock()): if self._stage_step(("add", "-A"), index, timeout=_STAGE_ADD_TIMEOUT_SECONDS) is None: result.set_result(None) @@ -597,19 +628,14 @@ async def read_blobs(self, tree: str, paths: Collection[str], *, max_bytes: int) at = end + 1 + size + 1 return found - def _stage_index(self) -> Path: - """This process's staging index, seeded from the shared one. + def _stage_path(self) -> Path: + """This process's staging index. Per process because two processes can run commands in one directory and - an index is a single file. Seeded so a first staging finds git's stat - cache already warm from the last turn's commit instead of hashing the - whole tree again. + an index is a single file. Seeded from the shared one (``_prepare_index``) + so a first staging finds git's stat cache already warm from the last + turn's commit instead of hashing the whole tree again. """ - index = self._stage_path() - self._prepare_index(index) - return index - - def _stage_path(self) -> Path: return self._git_dir / f"exec-{os.getpid()}.index" def _prepare_index(self, index: Path) -> None: diff --git a/raven/agent/loop/turn_path.py b/raven/agent/loop/turn_path.py index a8a4c8456..4448e512f 100644 --- a/raven/agent/loop/turn_path.py +++ b/raven/agent/loop/turn_path.py @@ -745,8 +745,10 @@ async def _run_agent_loop( # noqa: C901 (cc 100: pre-existing, above the ceilin removal_watch = RemovalWatch() # Behind the model's first reply rather than in front of the first # command, and only on a directory's first turn: see CheckpointService.warm. - if self.tools.get("exec") is not None and (repo := self._turn_checkpoint()) is not None: - await repo.warm() + if (repo := self._turn_checkpoint()) is not None: + repo.note_write() + if self.tools.get("exec") is not None: + await repo.warm() # Empty-response recovery state, local to the turn — the AgentLoop is a # long-lived singleton shared across sessions, so per-instance counters # would leak across turns; resetting here gives clean per-turn budgets. @@ -1387,12 +1389,18 @@ def _hook_rollback(decision) -> bool: except StagingTimeoutError: held_back = True exec_before = None + logger.info("exec held back: {} is still being snapshotted", exec_root) if held_back: result = _EXEC_NOT_STAGED_REPLY else: result = await self.tools.execute( tool_call.name, tool_call.arguments, run_meta=tool_call.run_meta ) + # Any tool may have written into the working + # directory, so no staging from before it is a later + # command's baseline. + if (repo := self._turn_checkpoint()) is not None: + repo.note_write() duration_ms = int((time.monotonic() - tool_t0) * 1000) if exec_before is not None: exec_after = await asyncio.to_thread(workdir_snapshot.take, exec_root) diff --git a/tests/test_agent_loop_session_stamps.py b/tests/test_agent_loop_session_stamps.py index 7c557e3c0..a5c1184ba 100644 --- a/tests/test_agent_loop_session_stamps.py +++ b/tests/test_agent_loop_session_stamps.py @@ -929,6 +929,78 @@ async def chat(self, *args: Any, **kwargs: Any) -> Any: assert runs == [1], "the held-back call must not have run the command" +@pytest.mark.asyncio +async def test_a_file_saved_between_turns_is_not_counted_as_the_commands_change(workspace): + """A staging from an earlier turn is reused only while nothing has written + since; a new turn counts as a write, because the user may have saved files + in between. Reused across the gap, the command would be shown the user's + edit as its own.""" + work = workspace / "work" + work.mkdir() + kept = work / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + + await _run_command_turn( + workspace, + work, + _command_script(), + _CommandTool(lambda: kept.write_text("one\ntwo\n", encoding="utf-8")), + checkpoint=True, + ) + kept.write_text("one\ntwo\nsaved by the user\n", encoding="utf-8") + completes = await _run_command_turn( + workspace, + work, + _command_script(), + _CommandTool(lambda: kept.write_text("one\ntwo\nsaved by the user\ncmd\n", encoding="utf-8")), + checkpoint=True, + ) + + written = completes[0]["file_written"] + assert (written[0]["added"], written[0]["removed"]) == (1, 0), written + + +@pytest.mark.asyncio +async def test_a_warm_up_from_a_turn_that_ran_nothing_is_not_reused_after_the_user_saves(workspace): + """One loop, two turns. The first only answers, so its warm-up is the latest + staging and no tool call has marked a write since. The user then saves a + file. Without the new turn counting as a write, the second turn's command + would be measured against the warm-up and shown the user's edit as its own.""" + work = workspace / "work" + work.mkdir() + kept = work / "kept.txt" + kept.write_text("one\n", encoding="utf-8") + script = [ + LLMResponse(content="nothing to run", finish_reason="stop"), + *_command_script(), + ] + agent = AgentLoop( + provider=ScriptedProvider(script), + workspace=workspace, + model="stub", + policy=TurnPolicy(max_iterations=4, interactive=False), + tools=ToolWiring(restrict_to_workspace=True), + engine=EngineWiring(runtime_config=RuntimeConfig(checkpoint=CheckpointConfig(policy="always"))), + ) + agent.tools.register(_CommandTool(lambda: kept.write_text("one\nsaved by the user\ncmd\n", encoding="utf-8"))) + completes: list[dict[str, Any]] = [] + + async def on_tool_event(phase: str, info: dict[str, Any]) -> None: + if phase == "complete": + completes.append(info) + + with workdir.bind(work): + await agent._process_message(_make_msg("just answer"), on_tool_event=on_tool_event) + for thread in threading.enumerate(): + if thread.name == "raven-stage": + thread.join(30) + kept.write_text("one\nsaved by the user\n", encoding="utf-8") + await agent._process_message(_make_msg("run it"), on_tool_event=on_tool_event) + + written = completes[0]["file_written"] + assert (written[0]["added"], written[0]["removed"]) == (1, 0), written + + @pytest.mark.asyncio async def test_a_command_outside_the_shadow_repo_is_reported_without_a_diff(workspace): """``working_dir`` can point anywhere, and the turn's shadow repo only holds diff --git a/tests/test_runtime_checkpoint.py b/tests/test_runtime_checkpoint.py index 989177d33..77d44c8c4 100644 --- a/tests/test_runtime_checkpoint.py +++ b/tests/test_runtime_checkpoint.py @@ -143,6 +143,7 @@ async def test_a_stale_staging_index_is_pruned_and_a_fresh_one_kept(workspace): os.utime(stale, (old, old)) again = CheckpointService(workspace) + again.note_write() assert await again.stage_tree() is not None assert not stale.exists() @@ -169,10 +170,12 @@ def _killed(cmd, **kwargs): raise subprocess.TimeoutExpired(cmd, 0.05) monkeypatch.setattr(cp_module.subprocess, "run", _killed) + svc.note_write() assert await svc.stage_tree() is None assert not lock.exists() monkeypatch.setattr(cp_module.subprocess, "run", real) + svc.note_write() assert await svc.stage_tree() is not None @@ -217,8 +220,9 @@ def _slow(cmd, **kwargs): async def test_a_stage_waits_out_the_warm_up_and_stages_again(workspace, monkeypatch): """The warm-up started with the turn may still be running when the first - command arrives. Its tree was taken before whatever the turn has done since, - so it is waited for and never used; the staging after it is the one read.""" + command arrives. When a tool call has written since it began, its tree may + miss that write, so it is waited for and never used; the staging after it + is the one read.""" import subprocess import raven.agent.loop.checkpoint as cp_module @@ -240,6 +244,7 @@ def _count(cmd, **kwargs): await svc.warm() (workspace / "a.txt").write_text("changed while warming\n", encoding="utf-8") + svc.note_write() tree = await svc.stage_tree() assert tree is not None @@ -362,9 +367,9 @@ def _count(cmd, **kwargs): async def test_two_commands_waiting_on_one_warm_up_share_the_staging_after_it(workspace, monkeypatch): - """Two sessions in one directory both find the warm-up running. The first to - wake starts the next staging; the second takes that one rather than starting - a third ``git add`` on the same index, which would fail on its lock.""" + """Two sessions in one directory both find a warm-up running that a write has + made stale. The first to wake starts the next staging; the second takes that + one rather than starting a third ``git add`` on the same index.""" import asyncio import subprocess @@ -386,6 +391,7 @@ def _count(cmd, **kwargs): monkeypatch.setattr(cp_module.subprocess, "run", _count) await svc.warm() + svc.note_write() first, second = await asyncio.gather(svc.stage_tree(), CheckpointService(workspace).stage_tree()) assert first is not None and second is not None @@ -416,6 +422,131 @@ def _spy(cmd, **kwargs): assert ceilings[0] > cp_module._GIT_TIMEOUT_SECONDS * 10 +async def test_retries_wait_out_a_staging_slower_than_the_budget(workspace, monkeypatch): + """A tree whose every staging takes longer than a command waits. Each retry + waits for the staging the refused call left running -- nothing has written + since it began -- instead of starting an equally slow one, so every command + runs within a retry or two however slow the tree, rather than being refused + for good.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + adds: list[object] = [] + + def _slow(cmd, **kwargs): + if "add" in cmd: + adds.append(cmd) + time.sleep(0.5) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_STARTED", {}) + monkeypatch.setattr(cp_module, "_WRITTEN_AT", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 0.3) + monkeypatch.setattr(cp_module.subprocess, "run", _slow) + svc = CheckpointService(workspace) + attempts: list[int] = [] + + for command in range(3): + for attempt in range(1, 11): + try: + tree = await svc.stage_tree() + except cp_module.StagingTimeoutError: + continue + assert tree is not None + attempts.append(attempt) + break + svc.note_write() + + assert len(attempts) == 3, "a command was refused on every retry" + assert all(n > 1 for n in attempts), "the staging was faster than the budget" + assert len(adds) == 3, "a retry must not start a staging of its own" + + +async def test_a_staging_under_the_budget_is_never_refused(workspace, monkeypatch): + """Under the budget no command is held back, even when a staging is already + running as it arrives: one the command can use is waited for inside the + budget, not waited out and then followed by a second.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + + def _slow(cmd, **kwargs): + if "add" in cmd: + time.sleep(0.8) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_STARTED", {}) + monkeypatch.setattr(cp_module, "_WRITTEN_AT", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 1.2) + monkeypatch.setattr(cp_module.subprocess, "run", _slow) + svc = CheckpointService(workspace) + + # The rest of the warm-up plus a fresh staging is past the budget; the + # warm-up alone and a fresh one alone are each inside it. + svc.note_write() + await svc.warm() + await asyncio.sleep(0.1) + for _ in range(4): + assert await svc.stage_tree() is not None + svc.note_write() + + +async def test_a_command_behind_a_stale_staging_past_the_budget_is_held_back(workspace, monkeypatch): + """A write made the running staging stale, and it has not finished within + the budget. The command is held back like any other that cannot be + measured yet -- not run unmeasured because the staging in its way was not + its own.""" + import subprocess + + import raven.agent.loop.checkpoint as cp_module + + real = subprocess.run + + def _slow(cmd, **kwargs): + if "add" in cmd: + time.sleep(0.5) + return real(cmd, **kwargs) + + monkeypatch.setattr(cp_module, "_STAGING", {}) + monkeypatch.setattr(cp_module, "_STAGE_WAIT_SECONDS", 0.1) + monkeypatch.setattr(cp_module.subprocess, "run", _slow) + svc = CheckpointService(workspace) + + await svc.warm() + warm_up = cp_module._STAGING[svc._stage_path()] + svc.note_write() + with pytest.raises(cp_module.StagingTimeoutError): + await svc.stage_tree() + + # Giving up on it must not have cancelled it: a later command reuses it. + assert await asyncio.wrap_future(warm_up) is not None + assert not warm_up.cancelled() + + +async def test_a_staging_from_before_a_write_is_not_reused(workspace): + """Reuse is what lets a retry succeed, and what must not hand a command a + tree from before a write: the old text of a file a tool changed since then + is not what the command changed.""" + svc = CheckpointService(workspace) + kept = workspace / "a.txt" + kept.write_text("one\n", encoding="utf-8") + first = await svc.stage_tree() + assert await svc.stage_tree() == first, "nothing written: the same tree" + + kept.write_text("two\n", encoding="utf-8") + svc.note_write() + second = await svc.stage_tree() + + assert second != first + assert await svc.read_blobs(second, [str(kept)], max_bytes=64) == {str(kept): b"two\n"} + + async def test_a_first_stage_starts_from_the_last_turns_index(workspace, monkeypatch): """The turn's commit has already hashed the tree into the shared index, and a staging index copied from it only has to stat what changed since. Started @@ -470,8 +601,14 @@ async def one_call() -> bool: return True return False + # On a thread of its own: asyncio.run on the main thread clears the default + # loop, and pytest-asyncio then makes one it never closes. + outcome: list[bool] = [] started = time.monotonic() - assert asyncio.run(one_call()) is True + runner = threading.Thread(target=lambda: outcome.append(asyncio.run(one_call()))) + runner.start() + runner.join(10) + assert outcome == [True] assert time.monotonic() - started < 5 release.set() diff --git a/ui-web/src/features/workspace/record.test.ts b/ui-web/src/features/workspace/record.test.ts index 1e6fff0c2..d94aba275 100644 --- a/ui-web/src/features/workspace/record.test.ts +++ b/ui-web/src/features/workspace/record.test.ts @@ -506,6 +506,19 @@ describe('recording the files a command left behind', () => { /* Past the event's budget the diff is dropped and the counts still arrive: the row says how big the change was even with no patch to open. */ + /* The runtime measured the change; the patch is for drawing. Where the two + disagree the runtime's numbers are the ones shown. */ + it('shows the runtime\'s counts over ones re-read from the patch', () => { + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ + path: '/w/log.json', created: false, size: 12, lines: null, added: 5, removed: 4, + diff: '--- log.json\n+++ log.json\n@@ -1,2 +1,2 @@\n a\n-b\n+B', + }]) + + const row = rowFor('/w/log.json') + expect(row?.hunks).toHaveLength(1) + expect([row?.add, row?.del]).toEqual([5, 4]) + }) + it('carries the counts of a measured change that came without its diff', () => { wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ path: '/w/log.json', created: false, size: 12, lines: null, added: 40, removed: 7 }]) diff --git a/ui-web/src/features/workspace/record.ts b/ui-web/src/features/workspace/record.ts index 00ca53870..7e4468765 100644 --- a/ui-web/src/features/workspace/record.ts +++ b/ui-web/src/features/workspace/record.ts @@ -164,24 +164,25 @@ function wsRecordWritten(w: FileWritten): void { const WS = record() const key = String(w.path) const hunk = w.diff ? hunks.fromUnified(w.diff) : null + /* The runtime's own counts when it sent them: it measured the change, and a + number re-read off the patch text is a second source for the same fact. */ + const add = w.added ?? hunk?.add ?? null + const del = w.removed ?? hunk?.del ?? null const had = WS.changes.find((x) => x.turn === WS.turn && sameFile(x.key, key)) if (had) { if (hunk && had.hunks.length && had.kind !== 'delete') { had.hunks.push(hunk) - had.add += hunk.add - had.del += hunk.del + had.add += add ?? 0 + had.del += del ?? 0 } return } const c = rowFor(key, w.created ? 'add' : 'write') c.listed = true - if (hunk) { - c.hunks.push(hunk) - c.add = hunk.add - c.del = hunk.del - } else if (w.added != null) { - c.add = w.added - c.del = w.removed == null ? 0 : w.removed + if (hunk) c.hunks.push(hunk) + if (add != null) { + c.add = add + c.del = del ?? 0 } else if (w.created) c.add = w.lines == null ? 0 : w.lines } diff --git a/ui-web/src/lib/hunks.test.ts b/ui-web/src/lib/hunks.test.ts index 0d706b656..31027229a 100644 --- a/ui-web/src/lib/hunks.test.ts +++ b/ui-web/src/lib/hunks.test.ts @@ -57,6 +57,19 @@ describe('diff hunk builders', () => { expect(hunk.rows[40]).toEqual(['gap', ['line 41']]) }) + it('keeps a removed line that reads like a file header', () => { + const hunk = fromUnified([ + '--- a/q.sql', '+++ b/q.sql', '@@ -1,2 +1,2 @@', + '--- old comment', '+++ new comment', ' select 1', + ]) + expect([hunk.add, hunk.del]).toEqual([1, 1]) + expect(hunk.rows.slice(1)).toEqual([ + ['del', '-- old comment', 1, null], + ['add', '++ new comment', null, 1], + ['ctx', 'select 1', 2, 2], + ]) + }) + it('drops file headers and numbers unified diff rows from the hunk header', () => { const hunk = fromUnified([ '--- a/file', '+++ b/file', '@@ -2,2 +2,3 @@', diff --git a/ui-web/src/lib/hunks.ts b/ui-web/src/lib/hunks.ts index f7000a68a..559c75d0b 100644 --- a/ui-web/src/lib/hunks.ts +++ b/ui-web/src/lib/hunks.ts @@ -145,8 +145,18 @@ export function fromUnified(lines: string | string[]): WsHunk { let del = 0 let oldLine: number | null = null let newLine: number | null = null - const source = typeof lines === 'string' ? lines.split('\n') : lines - source.filter((line) => !/^(---|\+\+\+)( |$)/.test(String(line))).forEach((raw) => { + const source = (typeof lines === 'string' ? lines.split('\n') : lines).map(String) + /* A file header is the `---`/`+++` pair in front of a hunk, not any line that + starts that way: a removed line reading `-- note` is written `--- note`, + and matching on the prefix alone dropped it, row and count both. */ + const header = new Set() + source.forEach((line, i) => { + const next = source[i + 2] + if (/^--- /.test(line) && /^\+\+\+ /.test(source[i + 1] ?? '') && (next === undefined || next.startsWith('@@'))) { + header.add(i).add(i + 1) + } + }) + source.filter((_, i) => !header.has(i)).forEach((raw) => { const line = String(raw) const header = line.match(/^@@ -(\d+)(?:,\d+)? \+(\d+)(?:,\d+)? @@/) if (header) { From 6e3687f68a58eb8c6abf18b73c3804c6c67fcacc Mon Sep 17 00:00:00 2001 From: arelchan <204152633+arelchan@users.noreply.github.com> Date: Tue, 29 Sep 2026 15:16:52 +0800 Subject: [PATCH 4/5] test(agent): mark the staging timing tests production_timing The four tests that slow a staging to sit either side of the wait budget wait by design, and one of them crossed the suite's 3s idle ceiling on CI. They prove a timing property, which is what production_timing exempts from the ceiling. Co-authored-by: Claude (claude-opus-5-5) --- tests/test_agent_loop_session_stamps.py | 2 ++ tests/test_runtime_checkpoint.py | 2 ++ 2 files changed, 4 insertions(+) diff --git a/tests/test_agent_loop_session_stamps.py b/tests/test_agent_loop_session_stamps.py index a5c1184ba..1d99af036 100644 --- a/tests/test_agent_loop_session_stamps.py +++ b/tests/test_agent_loop_session_stamps.py @@ -827,6 +827,7 @@ async def chat(self, *args: Any, **kwargs: Any) -> Any: @pytest.mark.asyncio +@pytest.mark.production_timing # a slow first reply against a slow first staging is the property async def test_the_first_command_in_a_cold_directory_is_measured_when_the_model_took_its_time(workspace, monkeypatch): """The first staging in a directory the shadow repo has never indexed hashes the whole tree -- seconds on a large one, far past what a command waits. It @@ -868,6 +869,7 @@ def _cold_first(cmd, **kwargs): @pytest.mark.asyncio +@pytest.mark.production_timing # a staging slower than the wait budget is the property async def test_a_command_is_held_back_until_the_snapshot_is_ready(workspace, monkeypatch): """Run unmeasured, a command in a directory still being snapshotted leaves a change nobody can show. It is not run instead: the call fails with a reply diff --git a/tests/test_runtime_checkpoint.py b/tests/test_runtime_checkpoint.py index 77d44c8c4..2ae6a3e23 100644 --- a/tests/test_runtime_checkpoint.py +++ b/tests/test_runtime_checkpoint.py @@ -422,6 +422,7 @@ def _spy(cmd, **kwargs): assert ceilings[0] > cp_module._GIT_TIMEOUT_SECONDS * 10 +@pytest.mark.production_timing # the stagings are slowed past the wait budget, which is the property async def test_retries_wait_out_a_staging_slower_than_the_budget(workspace, monkeypatch): """A tree whose every staging takes longer than a command waits. Each retry waits for the staging the refused call left running -- nothing has written @@ -465,6 +466,7 @@ def _slow(cmd, **kwargs): assert len(adds) == 3, "a retry must not start a staging of its own" +@pytest.mark.production_timing # the stagings are slowed to sit either side of the wait budget, which is the property async def test_a_staging_under_the_budget_is_never_refused(workspace, monkeypatch): """Under the budget no command is held back, even when a staging is already running as it arrives: one the command can use is waited for inside the From ed71c7ab035d74f2d23726f70f985b15d4151c3a Mon Sep 17 00:00:00 2001 From: arelchan <204152633+arelchan@users.noreply.github.com> Date: Tue, 29 Sep 2026 16:27:33 +0800 Subject: [PATCH 5/5] test: pin the runtime's counts on a merged row and restate when a command is staged The page adds a command's change to a row the turn already has using the runtime's added/removed counts, and nothing failed when that branch went back to the patch's own. A test now pins it. CONTEXT.md still said stage_tree stages the tree again before each command. Since stagings are reused, a command takes the latest one when it began after the last note_write and stages afresh only otherwise. Co-authored-by: Claude (claude-opus-5-5) --- CONTEXT.md | 9 +++++---- ui-web/src/features/workspace/record.test.ts | 19 +++++++++++++++++++ 2 files changed, 24 insertions(+), 4 deletions(-) diff --git a/CONTEXT.md b/CONTEXT.md index 479346522..02b770b3d 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -568,10 +568,11 @@ user's `.git`), so an interrupted or failed turn can be rolled back. One `Checkp per working directory, cached by `AgentLoop._turn_checkpoint()` and keyed on the directory the running turn is bound to. The same repo also answers what a file an `exec` command rewrote or removed held before it: `CheckpointService.warm` starts staging the tree when -a directory's first turn in this process starts, `stage_tree` stages it again just before -each command into a per-process index of its own (never the index the turn commit reads), -and `read_blobs` reads the old contents back for the command's `file_written` diff and -`file_removed` body. +a directory's first turn in this process starts, and a command's `stage_tree` takes the +latest staging if it began after the last `note_write` (every tool call and every turn +start) or stages the tree afresh otherwise, into a per-process index of its own (never the +index the turn commit reads). `read_blobs` reads the old contents back for the command's +`file_written` diff and `file_removed` body. _Avoid_: "shadow git" as the term — Checkpoint is the per-turn snapshot it produces. **Empty-Response Recovery** (`agent/loop/recovery.py`): diff --git a/ui-web/src/features/workspace/record.test.ts b/ui-web/src/features/workspace/record.test.ts index d94aba275..153583de2 100644 --- a/ui-web/src/features/workspace/record.test.ts +++ b/ui-web/src/features/workspace/record.test.ts @@ -519,6 +519,25 @@ describe('recording the files a command left behind', () => { expect([row?.add, row?.del]).toEqual([5, 4]) }) + /* The same preference where a second change is added to a row the turn + already has: this is where a command writing one file several times sums + its counts, and where a wrong number is least likely to be noticed. */ + it('adds the runtime\'s counts, not the patch\'s, to a row it already has', () => { + const args = { path: '/w/notes.md', content: 'a\n' } + wsOnTool('write_file', args) + wsOnToolDone('write_file', args, true, '', null, + '--- a/w/notes.md\n+++ b/w/notes.md\n@@ -0,0 +1,1 @@\n+a', { path: '/w/notes.md', after: 'a\n' }) + + wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ + path: '/w/notes.md', created: false, size: 4, lines: null, added: 5, removed: 3, + diff: '--- notes.md\n+++ notes.md\n@@ -1,1 +1,2 @@\n a\n+b', + }]) + + const row = rowFor('/w/notes.md') + expect(row?.hunks).toHaveLength(2) + expect([row?.add, row?.del]).toEqual([1 + 5, 0 + 3]) + }) + it('carries the counts of a measured change that came without its diff', () => { wsOnToolDone('exec', { command: 'python3 gen.py' }, true, '', null, undefined, undefined, undefined, [{ path: '/w/log.json', created: false, size: 12, lines: null, added: 40, removed: 7 }])