Compare commits

...

5 Commits

Author SHA1 Message Date
kevin 871b595954 CB-518: state the primary's orchestration as an explicit, ordered flow
CI / build (pull_request) Successful in 1m23s
CB-517 moved orchestration policy into CLAUDE.md but left it as a bullet
list, so the procedure was implicit: the order of operations had to be
reconstructed from a parallelisation bullet, and nothing said when to
review or when to tear down. A policy you have to reassemble on each
task is one you will reassemble differently each task. Restate the
primary's half as a numbered 0-8 flow — role check, split, gate, spawn
all, send all, collect, verify, review, adjudicate — so that following
it is checkable against the tool calls rather than a matter of recall.

Two steps carry the load. Spawn and send are separate on purpose:
folding them into one loop is what silently serialises work that was
meant to fan out. And review is now its own step ahead of the merge
rather than a clause inside it, because the two have opposite owners —
reviewers fan out over the diff (never the implementer of the scope
they review, and briefed from the diff rather than the author's
rationale, which carries the same blind spot), while adjudication, the
merge and teardown stay with the primary. Merging on a reviewer's word
is delegating the gate by proxy, so the step says so outright.

Nothing is dropped. The six bullets that trailed the tool table are
relocated into the step that owns each — profile explicitness into
spawn, playbook naming and self-containment into send, claim
verification into its own step, the ~60s blocking-send cap into a note
beneath the flow — and the table stays as the intent→tool lookup.

The wiki pointer moves with it. The block is canonical only if its
template matches byte for byte, so the template was produced by
splicing the block out of CLAUDE.md rather than by editing it in
parallel, and the sync check the repo documents passes. Bumping the
pointer in the same commit keeps charter and template versioned
together, as CB-517 did.

Co-Authored-By: Claude Opus 5 (1M context) <noreply@anthropic.com>
Claude-Session: https://claude.ai/code/session_017Kw1FosEt3Noix5GG9wJ2r
2026-08-04 22:32:19 +07:00
Dai Ha b67b1585c2 CB-517: deploy LavinMQ as a pinned, durable, self-restarting broker
CI / build (push) Successful in 1m23s
The CB-307 durable ReplyInbox needs an AMQP broker, but the one behind it
was run ad hoc and had simply vanished from the host — which takes the
whole daemon with it, since AmqpReplyInbox.open throws and Bridged.java:187
does not guard it. A missing broker is a hard startup failure, not a
degraded mode, so 'how the broker runs' is part of the system, not a local
detail.

Pinned to 2.9.1 (:latest would move the broker under a running daemon),
data on a named volume (held-but-unacked replies are the entire point of
Stage 2 — a plain 'compose down' would discard exactly what durability
protects), and restart: unless-stopped so it comes back after a reboot
instead of disappearing again.

Ports are bound to 127.0.0.1 deliberately: LavinMQ ships a default
guest/guest account, which is only acceptable while nothing off-host can
reach it.

Verified by driving the production AmqpReplyInbox against this deployment
(publish/peek/dedup/FIFO/ack, then reconnect): 8/8 including redelivery of
the unacked message. That pairing had never been exercised — the
@Tag("contract") test runs against a RabbitMQ container, and is excluded
from the default build, so mvn clean install covers the broker path zero
times.
2026-08-04 16:17:21 +02:00
Dai Ha 979b2b5632 CB-517: add bridge_whoami and make the bridge prompt a portable charter
CI / build (push) Successful in 1m25s
The communication rules lived only in two opt-in skills, so nothing
always-on told the primary how to orchestrate and nothing guaranteed a
worker loaded its playbook. Move protocol and policy into CLAUDE.md,
which a worker inherits for free (its worktree is a checkout of this
repo), and leave the skills as pure per-job procedure.

bridge_whoami closes the load-bearing gap: every tool already consumed
the caller identity ConnectionIdentity resolves from the connection, but
none reported it, so an agent had to infer its own role from side
channels the daemon does not control. Guessing fails asymmetrically — a
primary acting as a worker is refused by the authz gate and learns at
once, while a worker acting as the primary ends its turn without
bridge_reply and the sender silently receives nothing. The tool reuses
the same Principal the gate is built on, so the two cannot disagree; the
primary gets role only (handing it a sessionId it does not own would
invite the forged reply Authz refuses), and a worker missing from the
registry still gets role + sessionId rather than 'unknown'.

The CLAUDE.md block is written to be copied as-is into any project that
mounts the bridge: repo-local details (Authz paths, the .mcp.json/wiki
exclusions, the skill names) moved below it into a project addendum, and
every role-inference fallback is stated one-way — the mount-name signal
only holds for mcp__bridge__* (the launcher fixes it), not for the
primary's mount, which each project names itself. The wiki carries the
block verbatim as the template, with a sync check.

Because this repo IS the bridge, that block is shipped surface, not
documentation: the addendum adds a mandatory checklist mapping each part
of the code to the part of the prompt it can invalidate.

Also: delegate-by-default policy for the primary — the test is not 'could
I do this faster myself' but 'can I write a brief good enough for a
worker'.

mvn clean install: 356 tests green (353 + 3 for whoami); ide_diagnostics
clean on both changed files.
2026-08-04 16:01:16 +02:00
kevin cf4ad186ab CB-516: fail a delegation when its worker session is released
CI / build (push) Successful in 1m30s
Tearing a worker down left the send that was waiting on it stranded. The
rendezvous waiter stayed open, so a blocking bridge_send kept blocking and an
async one kept reporting PENDING until ASYNC_TIMEOUT_MS — thirty minutes —
even though the worker provably no longer existed and the delegation could
never complete.

Observed repeatedly this session: a task sitting at
{"phase":"pending","detail":"worker unknown"} long after I had personally
deleted the pane. The maddening part is that poll() already HAD the evidence —
it calls liveStatus() to build that detail string, gets back "unknown", and
reports PENDING anyway. It also never reached /metrics under any outcome, so a
stalled delegation was invisible to both the task view and the dashboard. The
only way I ever diagnosed one was reading the worker's pane over the herdr
socket by hand.

Fixed at the choke point rather than by guessing from status strings:
SessionManager.release() is the single path every teardown funnels through
(REST stop, MCP stop, idle-TTL reaper, recycle, shutdown drain), so it now
notifies a release listener with the terminalId, and Bridged.main wires that to
MessageService.abandon(). Abandon fails the open waiter with a real reason.

Deliberately NOT done by inferring "gone" from liveStatus(): that method
collapses a vanished pane, a wedged worker and a herdr hiccup into the same
"unknown" string, so acting on it would fail live delegations during a
transient blip. An explicit lifecycle signal cannot be ambiguous.

Notified before launcher.stop() so a blocked caller fails fast, and wrapped so
a listener failure can never prevent the teardown it is reacting to.

Resolving as a failure rather than letting it time out also means the outcome
is counted — a torn-down delegation now shows up as
sends_total{outcome="failed"} instead of nothing at all.

353 tests (was 346): abandon fails a waiting send / is a no-op with no waiter /
never clobbers a send the worker already answered, an abandoned async task
polls FAILED rather than PENDING, release notifies with the right terminal,
releasing an unknown pane notifies nobody, and a throwing listener does not
block teardown.

Verified live on the running daemon, reproducing the original scenario:
  async send        -> {"phase":"pending","detail":"worker working"}
  DELETE the worker -> 204
  poll              -> {"phase":"failed","detail":"the worker session was
                        released before it replied"}
  /metrics          -> bridged_sends_total{outcome="failed"} 1
Previously that poll returned PENDING for thirty minutes and the counter stayed
empty.

NOT addressed here, and worth deciding separately: a worker that simply never
replies (rather than being released) still rides out the full 30-minute
ASYNC_TIMEOUT_MS. That is a policy question about long-running delegations, not
a correctness bug.
2026-08-01 23:57:45 +07:00
kevin 5100f215cf CB-515: regression-protect the turn-attribution guards
CI / build (push) Successful in 1m33s
Five tests pinning the invariants that decide WHICH turn a reply belongs to.
These protect against a silent correctness bug — a reply attributed to the
wrong turn — not against a crash, which is why they were worth picking over
higher-percentage coverage gaps.

Chosen by blast radius, not by uncovered-line count. Both guards are compound
conditions with a side that never executed, i.e. exactly the shape where a
clause can be deleted as "redundant" and every existing test still passes.

CompletionResolver:
- The CB-115 misattribution guard suppresses a completion when the scrape is
  byte-identical to the pane at delivery. Its !scrapeFailed clause was
  unexercised: delete it and a FAILED read is misread as "no output change",
  so the send is suppressed and hangs to the caller's timeout instead of
  resolving. The new test sets the baseline to "" so the empty tail from a
  failed read would byte-match and wrongly suppress — built to die precisely
  when that clause dies.
- The fail() guard leaves an already-resolved waiter alone. The new test also
  asserts agent.read is never called, so the worker is not scraped for a send
  nobody is waiting on.

Rendezvous: a second resolution of an already-completed waiter returns false
and does not overwrite the first value, for both resolveCompletion and
resolveFailure.

Verified by sabotage, one guard at a time: removing !scrapeFailed reds
resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent; removing the
isDone() clause reds failLeavesAnAlreadyResolvedWaiterUntouchedAndSkipsTheScrape.
(The first attempt at the second sabotage reported a false pass — the patch hit
an identically-worded guard earlier in the file. Line-targeted and re-run.)

346 tests, was 341. Worker-implemented on the local-vLLM profile; it noticed
three of the eight cases I asked for already existed and said so with names
rather than duplicating them.

Also of note: the first delegation of this ticket wedged the worker — the pane
showed a zsh parse error and it went idle with an untouched worktree, task stuck
pending. The retry differed only in phrasing the same requirements as prose
instead of quoting Java boolean expressions. Filed as a bridge robustness
concern: injected content shares a channel with control, and a wedged turn is
invisible in both the task view and /metrics.
2026-08-01 23:48:05 +07:00
14 changed files with 681 additions and 113 deletions
+37 -54
View File
@@ -1,24 +1,21 @@
---
name: implementer
description: Implementer-role playbook for a bridged worker — you are in an isolated git worktree on a dedicated branch; implement the assigned task, commit, push, open your own PR to main, and hand off the PR URL via bridge_reply. You never merge. Load this when you have been delegated an implementation task over bridged.
description: Implementer-role procedure for a bridged worker — verify your worktree, implement the scope, commit, push, open your own PR, and hand off the PR URL. Load this when the lead delegates you an implementation task over bridged.
---
# Implementer worker
# Implementer worker — procedure
You are an **implementer** in the claude-bridge fleet. The lead delegated you one scoped task,
and you are running in an **isolated git worktree on your own branch** — a full peer of the
primary (same `CLAUDE.md`, skills, memory, MCP), differing only in the model behind you and the
branch you sit on. Your job for this turn: **implement the task, then hand off a PR the lead can
review and merge.** You do the work; the lead (or human) is the merge gate — you never merge.
The turn contract (one `bridge_reply`, `bridge_ask` for the lead's decisions, honest reporting,
never merge, never commit `.mcp.json` or `wiki/`) is in **`CLAUDE.md` → Bridge communication →
Worker** and already applies. This skill is only the *implement-and-hand-off procedure*.
Delivery mechanics (how the task reached you, how your reply resolves the lead's blocked send)
are in [`docs/MCP-Contract.md`](../../../docs/MCP-Contract.md); the worktree/PR model is in
[`docs/Worker-Git-Workflow.md`](../../../docs/Worker-Git-Workflow.md). You only need the steps
below.
You run in an **isolated git worktree on your own branch** — a full peer of the primary (same
repo, `CLAUDE.md`, skills, MCP), differing in the model behind you and the branch you sit on.
The worktree model is documented in [`docs/Worker-Git-Workflow.md`](../../../docs/Worker-Git-Workflow.md).
## 1. Confirm where you are — a worktree on a dedicated branch
## 1. Confirm where you are
Before touching anything, verify your ground truth:
Before touching anything:
```bash
git rev-parse --show-toplevel # your worktree root — NOT the primary's main tree
@@ -26,46 +23,40 @@ git branch --show-current # your dedicated branch: worker/<ticket>-<nonce>
git status # should be clean at the start
```
Do **all** work here, on this branch. **Never** switch to `main`, never `git checkout main`,
never rebase onto or push to `main` directly. The branch is your isolation — respect it.
Do **all** work here, on this branch. Never `git checkout main`, never rebase onto or push to
`main`. The branch is your isolation — respect it.
## 2. Implement the task
## 2. Implement
- Implement exactly the scope the lead named. Keep changes focused; if you notice something out
of scope, note it in your reply rather than expanding the diff.
- Match the surrounding code's style, naming, and idioms. Follow project `CLAUDE.md`.
- **You cannot run the IDE MCP tools** (intellij-index / jetbrains are the primary's, not yours).
So **never claim a file is "IDE-clean" or "diagnostics-clean"** — you cannot verify that. State
only what you actually ran (e.g. `mvn`, a test) and its real output. A fabricated clean claim is
worse than an honest "I could not verify inspections here."
- Run whatever build/test you can and **report the true result** — including failures.
- Implement exactly the scope the lead named. Keep the diff focused; note anything out of scope
in your reply instead of widening it.
- Match the surrounding code's style, naming, and idioms.
- Run whatever build/test you can — `mvn clean install` from the module root. Read its **full**
output; a piped `mvn ... | tail` hides failures.
## 3. Commit — focused, and never the excluded files
## 3. Commit
```bash
git add <the files you changed>
git add <the files you changed> # explicitly — never `git add -A` / `git add .`
git commit -m "<ticket>: <clear one-line summary>"
```
**Excluded from every commit, always:** `.mcp.json` (the primary's local, session-modified copy —
present only for parity) and `wiki/` (a separate submodule). Stage files explicitly; do **not**
`git add -A` / `git add .` blindly, or you risk staging them. If `.mcp.json` shows as modified,
leave it — it is flagged `--skip-worktree` and is not yours to commit.
`.mcp.json` will show as modified. Leave it — it is `--skip-worktree` and not yours to commit.
## 4. Push your branch
## 4. Push
```bash
git push -u origin HEAD
```
Push is over SSH as the same user — no extra credential needed. Push the branch as-is; do not
force-push over anything you did not create.
Push is over SSH as the same user — no extra credential needed. Never force-push over anything
you did not create.
## 5. Open your own PR to `main`
Open the PR via the gitea REST API. The daemon injected a **repo-scoped token** (`GITEA_TOKEN`)
and the forge host (`GITEA_HOST`) into your env for exactly this — the token can create a PR but
**cannot merge** (that stays the lead/human gate).
Via the gitea REST API. The daemon injected a **repo-scoped token** (`GITEA_TOKEN`) and the forge
host (`GITEA_HOST`) into your env for exactly this — the token can create a PR but **cannot
merge**.
```bash
API="${GITEA_HOST%/}/api/v1/repos/lms/claude-bridge/pulls"
@@ -81,16 +72,14 @@ JSON
)"
```
The response JSON includes `"html_url"` — that is your PR URL. If the call fails (non-2xx), read
the error body, fix the cause if it is yours (e.g. branch not pushed yet), and report the failure
honestly in your reply rather than inventing a URL. If `GITEA_TOKEN` is unset, your profile was
not granted PR-create — push the branch (step 4) and report the branch name so the lead opens the
PR.
The response JSON carries `"html_url"` — that is your PR URL. On a non-2xx, read the error body,
fix it if the cause is yours (e.g. branch not pushed yet), and report the failure rather than
inventing a URL. If `GITEA_TOKEN` is unset your profile was not granted PR-create: push the branch
and report its name so the lead opens the PR.
## 6. Reply via `bridge_reply` — the PR is the handoff
## 6. Hand off — what goes in `bridge_reply`
End your turn with **exactly one** `bridge_reply`. That reply is the entire handoff — the lead
cannot see your terminal. Include:
The reply is the entire handoff; the lead cannot see your terminal.
```
PR: <html_url from step 5, or "not created: <reason>" + branch name>
@@ -100,27 +89,21 @@ tests: <what you ran and its REAL result — or "not run: <why>">
summary: <2-3 lines: what you implemented and any caveat the reviewer needs>
```
Then stop. **Do not merge. Do not touch `.mcp.json` or `wiki/`.** One reply closes the turn.
```mermaid
sequenceDiagram
autonumber
participant L as Lead
participant B as bridged
participant I as Implementer (you)
participant G as git / gitea
L->>B: bridge_send(task) — blocks
B-->>I: your assignment (in a worktree on your branch)
L->>I: delegated task (you are in a worktree on your branch)
I->>I: implement + build/test here
I->>G: git commit (never .mcp.json / wiki)
I->>G: git push -u origin HEAD
I->>G: POST /pulls (GITEA_TOKEN) — open PR to main
G-->>I: html_url
I->>B: bridge_reply(PR url, branch, files, tests)
B-->>L: { outcome:"reply", text }
Note over L,G: lead reviews the PR, merges on green — you never merge
I->>L: bridge_reply(PR url, branch, files, tests)
Note over L,G: the lead reviews the PR and merges on green — you never merge
```
*The implement turn: work in the worktree, commit → push → open the PR, hand off the URL. The
lead is the merge gate.*
*The implement turn: work in the worktree, commit → push → open the PR, hand off the URL.*
+25 -58
View File
@@ -1,69 +1,39 @@
---
name: reviewer
description: Reviewer-role playbook for a bridged worker — read the assigned scope, find the real issues, ask the lead via bridge_ask when a decision is genuinely theirs, and report the finding via bridge_reply. Load this when you have been delegated a code review over bridged.
description: Reviewer-role procedure for a bridged worker — how to work a review scope and the exact shape of the finding to report. Load this when the lead delegates you a code review over bridged.
---
# Reviewer worker
# Reviewer worker — procedure
You are a **reviewer** in the claude-bridge fleet. The lead delegated you one scoped review
over `bridged`, and your whole job is **this single turn**: examine the scope it named, and
report back. You are not the owner of the code and you do not merge anything — you surface
what the owner needs to know, then hand the turn back.
Delivery mechanics (how the task reached you, how your reply resolves the lead's blocked
send) are in [`docs/MCP-Contract.md`](../../../docs/MCP-Contract.md); you only need the three
rules below.
The turn contract (one `bridge_reply`, `bridge_ask` for the lead's decisions, honest reporting,
never merge) is in **`CLAUDE.md` → Bridge communication → Worker** and already applies. This
skill is only the *review procedure*: how to work the scope, and the exact shape of what you
send back.
## 1. Read the whole scope before you judge
The delegation names your scope — a file, a diff, a PR, a function. **Read all of it first.**
A review that fires on a snippet misses the caller that makes it safe (or the one that makes
it a bug). Reviewing only part of the scope and guessing the rest is the most common way a
reviewer worker is wrong.
A review that fires on a snippet misses the caller that makes it safe (or the one that makes it
a bug). Reviewing part of the scope and guessing the rest is the most common way a reviewer is
wrong.
## 2. Stay in your lane
## 2. Stay in the scope
- Review **only** the assigned scope. If you notice something elsewhere, mention it in one
line — do **not** go hunt it. Wandering is how two workers end up reporting the same thing
and neither covers what it was given.
- Do **not** edit files, run the build, or spawn other workers. You review; the owner acts.
- You never set `ANTHROPIC_BASE_URL` and never touch herdr — you are a Claude Code process,
not part of the transport.
- Review **only** what you were assigned. Something elsewhere looks wrong? One line in your
reply — do not go hunt it. Wandering is how two reviewers report the same thing and neither
covers what it was given.
- Do **not** edit files or run the build. You review; the owner acts.
## 3. When the decision is the lead's — ask, don't guess
## 3. Reach for `bridge_ask` only for a genuine fork
Some things you cannot resolve from the code: an ambiguous requirement, a missing acceptance
criterion, "is this behavior intended or a bug?", or a choice between two defensible fixes.
Guessing there produces a confident-but-wrong finding. Instead **pause and ask the lead** with
`bridge_ask` — a single crisp question. The call blocks; when the lead answers you **resume
the same turn** with the answer and finish. Ask only when the answer changes your finding;
don't narrate options you could decide yourself.
Ambiguous requirement, a missing acceptance criterion, "intended or a bug?", or two defensible
fixes with different consequences — those are the lead's call, and guessing produces a
confident-but-wrong finding. Anything you could settle by reading more code is yours to settle.
```mermaid
sequenceDiagram
participant L as Lead
participant B as bridged
participant R as Reviewer (you)
## 4. The finding — what goes in `bridge_reply`
L->>B: bridge_send(review scope) — blocks
B-->>R: your assignment
R->>R: read the full scope
opt a decision only the lead can make
R->>B: bridge_ask("intended, or a bug?") — you block
B-->>L: { outcome:"question", turn_id }
L->>B: bridge_send(answer, turn_id)
B-->>R: { answer } — you resume the SAME turn
end
R->>B: bridge_reply(structured finding) — ends your turn
B-->>L: { outcome:"reply", text }
```
*The review turn, with the optional `bridge_ask` detour when the call is the lead's to make.*
## 4. Report with `bridge_reply` — one structured finding
End your turn with **exactly one** `bridge_reply`. Report the **single most important** real
issue in the scope, in these four lines, under ~90 words:
Report the **single most important** real issue in the scope, in these four lines, under
~90 words:
```
1. <path>:<line>
@@ -72,13 +42,10 @@ issue in the scope, in these four lines, under ~90 words:
4. severity: high | medium | low
```
- Found nothing real after reading? Reply `NO ISSUE` and one line saying why — a clean review
is a valid result, and a fabricated issue is worse than none.
- **Nothing real after reading?** Reply `NO ISSUE` and one line saying why. A clean review is a
valid result; a fabricated issue is worse than none.
- **Severity:** `high` = wrong result, data loss, security, or a hang/crash on a real path ·
`medium` = a real bug on an edge path, or a correctness risk under load/concurrency ·
`low` = clarity, a latent foot-gun, or a smell with no current failure.
- Be specific and verifiable: a line number and a one-line repro beat an adjective. If you
can't point to where it goes wrong, you haven't found it yet.
One reply closes the turn. If you asked mid-turn, the answer you got is already folded into
this finding — you do not ask again after replying.
- Be specific and verifiable: a line number and a one-line repro beat an adjective. If you can't
point at where it goes wrong, you haven't found it yet.
+189
View File
@@ -1,7 +1,196 @@
# claude-bridge — project instructions
## Bridge communication (enforced — read this first)
> **Canonical block.** Everything down to §Layering is the portable bridge charter, copied verbatim
> into every project that mounts the bridge MCP. Keep it byte-identical with the template in the
> wiki ([Use Cases](https://git.ltms.dev/lms/claude-bridge/wiki/7-Use-Cases) → *The portable
> CLAUDE.md block*); improvements go to the template first, then out to each project. Anything
> specific to *this* repo lives under §Project addendum below, never inline above it.
If no `bridge_*` MCP tools are mounted in this session, this section does not apply — skip it.
`bridged` is the **sole communication gateway** between agents here. The orchestrating session (the
**primary**) and every delegated peer (a **worker**) mount the *same* MCP server and talk only
through its `bridge_*` tools. No session addresses a peer, a broker, or the network directly.
### Which role am I? — settle this before acting
**Both roles read this file.** A worker runs in a git worktree of this same repo, so it inherits
this `CLAUDE.md` verbatim, and every rule below is role-conditional.
**Call `bridge_whoami`.** It returns `{"role":"primary"}` or `{"role":"worker","sessionId":…,
"profile":…,"worktree":…,"branch":…}`, resolved by the daemon from your connection — unforgeable,
and the same resolution its authorization gate uses. Don't infer what you can ask.
Only if that call is unavailable, fall back to these — each is one-way, so keep reading until one
fires: the reply charter in your system prompt (*"You are an off-subscription worker in the
claude-bridge fleet"*) ⇒ **worker**; bridge tools prefixed `mcp__bridge__*` ⇒ **worker** (the
launcher fixes that mount name; a primary's mount is named by whoever wrote its `.mcp.json`, so it
varies); `ANTHROPIC_BASE_URL` set ⇒ **worker** (Claude-model workers run on a clean env, so its
*absence* proves nothing). **Still unsure ⇒ act as a worker.** The two mistakes are not symmetric: a
primary acting as a worker is refused by the authorization gate — loud and self-correcting — while a
worker acting as the primary ends its turn with no `bridge_reply`, and the sender silently receives
nothing. Fail toward the recoverable error.
### Invariants — both roles, no exceptions
1. **Never set, export, or forward `ANTHROPIC_BASE_URL`** (or `ANTHROPIC_AUTH_TOKEN`). The primary
stays on subscription; only the bridge puts a worker off it, at spawn. Mounting the bridge must
never move a session across that boundary.
2. **The bridge is the only channel.** Text you print in your terminal reaches nobody — the other
side cannot see your screen. An answer that isn't in a `bridge_*` call is silently discarded.
3. **Identity comes from the connection, never an argument.** Workers never pass a target; you
cannot act as another session. Spawn/stop/send/drain are primary-only; reply/ask are
worker-only-and-only-as-itself. A call outside your role is refused, not queued.
4. **Delivery is status-gated: one message per turn.** Don't busy-poll a peer's terminal and don't
re-send because a call looks slow — the bridge delivers when the peer is `idle`/`blocked`.
5. **Never drive the terminal multiplexer directly** (no `herdr` CLI, no socket). The bridge owns
policy; the multiplexer owns PTYs. Going around the bridge bypasses every rule above.
### Primary (lead) — run this on every task, in order
**Delegate by default — that is the job.** With the bridge mounted you are an orchestrator on a
metered subscription, and workers are cheap, parallel, and disposable. The default answer to "who
does this?" is **a worker**, not you. Reach for `bridge_send` before you reach for `Edit`. The steps
below are the procedure — run them in order, every task, not only the big ones.
0. **Know your role** — `bridge_whoami`, once per session, before anything else.
1. **Split.** Write the unit list. Every unit carries: scope · the files or PR in question ·
acceptance criteria · exactly what to report back. A unit with no acceptance criteria is not
ready to delegate — refine it or keep it.
2. **Gate each unit** on one question: **"can I write a brief good enough for a worker to
succeed?"** — *not* "could I do this faster myself?" (usually you could; doing it yourself costs
your context and your subscription, while a wasted worker turn costs a worker turn). Yes ⇒
delegate. The keep-list is closed: the conversation with the user, decomposition and planning,
the final judgment call, verification, merges, and anything that depends on context only you
hold. Nothing else is yours by default.
3. **Spawn every delegated unit first** — `bridge_spawn{profile, worktree:true, ticket}`, one per
unit, *before* sending any. Pass `profile` explicitly: profiles differ in model and cost, not in
tier, so the default is rarely what you want.
4. **Then send them all** — `bridge_send{sessionId, content, wait:false}`. Line 1 of every brief is
`Load the <name> skill.` naming the worker's playbook; those skills are opt-in and that line is
what makes them reliable. Where the project ships no such skill, spell the procedure out in the
brief instead. The brief is self-contained — the worker sees your message and the repo, nothing
of your context, your plan, or your screen.
5. **Collect** — `bridge_poll{ticket}` → `bridge_ack{ticket, msgId}`. Answer a worker's `bridge_ask`
with `bridge_send{turnId, content}` — **not** `sessionId`. A worker gone quiet is diagnosed with
`bridge_status`, never by reading its terminal.
6. **Verify yourself.** Re-run the build and the checks. A worker mounts only the bridge MCP and
cannot run your other tooling, and a piped command (`… | tail`) hides failures behind a zero
exit — never promote a worker's "clean" to a fact.
7. **Review — fan out.** Spawn reviewers against the diff, one per dimension or per file, with
`wait:false`. Never the implementer of the scope it reviews, and brief them from the diff — not
from the implementer's rationale, which carries its own blind spot. Dispatch each PR's reviewers
as it lands; don't wait for the last implementer. Under ~50 changed lines, skip the fan-out and
read it yourself.
8. **Adjudicate, merge, tear down — yours alone.** Read the diff yourself: fully if it is small,
targeted at the reported findings and the risky paths if it is large. Reviewer findings direct
your attention; they never substitute for it. Then merge, then `bridge_stop{paneId}`.
**Steps 3 and 4 are separate on purpose** — spawning and sending in one loop is how parallel work
silently becomes serial, and it is the most common way this layer is wasted. For the same reason,
prefer `wait:false` + `bridge_poll` for anything non-trivial: a blocking `bridge_send` is capped by
*your own* MCP client call timeout (~60s), well below the task's real runtime.
**Delegating does not delegate responsibility.** Workers open PRs; you are the gate. Never delegate
the merge — and merging on a reviewer's word is delegating it by proxy.
| Intent | Tool |
|---|---|
| Confirm your own role | `bridge_whoami` |
| See backends available | `bridge_profiles` |
| Start a worker | `bridge_spawn{profile?, cwd?, worktree?, ticket?}` → `sessionId` + `paneId` |
| See the fleet | `bridge_list` · one worker's state: `bridge_status{sessionId}` |
| Delegate (blocking) | `bridge_send{sessionId, content}` |
| Delegate (long task) | `bridge_send{sessionId, content, wait:false}` → ticket → `bridge_poll{ticket}` |
| Answer a worker's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
| Collect a held reply | `bridge_poll{target}` · then `bridge_ack{target, msgId}` |
| Tear down | `bridge_stop{paneId}` |
### Worker — the turn contract
1. **Load the playbook skill the lead named** before doing anything else.
2. **Do the assigned scope only.** Note anything you spot outside it in one line; don't go hunt it.
3. **`bridge_ask{question}`** when a decision is genuinely the lead's (ambiguous requirement, two
defensible fixes, "bug or intended?"). It blocks and you resume the *same* turn with the answer.
Don't ask what you could decide yourself.
4. **End the turn with exactly one `bridge_reply{content}`**, carrying your complete answer. This is
the whole handoff. No `bridge_reply` ⇒ the sender gets nothing and the exchange stalls.
5. **Report honestly.** State only what you actually ran and its real output, including failures.
You mount **only** the bridge MCP — the primary's other servers (IDE, forge, docs) are not yours,
so never claim the result of a check you had no way to run.
6. **Never merge.** Stage files explicitly — never `git add -A` — and leave alone anything the
project marks as not-yours-to-commit.
### Where each rule lives (don't duplicate — extend the right layer)
| Layer | Scope | Reaches |
|---|---|---|
| the launcher's reply charter | the one rule that must survive with no repo: *end every turn with `bridge_reply`* | every worker, at launch, every peer kind |
| **this section** | protocol + orchestration policy | primary **and** every Claude worker — tracked in git, so worktrees inherit it |
| role playbook skills | per-job procedure (commit/PR recipe, finding format) | a worker told to load one |
| the bridge's own docs | design detail, flows, error model | on demand |
A rule belongs in **exactly one** layer — the outermost one that must obey it. Peers that don't read
`CLAUDE.md` (non-Claude adapters) get the charter only, so any rule *they* must obey belongs in the
charter, not here.
## Project addendum — claude-bridge (not part of the canonical block)
- **This repo is the bridge.** The daemon is `bridged`, its MCP mount is `http://127.0.0.1:8765/mcp`,
and the code behind the rules above is `mcp/BridgeMcp` (tools), `auth/Authz` (the role table),
`mcp/ConnectionIdentity` (connection→role), and `worker/*Launcher` (`REPLY_CHARTER`).
- **Skills available to delegate:** `implementer` (worktree → commit → push → own PR) and
`reviewer` (scoped review → one structured finding). Name one in every delegation.
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
(a submodule with its own remote).
- **Flows and the error model** — rendezvous, `bridge_ask`, detached delivery, turn-done fallback —
are diagrammed in `docs/MCP-Contract.md` §6, kept out of this file because it loads into every
session's context.
### The prompt is part of the product — update it with the code (mandatory)
This repo *is* the bridge, so the canonical block above is not documentation about someone else's
system: it is the instruction surface this codebase ships. **Every change here must end by asking
whether the block still tells the truth.** A code change that silently invalidates it is an
incomplete change — the agents reading it have no other source.
Before you call any work done, check the row that matches what you touched:
| You changed… | Re-read and update… |
|---|---|
| a `bridge_*` tool — added, removed, renamed, or its params/semantics | the primary's intent→tool table; any rule that names that tool |
| `Authz` / the role table | invariant 3, and the primary-only vs worker-only claims |
| `ConnectionIdentity` / how a caller is resolved | the `bridge_whoami` paragraph and the fallback ladder |
| `REPLY_CHARTER`, or a launcher's mount/flags | the fallback ladder (`mcp__bridge__*`), and the layering table's top row |
| the injector / status gating | invariant 4 |
| worktree provisioning or the parity overlay | the "both roles read this file" premise — it rests on the worker's worktree being a checkout of this repo |
| `.claude/skills/**` | the addendum's skill list, and the "name the playbook" rule |
| a new peer kind (non-Claude adapter) | what that peer can read — anything it must obey belongs in its charter, not in the block |
Then **propagate**: the block in this file and the template in the wiki
([Use Cases](https://git.ltms.dev/lms/claude-bridge/wiki/7-Use-Cases) → *The portable `CLAUDE.md`
block*) must stay byte-identical, and other projects carrying the block need the same edit. Verify
rather than trust:
```bash
python3 - <<'PY'
import pathlib
c = pathlib.Path("CLAUDE.md").read_text()
w = pathlib.Path("wiki/7-Use-Cases.md").read_text()
S, E = "## Bridge communication (enforced", "## Project addendum — claude-bridge"
block = c[c.index(S):c.index(E)].rstrip() + "\n"
i = w.index("```markdown\n") + len("```markdown\n")
print("in sync:", w[i:w.index("\n```\n", i) + 1] == block)
PY
```
## IDE MCP tools & validation workflow (enforced)
> **Primary only.** Workers have no IDE MCP mount — if you are a worker, skip this section and
> report the build/test output you actually ran (see §Bridge communication → Worker).
Two IDE MCP servers are connected: **intellij-index** (semantic code intelligence) and
**jetbrains** (file problems, reformat, debugger). IntelliJ has multiple projects open; our
module is **`bridged`**. Always pass these to IDE MCP tools:
@@ -209,6 +209,12 @@ public final class Bridged {
MessageService messages = new MessageService(agents, injector, rendezvous, replyInbox,
pushLoop, metrics);
// CB-516: releasing a worker must fail whatever send was waiting on it. Without this a
// torn-down delegation kept reporting PENDING until the 30-minute async timeout, and never
// reached /metrics — the delegation was unresolvable and nothing said so.
sessions.onRelease(terminal ->
messages.abandon(terminal, "the worker session was released before it replied"));
// MCP server face (CB-105): bridge_send/bridge_reply/bridge_status, mounted at /mcp.
// Caller identity is resolved from the connection (peer PID → herdr pane), not arguments.
ConnectionIdentity identity = new ConnectionIdentity(
@@ -198,6 +198,11 @@ public final class BridgeMcp {
if (denied != null) return denied;
return profiles(workers);
})
.toolCall(whoamiTool(), (exchange, _) -> {
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.READ, null);
if (denied != null) return denied;
return whoami(principal(exchange), sessions);
})
.build();
this.authz = callers;
this.metrics = metrics;
@@ -460,6 +465,49 @@ public final class BridgeMcp {
}
}
/**
* {@code bridge_whoami}: the caller's own identity, as the daemon already resolved it.
*
* <p>Every other tool <em>consumes</em> this identity — the authorization gate, the reply
* rendezvous, the cwd inherit — but none reported it, so an agent had to infer its own role
* from side channels the daemon does not control: a charter string in its system prompt, the
* name its MCP mount happens to carry, or {@code ANTHROPIC_BASE_URL} (which Claude-model
* workers do not set). The failure mode of guessing is asymmetric and silent: a primary that
* mistakes itself for a worker is refused by {@link Authz} and learns immediately, while a
* worker that mistakes itself for the primary ends its turn without {@code bridge_reply} and
* the sender simply receives nothing. This tool removes the guess.
*
* <p>For a worker the session registry adds what it knows about that session. A worker the
* registry has no record of — one that outlived a daemon restart — still gets its role and
* {@code sessionId}, which is the load-bearing part.
*/
static McpSchema.CallToolResult whoami(Principal caller, SessionManager sessions) {
Map<String, Object> m = new LinkedHashMap<>();
m.put("role", caller.role().name().toLowerCase());
if (!caller.isWorker()) {
return text(json(m));
}
m.put("sessionId", caller.terminal());
sessions.roster().stream()
.filter(s -> caller.terminal().equals(s.terminalId()))
.findFirst()
.ifPresent(s -> {
m.put("paneId", s.paneId());
m.put("profile", s.profile());
m.put("state", s.state().name().toLowerCase());
if (s.worktree() != null) {
m.put("worktree", s.worktree());
}
if (s.branch() != null) {
m.put("branch", s.branch());
}
if (s.ownerTerminal() != null) {
m.put("owner", s.ownerTerminal());
}
});
return text(json(m));
}
// --- fleet management logic (CB-108 / CB-301) --------------------------------------------
/** {@code bridge_spawn} without cwd/caller context (default resolution). */
@@ -689,6 +737,18 @@ public final class BridgeMcp {
List.of("sessionId")));
}
private static McpSchema.Tool whoamiTool() {
return tool("bridge_whoami",
"Report who YOU are on the bridge — your role is resolved from your connection "
+ "(unforgeable), never from anything you claim. Returns role 'primary' (you "
+ "orchestrate: spawn/send/stop, and you must never call bridge_reply) or "
+ "'worker' (you were delegated to: you must end every turn with exactly one "
+ "bridge_reply, and cannot spawn or send), plus your own sessionId, profile, "
+ "worktree and branch when you are a worker. Call this first when following "
+ "role-conditional instructions rather than guessing your role.",
objectSchema(Map.of(), List.of()));
}
// --- small helpers -------------------------------------------------------------------------
// The SDK 2.0.0 deprecates its own Tool builders without a stable replacement — isolate it here.
@@ -255,6 +255,33 @@ public final class MessageService {
};
}
/**
* Abandon any send still waiting on {@code target} because its session has gone away (CB-516).
*
* <p>Without this, tearing a worker down left its rendezvous waiter open: a blocking
* {@code bridge_send} kept blocking, and an async one kept reporting {@code PENDING} until
* {@link #ASYNC_TIMEOUT_MS} — thirty minutes — even though the worker provably no longer
* existed and the delegation could never complete. Worse, {@code poll} already had the evidence
* (it calls {@code liveStatus} to build its detail string and gets back {@code "unknown"}) and
* reported {@code PENDING} anyway.
*
* <p>Resolving the waiter as a failure — rather than letting it time out — also means the
* outcome is counted, so a torn-down delegation stops being invisible to {@code /metrics}.
*
* @return true if a live waiter was failed
*/
public boolean abandon(String target, String reason) {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.currentWaiter(target);
if (waiter == null || waiter.isDone()) {
return false; // nobody is blocked on this worker — nothing to abandon
}
boolean failed = rendezvous.resolveFailure(waiter, reason);
if (failed) {
log.debug("abandoned send to {}: {}", target, reason);
}
return failed;
}
/**
* Acknowledge a specific reply by {@code msgId} for {@code target}. Removes it from the inbox
* so that a subsequent drain or peek no longer returns it.
@@ -17,6 +17,7 @@ import java.util.Optional;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
import java.util.function.Consumer;
import java.util.function.LongSupplier;
/**
@@ -47,6 +48,9 @@ public final class SessionManager implements TurnListener {
private final LongSupplier nowNanos;
private final int contextCap;
/** CB-516: notified with a terminalId on every release; no-op until wired. */
private volatile Consumer<String> releaseListener = _ -> { };
/** Backward-compatible constructor: shared-tree sessions, production git seam. */
public SessionManager(PeerLauncher launcher) {
this(launcher, new GitWorktrees(), System::nanoTime, 0);
@@ -136,6 +140,10 @@ public final class SessionManager implements TurnListener {
if (removed != null) {
log.debug("releasing session pane={} terminal={} state={}",
removed.paneId(), removed.terminalId(), removed.state());
// CB-516: a send still waiting on this worker can never be answered now. Tell the
// listener BEFORE the pane is torn down, so a blocked caller fails fast with a real
// reason instead of sitting on a rendezvous nothing will ever resolve.
notifyReleased(removed.terminalId());
}
launcher.stop(paneId);
if (removed != null && removed.worktree() != null) {
@@ -143,6 +151,32 @@ public final class SessionManager implements TurnListener {
}
}
/**
* Register a callback invoked with a session's {@code terminalId} whenever it is released
* (CB-516). Every teardown path funnels through {@link #release}, so one hook covers the REST
* and MCP stop tools, the idle-TTL reaper, {@code recycle}, and shutdown drain alike.
*
* <p>Set rather than injected because {@code MessageService} — the intended listener — is
* constructed after this manager (it needs the injector and rendezvous, which need the session
* presence view this manager exposes). Wiring it at construction would require breaking that
* cycle for one callback.
*/
public void onRelease(Consumer<String> listener) {
this.releaseListener = (listener == null) ? _ -> { } : listener;
}
/** A listener failure must never prevent the teardown it is reacting to. */
private void notifyReleased(String terminalId) {
if (terminalId == null) {
return;
}
try {
releaseListener.accept(terminalId);
} catch (RuntimeException e) {
log.warn("release listener failed for terminal {}: {}", terminalId, e.toString());
}
}
private WorkerSession acquireWithWorktree(String profile, String requestedCwd, String callerCwd,
String ownerTerminal, WorktreeRequest wt) {
String resolvedProfile = (profile == null || profile.isBlank())
@@ -195,6 +195,66 @@ class CompletionResolverTest {
assertEquals("hello", waiter.getNow(null).text());
}
@Test
void resolvesWhenTheScrapeItselfFailsEvenWithABaselinePresent() {
// The most important branch of the CB-115 guard: a failed read means the resolver could not
// SEE the screen — "couldn't see", not "no change". It must still resolve the send (an empty
// tail beats hanging until the caller's timeout), even though a baseline was captured. The
// baseline here is "" (an empty pane at delivery), so without the !scrapeFailed clause the
// byte-identical guard would wrongly match the empty tail and suppress.
FakeHerdr herdr = new FakeHerdr().healthy(false); // agent.read throws HerdrException
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, ""); // empty pane baselined at delivery
resolver.resolve("term_a", turn);
assertTrue(waiter.isDone(),
"a failed scrape must still resolve the send, not hang until the caller's timeout");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("", waiter.getNow(null).text(), "the tail is empty because the screen was unreadable");
}
// --- CB-115/CB-116 fail guard: an already-done or absent waiter is left alone ---------
@Test
void failLeavesAnAlreadyResolvedWaiterUntouchedAndSkipsTheScrape() {
// The send was already resolved (e.g. by the worker's explicit reply) before fail fired.
// fail must not overwrite that value, and must not even scrape the worker — nobody needs it.
FakeHerdr herdr = new FakeHerdr().readText("an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a");
var turn = new CompletionResolver.InFlight(waiter, null);
assertTrue(rendezvous.resolveCompletion(waiter, "already replied"));
resolver.fail("term_a", turn);
assertFalse(herdr.called("agent.read"),
"fail must not scrape a waiter that is already done");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind(),
"fail must not overwrite the existing resolution");
assertEquals("already replied", waiter.getNow(null).text());
}
@Test
void failFallsBackToTheRegisteredWaiterWhenThereIsNoInFlightTurn() {
// A never-delivered readiness failure has no in-flight record but still has a blocked send;
// fail falls back to the waiter currently registered on the Rendezvous and fails it.
FakeHerdr herdr = new FakeHerdr().readText("stuck on an error screen");
Rendezvous rendezvous = new Rendezvous();
CompletionResolver resolver = new CompletionResolver(new AgentControl(herdr), rendezvous);
var waiter = rendezvous.open("term_a"); // send registered, but no captureBaseline ever ran
resolver.fail("term_a", null); // no in-flight turn → fall back to the registered waiter
assertTrue(waiter.isDone(), "fail falls back to the registered waiter when no turn is in flight");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertEquals("stuck on an error screen", waiter.getNow(null).text());
}
// --- CB-116 waiter identity: a late completion never crosses into the next turn ---------
@Test
@@ -1,5 +1,6 @@
package dev.ltms.bridged.mcp;
import dev.ltms.bridged.auth.Principal;
import dev.ltms.bridged.config.BridgedConfig;
import dev.ltms.bridged.guard.SubscriptionGuard;
import dev.ltms.bridged.herdr.AgentControl;
@@ -343,4 +344,58 @@ class BridgeMcpTest {
assertNotEquals(Boolean.TRUE, res.isError());
assertEquals("blocked", textOf(res));
}
// --- bridge_whoami: the caller's own identity, so an agent never has to guess its role -------
@Test
void whoamiReportsThePrimaryAsPrimaryAndNothingElse() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.whoami(
Principal.primary(100), sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"role\":\"primary\""), out);
// The primary owns no session — leaking a sessionId here would invite it to reply as one.
assertFalse(out.contains("sessionId"), out);
}
@Test
void whoamiReportsAWorkerWithItsRegisteredSession() {
FakeHerdr h = new FakeHerdr();
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
SessionManager sessions = new SessionManager(
workerService(h, "http://gx00.gw:8000", Set.of("gx00.gw")), worktrees);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller/proj", "term_primary",
new WorktreeRequest("cb-517", null));
McpSchema.CallToolResult res = BridgeMcp.whoami(Principal.worker(s.terminalId(), 200), sessions);
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"role\":\"worker\""), out);
assertTrue(out.contains("\"sessionId\":\"" + s.terminalId() + "\""), out);
assertTrue(out.contains("\"profile\":\"ltms-local\""), out);
assertTrue(out.contains("\"worktree\":\"" + s.worktree() + "\""), out);
assertTrue(out.contains("\"branch\":\"" + s.branch() + "\""), out);
assertTrue(out.contains("\"owner\":\"term_primary\""), out);
}
/**
* A worker the registry has no record of — it outlived a daemon restart — must still learn the
* load-bearing fact. Degrading to "I don't know who you are" would put it back to guessing,
* which is the failure this tool exists to remove.
*/
@Test
void whoamiStillReportsWorkerRoleWhenTheSessionIsUnregistered() {
FakeHerdr h = new FakeHerdr();
McpSchema.CallToolResult res = BridgeMcp.whoami(Principal.worker("term_orphan", 200),
sessionManager(h, "http://gx00.gw:8000", Set.of("gx00.gw")));
assertNotEquals(Boolean.TRUE, res.isError());
String out = textOf(res);
assertTrue(out.contains("\"role\":\"worker\""), out);
assertTrue(out.contains("\"sessionId\":\"term_orphan\""), out);
assertFalse(out.contains("profile"), out); // nothing invented for a session we don't track
}
}
@@ -401,4 +401,67 @@ class MessageServiceTest {
injector.onStatus(T, AgentStatus.IDLE); // deliver the task
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up
}
// --- CB-516: a released session must not leave a send hanging ------------------------------
/**
* The bug this fixes: tearing a worker down left its rendezvous waiter open, so a blocking send
* kept blocking and an async one kept reporting PENDING until the 30-minute async timeout —
* even though the worker provably no longer existed.
*/
@Test
void abandonFailsASendThatIsStillWaitingOnAReleasedSession() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
awaitWaiting();
assertTrue(messages.abandon(T, "session released"), "a live waiter is abandoned");
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
assertEquals(MessageService.Outcome.WORKER_FAILED, r.outcome(),
"an abandoned send fails rather than riding out its timeout");
assertEquals("session released", r.text(), "the caller is told why");
}
@Test
void abandonIsANoOpWhenNobodyIsWaiting() {
assertFalse(messages.abandon(T, "session released"),
"no open send ⇒ nothing to abandon");
}
@Test
void abandonDoesNotOverwriteAnAlreadyResolvedSend() throws Exception {
CompletableFuture<MessageService.Reply> send =
CompletableFuture.supplyAsync(() -> messages.send(T, "work", 30_000));
awaitWaiting();
assertTrue(rendezvous.resolve(T, "the real answer"));
assertFalse(messages.abandon(T, "session released"),
"a send already answered by the worker must not be clobbered");
MessageService.Reply r = send.get(5, TimeUnit.SECONDS);
assertEquals("the real answer", r.text());
}
/** The async path is the one that hung: poll must report FAILED, not PENDING forever. */
@Test
void anAbandonedAsyncTaskPollsAsFailedNotPending() throws Exception {
String ticket = messages.sendAsync(T, "long task");
awaitWaiting();
assertEquals(MessageService.Phase.PENDING, messages.poll(ticket).phase());
messages.abandon(T, "session released");
MessageService.TaskView view = null;
long deadline = System.currentTimeMillis() + 3000;
while (System.currentTimeMillis() < deadline) {
view = messages.poll(ticket);
if (view.phase() != MessageService.Phase.PENDING) break;
Thread.sleep(10);
}
assertNotNull(view);
assertEquals(MessageService.Phase.FAILED, view.phase(),
"a delegation whose worker is gone must not keep reporting PENDING");
assertTrue(view.detail() != null && view.detail().contains("released"),
"and the detail says why, rather than 'worker unknown'");
}
}
@@ -70,6 +70,28 @@ class RendezvousTest {
"no blocked send means no primary to surface the question to");
}
@Test
void resolveCompletionTwiceIsANoOpTheSecondTime() {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(W);
assertTrue(rendezvous.resolveCompletion(waiter, "first scrape"), "the first completion resolves");
assertFalse(rendezvous.resolveCompletion(waiter, "second scrape"),
"a second completion on an already-resolved waiter returns false");
assertEquals(Rendezvous.Kind.COMPLETION, waiter.getNow(null).kind());
assertEquals("first scrape", waiter.getNow(null).text(),
"the first resolution wins; the stored value is unchanged");
}
@Test
void resolveFailureTwiceIsANoOpTheSecondTime() {
CompletableFuture<Rendezvous.Resolution> waiter = rendezvous.open(W);
assertTrue(rendezvous.resolveFailure(waiter, "first reason"), "the first failure resolves");
assertFalse(rendezvous.resolveFailure(waiter, "second reason"),
"a second failure on an already-resolved waiter returns false");
assertEquals(Rendezvous.Kind.FAILED, waiter.getNow(null).kind());
assertEquals("first reason", waiter.getNow(null).text(),
"the first resolution wins; the stored value is unchanged");
}
@Test
void closeAskRemovesTheTurn() {
Rendezvous.AskTicket t = rendezvous.openAsk(W);
@@ -362,4 +362,46 @@ class SessionManagerTest {
assertTrue(sessions.roster().isEmpty(),
"no session is registered when spawn times out (roster empty)");
}
// --- CB-516: release must notify, so a blocked send can be failed --------------------------
@Test
void releaseNotifiesTheListenerWithTheReleasedTerminal() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(released::add);
WorkerSession s = sessions.acquire("ltms-local", null, "/caller", null);
sessions.release(s.paneId());
assertEquals(java.util.List.of(s.terminalId()), released,
"every teardown path funnels through release, so one hook must see the terminal");
}
@Test
void releasingAnUnknownPaneNotifiesNobody() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
java.util.List<String> released = new java.util.concurrent.CopyOnWriteArrayList<>();
sessions.onRelease(released::add);
sessions.release("w9:p404"); // idempotent teardown of something already gone
assertTrue(released.isEmpty(), "no session removed ⇒ no send was waiting on it");
}
@Test
void aThrowingReleaseListenerDoesNotBlockTheTeardown() {
FakeHerdr herdr = new FakeHerdr();
SessionManager sessions = sessionManager(herdr);
sessions.onRelease(_ -> {
throw new IllegalStateException("listener blew up");
});
WorkerSession s = sessions.acquire("ltms-local", null, "/caller", null);
assertDoesNotThrow(() -> sessions.release(s.paneId()),
"a listener failure must never prevent the teardown it is reacting to");
assertTrue(sessions.get(s.paneId()).isEmpty(), "and the session is still deregistered");
}
}
+60
View File
@@ -0,0 +1,60 @@
# LavinMQ — the AMQP broker behind bridged's durable ReplyInbox (CB-307 Stage 2).
#
# Why this file exists: the broker was previously run ad hoc and simply vanished from the host,
# which takes bridged down with it — AmqpReplyInbox.open throws on an unreachable broker and
# Bridged.java:187 does not guard it, so a missing broker is a hard startup failure, not a
# degraded mode. This pins the version, keeps the data, and brings itself back after a reboot.
#
# Usage:
# docker compose -f deploy/lavinmq/compose.yaml up -d
# docker compose -f deploy/lavinmq/compose.yaml ps
# docker compose -f deploy/lavinmq/compose.yaml logs -f
# docker compose -f deploy/lavinmq/compose.yaml down # keeps the volume
# docker compose -f deploy/lavinmq/compose.yaml down -v # DESTROYS held replies
#
# Management UI: http://127.0.0.1:15672 (guest / guest)
#
# This is bridged's OWN broker. Do not point bridged at any other AMQP server on this host —
# notably not the `local-rabbitmq` container, which belongs to a different project and would end
# up carrying this project's queues.
name: bridged-broker
services:
lavinmq:
# Pinned deliberately: :latest silently moves the broker under a running daemon.
image: cloudamqp/lavinmq:2.9.1
container_name: bridged-lavinmq
# The failure this deployment exists to prevent — survive reboots and Docker restarts, but
# stay down if it was stopped on purpose.
restart: unless-stopped
# Loopback-bound on purpose. LavinMQ ships a default guest/guest account, which is only
# acceptable because nothing off-host can reach it. bridged connects over 127.0.0.1, and
# binding 0.0.0.0 here would expose a broker with default credentials to the network.
ports:
- "127.0.0.1:5672:5672" # AMQP — bridged.yaml broker.uri points here
- "127.0.0.1:15672:15672" # HTTP management API + UI
# The whole point of Stage 2. Held-but-unacked replies live here; without a named volume a
# `docker compose down` would discard exactly what the durable inbox exists to protect.
volumes:
- lavinmq-data:/var/lib/lavinmq
healthcheck:
test: ["CMD", "lavinmqctl", "status"]
interval: 30s
timeout: 5s
retries: 3
start_period: 10s
logging:
driver: json-file
options:
max-size: "10m"
max-file: "3"
volumes:
lavinmq-data:
name: bridged-lavinmq-data
+1 -1
Submodule wiki updated: f4af2a1c22...0c896eb49b