From 3446010e932e0e899e7f937300a1cb02b291a68f Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 09:12:38 +0100 Subject: [PATCH 01/13] chore: start work on #79 From a00de1a05be06108d0f8870010da8af54e9fed6c Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:06:24 +0100 Subject: [PATCH 02/13] docs(specs): add tool-use orchestration loop design and plan --- ...9-25-tool-use-orchestration-loop-design.md | 329 +++++ .../2026-09-25-tool-use-orchestration-loop.md | 1302 +++++++++++++++++ 2 files changed, 1631 insertions(+) create mode 100644 .agents/specs/2026-09-25-tool-use-orchestration-loop-design.md create mode 100644 .agents/specs/2026-09-25-tool-use-orchestration-loop.md diff --git a/.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md b/.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md new file mode 100644 index 0000000..911510d --- /dev/null +++ b/.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md @@ -0,0 +1,329 @@ +# Design: Tool-Use Orchestration Loop — reason → act → observe + +- **Issue:** #79 — feat(agent): implement tool-use orchestration loop +- **Branch:** `feature/tool-use-orchestration-loop` +- **Draft PR:** [#103](https://github.com/Svagtlys/Octave/pull/103) +- **Date:** 2026-09-25 +- **Status:** Approved + +## Problem + +Issue #79 asks for the loop that turns a tool-capable LLM into an agent: +detect `finish_reason == "tool_calls"` in completions, execute the requested +calls against MCP servers, feed results back as tool messages, and re-invoke +the model for final synthesis. Every prerequisite seam has shipped and each +explicitly deferred its half of this bridge: + +| Capability | Status | Evidence | +|---|---|---| +| `CompletionRequest.tools` + request-side envelope | Done | `inference/types.py`, `OpenAIAdapter._chat_kwargs`, #78 | +| MCP tool inventory + `(server_id, tool_name)` call surface | Done | `ToolRegistry.call_tool`, #77 | +| Exposed-name → route reverse map (`ProviderToolset.routes`) | Done | `octave.tools.translate`, #78 | +| `finish_reason` surfaced on `CompletionResult` | Done | `inference/types.py` | +| `CompletionResult.tool_calls` parsing | **Gap** | `OpenAIAdapter.complete` drops `choice.message.tool_calls` | +| Assistant message carrying `tool_calls` | **Gap** | `Message` has no tool vocabulary | +| `role="tool"` result messages | **Gap** | `Role` lacks `"tool"`; no `tool_call_id` | +| Tool-call → MCP `tools/call` mapping + result wrapping | **Gap** | No composition layer exists | +| Multi-tool-call turns, error surfacing to the LLM | **Gap** | No loop anywhere | + +The #78 spec recorded this work item as the consumer that composes +`inventory() → translate_tools()` and named "response-side `tool_calls` +parsing" as its scope. This design ships the complete loop as a library +primitive. + +## Decisions from brainstorming + +| # | Topic | Decision | +|---|-------|----------| +| 1 | Scope | **Library-only.** Loop primitive + inference adapter changes + tests. Collaborators injected; no routes, no lifespan wiring, no session persistence — the Integration & Testing work item composes it with `ToolRegistry` + lifespans. Matches how #8/#15/#77/#78 shipped seams without consumers. | +| 2 | Loop termination | **`max_tool_rounds: int = 8` constructor option; hitting the limit raises `ToolLoopLimitError`** carrying `.messages` (partial transcript). Silent truncation is a footgun: the truncated result's `finish_reason == "tool_calls"` is indistinguishable from a clean answer without disciplined inspection. "Loud and deterministic beats silent" (#78 decision 4 precedent). | +| 3 | Error taxonomy | **Two-tier split.** Model-correctable failures become `is_error` tool messages and the loop continues; orchestration-fatal conditions raise. `ToolError(Exception)` added as package base in `octave.tools.errors`; `ToolTranslationError` re-based under it (backward-compatible); `ToolLoopLimitError(ToolError)` new. | +| 4 | Message model | **OpenAI-dialect minimal (additive).** `ToolCall(id, name, arguments: str)`; `Message` gains `tool_calls`/`tool_call_id`/`name`, `Role` gains `"tool"`; `CompletionResult` gains `tool_calls`. The whole stack (Ollama/vLLM/llama.cpp/LM Studio) speaks this wire shape; Anthropic-style content blocks would need re-translation back to it and break every existing `Message` caller. | +| 5 | Arguments representation | **Raw JSON string on `ToolCall.arguments`**, exactly as the wire provides — mirrors the `ToolInfo.input_schema` "never interpret" posture. The loop parses to a dict for the executor; malformed JSON is model-correctable → `is_error` tool message, never an exception. Keeps the adapter a lossless translator. | +| 6 | Package placement | **New `octave/agent/` package** — the composition layer and future Agent Manager home (`octave.agent.loop`, later `.manager`/`.registry`/`.router`). `octave.tools` stays pure (no I/O, `test_package.py` AST guard untouched); `octave.agent` is the only package importing both `octave.mcp` and `octave.inference`. Package docstring states which Agent-Manager components exist vs. are planned. | +| 7 | MCP access | **`ToolExecutor` Protocol + `McpToolExecutor`.** The loop depends only on the Protocol (`async call(server_id, tool_name, arguments) -> ToolOutcome`, never raises for tool-level failures); `McpToolExecutor` wraps `ToolRegistry` and owns `McpError → ToolOutcome(is_error=True)` conversion. Testability without a live fleet; the integration item swaps in richer executors (tagging, retries) without touching the loop. | +| 8 | Toolset input | **Loop consumes a pre-built `ProviderToolset`.** Honors #78's "the consumer composes `inventory() → translate_tools()` itself"; no new registry→translator glue service (waits for a second consumer). | +| 9 | Execution order | **Sequential within a turn, in model call order** (deterministic history). Concurrent fan-out deferred — follow-up issue drafted at [`plans/2026-09-25-issue-concurrent-tool-execution.md`](../../plans/2026-09-25-issue-concurrent-tool-execution.md); protocol invariant (all N results before next completion) is identical either way. | +| 10 | Streaming | **Deferred non-goal.** `CompletionChunk` unchanged; delta accumulation of `tool_calls` across stream chunks is its own concern (arrives with the streaming UI work item). | + +### Rejected alternatives + +- **Content-block message model (Anthropic-style).** Engine-portable in + theory; in practice every target engine speaks the OpenAI tool-calling wire + shape, so the adapter would translate blocks → OpenAI payloads anyway. + Breaking change to every `Message` construction site for zero present gain. +- **Side-channel tool-call structure (keep `Message` frozen).** Provider + history rules require assistant-with-`tool_calls` messages *in* the history + immediately followed by their results; splitting one ordered transcript + across two structures forces sync invariants on every consumer. +- **Loop inside `octave.tools`.** Dilutes the package's documented pure + functions / no-I/O contract (#78 decision 5) and drags `octave.mcp` + imports into a package whose test guards purity via AST scan. +- **Async generator loop (yield events).** Event streaming is a streaming/UI + concern; deferred with it. The class-based primitive leaves room to add it + later without signature churn on the terminal result. +- **Silent stop at round limit.** Returned transcript would be indistinguishable + from a clean final answer for callers that don't inspect `finish_reason`. +- **Raising on tool failures.** Defeats the loop's purpose: the model must see + failures as tool results to retry, rephrase, or give up gracefully. + +## Architecture + +```mermaid +graph LR + A[Caller future Integration item] -->|messages + ProviderToolset| B[octave.agent.loop.ToolLoop] + B -->|complete with toolset.tools| C[InferenceAdapter seam] + C -->|CompletionResult.tool_calls| B + B -->|routes lookup server_id + tool_name| D[ToolExecutor Protocol] + D -.->|implements| E[McpToolExecutor octave.agent.mcp_executor] + E -->|call_tool| F[ToolRegistry octave.mcp] + G[octave.tools.translate_tools] -->|ProviderToolset| A + B -->|ToolTurn| A +``` + +Dependency posture: + +- `octave.agent` imports `octave.inference` (types + `InferenceAdapter` ABC) + and `octave.tools` (types) — and `octave.mcp` **only** inside + `mcp_executor.py`. No SDK imports anywhere in the package (new + `tests/agent/test_package.py` AST guard, mirroring `tests/tools/test_package.py`). +- `octave.tools` unchanged except `errors.py` gaining the `ToolError` base. +- `octave.inference` changes are additive types + `OpenAIAdapter` response + parsing; the SDK quarantine rule is untouched. +- Responses-API compatibility: a future adapter maps `function_call` / + `function_call_output` wire items to/from the same Octave types internally; + nothing above the adapter seam changes. + +## Data model & API + +```python +# octave/inference/types.py (additive) +class ToolCall(BaseModel): + """One tool call as the model requested it. `arguments` is the raw JSON + string from the wire; Octave never interprets it (adapter stays lossless).""" + + id: str + name: str + arguments: str + +Role = Literal["system", "user", "assistant", "tool"] + +class Message(BaseModel): + role: Role + content: str + tool_calls: list[ToolCall] | None = None # assistant messages + tool_call_id: str | None = None # tool messages + name: str | None = None # tool messages: exposed tool name + +class CompletionResult(BaseModel): + ... + tool_calls: list[ToolCall] | None = None + +# octave/tools/errors.py (additive restructure) +class ToolError(Exception): + """Base for tool-plane failures (translation and orchestration).""" + +class ToolTranslationError(ToolError): ... # existing, re-based + +# octave/agent/errors.py +class ToolLoopLimitError(ToolError): + """max_tool_rounds exhausted without a final answer.""" + + def __init__(self, message: str, *, messages: list[Message]) -> None: + super().__init__(message) + self.messages = messages + +# octave/agent/types.py +class ToolOutcome(BaseModel): + """Flattened result of one tool execution.""" + + content: str + is_error: bool = False + +class ToolTurn(BaseModel): + """Result of one orchestrated turn.""" + + messages: list[Message] + """Input + appended assistant/tool messages + final assistant message.""" + + result: CompletionResult + """The final (no tool_calls) completion.""" + + tool_rounds: int + """Tool-execution rounds performed.""" + +# octave/agent/executor.py +class ToolExecutor(Protocol): + """Executes one resolved tool call. Must not raise for tool-level + failures — return ToolOutcome(is_error=True) instead.""" + + async def call( + self, server_id: str, tool_name: str, arguments: dict[str, Any] + ) -> ToolOutcome: ... + +# octave/agent/mcp_executor.py +class McpToolExecutor: + """ToolExecutor over ToolRegistry. Joins text content blocks with + newlines; empty content becomes '(no output)'. Catches McpError + (Rpc/Timeout/Connection/Config/NotConnected) -> ToolOutcome(is_error=True).""" + + def __init__(self, registry: ToolRegistry) -> None: ... + +# octave/agent/loop.py +class ToolLoop: + def __init__( + self, + *, + adapter: InferenceAdapter, + executor: ToolExecutor, + max_tool_rounds: int = 8, + ) -> None: ... + + async def run( + self, + messages: list[Message], + toolset: ProviderToolset, + *, + model: str | None = None, + ) -> ToolTurn: ... +``` + +## Loop mechanics + +```mermaid +sequenceDiagram + participant C as Caller + participant L as ToolLoop + participant A as InferenceAdapter + participant E as ToolExecutor + C->>L: run messages + toolset + L->>A: complete messages + toolset.tools + A-->>L: CompletionResult tool_calls non-empty + L->>L: append assistant message text + tool_calls + loop each tool call in model order + L->>L: toolset.routes lookup exposed name + alt no route or malformed arguments JSON + L->>L: append is_error tool message + else + L->>E: call server_id tool_name parsed arguments + E-->>L: ToolOutcome + L->>L: append tool message keyed by tool_call_id + end + end + L->>A: re-invoke with extended history + A-->>L: final CompletionResult tool_calls empty + L-->>C: ToolTurn messages + result + rounds +``` + +- **Trigger:** a tool round runs iff `result.tool_calls` is non-empty. + Content is the signal; `finish_reason` is advisory — some local engines + emit `"stop"` alongside populated tool calls, and engines that set + `"tool_calls"` always populate the field. `finish_reason == "tool_calls"` + with an empty/absent list is treated as final (engine bug, not a loop). +- **History order:** assistant message (with `tool_calls`) immediately + followed by exactly one `role="tool"` message per call, keyed by + `tool_call_id`, carrying `name` = exposed name. This is the provider + invariant; the loop constructs it unconditionally. +- **Round accounting:** `max_tool_rounds` bounds tool-execution rounds; the + synthesis completion always follows the last executed round, so a limit of + N permits at most N executions. `ToolTurn.tool_rounds` reports the count. +- **Tools on every request:** `toolset.tools` is sent each round, enabling + multi-round chains. +- **Returned transcript:** `ToolTurn.messages` is the complete turn + (input + everything appended, including the final assistant message), + ready for later persistence by the transcript/session work item. +- **`tool_call_id` provenance:** always the model's own id — echoed exactly, + never generated or rewritten (strict engines validate the pairing). + +## Error handling — the two-tier split + +**Tier 1 — model-correctable → `is_error` tool message, loop continues:** + +| Condition | Detected by | Tool message content | +|---|---|---| +| `arguments` not valid JSON object | loop (parse step) | `Tool call arguments are not a valid JSON object: ` | +| Exposed name not in `routes` (hallucinated call) | loop (route lookup) | `Unknown tool: ` | +| MCP-level tool failure (`ToolResult.is_error=True`) | `McpToolExecutor` | joined content, `is_error=True` | +| `McpRpcError` / `McpTimeoutError` / `McpConnectionError` / `McpConfigError` / `McpNotConnectedError` | `McpToolExecutor` | `Tool execution failed: ` | +| Empty tool output | `McpToolExecutor` | `(no output)` sentinel (some engines reject empty content) | + +**Tier 2 — orchestration-fatal → raise:** + +- `ToolLoopLimitError` with `.messages` — caller can persist the partial + transcript, show partial work, or retry. +- `AdapterError` (or subclass) from the inference call propagates untouched — + the adapter's contract; no retries here (adapter `max_retries` owns that). + +Rationale: the model can recover from Tier 1 by retrying the call, changing +tools, or answering around the failure — it must see failures *as tool +results*. Tier 2 conditions are invisible or meaningless to the model. + +## Scope reconciliation & deferrals + +| Deferred | Trigger / home | +|---|---| +| Streaming `tool_calls` delta accumulation (`CompletionChunk`) | Streaming chat work item (UI Chat #3) | +| Concurrent execution of multiple calls in one turn | Drafted issue: [`plans/2026-09-25-issue-concurrent-tool-execution.md`](../../plans/2026-09-25-issue-concurrent-tool-execution.md) | +| Registry→translator→loop convenience glue + lifespan/deps wiring | Integration & Testing #1 | +| Session/event persistence of `ToolTurn.messages` | Transcript persistence work item (sessions/events schema exists) | +| Agent Manager (lifecycle, registry, router) | Agent Manager #1–6 — same `octave.agent` package, later modules | +| `tool_choice` / `parallel_tool_calls` first-class fields | Ride `CompletionRequest.extra` until a caller needs them (#78 precedent) | +| Tool rename/re-describe overrides, tagging | TODO MCP Connector #8–10 — compose over translator input | +| Progress/observability events (UI shows in-flight tool calls) | Arrives with streaming; orthogonal to model history | + +## Testing + +**Inference — `tests/inference/test_types.py`:** `ToolCall` round-trip; +`Message` with/without tool fields defaults `None`; `CompletionResult.tool_calls` +default `None`; `role="tool"` validates. + +**Inference — `tests/inference/test_openai_adapter.py`** (`httpx.MockTransport`): + +1. Response carrying `choice.message.tool_calls` → parsed to + `CompletionResult.tool_calls` with id/name/arguments verbatim. +2. Response without tool calls → `tool_calls is None`. +3. Request with assistant `tool_calls` message → body carries + `tool_calls: [{id, type: "function", function: {name, arguments}}]`. +4. Request with `role="tool"` message → body carries `tool_call_id`, `name`. +5. Plain messages → body has **no** `tool_calls`/`tool_call_id`/`name` null + keys (`exclude_none=True` regression guard). +6. `arguments` passed through as opaque string (no re-serialization). + +**Agent loop — `tests/agent/test_loop.py`** (scripted fake adapter queue + +recording fake executor): + +7. No tool calls first round → passthrough; `tool_rounds == 0`; messages = + input + final assistant. +8. One round, multiple calls → assistant + N tool messages in call order; + second completion returns final; history ordering asserted exactly. +9. Multi-round chain (tool → tool → final) → two rounds, ordering across + rounds. +10. Limit exhausted → `max_tool_rounds = N`: the first N completions return + tool_calls (all executed); the (N+1)th synthesis completion also returns + tool_calls → `ToolLoopLimitError`; `.messages` ends on that unfulfilled + assistant message (no tool messages appended for it); total adapter calls + exactly N+1. +11. `finish_reason="tool_calls"` but empty list → treated as final. +12. Malformed arguments JSON → `is_error` tool message with the call's id; + executor not invoked for that call; other calls in the turn proceed. +13. Unknown exposed name → `is_error` tool message; executor not invoked. +14. Executor `is_error=True` outcome → tool message `is_error=True`, loop + continues. +15. `AdapterError` from any round propagates untouched. +16. `tool_call_id` echo: each tool message's id matches its call's id. + +**MCP executor — `tests/agent/test_mcp_executor.py`** (fake/stubbed +`ToolRegistry`): + +17. Text blocks joined with newlines; `is_error` passthrough. +18. Empty content → `"(no output)"`. +19. Each `McpError` subclass → `ToolOutcome(is_error=True)` with message + included; nothing raised. +20. Arguments dict forwarded to `registry.call_tool` verbatim. + +**Package hygiene — `tests/agent/test_package.py`:** AST guard — `octave.agent` +imports neither `openai` nor `mcp` SDK; re-export surface matches `__all__`. + +**Translator regression:** existing `tests/tools/*` green with `ToolError` +re-base (`ToolTranslationError` still catchable as before; `isinstance` +chain asserted). diff --git a/.agents/specs/2026-09-25-tool-use-orchestration-loop.md b/.agents/specs/2026-09-25-tool-use-orchestration-loop.md new file mode 100644 index 0000000..ca0993c --- /dev/null +++ b/.agents/specs/2026-09-25-tool-use-orchestration-loop.md @@ -0,0 +1,1302 @@ +# Tool-Use Orchestration Loop Implementation Plan + +> **For agentic workers:** REQUIRED SUB-SKILL: Use superpowers:subagent-driven-development (recommended) or superpowers:executing-plans to implement this plan task-by-task. Steps use checkbox (`- [ ]`) syntax for tracking. + +**Goal:** Implement the reason → act → observe loop: detect model tool calls, execute them against MCP servers, feed results back as provider-native tool messages, and re-invoke for final synthesis. + +**Architecture:** Additive OpenAI-dialect extensions to `octave.inference` (types + adapter parsing/serialization), then a new composition package `octave.agent` holding the loop (`ToolLoop`), a `ToolExecutor` Protocol, and `McpToolExecutor` wrapping `ToolRegistry`. `octave.tools` stays pure; no SDK imports cross the `octave.agent` boundary; no routes/lifespan wiring (library-only per spec). + +**Tech Stack:** Python 3.12, Pydantic v2, anyio, pytest (`asyncio_mode=auto`), `httpx.MockTransport`, uv. + +**Spec:** [`.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md`](./2026-09-25-tool-use-orchestration-loop-design.md) — issue #79, branch `feature/tool-use-orchestration-loop`, draft PR [#103](https://github.com/Svagtlys/Octave/pull/103). + +**Conventions for all tasks:** +- Work from `backend/` (all paths below relative to repo root; run commands from `backend/`). +- Test command form: `uv run pytest -v`; lint: `uv run ruff check src tests`; types: `uv run mypy src`. +- Commit messages: `type(scope): description`. +- The design doc + this plan are committed first (Task 1) — they are already on disk. +- **Error-content convention (spec clarification):** OpenAI-dialect tool messages carry no error flag; the loop conveys Tier-1 failures by prefixing content with `Error: `. All loop tests assert this shape. + +--- + +### Task 1: Commit design and plan documents + +**Files:** +- Commit (already exist): `.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md`, `.agents/specs/2026-09-25-tool-use-orchestration-loop.md` + +- [ ] **Step 1: Verify you are on the feature branch** + +Run: `git branch --show-current` +Expected: `feature/tool-use-orchestration-loop` (if not: `git checkout feature/tool-use-orchestration-loop`) + +- [ ] **Step 2: Commit the docs** + +```bash +git add .agents/specs/2026-09-25-tool-use-orchestration-loop-design.md .agents/specs/2026-09-25-tool-use-orchestration-loop.md +git commit -m "docs(specs): add tool-use orchestration loop design and plan" +``` + +--- + +### Task 2: Inference domain types — `ToolCall`, tool fields on `Message`, `CompletionResult.tool_calls` + +**Files:** +- Modify: `backend/src/octave/inference/types.py` +- Modify: `backend/src/octave/inference/__init__.py` +- Modify: `backend/tests/inference/test_types.py` + +- [ ] **Step 1: Write the failing tests** — append to `backend/tests/inference/test_types.py` + +Add `ToolCall` and `CompletionResult` to the existing import from `octave.inference.types` (lines 6–13), then append: + +```python +def test_tool_role_is_valid() -> None: + message = Message(role="tool", content="result", tool_call_id="call_1", name="mcp__fs__read") + assert message.tool_calls is None + + +def test_message_tool_fields_default_none() -> None: + message = Message(role="assistant", content="hi") + assert message.tool_calls is None + assert message.tool_call_id is None + assert message.name is None + + +def test_message_rejects_truly_unknown_role() -> None: + with pytest.raises(ValidationError): + Message(role="developer", content="hi") # type: ignore[arg-type] + + +def test_tool_call_round_trip() -> None: + call = ToolCall(id="call_1", name="mcp__fs__read", arguments='{"path": "/tmp/x"}') + assert ToolCall.model_validate(call.model_dump()) == call + + +def test_completion_result_tool_calls_defaults_none() -> None: + result = CompletionResult(text="hi", model="m") + assert result.tool_calls is None + + +def test_completion_result_round_trips_tool_calls() -> None: + call = ToolCall(id="c1", name="x", arguments="{}") + result = CompletionResult(text="", model="m", finish_reason="tool_calls", tool_calls=[call]) + assert CompletionResult.model_validate(result.model_dump()) == result +``` + +Also **replace** the existing test at lines 35–37 (`test_message_rejects_unknown_role` — it asserts `role="tool"` raises, which is now wrong): + +```python +# DELETE this test entirely (superseded by test_message_rejects_truly_unknown_role +# and test_tool_role_is_valid above): +# def test_message_rejects_unknown_role() -> None: +# with pytest.raises(ValidationError): +# Message(role="tool", content="hi") # type: ignore[arg-type] +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run pytest tests/inference/test_types.py -v` +Expected: FAIL — `ImportError: cannot import name 'ToolCall'` + +- [ ] **Step 3: Implement the type additions** in `backend/src/octave/inference/types.py` + +Replace `Role` (line 24) and `Message` (lines 27–31) and extend `CompletionResult` (lines 70–76): + +```python +Role = Literal["system", "user", "assistant", "tool"] + + +class ToolCall(BaseModel): + """One tool call as the model requested it. + + ``arguments`` is the raw JSON string from the wire; Octave never + interprets it (the adapter stays a lossless translator, the loop parses). + """ + + id: str + name: str + arguments: str + + +class Message(BaseModel): + """A single chat message. + + ``tool_calls`` populates assistant messages; ``tool_call_id`` and ``name`` + populate ``role="tool"`` result messages (keyed to the call they answer). + """ + + role: Role + content: str + tool_calls: list[ToolCall] | None = None + tool_call_id: str | None = None + name: str | None = None +``` + +```python +class CompletionResult(BaseModel): + """A full (non-streamed) chat completion.""" + + text: str + model: str + finish_reason: str | None = None + usage: Usage | None = None + tool_calls: list[ToolCall] | None = None +``` + +Add `"ToolCall"` to `__all__` (lines 12–22, alphabetical position before `"ToolDefinition"`). + +In `backend/src/octave/inference/__init__.py`: if `ToolDefinition` is re-exported there (import + `__all__`), add `ToolCall` alongside it in both places; if it is not, skip — types are importable from `octave.inference.types` directly. + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run pytest tests/inference/test_types.py -v` +Expected: all PASS (including conformance-adjacent tests) + +- [ ] **Step 5: Full inference suite still green** + +Run: `uv run pytest tests/inference -v` +Expected: all PASS + +- [ ] **Step 6: Commit** + +```bash +git add backend/src/octave/inference/types.py backend/src/octave/inference/__init__.py backend/tests/inference/test_types.py +git commit -m "feat(inference): add tool-call vocabulary to domain types" +``` + +--- + +### Task 3: `OpenAIAdapter` — parse `tool_calls` from completion responses + +**Files:** +- Modify: `backend/src/octave/inference/openai_adapter.py` (`complete`, lines 89–101) +- Modify: `backend/tests/inference/test_openai_adapter.py` + +- [ ] **Step 1: Write the failing tests** — append to `backend/tests/inference/test_openai_adapter.py` + +Add `ToolCall` to the import from `octave.inference.types` (line 19). Append: + +```python +CHAT_RESPONSE_TOOL_CALLS = { + "id": "chatcmpl-tc1", + "object": "chat.completion", + "created": 1725000000, + "model": "fake-model", + "choices": [ + { + "index": 0, + "message": { + "role": "assistant", + "content": None, + "tool_calls": [ + { + "id": "call_1", + "type": "function", + "function": { + "name": "mcp__fs__read", + "arguments": '{"path": "/tmp/x"}', + }, + }, + { + "id": "call_2", + "type": "function", + "function": {"name": "mcp__fs__list", "arguments": "{}"}, + }, + ], + }, + "finish_reason": "tool_calls", + } + ], + "usage": {"prompt_tokens": 9, "completion_tokens": 12, "total_tokens": 21}, +} + + +async def test_complete_parses_tool_calls() -> None: + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json=CHAT_RESPONSE_TOOL_CALLS) + + client = _mock_client(handler) + adapter = _adapter(client) + result = await adapter.complete( + CompletionRequest(model=None, messages=[Message(role="user", content="hi")]) + ) + assert result.finish_reason == "tool_calls" + assert result.text == "" # None content normalized to "" + assert result.tool_calls == [ + ToolCall(id="call_1", name="mcp__fs__read", arguments='{"path": "/tmp/x"}'), + ToolCall(id="call_2", name="mcp__fs__list", arguments="{}"), + ] + await adapter.aclose() + await client.aclose() + + +async def test_complete_without_tool_calls_is_none() -> None: + client = _mock_client() + adapter = _adapter(client) + result = await adapter.complete( + CompletionRequest(model=None, messages=[Message(role="user", content="hi")]) + ) + assert result.tool_calls is None + await adapter.aclose() + await client.aclose() +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run pytest tests/inference/test_openai_adapter.py -k tool_calls -v` +Expected: FAIL — `result.tool_calls` is `None` (field not populated) + +- [ ] **Step 3: Implement response parsing** in `backend/src/octave/inference/openai_adapter.py` + +Add `ToolCall` to the import from `octave.inference.types` (lines 29–37). Replace `complete` (lines 89–101): + +```python + async def complete(self, request: CompletionRequest) -> CompletionResult: + kwargs = self._chat_kwargs(request) + try: + response = await self._client.chat.completions.create(**kwargs) + except openai.APIError as exc: + raise _translate(exc) from exc + choice = response.choices[0] + tool_calls = [ + ToolCall( + id=call.id, + name=call.function.name, + arguments=call.function.arguments, + ) + for call in (choice.message.tool_calls or []) + ] + return CompletionResult( + text=choice.message.content or "", + model=response.model, + finish_reason=choice.finish_reason, + usage=_chat_usage(response.usage), + tool_calls=tool_calls or None, + ) +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run pytest tests/inference/test_openai_adapter.py -k tool_calls -v` +Expected: PASS + +- [ ] **Step 5: Commit** + +```bash +git add backend/src/octave/inference/openai_adapter.py backend/tests/inference/test_openai_adapter.py +git commit -m "feat(inference): parse tool_calls from completion responses" +``` + +--- + +### Task 4: `OpenAIAdapter` — serialize tool-carrying messages to provider payloads + +**Files:** +- Modify: `backend/src/octave/inference/openai_adapter.py` (`_chat_kwargs`, lines 158–178) +- Modify: `backend/tests/inference/test_openai_adapter.py` + +- [ ] **Step 1: Write the failing tests** — append to `backend/tests/inference/test_openai_adapter.py` + +Add `ToolCall` usage (already imported in Task 3). Append: + +```python +async def _capture_chat_body(messages: list[Message]) -> dict: + captured: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + captured.append(request) + return httpx.Response(200, json=CHAT_RESPONSE) + + client = _mock_client(handler) + adapter = _adapter(client) + await adapter.complete(CompletionRequest(model=None, messages=messages)) + await adapter.aclose() + await client.aclose() + return json.loads(captured[0].content) + + +async def test_plain_messages_emit_no_null_tool_fields() -> None: + body = await _capture_chat_body( + [ + Message(role="system", content="sys"), + Message(role="user", content="hi"), + Message(role="assistant", content="yo"), + ] + ) + for message in body["messages"]: + assert "tool_calls" not in message + assert "tool_call_id" not in message + assert "name" not in message + + +async def test_assistant_tool_calls_message_uses_provider_envelope() -> None: + body = await _capture_chat_body( + [ + Message(role="user", content="hi"), + Message( + role="assistant", + content="", + tool_calls=[ToolCall(id="call_1", name="mcp__fs__read", arguments='{"p": 1}')], + ), + ] + ) + assert body["messages"][1] == { + "role": "assistant", + "content": "", + "tool_calls": [ + { + "id": "call_1", + "type": "function", + "function": {"name": "mcp__fs__read", "arguments": '{"p": 1}'}, + } + ], + } + + +async def test_tool_result_message_carries_tool_call_id_and_name() -> None: + body = await _capture_chat_body( + [Message(role="tool", content="42", tool_call_id="call_1", name="mcp__fs__read")] + ) + assert body["messages"][0] == { + "role": "tool", + "content": "42", + "tool_call_id": "call_1", + "name": "mcp__fs__read", + } +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run pytest tests/inference/test_openai_adapter.py -k "null_tool_fields or provider_envelope or tool_call_id" -v` +Expected: FAIL — plain messages carry `"tool_calls": null` keys; assistant `tool_calls` serialize flat (no `{"type": "function", "function": ...}` envelope) + +- [ ] **Step 3: Implement message serialization** in `backend/src/octave/inference/openai_adapter.py` + +In `_chat_kwargs`, replace the `messages` entry (line 161): + +```python + "messages": [self._message_payload(message) for message in request.messages], +``` + +Add a method after `_chat_kwargs`: + +```python + def _message_payload(self, message: Message) -> dict[str, Any]: + """One message as a provider payload. exclude_none keeps tool fields + off plain messages (strict local engines reject null tool keys); + assistant tool_calls get the provider function-call envelope.""" + payload = message.model_dump(mode="json", exclude_none=True) + if message.tool_calls: + payload["tool_calls"] = [ + { + "id": call.id, + "type": "function", + "function": {"name": call.name, "arguments": call.arguments}, + } + for call in message.tool_calls + ] + return payload +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run pytest tests/inference/test_openai_adapter.py -v` +Expected: all PASS (existing tools-envelope and extra tests unaffected — `_handler` reads `body["messages"][-1]["content"]`, present on every plain message) + +- [ ] **Step 5: Commit** + +```bash +git add backend/src/octave/inference/openai_adapter.py backend/tests/inference/test_openai_adapter.py +git commit -m "feat(inference): serialize tool-carrying messages to provider payloads" +``` + +--- + +### Task 5: `octave.tools.errors` — introduce `ToolError` base + +**Files:** +- Modify: `backend/src/octave/tools/errors.py` +- Modify: `backend/tests/tools/test_package.py` + +- [ ] **Step 1: Write the failing test** — append to `backend/tests/tools/test_package.py` + +Update the import at line 6 usage by adding at top: `from octave.tools.errors import ToolError, ToolNameCollisionError, ToolTranslationError` and append: + +```python +def test_tool_error_hierarchy() -> None: + assert issubclass(ToolTranslationError, ToolError) + assert issubclass(ToolNameCollisionError, ToolTranslationError) +``` + +Also add `"ToolError"` to the name tuple in `test_public_names_are_exported` (lines 28–36). + +- [ ] **Step 2: Run test to verify it fails** + +Run: `uv run pytest tests/tools/test_package.py -v` +Expected: FAIL — `ImportError: cannot import name 'ToolError'` + +- [ ] **Step 3: Implement the restructure** — replace `backend/src/octave/tools/errors.py` entirely + +```python +"""Tool-plane errors (design specs #78, #79).""" + +__all__ = ["ToolError", "ToolNameCollisionError", "ToolTranslationError"] + + +class ToolError(Exception): + """Base for tool-plane failures (translation and orchestration).""" + + +class ToolTranslationError(ToolError): + """Base for translation failures.""" + + +class ToolNameCollisionError(ToolTranslationError): + """Two tools landed on the same exposed name after sanitization.""" +``` + +Add `ToolError` to the import and `__all__` in `backend/src/octave/tools/__init__.py`. + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run pytest tests/tools -v` +Expected: all PASS (re-base is backward-compatible: existing `except ToolTranslationError` sites catch identically) + +- [ ] **Step 5: Commit** + +```bash +git add backend/src/octave/tools/errors.py backend/src/octave/tools/__init__.py backend/tests/tools/test_package.py +git commit -m "refactor(tools): introduce ToolError base for the tool plane" +``` + +--- + +### Task 6: `octave.agent` package skeleton — types, errors, package guard + +**Files:** +- Create: `backend/src/octave/agent/__init__.py` +- Create: `backend/src/octave/agent/types.py` +- Create: `backend/src/octave/agent/errors.py` +- Create: `backend/tests/agent/__init__.py` (empty) +- Create: `backend/tests/agent/test_package.py` +- Create: `backend/tests/agent/test_types.py` + +- [ ] **Step 1: Write the failing tests** + +Create `backend/tests/agent/test_types.py`: + +```python +"""Agent loop vocabulary: defaults and provenance of ToolOutcome/ToolTurn.""" + +from octave.agent.types import ToolOutcome, ToolTurn +from octave.inference.types import CompletionResult, Message + + +def test_tool_outcome_defaults() -> None: + outcome = ToolOutcome(content="ok") + assert outcome.is_error is False + + +def test_tool_turn_round_trip() -> None: + turn = ToolTurn( + messages=[Message(role="user", content="hi"), Message(role="assistant", content="yo")], + result=CompletionResult(text="yo", model="m"), + tool_rounds=0, + ) + assert ToolTurn.model_validate(turn.model_dump()) == turn +``` + +Create `backend/tests/agent/test_package.py`: + +```python +"""Public API surface and SDK-quarantine posture of the agent package.""" + +import ast +from pathlib import Path + +import octave.agent as agent_pkg + +SDK_MODULES = {"openai", "mcp"} + + +def _imported_top_level_modules() -> set[str]: + names: set[str] = set() + for path in Path(agent_pkg.__file__).parent.rglob("*.py"): + tree = ast.parse(path.read_text()) + for node in ast.walk(tree): + if isinstance(node, ast.Import): + names.update(alias.name.split(".")[0] for alias in node.names) + elif isinstance(node, ast.ImportFrom) and node.module: + names.add(node.module.split(".")[0]) + return names + + +def test_no_sdk_imports() -> None: + """octave.agent composes Octave façades only, never the openai/mcp SDKs.""" + assert not _imported_top_level_modules() & SDK_MODULES + + +def test_public_names_are_exported() -> None: + for name in ( + "McpToolExecutor", + "ToolError", + "ToolExecutor", + "ToolLoop", + "ToolLoopLimitError", + "ToolOutcome", + "ToolTurn", + ): + assert hasattr(agent_pkg, name), name +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run pytest tests/agent -v` +Expected: FAIL — `ModuleNotFoundError: No module named 'octave.agent'` + +- [ ] **Step 3: Implement the package** + +Create `backend/src/octave/agent/types.py`: + +```python +"""Agent orchestration vocabulary (design spec #79).""" + +from pydantic import BaseModel + +from octave.inference.types import CompletionResult, Message + +__all__ = ["ToolOutcome", "ToolTurn"] + + +class ToolOutcome(BaseModel): + """Flattened result of one tool execution.""" + + content: str + is_error: bool = False + + +class ToolTurn(BaseModel): + """Result of one orchestrated turn.""" + + messages: list[Message] + """Input + appended assistant/tool messages + final assistant message.""" + + result: CompletionResult + """The final (no tool_calls) completion.""" + + tool_rounds: int + """Tool-execution rounds performed.""" +``` + +Create `backend/src/octave/agent/errors.py`: + +```python +"""Orchestration-fatal errors (design spec #79).""" + +from octave.inference.types import Message +from octave.tools.errors import ToolError + +__all__ = ["ToolLoopLimitError"] + + +class ToolLoopLimitError(ToolError): + """max_tool_rounds exhausted without a final answer. + + ``messages`` is the partial transcript (input + everything appended, + ending on the unfulfilled assistant tool_calls message) so callers can + persist partial work or retry. + """ + + def __init__(self, message: str, *, messages: list[Message]) -> None: + super().__init__(message) + self.messages = messages +``` + +Create `backend/src/octave/agent/__init__.py` (placeholder for Tasks 7–8; full surface now so the export test is honest — it will fail until Tasks 7–8 land, so temporarily this test is the only red one; run only `test_types.py` here): + +```python +"""Agent plane: tool-use orchestration loop (issue #79). + +Composition layer — the only package importing both octave.mcp and +octave.inference. Future Agent Manager components (manager, registry, +router) land here as additional modules; they do not exist yet. +""" + +from octave.agent.errors import ToolLoopLimitError +from octave.agent.types import ToolOutcome, ToolTurn +from octave.tools.errors import ToolError + +__all__ = ["ToolError", "ToolLoopLimitError", "ToolOutcome", "ToolTurn"] +``` + +- [ ] **Step 4: Run type tests to verify they pass** + +Run: `uv run pytest tests/agent/test_types.py -v` +Expected: PASS (`test_package.py::test_public_names_are_exported` is expected-red until Task 8 completes — do not run it green yet) + +- [ ] **Step 5: Commit** + +```bash +git add backend/src/octave/agent backend/tests/agent +git commit -m "feat(agent): scaffold agent package with orchestration vocabulary" +``` + +--- + +### Task 7: `ToolLoop` + `ToolExecutor` Protocol + +**Files:** +- Create: `backend/src/octave/agent/executor.py` +- Create: `backend/src/octave/agent/loop.py` +- Modify: `backend/src/octave/agent/__init__.py` +- Create: `backend/tests/agent/fakes.py` +- Create: `backend/tests/agent/test_loop.py` + +- [ ] **Step 1: Write the fakes** — create `backend/tests/agent/fakes.py` + +```python +"""Scripted adapter + recording executor for loop tests.""" + +from collections.abc import AsyncIterator + +from octave.agent.types import ToolOutcome +from octave.inference.adapter import InferenceAdapter +from octave.inference.config import AdapterConfig +from octave.inference.types import ( + CompletionChunk, + CompletionRequest, + CompletionResult, + EmbeddingRequest, + EmbeddingResult, + ModelInfo, + ToolCall, +) + + +class ScriptedAdapter(InferenceAdapter): + """Serves queued results (Exception items raise); records every request.""" + + def __init__(self, results: list[CompletionResult | Exception]) -> None: + super().__init__(AdapterConfig(adapter="fake", base_url="http://fake.test/v1")) + self._results = list(results) + self.complete_calls: list[CompletionRequest] = [] + + async def complete(self, request: CompletionRequest) -> CompletionResult: + self.complete_calls.append(request) + if not self._results: + raise AssertionError("ScriptedAdapter queue exhausted") + item = self._results.pop(0) + if isinstance(item, Exception): + raise item + return item + + async def stream(self, request: CompletionRequest) -> AsyncIterator[CompletionChunk]: + raise NotImplementedError + yield CompletionChunk() # satisfy async-generator typing + + async def embed(self, request: EmbeddingRequest) -> EmbeddingResult: + raise NotImplementedError + + async def list_models(self) -> list[ModelInfo]: + return [] + + +def tool_call(call_id: str, name: str, arguments: str = "{}") -> ToolCall: + return ToolCall(id=call_id, name=name, arguments=arguments) + + +def tool_result(*calls: ToolCall) -> CompletionResult: + return CompletionResult( + text="", model="fake", finish_reason="tool_calls", tool_calls=list(calls) + ) + + +def final_result(text: str = "final") -> CompletionResult: + return CompletionResult(text=text, model="fake", finish_reason="stop") + + +class RecordingExecutor: + """Returns queued outcomes per call (in order); records invocations.""" + + def __init__(self, outcomes: list[ToolOutcome] | None = None) -> None: + self._outcomes = list(outcomes) if outcomes is not None else [] + self.calls: list[tuple[str, str, dict]] = [] + + async def call( + self, server_id: str, tool_name: str, arguments: dict + ) -> ToolOutcome: + self.calls.append((server_id, tool_name, arguments)) + if self._outcomes: + return self._outcomes.pop(0) + return ToolOutcome(content="ok") +``` + +- [ ] **Step 2: Write the failing tests** — create `backend/tests/agent/test_loop.py` + +```python +"""ToolLoop: round mechanics, transcript ordering, two-tier error split.""" + +import pytest + +from octave.agent.errors import ToolLoopLimitError +from octave.agent.loop import ToolLoop +from octave.agent.types import ToolOutcome +from octave.inference.errors import AdapterConnectionError +from octave.inference.types import Message +from octave.tools.types import ProviderToolset, ToolRoute +from tests.agent.fakes import ( + RecordingExecutor, + ScriptedAdapter, + final_result, + tool_call, + tool_result, +) + +USER = [Message(role="user", content="hi")] +ROUTE = ToolRoute(server_id="srv", tool_name="read") +TOOLSET = ProviderToolset(tools=[], routes={"mcp__fs__read": ROUTE}) + + +def _loop(adapter: ScriptedAdapter, executor: RecordingExecutor, max_rounds: int = 8) -> ToolLoop: + return ToolLoop(adapter=adapter, executor=executor, max_tool_rounds=max_rounds) + + +async def test_passthrough_without_tool_calls() -> None: + adapter = ScriptedAdapter([final_result("yo")]) + loop = _loop(adapter, RecordingExecutor()) + turn = await loop.run(list(USER), TOOLSET) + assert turn.tool_rounds == 0 + assert turn.messages == [ + *USER, + Message(role="assistant", content="yo"), + ] + assert turn.result.text == "yo" + + +async def test_one_round_multiple_calls() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")), final_result()] + ) + executor = RecordingExecutor() + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert turn.tool_rounds == 1 + assert executor.calls == [("srv", "read", {}), ("srv", "read", {})] + assert turn.messages == [ + *USER, + Message( + role="assistant", + content="", + tool_calls=[tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")], + ), + Message(role="tool", content="ok", tool_call_id="c1", name="mcp__fs__read"), + Message(role="tool", content="ok", tool_call_id="c2", name="mcp__fs__read"), + Message(role="assistant", content="final"), + ] + + +async def test_multi_round_chain() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), tool_result(tool_call("c2", "mcp__fs__read")), final_result()] + ) + turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + assert turn.tool_rounds == 2 + assert [m.role for m in turn.messages] == [ + "user", "assistant", "tool", "assistant", "tool", "assistant", + ] + + +async def test_tools_sent_every_round() -> None: + from octave.inference.types import ToolDefinition + + populated = ProviderToolset( + tools=[ToolDefinition(name="mcp__fs__read")], routes={"mcp__fs__read": ROUTE} + ) + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), final_result()] + ) + await _loop(adapter, RecordingExecutor()).run(list(USER), populated) + assert [call.tools for call in adapter.complete_calls] == [populated.tools] * 2 + + +async def test_empty_toolset_omits_tools() -> None: + adapter = ScriptedAdapter([final_result()]) + await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + # empty toolset -> tools omitted (None), never an empty list + assert adapter.complete_calls[0].tools is None + + +async def test_loop_limit_raises_with_partial_transcript() -> None: + adapter = ScriptedAdapter([tool_result(tool_call("c1", "mcp__fs__read"))] * 3) + loop = _loop(adapter, RecordingExecutor(), max_rounds=2) + with pytest.raises(ToolLoopLimitError) as excinfo: + await loop.run(list(USER), TOOLSET) + assert len(adapter.complete_calls) == 3 + assert len(excinfo.value.messages) == 1 + 2 * 2 # user + (assistant+tool) x 2 + assert excinfo.value.messages[-1].role == "assistant" + assert excinfo.value.messages[-1].tool_calls is not None + + +async def test_finish_reason_tool_calls_with_empty_list_is_final() -> None: + from octave.inference.types import CompletionResult + + adapter = ScriptedAdapter( + [CompletionResult(text="done", model="fake", finish_reason="tool_calls", tool_calls=[])] + ) + turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + assert turn.tool_rounds == 0 + assert turn.messages[-1].content == "done" + + +async def test_malformed_arguments_is_error_tool_message() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read", "{not json")), final_result()] + ) + executor = RecordingExecutor() + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert executor.calls == [] + tool_message = turn.messages[2] + assert tool_message.role == "tool" + assert tool_message.tool_call_id == "c1" + assert tool_message.content.startswith("Error: ") + assert "not a valid JSON object" in tool_message.content + + +async def test_non_object_json_arguments_is_error_tool_message() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read", "[1, 2]")), final_result()] + ) + turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + assert turn.messages[2].content.startswith("Error: ") + + +async def test_unknown_route_is_error_tool_message() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__ghost__x")), final_result()] + ) + executor = RecordingExecutor() + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert executor.calls == [] + assert turn.messages[2].content == "Error: Unknown tool: mcp__ghost__x" + + +async def test_error_outcome_continues_loop() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), final_result()] + ) + executor = RecordingExecutor([ToolOutcome(content="boom", is_error=True)]) + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert turn.messages[2].content == "Error: boom" + assert turn.tool_rounds == 1 + + +async def test_adapter_error_propagates_untouched() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), AdapterConnectionError("engine down")] + ) + with pytest.raises(AdapterConnectionError): + await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + + +async def test_arguments_dict_forwarded_verbatim() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read", '{"path": "/x", "n": 3}')), final_result()] + ) + executor = RecordingExecutor() + await _loop(adapter, executor).run(list(USER), TOOLSET) + assert executor.calls == [("srv", "read", {"path": "/x", "n": 3})] +``` + +- [ ] **Step 3: Run tests to verify they fail** + +Run: `uv run pytest tests/agent/test_loop.py -v` +Expected: FAIL — `ModuleNotFoundError: No module named 'octave.agent.loop'` + +- [ ] **Step 4: Implement `executor.py`** — create `backend/src/octave/agent/executor.py` + +```python +"""Tool-execution seam (design spec #79 decision 7).""" + +from typing import Any, Protocol + +from octave.agent.types import ToolOutcome + +__all__ = ["ToolExecutor"] + + +class ToolExecutor(Protocol): + """Executes one resolved tool call. + + Implementations must not raise for tool-level failures — return + ``ToolOutcome(is_error=True)`` so the model can self-correct. + """ + + async def call( + self, server_id: str, tool_name: str, arguments: dict[str, Any] + ) -> ToolOutcome: ... +``` + +- [ ] **Step 5: Implement `loop.py`** — create `backend/src/octave/agent/loop.py` + +```python +"""The reason → act → observe loop (issue #79). + +Composes an InferenceAdapter, a ToolExecutor, and a pre-built +ProviderToolset (#78 routes). Tier-1 tool failures become error tool +messages the model can react to; the round limit and adapter failures +raise (spec decision 3). +""" + +import json +import logging +from typing import Any + +from octave.agent.errors import ToolLoopLimitError +from octave.agent.executor import ToolExecutor +from octave.agent.types import ToolOutcome, ToolTurn +from octave.inference.adapter import InferenceAdapter +from octave.inference.types import CompletionRequest, Message, ToolCall +from octave.tools.types import ProviderToolset + +__all__ = ["ToolLoop"] + +logger = logging.getLogger(__name__) + +_NO_OUTPUT = "(no output)" + + +class ToolLoop: + """Runs one orchestrated turn: completions, tool rounds, final synthesis.""" + + def __init__( + self, + *, + adapter: InferenceAdapter, + executor: ToolExecutor, + max_tool_rounds: int = 8, + ) -> None: + self._adapter = adapter + self._executor = executor + self._max_tool_rounds = max_tool_rounds + + async def run( + self, + messages: list[Message], + toolset: ProviderToolset, + *, + model: str | None = None, + ) -> ToolTurn: + """Drive messages to a final (no tool_calls) completion. + + Raises ToolLoopLimitError if the model still requests tools after + ``max_tool_rounds`` executions; AdapterError propagates untouched. + """ + history = list(messages) + rounds = 0 + tools = toolset.tools or None + while True: + result = await self._adapter.complete( + CompletionRequest(model=model, messages=history, tools=tools) + ) + if not result.tool_calls: + history.append(Message(role="assistant", content=result.text)) + return ToolTurn(messages=history, result=result, tool_rounds=rounds) + history.append( + Message( + role="assistant", content=result.text, tool_calls=result.tool_calls + ) + ) + if rounds >= self._max_tool_rounds: + logger.warning( + "tool loop limit reached | max_tool_rounds=%s calls=%s", + self._max_tool_rounds, + len(result.tool_calls), + ) + raise ToolLoopLimitError( + f"tool loop exhausted after {self._max_tool_rounds} rounds", + messages=history, + ) + rounds += 1 + logger.debug("tool round | round=%s calls=%s", rounds, len(result.tool_calls)) + for call in result.tool_calls: + history.append(await self._execute(call, toolset)) + + async def _execute(self, call: ToolCall, toolset: ProviderToolset) -> Message: + """One call -> one tool message; Tier-1 failures never raise.""" + route = toolset.routes.get(call.name) + if route is None: + return self._tool_message( + call, ToolOutcome(content=f"Unknown tool: {call.name}", is_error=True) + ) + arguments = self._parse_arguments(call) + if isinstance(arguments, str): # parse error message + return self._tool_message( + call, + ToolOutcome( + content=f"Tool call arguments are not a valid JSON object: {arguments}", + is_error=True, + ), + ) + outcome = await self._executor.call(route.server_id, route.tool_name, arguments) + if outcome.is_error: + logger.warning( + "tool call failed | tool=%s server=%s error=%s", + call.name, + route.server_id, + outcome.content, + ) + return self._tool_message(call, outcome) + + @staticmethod + def _parse_arguments(call: ToolCall) -> dict[str, Any] | str: + """Parsed dict, or an error description string.""" + try: + parsed = json.loads(call.arguments or "{}") + except json.JSONDecodeError as exc: + return str(exc) + if not isinstance(parsed, dict): + return f"expected an object, got {type(parsed).__name__}" + return parsed + + @staticmethod + def _tool_message(call: ToolCall, outcome: ToolOutcome) -> Message: + content = outcome.content or _NO_OUTPUT + if outcome.is_error: + content = f"Error: {content}" + return Message( + role="tool", content=content, tool_call_id=call.id, name=call.name + ) +``` + +Update `backend/src/octave/agent/__init__.py` imports and `__all__` to: + +```python +from octave.agent.errors import ToolLoopLimitError +from octave.agent.executor import ToolExecutor +from octave.agent.loop import ToolLoop +from octave.agent.types import ToolOutcome, ToolTurn +from octave.tools.errors import ToolError + +__all__ = [ + "ToolError", + "ToolExecutor", + "ToolLoop", + "ToolLoopLimitError", + "ToolOutcome", + "ToolTurn", +] +``` + +- [ ] **Step 6: Run tests to verify they pass** + +Run: `uv run pytest tests/agent/test_loop.py -v` +Expected: all PASS (13 tests) + +- [ ] **Step 7: Commit** + +```bash +git add backend/src/octave/agent backend/tests/agent +git commit -m "feat(agent): implement tool-use orchestration loop" +``` + +--- + +### Task 8: `McpToolExecutor` + +**Files:** +- Create: `backend/src/octave/agent/mcp_executor.py` +- Modify: `backend/src/octave/agent/__init__.py` +- Create: `backend/tests/agent/test_mcp_executor.py` + +- [ ] **Step 1: Write the failing tests** — create `backend/tests/agent/test_mcp_executor.py` + +```python +"""McpToolExecutor: content flattening and McpError -> outcome conversion.""" + +from typing import Any + +import pytest + +from octave.agent.mcp_executor import McpToolExecutor +from octave.mcp.errors import ( + McpConfigError, + McpConnectionError, + McpNotConnectedError, + McpRpcError, + McpTimeoutError, +) +from octave.mcp.types import ToolContent, ToolResult + + +class StubRegistry: + """Stands in for ToolRegistry; returns or raises what the test dictates.""" + + def __init__( + self, result: ToolResult | None = None, error: Exception | None = None + ) -> None: + self._result = result + self._error = error + self.calls: list[tuple[str, str, dict[str, Any]]] = [] + + async def call_tool( + self, id: str, name: str, arguments: dict[str, Any] | None = None + ) -> ToolResult: + self.calls.append((id, name, dict(arguments or {}))) + if self._error is not None: + raise self._error + assert self._result is not None + return self._result + + +async def test_joins_text_blocks_and_passes_is_error_through() -> None: + registry = StubRegistry( + result=ToolResult( + content=[ToolContent(kind="text", text="a"), ToolContent(kind="text", text="b")] + ) + ) + outcome = await McpToolExecutor(registry).call("srv", "read", {"p": 1}) # type: ignore[arg-type] + assert registry.calls == [("srv", "read", {"p": 1})] + assert outcome.content == "a\nb" + assert outcome.is_error is False + + +async def test_error_result_flags_outcome() -> None: + registry = StubRegistry( + result=ToolResult( + content=[ToolContent(kind="text", text="disk full")], is_error=True + ) + ) + outcome = await McpToolExecutor(registry).call("srv", "read", {}) # type: ignore[arg-type] + assert outcome.content == "disk full" + assert outcome.is_error is True + + +async def test_empty_content_becomes_sentinel() -> None: + registry = StubRegistry(result=ToolResult(content=[])) + outcome = await McpToolExecutor(registry).call("srv", "read", {}) # type: ignore[arg-type] + assert outcome.content == "(no output)" + + +@pytest.mark.parametrize( + "error", + [ + McpRpcError("bad params", code=-32602), + McpTimeoutError("timed out after 60s"), + McpConnectionError("server died"), + McpConfigError("unknown server id"), + McpNotConnectedError("not connected"), + ], +) +async def test_mcp_errors_become_error_outcomes(error: Exception) -> None: + registry = StubRegistry(error=error) + outcome = await McpToolExecutor(registry).call("srv", "read", {}) # type: ignore[arg-type] + assert outcome.is_error is True + assert str(error) in outcome.content +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `uv run pytest tests/agent/test_mcp_executor.py -v` +Expected: FAIL — `ModuleNotFoundError: No module named 'octave.agent.mcp_executor'` + +- [ ] **Step 3: Implement** — create `backend/src/octave/agent/mcp_executor.py` + +```python +"""ToolExecutor over ToolRegistry (issue #79). + +The only octave.agent module importing octave.mcp. Converts McpError +subclasses into error ToolOutcomes (Tier-1, model-correctable) so the +loop never aborts on a tool failure. +""" + +import logging +from typing import Any + +from octave.agent.types import ToolOutcome +from octave.mcp import McpError, ToolRegistry + +__all__ = ["McpToolExecutor"] + +logger = logging.getLogger(__name__) + +_NO_OUTPUT = "(no output)" + + +class McpToolExecutor: + """Structural implementation of ``octave.agent.executor.ToolExecutor``.""" + + def __init__(self, registry: ToolRegistry) -> None: + self._registry = registry + + async def call( + self, server_id: str, tool_name: str, arguments: dict[str, Any] + ) -> ToolOutcome: + """Invoke one tool; every failure mode becomes an error outcome.""" + try: + result = await self._registry.call_tool(server_id, tool_name, arguments) + except McpError as exc: + logger.error( + "MCP tool execution failed | server=%s tool=%s error=%s", + server_id, + tool_name, + exc, + ) + return ToolOutcome(content=f"Tool execution failed: {exc}", is_error=True) + content = "\n".join(block.text for block in result.content) + return ToolOutcome(content=content or _NO_OUTPUT, is_error=result.is_error) +``` + +Update `backend/src/octave/agent/__init__.py`: add `from octave.agent.mcp_executor import McpToolExecutor` and `"McpToolExecutor"` in `__all__` (alphabetical first). + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `uv run pytest tests/agent -v` +Expected: all PASS — **including `test_package.py`**, which now goes green (full export surface + no SDK imports). + +- [ ] **Step 5: Commit** + +```bash +git add backend/src/octave/agent backend/tests/agent +git commit -m "feat(agent): add McpToolExecutor over ToolRegistry" +``` + +--- + +### Task 9: Full-suite gate, lint, types, push + +**Files:** none new + +- [ ] **Step 1: Full backend suite** + +Run: `uv run pytest` +Expected: all PASS (0 failed). If a pre-existing test asserts old `Message` serialization or `role="tool"` rejection, it was missed — fix it to match the spec, do not skip. + +- [ ] **Step 2: Lint** + +Run: `uv run ruff check src tests` +Expected: `All checks passed!` (fix imports/ordering if flagged; `ruff format` not used by this repo) + +- [ ] **Step 3: Types** + +Run: `uv run mypy src` +Expected: `Success: no issues found in N source files` + +- [ ] **Step 4: Push** + +```bash +git push +``` + +Expected: branch `feature/tool-use-orchestration-loop` updated; draft PR #103 reflects the new commits. + +- [ ] **Step 5: Confirm the PR is still draft and leave it for human review** + +Run: `gh pr view 103 --json isDraft,title` +Expected: `isDraft: true` — do not mark ready or merge (release process owns that). + +--- + +## Self-review notes (author) + +- **Spec coverage:** §Data model → Tasks 2/6/7/8; adapter parse/serialize → Tasks 3/4; loop mechanics incl. limit semantics (N+1 completions, unfulfilled assistant message retained) → Task 7 tests; two-tier errors → Task 7 (Tier-1 in loop) + Task 8 (Tier-1 in executor) + Task 7 Step 5 (Tier-2 raise/propagate); package quarantine → Task 6 guard; deferrals untouched (streaming, parallel exec, wiring). Test 10's exact spec assertions mirrored in `test_loop_limit_raises_with_partial_transcript`. +- **Clarification applied:** OpenAI tool messages have no error flag; `Error: ` prefix conveys Tier-1 errors (spec intent: model sees failure as tool result). Spec §Testing item 14 phrasing "tool message is_error=True" maps to this prefix convention. +- **Type consistency:** `ToolOutcome(content, is_error)`, `ToolTurn(messages, result, tool_rounds)`, `ToolExecutor.call(server_id, tool_name, arguments)`, `ToolLoop(adapter=, executor=, max_tool_rounds=)`, `ToolLoop.run(messages, toolset, *, model=)`, `McpToolExecutor(registry)` — identical across all tasks and `__init__` exports. From 6233bc8fa2ba3837b7446b30cd878fc013b5ce3a Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:10:29 +0100 Subject: [PATCH 03/13] feat(inference): add tool-call vocabulary to domain types --- backend/src/octave/inference/__init__.py | 2 ++ backend/src/octave/inference/types.py | 25 +++++++++++++++-- backend/tests/inference/test_types.py | 34 ++++++++++++++++++++++-- 3 files changed, 57 insertions(+), 4 deletions(-) diff --git a/backend/src/octave/inference/__init__.py b/backend/src/octave/inference/__init__.py index 69b33b7..b72599b 100644 --- a/backend/src/octave/inference/__init__.py +++ b/backend/src/octave/inference/__init__.py @@ -27,6 +27,7 @@ EmbeddingResult, Message, ModelInfo, + ToolCall, ToolDefinition, Usage, ) @@ -52,6 +53,7 @@ "ModelInfo", "ModelNotFoundError", "OpenAIAdapter", + "ToolCall", "ToolDefinition", "UnknownAdapterError", "Usage", diff --git a/backend/src/octave/inference/types.py b/backend/src/octave/inference/types.py index 1a93685..25c8ee5 100644 --- a/backend/src/octave/inference/types.py +++ b/backend/src/octave/inference/types.py @@ -17,18 +17,38 @@ "EmbeddingResult", "Message", "ModelInfo", + "ToolCall", "ToolDefinition", "Usage", ] -Role = Literal["system", "user", "assistant"] +Role = Literal["system", "user", "assistant", "tool"] + + +class ToolCall(BaseModel): + """One tool call as the model requested it. + + ``arguments`` is the raw JSON string from the wire; Octave never + interprets it (the adapter stays a lossless translator, the loop parses). + """ + + id: str + name: str + arguments: str class Message(BaseModel): - """A single chat message.""" + """A single chat message. + + ``tool_calls`` populates assistant messages; ``tool_call_id`` and ``name`` + populate ``role="tool"`` result messages (keyed to the call they answer). + """ role: Role content: str + tool_calls: list[ToolCall] | None = None + tool_call_id: str | None = None + name: str | None = None class ToolDefinition(BaseModel): @@ -74,6 +94,7 @@ class CompletionResult(BaseModel): model: str finish_reason: str | None = None usage: Usage | None = None + tool_calls: list[ToolCall] | None = None class CompletionChunk(BaseModel): diff --git a/backend/tests/inference/test_types.py b/backend/tests/inference/test_types.py index 226ca39..d4b0871 100644 --- a/backend/tests/inference/test_types.py +++ b/backend/tests/inference/test_types.py @@ -6,9 +6,11 @@ from octave.inference.types import ( CompletionChunk, CompletionRequest, + CompletionResult, EmbeddingRequest, Message, ModelInfo, + ToolCall, ToolDefinition, ) @@ -32,9 +34,37 @@ def test_completion_request_extra_is_not_shared() -> None: assert b.extra == {} -def test_message_rejects_unknown_role() -> None: +def test_tool_role_is_valid() -> None: + message = Message(role="tool", content="result", tool_call_id="call_1", name="mcp__fs__read") + assert message.tool_calls is None + + +def test_message_tool_fields_default_none() -> None: + message = Message(role="assistant", content="hi") + assert message.tool_calls is None + assert message.tool_call_id is None + assert message.name is None + + +def test_message_rejects_truly_unknown_role() -> None: with pytest.raises(ValidationError): - Message(role="tool", content="hi") # type: ignore[arg-type] + Message(role="developer", content="hi") # type: ignore[arg-type] + + +def test_tool_call_round_trip() -> None: + call = ToolCall(id="call_1", name="mcp__fs__read", arguments='{"path": "/tmp/x"}') + assert ToolCall.model_validate(call.model_dump()) == call + + +def test_completion_result_tool_calls_defaults_none() -> None: + result = CompletionResult(text="hi", model="m") + assert result.tool_calls is None + + +def test_completion_result_round_trips_tool_calls() -> None: + call = ToolCall(id="c1", name="x", arguments="{}") + result = CompletionResult(text="", model="m", finish_reason="tool_calls", tool_calls=[call]) + assert CompletionResult.model_validate(result.model_dump()) == result def test_completion_chunk_defaults() -> None: From ae4e209a39721f0efed93214993edd0f6ee1a77e Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:16:21 +0100 Subject: [PATCH 04/13] feat(inference): parse tool_calls from completion responses --- .../src/octave/inference/openai_adapter.py | 10 +++ .../tests/inference/test_openai_adapter.py | 66 ++++++++++++++++++- 2 files changed, 75 insertions(+), 1 deletion(-) diff --git a/backend/src/octave/inference/openai_adapter.py b/backend/src/octave/inference/openai_adapter.py index 73e7160..540a6dd 100644 --- a/backend/src/octave/inference/openai_adapter.py +++ b/backend/src/octave/inference/openai_adapter.py @@ -33,6 +33,7 @@ EmbeddingRequest, EmbeddingResult, ModelInfo, + ToolCall, Usage, ) @@ -93,11 +94,20 @@ async def complete(self, request: CompletionRequest) -> CompletionResult: except openai.APIError as exc: raise _translate(exc) from exc choice = response.choices[0] + tool_calls = [ + ToolCall( + id=call.id, + name=call.function.name, + arguments=call.function.arguments, + ) + for call in (choice.message.tool_calls or []) + ] return CompletionResult( text=choice.message.content or "", model=response.model, finish_reason=choice.finish_reason, usage=_chat_usage(response.usage), + tool_calls=tool_calls or None, ) async def stream( diff --git a/backend/tests/inference/test_openai_adapter.py b/backend/tests/inference/test_openai_adapter.py index fff37bc..3ccd779 100644 --- a/backend/tests/inference/test_openai_adapter.py +++ b/backend/tests/inference/test_openai_adapter.py @@ -16,7 +16,7 @@ ModelNotFoundError, ) from octave.inference.openai_adapter import OpenAIAdapter -from octave.inference.types import CompletionRequest, Message, ToolDefinition +from octave.inference.types import CompletionRequest, Message, ToolCall, ToolDefinition from tests.inference.conformance import InferenceAdapterConformanceSuite BASE_URL = "http://engine.test/v1" @@ -348,3 +348,67 @@ def capture(request: httpx.Request) -> httpx.Response: assert all("tools" not in json.loads(r.content) for r in captured) await adapter.aclose() await client.aclose() + + +CHAT_RESPONSE_TOOL_CALLS = { + "id": "chatcmpl-tc1", + "object": "chat.completion", + "created": 1725000000, + "model": "fake-model", + "choices": [ + { + "index": 0, + "message": { + "role": "assistant", + "content": None, + "tool_calls": [ + { + "id": "call_1", + "type": "function", + "function": { + "name": "mcp__fs__read", + "arguments": '{"path": "/tmp/x"}', + }, + }, + { + "id": "call_2", + "type": "function", + "function": {"name": "mcp__fs__list", "arguments": "{}"}, + }, + ], + }, + "finish_reason": "tool_calls", + } + ], + "usage": {"prompt_tokens": 9, "completion_tokens": 12, "total_tokens": 21}, +} + + +async def test_complete_parses_tool_calls() -> None: + def handler(request: httpx.Request) -> httpx.Response: + return httpx.Response(200, json=CHAT_RESPONSE_TOOL_CALLS) + + client = _mock_client(handler) + adapter = _adapter(client) + result = await adapter.complete( + CompletionRequest(model=None, messages=[Message(role="user", content="hi")]) + ) + assert result.finish_reason == "tool_calls" + assert result.text == "" # None content normalized to "" + assert result.tool_calls == [ + ToolCall(id="call_1", name="mcp__fs__read", arguments='{"path": "/tmp/x"}'), + ToolCall(id="call_2", name="mcp__fs__list", arguments="{}"), + ] + await adapter.aclose() + await client.aclose() + + +async def test_complete_without_tool_calls_is_none() -> None: + client = _mock_client() + adapter = _adapter(client) + result = await adapter.complete( + CompletionRequest(model=None, messages=[Message(role="user", content="hi")]) + ) + assert result.tool_calls is None + await adapter.aclose() + await client.aclose() From ef42d8c4989c4012dce81c883402420adccc1566 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:19:38 +0100 Subject: [PATCH 05/13] feat(inference): serialize tool-carrying messages to provider payloads --- .../src/octave/inference/openai_adapter.py | 19 +++++- .../tests/inference/test_openai_adapter.py | 65 +++++++++++++++++++ 2 files changed, 83 insertions(+), 1 deletion(-) diff --git a/backend/src/octave/inference/openai_adapter.py b/backend/src/octave/inference/openai_adapter.py index 540a6dd..04b6a36 100644 --- a/backend/src/octave/inference/openai_adapter.py +++ b/backend/src/octave/inference/openai_adapter.py @@ -32,6 +32,7 @@ CompletionResult, EmbeddingRequest, EmbeddingResult, + Message, ModelInfo, ToolCall, Usage, @@ -168,7 +169,7 @@ async def aclose(self) -> None: def _chat_kwargs(self, request: CompletionRequest) -> dict[str, Any]: kwargs: dict[str, Any] = { "model": self._resolve_model(request.model), - "messages": [message.model_dump() for message in request.messages], + "messages": [self._message_payload(message) for message in request.messages], } if request.temperature is not None: kwargs["temperature"] = request.temperature @@ -187,6 +188,22 @@ def _chat_kwargs(self, request: CompletionRequest) -> dict[str, Any]: kwargs["extra_body"] = dict(request.extra) return kwargs + def _message_payload(self, message: Message) -> dict[str, Any]: + """One message as a provider payload. exclude_none keeps tool fields + off plain messages (strict local engines reject null tool keys); + assistant tool_calls get the provider function-call envelope.""" + payload = message.model_dump(mode="json", exclude_none=True) + if message.tool_calls: + payload["tool_calls"] = [ + { + "id": call.id, + "type": "function", + "function": {"name": call.name, "arguments": call.arguments}, + } + for call in message.tool_calls + ] + return payload + def _resolve_model(self, requested: str | None) -> str: model = requested or self._config.default_model if model is None: diff --git a/backend/tests/inference/test_openai_adapter.py b/backend/tests/inference/test_openai_adapter.py index 3ccd779..463b00b 100644 --- a/backend/tests/inference/test_openai_adapter.py +++ b/backend/tests/inference/test_openai_adapter.py @@ -412,3 +412,68 @@ async def test_complete_without_tool_calls_is_none() -> None: assert result.tool_calls is None await adapter.aclose() await client.aclose() + + +async def _capture_chat_body(messages: list[Message]) -> dict: + captured: list[httpx.Request] = [] + + def handler(request: httpx.Request) -> httpx.Response: + captured.append(request) + return httpx.Response(200, json=CHAT_RESPONSE) + + client = _mock_client(handler) + adapter = _adapter(client) + await adapter.complete(CompletionRequest(model=None, messages=messages)) + await adapter.aclose() + await client.aclose() + return json.loads(captured[0].content) + + +async def test_plain_messages_emit_no_null_tool_fields() -> None: + body = await _capture_chat_body( + [ + Message(role="system", content="sys"), + Message(role="user", content="hi"), + Message(role="assistant", content="yo"), + ] + ) + for message in body["messages"]: + assert "tool_calls" not in message + assert "tool_call_id" not in message + assert "name" not in message + + +async def test_assistant_tool_calls_message_uses_provider_envelope() -> None: + body = await _capture_chat_body( + [ + Message(role="user", content="hi"), + Message( + role="assistant", + content="", + tool_calls=[ToolCall(id="call_1", name="mcp__fs__read", arguments='{"p": 1}')], + ), + ] + ) + assert body["messages"][1] == { + "role": "assistant", + "content": "", + "tool_calls": [ + { + "id": "call_1", + "type": "function", + "function": {"name": "mcp__fs__read", "arguments": '{"p": 1}'}, + } + ], + } + + +async def test_tool_result_message_carries_tool_call_id_and_name() -> None: + body = await _capture_chat_body( + [Message(role="tool", content="42", tool_call_id="call_1", name="mcp__fs__read")] + ) + assert body["messages"][0] == { + "role": "tool", + "content": "42", + "tool_call_id": "call_1", + "name": "mcp__fs__read", + } From 129da2d68493b52424398fe4c6514070312c25e0 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:21:41 +0100 Subject: [PATCH 06/13] refactor(tools): introduce ToolError base for the tool plane --- backend/src/octave/tools/__init__.py | 3 ++- backend/src/octave/tools/errors.py | 12 ++++++++---- backend/tests/tools/test_package.py | 7 +++++++ 3 files changed, 17 insertions(+), 5 deletions(-) diff --git a/backend/src/octave/tools/__init__.py b/backend/src/octave/tools/__init__.py index 8fc9653..de24c20 100644 --- a/backend/src/octave/tools/__init__.py +++ b/backend/src/octave/tools/__init__.py @@ -5,12 +5,13 @@ this boundary (guarded by tests/tools/test_package.py). """ -from octave.tools.errors import ToolNameCollisionError, ToolTranslationError +from octave.tools.errors import ToolError, ToolNameCollisionError, ToolTranslationError from octave.tools.translate import translate_tools from octave.tools.types import ProviderToolset, ToolRoute __all__ = [ "ProviderToolset", + "ToolError", "ToolNameCollisionError", "ToolRoute", "ToolTranslationError", diff --git a/backend/src/octave/tools/errors.py b/backend/src/octave/tools/errors.py index 726f67e..821c2c1 100644 --- a/backend/src/octave/tools/errors.py +++ b/backend/src/octave/tools/errors.py @@ -1,10 +1,14 @@ -"""Translation-layer errors (design spec #78).""" +"""Tool-plane errors (design specs #78, #79).""" -__all__ = ["ToolNameCollisionError", "ToolTranslationError"] +__all__ = ["ToolError", "ToolNameCollisionError", "ToolTranslationError"] -class ToolTranslationError(Exception): - """Base for tool translation failures.""" +class ToolError(Exception): + """Base for tool-plane failures (translation and orchestration).""" + + +class ToolTranslationError(ToolError): + """Base for translation failures.""" class ToolNameCollisionError(ToolTranslationError): diff --git a/backend/tests/tools/test_package.py b/backend/tests/tools/test_package.py index a5fd768..6e80122 100644 --- a/backend/tests/tools/test_package.py +++ b/backend/tests/tools/test_package.py @@ -4,6 +4,7 @@ from pathlib import Path import octave.tools as tools_pkg +from octave.tools.errors import ToolError, ToolNameCollisionError, ToolTranslationError SDK_MODULES = {"openai", "mcp"} @@ -28,9 +29,15 @@ def test_no_sdk_imports() -> None: def test_public_names_are_exported() -> None: for name in ( "ProviderToolset", + "ToolError", "ToolRoute", "ToolTranslationError", "ToolNameCollisionError", "translate_tools", ): assert hasattr(tools_pkg, name), name + + +def test_tool_error_hierarchy() -> None: + assert issubclass(ToolTranslationError, ToolError) + assert issubclass(ToolNameCollisionError, ToolTranslationError) From 350c70a5dd77749881bb7ac7ef1c226ab5b1e811 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:24:26 +0100 Subject: [PATCH 07/13] feat(agent): scaffold agent package with orchestration vocabulary --- backend/src/octave/agent/__init__.py | 12 +++++++++ backend/src/octave/agent/errors.py | 19 ++++++++++++++ backend/src/octave/agent/types.py | 27 ++++++++++++++++++++ backend/tests/agent/__init__.py | 0 backend/tests/agent/test_package.py | 38 ++++++++++++++++++++++++++++ backend/tests/agent/test_types.py | 18 +++++++++++++ 6 files changed, 114 insertions(+) create mode 100644 backend/src/octave/agent/__init__.py create mode 100644 backend/src/octave/agent/errors.py create mode 100644 backend/src/octave/agent/types.py create mode 100644 backend/tests/agent/__init__.py create mode 100644 backend/tests/agent/test_package.py create mode 100644 backend/tests/agent/test_types.py diff --git a/backend/src/octave/agent/__init__.py b/backend/src/octave/agent/__init__.py new file mode 100644 index 0000000..abffb53 --- /dev/null +++ b/backend/src/octave/agent/__init__.py @@ -0,0 +1,12 @@ +"""Agent plane: tool-use orchestration loop (issue #79). + +Composition layer — the only package importing both octave.mcp and +octave.inference. Future Agent Manager components (manager, registry, +router) land here as additional modules; they do not exist yet. +""" + +from octave.agent.errors import ToolLoopLimitError +from octave.agent.types import ToolOutcome, ToolTurn +from octave.tools.errors import ToolError + +__all__ = ["ToolError", "ToolLoopLimitError", "ToolOutcome", "ToolTurn"] diff --git a/backend/src/octave/agent/errors.py b/backend/src/octave/agent/errors.py new file mode 100644 index 0000000..92fae05 --- /dev/null +++ b/backend/src/octave/agent/errors.py @@ -0,0 +1,19 @@ +"""Orchestration-fatal errors (design spec #79).""" + +from octave.inference.types import Message +from octave.tools.errors import ToolError + +__all__ = ["ToolLoopLimitError"] + + +class ToolLoopLimitError(ToolError): + """max_tool_rounds exhausted without a final answer. + + ``messages`` is the partial transcript (input + everything appended, + ending on the unfulfilled assistant tool_calls message) so callers can + persist partial work or retry. + """ + + def __init__(self, message: str, *, messages: list[Message]) -> None: + super().__init__(message) + self.messages = messages diff --git a/backend/src/octave/agent/types.py b/backend/src/octave/agent/types.py new file mode 100644 index 0000000..4760282 --- /dev/null +++ b/backend/src/octave/agent/types.py @@ -0,0 +1,27 @@ +"""Agent orchestration vocabulary (design spec #79).""" + +from pydantic import BaseModel + +from octave.inference.types import CompletionResult, Message + +__all__ = ["ToolOutcome", "ToolTurn"] + + +class ToolOutcome(BaseModel): + """Flattened result of one tool execution.""" + + content: str + is_error: bool = False + + +class ToolTurn(BaseModel): + """Result of one orchestrated turn.""" + + messages: list[Message] + """Input + appended assistant/tool messages + final assistant message.""" + + result: CompletionResult + """The final (no tool_calls) completion.""" + + tool_rounds: int + """Tool-execution rounds performed.""" diff --git a/backend/tests/agent/__init__.py b/backend/tests/agent/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/backend/tests/agent/test_package.py b/backend/tests/agent/test_package.py new file mode 100644 index 0000000..9378186 --- /dev/null +++ b/backend/tests/agent/test_package.py @@ -0,0 +1,38 @@ +"""Public API surface and SDK-quarantine posture of the agent package.""" + +import ast +from pathlib import Path + +import octave.agent as agent_pkg + +SDK_MODULES = {"openai", "mcp"} + + +def _imported_top_level_modules() -> set[str]: + names: set[str] = set() + for path in Path(agent_pkg.__file__).parent.rglob("*.py"): + tree = ast.parse(path.read_text()) + for node in ast.walk(tree): + if isinstance(node, ast.Import): + names.update(alias.name.split(".")[0] for alias in node.names) + elif isinstance(node, ast.ImportFrom) and node.module: + names.add(node.module.split(".")[0]) + return names + + +def test_no_sdk_imports() -> None: + """octave.agent composes Octave façades only, never the openai/mcp SDKs.""" + assert not _imported_top_level_modules() & SDK_MODULES + + +def test_public_names_are_exported() -> None: + for name in ( + "McpToolExecutor", + "ToolError", + "ToolExecutor", + "ToolLoop", + "ToolLoopLimitError", + "ToolOutcome", + "ToolTurn", + ): + assert hasattr(agent_pkg, name), name diff --git a/backend/tests/agent/test_types.py b/backend/tests/agent/test_types.py new file mode 100644 index 0000000..cc6b092 --- /dev/null +++ b/backend/tests/agent/test_types.py @@ -0,0 +1,18 @@ +"""Agent loop vocabulary: defaults and provenance of ToolOutcome/ToolTurn.""" + +from octave.agent.types import ToolOutcome, ToolTurn +from octave.inference.types import CompletionResult, Message + + +def test_tool_outcome_defaults() -> None: + outcome = ToolOutcome(content="ok") + assert outcome.is_error is False + + +def test_tool_turn_round_trip() -> None: + turn = ToolTurn( + messages=[Message(role="user", content="hi"), Message(role="assistant", content="yo")], + result=CompletionResult(text="yo", model="m"), + tool_rounds=0, + ) + assert ToolTurn.model_validate(turn.model_dump()) == turn From cd0d04319f759bbaf638dcc5a48b5dd8951db043 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:41:22 +0100 Subject: [PATCH 08/13] feat(agent): implement tool-use orchestration loop --- backend/src/octave/agent/__init__.py | 11 +- backend/src/octave/agent/executor.py | 19 +++ backend/src/octave/agent/loop.py | 127 ++++++++++++++++++++ backend/tests/agent/fakes.py | 74 ++++++++++++ backend/tests/agent/test_loop.py | 171 +++++++++++++++++++++++++++ 5 files changed, 401 insertions(+), 1 deletion(-) create mode 100644 backend/src/octave/agent/executor.py create mode 100644 backend/src/octave/agent/loop.py create mode 100644 backend/tests/agent/fakes.py create mode 100644 backend/tests/agent/test_loop.py diff --git a/backend/src/octave/agent/__init__.py b/backend/src/octave/agent/__init__.py index abffb53..80b4abf 100644 --- a/backend/src/octave/agent/__init__.py +++ b/backend/src/octave/agent/__init__.py @@ -6,7 +6,16 @@ """ from octave.agent.errors import ToolLoopLimitError +from octave.agent.executor import ToolExecutor +from octave.agent.loop import ToolLoop from octave.agent.types import ToolOutcome, ToolTurn from octave.tools.errors import ToolError -__all__ = ["ToolError", "ToolLoopLimitError", "ToolOutcome", "ToolTurn"] +__all__ = [ + "ToolError", + "ToolExecutor", + "ToolLoop", + "ToolLoopLimitError", + "ToolOutcome", + "ToolTurn", +] diff --git a/backend/src/octave/agent/executor.py b/backend/src/octave/agent/executor.py new file mode 100644 index 0000000..1c967ef --- /dev/null +++ b/backend/src/octave/agent/executor.py @@ -0,0 +1,19 @@ +"""Tool-execution seam (design spec #79 decision 7).""" + +from typing import Any, Protocol + +from octave.agent.types import ToolOutcome + +__all__ = ["ToolExecutor"] + + +class ToolExecutor(Protocol): + """Executes one resolved tool call. + + Implementations must not raise for tool-level failures — return + ``ToolOutcome(is_error=True)`` so the model can self-correct. + """ + + async def call( + self, server_id: str, tool_name: str, arguments: dict[str, Any] + ) -> ToolOutcome: ... diff --git a/backend/src/octave/agent/loop.py b/backend/src/octave/agent/loop.py new file mode 100644 index 0000000..f823b76 --- /dev/null +++ b/backend/src/octave/agent/loop.py @@ -0,0 +1,127 @@ +"""The reason → act → observe loop (issue #79). + +Composes an InferenceAdapter, a ToolExecutor, and a pre-built +ProviderToolset (#78 routes). Tier-1 tool failures become error tool +messages the model can react to; the round limit and adapter failures +raise (spec decision 3). +""" + +import json +import logging +from typing import Any + +from octave.agent.errors import ToolLoopLimitError +from octave.agent.executor import ToolExecutor +from octave.agent.types import ToolOutcome, ToolTurn +from octave.inference.adapter import InferenceAdapter +from octave.inference.types import CompletionRequest, Message, ToolCall +from octave.tools.types import ProviderToolset + +__all__ = ["ToolLoop"] + +logger = logging.getLogger(__name__) + +_NO_OUTPUT = "(no output)" + + +class ToolLoop: + """Runs one orchestrated turn: completions, tool rounds, final synthesis.""" + + def __init__( + self, + *, + adapter: InferenceAdapter, + executor: ToolExecutor, + max_tool_rounds: int = 8, + ) -> None: + self._adapter = adapter + self._executor = executor + self._max_tool_rounds = max_tool_rounds + + async def run( + self, + messages: list[Message], + toolset: ProviderToolset, + *, + model: str | None = None, + ) -> ToolTurn: + """Drive messages to a final (no tool_calls) completion. + + Raises ToolLoopLimitError if the model still requests tools after + ``max_tool_rounds`` executions; AdapterError propagates untouched. + """ + history = list(messages) + rounds = 0 + tools = toolset.tools or None + while True: + result = await self._adapter.complete( + CompletionRequest(model=model, messages=history, tools=tools) + ) + if not result.tool_calls: + history.append(Message(role="assistant", content=result.text)) + return ToolTurn(messages=history, result=result, tool_rounds=rounds) + history.append( + Message( + role="assistant", content=result.text, tool_calls=result.tool_calls + ) + ) + if rounds >= self._max_tool_rounds: + logger.warning( + "tool loop limit reached | max_tool_rounds=%s calls=%s", + self._max_tool_rounds, + len(result.tool_calls), + ) + raise ToolLoopLimitError( + f"tool loop exhausted after {self._max_tool_rounds} rounds", + messages=history, + ) + rounds += 1 + logger.debug("tool round | round=%s calls=%s", rounds, len(result.tool_calls)) + for call in result.tool_calls: + history.append(await self._execute(call, toolset)) + + async def _execute(self, call: ToolCall, toolset: ProviderToolset) -> Message: + """One call -> one tool message; Tier-1 failures never raise.""" + route = toolset.routes.get(call.name) + if route is None: + return self._tool_message( + call, ToolOutcome(content=f"Unknown tool: {call.name}", is_error=True) + ) + arguments = self._parse_arguments(call) + if isinstance(arguments, str): # parse error message + return self._tool_message( + call, + ToolOutcome( + content=f"Tool call arguments are not a valid JSON object: {arguments}", + is_error=True, + ), + ) + outcome = await self._executor.call(route.server_id, route.tool_name, arguments) + if outcome.is_error: + logger.warning( + "tool call failed | tool=%s server=%s error=%s", + call.name, + route.server_id, + outcome.content, + ) + return self._tool_message(call, outcome) + + @staticmethod + def _parse_arguments(call: ToolCall) -> dict[str, Any] | str: + """Parsed dict, or an error description string.""" + try: + parsed = json.loads(call.arguments or "{}") + except json.JSONDecodeError as exc: + return str(exc) + if not isinstance(parsed, dict): + return f"expected an object, got {type(parsed).__name__}" + return parsed + + @staticmethod + def _tool_message(call: ToolCall, outcome: ToolOutcome) -> Message: + content = outcome.content or _NO_OUTPUT + if outcome.is_error: + content = f"Error: {content}" + return Message( + role="tool", content=content, tool_call_id=call.id, name=call.name + ) diff --git a/backend/tests/agent/fakes.py b/backend/tests/agent/fakes.py new file mode 100644 index 0000000..1449332 --- /dev/null +++ b/backend/tests/agent/fakes.py @@ -0,0 +1,74 @@ +"""Scripted adapter + recording executor for loop tests.""" + +from collections.abc import AsyncIterator + +from octave.agent.types import ToolOutcome +from octave.inference.adapter import InferenceAdapter +from octave.inference.config import AdapterConfig +from octave.inference.types import ( + CompletionChunk, + CompletionRequest, + CompletionResult, + EmbeddingRequest, + EmbeddingResult, + ModelInfo, + ToolCall, +) + + +class ScriptedAdapter(InferenceAdapter): + """Serves queued results (Exception items raise); records every request.""" + + def __init__(self, results: list[CompletionResult | Exception]) -> None: + super().__init__(AdapterConfig(adapter="fake", base_url="http://fake.test/v1")) + self._results = list(results) + self.complete_calls: list[CompletionRequest] = [] + + async def complete(self, request: CompletionRequest) -> CompletionResult: + self.complete_calls.append(request) + if not self._results: + raise AssertionError("ScriptedAdapter queue exhausted") + item = self._results.pop(0) + if isinstance(item, Exception): + raise item + return item + + async def stream(self, request: CompletionRequest) -> AsyncIterator[CompletionChunk]: + raise NotImplementedError + yield CompletionChunk() # satisfy async-generator typing + + async def embed(self, request: EmbeddingRequest) -> EmbeddingResult: + raise NotImplementedError + + async def list_models(self) -> list[ModelInfo]: + return [] + + +def tool_call(call_id: str, name: str, arguments: str = "{}") -> ToolCall: + return ToolCall(id=call_id, name=name, arguments=arguments) + + +def tool_result(*calls: ToolCall) -> CompletionResult: + return CompletionResult( + text="", model="fake", finish_reason="tool_calls", tool_calls=list(calls) + ) + + +def final_result(text: str = "final") -> CompletionResult: + return CompletionResult(text=text, model="fake", finish_reason="stop") + + +class RecordingExecutor: + """Returns queued outcomes per call (in order); records invocations.""" + + def __init__(self, outcomes: list[ToolOutcome] | None = None) -> None: + self._outcomes = list(outcomes) if outcomes is not None else [] + self.calls: list[tuple[str, str, dict]] = [] + + async def call( + self, server_id: str, tool_name: str, arguments: dict + ) -> ToolOutcome: + self.calls.append((server_id, tool_name, arguments)) + if self._outcomes: + return self._outcomes.pop(0) + return ToolOutcome(content="ok") diff --git a/backend/tests/agent/test_loop.py b/backend/tests/agent/test_loop.py new file mode 100644 index 0000000..f71744d --- /dev/null +++ b/backend/tests/agent/test_loop.py @@ -0,0 +1,171 @@ +"""ToolLoop: round mechanics, transcript ordering, two-tier error split.""" + +import pytest + +from octave.agent.errors import ToolLoopLimitError +from octave.agent.loop import ToolLoop +from octave.agent.types import ToolOutcome +from octave.inference.errors import AdapterConnectionError +from octave.inference.types import Message +from octave.tools.types import ProviderToolset, ToolRoute +from tests.agent.fakes import ( + RecordingExecutor, + ScriptedAdapter, + final_result, + tool_call, + tool_result, +) + +USER = [Message(role="user", content="hi")] +ROUTE = ToolRoute(server_id="srv", tool_name="read") +TOOLSET = ProviderToolset(tools=[], routes={"mcp__fs__read": ROUTE}) + + +def _loop(adapter: ScriptedAdapter, executor: RecordingExecutor, max_rounds: int = 8) -> ToolLoop: + return ToolLoop(adapter=adapter, executor=executor, max_tool_rounds=max_rounds) + + +async def test_passthrough_without_tool_calls() -> None: + adapter = ScriptedAdapter([final_result("yo")]) + loop = _loop(adapter, RecordingExecutor()) + turn = await loop.run(list(USER), TOOLSET) + assert turn.tool_rounds == 0 + assert turn.messages == [ + *USER, + Message(role="assistant", content="yo"), + ] + assert turn.result.text == "yo" + + +async def test_one_round_multiple_calls() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")), final_result()] + ) + executor = RecordingExecutor() + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert turn.tool_rounds == 1 + assert executor.calls == [("srv", "read", {}), ("srv", "read", {})] + assert turn.messages == [ + *USER, + Message( + role="assistant", + content="", + tool_calls=[tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")], + ), + Message(role="tool", content="ok", tool_call_id="c1", name="mcp__fs__read"), + Message(role="tool", content="ok", tool_call_id="c2", name="mcp__fs__read"), + Message(role="assistant", content="final"), + ] + + +async def test_multi_round_chain() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), tool_result(tool_call("c2", "mcp__fs__read")), final_result()] + ) + turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + assert turn.tool_rounds == 2 + assert [m.role for m in turn.messages] == [ + "user", "assistant", "tool", "assistant", "tool", "assistant", + ] + + +async def test_tools_sent_every_round() -> None: + from octave.inference.types import ToolDefinition + + populated = ProviderToolset( + tools=[ToolDefinition(name="mcp__fs__read")], routes={"mcp__fs__read": ROUTE} + ) + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), final_result()] + ) + await _loop(adapter, RecordingExecutor()).run(list(USER), populated) + assert [call.tools for call in adapter.complete_calls] == [populated.tools] * 2 + + +async def test_empty_toolset_omits_tools() -> None: + adapter = ScriptedAdapter([final_result()]) + await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + # empty toolset -> tools omitted (None), never an empty list + assert adapter.complete_calls[0].tools is None + + +async def test_loop_limit_raises_with_partial_transcript() -> None: + adapter = ScriptedAdapter([tool_result(tool_call("c1", "mcp__fs__read"))] * 3) + loop = _loop(adapter, RecordingExecutor(), max_rounds=2) + with pytest.raises(ToolLoopLimitError) as excinfo: + await loop.run(list(USER), TOOLSET) + assert len(adapter.complete_calls) == 3 + # user + (assistant+tool) x 2 executed rounds + unfulfilled assistant + assert len(excinfo.value.messages) == 1 + 2 * 2 + 1 + assert excinfo.value.messages[-1].role == "assistant" + assert excinfo.value.messages[-1].tool_calls is not None + + +async def test_finish_reason_tool_calls_with_empty_list_is_final() -> None: + from octave.inference.types import CompletionResult + + adapter = ScriptedAdapter( + [CompletionResult(text="done", model="fake", finish_reason="tool_calls", tool_calls=[])] + ) + turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + assert turn.tool_rounds == 0 + assert turn.messages[-1].content == "done" + + +async def test_malformed_arguments_is_error_tool_message() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read", "{not json")), final_result()] + ) + executor = RecordingExecutor() + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert executor.calls == [] + tool_message = turn.messages[2] + assert tool_message.role == "tool" + assert tool_message.tool_call_id == "c1" + assert tool_message.content.startswith("Error: ") + assert "not a valid JSON object" in tool_message.content + + +async def test_non_object_json_arguments_is_error_tool_message() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read", "[1, 2]")), final_result()] + ) + turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + assert turn.messages[2].content.startswith("Error: ") + + +async def test_unknown_route_is_error_tool_message() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__ghost__x")), final_result()] + ) + executor = RecordingExecutor() + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert executor.calls == [] + assert turn.messages[2].content == "Error: Unknown tool: mcp__ghost__x" + + +async def test_error_outcome_continues_loop() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), final_result()] + ) + executor = RecordingExecutor([ToolOutcome(content="boom", is_error=True)]) + turn = await _loop(adapter, executor).run(list(USER), TOOLSET) + assert turn.messages[2].content == "Error: boom" + assert turn.tool_rounds == 1 + + +async def test_adapter_error_propagates_untouched() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read")), AdapterConnectionError("engine down")] + ) + with pytest.raises(AdapterConnectionError): + await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) + + +async def test_arguments_dict_forwarded_verbatim() -> None: + adapter = ScriptedAdapter( + [tool_result(tool_call("c1", "mcp__fs__read", '{"path": "/x", "n": 3}')), final_result()] + ) + executor = RecordingExecutor() + await _loop(adapter, executor).run(list(USER), TOOLSET) + assert executor.calls == [("srv", "read", {"path": "/x", "n": 3})] From d56f7b408327f9e3253580b2a3e2aed7de16aa34 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 13:46:57 +0100 Subject: [PATCH 09/13] feat(agent): add McpToolExecutor over ToolRegistry --- backend/src/octave/agent/__init__.py | 2 + backend/src/octave/agent/mcp_executor.py | 42 ++++++++++++ backend/tests/agent/test_mcp_executor.py | 81 ++++++++++++++++++++++++ 3 files changed, 125 insertions(+) create mode 100644 backend/src/octave/agent/mcp_executor.py create mode 100644 backend/tests/agent/test_mcp_executor.py diff --git a/backend/src/octave/agent/__init__.py b/backend/src/octave/agent/__init__.py index 80b4abf..589de2d 100644 --- a/backend/src/octave/agent/__init__.py +++ b/backend/src/octave/agent/__init__.py @@ -8,10 +8,12 @@ from octave.agent.errors import ToolLoopLimitError from octave.agent.executor import ToolExecutor from octave.agent.loop import ToolLoop +from octave.agent.mcp_executor import McpToolExecutor from octave.agent.types import ToolOutcome, ToolTurn from octave.tools.errors import ToolError __all__ = [ + "McpToolExecutor", "ToolError", "ToolExecutor", "ToolLoop", diff --git a/backend/src/octave/agent/mcp_executor.py b/backend/src/octave/agent/mcp_executor.py new file mode 100644 index 0000000..f4cbf8f --- /dev/null +++ b/backend/src/octave/agent/mcp_executor.py @@ -0,0 +1,42 @@ +"""ToolExecutor over ToolRegistry (issue #79). + +The only octave.agent module importing octave.mcp. Converts McpError +subclasses into error ToolOutcomes (Tier-1, model-correctable) so the +loop never aborts on a tool failure. +""" + +import logging +from typing import Any + +from octave.agent.types import ToolOutcome +from octave.mcp import McpError, ToolRegistry + +__all__ = ["McpToolExecutor"] + +logger = logging.getLogger(__name__) + +_NO_OUTPUT = "(no output)" + + +class McpToolExecutor: + """Structural implementation of ``octave.agent.executor.ToolExecutor``.""" + + def __init__(self, registry: ToolRegistry) -> None: + self._registry = registry + + async def call( + self, server_id: str, tool_name: str, arguments: dict[str, Any] + ) -> ToolOutcome: + """Invoke one tool; every failure mode becomes an error outcome.""" + try: + result = await self._registry.call_tool(server_id, tool_name, arguments) + except McpError as exc: + logger.error( + "MCP tool execution failed | server=%s tool=%s error=%s", + server_id, + tool_name, + exc, + ) + return ToolOutcome(content=f"Tool execution failed: {exc}", is_error=True) + content = "\n".join(block.text for block in result.content) + return ToolOutcome(content=content or _NO_OUTPUT, is_error=result.is_error) diff --git a/backend/tests/agent/test_mcp_executor.py b/backend/tests/agent/test_mcp_executor.py new file mode 100644 index 0000000..ad848c9 --- /dev/null +++ b/backend/tests/agent/test_mcp_executor.py @@ -0,0 +1,81 @@ +"""McpToolExecutor: content flattening and McpError -> outcome conversion.""" + +from typing import Any + +import pytest + +from octave.agent.mcp_executor import McpToolExecutor +from octave.mcp.errors import ( + McpConfigError, + McpConnectionError, + McpNotConnectedError, + McpRpcError, + McpTimeoutError, +) +from octave.mcp.types import ToolContent, ToolResult + + +class StubRegistry: + """Stands in for ToolRegistry; returns or raises what the test dictates.""" + + def __init__( + self, result: ToolResult | None = None, error: Exception | None = None + ) -> None: + self._result = result + self._error = error + self.calls: list[tuple[str, str, dict[str, Any]]] = [] + + async def call_tool( + self, id: str, name: str, arguments: dict[str, Any] | None = None + ) -> ToolResult: + self.calls.append((id, name, dict(arguments or {}))) + if self._error is not None: + raise self._error + assert self._result is not None + return self._result + + +async def test_joins_text_blocks_and_passes_is_error_through() -> None: + registry = StubRegistry( + result=ToolResult( + content=[ToolContent(kind="text", text="a"), ToolContent(kind="text", text="b")] + ) + ) + outcome = await McpToolExecutor(registry).call("srv", "read", {"p": 1}) # type: ignore[arg-type] + assert registry.calls == [("srv", "read", {"p": 1})] + assert outcome.content == "a\nb" + assert outcome.is_error is False + + +async def test_error_result_flags_outcome() -> None: + registry = StubRegistry( + result=ToolResult( + content=[ToolContent(kind="text", text="disk full")], is_error=True + ) + ) + outcome = await McpToolExecutor(registry).call("srv", "read", {}) # type: ignore[arg-type] + assert outcome.content == "disk full" + assert outcome.is_error is True + + +async def test_empty_content_becomes_sentinel() -> None: + registry = StubRegistry(result=ToolResult(content=[])) + outcome = await McpToolExecutor(registry).call("srv", "read", {}) # type: ignore[arg-type] + assert outcome.content == "(no output)" + + +@pytest.mark.parametrize( + "error", + [ + McpRpcError("bad params", code=-32602), + McpTimeoutError("timed out after 60s"), + McpConnectionError("server died"), + McpConfigError("unknown server id"), + McpNotConnectedError("not connected"), + ], +) +async def test_mcp_errors_become_error_outcomes(error: Exception) -> None: + registry = StubRegistry(error=error) + outcome = await McpToolExecutor(registry).call("srv", "read", {}) # type: ignore[arg-type] + assert outcome.is_error is True + assert str(error) in outcome.content From d82ccfd010d2eddcd862292b470b523e532791a5 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Fri, 25 Sep 2026 20:45:30 +0100 Subject: [PATCH 10/13] style: wrap long lines flagged by ruff E501 --- backend/src/octave/agent/loop.py | 9 ++- .../src/octave/inference/openai_adapter.py | 4 +- backend/tests/agent/fakes.py | 4 +- backend/tests/agent/test_loop.py | 33 ++++++++--- backend/tests/agent/test_mcp_executor.py | 5 +- backend/tests/agent/test_types.py | 5 +- .../tests/inference/test_openai_adapter.py | 10 +++- backend/tests/inference/test_types.py | 8 ++- ...6-09-25-issue-concurrent-tool-execution.md | 58 +++++++++++++++++++ 9 files changed, 117 insertions(+), 19 deletions(-) create mode 100644 plans/2026-09-25-issue-concurrent-tool-execution.md diff --git a/backend/src/octave/agent/loop.py b/backend/src/octave/agent/loop.py index f823b76..43a1040 100644 --- a/backend/src/octave/agent/loop.py +++ b/backend/src/octave/agent/loop.py @@ -76,7 +76,9 @@ async def run( messages=history, ) rounds += 1 - logger.debug("tool round | round=%s calls=%s", rounds, len(result.tool_calls)) + logger.debug( + "tool round | round=%s calls=%s", rounds, len(result.tool_calls) + ) for call in result.tool_calls: history.append(await self._execute(call, toolset)) @@ -92,7 +94,10 @@ async def _execute(self, call: ToolCall, toolset: ProviderToolset) -> Message: return self._tool_message( call, ToolOutcome( - content=f"Tool call arguments are not a valid JSON object: {arguments}", + content=( + "Tool call arguments are not a valid JSON object: " + f"{arguments}" + ), is_error=True, ), ) diff --git a/backend/src/octave/inference/openai_adapter.py b/backend/src/octave/inference/openai_adapter.py index 04b6a36..fd8d2d4 100644 --- a/backend/src/octave/inference/openai_adapter.py +++ b/backend/src/octave/inference/openai_adapter.py @@ -169,7 +169,9 @@ async def aclose(self) -> None: def _chat_kwargs(self, request: CompletionRequest) -> dict[str, Any]: kwargs: dict[str, Any] = { "model": self._resolve_model(request.model), - "messages": [self._message_payload(message) for message in request.messages], + "messages": [ + self._message_payload(message) for message in request.messages + ], } if request.temperature is not None: kwargs["temperature"] = request.temperature diff --git a/backend/tests/agent/fakes.py b/backend/tests/agent/fakes.py index 1449332..d15e94a 100644 --- a/backend/tests/agent/fakes.py +++ b/backend/tests/agent/fakes.py @@ -33,7 +33,9 @@ async def complete(self, request: CompletionRequest) -> CompletionResult: raise item return item - async def stream(self, request: CompletionRequest) -> AsyncIterator[CompletionChunk]: + async def stream( + self, request: CompletionRequest + ) -> AsyncIterator[CompletionChunk]: raise NotImplementedError yield CompletionChunk() # satisfy async-generator typing diff --git a/backend/tests/agent/test_loop.py b/backend/tests/agent/test_loop.py index f71744d..dc35672 100644 --- a/backend/tests/agent/test_loop.py +++ b/backend/tests/agent/test_loop.py @@ -21,7 +21,9 @@ TOOLSET = ProviderToolset(tools=[], routes={"mcp__fs__read": ROUTE}) -def _loop(adapter: ScriptedAdapter, executor: RecordingExecutor, max_rounds: int = 8) -> ToolLoop: +def _loop( + adapter: ScriptedAdapter, executor: RecordingExecutor, max_rounds: int = 8 +) -> ToolLoop: return ToolLoop(adapter=adapter, executor=executor, max_tool_rounds=max_rounds) @@ -38,9 +40,8 @@ async def test_passthrough_without_tool_calls() -> None: async def test_one_round_multiple_calls() -> None: - adapter = ScriptedAdapter( - [tool_result(tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")), final_result()] - ) + calls = [tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")] + adapter = ScriptedAdapter([tool_result(*calls), final_result()]) executor = RecordingExecutor() turn = await _loop(adapter, executor).run(list(USER), TOOLSET) assert turn.tool_rounds == 1 @@ -50,7 +51,7 @@ async def test_one_round_multiple_calls() -> None: Message( role="assistant", content="", - tool_calls=[tool_call("c1", "mcp__fs__read"), tool_call("c2", "mcp__fs__read")], + tool_calls=calls, ), Message(role="tool", content="ok", tool_call_id="c1", name="mcp__fs__read"), Message(role="tool", content="ok", tool_call_id="c2", name="mcp__fs__read"), @@ -60,7 +61,11 @@ async def test_one_round_multiple_calls() -> None: async def test_multi_round_chain() -> None: adapter = ScriptedAdapter( - [tool_result(tool_call("c1", "mcp__fs__read")), tool_result(tool_call("c2", "mcp__fs__read")), final_result()] + [ + tool_result(tool_call("c1", "mcp__fs__read")), + tool_result(tool_call("c2", "mcp__fs__read")), + final_result(), + ] ) turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) assert turn.tool_rounds == 2 @@ -105,7 +110,11 @@ async def test_finish_reason_tool_calls_with_empty_list_is_final() -> None: from octave.inference.types import CompletionResult adapter = ScriptedAdapter( - [CompletionResult(text="done", model="fake", finish_reason="tool_calls", tool_calls=[])] + [ + CompletionResult( + text="done", model="fake", finish_reason="tool_calls", tool_calls=[] + ) + ] ) turn = await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) assert turn.tool_rounds == 0 @@ -156,7 +165,10 @@ async def test_error_outcome_continues_loop() -> None: async def test_adapter_error_propagates_untouched() -> None: adapter = ScriptedAdapter( - [tool_result(tool_call("c1", "mcp__fs__read")), AdapterConnectionError("engine down")] + [ + tool_result(tool_call("c1", "mcp__fs__read")), + AdapterConnectionError("engine down"), + ] ) with pytest.raises(AdapterConnectionError): await _loop(adapter, RecordingExecutor()).run(list(USER), TOOLSET) @@ -164,7 +176,10 @@ async def test_adapter_error_propagates_untouched() -> None: async def test_arguments_dict_forwarded_verbatim() -> None: adapter = ScriptedAdapter( - [tool_result(tool_call("c1", "mcp__fs__read", '{"path": "/x", "n": 3}')), final_result()] + [ + tool_result(tool_call("c1", "mcp__fs__read", '{"path": "/x", "n": 3}')), + final_result(), + ] ) executor = RecordingExecutor() await _loop(adapter, executor).run(list(USER), TOOLSET) diff --git a/backend/tests/agent/test_mcp_executor.py b/backend/tests/agent/test_mcp_executor.py index ad848c9..e77ab55 100644 --- a/backend/tests/agent/test_mcp_executor.py +++ b/backend/tests/agent/test_mcp_executor.py @@ -38,7 +38,10 @@ async def call_tool( async def test_joins_text_blocks_and_passes_is_error_through() -> None: registry = StubRegistry( result=ToolResult( - content=[ToolContent(kind="text", text="a"), ToolContent(kind="text", text="b")] + content=[ + ToolContent(kind="text", text="a"), + ToolContent(kind="text", text="b"), + ] ) ) outcome = await McpToolExecutor(registry).call("srv", "read", {"p": 1}) # type: ignore[arg-type] diff --git a/backend/tests/agent/test_types.py b/backend/tests/agent/test_types.py index cc6b092..b0b62ef 100644 --- a/backend/tests/agent/test_types.py +++ b/backend/tests/agent/test_types.py @@ -11,7 +11,10 @@ def test_tool_outcome_defaults() -> None: def test_tool_turn_round_trip() -> None: turn = ToolTurn( - messages=[Message(role="user", content="hi"), Message(role="assistant", content="yo")], + messages=[ + Message(role="user", content="hi"), + Message(role="assistant", content="yo"), + ], result=CompletionResult(text="yo", model="m"), tool_rounds=0, ) diff --git a/backend/tests/inference/test_openai_adapter.py b/backend/tests/inference/test_openai_adapter.py index 463b00b..c11520a 100644 --- a/backend/tests/inference/test_openai_adapter.py +++ b/backend/tests/inference/test_openai_adapter.py @@ -450,7 +450,9 @@ async def test_assistant_tool_calls_message_uses_provider_envelope() -> None: Message( role="assistant", content="", - tool_calls=[ToolCall(id="call_1", name="mcp__fs__read", arguments='{"p": 1}')], + tool_calls=[ + ToolCall(id="call_1", name="mcp__fs__read", arguments='{"p": 1}') + ], ), ] ) @@ -469,7 +471,11 @@ async def test_assistant_tool_calls_message_uses_provider_envelope() -> None: async def test_tool_result_message_carries_tool_call_id_and_name() -> None: body = await _capture_chat_body( - [Message(role="tool", content="42", tool_call_id="call_1", name="mcp__fs__read")] + [ + Message( + role="tool", content="42", tool_call_id="call_1", name="mcp__fs__read" + ) + ] ) assert body["messages"][0] == { "role": "tool", diff --git a/backend/tests/inference/test_types.py b/backend/tests/inference/test_types.py index d4b0871..c56b82e 100644 --- a/backend/tests/inference/test_types.py +++ b/backend/tests/inference/test_types.py @@ -35,7 +35,9 @@ def test_completion_request_extra_is_not_shared() -> None: def test_tool_role_is_valid() -> None: - message = Message(role="tool", content="result", tool_call_id="call_1", name="mcp__fs__read") + message = Message( + role="tool", content="result", tool_call_id="call_1", name="mcp__fs__read" + ) assert message.tool_calls is None @@ -63,7 +65,9 @@ def test_completion_result_tool_calls_defaults_none() -> None: def test_completion_result_round_trips_tool_calls() -> None: call = ToolCall(id="c1", name="x", arguments="{}") - result = CompletionResult(text="", model="m", finish_reason="tool_calls", tool_calls=[call]) + result = CompletionResult( + text="", model="m", finish_reason="tool_calls", tool_calls=[call] + ) assert CompletionResult.model_validate(result.model_dump()) == result diff --git a/plans/2026-09-25-issue-concurrent-tool-execution.md b/plans/2026-09-25-issue-concurrent-tool-execution.md new file mode 100644 index 0000000..948c3f7 --- /dev/null +++ b/plans/2026-09-25-issue-concurrent-tool-execution.md @@ -0,0 +1,58 @@ +# Draft: follow-up work item — concurrent tool execution + +Architect mode cannot run `gh`. Run this from the repo root after the +orchestration-loop PR (#103) is planned/merged. Check open milestones first +(`gh milestone list --repo Svagtlys/Octave`) and substitute ``. + +```bash +gh issue create \ + --repo Svagtlys/Octave \ + --title "feat(agent): execute parallel tool calls concurrently" \ + --label "enhancement" \ + --label "area:agent" \ + --milestone "" \ + --body "$(cat plans/2026-09-25-issue-concurrent-tool-execution-body.md)" +``` + +Then add to the Milestone Tracker project (7) and set status per the +create-work-item skill. + +--- + +# Issue body (start) + +## Description + +The tool-use orchestration loop (`octave.agent.loop`, issue #79) executes +multiple tool calls returned in one completion sequentially, making a turn's +tool latency the sum of all call latencies. Execute them concurrently via an +anyio task group inside a turn, so latency becomes the slowest call. + +## Why + +MCP servers include slow operations (search, scraping, long jobs). A model +emitting 4 fast + 1 slow call pays the full sum today; users watch a stalled +turn. `McpClient` is already safe for concurrent use within one event loop +(SDK request-ID correlation), so the capability is unused. + +## Scope + +- [ ] `octave/agent/loop.py`: fan out tool calls of one turn via + `anyio.create_task_group`; collect `ToolOutcome`s; append `role="tool"` + messages in model call order (deterministic history) +- [ ] Preserve the two-tier error split: per-call failures stay + model-correctable tool messages; one failing call must not cancel + siblings (contain within the task group) +- [ ] Optional `parallel_tool_calls: bool = True` loop option (and pass-through + of the provider field via `CompletionRequest.extra` if desired) +- [ ] Tests: concurrent happy path (fake executor with staggered sleeps), + sibling isolation on failure, ordering of appended tool messages + +## Notes + +- Protocol invariant unchanged: all N results for an assistant `tool_calls` + message must be appended before the next completion. +- Sequential execution ships first in #79; this is the promotion path recorded + in the #79 design spec's deferral table. + +# Issue body (end) From dd605c4fbb9e99491fe430beb16d8387d2ce4855 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Sat, 26 Sep 2026 08:45:50 +0100 Subject: [PATCH 11/13] docs(architecture): document octave.agent orchestration loop --- docs/ARCHITECTURE.md | 15 ++++++++++++++- 1 file changed, 14 insertions(+), 1 deletion(-) diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 482e031..2f4ba02 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -70,7 +70,7 @@ Implements the Model Context Protocol client, enabling Octave to discover, manag - **Connection Lifecycle** — Start and manual restart (`McpClient.restart()`) with `is_connected` liveness; subprocess exit detected via transport-stream monitoring (fail-fast `McpConnectionError`). Auto-restart policy and health monitoring land with the server lifecycle manager (roadmap #4) *(shipped: `octave.mcp.manager` — per-server supervisors, auto-restart with backoff + crash-loop detection, probe-on-timeout health, PR #93)* - **Tool Discovery & Caching** — Fetches tool schemas and descriptions from servers; caches for fast lookup *(shipped: `octave.mcp.registry.ToolRegistry` — fleet-wide inventory cache with event-driven invalidation (restart drift, `tools/list_changed` notifications, warm-up + lazy refresh), PR #98)* - **Tool Execution Engine** — Invokes tools with arguments, handles responses and errors, enforces timeouts *(shipped: `ToolRegistry.call_tool` — `(server_id, tool_name)` surface forwarding to `McpClient.call_tool`; timeout/RPC/connection semantics stay the client's, PR #98)* -- **Schema Translation** — Converts MCP `inputSchema` into the provider-native `tools` array format with `mcp____` dedupe and reverse routing *(shipped: `octave.tools.translate_tools` — pure translation into `octave.inference` `ToolDefinition`s carried on `CompletionRequest.tools`; request-side only, response-side tool-call parsing lands with the agent loop, PR #100)* +- **Schema Translation** — Converts MCP `inputSchema` into the provider-native `tools` array format with `mcp____` dedupe and reverse routing *(shipped: `octave.tools.translate_tools` — pure translation into `octave.inference` `ToolDefinition`s carried on `CompletionRequest.tools`, PR #100; response-side `tool_calls` parsing shipped with the agent loop, PR #103)* - **Configuration Persistence** — Stores server connection configs in the unified database - **Tool Tagging System** — Labels tools with internal Octave tags (e.g., `context_retrieval`, `file_operations`) used by the Context Manager for vault population - **Tool Re-naming / Re-describing** — Maps custom agent-facing names and descriptions to underlying MCP tool identifiers, improving clarity for the agent without modifying the MCP server @@ -126,6 +126,19 @@ Orchestrates agent lifecycles, routes messages between agents and subsystems, an - Sends completed run context to the **Context Manager** for archival - Exposes agent status and results to the **Agent Manager View** in the frontend +**Implemented — tool-use orchestration loop (issue #79):** the `octave.agent` package is +the composition layer and future Agent Manager home. `ToolLoop` drives the +reason → act → observe cycle: detect `CompletionResult.tool_calls`, resolve exposed +names through `ProviderToolset.routes` (#78), execute sequentially via the +`ToolExecutor` Protocol, append provider-invariant tool messages, and re-invoke until a +final answer (round limit `max_tool_rounds` raises `ToolLoopLimitError` with the partial +transcript). `McpToolExecutor` wraps `ToolRegistry`, converting `McpError` failures into +model-correctable error tool messages. `octave.inference` gained the OpenAI-dialect +vocabulary for this (`ToolCall`, `role="tool"`, tool fields on `Message`). Library-only: +no routes/lifespan wiring yet — composition arrives with Integration & Testing #1. The +`openai`/`mcp` SDKs stay quarantined from the package (AST guard). Design: +[`.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md`](../.agents/specs/2026-09-25-tool-use-orchestration-loop-design.md). + --- ## Data Layer From 68e3ed869a04801602777128f3900b24ba755a47 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Sat, 26 Sep 2026 09:25:49 +0100 Subject: [PATCH 12/13] test(agent): add end-to-end smoke script for the orchestration loop --- backend/scripts/smoke_agent_loop.py | 187 ++++++++++++++++++++++++++++ 1 file changed, 187 insertions(+) create mode 100644 backend/scripts/smoke_agent_loop.py diff --git a/backend/scripts/smoke_agent_loop.py b/backend/scripts/smoke_agent_loop.py new file mode 100644 index 0000000..77e0ad2 --- /dev/null +++ b/backend/scripts/smoke_agent_loop.py @@ -0,0 +1,187 @@ +#!/usr/bin/env python +"""Manual smoke test: run the tool-use orchestration loop end-to-end. + +Wires the real seams — a stdio filesystem MCP server through +``ToolRegistry`` -> ``translate_tools`` -> ``ToolLoop`` -> a live +OpenAI-dialect engine — and prompts the model to create a file. Success +criterion: the loop caught at least one tool call and the file exists. + +Engine endpoint comes from the usual ``OCTAVE_INFERENCE_*`` env / ``.env`` +(see ``InferenceSettings``) unless overridden with flags. + +Usage (from backend/): + + # engine from env, temp sandbox, npx filesystem server + uv run scripts/smoke_agent_loop.py + + # explicit engine / directory / prompt + uv run scripts/smoke_agent_loop.py \\ + --base-url http://localhost:11434/v1 --model qwen2.5-coder \\ + --dir /tmp/octave-sandbox \\ + --prompt "Create hello.txt containing 'hi'" + +Requires Node/npx for the filesystem server. Prints the full transcript. +Exits non-zero when no tool call was caught or the file is missing. +Not part of the test suite. +""" + +import argparse +import asyncio +import sys +import tempfile +from pathlib import Path + +from octave.agent import McpToolExecutor, ToolLoop +from octave.inference import ( + AdapterConfig, + InferenceAdapter, + InferenceSettings, + OpenAIAdapter, +) +from octave.inference.errors import AdapterError +from octave.inference.types import Message +from octave.mcp import McpError, McpServerManager, StdioConfig, ToolRegistry +from octave.tools import translate_tools + +DEFAULT_PROMPT = ( + "Create a file named note.txt in the working directory containing " + "exactly the text: hello from octave" +) + + +async def run( + *, config: AdapterConfig, directory: Path, prompt: str, max_rounds: int +) -> int: + """Connect the fleet, run one orchestrated turn, verify the tool call.""" + manager = McpServerManager() + manager.register( + id="fs", + name="filesystem", + config=StdioConfig( + command="npx", + args=["-y", "@modelcontextprotocol/server-filesystem", str(directory)], + ), + ) + registry = ToolRegistry(manager=manager) + adapter: InferenceAdapter = OpenAIAdapter(config) + try: + await manager.start_all() + await registry.start() + inventory = await registry.inventory() + for server in inventory: + print(f"server: {server.server_id} — {len(server.tools)} tool(s)") + toolset = translate_tools(inventory) + print(f"tools: {', '.join(tool.name for tool in toolset.tools)}") + + loop = ToolLoop( + adapter=adapter, + executor=McpToolExecutor(registry), + max_tool_rounds=max_rounds, + ) + turn = await loop.run([Message(role="user", content=prompt)], toolset) + + print(f"\nrounds: {turn.tool_rounds}") + print("--- transcript ---") + tool_call_count = 0 + for message in turn.messages: + if message.tool_calls: + tool_call_count += len(message.tool_calls) + for call in message.tool_calls: + print( + f"assistant -> tool call: {call.name} " + f"arguments={call.arguments}" + ) + elif message.role == "tool": + preview = message.content[:120].replace("\n", " ") + print(f"tool[{message.name}] {preview}") + else: + preview = message.content[:120].replace("\n", " ") + print(f"{message.role}: {preview}") + print("--- final ---") + print(turn.result.text or "(empty)") + + created = directory / "note.txt" + if turn.tool_rounds == 0 or tool_call_count == 0: + print( + "\n✗ no tool call was caught (model answered without tools)", + file=sys.stderr, + ) + return 1 + if not created.exists(): + print( + f"\n✗ tool calls executed but {created} does not exist", + file=sys.stderr, + ) + return 1 + print(f"\n✓ caught {tool_call_count} tool call(s); {created} created") + return 0 + except (McpError, AdapterError) as exc: + print(f"\n✗ {type(exc).__name__}: {exc}", file=sys.stderr) + return 1 + finally: + await registry.stop() + await manager.stop_all() + await adapter.aclose() + + +def main() -> None: + parser = argparse.ArgumentParser( + description="Smoke-test the tool-use orchestration loop end-to-end." + ) + parser.add_argument( + "--base-url", default=None, help="OpenAI-dialect /v1 URL (default: env)" + ) + parser.add_argument( + "--model", default=None, help="Model id for completions (default: env)" + ) + parser.add_argument( + "--dir", + type=Path, + default=None, + help="Directory to sandbox the filesystem server (default: temp dir)", + ) + parser.add_argument("--prompt", default=DEFAULT_PROMPT, help="User message to send") + parser.add_argument("--max-rounds", type=int, default=8, help="Tool round limit") + parsed = parser.parse_args() + + directory: Path = parsed.dir or Path(tempfile.mkdtemp(prefix="octave-smoke-loop-")) + directory.mkdir(parents=True, exist_ok=True) + print(f"sandbox: {directory}") + + settings = InferenceSettings() + config = settings.to_adapter_config() + if parsed.base_url: + config = AdapterConfig( + adapter=config.adapter, + base_url=parsed.base_url, + api_key=config.api_key, + default_model=parsed.model or config.default_model, + timeout_seconds=config.timeout_seconds, + max_retries=config.max_retries, + extra=config.extra, + ) + elif parsed.model: + config = AdapterConfig( + adapter=config.adapter, + base_url=config.base_url, + api_key=config.api_key, + default_model=parsed.model, + timeout_seconds=config.timeout_seconds, + max_retries=config.max_retries, + extra=config.extra, + ) + + sys.exit( + asyncio.run( + run( + config=config, + directory=directory, + prompt=parsed.prompt, + max_rounds=parsed.max_rounds, + ) + ) + ) + + +if __name__ == "__main__": + main() From a76d4207fb058a536522ddc725f61e0e0482572e Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Sat, 26 Sep 2026 15:58:23 +0100 Subject: [PATCH 13/13] fix(agent): wait for supervisor settle and add CA-bundle flag to smoke script --- backend/scripts/smoke_agent_loop.py | 40 +++++++++++++++++++++++++++-- 1 file changed, 38 insertions(+), 2 deletions(-) diff --git a/backend/scripts/smoke_agent_loop.py b/backend/scripts/smoke_agent_loop.py index 77e0ad2..160a0b1 100644 --- a/backend/scripts/smoke_agent_loop.py +++ b/backend/scripts/smoke_agent_loop.py @@ -31,6 +31,9 @@ import tempfile from pathlib import Path +import anyio +import httpx + from octave.agent import McpToolExecutor, ToolLoop from octave.inference import ( AdapterConfig, @@ -50,7 +53,12 @@ async def run( - *, config: AdapterConfig, directory: Path, prompt: str, max_rounds: int + *, + config: AdapterConfig, + directory: Path, + prompt: str, + max_rounds: int, + ca_bundle: str | None = None, ) -> int: """Connect the fleet, run one orchestrated turn, verify the tool call.""" manager = McpServerManager() @@ -63,10 +71,29 @@ async def run( ), ) registry = ToolRegistry(manager=manager) - adapter: InferenceAdapter = OpenAIAdapter(config) + http_client = ( + httpx.AsyncClient(verify=ca_bundle) if ca_bundle else None + ) + adapter: InferenceAdapter = OpenAIAdapter(config, http_client=http_client) try: await manager.start_all() + # The registry's warm-up is a background task: wait until the + # supervisor reports the server actually connected, then refresh + # explicitly so the inventory below is populated. + for _ in range(30): + state = manager.status_of("fs").state + if state in ("connected", "crashed"): + break + await anyio.sleep(2.0) + if state != "connected": + status = manager.status_of("fs") + print( + f"✗ filesystem server state={state} error={status.last_error}", + file=sys.stderr, + ) + return 1 await registry.start() + await registry.refresh("fs") inventory = await registry.inventory() for server in inventory: print(f"server: {server.server_id} — {len(server.tools)} tool(s)") @@ -142,6 +169,14 @@ def main() -> None: ) parser.add_argument("--prompt", default=DEFAULT_PROMPT, help="User message to send") parser.add_argument("--max-rounds", type=int, default=8, help="Tool round limit") + parser.add_argument( + "--ca-bundle", + default=None, + help=( + "CA bundle for TLS verification (default: certifi). " + "For internal CAs try /etc/ssl/certs/ca-certificates.crt" + ), + ) parsed = parser.parse_args() directory: Path = parsed.dir or Path(tempfile.mkdtemp(prefix="octave-smoke-loop-")) @@ -178,6 +213,7 @@ def main() -> None: directory=directory, prompt=parsed.prompt, max_rounds=parsed.max_rounds, + ca_bundle=parsed.ca_bundle, ) ) )