Add real-time per-tool-call audit streaming (HARNESS-WP-0002-T03)

Claude Code executes its own tools internally in --print mode -- there
is no way for a caller to externally dispatch individual tool calls
without abandoning that self-contained agent model. What
--output-format stream-json --include-hook-events does allow: observing
each tool_use/tool_result/hook event in real time.

- adapter.py: AgenticClaudeCodeAdapter gains an optional on_tool_event
  callback; streaming mode (Popen + background reader thread) is used
  only when set, blocking subprocess.run path is unchanged otherwise.
- runner.py: run_task gains emit_tool_events/on_tool_event, collecting
  events onto RunResult.tool_events and posting a tool_call hub event
  per tool when reporting is enabled.
- cli.py: --stream-tool-events flag on `run`, prints each event as a
  tagged JSON line ahead of the unchanged final result block.

Live-verified against the real claude CLI: 5 real tool events streamed
correctly (2x Bash, 1x Write) plus Stop hook lifecycle events, real
commit landed, final result block unchanged. 13 new tests
(test_adapter.py + 2 in test_runner.py), all passing.

Co-Authored-By: Claude Sonnet 5 <noreply@anthropic.com>
This commit is contained in:
tegwick 2026-07-26 14:49:36 +02:00
parent c77a643393
commit 6c4f00e2c2
7 changed files with 384 additions and 18 deletions

View file

@ -9,12 +9,25 @@ LLMAdapter interface so a hosted adapter can be swapped in later, and adds:
- cwd pinned to the target repo
- --permission-mode acceptEdits
- allow-list from a named tool profile (default: green-commit-only)
- optional real-time per-tool-call audit events (HARNESS-WP-0002-T03)
Claude Code executes its own tools internally in `--print` mode — it is
not possible for a caller to externally dispatch individual tool calls
(that would require abandoning Claude Code's self-contained agent model
entirely). What `--output-format stream-json --include-hook-events` does
allow: observing each tool_use/tool_result/hook event as it happens. When
`on_tool_event` is supplied, this adapter runs in that streaming mode and
invokes the callback once per event, in real time, while still returning
one aggregate `LLMResponse` at the end for interface compatibility.
"""
from __future__ import annotations
import json
import subprocess
import threading
from pathlib import Path
from typing import Any, Callable
from llm_connect.claude_code import ClaudeCodeAdapter
from llm_connect.exceptions import LLMSubprocessError, LLMTimeoutError
@ -25,6 +38,24 @@ from rein_aharness.profiles import ToolProfile, get_profile
# Backward-compatible alias for the seed profile allow-list string.
ALLOWED_TOOLS = get_profile("green-commit-only").allowed_tools
ToolEventCallback = Callable[[dict[str, Any]], None]
def _is_tool_event(event: dict[str, Any]) -> bool:
"""True for tool_use/tool_result content blocks and hook lifecycle events.
Deliberately excludes plain assistant text messages — those aren't
tool audit events, just conversational output.
"""
event_type = event.get("type")
if event_type == "system" and str(event.get("subtype", "")).startswith("hook_"):
return True
if event_type in ("assistant", "user"):
for block in event.get("message", {}).get("content", []) or []:
if isinstance(block, dict) and block.get("type") in ("tool_use", "tool_result"):
return True
return False
class AgenticClaudeCodeAdapter(ClaudeCodeAdapter):
def __init__(
@ -32,10 +63,12 @@ class AgenticClaudeCodeAdapter(ClaudeCodeAdapter):
workdir: Path,
*,
tool_profile: str | ToolProfile = "green-commit-only",
on_tool_event: ToolEventCallback | None = None,
**kwargs,
):
super().__init__(**kwargs)
self._workdir = workdir
self._on_tool_event = on_tool_event
if isinstance(tool_profile, ToolProfile):
self._profile = tool_profile
else:
@ -54,6 +87,8 @@ class AgenticClaudeCodeAdapter(ClaudeCodeAdapter):
"--allowedTools",
self._profile.allowed_tools,
]
if self._on_tool_event is not None:
cmd += ["--output-format", "stream-json", "--include-hook-events", "--verbose"]
if self._model:
cmd.extend(["--model", self._model])
return cmd
@ -62,6 +97,14 @@ class AgenticClaudeCodeAdapter(ClaudeCodeAdapter):
self._preflight_budget(config)
cmd = self._build_command(config)
timeout = config.timeout_seconds or self._config.timeout_seconds
if self._on_tool_event is not None:
response = self._execute_streaming(cmd, prompt, timeout)
else:
response = self._execute_blocking(cmd, prompt, timeout)
self._consume_budget(config, response)
return response
def _execute_blocking(self, cmd: list[str], prompt: str, timeout: int) -> LLMResponse:
try:
result = subprocess.run(
cmd,
@ -81,7 +124,7 @@ class AgenticClaudeCodeAdapter(ClaudeCodeAdapter):
return_code=result.returncode,
stderr=result.stderr,
)
response = LLMResponse(
return LLMResponse(
content=result.stdout,
model=self._model or "claude-code-cli",
usage={},
@ -93,5 +136,74 @@ class AgenticClaudeCodeAdapter(ClaudeCodeAdapter):
"tool_profile": self._profile.name,
},
)
self._consume_budget(config, response)
return response
def _execute_streaming(self, cmd: list[str], prompt: str, timeout: int) -> LLMResponse:
proc = subprocess.Popen(
cmd,
stdin=subprocess.PIPE,
stdout=subprocess.PIPE,
stderr=subprocess.PIPE,
text=True,
cwd=self._workdir,
)
text_parts: list[str] = []
tool_event_count = 0
def reader() -> None:
nonlocal tool_event_count
assert proc.stdout is not None
for line in proc.stdout:
line = line.strip()
if not line:
continue
try:
event = json.loads(line)
except json.JSONDecodeError:
continue
self._handle_stream_event(event, text_parts)
if _is_tool_event(event):
tool_event_count += 1
self._on_tool_event(event)
reader_thread = threading.Thread(target=reader, daemon=True)
assert proc.stdin is not None
proc.stdin.write(prompt)
proc.stdin.close()
reader_thread.start()
try:
returncode = proc.wait(timeout=timeout)
except subprocess.TimeoutExpired as exc:
proc.kill()
proc.wait()
raise LLMTimeoutError(f"claude CLI timed out after {timeout}s", cause=exc) from exc
reader_thread.join(timeout=5)
stderr = proc.stderr.read() if proc.stderr else ""
if returncode != 0:
raise LLMSubprocessError(
f"claude CLI exited with code {returncode}",
return_code=returncode,
stderr=stderr,
)
return LLMResponse(
content="".join(text_parts),
model=self._model or "claude-code-cli",
usage={},
finish_reason="stop",
metadata={
"provider": "claude-code-agentic",
"cli_path": self._cli_path,
"workdir": str(self._workdir),
"tool_profile": self._profile.name,
"tool_event_count": tool_event_count,
},
)
@staticmethod
def _handle_stream_event(event: dict[str, Any], text_parts: list[str]) -> None:
if event.get("type") != "assistant":
return
for block in event.get("message", {}).get("content", []) or []:
if isinstance(block, dict) and block.get("type") == "text":
text_parts.append(block["text"])

View file

@ -130,10 +130,18 @@ def _cmd_run(args: argparse.Namespace) -> int:
print(f"invalid task spec: {exc}", file=sys.stderr)
return 2
def _print_event(event: dict) -> None:
# Tagged + single-line so a consumer reading stdout line-by-line
# (e.g. glas-harness's ReinAharness) can tell an event line apart
# from the pretty-printed final result block below.
print(json.dumps({"stream_event": event}), flush=True)
result = run_task(
spec,
report_to_hub=not args.no_hub,
write_metrics=not args.no_metrics,
emit_tool_events=args.stream_tool_events,
on_tool_event=_print_event if args.stream_tool_events else None,
)
closed = False
@ -191,6 +199,16 @@ def main(argv: list[str] | None = None) -> int:
action="store_true",
help="Skip writing .kaizen/metrics in the target repo",
)
run.add_argument(
"--stream-tool-events",
action="store_true",
help=(
"Run claude with --output-format stream-json --include-hook-events "
"and print each tool_use/tool_result/hook event as its own JSON "
"line while running (real-time audit, not external tool dispatch "
"-- HARNESS-WP-0002-T03)"
),
)
scan = sub.add_parser(
"mail-scan", help="Deterministic company-mailbox scan (no LLM session)"

View file

@ -12,8 +12,9 @@ from __future__ import annotations
import subprocess
import time
from dataclasses import dataclass
from dataclasses import dataclass, field
from pathlib import Path
from typing import Any, Callable
from rein_aharness import hub, metrics
from rein_aharness.manifest import resolve_run_policy
@ -53,6 +54,10 @@ class RunResult:
budget_tokens: int | None = None
tokens_spent: int | None = None
execution_time_s: float = 0.0
# Real-time per-tool-call audit events, populated only when run_task is
# called with emit_tool_events=True. See adapter.py's module docstring
# for why this is observation, not external tool dispatch.
tool_events: list[dict[str, Any]] = field(default_factory=list)
def _git(repo: Path, *args: str) -> str:
@ -72,6 +77,8 @@ def run_task(
adapter=None,
report_to_hub: bool = True,
write_metrics: bool = True,
emit_tool_events: bool = False,
on_tool_event: Callable[[dict[str, Any]], None] | None = None,
) -> RunResult:
try:
profile_name, budget_tokens, lane, blueprint = resolve_run_policy(
@ -103,12 +110,27 @@ def run_task(
budget_tokens=None,
)
collected_events: list[dict[str, Any]] = []
def _on_event(event: dict[str, Any]) -> None:
collected_events.append(event)
if report_to_hub:
hub.post_progress_event(
summary=f"tool event: {spec.title}",
event_type="tool_call",
detail={"repo": spec.target_repo.name, "agent": spec.agent, "event": event},
task_id=spec.hub_task_id,
)
if on_tool_event is not None:
on_tool_event(event)
if adapter is None:
from rein_aharness.adapter import AgenticClaudeCodeAdapter
adapter = AgenticClaudeCodeAdapter(
workdir=spec.target_repo,
tool_profile=profile,
on_tool_event=_on_event if (emit_tool_events or on_tool_event) else None,
)
head_before = _git(spec.target_repo, "rev-parse", "HEAD")
@ -162,6 +184,7 @@ def run_task(
budget_tokens=budget_tokens,
tokens_spent=tokens_spent,
execution_time_s=execution_time_s,
tool_events=collected_events,
)
if write_metrics: