From 8e5fd01ac9fcd123d8955404edde4616bcfd5227 Mon Sep 17 00:00:00 2001 From: Dai Ha Date: Thu, 16 Jul 2026 16:45:12 +0200 Subject: [PATCH] =?UTF-8?q?wiki:=20add=20page=209=20=E2=80=94=20as-built?= =?UTF-8?q?=20implementation=20architecture?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Fan-out code audit (one worker per layer, off-subscription) synthesized into a verified as-built map: component layers, per-package class reference, REST route table, bootstrap wiring, forward + reverse rendezvous sequence flows, and two state machines (AgentStatus; Injector per-target turn lifecycle with the real grace constants 8/120/240). Enums, routes, and the guard allowlist verified against source at aa0cf81. All 6 mermaid diagrams mmdc-validated and theme-safe. --- 9-Implementation.md | 357 ++++++++++++++++++++++++++++++++++++++++++++ _Sidebar.md | 1 + 2 files changed, 358 insertions(+) create mode 100644 9-Implementation.md diff --git a/9-Implementation.md b/9-Implementation.md new file mode 100644 index 0000000..aaea5d6 --- /dev/null +++ b/9-Implementation.md @@ -0,0 +1,357 @@ +# 9. Implementation Architecture (as-built) + +> **Scope.** This is the *as-built* code map of the `bridged` module — the actual packages, +> classes, flows, and state machines in the source tree, as a companion to the design-level +> [1. Architecture](1-Architecture) and [2. Message Server](2-Message-Server). Every enum, +> constant, and route below was verified against source at commit `aa0cf81`. + +`bridged` is a single-host Java 25 / Maven daemon: the sole gateway between an on-subscription +**primary** (Opus) and off-subscription **workers**, speaking to the **herdr** PTY manager +(protocol 14, herdr 0.7.0) over a Unix-domain socket. It presents two equivalent faces — a REST +server (the testability seam) and an MCP server — over one shared service core. + +## Component map + +The daemon is layered. The two north faces (REST + MCP) are thin adapters over one service core +(`msg`); the core drives delivery through `inject`, which speaks to the outside world only through +`herdr`. `guard` sits on the spawn path; `config` and the `Bridged` entry point wire it all. + +```mermaid +flowchart TB + primary["Primary (Opus)
on-subscription"] + worker["Worker (Claude Code)
off-subscription"] + + subgraph bridged["bridged daemon"] + direction TB + subgraph faces["North faces (thin adapters)"] + rest["rest.BridgedApp
REST · Javalin"] + mcp["mcp.BridgeMcp
MCP · /mcp servlet"] + end + msg["msg.MessageService + Rendezvous
service core · per-session rendezvous"] + inject["inject.Injector + StatusPoller
single-writer, status-gated delivery"] + guard["guard.SubscriptionGuard
boundary enforcement"] + worksvc["worker.WorkerService
spawn · reap · teardown"] + herdr["herdr.AgentControl / WorkspaceControl
JSON-RPC over UNIX socket"] + end + + herdrd["herdr daemon
(PTY manager)"] + + primary -->|"bridge_send / answer"| rest + primary -->|"MCP tool calls"| mcp + worker -->|"bridge_reply / bridge_ask"| mcp + rest --> msg + mcp --> msg + mcp --> worksvc + worksvc --> guard + worksvc --> herdr + msg --> inject + inject --> herdr + herdr --> herdrd + herdrd -.->|"drives PTYs"| worker + + classDef core fill:#2b6cb0,stroke:#1a365d,color:#ffffff; + classDef gate fill:#b7791f,stroke:#7b341e,color:#ffffff; + class msg,inject core + class guard gate +``` + +*Figure 1 — component layers. The service core (`msg`) is reached identically from either face; +delivery reaches workers only via `inject → herdr`. The subscription guard (amber) gates spawn.* + +## Package & class reference + +Nine packages under `dev.ltms.bridged`. Below, per layer: the classes, their kind, and their +role. Method signatures are abbreviated; see source for full contracts. + +### `herdr` — the wire layer (11 classes) + +The only speaker of the herdr protocol. Encodes newline-delimited JSON-RPC (string request ids, +**one request per connection**) and projects herdr's workspace/tab/pane/agent nodes into typed +Java records. + +| Class | Kind | Role | +|---|---|---| +| `HerdrClient` | interface | The client face — the only herdr speaker in `bridged` (`call(method, params)`). | +| `UnixSocketHerdrClient` | class | JDK Unix-socket transport; fresh connection per `call`, one frame, one line, close. | +| `HerdrCodec` | class | Encodes/decodes JSON-RPC frames (`encode`, `decodeResult`). | +| `HerdrException` | class | Runtime failure carrying an optional herdr protocol `code()`. | +| `AgentControl` | class | Domain wrapper over `agent.*`: `start`, `send`, `read`, `status`, `close`. | +| `Agent` | record | Projects an `agent` node into a worker handle (status/session metadata). | +| `AgentStatus` | enum | Worker lifecycle state (see [state machine](#state-machine-agentstatus)). | +| `WorkspaceControl` | class | Placement over `workspace.*`/`tab.*`: `ensureWorkspace` (synchronized), `createTab`, `renameTab`, `closeTab`. | +| `Workspace` | record | Projects a herdr workspace node. | +| `Tab` | record | Projects a tab plus the placeholder pane seeded at creation. | +| `PaneLocator` | class | Resolves a PID to a herdr `terminal_id` (`terminalForPid`). | + +**Invariant:** a prompt is submitted only by a standalone `"\r"` `agent.send` after the message +text (`AgentControl.SUBMIT_KEY`) — never a shell prefix. `closeTab`/`PaneLocator` swallow only +"already gone" errors; other failures propagate. + +### `inject` — status-gated delivery + turn lifecycle (6 classes) + +The single-writer delivery layer. Polls worker status, queues per target, injects only when safe, +and **synthesizes turn boundaries** so a blocked `bridge_send` resolves even when a worker never +calls `bridge_reply`. This is where the turn state machine lives. + +| Class | Kind | Role | +|---|---|---| +| `StatusPoller` | class | Virtual-thread loop sampling active workers; feeds refined status to `Injector` (`start`/`stop`, idempotent). | +| `StatusRefiner` | class | Reclassifies `UNKNOWN` by reading pane content — pane text is authoritative (`refine`, `classify`). | +| `Injector` | class | Queues + delivers at most one message per turn, serialized per target (`enqueue`, `onStatus`, `drop`, `activeTargets`). | +| `CompletionResolver` | class | `TurnListener` impl; scrapes the pane to resolve `Rendezvous` on completion/failure. | +| `TurnListener` | interface | Turn-lifecycle callbacks: `onDelivered`, `onTurnComplete`, `onTurnFailed`. | +| `WorkerPresence` | class | Tracks which workers have connected the bridge MCP (`markPresent`, `isPresent`, `forget`). | + +### `msg` — the service core (2 classes) + +Owns the forward rendezvous (`bridge_send` → `bridge_reply`) and the reverse rendezvous +(`bridge_ask` → answer), plus async fire-and-poll. + +| Class | Kind | Role | +|---|---|---| +| `MessageService` | class | Orchestrates send/reply/ask/answer + async dispatch (`send`, `answer`, `ask`, `sendAsync`, `poll`). | +| `Rendezvous` | class | Low-level registry of forward waiters + reverse-ask futures (`open`, `resolve`, `resolveQuestion`, `openAsk`, `answerAsk`, `closeAsk`, `resolveCompletion`, `resolveFailure`). | + +**Enums (verbatim).** `Rendezvous.Kind`: `REPLY`, `COMPLETION`, `FAILED`, `QUESTION`. +`MessageService.Outcome`: `REPLIED`, `COMPLETED_UNREPLIED`, `WORKER_FAILED`, `QUESTION`, +`TIMED_OUT_WORKING`, `TIMED_OUT_QUEUED`, `BUSY`, `STALE_TURN`. `AskOutcome`: `ANSWERED`, +`NO_WAITER`, `TIMED_OUT`. `Phase`: `PENDING`, `DONE`, `FAILED`. + +### `mcp` — the MCP north face (6 classes) + +Exposes `bridged` as a Streamable-HTTP MCP endpoint and resolves caller identity from the +**connection**, not from tool arguments (unspoofable). + +| Class | Kind | Role | +|---|---|---| +| `BridgeMcp` | class | Builds the MCP server, registers `bridge_*` tools, holds thin tool adapters (`servlet`, `send`, `answer`, `ask`, `reply`, `spawn`). | +| `ConnectionIdentity` | class | Resolves *who is calling*: peer PID → herdr pane → `terminal_id` (`resolve`, `callerTerminal`, `cwdForPid`). | +| `PeerPidLookup` | interface | Abstracts OS peer-PID lookup for a loopback source port. | +| `LsofPeerPidLookup` | class | `lsof` impl; excludes bridged's own PID. | +| `ProcessCwdLookup` | interface | Abstracts PID → cwd lookup. | +| `LsofProcessCwdLookup` | class | `lsof`-based cwd lookup for spawn inheritance. | + +**Tool surface:** `bridge_send` (delegate + block; with `turnId`, answers an ask) · `bridge_reply` +(worker → structured answer) · `bridge_ask` (worker pauses to ask primary) · `bridge_status` · +`bridge_poll` (async ticket) · `bridge_spawn` · `bridge_list` · `bridge_stop` · `bridge_profiles`. +**Invariant:** `bridge_reply`/`bridge_ask` accept *no* identity argument; it comes only from the +transport context. + +### `rest` · `worker` · `guard` · `config` — the edge (6 classes) + +| Class | Kind | Role | +|---|---|---| +| `rest.BridgedApp` | class | Javalin routes; validates bodies, maps `Outcome` → HTTP status. | +| `worker.WorkerService` | class | Guard-checked spawn, orphan-pane reaping at boot, teardown (`spawn`, `reapOrphanWorkers`, `stop`, `list`, `profiles`). | +| `guard.SubscriptionGuard` | class | Host-allowlist + primary-cleanliness enforcement (`assertWorker`, `assertPrimaryClean`). | +| `guard.GuardException` | class | Thrown on any subscription-boundary violation. | +| `config.BridgedConfig` | record | YAML config with defaults; single legacy worker or named `workers` map (`load`, `workerProfiles`, `defaultProfile`). | +| `Bridged` | class | Static `main` — wires real collaborators and starts the server. | + +**REST routes** (the acceptance surface — every capability is reachable here without MCP): + +| Method + path | Purpose | +|---|---| +| `GET /healthz` | Liveness + herdr reachability | +| `GET /sessions` | herdr workspace list | +| `GET /agents` | herdr-tracked agents | +| `GET /profiles` | Configured worker profiles + default | +| `POST /workers` | Spawn a worker (`?profile=`, `?cwd=`, or JSON body) | +| `DELETE /workers/{paneId}` | Stop a worker | +| `POST /sessions/{id}/message` | `bridge_send` — blocking, `wait:false`, or answer via `turnId` | +| `POST /sessions/{id}/reply` | `bridge_reply` (worker) | +| `POST /sessions/{id}/ask` | `bridge_ask` (worker → primary) | +| `GET /sessions/{id}/status` | Worker lifecycle + MCP readiness | +| `GET /tasks/{ticket}` | Poll an async (`wait:false`) send | + +## Bootstrap wiring + +`Bridged.main` assembles the object graph in dependency order, then starts both faces. + +```mermaid +flowchart LR + cfg["load config
+ assertPrimaryClean"] --> guard["SubscriptionGuard"] + guard --> herdr["connect UnixSocketHerdrClient
→ AgentControl / WorkspaceControl"] + herdr --> wsvc["WorkerService
→ reapOrphanWorkers()"] + wsvc --> rv["Rendezvous → CompletionResolver
→ Injector + StatusPoller"] + rv --> ms["MessageService"] + ms --> mcp["BridgeMcp
(connection identity)"] + ms --> app["BridgedApp
(start Javalin)"] + mcp --> app +``` + +*Figure 2 — startup wiring in `Bridged.main`. The guard asserts the primary env is clean before +anything else; `reapOrphanWorkers()` clears stale panes from a prior daemon restart.* + +## Flow: forward rendezvous (`bridge_send` → `bridge_reply`) + +The main path. A primary's send blocks on a per-session waiter that resolves on the worker's +explicit reply — **or**, as a fallback, on a confirmed `working → idle` turn boundary the injector +observes (so a worker that finishes without calling `bridge_reply` still unblocks the caller). + +```mermaid +sequenceDiagram + autonumber + participant P as Primary + participant MS as MessageService + participant RV as Rendezvous + participant IJ as Injector + participant W as Worker + + P->>MS: send(target, content) + Note over MS: acquire per-session ReentrantLock + MS->>IJ: enqueue(target, content) + MS->>RV: open(target) → CompletableFuture + IJ->>W: inject when injectable (one msg/turn) + IJ-->>MS: onDelivered (capture waiter + baseline) + alt worker replies explicitly + W->>RV: bridge_reply → resolve(REPLY) + else turn completes without reply + IJ-->>RV: onTurnComplete → resolveCompletion(COMPLETION) + else worker fails / wedges + IJ-->>RV: onTurnFailed → resolveFailure(FAILED) + end + RV-->>MS: Resolution + MS-->>P: Reply (REPLIED / COMPLETED_UNREPLIED / WORKER_FAILED / TIMED_OUT_*) +``` + +*Figure 3 — forward send. If the timeout fires first, the outcome is `TIMED_OUT_WORKING` +(delivered) or `TIMED_OUT_QUEUED` (never delivered).* + +## Flow: reverse rendezvous (`bridge_ask` → answer) + +A worker pauses its own turn to ask the primary a question; the question surfaces on the primary's +open forward send, and the answer resumes the *same* worker turn. Duplicate asks from one session +coalesce onto a single `turnId` (only the `fresh` owner surfaces the question and tears the turn +down). + +```mermaid +sequenceDiagram + autonumber + participant W as Worker + participant MS as MessageService + participant RV as Rendezvous + participant P as Primary + + W->>MS: ask(session, question) + MS->>RV: openAsk(session) + Note over RV: coalesce by session — mint turnId (fresh=true)
or ride the open one (fresh=false) + alt fresh owner, a forward send is waiting + RV->>RV: resolveQuestion(session, question, turnId) + RV-->>P: forward send returns QUESTION + turnId + P->>MS: answer(turnId, content) + MS->>RV: answerAsk(turnId, content) → unblock worker + MS->>RV: open(workerSession) — new forward waiter + W-->>RV: resumes turn → bridge_reply + RV-->>P: Reply + else no forward send open + RV-->>W: AskOutcome.NO_WAITER (closeAsk) + end +``` + +*Figure 4 — reverse ask. Worker session is derived from `turnId` (`askSession`), never a caller +argument. If the `turnId` is unknown, `answer` returns `Outcome.STALE_TURN`.* + +## State machine: `AgentStatus` + +The worker's herdr-reported lifecycle. `fromWire` maps `idle|working|blocked|done` and folds any +unrecognized string to `UNKNOWN`. A status is **injectable** (safe to deliver into) only when +`IDLE`, `BLOCKED`, or `DONE` — never mid-turn (`WORKING`). `UNKNOWN` both wedges delivery and is +reclassified by `StatusRefiner` from pane content. + +```mermaid +stateDiagram-v2 + [*] --> IDLE + IDLE --> WORKING: message injected, turn starts + WORKING --> IDLE: turn completes + WORKING --> BLOCKED: worker asks / awaits input + WORKING --> DONE: worker exits turn + BLOCKED --> WORKING: unblocked + DONE --> IDLE: next turn + IDLE --> UNKNOWN: herdr can't classify + WORKING --> UNKNOWN: herdr can't classify + UNKNOWN --> IDLE: StatusRefiner reads pane + UNKNOWN --> WORKING: StatusRefiner reads pane + + note right of WORKING + not injectable + end note + note right of DONE + injectable; working to done + is a turn-boundary like IDLE + end note +``` + +*Figure 5 — `AgentStatus`. Injectable states: `IDLE`, `BLOCKED`, `DONE`.* + +## State machine: `Injector` per-target turn lifecycle + +The heart of delivery reliability. Each target moves through a lifecycle driven by +`Injector.onStatus`; a turn is judged complete **only** from a confirmed `WORKING → injectable` +boundary (a real `WORKING` sample seen, then an injectable one). Three grace counters bound the +failure paths: `PICKUP_GRACE_POLLS = 8`, `TURN_STALL_GRACE_POLLS = 120`, +`READINESS_GRACE_POLLS = 240`. + +```mermaid +stateDiagram-v2 + [*] --> Quiescent + Quiescent --> Queued: enqueue(text) + Queued --> AwaitingPickup: injectable and ready and send ok + Queued --> Failed: 240 polls not ready + + AwaitingPickup --> AwaitingCompletion: WORKING observed + AwaitingPickup --> Quiescent: 8 injectable polls, no WORKING (missed fast turn) + + AwaitingCompletion --> Complete: next injectable after WORKING + AwaitingCompletion --> Failed: 120 UNKNOWN polls (stall) + + Complete --> Quiescent: onTurnComplete fires + Failed --> Quiescent: onTurnFailed fires + Queued --> Failed: drop(cause) + AwaitingCompletion --> Failed: drop(cause) + + Complete --> [*] + Failed --> [*] +``` + +*Figure 6 — per-target turn lifecycle in `Injector`. `Complete` fires `TurnListener.onTurnComplete` +(→ `CompletionResolver` scrapes the pane and resolves the captured waiter); `Failed` fires +`onTurnFailed`. `CompletionResolver` captures the exact waiter at delivery time and suppresses any +scrape matching the pre-turn baseline, so a late completion for turn N can never resolve turn +N+1's waiter (CB-116 safety).* + +## Subscription boundary + +The one non-negotiable invariant, enforced in code. `SubscriptionGuard.assertWorker(baseUrl)` runs +in `WorkerService.spawn()` **before any herdr call**: the worker's `ANTHROPIC_BASE_URL` host must +be on the allowlist (Stage-1: `gx00.gw`, `ollama.ltms.dev`). `assertPrimaryClean` (called at +startup) hard-stops if the primary env carries any `ANTHROPIC_BASE_URL`. Spawn injects +`ANTHROPIC_BASE_URL`/`ANTHROPIC_MODEL`/`CLAUDE_CONFIG_DIR`/`ANTHROPIC_AUTH_TOKEN` into the +**worker's** env only — `bridged`'s own env is never mutated. CWD resolves +`requestedCwd → profile cwd → caller cwd → daemon user.dir → "."`. + +## Concurrency model (at a glance) + +- **`herdr`** — lock-free: each `UnixSocketHerdrClient.call` owns its own socket; safe to call + concurrently from virtual threads. `WorkspaceControl.ensureWorkspace` is `synchronized`. +- **`inject`** — one virtual-thread poller loop; each target has its own monitor so + check-and-send is serialized; futures/listeners complete *after* the monitor is released to + keep the poller unblocked. +- **`msg`** — per-session `ReentrantLock` serializes forward sends (at most one forward waiter + per session); `Rendezvous` uses `ConcurrentHashMap` + `CompletableFuture` for cross-thread + resolution; async sends run on a dedicated virtual-thread executor. At most one open reverse-ask + turn per session (coalesced via `openAsksBySession`). +- **`mcp`** — identity resolution shells out to `lsof` once per worker-identity request; shared + services own their own concurrency. +- **`worker`** — `AtomicLong` name sequence + per-process nonce avoid herdr name collisions + without locks; the orphan reaper is best-effort and never aborts startup. + +## Related pages + +- [1. Architecture](1-Architecture) — system, invariants, the two modes (design level). +- [2. Message Server](2-Message-Server) — the `bridged` design rationale, herdr control contract. +- [8. Roadmap](8-Roadmap) — stages, tickets, and the feature ⇄ endpoint ⇄ test map. + +--- +*Generated from a fan-out code audit (one worker per layer) and verified against source at +`aa0cf81`.* diff --git a/_Sidebar.md b/_Sidebar.md index ccbc20f..f33807e 100644 --- a/_Sidebar.md +++ b/_Sidebar.md @@ -12,6 +12,7 @@ 6. [Team](6-Team) — orchestrating a mixed fleet 7. [Use Cases](7-Use-Cases) — the review scenario + mechanisms 8. [Roadmap](8-Roadmap) — stages, tech stack, tickets +9. [Implementation](9-Implementation) — as-built code map · classes · flows · state machines --- 🟢 herdr-centric `bridged` · AgentAPI = fallback