wiki: add page 9 — as-built implementation architecture
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.
+357
@@ -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)<br/>on-subscription"]
|
||||
worker["Worker (Claude Code)<br/>off-subscription"]
|
||||
|
||||
subgraph bridged["bridged daemon"]
|
||||
direction TB
|
||||
subgraph faces["North faces (thin adapters)"]
|
||||
rest["rest.BridgedApp<br/>REST · Javalin"]
|
||||
mcp["mcp.BridgeMcp<br/>MCP · /mcp servlet"]
|
||||
end
|
||||
msg["msg.MessageService + Rendezvous<br/>service core · per-session rendezvous"]
|
||||
inject["inject.Injector + StatusPoller<br/>single-writer, status-gated delivery"]
|
||||
guard["guard.SubscriptionGuard<br/>boundary enforcement"]
|
||||
worksvc["worker.WorkerService<br/>spawn · reap · teardown"]
|
||||
herdr["herdr.AgentControl / WorkspaceControl<br/>JSON-RPC over UNIX socket"]
|
||||
end
|
||||
|
||||
herdrd["herdr daemon<br/>(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<br/>+ assertPrimaryClean"] --> guard["SubscriptionGuard"]
|
||||
guard --> herdr["connect UnixSocketHerdrClient<br/>→ AgentControl / WorkspaceControl"]
|
||||
herdr --> wsvc["WorkerService<br/>→ reapOrphanWorkers()"]
|
||||
wsvc --> rv["Rendezvous → CompletionResolver<br/>→ Injector + StatusPoller"]
|
||||
rv --> ms["MessageService"]
|
||||
ms --> mcp["BridgeMcp<br/>(connection identity)"]
|
||||
ms --> app["BridgedApp<br/>(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)<br/>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`.*
|
||||
+1
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user