diff --git a/CONTEXT.md b/CONTEXT.md index ed55a7c61..02b770b3d 100644 --- a/CONTEXT.md +++ b/CONTEXT.md @@ -566,7 +566,13 @@ _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, 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/raven/agent/loop/_shared.py b/raven/agent/loop/_shared.py index cc86794cc..18721ea47 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. Run the same command again: " + "it waits for the snapshot already under way rather than starting another." +) + + #: 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..1a17fa793 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 @@ -129,13 +137,64 @@ _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 +# 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]"] = {} +# 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] = {} + + +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 +233,16 @@ 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 + 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.""" + 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 +252,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 +264,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,14 +302,35 @@ 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.""" + """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: @@ -305,6 +415,242 @@ 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. 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 + 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. + + 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 = 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() + + 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._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_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 + _STAGE_STARTED[index] = time.monotonic() + 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. + """ + 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_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 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. + """ + 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 +672,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..4448e512f 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,12 @@ 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 (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. @@ -1301,6 +1325,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 +1382,25 @@ 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 + 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) @@ -1411,15 +1452,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/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 0b3a40c06..1d99af036 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 @@ -574,20 +577,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 +651,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 +673,358 @@ 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 +@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 + 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 +@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 + 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_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 + 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 +1313,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..2ae6a3e23 100644 --- a/tests/test_runtime_checkpoint.py +++ b/tests/test_runtime_checkpoint.py @@ -11,8 +11,10 @@ from __future__ import annotations +import asyncio import subprocess import tempfile +import time from pathlib import Path import pytest @@ -80,6 +82,539 @@ 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) + again.note_write() + 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) + 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 + + +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. 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 + + 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") + svc.note_write() + 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_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 + 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 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 + + 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() + 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 + 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 + + +@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 + 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" + + +@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 + 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 + 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 + + # 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() + 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() + + def _os_environ(): import os 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..153583de2 100644 --- a/ui-web/src/features/workspace/record.test.ts +++ b/ui-web/src/features/workspace/record.test.ts @@ -488,6 +488,136 @@ 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. */ + /* 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]) + }) + + /* 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 }]) + + 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..7e4468765 100644 --- a/ui-web/src/features/workspace/record.ts +++ b/ui-web/src/features/workspace/record.ts @@ -132,33 +132,58 @@ 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 + /* 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 += add ?? 0 + had.del += del ?? 0 + } + 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) + if (add != null) { + c.add = add + c.del = del ?? 0 + } 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/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) { 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.