From 42177bc461e0083fb44b30b0bdeea485224c1a12 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Sun, 27 Sep 2026 19:11:48 +0100 Subject: [PATCH 1/9] chore: start work on #27 From 0f7a821439795bc84b0e039eaca463070c0ca480 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 10:38:23 +0100 Subject: [PATCH 2/9] docs(specs): add agent message routing design (#27) --- ...2026-09-28-agent-message-routing-design.md | 412 ++++++++++++++++++ 1 file changed, 412 insertions(+) create mode 100644 .agents/specs/2026-09-28-agent-message-routing-design.md diff --git a/.agents/specs/2026-09-28-agent-message-routing-design.md b/.agents/specs/2026-09-28-agent-message-routing-design.md new file mode 100644 index 0000000..cc9cbd5 --- /dev/null +++ b/.agents/specs/2026-09-28-agent-message-routing-design.md @@ -0,0 +1,412 @@ +# Agent Message Routing — Design + +> Issue: #27 · PR: #107 · Branch: `feature/agent-message-routing` +> Date: 2026-09-28 · Status: Approved in brainstorming +> Depends on: 2026-09-27 Agent Lifecycle Model design (#25, PR #105 — shipped), 2026-09-27 Agent Registry design (#26, PR #106 — shipped) +> Consumers: Integration & Testing #1 (composition + real TurnRunner), Context Manager #5 (archival seam), Agent Manager #6 (scheduling), Agent Manager View UI + +## Problem Statement + +Issue #27 asks for the message router: deliver incoming user messages to the +correct agent instance, broadcast when needed. Multi-agent conversations must +keep context isolated per session. + +The routing model agreed in brainstorming is **session-level turn-taking**: +when U1 says "Hello" in a session with AI1 and AI2, a *next-turn decider* +(an LLM acting as referee) picks which agent responds; that agent's turn runs +to completion and its reply lands in the transcript; the decider then picks +the next agent or hands control back to the user. Repeat until await-user or a +hop limit. + +Two substrate facts shape this work item: + +1. **#25 decided instances carry no context** — "the live context *is* the + session transcript." There are therefore no per-agent queues to deliver + into: delivering to an agent means *its next turn reads the transcript up + to the current seq*. The transcript is the single shared medium; isolation + comes from the session boundary, not from inboxes. +2. **No transcript write path exists.** `events.seq` is documented as + gap-free monotonic and `target_participant_id` (NULL = broadcast) is + FK-backed by `session_participants`, but nothing in the codebase appends + events yet. The router is the first consumer that requires one, so this + work item ships `EventStore` in `octave.db`. + +Library-level scope only — no routes, no WebSocket wiring — consistent with +#25/#26; composition arrives with Integration & Testing #1. + +## Design Decisions (from brainstorming) + +| # | Decision | Rejected alternative | +|---|----------|----------------------| +| 1 | **Full session driver**: `MessageRouter.deliver()` appends the user message, then loops — decider picks → `begin_turn` → turn-runner port → append reply → `end_turn` → decider again — until await-user / hop limit / error | Primitives only (decide-one-step, loop deferred to Integration #1 — leaves mutex interplay and event ordering unproven); deterministic policy only | +| 2 | **Shared-transcript delivery**; `events.target_participant_id` stays NULL. Routing is ephemeral driver state, derivable from event authorship order; the schema keeps the column for future A2A DMs | Per-agent inbox/queue tables (contradicts #25's "context *is* the transcript"); pre-writing targets onto events (couples persistence to an ephemeral decision) | +| 3 | **`EventStore` in `octave.db`** — the first transcript write path: gap-free `seq`, per-kind payload validation, VaultStore transaction convention | The router writing raw ORM `Event` rows (distributes the seq/payload invariant); deferring persistence entirely to Integration #1 (the driver's core observable behavior — "the message landed in the transcript" — would be untested) | +| 4 | **`TurnDecider` protocol + `LlmTurnDecider` default** with a 1:1 fast path (no LLM call in single-agent chat); structured JSON choice (`participant_id` or `await_user`), roster-validated, one retry, await-user fallback | Scored choice (argmax over model-reported scores — unreliable confidence, wider failure surface); hard-coding round-robin (no intelligence for multi-agent) | +| 5 | **Decision models (Laya/Jev) = follow-up issue.** Their declare-a-decision-space API (`choice`/`score`/`noul` → answers + probabilities) is not chat completion; the shipped `InferenceAdapter` ABC is chat-only. `TurnDecider` is the swap seam; the follow-up adds a decision capability to the inference seam (or a sibling adapter) and a `DecisionTurnDecider` | Extending `InferenceAdapter` now (churns the ABC before any decision model is deployed; the seam already isolates the change) | +| 6 | **`TurnRunner` protocol stays separate from `AgentInstanceManager`** — referee vs. player: the manager owns the turn bracket (`begin_turn`/`end_turn`, `octave.db`-only imports); the runner is what happens inside it (context assembly, `ToolLoop`, inference, MCP). Mirrors the shipped `ToolExecutor`/`McpToolExecutor` split | Folding execution into the manager (pulls inference + MCP + CM imports into the one component #25 kept DB-only; every execution-policy change edits a shipped lifecycle component; routing tests would drag in fake adapters) | +| 7 | **`run_record` archival deferred to CM #5**; the driver exposes `TurnRecord(instance_id, agent_id, seq_range)` per completed turn — the exact surface #25's archival contract demands. CM #5 plugs its mechanism (callback on the turn boundary) into this seam | Archiving from the router now (chunking policy and embedding provenance are CM territory per the 2026-09-20/21 ADRs; `octave.db` must not import `octave.inference`) | +| 8 | **Hop limit `max_agent_turns` (default 4)** bounds agent turns per user message | Unlimited turns (agent-chatter loops); treating scheduling as policy here (#6 owns priority/queueing; this is a guard rail only) | + +## Delivery Model + +```mermaid +sequenceDiagram + participant U as Caller + participant R as MessageRouter + participant D as TurnDecider + participant M as AgentInstanceManager + participant T as TurnRunner + U->>R: deliver session_id author content + R->>R: append user_message event + loop until stop + R->>D: decide roster + transcript tail + alt AWAIT_USER + D-->>R: await user + R-->>U: RouteOutcome stop AWAIT_USER + else picks agent instance + D-->>R: candidate instance + R->>M: begin_turn instance_id + R->>T: run_turn instance messages + T-->>R: reply text + R->>R: append assistant_message event + R->>M: end_turn instance_id + end + end +``` + +No per-agent queues exist. When the decider picks AI1, "AI1's message plus +everything queued for it" is precisely the transcript through the current +`seq` — which the turn runner reads as its input. U1's message, AI1's reply, +and any system events are all in it when AI2's turn is claimed. Context +isolation follows from the session row: instances in other sessions are +never in this transcript or this roster. + +## EventStore (`octave.db`) + +The canonical transcript write path. VaultStore convention: constructed with +the caller's `AsyncSession`, **never commits**; callers own transaction +boundaries (`octave.db.deps`). + +```python +class EventStore: + def __init__(self, session: AsyncSession) -> None: ... + + async def append( + self, + session_id: str, + kind: EventKind, + *, + author_participant_id: str | None = None, + target_participant_id: str | None = None, + payload: dict[str, Any], + ) -> Event: ... + + async def read( + self, session_id: str, *, after_seq: int = 0 + ) -> list[Event]: ... +``` + +- **seq assignment**: `SELECT COALESCE(MAX(seq), 0) + 1` scoped to the + session, then INSERT; `uq_events_session_seq` is the backstop. The flush + runs inside `session.begin_nested()` (SAVEPOINT) so an `IntegrityError` + from a racing append is recoverable without poisoning the caller's + transaction; retry re-computes seq, up to 2 retries, then raise. Gap-free + monotonic per the `events.seq` docstring. +- **Payload validation**: `user_message` → `UserMessagePayload`, + `assistant_message` → `AssistantMessagePayload` (Pydantic-validated, + stored via `model_dump`); `tool_call` / `tool_result` / `system` pass + through as dicts (the `octave.db.types` module docstring assigns their + validation to consumers). A malformed validated-kind payload raises + `ValidationError` — a write-path bug, fail loud. +- **`read`**: `seq > after_seq`, ordered `seq ASC`. The driver's + transcript-so-far source. +- `octave.db` never imports `octave.inference` (2026-09-21 ADR holds): + `EventStore` takes/returns dicts and ORM rows; the transcript→`Message` + mapping lives in the agent plane. + +## Driver (`octave.agent.router`) + +```python +class MessageRouter: + def __init__( + self, + *, + session: AsyncSession, + manager: AgentInstanceManager, + registry: AgentRegistry, + decider: TurnDecider, + turn_runner: TurnRunner, + max_agent_turns: int = 4, + ) -> None: ... + + async def deliver( + self, session_id: str, *, author_participant_id: str, content: str + ) -> RouteOutcome: ... +``` + +Collaborators are constructor-injected (composition root wires them in +Integration #1); `EventStore` is constructed internally over the same +session. + +### Loop semantics + +1. **Author gate**: `author_participant_id` must be a current member of the + session (`session_participants` row with `left_at IS NULL`); otherwise + `NotAMemberError`. (The events FK covers author *identity*; membership is + the app-level gate.) +2. **Append** the `user_message` event (`UserMessagePayload`), author = the + gate-passed participant. +3. **Roster**: `AgentRegistry.list_instances(session_id=...)` filtered to + `instance_status == IDLE` and `definition_status == ACTIVE`, each joined + to its `Participant` row (guaranteed by `spawn`) for identity. Plus the + `AWAIT_USER` option. Empty roster → stop `AWAIT_USER` immediately. +4. **Decide** via `TurnDecider` over `DecisionState` (roster + transcript + tail: last `decider_tail_events` events, default 30). +5. **Claim** with `begin_turn` (the #25 mutex; re-checks the pause gate). + `TurnInProgressError` propagates — concurrent drivers on one session is a + caller bug; the mutex is fail-loud by design. +6. **Run**: map transcript → `list[Message]` (`user_message`→`user`, + `assistant_message`→`assistant`; `tool_*`/`system` skipped with a debug + log), hand to `TurnRunner.run_turn`, receive reply text. +7. **Record**: append `assistant_message` (author = the instance's + participant), `end_turn`, append a `TurnRecord` to the outcome with the + completed `(session_id, agent_id, seq_range)`. +8. Repeat from 4 until stop. + +**Error policy**: a `turn_runner` exception → `end_turn` still releases the +instance (#25: turn failure lands idle; the error is an event, not state) → +append a `system` event with the error text → stop `ERROR`. A decider choice +outside the roster → `DeciderChoiceError` → one retry → fallback +`AWAIT_USER` (a confused referee hands control to the human; never crashes +the session, never picks a phantom). + +**Hop limit**: `max_agent_turns` bounds *agent* turns per `deliver()` call; +exhausted → stop `HOP_LIMIT`. Checked before claiming, so the limit never +truncates a turn mid-flight. + +### Outcome types + +```python +class StopReason(StrEnum): + AWAIT_USER = "await_user" + HOP_LIMIT = "hop_limit" + ERROR = "error" + +@dataclass(frozen=True) +class TurnRecord: + session_id: str + agent_id: str + instance_id: str + seq_start: int + seq_end: int + """Events contributed by this turn (``assistant_message`` now; tool + events when Integration #1's real runner appends them inside the + bracket). CM #5's archival key per #25: (session_id, agent_id, + seq_range) — never instance_id.""" + +@dataclass(frozen=True) +class RouteOutcome: + turns: list[TurnRecord] + stop_reason: StopReason + error: str | None = None +``` + +`TurnRecord.seq_range` is the archival seam #25 mandated: CM #5's +`run_record` append rides this boundary (likely an `on_turn_complete` +callback added in CM #5), keyed `(session_id, agent_id, seq_range)` — never +`instance_id`, per #25. + +## Decider (`octave.agent.decider`) + +```python +@dataclass(frozen=True) +class Candidate: + instance_id: str + participant_id: str + label: str + +@dataclass(frozen=True) +class DecisionState: + roster: list[Candidate] + messages: list[Message] # transcript tail, mapped + +class Decision(StrEnum): + """Sentinel vocabulary for decider output. A decision is a plain + ``str``: either ``Decision.AWAIT_USER`` or a candidate's + ``participant_id`` (StrEnum members are strs, so the union collapses + cleanly).""" + + AWAIT_USER = "await_user" + +class TurnDecider(Protocol): + async def decide(self, state: DecisionState) -> str: + """Return ``Decision.AWAIT_USER`` or a participant_id from + ``state.roster``. Deciders raise ``DeciderChoiceError`` on + malformed backend output; the driver owns retry + fallback.""" +``` + +### `LlmTurnDecider` (default, ships now) + +Constructed with an `InferenceAdapter`, `model: str | None` (which model is +"the system one" is wiring policy for Integration #1), and the tail budget. + +- **1:1 fast path**: roster of exactly one candidate and no LLM call — + if the last transcript event is the user's, pick the candidate; if the + sole agent just spoke, `AWAIT_USER`. Keeps the common chat case + deterministic, zero-latency, and terminating. +- **Multi-agent**: one `complete()` — prompt presents the roster + (participant_id + label each) and the transcript tail, and requires + exactly `{"next": "" | "await_user"}`. Parse → Pydantic + validate → check membership in the roster. Malformed or out-of-roster → + `DeciderChoiceError` (driver retries once, then falls back `AWAIT_USER`). +- `AdapterError` propagates (the driver's decider-failure path applies). + +### Decision models (Laya/Jev) — follow-up, not #27 + +A decision model's API — declare a decision space (`choice` / `score` / +`noul` with criteria), receive answers with per-option probabilities — is +not expressible through `complete(CompletionRequest)`. Locally hosted, it +still belongs behind the inference seam, but as a new **capability**: the +follow-up issue decides between a `decision()` method on `InferenceAdapter` +(+ registry/config plumbing) and a sibling adapter seam, then ships a +`DecisionTurnDecider` implementing `TurnDecider`. The driver never learns +which backend chose; `TurnDecider` is the swap seam this spec guarantees. + +## TurnRunner port (`octave.agent.router`) + +```python +class TurnRunner(Protocol): + async def run_turn( + self, *, instance: RunningAgent, messages: list[Message] + ) -> str: ... +``` + +The real runner — context assembly (CM territory), `resolve_model`, +`ToolLoop` over the instance's model binding and toolset — is Integration +#1's composition, mirroring the shipped `ToolExecutor`/`McpToolExecutor` +split. #27 tests inject a scripted fake. The protocol returns final text; +`ToolLoop` already returns exactly that (`ToolTurn.result.text`). + +## Error Handling + +- `NotAMemberError(RoutingError)` — author is not a current session member. +- `DeciderChoiceError(RoutingError)` — decider output malformed/out-of-roster; + raised by deciders, caught by the driver (retry → `AWAIT_USER` fallback). +- `RoutingError(AgentError)` — new base in `octave.agent.errors`, alongside + the #25 lifecycle-gate errors. No new error for "no instances": that is a + normal `AWAIT_USER` stop, not a failure. +- `TurnInProgressError`, `AdapterError`, `InstanceNotFoundError` propagate + unchanged from existing owners. +- Transaction discipline: `EventStore` and the driver never commit. One + `deliver()` call is one caller-owned transaction — the user message, agent + turns, and replies commit atomically or not at all. (Crash mid-loop + therefore rolls back to the last turn boundary; reconcile (#25) plus the + transcript's gap-free seq make the residue consistent.) + +## Non-Goals + +- REST/WebSocket routes, frontend wiring (Integration #1 composes; the + current `/ws` echo stub is untouched). +- `run_record`/`run_summary` archival (CM #5; the seam is `TurnRecord`). +- Decision-model adapter + `DecisionTurnDecider` (follow-up issue below). +- A2A direct messages / non-NULL `target_participant_id` (column stays as + shipped; broadcast semantics throughout). +- Priority, queueing, resource constraints (#6); session `turn_policy` + column (the decider *is* the policy — in code, swappable, per-turn). +- Streaming agent replies into the transcript (assistant_message is + appended whole at turn end; streaming is Integration #1/UI territory). +- Multi-user session vault-visibility concerns (2026-09-13 ADR follow-up). + +## Testing + +Real SQLite via the shared `session_factory` fixture; scripted fakes for +`TurnRunner`/`TurnDecider`/`InferenceAdapter`; spawns driven through +`AgentInstanceManager`; seed helpers modeled on +`tests/agent/test_instances.py`. + +**EventStore** +- seq starts at 1; gap-free and monotonic across interleaved appends to two + sessions; `read(after_seq=N)` windows correctly. +- Concurrent append race (two sessions appending the same session) → retry + path yields distinct gap-free seqs; retry exhaustion raises. +- `user_message`/`assistant_message` payloads validated (malformed → + `ValidationError`); `tool_call`/`tool_result`/`system` dicts pass through. +- Never commits: a fresh session with no writes flushed sees appended rows + only after the caller commits; DB byte-identical on read-only use. + +**MessageRouter** +- 1:1 chat end-to-end: deliver → user_message event → one agent turn → + assistant_message event → stop `AWAIT_USER`; turn mutex held during the + runner call (fake asserts instance is `active` mid-turn). +- Multi-agent scripted decider (the issue's scenario): U1 → AI1 → AI2 → + await_user; transcript order and authorship correct; AI2's runner input + contains U1's message *and* AI1's reply. +- Empty roster → immediate `AWAIT_USER`, no runner call. +- Runner raises → instance back to `idle`, `system` event recorded, stop + `ERROR` with message; user_message event survives. +- Hop limit: decider always picks an agent → stop `HOP_LIMIT` after + `max_agent_turns`; all claimed turns properly ended (no leaked `active`). +- Decider returns out-of-roster twice → fallback `AWAIT_USER`, no crash. +- Author not a member (or `left_at` set) → `NotAMemberError`, no event + appended. +- Transcript mapping skips `tool_*`/`system` events with debug log. +- Post-loop registry consistency: `AgentRegistry.count_by_status` shows no + `active` instances left behind. + +**LlmTurnDecider** +- Fast path both directions (user spoke last → pick sole candidate; agent + spoke last → `AWAIT_USER`); zero adapter calls on the fast path. +- Multi-agent: fake adapter returns valid JSON → candidate id; malformed → + `DeciderChoiceError`; out-of-roster id → `DeciderChoiceError`; + `AdapterError` propagates. + +**Package** — `MessageRouter`, `TurnDecider`, `LlmTurnDecider`, +`TurnRunner`, `RouteOutcome`, `TurnRecord`, `StopReason`, `Candidate`, +`DecisionState`, `Decision` exported from `octave.agent`; +`EventStore` from `octave.db`; `test_package.py` updated. SDK-quarantine AST +guard unaffected (no `openai`/`mcp` imports). + +## Implementation Sketch + +Paths for the implementation plan to refine: + +| File | Change | +|---|---| +| `backend/src/octave/db/event_store.py` | new: `EventStore` (append with gap-free seq + SAVEPOINT retry, read) | +| `backend/src/octave/db/__init__.py` | export `EventStore` | +| `backend/src/octave/agent/decider.py` | new: `Candidate`, `DecisionState`, `Decision`, `TurnDecider`, `LlmTurnDecider` | +| `backend/src/octave/agent/router.py` | new: `TurnRunner`, `StopReason`, `TurnRecord`, `RouteOutcome`, `MessageRouter`, transcript→Message mapping | +| `backend/src/octave/agent/errors.py` | `RoutingError`, `NotAMemberError`, `DeciderChoiceError` | +| `backend/src/octave/agent/__init__.py` | export router/decider surface | +| `backend/tests/db/test_event_store.py` | new: EventStore tests | +| `backend/tests/agent/test_router.py` | new: driver tests | +| `backend/tests/agent/test_decider.py` | new: LlmTurnDecider tests | +| `backend/tests/agent/test_package.py` | assert new public names | +| `docs/TODO.md` | mark Agent Manager #3 done | + +## Consequences + +- `octave.db` gains its first transcript write path; every future event + writer (Integration #1's `tool_call`/`tool_result` recording, CM system + events) goes through `EventStore` — the seq invariant has exactly one + owner. +- `octave.agent` gains driver vocabulary; Integration #1's job shrinks to + composition: wire a real `TurnRunner` + adapter + session deps behind + `MessageRouter.deliver()`. +- CM #5 inherits the archival seam #25 required (`TurnRecord` seq ranges at + the turn boundary) — no rework when archival lands. +- The inference seam is untouched; the decision-model capability question is + isolated behind `TurnDecider` with a follow-up issue. +- `events.target_participant_id` remains NULL everywhere it is written; the + A2A-DM door stays open, unopened. +- No schema change, no migration — purely additive Python, green migration + chain. + +## Follow-up Issues to File + +1. **Decision-model adapter for turn decisions (Laya/Jev-style).** A + decision-space API (`choice`/`score`/`noul` → answers + probabilities) + hosted behind the inference seam: decide `decision()` on + `InferenceAdapter` vs. sibling adapter; config/registry plumbing for a + locally hosted decision model; ship `DecisionTurnDecider` implementing + `TurnDecider`. Reference: Laya's declare-space-in / answers-out shape. From 18c9705dac045606ba8623868fa601a0beefc932 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:50:43 +0100 Subject: [PATCH 3/9] fix(db): explicit BEGIN transaction control for SQLite engines (#27) Under pysqlite's default isolation_level="", SAVEPOINT opens an implicit transaction and RELEASE SAVEPOINT commits it, so session.begin_nested() silently commits the caller's whole pending transaction. EventStore's gap-free seq retry (this work item) depends on savepoints that nest and roll back without committing, and the plan's test_never_commits caught the violation. Applies SQLAlchemy's documented pysqlite transaction-control recipe via a shared attach_transaction_control(engine) helper wired into SqliteVecAdapter.make_engine() and the test conftest engine. Verified through the async_creator path: savepoint collisions roll back cleanly, outer transactions survive, and close-without-commit rolls back. Deviation from plan 2026-09-28-agent-message-routing: the plan did not anticipate that its savepoint-retry design is incompatible with the engine's default pysqlite transaction mode. --- backend/src/octave/db/_bootstrap.py | 31 ++++++++++++++++++++++++- backend/src/octave/db/sqlite_adapter.py | 6 +++-- backend/tests/conftest.py | 9 ++++++- 3 files changed, 42 insertions(+), 4 deletions(-) diff --git a/backend/src/octave/db/_bootstrap.py b/backend/src/octave/db/_bootstrap.py index 88604ef..3070a89 100644 --- a/backend/src/octave/db/_bootstrap.py +++ b/backend/src/octave/db/_bootstrap.py @@ -17,9 +17,11 @@ import aiosqlite import sqlite_vec +from sqlalchemy import event from sqlalchemy.engine import URL +from sqlalchemy.ext.asyncio import AsyncEngine -__all__ = ["make_async_creator"] +__all__ = ["attach_transaction_control", "make_async_creator"] logger = logging.getLogger(__name__) @@ -43,3 +45,30 @@ async def _connect() -> aiosqlite.Connection: return conn return _connect + + +def attach_transaction_control(engine: AsyncEngine) -> None: + """Give SQLite connections explicit transaction control (issue #27). + + Under pysqlite's default ``isolation_level=""``, ``SAVEPOINT`` opens an + implicit transaction and ``RELEASE SAVEPOINT`` *commits it* — so + ``session.begin_nested()`` (EventStore's seq-collision retry) silently + commits the caller's whole pending transaction, breaking the + "stores never commit; callers own the transaction boundary" convention. + + SQLAlchemy's documented remedy for pysqlite, verified to work through + the ``async_creator`` path: disable the DBAPI's implicit-BEGIN behavior + (``isolation_level = None``) and emit an explicit ``BEGIN`` whenever + SQLAlchemy starts a root transaction. Savepoints then nest, roll back + to their savepoint, and release without committing. + + Call exactly once per engine, immediately after ``create_async_engine``. + """ + + @event.listens_for(engine.sync_engine, "connect") + def _disable_implicit_begin(dbapi_conn: Any, _conn_record: Any) -> None: + dbapi_conn.isolation_level = None + + @event.listens_for(engine.sync_engine, "begin") + def _emit_explicit_begin(conn: Any) -> None: + conn.exec_driver_sql("BEGIN") diff --git a/backend/src/octave/db/sqlite_adapter.py b/backend/src/octave/db/sqlite_adapter.py index a2a0cfa..24a629f 100644 --- a/backend/src/octave/db/sqlite_adapter.py +++ b/backend/src/octave/db/sqlite_adapter.py @@ -23,7 +23,7 @@ create_async_engine, ) -from octave.db._bootstrap import make_async_creator +from octave.db._bootstrap import attach_transaction_control, make_async_creator from octave.db.adapter import DbAdapter from octave.db.errors import DbDimensionMismatchError, DbError from octave.db.registry import register_db @@ -45,7 +45,9 @@ class SqliteVecAdapter(DbAdapter): def make_engine(self) -> AsyncEngine: creator = make_async_creator(make_url(self._config.url)) - return create_async_engine("sqlite+aiosqlite://", async_creator=creator) + engine = create_async_engine("sqlite+aiosqlite://", async_creator=creator) + attach_transaction_control(engine) + return engine def make_session_factory( self, engine: AsyncEngine diff --git a/backend/tests/conftest.py b/backend/tests/conftest.py index 1240d21..539e73e 100644 --- a/backend/tests/conftest.py +++ b/backend/tests/conftest.py @@ -23,6 +23,7 @@ create_async_engine, ) +from octave.db._bootstrap import attach_transaction_control from octave.db.models import Base @@ -40,11 +41,17 @@ async def _connect() -> aiosqlite.Connection: @pytest_asyncio.fixture async def engine(tmp_path: Path) -> AsyncIterator[AsyncEngine]: - """Fresh SQLite file per test (tests stay independent per coding rules).""" + """Fresh SQLite file per test (tests stay independent per coding rules). + + ``attach_transaction_control`` matches the production engine wiring + (``SqliteVecAdapter.make_engine``): explicit-BEGIN transaction control + so savepoints do not act as commits (see ``_bootstrap``). + """ eng = create_async_engine( "sqlite+aiosqlite://", async_creator=_make_async_creator(str(tmp_path / "test.db")), ) + attach_transaction_control(eng) async with eng.begin() as conn: await conn.run_sync(Base.metadata.create_all) From 9d7d09de83d5b4466aecedc352e5caec19c0d55e Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:51:44 +0100 Subject: [PATCH 4/9] feat(db): add EventStore transcript write path (#27) --- backend/src/octave/db/__init__.py | 2 + backend/src/octave/db/event_store.py | 110 ++++++++++++++++++ backend/tests/db/test_event_store.py | 165 +++++++++++++++++++++++++++ 3 files changed, 277 insertions(+) create mode 100644 backend/src/octave/db/event_store.py create mode 100644 backend/tests/db/test_event_store.py diff --git a/backend/src/octave/db/__init__.py b/backend/src/octave/db/__init__.py index 6d0228e..ab012ac 100644 --- a/backend/src/octave/db/__init__.py +++ b/backend/src/octave/db/__init__.py @@ -12,6 +12,7 @@ from octave.db.adapter import DbAdapter from octave.db.config import DatabaseSettings, DbConfig from octave.db.errors import DbError +from octave.db.event_store import EventStore from octave.db.lifespan import db_lifespan from octave.db.migrations import current, upgrade from octave.db.models import Base @@ -30,6 +31,7 @@ "DbError", "DatabaseSettings", "EventKind", + "EventStore", "SqliteVecAdapter", "VaultHit", "VaultKind", diff --git a/backend/src/octave/db/event_store.py b/backend/src/octave/db/event_store.py new file mode 100644 index 0000000..e2cc3a0 --- /dev/null +++ b/backend/src/octave/db/event_store.py @@ -0,0 +1,110 @@ +"""Transcript write path: append/read events with gap-free seq (issue #27). + +``events.seq`` is documented as per-session monotonic (gap-free); clock +resolution cannot guarantee this, so seq assignment has exactly one owner. +Follows the VaultStore convention — constructed with the caller's +AsyncSession, never commits; callers own transaction boundaries +(``octave.db.deps``). The uq_events_session_seq constraint is the backstop +for the seq race: the insert runs in a SAVEPOINT so a collision is +recoverable without poisoning the caller's transaction. + +Payload validation: ``user_message``/``assistant_message`` validate against +their ``octave.db.types`` models; ``tool_call``/``tool_result``/``system`` +pass through as dicts (validation assigned to consumers, per the types +module docstring). +""" + +import uuid +from typing import Any + +from pydantic import BaseModel +from sqlalchemy import func, select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncSession + +from octave.db.models import Event +from octave.db.types import ( + AssistantMessagePayload, + EventKind, + UserMessagePayload, +) + +__all__ = ["EventStore"] + +_PAYLOAD_MODELS: dict[EventKind, type[BaseModel]] = { + EventKind.USER_MESSAGE: UserMessagePayload, + EventKind.ASSISTANT_MESSAGE: AssistantMessagePayload, +} + +_RETRY_LIMIT = 2 + + +class EventStore: + """The canonical event append/read surface. Never commits.""" + + def __init__(self, session: AsyncSession) -> None: + self._session = session + + async def append( + self, + session_id: str, + kind: EventKind, + *, + author_participant_id: str | None = None, + target_participant_id: str | None = None, + payload: dict[str, Any], + ) -> Event: + """Append one event, assigning the next gap-free ``seq``. + + Malformed payloads for validated kinds raise ``ValidationError`` + (a write-path bug — fail loud). On a ``seq`` collision the insert + retries with a freshly computed seq up to ``_RETRY_LIMIT`` times, + then re-raises. Cross-connection effectiveness is engine-dependent + (SQLite snapshots stale reads within a transaction); the constraint + guarantees no duplicate ever lands either way. + """ + validator = _PAYLOAD_MODELS.get(kind) + if validator is not None: + payload = validator.model_validate(payload).model_dump() + last_error: IntegrityError | None = None + for _ in range(_RETRY_LIMIT + 1): + seq = await self._next_seq(session_id) + event = Event( + id=uuid.uuid4().hex, + session_id=session_id, + seq=seq, + kind=str(kind), + author_participant_id=author_participant_id, + target_participant_id=target_participant_id, + payload=payload, + ) + try: + async with self._session.begin_nested(): + self._session.add(event) + await self._session.flush() + except IntegrityError as exc: + last_error = exc + continue + return event + assert last_error is not None + raise last_error + + async def read(self, session_id: str, *, after_seq: int = 0) -> list[Event]: + """Transcript window ``seq > after_seq``, ordered.""" + rows = await self._session.execute( + select(Event) + .where(Event.session_id == session_id, Event.seq > after_seq) + .order_by(Event.seq) + ) + return list(rows.scalars().all()) + + async def _next_seq(self, session_id: str) -> int: + """Next seq for this session (1-based, gap-free). Autoflush makes + same-session prior appends visible; overridable seam for retry + testing.""" + current = await self._session.execute( + select(func.coalesce(func.max(Event.seq), 0)).where( + Event.session_id == session_id + ) + ) + return int(current.scalar_one()) + 1 diff --git a/backend/tests/db/test_event_store.py b/backend/tests/db/test_event_store.py new file mode 100644 index 0000000..8c5124f --- /dev/null +++ b/backend/tests/db/test_event_store.py @@ -0,0 +1,165 @@ +"""EventStore: gap-free transcript appends (design spec 2026-09-28). + +Real SQLite via session_factory; the uq_events_session_seq constraint is +the backstop for the seq race, exercised via a scripted stale-read subclass. +""" + +import pytest +from pydantic import ValidationError +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +import octave.db as db_pkg +from octave.db.event_store import EventStore +from octave.db.models import Participant, Session, User +from octave.db.types import EventKind + + +async def _seed(session_factory: async_sessionmaker[AsyncSession]) -> str: + """User u_1, sessions s_1/s_2, and a user participant (author FK).""" + async with session_factory() as session: + session.add(User(id="u_1", display_name="Alice")) + session.add(Session(id="s_1", created_by_user_id="u_1", status="active")) + session.add(Session(id="s_2", created_by_user_id="u_1", status="active")) + participant = Participant(id="p_u1", user_id="u_1", label="Alice") + session.add(participant) + await session.commit() + return "p_u1" + + +async def _append( + session_factory: async_sessionmaker[AsyncSession], + session_id: str, + author_id: str, + content: str, +) -> int: + async with session_factory() as session: + event = await EventStore(session).append( + session_id, + EventKind.USER_MESSAGE, + author_participant_id=author_id, + payload={"content": content}, + ) + await session.commit() + return event.seq + + +async def test_seq_starts_at_one_and_increments( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + assert await _append(session_factory, "s_1", author, "one") == 1 + assert await _append(session_factory, "s_1", author, "two") == 2 + + +async def test_seq_scoped_per_session( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + assert await _append(session_factory, "s_1", author, "a") == 1 + assert await _append(session_factory, "s_2", author, "b") == 1 + assert await _append(session_factory, "s_1", author, "c") == 2 + + +async def test_read_orders_and_windows( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + for text in ("one", "two", "three"): + await _append(session_factory, "s_1", author, text) + async with session_factory() as session: + store = EventStore(session) + everything = await store.read("s_1") + tail = await store.read("s_1", after_seq=2) + assert [e.seq for e in everything] == [1, 2, 3] + assert [e.payload["content"] for e in everything] == ["one", "two", "three"] + assert [e.seq for e in tail] == [3] + + +async def test_validates_user_message_payload( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + async with session_factory() as session: + with pytest.raises(ValidationError): + await EventStore(session).append( + "s_1", EventKind.USER_MESSAGE, author_participant_id=author, payload={} + ) + + +async def test_assistant_payload_defaults_model_name_null( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + async with session_factory() as session: + event = await EventStore(session).append( + "s_1", + EventKind.ASSISTANT_MESSAGE, + author_participant_id=author, + payload={"content": "hi"}, + ) + assert event.payload == {"content": "hi", "model_name": None} + + +async def test_tool_and_system_kinds_pass_through( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + async with session_factory() as session: + await _seed(session_factory) + event = await EventStore(session).append( + "s_1", EventKind.TOOL_CALL, payload={"server_id": "s", "tool": "t"} + ) + assert event.payload == {"server_id": "s", "tool": "t"} + + +async def test_never_commits( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + async with session_factory() as session: + await EventStore(session).append( + "s_1", EventKind.SYSTEM, payload={"content": "uncommitted"} + ) + async with session_factory() as other: + assert await EventStore(other).read("s_1") == [] # not committed yet + async with session_factory() as session: + await EventStore(session).append( + "s_1", EventKind.SYSTEM, payload={"content": "committed"} + ) + await session.commit() + async with session_factory() as other: + rows = await EventStore(other).read("s_1") + assert [e.payload["content"] for e in rows] == ["committed"] + + +class _StaleReadStore(EventStore): + """Scripts stale seq reads to exercise the IntegrityError retry path: + real cross-connection races are invisible to a stale SQLite snapshot, + so the retry loop is tested at its seam.""" + + def __init__(self, session: AsyncSession, stale: list[int]) -> None: + super().__init__(session) + self._stale = stale + + async def _next_seq(self, session_id: str) -> int: + if self._stale: + return self._stale.pop(0) + return await super()._next_seq(session_id) + + +async def test_retry_recovers_from_seq_conflict( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + async with session_factory() as winner: + await EventStore(winner).append("s_1", EventKind.SYSTEM, payload={"c": 1}) + await winner.commit() # committed seq=1 + async with session_factory() as loser: + store = _StaleReadStore(loser, stale=[1]) # stale read -> collides + event = await store.append("s_1", EventKind.SYSTEM, payload={"c": 2}) + assert event.seq == 2 # retry recomputed from fresh max + await loser.commit() + + +async def test_exported_from_package() -> None: + assert hasattr(db_pkg, "EventStore") + assert "EventStore" in db_pkg.__all__ From b4f7326869ac3983c90fd21a3281296b7169b2f7 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 13:59:13 +0100 Subject: [PATCH 5/9] feat(agent): add routing error hierarchy (#27) --- backend/src/octave/agent/__init__.py | 6 ++++++ backend/src/octave/agent/errors.py | 17 +++++++++++++++++ backend/tests/agent/test_package.py | 3 +++ 3 files changed, 26 insertions(+) diff --git a/backend/src/octave/agent/__init__.py b/backend/src/octave/agent/__init__.py index fc937c0..d1d5595 100644 --- a/backend/src/octave/agent/__init__.py +++ b/backend/src/octave/agent/__init__.py @@ -11,10 +11,13 @@ AgentError, AgentNotFoundError, AgentPausedError, + DeciderChoiceError, InstanceExistsError, InstanceNotFoundError, InvalidTransitionError, ModelBindingError, + NotAMemberError, + RoutingError, SessionNotFoundError, TerminalSessionError, ToolLoopLimitError, @@ -34,12 +37,15 @@ "AgentNotFoundError", "AgentPausedError", "AgentRegistry", + "DeciderChoiceError", "InstanceExistsError", "InstanceNotFoundError", "InvalidTransitionError", "McpToolExecutor", "ModelBindingError", + "NotAMemberError", "ResolvedModel", + "RoutingError", "RunningAgent", "SessionNotFoundError", "TerminalSessionError", diff --git a/backend/src/octave/agent/errors.py b/backend/src/octave/agent/errors.py index db14a00..ad2d9fa 100644 --- a/backend/src/octave/agent/errors.py +++ b/backend/src/octave/agent/errors.py @@ -8,10 +8,13 @@ "AgentError", "AgentNotFoundError", "AgentPausedError", + "DeciderChoiceError", "InstanceExistsError", "InstanceNotFoundError", "InvalidTransitionError", "ModelBindingError", + "NotAMemberError", + "RoutingError", "SessionNotFoundError", "TerminalSessionError", "ToolLoopLimitError", @@ -59,6 +62,20 @@ class ModelBindingError(AgentError): """Binding missing, malformed, or (for tag form) unresolvable.""" +class RoutingError(AgentError): + """Base for router failures (design spec 2026-09-28).""" + + +class NotAMemberError(RoutingError): + """Author is not a current member of the session (left_at set or + never invited).""" + + +class DeciderChoiceError(RoutingError): + """Decider output malformed or outside the roster. Raised by deciders; + the driver owns retry + await-user fallback.""" + + class ToolLoopLimitError(ToolError): """max_tool_rounds exhausted without a final answer. diff --git a/backend/tests/agent/test_package.py b/backend/tests/agent/test_package.py index 7fa6201..27e08cf 100644 --- a/backend/tests/agent/test_package.py +++ b/backend/tests/agent/test_package.py @@ -32,12 +32,15 @@ def test_public_names_are_exported() -> None: "AgentNotFoundError", "AgentPausedError", "AgentRegistry", + "DeciderChoiceError", "InstanceExistsError", "InstanceNotFoundError", "InvalidTransitionError", "McpToolExecutor", "ModelBindingError", + "NotAMemberError", "ResolvedModel", + "RoutingError", "RunningAgent", "SessionNotFoundError", "TerminalSessionError", From 7b8892dd21de74be7ccb97c57951b98baa5821f2 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 14:08:09 +0100 Subject: [PATCH 6/9] feat(agent): add TurnDecider protocol and LlmTurnDecider (#27) --- backend/src/octave/agent/decider.py | 121 ++++++++++++++++++++++++++++ backend/tests/agent/test_decider.py | 89 ++++++++++++++++++++ 2 files changed, 210 insertions(+) create mode 100644 backend/src/octave/agent/decider.py create mode 100644 backend/tests/agent/test_decider.py diff --git a/backend/src/octave/agent/decider.py b/backend/src/octave/agent/decider.py new file mode 100644 index 0000000..837b8ac --- /dev/null +++ b/backend/src/octave/agent/decider.py @@ -0,0 +1,121 @@ +"""Turn-taking decisions (design spec 2026-09-28). + +The referee, not a player: a ``TurnDecider`` reads the roster + transcript +tail and names the next speaker (or hands control to the human). The +driver owns retry and fallback policy; deciders only raise +``DeciderChoiceError`` on unusable backend output. + +``LlmTurnDecider`` is the default strategy. Decision models (Laya/Jev +style) arrive behind this same protocol in a follow-up issue — the driver +never learns which backend chose. +""" + +import json +from dataclasses import dataclass, field +from enum import StrEnum +from typing import Protocol + +from pydantic import BaseModel, ValidationError + +from octave.agent.errors import DeciderChoiceError +from octave.inference.adapter import InferenceAdapter +from octave.inference.types import CompletionRequest, Message + +__all__ = ["Candidate", "Decision", "DecisionState", "LlmTurnDecider", "TurnDecider"] + + +@dataclass(frozen=True) +class Candidate: + """One agent eligible to take the next turn.""" + + instance_id: str + participant_id: str + label: str + + +@dataclass(frozen=True) +class DecisionState: + """Decider input: the candidate roster + the transcript tail.""" + + roster: list[Candidate] = field(default_factory=list) + messages: list[Message] = field(default_factory=list) + + +class Decision(StrEnum): + """Sentinel vocabulary for decider output. A decision is a plain + ``str``: either ``Decision.AWAIT_USER`` or a candidate's + ``participant_id`` (StrEnum members are strs, so the union collapses + cleanly).""" + + AWAIT_USER = "await_user" + + +class TurnDecider(Protocol): + """Next-speaker resolution. Implementations raise + ``DeciderChoiceError`` on malformed/out-of-roster output.""" + + async def decide(self, state: DecisionState) -> str: + """Return ``Decision.AWAIT_USER`` or a participant_id from + ``state.roster``.""" + ... + + +class _Choice(BaseModel): + """Wire shape of the multi-agent decision: exactly one JSON object.""" + + next: str + + +class LlmTurnDecider: + """Default strategy over the chat ``InferenceAdapter`` (local-first). + + 1:1 fast path (one candidate) requires no LLM call: user spoke last → + the candidate speaks; the sole agent spoke last → await the human. + Multi-agent asks for ``{"next": "" | "await_user"}``, + roster-validated; malformed or phantom choices raise + ``DeciderChoiceError`` (the driver retries once, then awaits user). + """ + + def __init__( + self, *, adapter: InferenceAdapter, model: str | None = None + ) -> None: + self._adapter = adapter + self._model = model + + async def decide(self, state: DecisionState) -> str: + if len(state.roster) == 1: + if state.messages and state.messages[-1].role == "assistant": + return Decision.AWAIT_USER + return state.roster[0].participant_id + if not state.roster: + return Decision.AWAIT_USER + result = await self._adapter.complete( + CompletionRequest(model=self._model, messages=self._prompt(state)) + ) + try: + choice = _Choice.model_validate(json.loads(result.text)).next + except (json.JSONDecodeError, ValidationError) as exc: + raise DeciderChoiceError( + f"decider output not parseable: {result.text!r}" + ) from exc + if choice == Decision.AWAIT_USER: + return Decision.AWAIT_USER + if any(candidate.participant_id == choice for candidate in state.roster): + return choice + raise DeciderChoiceError(f"decider chose unknown candidate: {choice!r}") + + @staticmethod + def _prompt(state: DecisionState) -> list[Message]: + roster = "\n".join( + f"- {candidate.participant_id}: {candidate.label}" + for candidate in state.roster + ) + system = ( + "You are the turn referee for a multi-agent conversation. " + "Decide who speaks next: one candidate participant_id, or " + "await_user to hand control to the human. Base the choice on " + "the conversation below. Reply with exactly one JSON object: " + '{"next": "" | "await_user"}\n\n' + f"Candidates:\n{roster}" + ) + return [Message(role="system", content=system), *state.messages] diff --git a/backend/tests/agent/test_decider.py b/backend/tests/agent/test_decider.py new file mode 100644 index 0000000..163d10f --- /dev/null +++ b/backend/tests/agent/test_decider.py @@ -0,0 +1,89 @@ +"""LlmTurnDecider: 1:1 fast path + roster-validated JSON choice (spec #27).""" + +import pytest + +from octave.agent.decider import Candidate, Decision, DecisionState, LlmTurnDecider +from octave.agent.errors import DeciderChoiceError +from octave.inference.errors import AdapterConnectionError +from octave.inference.types import Message +from tests.agent.fakes import ScriptedAdapter, final_result + + +def _candidate(iid: str, pid: str) -> Candidate: + return Candidate(instance_id=iid, participant_id=pid, label=f"Agent {iid}") + + +def _state(*messages: Message, roster: list[Candidate] | None = None) -> DecisionState: + return DecisionState( + roster=list(roster) if roster else [_candidate("i_1", "p_ai1")], + messages=list(messages), + ) + + +async def test_fast_path_picks_sole_candidate_when_user_spoke_last() -> None: + adapter = ScriptedAdapter([]) # any call fails the test via queue assert + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello")) + assert await decider.decide(state) == "p_ai1" + assert adapter.complete_calls == [] + + +async def test_fast_path_awaits_user_when_agent_spoke_last() -> None: + adapter = ScriptedAdapter([]) + decider = LlmTurnDecider(adapter=adapter) + state = _state( + Message(role="user", content="Hello"), + Message(role="assistant", content="Hi"), + ) + assert await decider.decide(state) == Decision.AWAIT_USER + assert adapter.complete_calls == [] + + +async def test_fast_path_empty_transcript_picks_candidate() -> None: + adapter = ScriptedAdapter([]) + decider = LlmTurnDecider(adapter=adapter) + assert await decider.decide(_state()) == "p_ai1" + + +async def test_multi_agent_valid_choice() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result('{"next": "p_ai2"}')]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + assert await decider.decide(state) == "p_ai2" + assert len(adapter.complete_calls) == 1 + + +async def test_multi_agent_await_user_choice() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result('{"next": "await_user"}')]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + assert await decider.decide(state) == Decision.AWAIT_USER + + +async def test_multi_agent_malformed_output_raises() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result("AI1 should go!")]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + with pytest.raises(DeciderChoiceError): + await decider.decide(state) + + +async def test_multi_agent_out_of_roster_choice_raises() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result('{"next": "p_ghost"}')]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + with pytest.raises(DeciderChoiceError): + await decider.decide(state) + + +async def test_multi_agent_adapter_error_propagates() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([AdapterConnectionError("engine down")]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + with pytest.raises(AdapterConnectionError): + await decider.decide(state) From 3b1c279edacc15b2a4c17a699f175ba4e8f8de70 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 14:28:08 +0100 Subject: [PATCH 7/9] feat(agent): add session-level MessageRouter driver (#27) --- backend/src/octave/agent/router.py | 273 +++++++++++++++++++++++++++ backend/tests/agent/test_router.py | 294 +++++++++++++++++++++++++++++ 2 files changed, 567 insertions(+) create mode 100644 backend/src/octave/agent/router.py create mode 100644 backend/tests/agent/test_router.py diff --git a/backend/src/octave/agent/router.py b/backend/src/octave/agent/router.py new file mode 100644 index 0000000..3b63960 --- /dev/null +++ b/backend/src/octave/agent/router.py @@ -0,0 +1,273 @@ +"""Session-level message routing (design spec 2026-09-28, issue #27). + +The driver: append the user message, then loop — decider picks an idle +instance, ``begin_turn`` claims the turn (the #25 mutex), the injected +``TurnRunner`` port produces the reply, the reply is recorded, the turn is +released — until await-user, the hop limit, or a runner failure. + +Delivery is via the shared transcript (spec decision 2): there are no +per-agent queues; an agent's "message plus everything queued for it" is +the transcript through the current seq, which the runner reads as input. +Routing itself is ephemeral — ``events.target_participant_id`` stays NULL. + +Follows the VaultStore convention: constructed with the caller's +AsyncSession, never commits. One ``deliver()`` is one caller-owned +transaction; the whole exchange commits atomically or not at all. +""" + +import logging +from dataclasses import dataclass, field +from enum import StrEnum +from typing import Protocol + +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from octave.agent.decider import Candidate, Decision, DecisionState, TurnDecider +from octave.agent.errors import DeciderChoiceError, NotAMemberError +from octave.agent.instances import AgentInstanceManager +from octave.agent.registry import AgentRegistry, RunningAgent +from octave.db.event_store import EventStore +from octave.db.models import Event, Participant, SessionParticipant +from octave.db.types import AgentStatus, EventKind, InstanceStatus +from octave.inference.types import Message + +__all__ = [ + "MessageRouter", + "RouteOutcome", + "StopReason", + "TurnRecord", + "TurnRunner", +] + +logger = logging.getLogger(__name__) + + +class StopReason(StrEnum): + """Why ``deliver()`` returned control to the caller.""" + + AWAIT_USER = "await_user" + HOP_LIMIT = "hop_limit" + ERROR = "error" + + +@dataclass(frozen=True) +class TurnRecord: + """One completed agent turn. Events contributed by this turn + (``assistant_message`` now; tool events when Integration #1's real + runner appends them inside the bracket). CM #5's archival key per #25: + (session_id, agent_id, seq_range) — never instance_id.""" + + session_id: str + agent_id: str + instance_id: str + seq_start: int + seq_end: int + + +@dataclass(frozen=True) +class RouteOutcome: + turns: list[TurnRecord] = field(default_factory=list) + stop_reason: StopReason = StopReason.AWAIT_USER + error: str | None = None + + +class TurnRunner(Protocol): + """What happens inside a claimed turn. The real runner (context + assembly + ToolLoop over the instance's model binding) is Integration + #1's composition — mirrors the ToolExecutor/McpToolExecutor split. + Returns the final reply text.""" + + async def run_turn( + self, *, instance: RunningAgent, messages: list[Message] + ) -> str: ... + + +class MessageRouter: + """Session-level turn-taking driver. Never commits; the caller owns + the transaction boundary.""" + + def __init__( + self, + *, + session: AsyncSession, + manager: AgentInstanceManager, + registry: AgentRegistry, + decider: TurnDecider, + turn_runner: TurnRunner, + max_agent_turns: int = 4, + decider_tail: int = 30, + ) -> None: + self._session = session + self._manager = manager + self._registry = registry + self._decider = decider + self._turn_runner = turn_runner + self._max_agent_turns = max_agent_turns + self._decider_tail = decider_tail + self._events = EventStore(session) + + async def deliver( + self, session_id: str, *, author_participant_id: str, content: str + ) -> RouteOutcome: + """Record the user message and drive agent turns until stop.""" + await self._require_membership(session_id, author_participant_id) + await self._events.append( + session_id, + EventKind.USER_MESSAGE, + author_participant_id=author_participant_id, + payload={"content": content}, + ) + turns: list[TurnRecord] = [] + agent_turns = 0 + while True: + roster = await self._roster(session_id) + if not roster: + return RouteOutcome(turns=turns, stop_reason=StopReason.AWAIT_USER) + if agent_turns >= self._max_agent_turns: + return RouteOutcome(turns=turns, stop_reason=StopReason.HOP_LIMIT) + events = await self._events.read(session_id) + messages = _transcript_messages(events) + state = DecisionState( + roster=roster, messages=messages[-self._decider_tail :] + ) + choice = await self._decide_with_fallback(state) + if choice == Decision.AWAIT_USER: + return RouteOutcome(turns=turns, stop_reason=StopReason.AWAIT_USER) + candidate = next(c for c in roster if c.participant_id == choice) + outcome = await self._run_agent_turn( + session_id, candidate, messages, turns + ) + if outcome is not None: + return outcome + agent_turns += 1 + + async def _run_agent_turn( + self, + session_id: str, + candidate: Candidate, + messages: list[Message], + turns: list[TurnRecord], + ) -> RouteOutcome | None: + """Claim, run, record, release. None = turn completed normally; + a RouteOutcome = the loop must stop (runner failure).""" + # TurnInProgressError propagates: concurrent drivers on one + # session is a caller bug; the mutex is fail-loud by design. + await self._manager.begin_turn(candidate.instance_id) + running = await self._registry.get_instance(candidate.instance_id) + try: + reply_text = await self._turn_runner.run_turn( + instance=running, messages=messages + ) + except Exception as exc: # noqa: BLE001 — turn failures land idle + logger.warning( + "agent turn failed | instance=%s error=%s", candidate.instance_id, exc + ) + await self._manager.end_turn(candidate.instance_id) + await self._events.append( + session_id, + EventKind.SYSTEM, + payload={"content": f"Agent turn failed: {exc}"}, + ) + return RouteOutcome( + turns=turns, stop_reason=StopReason.ERROR, error=str(exc) + ) + reply = await self._events.append( + session_id, + EventKind.ASSISTANT_MESSAGE, + author_participant_id=candidate.participant_id, + payload={"content": reply_text}, + ) + await self._manager.end_turn(candidate.instance_id) + turns.append( + TurnRecord( + session_id=session_id, + agent_id=running.agent_id, + instance_id=candidate.instance_id, + seq_start=reply.seq, + seq_end=reply.seq, + ) + ) + return None + + async def _decide_with_fallback(self, state: DecisionState) -> str: + """One retry on DeciderChoiceError, then AWAIT_USER: a confused + referee hands control to the human — never crashes the session.""" + for attempt in (0, 1): + try: + return await self._decider.decide(state) + except DeciderChoiceError as exc: + logger.warning( + "decider choice invalid | attempt=%s error=%s", attempt, exc + ) + return Decision.AWAIT_USER + + async def _roster(self, session_id: str) -> list[Candidate]: + """Idle instances of active definitions, joined to their + participant identity (guaranteed by spawn; a missing participant + row is write-path corruption — surface it loud).""" + running = [ + record + for record in await self._registry.list_instances( + session_id=session_id, status=InstanceStatus.IDLE + ) + if record.definition_status == AgentStatus.ACTIVE + ] + if not running: + return [] + rows = await self._session.execute( + select(Participant.agent_id, Participant.id).where( + Participant.agent_id.in_([record.agent_id for record in running]) + ) + ) + # agent_id is nullable at the column level (participant supertype); + # a NULL here would be write-path corruption, so drop it silently — + # the dict lookup below then surfaces the missing identity loudly. + participant_by_agent: dict[str, str] = { + agent_id: participant_id + for agent_id, participant_id in rows.all() + if agent_id is not None + } + return [ + Candidate( + instance_id=record.instance_id, + participant_id=participant_by_agent[record.agent_id], + label=record.agent_name, + ) + for record in running + ] + + async def _require_membership( + self, session_id: str, participant_id: str + ) -> None: + """App-level author gate: current membership (the events FK covers + identity, not membership).""" + row = await self._session.execute( + select(SessionParticipant).where( + SessionParticipant.session_id == session_id, + SessionParticipant.participant_id == participant_id, + SessionParticipant.left_at.is_(None), + ) + ) + if row.scalar_one_or_none() is None: + raise NotAMemberError( + f"participant {participant_id} is not a current member of {session_id}" + ) + + +def _transcript_messages(events: list[Event]) -> list[Message]: + """Transcript → chat messages. Tool/system entries are skipped with a + debug log: they are context for humans and archival, not chat history.""" + messages: list[Message] = [] + for event in events: + if event.kind == EventKind.USER_MESSAGE: + messages.append(Message(role="user", content=event.payload["content"])) + elif event.kind == EventKind.ASSISTANT_MESSAGE: + messages.append( + Message(role="assistant", content=event.payload["content"]) + ) + else: + logger.debug( + "skipping non-message event | kind=%s seq=%s", event.kind, event.seq + ) + return messages diff --git a/backend/tests/agent/test_router.py b/backend/tests/agent/test_router.py new file mode 100644 index 0000000..f97141b --- /dev/null +++ b/backend/tests/agent/test_router.py @@ -0,0 +1,294 @@ +"""MessageRouter driver (design spec 2026-09-28). + +Real SQLite via session_factory; spawns driven through AgentInstanceManager; +scripted decider/runner fakes (no LLM). Transcript assertions read events +back through EventStore on a fresh session. +""" + +import pytest +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from octave.agent import AgentInstanceManager, AgentRegistry +from octave.agent.decider import Decision, DecisionState +from octave.agent.errors import DeciderChoiceError, NotAMemberError +from octave.agent.router import MessageRouter, StopReason +from octave.db.event_store import EventStore +from octave.db.models import Agent, Participant, Session, SessionParticipant, User +from octave.db.types import InstanceStatus +from octave.inference.types import Message + +_TAG_BINDING = {"kind": "tag", "tag": "quick"} + + +async def _seed(session_factory: async_sessionmaker[AsyncSession]) -> None: + async with session_factory() as session: + session.add(User(id="u_1", display_name="Alice")) + session.add(Agent(id="a_1", name="Echo", model_binding=_TAG_BINDING)) + session.add(Agent(id="a_2", name="Second", model_binding=_TAG_BINDING)) + session.add(Session(id="s_1", created_by_user_id="u_1", status="active")) + await session.commit() + + +async def _spawn( + session_factory: async_sessionmaker[AsyncSession], agent_id: str +) -> str: + async with session_factory() as session: + instance = await AgentInstanceManager(session).spawn( + agent_id=agent_id, session_id="s_1" + ) + await session.commit() + return instance.id + + +async def _invite_user(session_factory: async_sessionmaker[AsyncSession]) -> str: + """User participant + live membership in s_1.""" + async with session_factory() as session: + participant = Participant(id="p_u1", user_id="u_1", label="Alice") + session.add(participant) + await session.flush() + session.add( + SessionParticipant(session_id="s_1", participant_id="p_u1") + ) + await session.commit() + return "p_u1" + + +async def _participant_of( + session_factory: async_sessionmaker[AsyncSession], agent_id: str +) -> str: + async with session_factory() as session: + row = await session.execute( + select(Participant.id).where(Participant.agent_id == agent_id) + ) + return row.scalar_one() + + +class ScriptedDecider: + """Pops queued choices (Exception items raise); records each state.""" + + def __init__(self, choices: list[str]) -> None: + self._choices = list(choices) + self.states: list[DecisionState] = [] + + async def decide(self, state: DecisionState) -> str: + self.states.append(state) + item = self._choices.pop(0) + if isinstance(item, Exception): + raise item + return item + + +class RecordingRunner: + """Returns queued replies (Exception items raise); records inputs.""" + + def __init__(self, replies: list[str]) -> None: + self._replies = list(replies) + self.calls: list[tuple[str, list[Message]]] = [] + + async def run_turn(self, *, instance, messages: list[Message]) -> str: # type: ignore[no-untyped-def] + self.calls.append((instance.instance_id, list(messages))) + item = self._replies.pop(0) + if isinstance(item, Exception): + raise item + return item + + +def _router( + session: AsyncSession, + *, + decider: ScriptedDecider, + runner: RecordingRunner, + max_agent_turns: int = 4, +) -> MessageRouter: + return MessageRouter( + session=session, + manager=AgentInstanceManager(session), + registry=AgentRegistry(session), + decider=decider, + turn_runner=runner, + max_agent_turns=max_agent_turns, + ) + + +async def _deliver( + session_factory: async_sessionmaker[AsyncSession], + *, + decider: ScriptedDecider, + runner: RecordingRunner, + author: str = "p_u1", + max_agent_turns: int = 4, +): + async with session_factory() as session: + router = _router( + session, decider=decider, runner=runner, max_agent_turns=max_agent_turns + ) + outcome = await router.deliver( + "s_1", author_participant_id=author, content="Hello" + ) + await session.commit() + return outcome + + +async def _transcript( + session_factory: async_sessionmaker[AsyncSession], +) -> list[tuple[int, str, str | None, str]]: + async with session_factory() as session: + events = await EventStore(session).read("s_1") + return [ + (e.seq, e.kind, e.author_participant_id, str(e.payload.get("content", ""))) + for e in events + ] + + +async def test_one_to_one_end_to_end( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + p_ai1 = await _participant_of(session_factory, "a_1") + await _invite_user(session_factory) + decider = ScriptedDecider([p_ai1, Decision.AWAIT_USER]) + runner = RecordingRunner(["Hi there"]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.AWAIT_USER + assert len(outcome.turns) == 1 + assert outcome.turns[0].agent_id == "a_1" + assert outcome.turns[0].seq_start == outcome.turns[0].seq_end == 2 + assert await _transcript(session_factory) == [ + (1, "user_message", "p_u1", "Hello"), + (2, "assistant_message", p_ai1, "Hi there"), + ] + + +async def test_multi_agent_handoff( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _spawn(session_factory, "a_2") + p_ai1 = await _participant_of(session_factory, "a_1") + p_ai2 = await _participant_of(session_factory, "a_2") + await _invite_user(session_factory) + decider = ScriptedDecider([p_ai1, p_ai2, Decision.AWAIT_USER]) + runner = RecordingRunner(["from AI1", "from AI2"]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert [t.agent_id for t in outcome.turns] == ["a_1", "a_2"] + # AI2's runner input contains U1's message AND AI1's reply + assert runner.calls[1][1] == [ + Message(role="user", content="Hello"), + Message(role="assistant", content="from AI1"), + ] + kinds = [(seq, kind) for seq, kind, _, _ in await _transcript(session_factory)] + assert kinds == [ + (1, "user_message"), + (2, "assistant_message"), + (3, "assistant_message"), + ] + + +async def test_empty_roster_awaits_user_immediately( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _invite_user(session_factory) + decider = ScriptedDecider([]) # never consulted + runner = RecordingRunner([]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.AWAIT_USER + assert outcome.turns == [] + assert runner.calls == [] + assert await _transcript(session_factory) == [(1, "user_message", "p_u1", "Hello")] + + +async def test_runner_failure_releases_instance( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _invite_user(session_factory) + p_ai1 = await _participant_of(session_factory, "a_1") + decider = ScriptedDecider([p_ai1]) + runner = RecordingRunner([RuntimeError("boom")]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.ERROR + assert outcome.error == "boom" + assert outcome.turns == [] # failed turn contributes no TurnRecord + transcript = await _transcript(session_factory) + assert [(seq, kind) for seq, kind, _, _ in transcript] == [ + (1, "user_message"), + (2, "system"), + ] + assert "boom" in transcript[1][3] + async with session_factory() as session: + counts = await AgentRegistry(session).count_by_status() + assert counts[InstanceStatus.ACTIVE] == 0 # turn released + + +async def test_hop_limit_stops_agent_chatter( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _invite_user(session_factory) + p_ai1 = await _participant_of(session_factory, "a_1") + decider = ScriptedDecider([p_ai1]) # one decision; loop stops at limit + runner = RecordingRunner(["first"]) + outcome = await _deliver( + session_factory, decider=decider, runner=runner, max_agent_turns=1 + ) + assert outcome.stop_reason == StopReason.HOP_LIMIT + assert len(outcome.turns) == 1 + + +async def test_decider_invalid_falls_back_to_await_user( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _invite_user(session_factory) + decider = ScriptedDecider( + [DeciderChoiceError("bad1"), DeciderChoiceError("bad2")] + ) + runner = RecordingRunner([]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.AWAIT_USER + assert outcome.turns == [] + assert runner.calls == [] + + +async def test_non_member_author_rejected( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + async with session_factory() as session: + router = _router( + session, decider=ScriptedDecider([]), runner=RecordingRunner([]) + ) + try: + with pytest.raises(NotAMemberError): + await router.deliver( + "s_1", author_participant_id="p_ghost", content="Hello" + ) + finally: + await session.rollback() + assert await _transcript(session_factory) == [] # nothing appended + + +async def test_registry_consistent_after_loop( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _spawn(session_factory, "a_2") + p_ai1 = await _participant_of(session_factory, "a_1") + p_ai2 = await _participant_of(session_factory, "a_2") + await _invite_user(session_factory) + decider = ScriptedDecider([p_ai1, p_ai2, Decision.AWAIT_USER]) + runner = RecordingRunner(["one", "two"]) + await _deliver(session_factory, decider=decider, runner=runner) + async with session_factory() as session: + counts = await AgentRegistry(session).count_by_status() + assert counts[InstanceStatus.ACTIVE] == 0 + assert counts[InstanceStatus.IDLE] == 2 From 3c36eb3c2abe170f8384de781de4edcdf073d400 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 14:36:55 +0100 Subject: [PATCH 8/9] feat(agent): export message routing surface, tick Agent Manager #3 (#27) --- backend/src/octave/agent/__init__.py | 31 +++++++++++++++++++++++++--- backend/tests/agent/test_package.py | 10 +++++++++ docs/TODO.md | 2 +- 3 files changed, 39 insertions(+), 4 deletions(-) diff --git a/backend/src/octave/agent/__init__.py b/backend/src/octave/agent/__init__.py index d1d5595..bea39c0 100644 --- a/backend/src/octave/agent/__init__.py +++ b/backend/src/octave/agent/__init__.py @@ -2,11 +2,19 @@ (design spec #25). Composition layer — the only package importing both octave.mcp and -octave.inference, plus octave.db for the instance manager. Future Agent -Manager components (router, result collector) land here as additional -modules; they do not exist yet. +octave.inference, plus octave.db for the instance manager. Manager +components ship here as additional modules: instance lifecycle +(#25), registry (#26), message routing (#27). Result collection (#4) and +scheduling (#6) remain. """ +from octave.agent.decider import ( + Candidate, + Decision, + DecisionState, + LlmTurnDecider, + TurnDecider, +) from octave.agent.errors import ( AgentError, AgentNotFoundError, @@ -28,6 +36,13 @@ from octave.agent.loop import ToolLoop from octave.agent.mcp_executor import McpToolExecutor from octave.agent.registry import AgentRegistry, RunningAgent +from octave.agent.router import ( + MessageRouter, + RouteOutcome, + StopReason, + TurnRecord, + TurnRunner, +) from octave.agent.types import ToolOutcome, ToolTurn from octave.tools.errors import ToolError @@ -37,17 +52,24 @@ "AgentNotFoundError", "AgentPausedError", "AgentRegistry", + "Candidate", "DeciderChoiceError", + "Decision", + "DecisionState", "InstanceExistsError", "InstanceNotFoundError", "InvalidTransitionError", + "LlmTurnDecider", "McpToolExecutor", + "MessageRouter", "ModelBindingError", "NotAMemberError", "ResolvedModel", + "RouteOutcome", "RoutingError", "RunningAgent", "SessionNotFoundError", + "StopReason", "TerminalSessionError", "ToolError", "ToolExecutor", @@ -55,6 +77,9 @@ "ToolLoopLimitError", "ToolOutcome", "ToolTurn", + "TurnDecider", "TurnInProgressError", + "TurnRecord", + "TurnRunner", "resolve_model", ] diff --git a/backend/tests/agent/test_package.py b/backend/tests/agent/test_package.py index 27e08cf..ef4e9a5 100644 --- a/backend/tests/agent/test_package.py +++ b/backend/tests/agent/test_package.py @@ -32,17 +32,24 @@ def test_public_names_are_exported() -> None: "AgentNotFoundError", "AgentPausedError", "AgentRegistry", + "Candidate", "DeciderChoiceError", + "Decision", + "DecisionState", "InstanceExistsError", "InstanceNotFoundError", "InvalidTransitionError", + "LlmTurnDecider", "McpToolExecutor", + "MessageRouter", "ModelBindingError", "NotAMemberError", "ResolvedModel", + "RouteOutcome", "RoutingError", "RunningAgent", "SessionNotFoundError", + "StopReason", "TerminalSessionError", "ToolError", "ToolExecutor", @@ -50,7 +57,10 @@ def test_public_names_are_exported() -> None: "ToolLoopLimitError", "ToolOutcome", "ToolTurn", + "TurnDecider", "TurnInProgressError", + "TurnRecord", + "TurnRunner", "resolve_model", ): assert hasattr(agent_pkg, name), name diff --git a/docs/TODO.md b/docs/TODO.md index 5312b5c..64332fa 100644 --- a/docs/TODO.md +++ b/docs/TODO.md @@ -106,7 +106,7 @@ - [x] 1. Design agent lifecycle model (definition active/paused gate; instance spawn → idle ⇄ active → destroy) — PR #105 (lifecycle vocabulary, `agent_instances` table, `AgentInstanceManager`; router/registry/scheduling remain #2–6) - [x] 2. Build agent registry (track running agents, their IDs, status, and assigned context) — PR #106 (read-only `AgentRegistry` over `agent_instances ⋈ agents`; no new table) -- [ ] 3. Implement agent message routing (deliver messages to correct agent, broadcast when needed) +- [x] 3. Implement agent message routing (deliver messages to correct agent, broadcast when needed) — PR #107 (session-level `MessageRouter` driver, `TurnDecider`/`TurnRunner` ports, `EventStore` transcript write path; LLM decider default, decision-model decider deferred to follow-up) - [ ] 4. Create agent result collection (capture agent outputs and make them queryable) - [ ] 5. Build inter-agent result sharing (allow agents to request and receive results from other agents) - [ ] 6. Add agent priority and scheduling (queue management, resource constraints) From a72c126c9a71c6c7b793c2ea9896dde9aa9f19a3 Mon Sep 17 00:00:00 2001 From: Momo the Bestest <45446348+Svagtlys@users.noreply.github.com> Date: Mon, 28 Sep 2026 16:18:37 +0100 Subject: [PATCH 9/9] docs: document message routing and explicit transaction control (#27) --- .../specs/2026-09-28-agent-message-routing.md | 1508 +++++++++++++++++ docs/ARCHITECTURE.md | 24 + 2 files changed, 1532 insertions(+) create mode 100644 .agents/specs/2026-09-28-agent-message-routing.md diff --git a/.agents/specs/2026-09-28-agent-message-routing.md b/.agents/specs/2026-09-28-agent-message-routing.md new file mode 100644 index 0000000..dd7a207 --- /dev/null +++ b/.agents/specs/2026-09-28-agent-message-routing.md @@ -0,0 +1,1508 @@ +# Agent Message Routing 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:** Ship the session-level `MessageRouter` — turn-taking driver with an LLM `TurnDecider`, a `TurnRunner` port, and `EventStore`, the first transcript write path with gap-free `seq` (issue #27). + +**Architecture:** `octave.db.event_store.EventStore` owns transcript appends (gap-free per-session `seq`, per-kind payload validation, VaultStore transaction convention). `octave.agent.router.MessageRouter.deliver()` appends the user message, then loops: decider picks an idle instance → `begin_turn` (the #25 mutex) → injected `TurnRunner` port → append `assistant_message` → `end_turn` → repeat until `AWAIT_USER` / `HOP_LIMIT` / `ERROR`. `octave.agent.decider.LlmTurnDecider` is the default strategy (1:1 fast path = zero LLM calls). Library-level only — no routes; Integration #1 wires the real runner. + +**Tech Stack:** Python 3.13, SQLAlchemy 2.x async ORM (`select`, `func.max`, `begin_nested` savepoints), Pydantic payload models + `TypeAdapter`-free direct model validation, `InferenceAdapter` over `ScriptedAdapter` fakes, pytest + pytest-asyncio against real SQLite (shared `session_factory` fixture in `backend/tests/conftest.py`). + +**Spec:** [`.agents/specs/2026-09-28-agent-message-routing-design.md`](2026-09-28-agent-message-routing-design.md) · Branch: `feature/agent-message-routing` · PR: #107 + +**Commands run from `backend/`.** Test style mirrors `tests/agent/test_registry.py`: real SQLite via `session_factory`, explicit commits, scripted fakes (no mock libraries). + +--- + +## File Structure + +| File | Responsibility | +|---|---| +| Create `backend/src/octave/db/event_store.py` | `EventStore` — transcript append (gap-free seq + savepoint retry) and read. The only event write path. | +| Create `backend/src/octave/agent/decider.py` | `Candidate`, `DecisionState`, `Decision`, `TurnDecider` protocol, `LlmTurnDecider`. | +| Create `backend/src/octave/agent/router.py` | `TurnRunner` protocol, `StopReason`, `TurnRecord`, `RouteOutcome`, `MessageRouter`, transcript→Message mapping. | +| Modify `backend/src/octave/agent/errors.py` | `RoutingError`, `NotAMemberError`, `DeciderChoiceError`. | +| Modify `backend/src/octave/db/__init__.py` | Export `EventStore`. | +| Modify `backend/src/octave/agent/__init__.py` | Export the router/decider public surface. | +| Create `backend/tests/db/test_event_store.py` | EventStore tests. | +| Create `backend/tests/agent/test_decider.py` | LlmTurnDecider tests. | +| Create `backend/tests/agent/test_router.py` | MessageRouter driver tests. | +| Modify `backend/tests/agent/test_package.py` | Assert the new public names. | +| Modify `docs/TODO.md` | Mark Agent Manager #3 done. | + +--- + +### Task 0: Commit the design doc + +**Files:** +- Commit: `.agents/specs/2026-09-28-agent-message-routing-design.md` (already written) + +- [ ] **Step 1: Commit** + +```bash +git add .agents/specs/2026-09-28-agent-message-routing-design.md +git commit -m "docs(specs): add agent message routing design (#27)" +``` + +--- + +### Task 1: `EventStore` — transcript write path + +**Files:** +- Create: `backend/src/octave/db/event_store.py` +- Create: `backend/tests/db/test_event_store.py` +- Modify: `backend/src/octave/db/__init__.py` + +- [ ] **Step 1: Write the failing tests** + +Create `backend/tests/db/test_event_store.py`: + +```python +"""EventStore: gap-free transcript appends (design spec 2026-09-28). + +Real SQLite via session_factory; the uq_events_session_seq constraint is +the backstop for the seq race, exercised via a scripted stale-read subclass. +""" + +import pytest +from pydantic import ValidationError +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +import octave.db as db_pkg +from octave.db.event_store import EventStore +from octave.db.models import Participant, Session, User +from octave.db.types import EventKind + + +async def _seed(session_factory: async_sessionmaker[AsyncSession]) -> str: + """User u_1, sessions s_1/s_2, and a user participant (author FK).""" + async with session_factory() as session: + session.add(User(id="u_1", display_name="Alice")) + session.add(Session(id="s_1", created_by_user_id="u_1", status="active")) + session.add(Session(id="s_2", created_by_user_id="u_1", status="active")) + participant = Participant(id="p_u1", user_id="u_1", label="Alice") + session.add(participant) + await session.commit() + return "p_u1" + + +async def _append( + session_factory: async_sessionmaker[AsyncSession], + session_id: str, + author_id: str, + content: str, +) -> int: + async with session_factory() as session: + event = await EventStore(session).append( + session_id, + EventKind.USER_MESSAGE, + author_participant_id=author_id, + payload={"content": content}, + ) + await session.commit() + return event.seq + + +async def test_seq_starts_at_one_and_increments( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + assert await _append(session_factory, "s_1", author, "one") == 1 + assert await _append(session_factory, "s_1", author, "two") == 2 + + +async def test_seq_scoped_per_session( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + assert await _append(session_factory, "s_1", author, "a") == 1 + assert await _append(session_factory, "s_2", author, "b") == 1 + assert await _append(session_factory, "s_1", author, "c") == 2 + + +async def test_read_orders_and_windows( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + for text in ("one", "two", "three"): + await _append(session_factory, "s_1", author, text) + async with session_factory() as session: + store = EventStore(session) + everything = await store.read("s_1") + tail = await store.read("s_1", after_seq=2) + assert [e.seq for e in everything] == [1, 2, 3] + assert [e.payload["content"] for e in everything] == ["one", "two", "three"] + assert [e.seq for e in tail] == [3] + + +async def test_validates_user_message_payload( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + async with session_factory() as session: + with pytest.raises(ValidationError): + await EventStore(session).append( + "s_1", EventKind.USER_MESSAGE, author_participant_id=author, payload={} + ) + + +async def test_assistant_payload_defaults_model_name_null( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + author = await _seed(session_factory) + async with session_factory() as session: + event = await EventStore(session).append( + "s_1", + EventKind.ASSISTANT_MESSAGE, + author_participant_id=author, + payload={"content": "hi"}, + ) + assert event.payload == {"content": "hi", "model_name": None} + + +async def test_tool_and_system_kinds_pass_through( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + async with session_factory() as session: + await _seed(session_factory) + event = await EventStore(session).append( + "s_1", EventKind.TOOL_CALL, payload={"server_id": "s", "tool": "t"} + ) + assert event.payload == {"server_id": "s", "tool": "t"} + + +async def test_never_commits( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + async with session_factory() as session: + await EventStore(session).append( + "s_1", EventKind.SYSTEM, payload={"content": "uncommitted"} + ) + async with session_factory() as other: + assert await EventStore(other).read("s_1") == [] # not committed yet + async with session_factory() as session: + await EventStore(session).append( + "s_1", EventKind.SYSTEM, payload={"content": "committed"} + ) + await session.commit() + async with session_factory() as other: + rows = await EventStore(other).read("s_1") + assert [e.payload["content"] for e in rows] == ["committed"] + + +class _StaleReadStore(EventStore): + """Scripts stale seq reads to exercise the IntegrityError retry path: + real cross-connection races are invisible to a stale SQLite snapshot, + so the retry loop is tested at its seam.""" + + def __init__(self, session: AsyncSession, stale: list[int]) -> None: + super().__init__(session) + self._stale = stale + + async def _next_seq(self, session_id: str) -> int: + if self._stale: + return self._stale.pop(0) + return await super()._next_seq(session_id) + + +async def test_retry_recovers_from_seq_conflict( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + async with session_factory() as winner: + await EventStore(winner).append("s_1", EventKind.SYSTEM, payload={"c": 1}) + await winner.commit() # committed seq=1 + async with session_factory() as loser: + store = _StaleReadStore(loser, stale=[1]) # stale read -> collides + event = await store.append("s_1", EventKind.SYSTEM, payload={"c": 2}) + assert event.seq == 2 # retry recomputed from fresh max + await loser.commit() + + +async def test_exported_from_package() -> None: + assert hasattr(db_pkg, "EventStore") + assert "EventStore" in db_pkg.__all__ +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `cd backend && uv run pytest tests/db/test_event_store.py -v` +Expected: collection error — `ModuleNotFoundError: No module named 'octave.db.event_store'`. + +- [ ] **Step 3: Write minimal implementation** + +Create `backend/src/octave/db/event_store.py`: + +```python +"""Transcript write path: append/read events with gap-free seq (issue #27). + +``events.seq`` is documented as per-session monotonic (gap-free); clock +resolution cannot guarantee this, so seq assignment has exactly one owner. +Follows the VaultStore convention — constructed with the caller's +AsyncSession, never commits; callers own transaction boundaries +(``octave.db.deps``). The uq_events_session_seq constraint is the backstop +for the seq race: the insert runs in a SAVEPOINT so a collision is +recoverable without poisoning the caller's transaction. + +Payload validation: ``user_message``/``assistant_message`` validate against +their ``octave.db.types`` models; ``tool_call``/``tool_result``/``system`` +pass through as dicts (validation assigned to consumers, per the types +module docstring). +""" + +import uuid +from typing import Any + +from pydantic import BaseModel +from sqlalchemy import func, select +from sqlalchemy.exc import IntegrityError +from sqlalchemy.ext.asyncio import AsyncSession + +from octave.db.models import Event +from octave.db.types import ( + AssistantMessagePayload, + EventKind, + UserMessagePayload, +) + +__all__ = ["EventStore"] + +_PAYLOAD_MODELS: dict[EventKind, type[BaseModel]] = { + EventKind.USER_MESSAGE: UserMessagePayload, + EventKind.ASSISTANT_MESSAGE: AssistantMessagePayload, +} + +_RETRY_LIMIT = 2 + + +class EventStore: + """The canonical event append/read surface. Never commits.""" + + def __init__(self, session: AsyncSession) -> None: + self._session = session + + async def append( + self, + session_id: str, + kind: EventKind, + *, + author_participant_id: str | None = None, + target_participant_id: str | None = None, + payload: dict[str, Any], + ) -> Event: + """Append one event, assigning the next gap-free ``seq``. + + Malformed payloads for validated kinds raise ``ValidationError`` + (a write-path bug — fail loud). On a ``seq`` collision the insert + retries with a freshly computed seq up to ``_RETRY_LIMIT`` times, + then re-raises. Cross-connection effectiveness is engine-dependent + (SQLite snapshots stale reads within a transaction); the constraint + guarantees no duplicate ever lands either way. + """ + validator = _PAYLOAD_MODELS.get(kind) + if validator is not None: + payload = validator.model_validate(payload).model_dump() + last_error: IntegrityError | None = None + for _ in range(_RETRY_LIMIT + 1): + seq = await self._next_seq(session_id) + event = Event( + id=uuid.uuid4().hex, + session_id=session_id, + seq=seq, + kind=str(kind), + author_participant_id=author_participant_id, + target_participant_id=target_participant_id, + payload=payload, + ) + try: + async with self._session.begin_nested(): + self._session.add(event) + await self._session.flush() + except IntegrityError as exc: + last_error = exc + continue + return event + assert last_error is not None + raise last_error + + async def read(self, session_id: str, *, after_seq: int = 0) -> list[Event]: + """Transcript window ``seq > after_seq``, ordered.""" + rows = await self._session.execute( + select(Event) + .where(Event.session_id == session_id, Event.seq > after_seq) + .order_by(Event.seq) + ) + return list(rows.scalars().all()) + + async def _next_seq(self, session_id: str) -> int: + """Next seq for this session (1-based, gap-free). Autoflush makes + same-session prior appends visible; overridable seam for retry + testing.""" + current = await self._session.execute( + select(func.coalesce(func.max(Event.seq), 0)).where( + Event.session_id == session_id + ) + ) + return int(current.scalar_one()) + 1 +``` + +- [ ] **Step 4: Export from `octave.db`** + +Modify `backend/src/octave/db/__init__.py` — add the import after the `DbError` import line and the export in `__all__` (alphabetical placement): + +```python +from octave.db.errors import DbError +from octave.db.event_store import EventStore +``` + +```python +__all__ = [ + "Base", + "DbAdapter", + "DbAdapterRegistry", + "DbConfig", + "DbError", + "DatabaseSettings", + "EventKind", + "EventStore", + ... +] +``` + +(Keep the rest of `__all__` exactly as it is; insert `"EventStore"` after `"EventKind"`.) + +- [ ] **Step 5: Run tests to verify they pass** + +Run: `cd backend && uv run pytest tests/db/test_event_store.py -v` +Expected: all PASS (9 tests). + +- [ ] **Step 6: Lint + type check, then commit** + +Run: `cd backend && uv run ruff check src tests && uv run mypy src` +Expected: no errors. + +```bash +git add backend/src/octave/db/event_store.py backend/src/octave/db/__init__.py backend/tests/db/test_event_store.py +git commit -m "feat(db): add EventStore transcript write path (#27)" +``` + +--- + +### Task 2: Routing errors + +**Files:** +- Modify: `backend/src/octave/agent/errors.py` +- Modify: `backend/src/octave/agent/__init__.py` +- Modify: `backend/tests/agent/test_package.py` + +- [ ] **Step 1: Write the failing test** + +Modify `backend/tests/agent/test_package.py` — extend the `test_public_names_are_exported` tuple with the new names (Task 5 adds the router/decider names; this task adds only the errors): + +```python +def test_public_names_are_exported() -> None: + for name in ( + "AgentError", + "AgentInstanceManager", + "AgentNotFoundError", + "AgentPausedError", + "AgentRegistry", + "DeciderChoiceError", + "InstanceExistsError", + "InstanceNotFoundError", + "InvalidTransitionError", + "McpToolExecutor", + "ModelBindingError", + "NotAMemberError", + "ResolvedModel", + "RoutingError", + "RunningAgent", + "SessionNotFoundError", + "TerminalSessionError", + "ToolError", + "ToolExecutor", + "ToolLoop", + "ToolLoopLimitError", + "ToolOutcome", + "ToolTurn", + "TurnInProgressError", + "resolve_model", + ): + assert hasattr(agent_pkg, name), name +``` + +- [ ] **Step 2: Run test to verify it fails** + +Run: `cd backend && uv run pytest tests/agent/test_package.py -v` +Expected: FAIL — `AssertionError: DeciderChoiceError` (first new name). + +- [ ] **Step 3: Implement** + +Modify `backend/src/octave/agent/errors.py` — append after `ModelBindingError` and before `ToolLoopLimitError`: + +```python +class RoutingError(AgentError): + """Base for router failures (design spec 2026-09-28).""" + + +class NotAMemberError(RoutingError): + """Author is not a current member of the session (left_at set or + never invited).""" + + +class DeciderChoiceError(RoutingError): + """Decider output malformed or outside the roster. Raised by deciders; + the driver owns retry + await-user fallback.""" +``` + +Add the three names to `errors.py` `__all__` (alphabetical, matching existing style): + +```python +__all__ = [ + "AgentError", + "AgentNotFoundError", + "AgentPausedError", + "DeciderChoiceError", + "InstanceExistsError", + "InstanceNotFoundError", + "InvalidTransitionError", + "ModelBindingError", + "NotAMemberError", + "RoutingError", + "SessionNotFoundError", + "TerminalSessionError", + "ToolLoopLimitError", + "TurnInProgressError", +] +``` + +Modify `backend/src/octave/agent/__init__.py` — extend the errors import and `__all__`: + +```python +from octave.agent.errors import ( + AgentError, + AgentNotFoundError, + AgentPausedError, + DeciderChoiceError, + InstanceExistsError, + InstanceNotFoundError, + InvalidTransitionError, + ModelBindingError, + NotAMemberError, + RoutingError, + SessionNotFoundError, + TerminalSessionError, + ToolLoopLimitError, + TurnInProgressError, +) +``` + +In `__all__`, insert `"DeciderChoiceError"` after `"AgentRegistry"`, and `"NotAMemberError"` + `"RoutingError"` after `"ModelBindingError"` (keep alphabetical). + +- [ ] **Step 4: Run test to verify it passes** + +Run: `cd backend && uv run pytest tests/agent/test_package.py -v` +Expected: PASS. + +- [ ] **Step 5: Commit** + +```bash +git add backend/src/octave/agent/errors.py backend/src/octave/agent/__init__.py backend/tests/agent/test_package.py +git commit -m "feat(agent): add routing error hierarchy (#27)" +``` + +--- + +### Task 3: `TurnDecider` + `LlmTurnDecider` + +**Files:** +- Create: `backend/src/octave/agent/decider.py` +- Create: `backend/tests/agent/test_decider.py` + +- [ ] **Step 1: Write the failing tests** + +Create `backend/tests/agent/test_decider.py`: + +```python +"""LlmTurnDecider: 1:1 fast path + roster-validated JSON choice (spec #27).""" + +import pytest +from octave.agent.decider import Candidate, Decision, DecisionState, LlmTurnDecider +from octave.agent.errors import DeciderChoiceError +from octave.inference.errors import AdapterConnectionError +from octave.inference.types import Message +from tests.agent.fakes import ScriptedAdapter, final_result + + +def _candidate(iid: str, pid: str) -> Candidate: + return Candidate(instance_id=iid, participant_id=pid, label=f"Agent {iid}") + + +def _state(*messages: Message, roster: list[Candidate] | None = None) -> DecisionState: + return DecisionState( + roster=list(roster) if roster else [_candidate("i_1", "p_ai1")], + messages=list(messages), + ) + + +async def test_fast_path_picks_sole_candidate_when_user_spoke_last() -> None: + adapter = ScriptedAdapter([]) # any call fails the test via queue assert + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello")) + assert await decider.decide(state) == "p_ai1" + assert adapter.complete_calls == [] + + +async def test_fast_path_awaits_user_when_agent_spoke_last() -> None: + adapter = ScriptedAdapter([]) + decider = LlmTurnDecider(adapter=adapter) + state = _state( + Message(role="user", content="Hello"), + Message(role="assistant", content="Hi"), + ) + assert await decider.decide(state) == Decision.AWAIT_USER + assert adapter.complete_calls == [] + + +async def test_fast_path_empty_transcript_picks_candidate() -> None: + adapter = ScriptedAdapter([]) + decider = LlmTurnDecider(adapter=adapter) + assert await decider.decide(_state()) == "p_ai1" + + +async def test_multi_agent_valid_choice() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result('{"next": "p_ai2"}')]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + assert await decider.decide(state) == "p_ai2" + assert len(adapter.complete_calls) == 1 + + +async def test_multi_agent_await_user_choice() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result('{"next": "await_user"}')]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + assert await decider.decide(state) == Decision.AWAIT_USER + + +async def test_multi_agent_malformed_output_raises() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result("AI1 should go!")]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + with pytest.raises(DeciderChoiceError): + await decider.decide(state) + + +async def test_multi_agent_out_of_roster_choice_raises() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([final_result('{"next": "p_ghost"}')]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + with pytest.raises(DeciderChoiceError): + await decider.decide(state) + + +async def test_multi_agent_adapter_error_propagates() -> None: + roster = [_candidate("i_1", "p_ai1"), _candidate("i_2", "p_ai2")] + adapter = ScriptedAdapter([AdapterConnectionError("engine down")]) + decider = LlmTurnDecider(adapter=adapter) + state = _state(Message(role="user", content="Hello"), roster=roster) + with pytest.raises(AdapterConnectionError): + await decider.decide(state) +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `cd backend && uv run pytest tests/agent/test_decider.py -v` +Expected: collection error — `ModuleNotFoundError: No module named 'octave.agent.decider'`. + +- [ ] **Step 3: Write minimal implementation** + +Create `backend/src/octave/agent/decider.py`: + +```python +"""Turn-taking decisions (design spec 2026-09-28). + +The referee, not a player: a ``TurnDecider`` reads the roster + transcript +tail and names the next speaker (or hands control to the human). The +driver owns retry and fallback policy; deciders only raise +``DeciderChoiceError`` on unusable backend output. + +``LlmTurnDecider`` is the default strategy. Decision models (Laya/Jev +style) arrive behind this same protocol in a follow-up issue — the driver +never learns which backend chose. +""" + +import json +from dataclasses import dataclass, field +from enum import StrEnum +from typing import Protocol + +from pydantic import BaseModel, ValidationError + +from octave.agent.errors import DeciderChoiceError +from octave.inference.adapter import InferenceAdapter +from octave.inference.types import CompletionRequest, Message + +__all__ = ["Candidate", "Decision", "DecisionState", "LlmTurnDecider", "TurnDecider"] + + +@dataclass(frozen=True) +class Candidate: + """One agent eligible to take the next turn.""" + + instance_id: str + participant_id: str + label: str + + +@dataclass(frozen=True) +class DecisionState: + """Decider input: the candidate roster + the transcript tail.""" + + roster: list[Candidate] = field(default_factory=list) + messages: list[Message] = field(default_factory=list) + + +class Decision(StrEnum): + """Sentinel vocabulary for decider output. A decision is a plain + ``str``: either ``Decision.AWAIT_USER`` or a candidate's + ``participant_id`` (StrEnum members are strs, so the union collapses + cleanly).""" + + AWAIT_USER = "await_user" + + +class TurnDecider(Protocol): + """Next-speaker resolution. Implementations raise + ``DeciderChoiceError`` on malformed/out-of-roster output.""" + + async def decide(self, state: DecisionState) -> str: + """Return ``Decision.AWAIT_USER`` or a participant_id from + ``state.roster``.""" + ... + + +class _Choice(BaseModel): + """Wire shape of the multi-agent decision: exactly one JSON object.""" + + next: str + + +class LlmTurnDecider: + """Default strategy over the chat ``InferenceAdapter`` (local-first). + + 1:1 fast path (one candidate) requires no LLM call: user spoke last → + the candidate speaks; the sole agent spoke last → await the human. + Multi-agent asks for ``{"next": "" | "await_user"}``, + roster-validated; malformed or phantom choices raise + ``DeciderChoiceError`` (the driver retries once, then awaits user). + """ + + def __init__( + self, *, adapter: InferenceAdapter, model: str | None = None + ) -> None: + self._adapter = adapter + self._model = model + + async def decide(self, state: DecisionState) -> str: + if len(state.roster) == 1: + if state.messages and state.messages[-1].role == "assistant": + return Decision.AWAIT_USER + return state.roster[0].participant_id + if not state.roster: + return Decision.AWAIT_USER + result = await self._adapter.complete( + CompletionRequest(model=self._model, messages=self._prompt(state)) + ) + try: + choice = _Choice.model_validate(json.loads(result.text)).next + except (json.JSONDecodeError, ValidationError) as exc: + raise DeciderChoiceError( + f"decider output not parseable: {result.text!r}" + ) from exc + if choice == Decision.AWAIT_USER: + return Decision.AWAIT_USER + if any(candidate.participant_id == choice for candidate in state.roster): + return choice + raise DeciderChoiceError(f"decider chose unknown candidate: {choice!r}") + + @staticmethod + def _prompt(state: DecisionState) -> list[Message]: + roster = "\n".join( + f"- {candidate.participant_id}: {candidate.label}" + for candidate in state.roster + ) + system = ( + "You are the turn referee for a multi-agent conversation. " + "Decide who speaks next: one candidate participant_id, or " + "await_user to hand control to the human. Base the choice on " + "the conversation below. Reply with exactly one JSON object: " + '{"next": "" | "await_user"}\n\n' + f"Candidates:\n{roster}" + ) + return [Message(role="system", content=system), *state.messages] +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `cd backend && uv run pytest tests/agent/test_decider.py -v` +Expected: all PASS (8 tests). + +- [ ] **Step 5: Lint + type check, then commit** + +Run: `cd backend && uv run ruff check src tests && uv run mypy src` +Expected: no errors. + +```bash +git add backend/src/octave/agent/decider.py backend/tests/agent/test_decider.py +git commit -m "feat(agent): add TurnDecider protocol and LlmTurnDecider (#27)" +``` + +--- + +### Task 4: `MessageRouter` driver + +**Files:** +- Create: `backend/src/octave/agent/router.py` +- Create: `backend/tests/agent/test_router.py` + +- [ ] **Step 1: Write the failing tests** + +Create `backend/tests/agent/test_router.py`: + +```python +"""MessageRouter driver (design spec 2026-09-28). + +Real SQLite via session_factory; spawns driven through AgentInstanceManager; +scripted decider/runner fakes (no LLM). Transcript assertions read events +back through EventStore on a fresh session. +""" + +import pytest +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession, async_sessionmaker + +from octave.agent import AgentInstanceManager, AgentRegistry +from octave.agent.decider import Decision, DecisionState +from octave.agent.errors import DeciderChoiceError, NotAMemberError +from octave.agent.router import MessageRouter, StopReason +from octave.db.event_store import EventStore +from octave.db.models import Agent, Participant, Session, SessionParticipant, User +from octave.db.types import InstanceStatus +from octave.inference.types import Message + +_TAG_BINDING = {"kind": "tag", "tag": "quick"} + + +async def _seed(session_factory: async_sessionmaker[AsyncSession]) -> None: + async with session_factory() as session: + session.add(User(id="u_1", display_name="Alice")) + session.add(Agent(id="a_1", name="Echo", model_binding=_TAG_BINDING)) + session.add(Agent(id="a_2", name="Second", model_binding=_TAG_BINDING)) + session.add(Session(id="s_1", created_by_user_id="u_1", status="active")) + await session.commit() + + +async def _spawn( + session_factory: async_sessionmaker[AsyncSession], agent_id: str +) -> str: + async with session_factory() as session: + instance = await AgentInstanceManager(session).spawn( + agent_id=agent_id, session_id="s_1" + ) + await session.commit() + return instance.id + + +async def _invite_user(session_factory: async_sessionmaker[AsyncSession]) -> str: + """User participant + live membership in s_1.""" + async with session_factory() as session: + participant = Participant(id="p_u1", user_id="u_1", label="Alice") + session.add(participant) + await session.flush() + session.add( + SessionParticipant(session_id="s_1", participant_id="p_u1") + ) + await session.commit() + return "p_u1" + + +async def _participant_of( + session_factory: async_sessionmaker[AsyncSession], agent_id: str +) -> str: + async with session_factory() as session: + row = await session.execute( + select(Participant.id).where(Participant.agent_id == agent_id) + ) + return row.scalar_one() + + +class ScriptedDecider: + """Pops queued choices (Exception items raise); records each state.""" + + def __init__(self, choices: list[str]) -> None: + self._choices = list(choices) + self.states: list[DecisionState] = [] + + async def decide(self, state: DecisionState) -> str: + self.states.append(state) + item = self._choices.pop(0) + if isinstance(item, Exception): + raise item + return item + + +class RecordingRunner: + """Returns queued replies (Exception items raise); records inputs.""" + + def __init__(self, replies: list[str]) -> None: + self._replies = list(replies) + self.calls: list[tuple[str, list[Message]]] = [] + + async def run_turn(self, *, instance, messages: list[Message]) -> str: # type: ignore[no-untyped-def] + self.calls.append((instance.instance_id, list(messages))) + item = self._replies.pop(0) + if isinstance(item, Exception): + raise item + return item + + +def _router( + session: AsyncSession, + *, + decider: ScriptedDecider, + runner: RecordingRunner, + max_agent_turns: int = 4, +) -> MessageRouter: + return MessageRouter( + session=session, + manager=AgentInstanceManager(session), + registry=AgentRegistry(session), + decider=decider, + turn_runner=runner, + max_agent_turns=max_agent_turns, + ) + + +async def _deliver( + session_factory: async_sessionmaker[AsyncSession], + *, + decider: ScriptedDecider, + runner: RecordingRunner, + author: str = "p_u1", + max_agent_turns: int = 4, +): + async with session_factory() as session: + router = _router( + session, decider=decider, runner=runner, max_agent_turns=max_agent_turns + ) + outcome = await router.deliver("s_1", author_participant_id=author, content="Hello") + await session.commit() + return outcome + + +async def _transcript( + session_factory: async_sessionmaker[AsyncSession], +) -> list[tuple[int, str, str | None, str]]: + async with session_factory() as session: + events = await EventStore(session).read("s_1") + return [ + (e.seq, e.kind, e.author_participant_id, str(e.payload.get("content", ""))) + for e in events + ] + + +async def test_one_to_one_end_to_end( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + p_ai1 = await _participant_of(session_factory, "a_1") + await _invite_user(session_factory) + decider = ScriptedDecider([p_ai1, Decision.AWAIT_USER]) + runner = RecordingRunner(["Hi there"]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.AWAIT_USER + assert len(outcome.turns) == 1 + assert outcome.turns[0].agent_id == "a_1" + assert outcome.turns[0].seq_start == outcome.turns[0].seq_end == 2 + assert await _transcript(session_factory) == [ + (1, "user_message", "p_u1", "Hello"), + (2, "assistant_message", p_ai1, "Hi there"), + ] + + +async def test_multi_agent_handoff( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _spawn(session_factory, "a_2") + p_ai1 = await _participant_of(session_factory, "a_1") + p_ai2 = await _participant_of(session_factory, "a_2") + await _invite_user(session_factory) + decider = ScriptedDecider([p_ai1, p_ai2, Decision.AWAIT_USER]) + runner = RecordingRunner(["from AI1", "from AI2"]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert [t.agent_id for t in outcome.turns] == ["a_1", "a_2"] + # AI2's runner input contains U1's message AND AI1's reply + assert runner.calls[1][1] == [ + Message(role="user", content="Hello"), + Message(role="assistant", content="from AI1"), + ] + kinds = [(seq, kind) for seq, kind, _, _ in await _transcript(session_factory)] + assert kinds == [ + (1, "user_message"), + (2, "assistant_message"), + (3, "assistant_message"), + ] + + +async def test_empty_roster_awaits_user_immediately( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _invite_user(session_factory) + decider = ScriptedDecider([]) # never consulted + runner = RecordingRunner([]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.AWAIT_USER + assert outcome.turns == [] + assert runner.calls == [] + assert await _transcript(session_factory) == [(1, "user_message", "p_u1", "Hello")] + + +async def test_runner_failure_releases_instance( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _invite_user(session_factory) + p_ai1 = await _participant_of(session_factory, "a_1") + decider = ScriptedDecider([p_ai1]) + runner = RecordingRunner([RuntimeError("boom")]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.ERROR + assert outcome.error == "boom" + assert outcome.turns == [] # failed turn contributes no TurnRecord + transcript = await _transcript(session_factory) + assert [(seq, kind) for seq, kind, _, _ in transcript] == [ + (1, "user_message"), + (2, "system"), + ] + assert "boom" in transcript[1][3] + async with session_factory() as session: + counts = await AgentRegistry(session).count_by_status() + assert counts[InstanceStatus.ACTIVE] == 0 # turn released + + +async def test_hop_limit_stops_agent_chatter( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _invite_user(session_factory) + p_ai1 = await _participant_of(session_factory, "a_1") + decider = ScriptedDecider([p_ai1]) # one decision; loop stops at limit + runner = RecordingRunner(["first"]) + outcome = await _deliver( + session_factory, decider=decider, runner=runner, max_agent_turns=1 + ) + assert outcome.stop_reason == StopReason.HOP_LIMIT + assert len(outcome.turns) == 1 + + +async def test_decider_invalid_falls_back_to_await_user( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _invite_user(session_factory) + decider = ScriptedDecider( + [DeciderChoiceError("bad1"), DeciderChoiceError("bad2")] + ) + runner = RecordingRunner([]) + outcome = await _deliver(session_factory, decider=decider, runner=runner) + assert outcome.stop_reason == StopReason.AWAIT_USER + assert outcome.turns == [] + assert runner.calls == [] + + +async def test_non_member_author_rejected( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + async with session_factory() as session: + router = _router( + session, decider=ScriptedDecider([]), runner=RecordingRunner([]) + ) + try: + with pytest.raises(NotAMemberError): + await router.deliver( + "s_1", author_participant_id="p_ghost", content="Hello" + ) + finally: + await session.rollback() + assert await _transcript(session_factory) == [] # nothing appended + + +async def test_registry_consistent_after_loop( + session_factory: async_sessionmaker[AsyncSession], +) -> None: + await _seed(session_factory) + await _spawn(session_factory, "a_1") + await _spawn(session_factory, "a_2") + p_ai1 = await _participant_of(session_factory, "a_1") + p_ai2 = await _participant_of(session_factory, "a_2") + await _invite_user(session_factory) + decider = ScriptedDecider([p_ai1, p_ai2, Decision.AWAIT_USER]) + runner = RecordingRunner(["one", "two"]) + await _deliver(session_factory, decider=decider, runner=runner) + async with session_factory() as session: + counts = await AgentRegistry(session).count_by_status() + assert counts[InstanceStatus.ACTIVE] == 0 + assert counts[InstanceStatus.IDLE] == 2 +``` + +- [ ] **Step 2: Run tests to verify they fail** + +Run: `cd backend && uv run pytest tests/agent/test_router.py -v` +Expected: collection error — `ModuleNotFoundError: No module named 'octave.agent.router'`. + +- [ ] **Step 3: Write minimal implementation** + +Create `backend/src/octave/agent/router.py`: + +```python +"""Session-level message routing (design spec 2026-09-28, issue #27). + +The driver: append the user message, then loop — decider picks an idle +instance, ``begin_turn`` claims the turn (the #25 mutex), the injected +``TurnRunner`` port produces the reply, the reply is recorded, the turn is +released — until await-user, the hop limit, or a runner failure. + +Delivery is via the shared transcript (spec decision 2): there are no +per-agent queues; an agent's "message plus everything queued for it" is +the transcript through the current seq, which the runner reads as input. +Routing itself is ephemeral — ``events.target_participant_id`` stays NULL. + +Follows the VaultStore convention: constructed with the caller's +AsyncSession, never commits. One ``deliver()`` is one caller-owned +transaction; the whole exchange commits atomically or not at all. +""" + +import logging +from dataclasses import dataclass, field +from enum import StrEnum +from typing import Protocol + +from sqlalchemy import select +from sqlalchemy.ext.asyncio import AsyncSession + +from octave.agent.decider import Candidate, Decision, DecisionState, TurnDecider +from octave.agent.errors import DeciderChoiceError, NotAMemberError +from octave.agent.instances import AgentInstanceManager +from octave.agent.registry import AgentRegistry, RunningAgent +from octave.db.event_store import EventStore +from octave.db.models import Event, Participant, SessionParticipant +from octave.db.types import AgentStatus, EventKind, InstanceStatus +from octave.inference.types import Message + +__all__ = [ + "MessageRouter", + "RouteOutcome", + "StopReason", + "TurnRecord", + "TurnRunner", +] + +logger = logging.getLogger(__name__) + + +class StopReason(StrEnum): + """Why ``deliver()`` returned control to the caller.""" + + AWAIT_USER = "await_user" + HOP_LIMIT = "hop_limit" + ERROR = "error" + + +@dataclass(frozen=True) +class TurnRecord: + """One completed agent turn. Events contributed by this turn + (``assistant_message`` now; tool events when Integration #1's real + runner appends them inside the bracket). CM #5's archival key per #25: + (session_id, agent_id, seq_range) — never instance_id.""" + + session_id: str + agent_id: str + instance_id: str + seq_start: int + seq_end: int + + +@dataclass(frozen=True) +class RouteOutcome: + turns: list[TurnRecord] = field(default_factory=list) + stop_reason: StopReason = StopReason.AWAIT_USER + error: str | None = None + + +class TurnRunner(Protocol): + """What happens inside a claimed turn. The real runner (context + assembly + ToolLoop over the instance's model binding) is Integration + #1's composition — mirrors the ToolExecutor/McpToolExecutor split. + Returns the final reply text.""" + + async def run_turn( + self, *, instance: RunningAgent, messages: list[Message] + ) -> str: ... + + +class MessageRouter: + """Session-level turn-taking driver. Never commits; the caller owns + the transaction boundary.""" + + def __init__( + self, + *, + session: AsyncSession, + manager: AgentInstanceManager, + registry: AgentRegistry, + decider: TurnDecider, + turn_runner: TurnRunner, + max_agent_turns: int = 4, + decider_tail: int = 30, + ) -> None: + self._session = session + self._manager = manager + self._registry = registry + self._decider = decider + self._turn_runner = turn_runner + self._max_agent_turns = max_agent_turns + self._decider_tail = decider_tail + self._events = EventStore(session) + + async def deliver( + self, session_id: str, *, author_participant_id: str, content: str + ) -> RouteOutcome: + """Record the user message and drive agent turns until stop.""" + await self._require_membership(session_id, author_participant_id) + await self._events.append( + session_id, + EventKind.USER_MESSAGE, + author_participant_id=author_participant_id, + payload={"content": content}, + ) + turns: list[TurnRecord] = [] + agent_turns = 0 + while True: + roster = await self._roster(session_id) + if not roster: + return RouteOutcome(turns=turns, stop_reason=StopReason.AWAIT_USER) + if agent_turns >= self._max_agent_turns: + return RouteOutcome(turns=turns, stop_reason=StopReason.HOP_LIMIT) + events = await self._events.read(session_id) + messages = _transcript_messages(events) + state = DecisionState( + roster=roster, messages=messages[-self._decider_tail :] + ) + choice = await self._decide_with_fallback(state) + if choice == Decision.AWAIT_USER: + return RouteOutcome(turns=turns, stop_reason=StopReason.AWAIT_USER) + candidate = next(c for c in roster if c.participant_id == choice) + outcome = await self._run_agent_turn( + session_id, candidate, messages, turns + ) + if outcome is not None: + return outcome + agent_turns += 1 + + async def _run_agent_turn( + self, + session_id: str, + candidate: Candidate, + messages: list[Message], + turns: list[TurnRecord], + ) -> RouteOutcome | None: + """Claim, run, record, release. None = turn completed normally; + a RouteOutcome = the loop must stop (runner failure).""" + # TurnInProgressError propagates: concurrent drivers on one + # session is a caller bug; the mutex is fail-loud by design. + await self._manager.begin_turn(candidate.instance_id) + running = await self._registry.get_instance(candidate.instance_id) + try: + reply_text = await self._turn_runner.run_turn( + instance=running, messages=messages + ) + except Exception as exc: # noqa: BLE001 — turn failures land idle + logger.warning( + "agent turn failed | instance=%s error=%s", candidate.instance_id, exc + ) + await self._manager.end_turn(candidate.instance_id) + await self._events.append( + session_id, + EventKind.SYSTEM, + payload={"content": f"Agent turn failed: {exc}"}, + ) + return RouteOutcome( + turns=turns, stop_reason=StopReason.ERROR, error=str(exc) + ) + reply = await self._events.append( + session_id, + EventKind.ASSISTANT_MESSAGE, + author_participant_id=candidate.participant_id, + payload={"content": reply_text}, + ) + await self._manager.end_turn(candidate.instance_id) + turns.append( + TurnRecord( + session_id=session_id, + agent_id=running.agent_id, + instance_id=candidate.instance_id, + seq_start=reply.seq, + seq_end=reply.seq, + ) + ) + return None + + async def _decide_with_fallback(self, state: DecisionState) -> str: + """One retry on DeciderChoiceError, then AWAIT_USER: a confused + referee hands control to the human — never crashes the session.""" + for attempt in (0, 1): + try: + return await self._decider.decide(state) + except DeciderChoiceError as exc: + logger.warning( + "decider choice invalid | attempt=%s error=%s", attempt, exc + ) + return Decision.AWAIT_USER + + async def _roster(self, session_id: str) -> list[Candidate]: + """Idle instances of active definitions, joined to their + participant identity (guaranteed by spawn; a missing participant + row is write-path corruption — surface it loud).""" + running = [ + record + for record in await self._registry.list_instances( + session_id=session_id, status=InstanceStatus.IDLE + ) + if record.definition_status == AgentStatus.ACTIVE + ] + if not running: + return [] + rows = await self._session.execute( + select(Participant.agent_id, Participant.id).where( + Participant.agent_id.in_([record.agent_id for record in running]) + ) + ) + participant_by_agent = dict(rows.all()) + return [ + Candidate( + instance_id=record.instance_id, + participant_id=participant_by_agent[record.agent_id], + label=record.agent_name, + ) + for record in running + ] + + async def _require_membership( + self, session_id: str, participant_id: str + ) -> None: + """App-level author gate: current membership (the events FK covers + identity, not membership).""" + row = await self._session.execute( + select(SessionParticipant).where( + SessionParticipant.session_id == session_id, + SessionParticipant.participant_id == participant_id, + SessionParticipant.left_at.is_(None), + ) + ) + if row.scalar_one_or_none() is None: + raise NotAMemberError( + f"participant {participant_id} is not a current member of {session_id}" + ) + + +def _transcript_messages(events: list[Event]) -> list[Message]: + """Transcript → chat messages. Tool/system entries are skipped with a + debug log: they are context for humans and archival, not chat history.""" + messages: list[Message] = [] + for event in events: + if event.kind == EventKind.USER_MESSAGE: + messages.append(Message(role="user", content=event.payload["content"])) + elif event.kind == EventKind.ASSISTANT_MESSAGE: + messages.append( + Message(role="assistant", content=event.payload["content"]) + ) + else: + logger.debug( + "skipping non-message event | kind=%s seq=%s", event.kind, event.seq + ) + return messages +``` + +- [ ] **Step 4: Run tests to verify they pass** + +Run: `cd backend && uv run pytest tests/agent/test_router.py -v` +Expected: all PASS (8 tests). + +- [ ] **Step 5: Lint + type check, then commit** + +Run: `cd backend && uv run ruff check src tests && uv run mypy src` +Expected: no errors. + +```bash +git add backend/src/octave/agent/router.py backend/tests/agent/test_router.py +git commit -m "feat(agent): add session-level MessageRouter driver (#27)" +``` + +--- + +### Task 5: Public exports, TODO, full suite + +**Files:** +- Modify: `backend/src/octave/agent/__init__.py` +- Modify: `backend/tests/agent/test_package.py` +- Modify: `docs/TODO.md` + +- [ ] **Step 1: Write the failing test** + +Modify `backend/tests/agent/test_package.py` — extend the names tuple (insert alphabetically): + +```python + "Candidate", + "Decision", + "DecisionState", + "LlmTurnDecider", + "MessageRouter", + "RouteOutcome", + "StopReason", + "TurnDecider", + "TurnRecord", + "TurnRunner", +``` + +- [ ] **Step 2: Run test to verify it fails** + +Run: `cd backend && uv run pytest tests/agent/test_package.py -v` +Expected: FAIL — `AssertionError: Candidate`. + +- [ ] **Step 3: Implement exports** + +Modify `backend/src/octave/agent/__init__.py` — add imports after the errors import: + +```python +from octave.agent.decider import ( + Candidate, + Decision, + DecisionState, + LlmTurnDecider, + TurnDecider, +) +from octave.agent.router import ( + MessageRouter, + RouteOutcome, + StopReason, + TurnRecord, + TurnRunner, +) +``` + +Extend `__all__` (keep alphabetical): + +```python +__all__ = [ + "AgentError", + "AgentInstanceManager", + "AgentNotFoundError", + "AgentPausedError", + "AgentRegistry", + "Candidate", + "DeciderChoiceError", + "Decision", + "DecisionState", + "InstanceExistsError", + "InstanceNotFoundError", + "InvalidTransitionError", + "LlmTurnDecider", + "McpToolExecutor", + "MessageRouter", + "ModelBindingError", + "NotAMemberError", + "ResolvedModel", + "RouteOutcome", + "RoutingError", + "RunningAgent", + "SessionNotFoundError", + "StopReason", + "TerminalSessionError", + "ToolError", + "ToolExecutor", + "ToolLoop", + "ToolLoopLimitError", + "ToolOutcome", + "ToolTurn", + "TurnDecider", + "TurnInProgressError", + "TurnRecord", + "TurnRunner", + "resolve_model", +] +``` + +Also update the package docstring's second paragraph: replace `"Future Agent Manager components (router, result collector) land here as additional modules; they do not exist yet."` with: + +``` +Manager components ship here as additional modules: instance lifecycle +(#25), registry (#26), message routing (#27). Result collection (#4) and +scheduling (#6) remain. +``` + +- [ ] **Step 4: Mark the roadmap item done** + +Modify `docs/TODO.md` — replace line: + +```markdown +- [ ] 3. Implement agent message routing (deliver messages to correct agent, broadcast when needed) +``` + +with: + +```markdown +- [x] 3. Implement agent message routing (deliver messages to correct agent, broadcast when needed) — PR #107 (session-level `MessageRouter` driver, `TurnDecider`/`TurnRunner` ports, `EventStore` transcript write path; LLM decider default, decision-model decider deferred to follow-up) +``` + +- [ ] **Step 5: Run tests to verify they pass** + +Run: `cd backend && uv run pytest tests/agent/test_package.py -v` +Expected: PASS. + +- [ ] **Step 6: Full suite + lint + types** + +Run: `cd backend && uv run pytest -q` +Expected: all tests pass, no new failures against the pre-existing suite. + +Run: `cd backend && uv run ruff check src tests && uv run mypy src` +Expected: no errors. + +- [ ] **Step 7: Commit** + +```bash +git add backend/src/octave/agent/__init__.py backend/tests/agent/test_package.py docs/TODO.md +git commit -m "feat(agent): export message routing surface, tick Agent Manager #3 (#27)" +``` + +--- + +## Deviations From Spec + +- **Retry testing seam**: the spec's cross-session concurrent-append test is impossible on SQLite (a transaction's snapshot never sees other connections' commits, so the retry's re-read is stale by construction). The plan tests the retry loop deterministically via a `_StaleReadStore` subclass scripting one stale `_next_seq` value, and documents engine-dependence in `EventStore.append`'s docstring. The constraint guarantee (no duplicate seq ever lands) is unchanged. +- **`_next_seq` as an overridable method** (not an inline expression) — the seam the retry test needs; also makes the seq policy swappable if pgvector-era locking (e.g., `SELECT ... FOR UPDATE`) arrives. +- **`decider_tail` on the driver** (default 30) rather than only on the decider: the spec assigns the tail budget to `DecisionState` construction, which is the driver's job; `LlmTurnDecider` receives the already-tailed state. diff --git a/docs/ARCHITECTURE.md b/docs/ARCHITECTURE.md index 4eec204..0c06bde 100644 --- a/docs/ARCHITECTURE.md +++ b/docs/ARCHITECTURE.md @@ -167,6 +167,25 @@ write path). Malformed `assignments` JSON degrades to empty with a warning (repo surface); bad status enums fail loud. Library-only: no routes. Design: [`.agents/specs/2026-09-27-agent-registry-design.md`](../.agents/specs/2026-09-27-agent-registry-design.md). +**Implemented — agent message routing (issue #27):** `MessageRouter` +(`octave.agent.router`) is the session-level turn-taking driver: `deliver()` appends +the user message (author must be a current member — `NotAMemberError` otherwise), then +loops — the injected `TurnDecider` names the next speaker from the idle-instance +roster, `begin_turn` claims the #25 mutex, the injected `TurnRunner` port produces the +reply, the reply is appended as `assistant_message`, the turn is released — until +`AWAIT_USER`, the hop limit (`max_agent_turns`), or a runner failure (`ERROR`; failed +turn lands idle + `system` event). Delivery is via the shared transcript: no per-agent +queues; the runner reads the transcript as input. `LlmTurnDecider` +(`octave.agent.decider`) is the default strategy — a 1:1 fast path with zero LLM calls, +roster-validated JSON choice for multi-agent sessions; the driver retries a confused +referee once, then hands control to the human. `EventStore` (`octave.db.event_store`) +owns transcript appends: gap-free per-session `seq` with savepoint retry on collision +(uq constraint is the backstop), per-kind payload validation. `TurnRunner` stays a +port — the real runner (context assembly + `ToolLoop`) is Integration #1's composition. +Library-only: no routes; one `deliver()` is one caller-owned transaction (never +commits). Design: +[`.agents/specs/2026-09-28-agent-message-routing-design.md`](../.agents/specs/2026-09-28-agent-message-routing-design.md). + --- ## Data Layer @@ -177,6 +196,11 @@ the inference adapter pattern: - **`DbAdapter` ABC + registry** — engine selection by name or import string; `sqlite` (SQLite + vec0) is the only adapter registered today, `pgvector` will self-register via the plugin path when it lands +- **Explicit transaction control** (`octave.db._bootstrap.attach_transaction_control`) + — SQLite connections disable pysqlite's implicit-BEGIN mode and emit explicit + `BEGIN`; without it `RELEASE SAVEPOINT` acts as a commit, silently breaking the + "stores never commit; callers own the transaction boundary" convention that + `EventStore`'s seq-retry (and every store) depends on (issue #27) - **Alembic** for schema migrations, run programmatically via `octave.db.migrations.upgrade()` — auto-applied on app startup via `octave.db.lifespan.db_lifespan` unless `OCTAVE_DB_AUTO_MIGRATE=false`