Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
9 changes: 7 additions & 2 deletions src/context_system/builder.py
Original file line number Diff line number Diff line change
Expand Up @@ -101,8 +101,8 @@ def _build_workspace_section(root: Path, current: Path) -> str:
f"- Today's date: {date.today().isoformat()}",
f"- Workspace root: {workspace.workspace_root}",
f"- Current directory: {workspace.current_directory}",
f"- Python files: {workspace.python_file_count}",
f"- Test files: {workspace.test_file_count}",
f"- Python files: {_count_text(workspace.python_file_count, workspace.counts_partial)}",
f"- Test files: {_count_text(workspace.test_file_count, workspace.counts_partial)}",
]
if workspace.key_files:
lines.append(f"- Key files: {', '.join(workspace.key_files)}")
Expand All @@ -111,6 +111,11 @@ def _build_workspace_section(root: Path, current: Path) -> str:
return "\n".join(lines)


def _count_text(count: int, partial: bool) -> str:
"""``1234``, or ``1234+ (partial scan)`` when the walk hit its budget."""
return f"{count}+ (partial scan)" if partial else str(count)


def _build_git_section(cwd: str) -> str:
from .git_context import collect_git_context, format_git_status
try:
Expand Down
3 changes: 3 additions & 0 deletions src/context_system/models.py
Original file line number Diff line number Diff line change
Expand Up @@ -105,3 +105,6 @@ class WorkspaceSnapshot:
key_files: tuple[str, ...]
python_file_count: int
test_file_count: int
#: True when the file walk hit its directory budget, so the counts are
#: lower bounds — rendered as ``N+ (partial scan)``.
counts_partial: bool = False
101 changes: 95 additions & 6 deletions src/context_system/workspace_snapshot.py
Original file line number Diff line number Diff line change
@@ -1,18 +1,61 @@
"""The ``## Runtime Context`` facts about a workspace: shape, key files, counts.

The Python/test file counts come from ONE bounded ``os.scandir`` walk that
never descends into vendored or generated directories (``node_modules``,
``.git``, virtualenvs, caches, ``.clawcodex`` worktree checkouts …) and stops
after :data:`MAX_SCANNED_DIRS` directories. The previous implementation used
``Path.rglob`` twice, which walks *everything* and filters afterwards: on a
workspace with a few front-end packages (hundreds of thousands of
``node_modules`` directories) each walk took ~20 s, and the system prompt is
built at spawn and again on every resume/clear — so opening a saved session
from the web sidebar sat on this for ~45 s. The counts are a hint for the
model, not an inventory: a bounded scan that says "N+ (partial)" is the right
trade, and it keeps the cost proportional to the project rather than to
whatever a package manager left on disk.
"""

from __future__ import annotations

import os
from collections import deque
from pathlib import Path

from .models import WorkspaceSnapshot

_IGNORED_NAMES = {
#: Directory names the walk never enters. Vendored trees, VCS internals,
#: virtualenvs, tool caches, build output and the per-repo ``.clawcodex``
#: directory (its ``worktrees/`` holds whole extra checkouts of the repo,
#: which would count every file again per worktree).
_IGNORED_NAMES = frozenset({
".git",
".hg",
".svn",
".venv",
"venv",
".tox",
".nox",
".eggs",
".cache",
".pytest_cache",
".mypy_cache",
".ruff_cache",
"__pycache__",
"node_modules",
}
".clawcodex",
".claude",
"site-packages",
})

#: Upper bound on directories one snapshot visits (after pruning). Roughly a
#: few tens of milliseconds on a warm disk; past it the counts are reported as
#: partial rather than the walk running on.
MAX_SCANNED_DIRS = 2500

#: Upper bound on directory entries one snapshot reads, so a single flat
#: directory of a million files (a dataset dump next to the code) is bounded
#: the same way a deep tree is.
MAX_SCANNED_ENTRIES = 200_000

_KEY_FILE_CANDIDATES = (
"README.md",
"CLAWCODEX.md",
Expand All @@ -29,6 +72,7 @@ def build_workspace_snapshot(
*,
cwd: str | Path | None = None,
top_level_limit: int = 12,
max_dirs: int = MAX_SCANNED_DIRS,
) -> WorkspaceSnapshot:
root = Path(workspace_root).expanduser().resolve()
current = Path(cwd).expanduser().resolve() if cwd is not None else root
Expand All @@ -49,8 +93,7 @@ def build_workspace_snapshot(
break

key_files = tuple(name for name in _KEY_FILE_CANDIDATES if (root / name).exists())
python_file_count = sum(1 for path in root.rglob("*.py") if _is_countable(path))
test_file_count = sum(1 for path in root.rglob("test_*.py") if _is_countable(path))
python_file_count, test_file_count, partial = count_python_files(root, max_dirs=max_dirs)

return WorkspaceSnapshot(
workspace_root=root,
Expand All @@ -59,9 +102,56 @@ def build_workspace_snapshot(
key_files=key_files,
python_file_count=python_file_count,
test_file_count=test_file_count,
counts_partial=partial,
)


def count_python_files(root: Path, *, max_dirs: int = MAX_SCANNED_DIRS) -> tuple[int, int, bool]:
"""``(python_files, test_files, partial)`` under ``root``, pruned and bounded.

One breadth-first ``scandir`` walk: ignored directory names are skipped
*before* they are entered (this is what makes the walk cheap — filtering
``rglob`` output afterwards still pays for every directory it crawled),
symlinks are never followed (nor counted: a symlinked ``.py`` is not a
file of this workspace), and once ``max_dirs`` directories or
:data:`MAX_SCANNED_ENTRIES` entries have been scanned the walk stops and
``partial`` is True. Breadth-first so a partial
scan still covers the project's shallow structure rather than one deep
corner of it. Entries are visited in name order, so a partial count is
deterministic for a given tree.
"""
python_files = 0
test_files = 0
scanned = 0
entries_seen = 0
pending: deque[str] = deque([str(root)])
while pending:
if scanned >= max_dirs or entries_seen >= MAX_SCANNED_ENTRIES:
return python_files, test_files, True
directory = pending.popleft()
scanned += 1
try:
with os.scandir(directory) as scan:
items = sorted(scan, key=lambda entry: entry.name)
except OSError:
continue
entries_seen += len(items)
for entry in items:
name = entry.name
try:
if entry.is_dir(follow_symlinks=False):
if name not in _IGNORED_NAMES:
pending.append(entry.path)
continue
if name.endswith(".py") and entry.is_file(follow_symlinks=False):
python_files += 1
if name.startswith("test_"):
test_files += 1
except OSError:
continue
return python_files, test_files, False


def _is_within(child: Path, parent: Path) -> bool:
try:
child.relative_to(parent)
Expand All @@ -70,5 +160,4 @@ def _is_within(child: Path, parent: Path) -> bool:
return False


def _is_countable(path: Path) -> bool:
return path.is_file() and not any(part in _IGNORED_NAMES for part in path.parts)
__all__ = ["MAX_SCANNED_DIRS", "MAX_SCANNED_ENTRIES", "build_workspace_snapshot", "count_python_files"]
43 changes: 42 additions & 1 deletion src/server/agent_server.py
Original file line number Diff line number Diff line change
Expand Up @@ -404,7 +404,10 @@ async def _handle_control_request(self, msg: dict) -> None:
# refuse (critic C8). ``interrupt`` is exempt (a benign abort of a
# non-existent turn). Mirrors the ``session not ready`` pattern the
# permission-mode handlers already use.
if self.init_error is not None and subtype != "interrupt":
# ``get_activity`` is exempt too: a runtime that refused to start
# has no work, and answering "refused" would read as busy — leaving
# it unreleasable by a conditional close.
if self.init_error is not None and subtype not in ("interrupt", "get_activity"):
self._reply(request_id, {"ok": False, "error": self.init_error})
return
if subtype == "interrupt":
Expand Down Expand Up @@ -800,6 +803,9 @@ async def _handle_control_request(self, msg: dict) -> None:
if subtype == "resume":
self._do_resume(request_id, inner.get("session_id"))
return
if subtype == "get_activity":
self._reply(request_id, self._activity_snapshot())
return
if subtype == "branch":
self._do_branch(request_id)
return
Expand Down Expand Up @@ -3096,6 +3102,41 @@ def _do_branch(self, request_id: object) -> None:
logger.exception("[agent-server] branch failed")
self._reply(request_id, {"ok": False, "error": str(exc)})

def _activity_snapshot(self) -> dict[str, Any]:
"""What this session is doing that a gateway cannot see from outside.

A gateway knows the turns IT submitted. It does not know about a
/goal continuation the agent queued for itself, a /loop or cron job
waiting to fire, a prompt still in the inbox, or a background shell
the agent started — all of which an idle-looking session may carry,
and all of which ``session.close`` with ``if_idle`` must not kill.
"""
with self._lock:
turn_active = self._current_abort is not None
goal_active = self._goal_mgr is not None and self._goal_mgr.is_active()
queued = not self._inbox.empty()
scheduled = False
try:
scheduled = bool(self.cron_scheduler.list_jobs()) or self.cron_scheduler.wakeup_info() is not None
except Exception: # noqa: BLE001 — the scheduler is optional state
logger.debug("[agent-server] scheduler probe failed", exc_info=True)
background = False
try:
background = self._bgtasks is not None and any(
getattr(task, "status", "") == "running" for task in self._bgtasks.list()
)
except Exception: # noqa: BLE001
logger.debug("[agent-server] background task probe failed", exc_info=True)
return {
"ok": True,
"turn_active": turn_active,
"queued": queued,
"goal_active": goal_active,
"scheduled": scheduled,
"background": background,
"busy": turn_active or queued or goal_active or scheduled or background,
}

def _do_resume(self, request_id: object, session_id: object) -> None:
"""Load a saved conversation into this session (the original's /resume).
Idle-only — replacing the conversation mid-turn would race the worker."""
Expand Down
20 changes: 20 additions & 0 deletions src/server/desktop_gateway.py
Original file line number Diff line number Diff line change
Expand Up @@ -51,6 +51,16 @@ async def send_event(
await _send_json(websocket, {"method": "event", "params": params})


#: Pure reads that answer beside whatever the socket is doing, instead of
#: queueing behind it. Everything else runs one call at a time per socket,
#: in order — a prompt after an interrupt, a close after a submit — which is
#: what a client relying on the answer order needs. These two are file reads
#: off the event loop: the sidebar click's transcript must not sit behind
#: the runtime attach the previous click started, or the teardown of the
#: session being left.
CONCURRENT_METHODS = frozenset({"session.history", "projects.tree"})


async def handle_gateway_socket(websocket: WebSocket, state: DesktopServeState) -> None:
"""Accept one gateway socket and pump it until disconnect."""
await websocket.accept()
Expand All @@ -59,6 +69,7 @@ async def handle_gateway_socket(websocket: WebSocket, state: DesktopServeState)

conn = GatewayConnection(websocket=websocket, state=state)
await conn.on_open()
concurrent: set[asyncio.Task] = set()
try:
while True:
try:
Expand All @@ -69,8 +80,15 @@ async def handle_gateway_socket(websocket: WebSocket, state: DesktopServeState)
continue
if not isinstance(frame, dict):
continue
if frame.get("method") in CONCURRENT_METHODS:
task = asyncio.create_task(_dispatch(conn, frame))
concurrent.add(task)
task.add_done_callback(concurrent.discard)
continue
await _dispatch(conn, frame)
finally:
for task in concurrent:
task.cancel()
await conn.on_close()


Expand All @@ -92,6 +110,8 @@ async def _dispatch(conn: Any, frame: dict[str, Any]) -> None:

try:
result = await handler(params)
except asyncio.CancelledError:
raise
except Exception as exc: # noqa: BLE001 — one bad call must not drop the socket
logger.warning("gateway: %s failed", method, exc_info=True)
if request_id is not None:
Expand Down
Loading
Loading