ive harnessed the harness
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736"""A small, root-bound filesystem and process vocabulary for workbench cells.
This is a convenience layer over the host-provided source checkout, not a sandbox.Every relative path is checked beneath :attr:`Workspace.root`, including pathsthat pass through symlinks. For example::
from klbr import open_file, run
await (await open_file("README.md")).replace("old heading", "new heading") result = await run(("git", "status", "--short")) print(result.stdout.decode())
Every cell already has these bound under those names."""from __future__ import annotations
import reimport shutilimport zlib
import asyncioimport base64import mathimport osimport fnmatchfrom collections.abc import Mapping, Sequencefrom dataclasses import dataclassfrom pathlib import Path
from . import runtimefrom klbr_runtime import workspace_trace as trace
@dataclass(frozen=True, slots=True)class Entry: """One thing below a directory: what it is, not just where it is.
``is_dir`` and ``size`` are read here because a caller that has to stat to tell a file from a directory ends up writing a walker, and a walker that guesses wrong calls `ls` on files. ``size`` is ``None`` for a directory: a directory's size is not one number. """
path: Path relative: str is_dir: bool size: int | None
@property def name(self) -> str: """The last path component.""" return self.path.name
@dataclass(frozen=True, slots=True)class Listing: """Entries found under a path, and what the walk deliberately did not enter.
Iterable, so ``[e for e in await ls(...)]`` reads as a sequence of entries. ``pruned`` names the directories skipped as build or version-control noise, and ``truncated`` says a cap stopped the walk early - a short listing and a listing that stopped looking are different answers, and this one cannot be mistaken for the other. """
root: Path path: Path entries: tuple[Entry, ...] pruned: tuple[str, ...] = () truncated: bool = False
def __iter__(self): return iter(self.entries)
def __len__(self) -> int: return len(self.entries)
ANCHOR = re.compile( r"^(?P<start>[1-9]\d*):(?P<start_hash>[0-9a-fA-F]{2})" r"(?:\.\.(?P<end>[1-9]\d*):(?P<end_hash>[0-9a-fA-F]{2}))?$")
def line_hash(text: str) -> str: """Two hex characters of a line's content, whitespace-stripped: CRC-32, low byte.
The fingerprint is of content, not position, so reformatting a line leaves its anchor alone while changing it does not. This is the rule the rest of the harness uses, so an anchor means the same thing in every harness that prints one.
The low byte of a CRC is affine in the input, so two characters collide either never or always for a given kind of change - correlated across lines, where a mixing hash would fail one line at a time. That is the price of a checksum cheap enough to run over every line of every file shown; the undetected patterns need four or more correlated bit flips to reach it. """ return f"{zlib.crc32(text.strip().encode('utf-8')) % 256:02x}"
@dataclass(frozen=True, slots=True)class File: """One file in the session's tree, opened once and acted on through the handle.
Opening is not reading: it resolves the path and nothing else, so a handle for a file that does not exist yet is legitimate - ``write`` creates it. Every method sees the file as it is when the call is made, and the hashline anchors ``edit`` takes come from the same object the anchors were printed from, so an edit on a file that was never opened cannot be written down. """
path: Path relative: str
async def text(self, encoding: str = "utf-8") -> str: """The file as text, decoded with ``encoding``.""" result = await asyncio.to_thread(self.path.read_text, encoding=encoding) trace.activity("read", self.relative) return result
async def read(self) -> bytes: """The file as bytes, with no decoding.""" result = await asyncio.to_thread(self.path.read_bytes) trace.activity("read", self.relative) return result
async def exists(self) -> bool: """Whether anything is at this path right now.
Opening does not answer this, on purpose: a handle is a subject, and a subject can be created. Ask here when the difference matters. """ return await asyncio.to_thread(self.path.exists)
async def hashlines(self) -> str: """The file with an ``N:hash|`` prefix per line, for anchoring edits.""" text = await self.text() return "\n".join(f"{number}:{line_hash(line)}|{line}" for number, line in enumerate(text.splitlines(), 1))
async def write(self, content: str | bytes, *, encoding: str = "utf-8") -> None: """Replace the whole file, creating parent directories. ``str`` or ``bytes``.""" if not isinstance(content, (str, bytes)): raise TypeError("workspace content must be str or bytes")
def write() -> None: with trace.edit(self.path, self.relative): self.path.parent.mkdir(parents=True, exist_ok=True) if isinstance(content, str): self.path.write_text(content, encoding=encoding) else: self.path.write_bytes(content)
await asyncio.to_thread(write)
async def replace(self, old: str, new: str, *, encoding: str = "utf-8") -> None: """Replace exactly one literal occurrence, or raise without writing.
Both arguments are required, and both must be non-empty: a caller that left ``new`` out means something, and guessing which is how text disappears. """ if not old: raise ValueError("replacement text must not be empty")
def replace() -> None: with trace.edit(self.path, self.relative): text = self.path.read_text(encoding=encoding) count = text.count(old) if count != 1: raise ValueError(f"expected one occurrence in {self.path}, found {count}") self.path.write_text(text.replace(old, new, 1), encoding=encoding)
await asyncio.to_thread(replace)
async def edit( self, anchor: str, new: str | None = None, *, after: bool = False, encoding: str = "utf-8" ) -> None: """Edit by line anchor from ``hashlines()``, or raise without writing.
``"12:f1"`` names one line, ``"12:f1..14:9c"`` a range, ``after=True`` inserts after the anchored line instead of replacing it, and leaving ``new`` out deletes. Both ends of a range are verified against the file as it is now, so an edit aimed at a file that moved under the model is refused rather than applied somewhere else. """ if after and new is None: raise ValueError("after=True has nothing to insert: pass new=... or drop after") parsed = ANCHOR.match(anchor.strip()) if parsed is None: raise ValueError( f'anchor {anchor!r} is not "N:hash" or "N:hash..M:hash"; copy one from hashlines()' ) start = int(parsed["start"]) end = int(parsed["end"]) if parsed["end"] else start if after and end != start: raise ValueError("after=True anchors one line; a range has no after") if end < start: raise ValueError(f"anchor range {start}..{end} runs backwards")
def apply() -> None: with trace.edit(self.path, self.relative): text = self.path.read_text(encoding=encoding) lines = text.splitlines() for number, expected in ((start, parsed["start_hash"]), (end, parsed["end_hash"] or "")): if number > len(lines): raise ValueError( f"line {number} is past the end of {self.path} ({len(lines)} lines)" ) actual = line_hash(lines[number - 1]) if expected and actual != expected.lower(): raise ValueError( f"line {number} now hashes {actual}, not {expected.lower()}: " f"{self.path} changed since it was read; read it again" ) body = [] if new is None else new.splitlines() if after: lines[start:start] = body else: lines[start - 1 : end] = body newline = "\r\n" if "\r\n" in text else "\n" joined = newline.join(lines) self.path.write_text( joined + (newline if text.endswith(("\n", "\r\n")) else ""), encoding=encoding )
await asyncio.to_thread(apply)
def __str__(self) -> str: return f"File({self.relative})"
@dataclass(frozen=True, slots=True)class TextMatch: """One literal text match, with a one-based line and column."""
path: Path relative: str line: int column: int text: str
@dataclass(frozen=True, slots=True)class SearchResult: """Matches, and whether the search stopped before the tree ran out.
Iterable over its matches so an ordinary lookup reads as a sequence. ``truncated`` is the part that keeps a partial answer from being read as the whole answer: a search that hit its cap found *these*, not *all*. """
matches: tuple[TextMatch, ...] truncated: bool files_searched: int files_unreadable: int pruned: tuple[str, ...]
def __iter__(self): return iter(self.matches)
def __len__(self) -> int: return len(self.matches)
def __bool__(self) -> bool: return bool(self.matches)
@dataclass(frozen=True, slots=True)class RunResult: """The complete, bounded outcome of :meth:`Workspace.run`.
``completed`` means the child exited without timeout or cancellation; it does not imply a zero exit code. ``stdout`` and ``stderr`` retain at most ``max_output_bytes`` together. Their truncation flags say when further output was discarded while the pipes were still drained. """
argv: tuple[str, ...] cwd: Path exit_code: int | None stdout: bytes stderr: bytes duration: float completed: bool timed_out: bool cancelled: bool stdout_truncated: bool = False stderr_truncated: bool = False
@property def ok(self) -> bool: """Whether the command completed successfully with exit code zero.""" return self.completed and self.exit_code == 0
@dataclass(frozen=True, slots=True)class Workspace: """A path vocabulary rooted at one host-supplied source checkout.
Paths returned by this class are absolute resolved paths. Pass relative paths to keep reads, edits, and command working directories within ``root``. """
root: Path
def __post_init__(self) -> None: root = Path(self.root).resolve(strict=True) if not root.is_dir(): raise NotADirectoryError(root) object.__setattr__(self, "root", root)
def path(self, path: str | Path = ".") -> Path: """Resolve a relative path below this workspace, rejecting escapes.""" candidate = Path(path) if candidate.is_absolute(): raise ValueError(f"workspace paths must be relative: {path!r}") resolved = (self.root / candidate).resolve(strict=False) try: resolved.relative_to(self.root) except ValueError as error: raise ValueError(f"workspace path escapes source root: {path!r}") from error return resolved
async def open_file(self, path: str | Path) -> File: """Resolve one path under the root into a handle. Not a read, and not a stat.
A handle for a file that does not exist yet is legitimate: ``write`` creates it. What is refused here is escaping the tree or naming a directory. """ target = self.path(path) if target.is_dir(): raise IsADirectoryError(target) return File(target, _relative(self.root, target))
async def ls( self, path: str | Path = ".", *, recursive: bool = False, match: str | None = None, include_pruned: bool = False, max_entries: int = 20_000, ) -> Listing: """List entries: one directory by default, a whole tree with ``recursive=True``.
One level is the default because that is what ``ls`` means, and because a walk that surprises a caller is a walk that reads build output for seconds. ``match`` is a shell glob on the name (``"*.rs"``) and applies to the walk, not to a one-level listing. Build and version-control directories are not entered while recursing unless ``include_pruned`` is set; which were skipped is in ``pruned``, and a walk that reached ``max_entries`` says ``truncated`` rather than looking complete. A file lists itself. """ start = self.path(path) if max_entries < 1: raise ValueError("max_entries must be positive")
def one_level() -> Listing: if start.is_file(): return Listing(self.root, start.parent, (entry_of(self.root, start),)) if not start.is_dir(): raise FileNotFoundError(start) entries = tuple(entry_of(self.root, item) for item in sorted(start.iterdir())) return Listing(self.root, start, entries)
def tree() -> Listing: if start.is_file(): return Listing(self.root, start.parent, (entry_of(self.root, start),)) if not start.is_dir(): raise FileNotFoundError(start) found: list[Entry] = [] pruned: set[str] = set() truncated = False for parent, directories, files in os.walk(start, followlinks=False): parent_path = Path(parent) kept: list[str] = [] for name in directories: item = parent_path / name if not _within(self.root, item.resolve(strict=False)): continue if name in PRUNED_DIRECTORY_NAMES and not include_pruned: pruned.add(name) continue kept.append(name) directories[:] = kept for name in (*directories, *files): item = (parent_path / name).resolve(strict=False) if not _within(self.root, item): continue if match is not None and not fnmatch.fnmatch(item.name, match): continue found.append(entry_of(self.root, item)) if len(found) >= max_entries: truncated = True break if truncated: break entries = tuple(sorted(found, key=lambda entry: entry.relative)) return Listing(self.root, start, entries, tuple(sorted(pruned)), truncated)
result = await asyncio.to_thread(tree if recursive else one_level) trace.activity("list", str(path)) return result
async def search( self, needle: str, path: str | Path = ".", *, encoding: str = "utf-8", max_results: int = 100, include_pruned: bool = False, match: str | None = None, ) -> SearchResult: """Find literal text below ``path``, and say whether the search ran out of room.
Literal on purpose: a model looking for exact source text should not have to quote regular-expression syntax to find it, and a literal cannot silently match something else. ``match`` narrows the files by name glob (``"*.rs"``).
The result is iterable over its matches. ``truncated`` is true when the cap stopped the walk before the tree did - a hit list that stopped early is a different answer from a short one, and the caller has to be able to tell them apart. """ if not needle: raise ValueError("search needle must not be empty") if max_results < 1: raise ValueError("max_results must be positive") listing = await self.ls( path, recursive=True, match=match, include_pruned=include_pruned ) files = [entry for entry in listing.entries if not entry.is_dir]
def find() -> SearchResult: matches: list[TextMatch] = [] unreadable = 0 searched = 0 for entry in files: try: text = entry.path.read_text(encoding=encoding) except (OSError, UnicodeError): unreadable += 1 continue searched += 1 for line_number, line in enumerate(text.splitlines(), 1): offset = 0 while (position := line.find(needle, offset)) >= 0: matches.append( TextMatch( entry.path, entry.relative, line_number, position + 1, line, ) ) if len(matches) >= max_results: return SearchResult( tuple(matches), True, searched, unreadable, listing.pruned ) offset = position + len(needle) return SearchResult(tuple(matches), False, searched, unreadable, listing.pruned)
result = await asyncio.to_thread(find) trace.activity("search", f"{needle} in {path}") return result
async def delete( self, path: str | Path, *, force: bool = False, recursive: bool = False ) -> None: """Delete one file, or a directory with ``recursive=True``.
``force`` is narrow, and narrower than the shell's: it means a path that is already gone is not an error. It does not override the refusals - the root of the tree is never deletable, symlinks are refused rather than deleting whatever they resolve to, and a directory without ``recursive`` must be empty, with the refusal naming how many entries are in the way. """ target = self.path(path) if (self.root / Path(path)).is_symlink(): raise ValueError("workspace.delete refuses symlinks") if target == self.root: raise ValueError("workspace.delete refuses the root of the tree") if not target.is_dir(): await asyncio.to_thread(target.unlink, missing_ok=force) return entries = await asyncio.to_thread(lambda: sorted(target.iterdir())) if entries and not recursive: count = f"{len(entries)} " + ("entry" if len(entries) == 1 else "entries") raise ValueError( f"{path} is a directory with {count}; " f"pass recursive=True to delete it and everything in it" ) if recursive: await asyncio.to_thread(shutil.rmtree, target, True) else: await asyncio.to_thread(target.rmdir)
async def run( self, argv: Sequence[str | os.PathLike[str]], *, cwd: str | Path = ".", env: Mapping[str, str] | None = None, stdin: bytes | str | None = None, timeout: float | None = None, max_output_bytes: int = 1_048_576, ) -> RunResult: """Ask the host to run argv under its immutable source root.
The host owns the process group, drains its bounded pipes, and kills it if this request times out or its workbench generation is discarded. """ command = tuple(os.fspath(part) for part in argv) if not command or any(not isinstance(part, str) or not part for part in command): raise ValueError("argv must contain non-empty strings") if timeout is not None and (not math.isfinite(timeout) or timeout <= 0): raise ValueError("timeout must be positive and finite") if max_output_bytes < 0: raise ValueError("max_output_bytes must not be negative") input_bytes = stdin.encode() if isinstance(stdin, str) else stdin if input_bytes is not None and not isinstance(input_bytes, bytes): raise TypeError("stdin must be str, bytes, or None") working_directory = self.path(cwd) if not working_directory.is_dir(): raise NotADirectoryError(working_directory) overlay = dict(env or {}) if any(not isinstance(key, str) or not isinstance(value, str) for key, value in overlay.items()): raise TypeError("environment keys and values must be strings") if any(not key or "\0" in key or "=" in key or "\0" in value for key, value in overlay.items()): raise ValueError("environment contains an invalid key or value") requested_timeout = 3_600_000 if timeout is None else math.ceil(timeout * 1000) reply = await runtime.scope().request({ "method": "workspace.run", "argv": tuple(command), "cwd": os.fspath(cwd), "env": overlay, "stdin_base64": base64.b64encode(input_bytes).decode() if input_bytes is not None else None, "timeout_ms": requested_timeout, # The frame has to fit base64 output plus its JSON envelope. This # remains a total-output bound and reports truncation truthfully. "max_output_bytes": min(max_output_bytes, 750_000), }) try: return RunResult( tuple(reply["argv"]), Path(reply["cwd"]), reply["exit_code"], base64.b64decode(reply["stdout_base64"], validate=True), base64.b64decode(reply["stderr_base64"], validate=True), reply["duration_ms"] / 1000, reply["completed"], reply["timed_out"], reply["cancelled"], reply["stdout_truncated"], reply["stderr_truncated"], ) except (KeyError, TypeError, ValueError) as error: raise RuntimeError("host returned an invalid workspace.run reply") from error
# Directory names that are never the subject of "where is this text": build output, dependency# trees and version-control stores. Walking them costs seconds and can fill the match cap before# the walk reaches source, which turns a small question into a partial answer.PRUNED_DIRECTORY_NAMES = frozenset( {"target", "target-parent", "node_modules", ".git", ".jj", ".venv", "venv", "__pycache__", "dist", "build"})
def entry_of(root: Path, item: Path) -> Entry: """One entry, with what it is read here rather than left for the caller to stat.""" resolved = item.resolve(strict=False) if resolved.is_dir(): return Entry(resolved, _relative(root, resolved), True, None) try: size = resolved.stat().st_size except OSError: size = None return Entry(resolved, _relative(root, resolved), False, size)
def _relative(root: Path, path: Path) -> str: try: return path.relative_to(root).as_posix() except ValueError: return path.as_posix()
def _within(root: Path, path: Path) -> bool: try: path.relative_to(root) return True except ValueError: return False
async def current() -> Workspace: """Construct the workspace from the host's authoritative ``source_root``.
``await workspace.current()`` never uses the worker's current directory as a fallback, because that would silently target the wrong checkout. """ inspected = await runtime.inspect() try: source_root = inspected["source_root"]["path"] except (KeyError, TypeError) as error: raise RuntimeError("runtime.inspect did not provide source_root.path") from error if not isinstance(source_root, str) or not source_root: raise RuntimeError("runtime.inspect provided an invalid source_root.path") return Workspace(Path(source_root))
_current: Workspace | None = None
async def here() -> Workspace: """The workspace this kernel generation is already working in, resolved once.
A generation is one worker process holding one session's target, so the resolved root is valid for the life of the process; a target that changed would have reset the kernel. This is what the module-level verbs below act on, and it is why they cost no host round trip after the first call. """ global _current if _current is None: _current = await current() return _current
# The same operations as the methods above, pointed at the workspace this cell is already in.# A model reaches for `klbr.run` before it constructs an object, and making it fetch the obvious# workspace first is how a cell ends up shelling out through the kernel instead. A cell that# means a *different* tree uses `(await workspace.current())` or a Workspace it holds.async def open_file(path: str | Path) -> File: """Open one file in the current workspace: a handle to read and write through.
Opening is not reading, so this succeeds for a file that does not exist yet. """ return await (await here()).open_file(path)
async def ls( path: str | Path = ".", *, recursive: bool = False, match: str | None = None, include_pruned: bool = False, max_entries: int = 20_000,) -> Listing: """List the current workspace: one directory, or a tree with ``recursive=True``.""" return await (await here()).ls( path, recursive=recursive, match=match, include_pruned=include_pruned, max_entries=max_entries, )
async def search( needle: str, path: str | Path = ".", *, encoding: str = "utf-8", max_results: int = 100, include_pruned: bool = False, match: str | None = None,) -> SearchResult: """Find literal text below ``path`` in the current workspace; the result says if it was capped.""" return await (await here()).search( needle, path, encoding=encoding, max_results=max_results, include_pruned=include_pruned, match=match, )
async def delete( path: str | Path, *, force: bool = False, recursive: bool = False) -> None: """Delete one file, or a directory with ``recursive=True``; ``force`` forgives a missing path.""" await (await here()).delete(path, force=force, recursive=recursive)
async def run( argv: Sequence[str | os.PathLike[str]], *, cwd: str | Path = ".", env: Mapping[str, str] | None = None, stdin: bytes | str | None = None, timeout: float | None = None, max_output_bytes: int = 1048576,) -> RunResult: """Ask the host to run ``argv`` in the current workspace.""" return await (await here()).run( argv, cwd=cwd, env=env, stdin=stdin, timeout=timeout, max_output_bytes=max_output_bytes, )
# What every cell already has bound, in one place: the entry point builds the namespace from# this, so the vocabulary a cell sees and the vocabulary the package exports cannot drift.PRELUDE_NAMES = ( "delete", "ls", "open_file", "run", "search",)
def prelude() -> dict[str, object]: """The tree verbs as a namespace mapping.""" return {name: globals()[name] for name in PRELUDE_NAMES}
__all__ = [ "ANCHOR", "PRELUDE_NAMES", "PRUNED_DIRECTORY_NAMES", "Entry", "File", "Listing", "RunResult", "SearchResult", "TextMatch", "Workspace", "current", "delete", "here", "line_hash", "ls", "open_file", "prelude", "run", "search",]