diff --git a/configs/pi/extensions/workflow-context.ts b/configs/pi/extensions/workflow-context.ts index 7f042a97..9340cbdb 100644 --- a/configs/pi/extensions/workflow-context.ts +++ b/configs/pi/extensions/workflow-context.ts @@ -44,9 +44,16 @@ function runInBackground(args: string[], cwd?: string): void { export default function workflowContextExtension(pi: ExtensionAPI): void { pi.on("before_agent_start", async (event: any, ctx: any) => { try { + // 45s, matching the Claude hook `witan setup` installs. 5s cleared + // the warm output-cache hit (0.6-0.9s measured in agent-kit#349) + // but sat far below the cold path (16-23s there), so every prompt + // that missed the 30s cache silently contributed no block. A + // timeout is still wanted — a hung read must degrade to no + // context rather than stall the turn — but it has to sit above + // the cold path, not inside it. const r = spawnSync("witan", ["inject-context"], { encoding: "utf8", - timeout: 5000, + timeout: 45000, cwd: ctx?.cwd, }); const text = (r.stdout ?? "").trim(); diff --git a/mcp/servers/witan/CHANGELOG.md b/mcp/servers/witan/CHANGELOG.md index b3612a94..62751935 100644 --- a/mcp/servers/witan/CHANGELOG.md +++ b/mcp/servers/witan/CHANGELOG.md @@ -8,6 +8,40 @@ a MINOR bump may include breaking changes). ## [Unreleased] +### Fixed + +- **`task_ready` no longer re-reads blockers it has already fetched.** Its + repo-scoped branch scans every Task and then narrows to the candidates, but + resolved blocker statuses against the narrowed set — so every blocker living + in another repo was unknown and fetched again, one `get_task` per blocker. + Measured against a deployment that was ~3.0s of a 4.2s call, on identical + scans to a 1.2s `task_list`. + +- **The context hook issues its independent reads together.** `witan + inject-context` against a deployment made up to ten sequential tool calls, so + the cold path was their sum. It now runs in two waves (the second needs the + first's answers), which measured 11.3s -> 6.5s on the same graph, byte-identical + output. Each read keeps its own failure isolation, and a machine that cannot + start threads falls back to running them in line. The wave primes the proxy's + tool schema once first, so the fan-out does not have every worker discover it + (4 `tools/list` cold, against 1 for the sequential path it replaced); against + a `witan-core` predating `prime_tool_schema` this is skipped and the old + behaviour stands. + +- **The context hook's ready list respects cross-repo blockers.** It offers a + repo-scoped slice of an all-Task scan, and an unresolvable blocker counts as + closed — so a task blocked by an open task in another repo was advertised as + ready to work. `readiness.filter_ready` now takes the wider row set, which + the hook already had in hand. + +- **The prompt and stop hooks get timeouts above what they are timing.** Both + were 15s. `witan inject-context` was measured at 16-23s cold on a large graph, + so the hook was killed mid-read and the user paid the full wait for no block; + it is now 45s. `witan session-checkpoint` writes, and a write has been measured + at up to 51s, so it is now 60s — being killed there leaves a session open with + no handoff summary. The pi `workflow-context` extension was tighter still at + 5s and now matches at 45s. + ## [0.36.0] - 2026-09-18 ### Added diff --git a/mcp/servers/witan/tests/test_context.py b/mcp/servers/witan/tests/test_context.py index 82232bbe..a59d8f3e 100644 --- a/mcp/servers/witan/tests/test_context.py +++ b/mcp/servers/witan/tests/test_context.py @@ -1336,3 +1336,176 @@ def test_inject_context_remote_asks_for_held_tasks_across_all_repos( "task_list", {"assignee": "@me", "status": "in_progress", "repo": ""}, ) in server.calls + + +def test_inject_context_remote_issues_its_independent_reads_concurrently( + tmp_path, monkeypatch +): + """The cold path is the slowest read, not the sum of them. + + Issued one at a time, the hook's ten round trips against a deployment took + ~11s measured and blew the 15s timeout it installs (agent-kit#349). Each + read here blocks on a barrier that only releases once every first-wave + call has arrived, so a sequential implementation deadlocks and times out + rather than quietly passing slower. + """ + import threading + + from witan import context as ctx_module + + monkeypatch.setenv("TMPDIR", str(tmp_path)) + monkeypatch.setenv("WITAN_CONTEXT_TTL", "0") + import tempfile + + monkeypatch.setattr(tempfile, "tempdir", None) + repo = "https://github.com/test/ctx-concurrent" + monkeypatch.setenv("WITAN_REPO", repo) + monkeypatch.setattr(ctx_module, "_current_branch", lambda: "main") + + # The four first-wave reads. The second wave (sessions, comments) depends + # on their answers, so it is deliberately not held here. + first_wave = { + "workflow_project_list", + "task_ready", + "task_for_branch", + "task_list", + } + barrier = threading.Barrier(len(first_wave), timeout=10) + + class _Barriered(_FakeRemoteServer): + def _guard(self, name): + if name in first_wave: + barrier.wait() + super()._guard(name) + + server = _Barriered( + projects=[{"slug": "wp-x", "title": "X", "phase": "spec"}], + ready=[], + sessions_by_project={}, + branch_tasks=[], + held=[], + ) + text = ctx_module.inject_context_remote( + server, "https://witan.example.org/mcp", debug=True + ) + + assert barrier.broken is False + assert "## Active Workflow Projects" in text + + +def test_inject_context_remote_primes_the_tool_schema_once_before_fanning_out( + tmp_path, monkeypatch +): + """One schema resolution per wave, and only when the proxy offers one. + + Without it each worker finds the param-name cache unset and lists the + surface itself (measured: 4 `tools/list` cold, where the sequential path + issued 1). The catch-all `__getattr__` case is the trap: a proxy predating + `prime_tool_schema` must be detected as NOT having it, rather than having + a bogus tool of that name called on the deployment. + """ + from witan import context as ctx_module + + monkeypatch.setenv("TMPDIR", str(tmp_path)) + monkeypatch.setenv("WITAN_CONTEXT_TTL", "0") + import tempfile + + monkeypatch.setattr(tempfile, "tempdir", None) + monkeypatch.setenv("WITAN_REPO", "https://github.com/test/ctx-prime") + + class _Priming(_FakeRemoteServer): + primed = 0 + + def prime_tool_schema(self): + type(self).primed += 1 + + server = _Priming(projects=[], ready=[], sessions_by_project={}) + ctx_module.inject_context_remote(server, "https://witan.example.org/mcp") + assert _Priming.primed == 1 + assert not [c for c in server.calls if c[0] == "prime_tool_schema"] + + # A proxy-shaped object whose __getattr__ answers for ANY name, i.e. one + # from a witan-core that predates the method. + class _CatchAll(_FakeRemoteServer): + def __getattr__(self, name): + def _tool_call(**kwargs): + self.calls.append((name, kwargs)) + raise RuntimeError(f"no such tool: {name}") + + return _tool_call + + legacy = _CatchAll(projects=[], ready=[], sessions_by_project={}) + ctx_module.inject_context_remote(legacy, "https://witan.example.org/mcp") + assert not [c for c in legacy.calls if c[0] == "prime_tool_schema"] + + +def test_gather_falls_back_to_serial_when_the_pool_cannot_start(monkeypatch): + """A machine at its thread limit still gets its answers, just slower. + + ``_gather`` is called outside the remote path's own try/except, so raising + here would crash a hook whose entire contract is that it degrades quietly. + """ + from witan import context as ctx_module + + def _no_threads(*_args, **_kwargs): + raise RuntimeError("can't start new thread") + + monkeypatch.setattr(ctx_module, "ThreadPoolExecutor", _no_threads) + + def _boom(): + raise ValueError("read failed") + + got = ctx_module._gather({"a": lambda: 1, "b": lambda: 2, "c": _boom}) + + assert got["a"] == 1 + assert got["b"] == 2 + # And a genuine read failure is still reported as a value, not raised. + assert isinstance(got["c"], ValueError) + + +def test_gather_reports_each_failure_against_its_own_key(): + from witan import context as ctx_module + + def _boom(): + raise ValueError("nope") + + got = ctx_module._gather({"ok": lambda: "fine", "bad": _boom}) + + assert got["ok"] == "fine" + assert isinstance(got["bad"], ValueError) + + +def test_inject_context_remote_one_failed_read_costs_only_its_own_block( + tmp_path, monkeypatch +): + """Per-read isolation must survive the move into worker threads. + + A read that fails inside a worker comes back as a value rather than a live + exception, so the block it feeds is the only thing that may disappear. + """ + from witan import context as ctx_module + + monkeypatch.setenv("TMPDIR", str(tmp_path)) + monkeypatch.setenv("WITAN_CONTEXT_TTL", "0") + import tempfile + + monkeypatch.setattr(tempfile, "tempdir", None) + repo = "https://github.com/test/ctx-isolated" + monkeypatch.setenv("WITAN_REPO", repo) + monkeypatch.setattr(ctx_module, "_current_branch", lambda: "main") + + server = _FakeRemoteServer( + projects=[{"slug": "wp-x", "title": "X", "phase": "spec"}], + ready=[{"slug": "tk-r", "title": "Ready one", "priority": "p1"}], + sessions_by_project={}, + branch_tasks=[{"slug": "tk-b", "title": "Branch one", "status": "open"}], + raises=("task_for_branch",), + ) + text = ctx_module.inject_context_remote( + server, "https://witan.example.org/mcp", debug=True + ) + + assert "## In-Flight Branch" not in text + assert "## Active Workflow Projects" in text + assert "## Ready Tasks" in text + assert "Ready one" in text diff --git a/mcp/servers/witan/tests/test_readiness.py b/mcp/servers/witan/tests/test_readiness.py index 792613b5..89063b97 100644 --- a/mcp/servers/witan/tests/test_readiness.py +++ b/mcp/servers/witan/tests/test_readiness.py @@ -128,3 +128,41 @@ def test_filter_ready_orders_by_priority_and_reclaims_expired(): def test_filter_ready_unknown_blocker_treated_closed(): tasks = [{"slug": "tk-x", "status": "open", "blocked_by": ["tk-gone"]}] assert [t["slug"] for t in readiness.filter_ready(tasks)] == ["tk-x"] + + +def test_filter_ready_resolves_blockers_from_the_wider_row_set(): + """A blocker outside the candidate slice holds its task back. + + The context hook offers a repo-scoped slice of an all-Task scan. Without + the wider set the open blocker below is simply unknown, and unknown reads + as closed — so the task would be advertised as ready to work while the + thing blocking it is wide open (agent-kit#349). + """ + tasks = [{"slug": "tk-x", "status": "open", "blocked_by": ["tk-other-repo"]}] + others = [ + {"slug": "tk-other-repo", "status": "open", "repo": "https://example.com/b"} + ] + + assert [t["slug"] for t in readiness.filter_ready(tasks)] == ["tk-x"] + assert readiness.filter_ready(tasks, blocker_rows=others) == [] + + others[0]["status"] = "closed" + assert [t["slug"] for t in readiness.filter_ready(tasks, blocker_rows=others)] == [ + "tk-x" + ] + + +def test_filter_ready_candidate_status_wins_over_a_stale_wider_row(): + """Where a slug is in both sets, the candidate row decides. + + ``tk-b`` is closed among the candidates and open in the wider set. Reading + the wider copy would hold ``tk-a`` back against data the caller has already + superseded. (``tk-b`` itself is absent either way — closed is not pickable.) + """ + tasks = [ + {"slug": "tk-a", "status": "open", "blocked_by": ["tk-b"]}, + {"slug": "tk-b", "status": "closed"}, + ] + stale = [{"slug": "tk-b", "status": "open"}] + got = [t["slug"] for t in readiness.filter_ready(tasks, blocker_rows=stale)] + assert got == ["tk-a"] diff --git a/mcp/servers/witan/tests/test_setup.py b/mcp/servers/witan/tests/test_setup.py index 3ea65370..269fb72e 100644 --- a/mcp/servers/witan/tests/test_setup.py +++ b/mcp/servers/witan/tests/test_setup.py @@ -48,8 +48,13 @@ def test_witan_bundle_registers_witan_mcp_server_and_hooks(tmp_path, monkeypatch ] assert inject and checkpoint # Both prompt-path hooks carry a timeout so a hung git/graph can't stall. - assert inject[0]["timeout"] == 15 - assert checkpoint[0]["timeout"] == 15 + # It has to sit ABOVE the cold-path read cost, not inside it: at 15s the + # hook was killed mid-read on a large graph, so the user paid the full + # wait and got no block (agent-kit#349). + assert inject[0]["timeout"] == 45 + # Longer than the read hook: this one is a write, and a write against a + # deployment has been measured at up to 51s. + assert checkpoint[0]["timeout"] == 60 def test_mcp_only_platforms_keep_the_self_contained_uvx_entry(tmp_path, monkeypatch): diff --git a/mcp/servers/witan/tests/test_tasks.py b/mcp/servers/witan/tests/test_tasks.py index 8e438b74..1dae7ca7 100644 --- a/mcp/servers/witan/tests/test_tasks.py +++ b/mcp/servers/witan/tests/test_tasks.py @@ -217,6 +217,49 @@ def test_create_with_already_closed_blocker_is_open(server): assert b["slug"] in {t["slug"] for t in server.task_ready()} +@requires_omnigraph +def test_ready_resolves_a_cross_repo_blocker_without_re_reading_it(server, monkeypatch): + """A blocker in another repo is already in the all-Task scan — use it. + + ``task_ready``'s repo-scoped branch reads every Task and then narrows to + the candidates. Resolving blocker statuses against the NARROW set left + every cross-repo blocker unknown, so each one was fetched again one + ``get_task`` at a time: ~3.0s of a 4.2s call against the deployed service + and most of what blew the context hook's 15s timeout (agent-kit#349). + """ + from witan import server as srv + + other = "https://github.com/test/other" + blocker = server.task_create(title="blocker elsewhere", description="x", repo=other) + dependent = server.task_create( + title="dependent here", description="x", blocked_by=[blocker["slug"]] + ) + + reads: list[str] = [] + real_read = srv.client.read + + def counting_read(queries, name, params): + reads.append(name) + return real_read(queries, name, params) + + monkeypatch.setattr(srv.client, "read", counting_read) + ready = {t["slug"] for t in server.task_ready(repo="https://github.com/test/repo")} + + # The open blocker lives in another repo, so it is not a candidate — but it + # still holds its dependent back. + assert blocker["slug"] not in ready + assert dependent["slug"] not in ready + assert "get_task" not in reads + + # And the answer is not merely "everything is excluded": closing the + # blocker releases the dependent, still off the same two scans. + server.task_close(blocker["slug"]) + reads.clear() # after the close, whose own reads are not what is counted + ready = {t["slug"] for t in server.task_ready(repo="https://github.com/test/repo")} + assert dependent["slug"] in ready + assert "get_task" not in reads + + @requires_omnigraph def test_update_to_closed_unblocks_dependents(server): a = server.task_create(title="blocker", description="x") diff --git a/mcp/servers/witan/witan/context.py b/mcp/servers/witan/witan/context.py index 9be29043..f5ece416 100644 --- a/mcp/servers/witan/witan/context.py +++ b/mcp/servers/witan/witan/context.py @@ -14,6 +14,7 @@ import os import time from collections.abc import Callable +from concurrent.futures import ThreadPoolExecutor, as_completed from pathlib import Path from witan_core.observability import get_logger @@ -249,6 +250,75 @@ def _dbg_exc(enabled: bool, msg: str) -> None: ) +def _dbg_exc_from(enabled: bool, exc: BaseException, msg: str) -> None: + """:func:`_dbg_exc` for an exception that is not the *active* one. + + A read that failed inside a worker thread is caught and carried back as a + value, so by the time its block is skipped there is no live exception for + ``exc_info=True`` to resolve. Passing the object keeps the traceback in the + record instead of losing it to an empty ``sys.exc_info()``. + """ + if enabled: + logger.debug("witan.context.debug", detail=msg, exc_info=exc) + + +# The hook's reads are network-bound against a deployment and mostly +# independent of one another, so issuing them one at a time made the cold path +# the SUM of ten round trips (~11s measured) rather than the slowest of two +# waves. Threads rather than asyncio: ``server`` is duck-typed on plain +# attribute-style methods (``RemoteMCPProxy.__getattr__`` wraps each call in +# its own ``asyncio.run``), and every call is I/O-bound, so a pool parallelises +# them without changing that interface. +_HOOK_READ_WORKERS = 6 + + +def _sequentially(calls: dict[str, Callable[[], object]]) -> dict[str, object]: + """:func:`_gather`'s contract, one call at a time.""" + out: dict[str, object] = {} + for key, fn in calls.items(): + try: + out[key] = fn() + except Exception as exc: # noqa: BLE001 — reported to the caller as a value + out[key] = exc + return out + + +def _gather(calls: dict[str, Callable[[], object]]) -> dict[str, object]: + """Run every thunk concurrently; map each key to its result OR its exception. + + Returning failures as values rather than raising is what lets each caller + keep the per-read isolation the hook depends on: one read failing must cost + its own block and nothing else, and a shared ``gather`` that propagated the + first exception would take the whole prompt's context with it. + + NEVER RAISES, including when the pool itself cannot be built — a machine at + its thread limit falls back to running the calls in line rather than taking + down a hook whose whole contract is that it degrades quietly. Only the + speed-up is optional; the answers are not. + """ + if len(calls) <= 1: + return _sequentially(calls) + out: dict[str, object] = {} + try: + with ThreadPoolExecutor( + max_workers=min(len(calls), _HOOK_READ_WORKERS), + thread_name_prefix="witan-ctx", + ) as pool: + futures = {pool.submit(fn): key for key, fn in calls.items()} + for future in as_completed(futures): + key = futures[future] + try: + out[key] = future.result() + except Exception as exc: # noqa: BLE001 — see docstring + out[key] = exc + except Exception: # noqa: BLE001 — pool/thread failure, not a read failure + logger.debug("witan.context.gather_fell_back_to_serial", exc_info=True) + # Keep whatever already came back; only re-run what never answered. + pending = {k: fn for k, fn in calls.items() if k not in out} + return {**out, **_sequentially(pending)} + return out + + def inject_context( graph_uri: str, queries_dir: Path, @@ -381,7 +451,10 @@ def comments_for(slug: str) -> list[dict]: # Shared with ``task_ready`` so the injected list and the tool agree — # including the reclaim of ``in_progress`` tasks whose lease has lapsed. - ready = readiness.filter_ready(tasks) + # ``all_rows`` is the all-Task scan ``tasks`` was sliced out of, so passing + # it costs no read and lets a cross-repo blocker resolve to its real status + # instead of defaulting to closed. + ready = readiness.filter_ready(tasks, blocker_rows=all_rows) _dbg( debug, f"ready={len(ready)} open_branch_tasks={len(open_branch_tasks)} " @@ -706,45 +779,83 @@ def inject_context_remote(server, remote_url: str, debug: bool = False) -> str: # list_unscoped_tasks locally). limit=10000 matches the local path's # list_unscoped_tasks cap so a busy graph's header count isn't # silently truncated by task_ready's own default limit=20. - projects = ( - server.workflow_project_list(repo=repo, status="active") if repo else [] - ) - ready = server.task_ready(repo=(repo or ""), limit=10000) + # + # `repo=""` on the held-task read is load-bearing, not decoration. An + # OMITTED repo is filled in client-side by `RemoteMCPProxy._map_args` + # (`task_list` is in `_REPO_IS_SCOPE_OR_STAMP`), which would scope this + # to whichever checkout the hook fired in — while the local path + # answers from an all-repos `list_unscoped_tasks` scan. The sentinel + # makes both transports answer the same question: which tasks do I + # hold, anywhere. + first: dict[str, Callable[[], object]] = { + "ready": lambda: server.task_ready(repo=(repo or ""), limit=10000), + "held": lambda: server.task_list( + assignee="@me", status="in_progress", repo="" + ), + } + if repo: + first["projects"] = lambda: server.workflow_project_list( + repo=repo, status="active" + ) + if repo and branch: + first["branch"] = lambda: server.task_for_branch(branch=branch, repo=repo) + + # Resolve the deployment's tool surface ONCE before fanning out. Each + # worker otherwise finds the proxy's param-name cache unset (none of + # them has finished listing yet) and lists it itself: measured at four + # `tools/list` sequences on a cold process where the sequential path + # issued one. The listings overlap, so this is server load rather than + # latency, and it is load every agent session repeats. + # + # Feature-detected ON THE CLASS, not the instance. `RemoteMCPProxy` + # has a catch-all `__getattr__` that turns any unknown attribute into + # a TOOL CALL, so an instance-level `getattr(server, ..., None)` never + # returns None — against a witan-core predating this method it would + # hand back a closure that tries to invoke a `prime_tool_schema` tool + # on the deployment. A class lookup is not intercepted (`__getattr__` + # is an instance hook), so this is absent exactly when the method is. + # That keeps the `witan-core>=0.37` floor honest: older cores simply + # skip priming and behave as they did before. + # + # Best-effort: a failure here is not worth losing the block over, + # since each worker resolves the schema itself anyway. + if getattr(type(server), "prime_tool_schema", None) is not None: + try: + server.prime_tool_schema() + except Exception: # noqa: BLE001 — the workers still resolve it + _dbg_exc(debug, "tool-schema prime failed (workers will resolve it)") + + done = _gather(first) + + # The projects/ready pair is what the block exists to show, so a + # failure in either still returns "" — same contract as before, just + # re-raised out of the worker that hit it rather than raised inline. + for key in ("projects", "ready"): + if isinstance(done.get(key), Exception): + raise done[key] + projects: list[dict] = done.get("projects") or [] + ready: list[dict] = done["ready"] _dbg(debug, f"remote reads OK: projects={len(projects)} ready={len(ready)}") except Exception: # noqa: BLE001 _dbg_exc(debug, "FAILED building remote context (returning empty block)") return "" - sessions_by_project: dict[str, list[dict]] = {} - for p in projects[:3]: - try: - sessions_by_project[p["slug"]] = server.workflow_session_list( - project_slug=p["slug"] - ) - except Exception: # noqa: BLE001 - _dbg_exc( - debug, f"sessions read failed for {p['slug']!r} (skipping resume lines)" - ) - # Isolated exactly as the local path isolates its own CodeBranch read # (see :func:`inject_context`): a deployment that predates # ``task_for_branch`` answers with an unknown-tool error, and that must # cost the branch block only — never the projects/ready-tasks block that # already rendered above it. open_branch_tasks: list[dict] = [] - if repo and branch: - try: - open_branch_tasks = [ - t - for t in server.task_for_branch(branch=branch, repo=repo) - if t.get("status") != "closed" - ] - except Exception: # noqa: BLE001 - _dbg_exc( - debug, - "task_for_branch failed (deployment may predate it) — " - "no In-Flight Branch block", - ) + branch_result = done.get("branch") + if isinstance(branch_result, Exception): + _dbg_exc_from( + debug, + branch_result, + "task_for_branch failed (deployment may predate it) — " + "no In-Flight Branch block", + ) + elif branch_result: + open_branch_tasks = [t for t in branch_result if t.get("status") != "closed"] _dbg( debug, @@ -754,21 +865,57 @@ def inject_context_remote(server, remote_url: str, debug: bool = False) -> str: # Isolated like the branch read above: a deployment predating `@me` or # `TaskComment` costs this block and nothing else. + held_result = done.get("held") + if isinstance(held_result, Exception): + _dbg_exc_from(debug, held_result, "held-task read failed — no comment block") + held: list[dict] = [] + else: + held = held_result or [] + + # Second wave. Both sets depend on a first-wave answer — sessions on + # `projects`, comments on `held` — so they cannot join the batch above, but + # they are independent of each other and go out together. # - # `repo=""` is load-bearing, not decoration. An OMITTED repo is filled in - # client-side by `RemoteMCPProxy._map_args` (`task_list` is in - # `_REPO_IS_SCOPE_OR_STAMP`), which would scope this to whichever checkout - # the hook fired in — while the local path answers from an all-repos - # `list_unscoped_tasks` scan. The sentinel makes both transports answer the - # same question: which tasks do I hold, anywhere. - try: - held = server.task_list(assignee="@me", status="in_progress", repo="") - except Exception: # noqa: BLE001 - _dbg_exc(debug, "held-task read failed — no comment block") - held = [] + # The comment prefetch slices `held` exactly as `_held_task_comments` does, + # so every fetched thread is one that function will look at and no task is + # read speculatively. + second: dict[str, Callable[[], object]] = { + f"sessions:{p['slug']}": ( + lambda slug=p["slug"]: server.workflow_session_list(project_slug=slug) + ) + for p in projects[:3] + } + second.update( + { + f"comments:{t['slug']}": (lambda slug=t["slug"]: server.task_get(slug=slug)) + for t in held[:_HELD_TASK_LIMIT] + } + ) + fetched = _gather(second) + + sessions_by_project: dict[str, list[dict]] = {} + for p in projects[:3]: + result = fetched.get(f"sessions:{p['slug']}") + if isinstance(result, Exception): + _dbg_exc_from( + debug, + result, + f"sessions read failed for {p['slug']!r} (skipping resume lines)", + ) + else: + sessions_by_project[p["slug"]] = result or [] def comments_for(slug: str) -> list[dict]: - return (server.task_get(slug=slug) or {}).get("comments") or [] + """Hand back the prefetched thread, re-raising whatever fetching it hit. + + Re-raising rather than swallowing keeps ``_held_task_comments``'s own + per-task isolation (and its debug line) in charge of a failed read, + exactly as when it did the fetching itself. + """ + result = fetched.get(f"comments:{slug}") + if isinstance(result, Exception): + raise result + return (result or {}).get("comments") or [] commented = _held_task_comments(held, comments_for, remote_url, debug=debug) diff --git a/mcp/servers/witan/witan/extensions/pi/workflow-context.ts b/mcp/servers/witan/witan/extensions/pi/workflow-context.ts index 7f042a97..9340cbdb 100644 --- a/mcp/servers/witan/witan/extensions/pi/workflow-context.ts +++ b/mcp/servers/witan/witan/extensions/pi/workflow-context.ts @@ -44,9 +44,16 @@ function runInBackground(args: string[], cwd?: string): void { export default function workflowContextExtension(pi: ExtensionAPI): void { pi.on("before_agent_start", async (event: any, ctx: any) => { try { + // 45s, matching the Claude hook `witan setup` installs. 5s cleared + // the warm output-cache hit (0.6-0.9s measured in agent-kit#349) + // but sat far below the cold path (16-23s there), so every prompt + // that missed the 30s cache silently contributed no block. A + // timeout is still wanted — a hung read must degrade to no + // context rather than stall the turn — but it has to sit above + // the cold path, not inside it. const r = spawnSync("witan", ["inject-context"], { encoding: "utf8", - timeout: 5000, + timeout: 45000, cwd: ctx?.cwd, }); const text = (r.stdout ?? "").trim(); diff --git a/mcp/servers/witan/witan/readiness.py b/mcp/servers/witan/witan/readiness.py index 798df37f..81b574ce 100644 --- a/mcp/servers/witan/witan/readiness.py +++ b/mcp/servers/witan/witan/readiness.py @@ -106,15 +106,30 @@ def is_ready( return all(blocker_status(b) == "closed" for b in (task.get("blocked_by") or [])) -def filter_ready(tasks: list[dict], *, now: datetime | None = None) -> list[dict]: +def filter_ready( + tasks: list[dict], + *, + blocker_rows: list[dict] | None = None, + now: datetime | None = None, +) -> list[dict]: """Ready tasks from a self-contained list, priority-ordered (p0 first). Blocker statuses are resolved within ``tasks`` — an unknown blocker is treated as closed. Use this when the full candidate set is already in hand (the context hook); ``task_ready`` supplies its own resolver that can fetch blockers outside the list. + + ``blocker_rows`` widens only the *lookup*, never the candidate set: a + caller that scanned every Task but is offering a repo-scoped slice of them + passes the whole scan here. Without it a blocker in another repo is + unknown, and unknown reads as closed — so a task whose cross-repo blocker + is wide open is shown as ready. That is the expensive direction to fail in + for a signal whose job is to stop two actors colliding, and the rows are + already in hand either way. """ - status_by_slug = {t["slug"]: t.get("status") for t in tasks} + status_by_slug = { + t["slug"]: t.get("status") for t in (*(blocker_rows or ()), *tasks) + } def blocker_status(slug: str) -> str: return status_by_slug.get(slug) or "closed" diff --git a/mcp/servers/witan/witan/server.py b/mcp/servers/witan/witan/server.py index 20bc3205..48342098 100644 --- a/mcp/servers/witan/witan/server.py +++ b/mcp/servers/witan/witan/server.py @@ -6939,6 +6939,14 @@ def task_ready( limit: Maximum tasks to return. Defaults to 20. """ + # The candidate set and the blocker-status lookup are built from different + # row sets on purpose. A task in ANOTHER repo is never a candidate here, + # but it is routinely a blocker of one — and ``list_unscoped_tasks`` is an + # all-Task scan, so it has already been fetched. Narrowing the lookup to + # the candidates threw those rows away and then re-read them one blocker at + # a time: ~3.0s of a 4.2s ``task_ready`` against the deployed service, + # which is most of what blows the context hook's timeout (agent-kit#349). + known: list[dict] = [] if project_slug: rows = client.read( "read.gq", "list_tasks_by_project", {"project_slug": project_slug} @@ -6953,10 +6961,13 @@ def task_ready( r for r in all_rows if not r.get("repo") and r["slug"] not in seen ] rows = repo_rows + unscoped + known = all_rows else: rows = client.read("read.gq", "list_unscoped_tasks", {}) - status_by_slug = {r["slug"]: r.get("status") for r in rows} + # ``rows`` last: where a slug appears in both, the candidate row is the one + # the readiness decision is about, so its status wins. + status_by_slug = {r["slug"]: r.get("status") for r in (*known, *rows)} def blocker_status(blocker_slug: str) -> str: if blocker_slug in status_by_slug: diff --git a/mcp/servers/witan/witan/setup.py b/mcp/servers/witan/witan/setup.py index 17f71415..6171d1a9 100644 --- a/mcp/servers/witan/witan/setup.py +++ b/mcp/servers/witan/witan/setup.py @@ -62,19 +62,37 @@ def witan_bundle(pkg_dir: Path, author: str) -> RegistrationBundle: # These hooks run on the prompt/stop critical path and do git + graph I/O, so # they carry a timeout: a hung git or graph read must degrade to no context, # never stall the agent. The first prompt in a cache window does several - # full-store reads, which on a large graph can take ~10s — 15s gives that - # cold path headroom to finish (and populate the on-disk cache) instead of - # being killed, which would leave the cache empty and every prompt cold. + # full-store reads, which on a large graph can take ~10s — the timeout has + # to give that cold path headroom to finish (and populate the on-disk cache) + # instead of being killed, which leaves the cache empty and every prompt + # cold. + # + # Both were 15s, which sat INSIDE the cost distribution of the work they + # were timing rather than above it — so the hook was killed mid-read, the + # user paid the full wait, and the output was discarded (agent-kit#349). + # The two differ now because they are timing different things: + # + # inject-context is a read path, measured at 16-23s cold on a graph with 19 + # active projects and 184 ready tasks. The reads behind that are fixed, so + # 45s is headroom for a graph bigger than the one measured, not a budget + # anything is expected to use. + # + # session-checkpoint is a WRITE path (`workflow_session_end`), and a single + # write against a deployment has been measured at up to 51s — see the + # credential-refresh note in `witan_core.remote.proxy._invoke`. Killing it + # is not a dropped block but a session left open with no handoff summary, + # which is invisible until someone resumes and finds nothing recorded. 60s + # covers that recorded worst case. hooks: list[Hook] = [ DeclarativeHook( event=HookEvent.USER_PROMPT_SUBMIT, command="witan inject-context", - timeout_seconds=15, + timeout_seconds=45, ), DeclarativeHook( event=HookEvent.STOP, command="witan session-checkpoint", - timeout_seconds=15, + timeout_seconds=60, ), ] if pi_ext_dir.is_dir(): diff --git a/packages/witan-core/CHANGELOG.md b/packages/witan-core/CHANGELOG.md index d2b28424..7bfafcb8 100644 --- a/packages/witan-core/CHANGELOG.md +++ b/packages/witan-core/CHANGELOG.md @@ -8,6 +8,18 @@ a MINOR bump may include breaking changes). ## [Unreleased] +### Added + +- **`RemoteMCPProxy.ensure_tool_schema()` / `prime_tool_schema()`**, for a + caller about to issue several tool calls concurrently. `_invoke_once` + resolves the tool surface lazily, which is right for sequential use but + means every worker in a fan-out lists it: measured at four `tools/list` + sequences for a four-call wave on a cold proxy, where the sequential shape + issued one. Priming first trades one connection for N-1 fewer listings. + Idempotent and connection-free when the cache is warm. Deliberately not a + lock inside `_invoke_once`, which would deadlock the shared-event-loop + caller (`witan.remote.serve`). + ### Changed - **`configure_sentry` no longer turns a dropped OTLP batch into a Sentry diff --git a/packages/witan-core/tests/test_remote_proxy.py b/packages/witan-core/tests/test_remote_proxy.py index 235b1891..65e8489f 100644 --- a/packages/witan-core/tests/test_remote_proxy.py +++ b/packages/witan-core/tests/test_remote_proxy.py @@ -961,6 +961,83 @@ def test_declared_ttl_bounds_how_long_the_list_is_held(monkeypatch): assert proxy.lists == 2 +def test_priming_the_schema_keeps_a_concurrent_wave_to_one_tools_list(): + """A fan-out must not make every worker discover the surface itself. + + Each worker checks ``_param_names is None`` before any of them has + finished listing, so without priming they all list. Measured on a cold + proxy against a real deployment: four ``tools/list`` for a four-call wave + where the sequential shape issued one (agent-kit#372 review). + """ + import threading + from concurrent.futures import ThreadPoolExecutor + + from witan_core import caching + + # A server that permits caching, which is the only case priming can help: + # against one declaring ttlMs=0 ("do not cache") every call re-lists by + # design, primed or not. + def cacheable(): + return _echo_server(**caching.hint_kwargs(ttl_seconds=300)) + + def wave(proxy): + with ThreadPoolExecutor(max_workers=4) as pool: + futures = [pool.submit(proxy.echo, value=str(i)) for i in range(4)] + return sorted(f.result() for f in futures) + + class _Barriered(_CountingProxy): + """Holds every worker at the listing until all four have reached it. + + The unprimed duplication is a RACE: against a fast in-memory server one + worker can finish listing before the others look, and then they read + its cache. Left to chance the "all four list" assertion passes or fails + with machine load. The barrier pins the interleaving that the real cold + path hits against a deployment, where the listing is a round trip. + """ + + def __init__(self, server, barrier): + super().__init__(server) + self._barrier = barrier + + def _new_client(self, token): + client = super()._new_client(token) + counted = client.list_tools_mcp + + async def _held(*args, **kwargs): + self._barrier.wait() + return await counted(*args, **kwargs) + + client.list_tools_mcp = _held + return client + + unprimed = _Barriered(cacheable(), threading.Barrier(4, timeout=10)) + assert wave(unprimed) == ["0", "1", "2", "3"] + assert unprimed.lists == 4 + + primed = _CountingProxy(cacheable()) + primed.prime_tool_schema() + assert primed.lists == 1 + assert wave(primed) == ["0", "1", "2", "3"] + # The four workers read the primed list rather than each fetching one. + assert primed.lists == 1 + + +def test_priming_a_warm_schema_opens_no_connection(): + """Idempotent, so a caller can prime unconditionally before every wave.""" + from witan_core import caching + + proxy = _CountingProxy(_echo_server(**caching.hint_kwargs(ttl_seconds=300))) + proxy.prime_tool_schema() + assert proxy.lists == 1 + + def _fail(_token): + raise AssertionError("prime_tool_schema opened a connection when warm") + + proxy._new_client = _fail + proxy.prime_tool_schema() + assert proxy.lists == 1 + + def test_a_zero_ttl_is_an_instruction_not_a_missing_value(): # A 2026-07-28 server that sets no hint still sends ttlMs=0, which says # "do not cache this" — so the proxy re-reads rather than treating the diff --git a/packages/witan-core/witan_core/remote/proxy.py b/packages/witan-core/witan_core/remote/proxy.py index e88cb289..21d7cb13 100644 --- a/packages/witan-core/witan_core/remote/proxy.py +++ b/packages/witan-core/witan_core/remote/proxy.py @@ -816,6 +816,65 @@ async def dispatch(self, name: str, /, **kwargs: Any) -> Any: raise RemoteToolUnavailable(self._admin_error(name)) return await self._invoke(name, (), kwargs) + async def ensure_tool_schema(self) -> None: + """Resolve the tool surface once, for a caller about to FAN OUT. + + ``_invoke_once`` resolves it lazily, which is right for sequential use: + the first call lists, every later one reads the cache. A concurrent + batch defeats that, because each worker checks ``_param_names is None`` + before any of them has finished listing, so they all list. Measured on + a cold proxy: a four-call wave issued four ``tools/list`` sequences + where the sequential shape issued one (agent-kit#372 review). + + The cost is server load, not latency — the redundant listings overlap, + so they add roughly one listing's wall time either way. Calling this + first trades one extra connection for N-1 fewer listings against a + deployment that every agent session is hitting. + + Deliberately NOT a lock inside ``_invoke_once``. The refresh happens + across an ``await``, and this proxy is driven both from one loop per + thread (the CLI/hook path) and from several coroutines on ONE shared + loop (``witan.remote.serve``); a ``threading.Lock`` held across that + await would deadlock the second case outright. Priming from the caller + that knows it is about to fan out needs no cross-context locking at + all. + + Idempotent and cheap when the cache is warm: no connection is opened. + :meth:`prime_tool_schema` is the synchronous form, for a caller that is + not already in a loop — the same pairing as :meth:`dispatch` and the + ``__getattr__`` call wrapper. + """ + if self._param_names is not None and time.monotonic() < ( + self._param_names_expiry + ): + return + token = self._token_provider() + async with ( + self._reclassifying("list_tools"), + AsyncExitStack() as stack, + ): + try: + client = await stack.enter_async_context(self._new_client(token)) + except Exception as exc: # noqa: BLE001 — see _invoke_once + rejected = auth_failure(exc) + if rejected is not None: + raise RemoteCredentialRejected( + self._credential_rejected_error( + "list_tools", rejected.response.status_code + ) + ) from exc + raise RemoteUnreachable(self._unreachable_error(exc)) from exc + await self._refresh_param_names(client) + + def prime_tool_schema(self) -> None: + """Synchronous :meth:`ensure_tool_schema`, for a caller outside a loop. + + Defined as a real method rather than left to ``__getattr__``, which + would otherwise treat the name as a TOOL and try to call one by that + name on the deployment. + """ + asyncio.run(self.ensure_tool_schema()) + async def remote_tools(self) -> list[Any]: """The deployment's advertised tools, as listed MCP tool objects.