Compare commits

...

13 Commits

Author SHA1 Message Date
Dai Ha 4fa6553db5 CB-527/CB-528: bound AMQP prefetch and confirm publishes before claiming durable
CI / build (pull_request) Successful in 1m3s
CI / contract (pull_request) Successful in 1m15s
CB-527: basicQos(prefetch) on the consume channel before basicConsume, configurable
via broker.prefetch (default 32), so an undrained inbox backlog stays on the broker
instead of growing the JVM heap without limit.

CB-528: publish moves to its own confirm-mode channel with mandatory=true and a
return listener, so an unroutable or unconfirmed reply now throws instead of
vanishing silently. The confirm callback checks the per-message returned flag
(set by the return listener, which the broker always fires before the matching
confirm) so an acked-but-returned publish is still reported as a failure. The ack
path stays on its own channel/lock and never waits on a publish confirm.
2026-08-16 16:55:40 +02:00
Dai Ha 2124e043ce CB-593: correct the member MCP claim — measured, not assumed
CI / contract (push) Successful in 43s
CI / build (push) Successful in 1m19s
CLAUDE.md told every member 'You mount only the bridge MCP' and told the lead
'a worker mounts only the bridge MCP and cannot run your other tooling'. Both
were false for Claude Code members.

Measured by spawning one member per backend and asking each what it actually has:

  opencode (gx)      11 bridge tools only                     claim TRUE
  claude-code (local) 11 bridge + 45 gitea + 2 context7       claim FALSE

Source is ~/.claude.json user-scope mcpServers; --mcp-config adds to that scope
rather than replacing it, so the worktree parity overlay (which correctly
neutralises .mcp.json and opencode.json) cannot see or stop it.

The forge tools are mounted but not usable: CB-592's blocked sentinel means
get_me and list_issues both fail with 'invalid username, password or token'.
That is defence in depth working in a path it was not designed for, so the text
now says a mounted tool is not a working tool rather than pretending the tools
are absent.

Block propagated byte-identically to wiki 7-Use-Cases.md (wiki 1d95e3f).
2026-08-16 09:46:42 +02:00
Dai Ha e4f3620acb CB-591: record the final 86400s route timeout and the request:0s trap
CI / contract (push) Successful in 44s
CI / build (push) Successful in 1m19s
systems/vms moved the LLM route timeout again, 1800s -> 86400s (24h), after
the silent-truncation risk was discussed. They tried request: 0s first: it
removes the total-duration timer, but on an AIGatewayRoute the idle timeout
is derived from the request timeout, so 0s also removed any bound on a
stalled connection.

At 86400s our own MessageService.ASYNC_TIMEOUT_MS (30 min) binds first, so a
runaway request now ends as a clean FAILED ticket we raised instead of a
silently truncated 200. While the gateway sat at 1800s the two numbers were
equal and did not nest.
2026-08-15 21:29:49 +02:00
Dai Ha 032a59a34d CB-591: fleet moved onto the gateway — both ceilings fixed and re-verified
CI / contract (push) Successful in 40s
CI / build (push) Successful in 1m19s
`local` runs on /anthropic and `gx` on /v1, both weight 100; `local-direct`
stays weight 0 as the escape hatch.

systems/vms fixed both blockers, and each was re-checked from this side rather
than taken on trust:

    listener buffer    32 KiB -> 32 Mi   ours: 1.2 MB body -> 200 (was 413)
    LLM route timeout  60s    -> 1800s   ours: 101s stream -> 200,
                                               message_stop present, 4000/4000

Neither was deliberate: 32 KiB was Envoy Gateway's default
per_connection_buffer_limit_bytes, and 60s was Envoy AI Gateway's own default.
The 60s bounded GENERATION as well as prompt size — a tiny prompt with a long
answer returned 504 at 60.05s.

Verified with real workloads, not liveness probes. A `local` member read this
document and CLAUDE.md in full — 48,344 bytes of file content, comfortably past
the old 32,768 ceiling — and answered four questions correctly, including the
document's length (said ~456, actual 455). A `gx` member did the same. The
trivial 3-question probe is what hid the 32 KiB ceiling for an afternoon, so it
no longer counts as proof here.

§7.2 is new and is the part that matters later. One risk is ACCEPTED, not
solved: on a mid-response timeout over chunked HTTP/1.1, Envoy ends the chunked
encoding cleanly instead of resetting, so a truncated answer arrives as HTTP 200
with no error and no terminator (envoyproxy/envoy#17186, acknowledged 2021,
never fixed; the Dec 2025 fix #42269 is HTTP/2 only and SSE here is HTTP/1.1).
Measured at the old 60s: 200, 61.07s, 2473 of 4000 emitted, message_stop 0,
error events 0, ending on a well-formed frame.

The recommended defence — reject a stream with no terminator — does NOT
transfer to us: Claude Code and opencode are third-party clients and we do not
own their SSE parsing. So this is acceptable because a request would have to run
1800s to trip it, not because we could detect it. If a member ever returns a
confident but truncated answer, suspect this before anything in our own code.

Also recorded, from the upstream bisection: ClientTrafficPolicy is honoured in
standalone `aigw run` but BackendTrafficPolicy is silently ignored, and nothing
external distinguishes them (envoyproxy/gateway#9513). Same silent-default shape
this repo keeps hitting.

bridged.yaml carries the same notes inline (gitignored, so not in this commit).

Refs: gitea #76
2026-08-15 20:47:41 +02:00
Dai Ha e689090024 CB-591: correct the root cause — Envoy's buffer limit, not Caddy
CI / contract (push) Successful in 1m7s
CI / build (push) Successful in 1m35s
I wrote "Caddy request_body max_size and/or Envoy's own" and marked it
unverified. The Caddy half was wrong, and an unverified guess still points the
next reader at the wrong component.

Confirmed by the systems/vms side: Envoy Gateway defaults a listener's
per_connection_buffer_limit_bytes to 32768, and aigw buffers the WHOLE request
body before it can route on the model name. So that default is not a network
tuning knob — it is a hard ceiling on prompt size. From the live config_dump:

    listener default/llm/http    per_connection_buffer_limit_bytes: 32768

Nobody chose 32 KiB; it was inherited.

Both TLS edges are innocent, and the technique that showed it is better than
mine: both 413s carry x-llm-consumer, a header their auth proxy sets only AFTER
authenticating, so the body cleared both edges and the auth. On llm.vm, aigw
413s at 39 KB while the vLLM backend answers 200 at the same size. I found the
boundary; they found the component, by reading the failure's response headers.

Consequences recorded in the doc:

  * DO NOT plan around 32 KiB. The intended ceiling is far higher, so sizing our
    profiles to it would be designing around a bug.
  * Their fix (ClientTrafficPolicy, bufferLimit: 8Mi) is written but NOT
    deployed, pending their operator's approval. We do not re-test until they
    confirm — a half-changed system gives a number neither side can trust.
  * In standalone `aigw run` a SecurityPolicy is accepted and then silently
    ignored, so "the config was accepted" proves nothing there. They will verify
    by re-reading the live config_dump and sending a large request. Same
    silent-default shape this repo keeps hitting, one layer down.

bridged.yaml carries the same correction (gitignored, so not in this commit).

Refs: gitea #76
2026-08-15 19:47:01 +02:00
Dai Ha 1cc34888fd CB-591: record the live result — blocked by a 32 KiB body limit at the gateway
CI / contract (push) Successful in 44s
CI / build (push) Successful in 55s
Deployed U1-U2c, restarted, spawned both new profiles for real, then reverted.

llm.ltms.dev answers HTTP 413 above 32 KiB (32768 bytes), on BOTH surfaces:

    /v1        32695 bytes -> 200        /anthropic  32095 bytes -> 200
    /v1        32795 bytes -> 413        /anthropic  32855 bytes -> 413

That is far below one agent turn. It is an edge limit (Caddy request_body
max_size, and/or Envoy), so the fix is in systems/vms, not here.

The part worth recording is how it nearly passed. Two members, same message,
same moment: `local` finished in 66s, `gx` never finished at all. `local`
passed only because the probe was three trivial questions in a fresh session,
so the request fit under 32 KiB — the profile looked healthy and was a
landmine set to fire on the first turn that reads a file. So §7's checklist
was not wrong, it was too easy; it now says to use a file-reading task.

opencode's failure mode is worse than a crash: it catches the 413, compacts
its context, retries, and loops. Observed 10+ minutes BUSY with no reply. From
the lead's side that is indistinguishable from a slow worker. Reproduced
outside the bridge with the launcher's own generated config, which is how it
became a one-line error instead of a hang; §7.1 records that procedure.

Everything else about the migration checked out and is recorded so it is not
re-tested: token accepted on both surfaces, unauthenticated 401 (the Caddy
proxy does gate, whatever the gateway's own fail-open policy does),
/v1/models exactly ["deepseek-v4-flash"], the guard allowlist accepted
llm.ltms.dev, and the generated opencode provider block is correct with a real
llmk- key.

Also answers §3b's open question: reasoning survives BOTH surfaces —
/anthropic returns a real "type":"thinking" block and /v1 returns a populated
reasoning_content. The feared /v1 translation loss did not happen.

Config state (bridged.yaml is gitignored, so it is described rather than
committed): `local` back on http://gx00.gw:8000, `gx` kept at weight 0,
`local-direct` kept, llm.ltms.dev left in the guard allowlist. The file
carries these numbers and the exact two-key edit to switch back.

Verified after the revert with a task that reads two large files: correct on
all three questions. Daemon pid 66745, jar f1fd659423e6.

Refs: gitea #76
2026-08-15 19:27:45 +02:00
Dai Ha 0331ecd5d3 CB-592: add the BRIDGED_MEMBER marker — the sentinel alone cannot hold
CI / contract (push) Successful in 1m5s
CI / build (push) Successful in 1m39s
Live check on a member pane showed the CB-592 shadow did NOT take effect:
GITEA_ACCESS_TOKEN inside the pane was still the real admin token.

Measured cause. The overlay itself works — GITEA_TOKEN is injected the same
way, is exported by no shell file, and does reach the pane. The sentinel loses
one step later. A herdr pane runs a LOGIN shell, ~/.zprofile line 41 sources
${SHARED_ENV}/tools/secrets.sh, and that file does a plain unconditional
`export GITEA_ACCESS_TOKEN=...`. A login shell overwrites a value already in
the environment, so the real token is put back before the member starts.
Confirmed directly:

    GITEA_ACCESS_TOKEN=cb592-sentinel zsh -lc ...
    -> RESULT: sentinel was OVERWRITTEN by the login shell

This defeats any launcher-side overlay for any name secrets.sh exports. No
change in this repo can win it alone.

So this adds the half that does survive: BRIDGED_MEMBER=1, a name secrets.sh
never exports. It is a no-op until the operator guards the export:

    [ -n "${BRIDGED_MEMBER:-}" ] || export GITEA_ACCESS_TOKEN=...

Setting it now costs nothing and makes that one line the whole remaining fix.
The sentinel stays: it is correct for any peer kind whose pane does not start
a login shell, and it keeps the intent explicit where every adapter passes.

Also corrects the javadoc and the test javadoc, which both claimed a
protection that was measured not to hold.

The other reported failure was my own bad test, not a regression. The probe
called /api/v1/user, which a minimal write:repository token cannot read. Same
token on the repo endpoint answers 200, so CB-302 is intact:

    GITEA_ACCESS_TOKEN: /user=200  /repos/lms/claude-bridge=200
    WORKER_GITEA_TOKEN: /user=403  /repos/lms/claude-bridge=200

Tests 805 -> 807. Both new tests proved to discriminate by reverting the
marker: everySpawnMarksThePaneAsAMember and
aProfileEnvEntryCannotClearTheMemberMarker both fail without it.

Refs: gitea #77
2026-08-15 18:37:52 +02:00
Dai Ha 831a918c30 Merge CB-592: shadow the admin GITEA_ACCESS_TOKEN in every member's environment
CI / contract (push) Successful in 41s
CI / build (push) Successful in 1m22s
A live probe showed every spawned member carried the admin GITEA_ACCESS_TOKEN: 108
environment variables in a member's pane against 99 in the primary's. The operator's
rule is that only the leader and architects may use it; everyone else uses
WORKER_GITEA_TOKEN. We were not enforcing that at all.

The cause is invisible from inside the launcher. baseEnv builds a fresh map holding only
PATH and the profile's env:, so a member looks like it gets a small explicit environment.
That map is an OVERLAY: WorkspaceControl.createTab/splitPane send only the keys it
contains, and herdr spawns the pane from its own login-shell environment, so every key we
never mention passes straight through — admin token included.

The fix puts a non-blank sentinel over the key in baseEnv, applied AFTER the profile's
env: so no profile, present or future, can restore the real token by naming it in config.
One place, every adapter, including peer kinds not yet written — deliberately not a
per-profile bridged.yaml entry, which is the silent-default shape this repo has shipped
nine times.

A non-blank sentinel rather than the empty string, on purpose: whether an empty overlay
value overrides an inherited variable or is skipped as blank cannot be settled from this
repo, because herdr's merge happens in an external process. baseEnv's own PATH seeding
(CB-511) already relies on a non-blank value replacing an inherited one, so this reuses
the shape that is demonstrated to work rather than the one that is merely plausible.

CB-302's repo-scoped GITEA_TOKEN grant is untouched — a worker can still open its own PR.
The subscription boundary was checked and is unaffected: the primary's pane carries no
ANTHROPIC_* at all, so nothing is inherited there.

Closes gitea #77. Live verification follows separately: the daemon must be redeployed
before this reaches any pane.
2026-08-15 18:27:05 +02:00
Dai Ha 3db5277ae8 CB-592: shadow the admin GITEA_ACCESS_TOKEN in every member's herdr overlay
CI / contract (pull_request) Successful in 1m12s
CI / build (pull_request) Successful in 1m37s
herdr spawns a pane from its own login-shell process env and layers our map on
top, so any key baseEnv never mentions passes straight through — including the
admin forge token. baseEnv now puts a non-blank sentinel for
GITEA_ACCESS_TOKEN, applied after the profile's own env: so no profile can
restore it. One place, every adapter, every profile including future ones.
CB-302's GITEA_TOKEN grant (applyGitToken) is untouched.
2026-08-15 18:25:16 +02:00
Dai Ha 6939e0cbbc CB-592: the tracked opencode.json must name the worker forge token, not the admin one
Operator's rule, 2026-08-15: only the leader and architects may use GITEA_ACCESS_TOKEN;
everyone else uses WORKER_GITEA_TOKEN.

opencode.json is TRACKED, so it ships in every worker worktree, and it mounted the gitea
MCP with {env:GITEA_ACCESS_TOKEN}. A live probe confirmed that variable actually resolves
inside a member: herdr spawns each pane from its own login-shell environment and layers
the launcher's map on top, so a member sees 108 variables rather than the small explicit
set baseEnv appears to build. That gave an opencode member admin forge TOOLS — enough to
merge its own PR, which both CLAUDE.md and the member contract forbid.

This is the narrow half of the fix: it removes the tooling. The admin token is still
present as a string in every member's environment, which is the real defect and is
tracked as CB-592 (gitea #77) — that fix belongs in the launcher, in one place, not
per-profile in bridged.yaml where a sixth profile would silently reopen it.

.mcp.json keeps GITEA_ACCESS_TOKEN and is correct to: it is skip-worktree, the primary's
own local copy, and the primary is the lead. That is the pattern this change follows —
the shared tracked file grants least privilege, and anything needing more overrides
locally.
2026-08-15 18:20:55 +02:00
Dai Ha f0095bf8b2 CB-591: plan the move onto the LLM/MCP gateway, and check AI_GATEWAY_TOKEN
CI / contract (push) Successful in 1m5s
CI / build (push) Successful in 1m38s
The gateway (llm.ltms.dev) replaced Bifrost on 2026-08-15 and serves an Anthropic
surface and an OpenAI surface, so both member kinds can point at it. The plan is in
docs/CB-591-Gateway-Migration.md; gitea #76 tracks the work.

The opencode half needs no code: OpenCodeLauncher already pins an OpenAI-compatible
endpoint (CB-508), so baseUrl + tokenEnv + provider/model is a config change. That
matters more than it looks — every opencode member today is sol or terra, and both sit
on one OpenAI account via credentialId: openai-shared, so an exhaustion on either locks
out both. A gateway-backed opencode profile is free and off that credential, which
retires a single point of failure rather than only adding capacity.

Also extends the redeploy script's --check to AI_GATEWAY_TOKEN. A profile's tokenEnv is
resolved from the DAEMON's own environment by HerdrPeerLauncher.resolveEnv, so a token
added to secrets.sh after the daemon started is simply absent: the launcher injects an
empty token and the gateway answers 401, long after the restart and with nothing tying
the two together. That is the same trap as WORKER_GITEA_TOKEN, and it gets the same
login-shell check that never prints the value.
2026-08-15 17:06:31 +02:00
Dai Ha 5206679efd Merge CB-588: nudge the lead when an async ticket goes terminal
An async delegation ticket (bridge_send wait:false) resolves on MessageService.reply's
rendezvous fast path, which returns before onReplyQueued. So CB-307's push loop only ever
heard about the durable-inbox case, and the mode CLAUDE.md tells leads to prefer never
nudged anyone. Closes gitea #72.

Adds a second, independent reminder schedule keyed by the lead terminal, so several
tickets finishing together coalesce into one nudge. The CB-307 path is untouched.

Three defects were found in review and fixed before merge:
 * a pendingTickets entry outlived the ticket it named. poll() returns null once
   pruneTerminalTickets drops a ticket, so ticketCollected was never reached and the
   entry leaked for the daemon's life, riding along on every later nudge and sending
   the lead after a ticket bridge_poll can no longer find.
 * a lost nudge: a ticket landing between decideTickets returning STOP and
   activeLeads.remove coalesced onto a schedule that was about to die. That is the
   exact failure this ticket exists to remove, reintroduced in a narrow window.
 * the success direction was unpinned in tests, and the comment listing the paths that
   complete the future was short by several.

The obvious fix for the second one was wrong: restarting on any pending ticket defeats
the reminder cap, because a never-collected ticket at cap is expected to still be there.
The fix diffs against a snapshot taken before the decision, so only a ticket that truly
arrived during the window restarts the schedule.

Verified on my own unpiped build: 802 tests, 0 failures, BUILD SUCCESS.
Two reviewers on the diff; the loop-gating finding they raised is split out as CB-590.
2026-08-15 17:06:12 +02:00
Dai Ha 6d0c94dbdb Correct enforceMaxLoad's comment after CB-585
CI / contract (push) Successful in 42s
CI / build (push) Successful in 1m16s
The comment said non-positive means unlimited at load. That stopped being
true when CB-585 made an explicit maxLoad: 0 survive as a real cap of zero
and made a negative value refuse config load. The code below it was already
right — only the comment described the old normalisation. Flagged by the
CB-585 worker, which correctly stayed out of a file not on its list.
2026-08-15 16:02:35 +02:00
13 changed files with 936 additions and 29 deletions
+10 -6
View File
@@ -81,9 +81,9 @@ below are the procedure — run them in order, every task, not only the big ones
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.
6. **Verify yourself.** Re-run the build and the checks. A worker cannot run your IDE tooling, any
forge tools it appears to have hold a blocked credential and fail, 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
@@ -155,9 +155,13 @@ you.
without replying, the bridge scrapes your pane, and it can return only the last 4000 characters.
A clipped scrape is marked as partial, but the missing text is gone — your report reaches the
lead with its end cut off.
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.
5. **Report honestly.** State only what you actually ran and its real output, including failures,
and never claim the result of a check you had no way to run. **Measure your own tools; do not
assume them.** What you mount depends on your backend: an opencode member gets the bridge and
nothing else, while a Claude Code member also inherits the operator's user-scope MCP servers,
which the bridge never chose for you. Two rules follow. The primary's IDE tooling is still not
yours, whatever you see. And **a mounted tool is not a working tool** — the forge server you may
find there holds a deliberately blocked credential and fails every call, by design.
6. **Never merge.** Stage files explicitly — never `git add -A` — and leave alone anything the
project marks as not-yours-to-commit.
+6 -2
View File
@@ -434,10 +434,14 @@ guard:
# on a durable per-target queue (agent.<target>.inbox) and survive a restart — the broker
# redelivers anything the primary had not yet drained. Production default is LavinMQ; a stock
# RabbitMQ speaks the same AMQP 0-9-1, so it is a URI-only swap.
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
# prefetch → CB-527: consumer basicQos, capping how many unacked messages the inbox holds
# in-heap per owned target (the rest sits on the broker's durable queue instead of
# growing the JVM heap). Default 32 when omitted.
# broker:
# uri: amqp://guest:guest@127.0.0.1:5672
# prefetch: 32
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open bridge_send,
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
@@ -348,8 +348,9 @@ public final class Bridged {
// connection, so keep the reference to close it in the ordered shutdown hook.
final ReplyInbox replyInbox;
if (cfg.broker() != null && cfg.broker().isConfigured()) {
replyInbox = AmqpReplyInbox.open(cfg.broker().uri());
log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri());
replyInbox = AmqpReplyInbox.open(cfg.broker().uri(), cfg.broker().prefetchOrDefault());
log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})",
cfg.broker().uri(), cfg.broker().prefetchOrDefault());
} else {
replyInbox = new InMemoryReplyInbox();
log.info("reply inbox: in-memory (soft-state)");
@@ -5,6 +5,7 @@ import com.fasterxml.jackson.core.JsonParser;
import com.fasterxml.jackson.core.JsonToken;
import com.fasterxml.jackson.databind.ObjectMapper;
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
import dev.ltms.bridged.msg.AmqpReplyInbox;
import dev.ltms.bridged.peer.MemberRole;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -532,16 +533,24 @@ public record BridgedConfig(
* stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1
* and is a URI-only swap.
*
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/{@code null}
* ⇒ the broker block is treated as absent (in-memory adapter).
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/
* {@code null} ⇒ the broker block is treated as absent (in-memory adapter).
* @param prefetch CB-527: the consumer's {@code basicQos} prefetch count, bounding how many
* unacked messages the AMQP inbox holds in-heap per owned target. {@code null}/
* non-positive ⇒ {@link AmqpReplyInbox#DEFAULT_PREFETCH}.
*/
@JsonIgnoreProperties(ignoreUnknown = true)
public record Broker(String uri) {
public record Broker(String uri, Integer prefetch) {
/** True when a usable broker URI is configured (an empty block does not enable AMQP). */
public boolean isConfigured() {
return uri != null && !uri.isBlank();
}
/** The prefetch to use, defaulting to {@link AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
public int prefetchOrDefault() {
return (prefetch != null && prefetch > 0) ? prefetch : AmqpReplyInbox.DEFAULT_PREFETCH;
}
}
/**
@@ -372,8 +372,11 @@ public final class CompositePeerLauncher implements PeerLauncher {
}
private void enforceMaxLoad(String profile) {
// Absent config, or a config whose maxLoad normalized to null (non-positive ⇒ unlimited at
// load), means no cap — never cap what wasn't configured.
// Absent config, or a config whose maxLoad normalized to null (ABSENT ⇒ unlimited at load),
// means no cap — never cap what wasn't configured. Note "non-positive ⇒ unlimited" was true
// until CB-585: an explicit `maxLoad: 0` now survives as 0 and is a real cap of zero, so the
// check below refuses every spawn on that profile, and a negative value is refused at config
// load rather than normalized away.
BridgedConfig.Profile cfg = profiles0().get(profile);
Integer cap = (cfg == null) ? null : cfg.maxLoad();
if (cap == null) {
@@ -766,10 +766,61 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
}
}
/** A fresh mutable env map — the conventional starting point for {@link #buildLaunch}. */
/**
* CB-592: overlay value that shadows the admin {@code GITEA_ACCESS_TOKEN} a herdr pane
* otherwise inherits from herdr's own login-shell process environment (gitea issue #77).
* herdr spawns a pane from its <em>own</em> process environment and layers our map on top —
* {@link dev.ltms.bridged.herdr.WorkspaceControl#createTab} and {@code #splitPane} send only
* the keys we put in that map, so any key we never mention passes straight through from
* herdr's own shell, admin token included.
*
* <p>Deliberately a non-blank sentinel, not {@code ""}. Whether an empty-string overlay value
* overrides an inherited variable or is skipped as blank could not be settled by reading this
* codebase — herdr's server-side merge is an external process, not something in this repo.
* A non-blank replacement sidesteps that ambiguity entirely: {@link #baseEnv}'s own {@code
* PATH} seeding already depends on the overlay reliably replacing an inherited value (see its
* javadoc), and that is only demonstrated for a non-blank value, so this reuses the same,
* proven-reliable shape rather than the unverified one.
*
* <p><b>MEASURED ON A LIVE PANE, 2026-08-15: this sentinel alone does NOT hold.</b> The overlay
* itself works — {@code GITEA_TOKEN} is injected here, is exported by no shell file, and does
* reach the pane. The sentinel loses one step later. A herdr pane runs a <em>login</em> shell,
* {@code ~/.zprofile} sources {@code ${SHARED_ENV}/tools/secrets.sh}, and that file does a plain
* unconditional {@code export GITEA_ACCESS_TOKEN=...}. A login shell overwrites a value already
* in the environment, so the real admin token is put back over this sentinel before the member
* process ever starts. That defeat applies to <em>every</em> name {@code secrets.sh} exports,
* and no launcher-side overlay can win against it.
*
* <p>So this constant is not the control on its own — {@link #MEMBER_MARKER} is the other half.
* Keeping the sentinel is still worth it: it is correct for any peer kind whose pane does not
* start a login shell, and it makes the intent explicit at the one place every adapter passes.
*/
private static final String BLOCKED_GITEA_ACCESS_TOKEN =
"blocked-by-bridged-cb592-see-gitea-issue-77";
/**
* CB-592: marks a pane as a bridged member so a shell startup file can decline to export
* operator-only credentials into it (gitea issue #77).
*
* <p>This name is deliberately one that {@code secrets.sh} never exports, which is exactly why
* it survives the login shell that wipes {@link #BLOCKED_GITEA_ACCESS_TOKEN}. The mechanism is
* measured, not assumed: {@code GITEA_TOKEN} is injected the same way, is absent from a login
* shell of its own, and was observed set inside a live member pane.
*
* <p>It is a no-op until the operator guards the export, which is a one-line change in a file
* this repo does not own and must not edit unasked:
*
* <pre>{@code
* [ -n "${BRIDGED_MEMBER:-}" ] || export GITEA_ACCESS_TOKEN=...
* }</pre>
*
* <p>Setting the marker now costs nothing and means that edit is the whole remaining fix.
*/
static final String MEMBER_MARKER = "BRIDGED_MEMBER";
/**
* Seed a worker's environment (CB-511): the daemon's own {@code PATH}, then the profile's
* {@code env:} entries.
* {@code env:} entries, then the CB-592 admin-token shadow.
*
* <p>Why this exists: bridged passes herdr an explicit env map, and herdr merges it into
* <em>its own</em> process environment. So before this, a worker inherited whatever PATH the
@@ -783,6 +834,12 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
* overriding {@code ANTHROPIC_BASE_URL} and slipping past {@link
* dev.ltms.bridged.guard.SubscriptionGuard}, which is checked against the profile's
* {@code baseUrl} and nothing else.
*
* <p>The CB-592 shadow and marker are put in <em>last</em>, after the profile's own
* {@code env:}, so no profile — present or future — can restore the admin token, or hide that
* the pane is a member, by naming either in config. This is the one place both are applied:
* every {@code buildLaunch} in every adapter calls this first, so a new profile, and a peer
* kind not yet written, gets them for free.
*/
protected Map<String, String> baseEnv(BridgedConfig.Profile cfg) {
Map<String, String> workerEnv = new LinkedHashMap<>();
@@ -793,6 +850,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
if (cfg != null && cfg.env() != null) {
workerEnv.putAll(cfg.env());
}
workerEnv.put("GITEA_ACCESS_TOKEN", BLOCKED_GITEA_ACCESS_TOKEN);
workerEnv.put(MEMBER_MARKER, "1");
return workerEnv;
}
@@ -7,6 +7,7 @@ import com.rabbitmq.client.ConnectionFactory;
import com.rabbitmq.client.DeliverCallback;
import com.rabbitmq.client.Recoverable;
import com.rabbitmq.client.RecoveryListener;
import com.rabbitmq.client.Return;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
@@ -14,7 +15,13 @@ import java.io.IOException;
import java.nio.charset.StandardCharsets;
import java.util.LinkedHashMap;
import java.util.List;
import java.util.NavigableMap;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ConcurrentSkipListMap;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.TimeoutException;
/**
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
@@ -29,6 +36,27 @@ import java.util.concurrent.ConcurrentHashMap;
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
* adapter cannot give, with the port contract preserved.
*
* <p><strong>Prefetch bounds the held backlog (CB-527).</strong> The consumer channel calls
* {@code basicQos} with a configurable prefetch count ({@link #DEFAULT_PREFETCH} unless the caller
* passes another value to {@link #open(String, int)}) before starting any consumer. Without a bound,
* the broker pushes its entire queue into {@link #held} the instant a target is {@link #own owned},
* so an undrained primary grows the JVM heap without limit and any queue-level control
* ({@code x-max-length}, per-message TTL) never fires because the queue never actually holds a
* backlog. Prefetch keeps the backlog where it is visible — on the broker — until the owner drains it.
*
* <p><strong>Publishes require a confirmed, routable delivery (CB-528).</strong> {@link #publish}
* runs on a channel separate from the consume/ack channel ({@link #channel}), so a slow or blocked
* publish confirm can never hold {@link #channelLock} and stall an ack — the ack path never waits on
* a publish confirm. That publish channel is in publisher-confirm mode and every publish sets the
* {@code mandatory} flag, so an unroutable publish (queue not declared, e.g. the owner never called
* {@link #own}) is returned by the broker instead of silently dropped. The broker sends the
* <em>return</em> for an unroutable message before the <em>confirm</em> that covers it — the ack/nack
* callback checks the returned-set at confirm time rather than assuming an ack means routed — so
* "confirmed" here means "durably queued", not merely "accepted by the broker". A returned or nacked
* (or un-confirmed within the timeout) publish surfaces as an {@link IllegalStateException} on the
* caller's thread; the caller — {@link MessageService#reply} — must not report success for a
* black-holed reply.
*
* <p><strong>Ownership is explicit.</strong> {@link #own} declares the queue and starts the consumer;
* {@link #release} cancels it. {@link #publish} sends to the queue but does <em>not</em> imply ownership
* and does not attach a consumer. This split is required by CB-308 federation, where one gateway may
@@ -54,6 +82,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
private static final String QUEUE_PREFIX = "agent.";
private static final String QUEUE_SUFFIX = ".inbox";
/** CB-527: the prefetch used when a caller does not pass an explicit value to {@link #open(String, int)}. */
public static final int DEFAULT_PREFETCH = 32;
/** How long {@link #publish} waits for its publisher confirm before failing the call (CB-528). */
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
private final Connection connection;
private final Channel channel;
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
@@ -63,39 +97,81 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
/**
* CB-528: a dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume
* + ack) so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
*/
private final Channel publishChannel;
private final Object publishChannelLock = new Object();
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
/** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */
private record Held(long deliveryTag, InboxMessage message) {}
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
private static final class Pending {
final String msgId;
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
volatile boolean returned;
Pending(String msgId) {
this.msgId = msgId;
}
}
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) with {@link #DEFAULT_PREFETCH}. */
public static AmqpReplyInbox open(String uri) {
return open(uri, DEFAULT_PREFETCH);
}
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
public static AmqpReplyInbox open(String uri, int prefetch) {
try {
ConnectionFactory factory = new ConnectionFactory();
factory.setUri(uri);
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
factory.setAutomaticRecoveryEnabled(true);
factory.setTopologyRecoveryEnabled(true);
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"));
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"), prefetch);
} catch (Exception e) {
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
}
}
/** Wrap an already-open connection (injection seam for the contract test). */
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
AmqpReplyInbox(Connection connection) {
this(connection, DEFAULT_PREFETCH);
}
/** As above, with an explicit prefetch (injection seam for the contract test). */
AmqpReplyInbox(Connection connection, int prefetch) {
this.connection = connection;
try {
this.channel = connection.createChannel();
// CB-527: bound the held backlog per owned target — must be set before any own()/basicConsume.
this.channel.basicQos(prefetch);
this.publishChannel = connection.createChannel();
this.publishChannel.confirmSelect();
this.publishChannel.addReturnListener(this::onReturn);
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
} catch (IOException e) {
throw new IllegalStateException("cannot open AMQP channel", e);
}
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
// repopulates it with valid tags (dedup by msgId still prevents any double-queue).
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
// confirm still in flight when the connection dropped is equally stale — its sequence number
// meant nothing on the old channel and means nothing on the recovered one, so fail it now
// rather than let it silently ride out CONFIRM_TIMEOUT_MS.
if (connection instanceof Recoverable recoverable) {
recoverable.addRecoveryListener(new RecoveryListener() {
@Override
public void handleRecovery(Recoverable recoverable) {
held.clear();
failPendingPublishesOnRecovery();
log.info("AMQP connection recovered; cleared held replies for fresh redelivery");
}
@@ -141,6 +217,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
}
}
/**
* Publish {@code content} and block until the broker's publisher confirm for it lands (CB-528).
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
* confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish,
* not a lost-and-forgotten one.
*/
@Override
public void publish(String target, String msgId, String content) {
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
@@ -148,12 +230,35 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
.deliveryMode(2) // persistent — survives a broker restart
.contentType("text/plain")
.build();
try {
synchronized (channelLock) {
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
Pending pending = new Pending(msgId);
long seq;
synchronized (publishChannelLock) {
seq = publishChannel.getNextPublishSeqNo();
pendingBySeq.put(seq, pending);
pendingByMsgId.put(msgId, pending);
try {
publishChannel.basicPublish("", queueName(target), true, props,
content.getBytes(StandardCharsets.UTF_8));
} catch (IOException e) {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msgId, pending);
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
}
} catch (IOException e) {
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
}
try {
pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS);
} catch (ExecutionException e) {
Throwable cause = e.getCause();
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
} catch (TimeoutException e) {
throw new IllegalStateException("publish confirm for reply " + msgId + " to " + queueName(target)
+ " timed out after " + CONFIRM_TIMEOUT_MS + "ms — broker may be unreachable or overloaded", e);
} catch (InterruptedException e) {
Thread.currentThread().interrupt();
throw new IllegalStateException("interrupted awaiting publish confirm for " + msgId, e);
} finally {
pendingBySeq.remove(seq, pending);
pendingByMsgId.remove(msgId, pending);
}
}
@@ -222,6 +327,65 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
};
}
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm (CB-528). */
private void onReturn(Return r) {
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
if (pending != null) {
pending.returned = true;
} else {
log.warn("AMQP return for reply {} (routingKey={}, {} {}) with no matching in-flight publish"
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
r.getReplyText());
}
}
private void onAck(long seq, boolean multiple) {
resolveConfirm(seq, multiple, true);
}
private void onNack(long seq, boolean multiple) {
resolveConfirm(seq, multiple, false);
}
/**
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
* a return means "confirmed but never routed", not "durably queued".
*/
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
NavigableMap<Long, Pending> covered = multiple
? pendingBySeq.headMap(seq, true)
: pendingBySeq.subMap(seq, true, seq, true);
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
if (ack && !pending.returned) {
pending.confirmed.complete(null);
} else if (ack) {
pending.confirmed.completeExceptionally(new IllegalStateException(
"reply " + pending.msgId + " was returned as unroutable (queue not declared/owned)"));
} else {
pending.confirmed.completeExceptionally(new IllegalStateException(
"broker nacked publish of reply " + pending.msgId));
}
}
}
/** Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery. */
private void failPendingPublishesOnRecovery() {
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
Pending pending = it.next().getValue();
it.remove();
pendingByMsgId.remove(pending.msgId, pending);
pending.confirmed.completeExceptionally(new IllegalStateException(
"AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId
+ " is unknown"));
}
}
private static String queueName(String target) {
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
}
@@ -233,6 +397,11 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
} catch (Exception e) {
log.debug("AMQP channel close: {}", e.toString());
}
try {
publishChannel.close();
} catch (Exception e) {
log.debug("AMQP publish channel close: {}", e.toString());
}
try {
connection.close();
} catch (Exception e) {
@@ -676,6 +676,85 @@ class ClaudeCodeLauncherTest {
"the guard-checked baseUrl must win over any env: entry, or the boundary is bypassable");
}
// --- CB-592: the admin GITEA_ACCESS_TOKEN never reaches a member -----------------------------
/**
* herdr's env map is an overlay onto its own (login-shell) process environment, so a worker
* inherits whatever the daemon's shell carries — including the admin GITEA_ACCESS_TOKEN — for
* every key baseEnv does not explicitly shadow. This pins that the launcher DOES send an
* explicit (non-blank) GITEA_ACCESS_TOKEN to herdr on every spawn, whatever the profile is, so
* a future baseEnv refactor cannot silently drop it and reopen the leak. Asserted against what
* tab.create's params actually carry, not an internal map built in the test (gitea #77).
*
* <p>Scope, measured on a live pane 2026-08-15: this pins what the launcher SENDS, and that is
* all it can pin. It does not prove the value survives, and it does not: the pane runs a login
* shell, ~/.zprofile sources secrets.sh, and its unconditional `export GITEA_ACCESS_TOKEN=...`
* puts the real token back over this sentinel. Closing that needs the operator to guard the
* export on BRIDGED_MEMBER — see everySpawnMarksThePaneAsAMember below.
*/
@Test
void everySpawnShadowsTheAdminGiteaAccessToken() {
FakeHerdr herdr = new FakeHerdr();
service(herdr, List.of("claude"), null).spawn();
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
assertNotNull(shadowed, "GITEA_ACCESS_TOKEN must be explicitly overlaid, not left unmentioned");
assertFalse(shadowed.isBlank(), "a blank overlay value's override behaviour is unverified — must be non-blank");
}
/**
* No profile — present or future — may restore the admin token by naming it in {@code env:}.
* The shadow is applied after the profile's own env in {@link HerdrPeerLauncher#baseEnv}
* precisely so this can never happen; this test pins that ordering.
*/
@Test
void aProfileEnvEntryCannotRestoreTheAdminGiteaAccessToken() {
FakeHerdr herdr = new FakeHerdr();
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null, null, null,
null, Map.of("GITEA_ACCESS_TOKEN", "admin-secret-from-profile-config"), null, null);
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> null).spawn();
assertNotEquals("admin-secret-from-profile-config", startEnv(herdr).get("GITEA_ACCESS_TOKEN"),
"a profile's own env: must not be able to smuggle the admin token back in");
}
/**
* The half of CB-592 that can actually survive the pane's login shell. BRIDGED_MEMBER is a name
* secrets.sh never exports, so nothing overwrites it — measured: GITEA_TOKEN is injected the
* same way, is absent from a login shell of its own, and was observed set inside a live member
* pane. It lets the operator guard the admin export with
* `[ -n "${BRIDGED_MEMBER:-}" ] || export GITEA_ACCESS_TOKEN=...`, which is the whole fix.
* Pinned here so a refactor cannot drop the marker and quietly un-guard every member (#77).
*/
@Test
void everySpawnMarksThePaneAsAMember() {
FakeHerdr herdr = new FakeHerdr();
service(herdr, List.of("claude"), null).spawn();
assertEquals("1", startEnv(herdr).get("BRIDGED_MEMBER"),
"every member pane must be marked, or a shell file cannot tell it apart from the operator's");
}
/** A profile must not be able to hide that its pane is a member, for the same reason as above. */
@Test
void aProfileEnvEntryCannotClearTheMemberMarker() {
FakeHerdr herdr = new FakeHerdr();
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
List.of("claude"), "tab", "bridged-workers", "w #{n}", null, null, null, null, null,
null, Map.of("BRIDGED_MEMBER", ""), null, null);
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
_ -> null).spawn();
assertEquals("1", startEnv(herdr).get("BRIDGED_MEMBER"),
"a profile's own env: must not be able to unmark its pane");
}
// ── CB-533: the model is pinned on the command line, not only in the environment ────────────
/** A launcher for a profile identical but for its {@code model:} — the only variable here. */
@@ -208,6 +208,21 @@ class OpenCodeLauncherTest {
"a git-token profile gets the peer-neutral GITEA_TOKEN grant, same as Claude");
}
/**
* CB-592: the shadow lives in {@link HerdrPeerLauncher#baseEnv}, shared by every adapter — this
* pins that the opencode path gets it too, not just Claude's. See the matching test in
* {@code ClaudeCodeLauncherTest} for the full rationale (gitea #77).
*/
@Test
void everySpawnShadowsTheAdminGiteaAccessToken(@TempDir Path root) {
FakeHerdr herdr = new FakeHerdr();
service(herdr, root, opencodeCfg(null, null, null)).spawn();
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
assertNotNull(shadowed, "GITEA_ACCESS_TOKEN must be explicitly overlaid, not left unmentioned");
assertFalse(shadowed.isBlank(), "a blank overlay value's override behaviour is unverified — must be non-blank");
}
@Test
void capabilitiesDeclareOrphanReapAndMcpAskAndConditionalSelfPr(@TempDir Path root) {
FakeHerdr herdr = new FakeHerdr();
@@ -16,6 +16,7 @@ import java.util.List;
import java.util.concurrent.TimeUnit;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
/**
@@ -166,6 +167,88 @@ class AmqpReplyInboxContractTest {
}
}
@Test
void prefetchBoundsTheHeldBacklog() throws Exception {
String target = "worker-prefetch-" + System.nanoTime();
int prefetch = 4;
int published = 10;
try (AmqpReplyInbox inbox = new AmqpReplyInbox(newConnection(), prefetch);
Connection inspect = newConnection()) {
inbox.own(target);
for (int i = 0; i < published; i++) {
inbox.publish(target, "m" + i, "payload " + i);
}
awaitHeldAtLeast(inbox, target, prefetch);
long depth;
try (Channel ch = inspect.createChannel()) {
depth = ch.queueDeclarePassive(queueName(target)).getMessageCount();
}
assertTrue(depth >= published - prefetch,
"broker should still hold at least " + (published - prefetch)
+ " undelivered messages behind a prefetch of " + prefetch + ", saw " + depth);
drainUntilEmpty(inbox, target, inspect);
}
}
@Test
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
String target = "worker-unroutable-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
// Deliberately never own(target): the queue is never declared, so the default-exchange
// route to agent.<target>.inbox does not exist and the broker must return the publish.
IllegalStateException ex = assertThrows(IllegalStateException.class,
() -> inbox.publish(target, "m1", "nobody home"));
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
"expected an unroutable-publish failure, got: " + ex.getMessage());
}
}
@Test
void confirmedPublishDeliversNormally() throws Exception {
String target = "worker-confirm-" + System.nanoTime();
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
inbox.own(target);
inbox.publish(target, "m1", "confirmed delivery"); // must return normally: routed and confirmed
List<ReplyInbox.InboxMessage> got = awaitPeek(inbox, target);
assertEquals(1, got.size());
assertEquals("confirmed delivery", got.getFirst().content());
inbox.ack(target, "m1");
}
}
/** Poll peek until at least {@code n} replies for {@code target} are held, or ~10s elapse. */
@SuppressWarnings("BusyWait")
private static void awaitHeldAtLeast(AmqpReplyInbox inbox, String target, int n) throws InterruptedException {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
while (inbox.peek(target).size() < n && System.nanoTime() < deadline) {
Thread.sleep(50);
}
}
/**
* Repeatedly ack whatever is currently held (freeing prefetch slots for the next delivery) until
* both the local snapshot and the broker's own queue depth are empty, or ~10s elapse.
*/
@SuppressWarnings("BusyWait")
private static void drainUntilEmpty(AmqpReplyInbox inbox, String target, Connection inspect) throws Exception {
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
while (System.nanoTime() < deadline) {
for (ReplyInbox.InboxMessage msg : inbox.peek(target)) {
inbox.ack(target, msg.msgId());
}
try (Channel ch = inspect.createChannel()) {
if (ch.queueDeclarePassive(queueName(target)).getMessageCount() == 0 && inbox.peek(target).isEmpty()) {
return;
}
}
Thread.sleep(50);
}
throw new AssertionError("did not drain " + target + " to empty within the deadline");
}
/** Poll peek (broker delivery is async) until a reply for {@code target} appears or ~10s elapse. */
@SuppressWarnings("BusyWait") // deliberate poll for async broker delivery, bounded by the deadline
private static List<ReplyInbox.InboxMessage> awaitPeek(AmqpReplyInbox inbox, String target)
+467
View File
@@ -0,0 +1,467 @@
# CB-591 — move the fleet onto the LLM and MCP gateway
**Status: DONE — the fleet is on the gateway as of 2026-08-15.** `local` runs on `/anthropic` and
`gx` on `/v1`, both at `weight: 100`; `local-direct` stays at `weight: 0` as the escape hatch. Getting
here took a revert and two upstream fixes — see §7.1, which is the useful part of this document. One
risk is **accepted rather than solved**: a stream cut by any mid-response timer arrives as HTTP 200
with no terminator, and our third-party members cannot detect it (§7.2).
· **Upstream:** [systems/vms wiki → LLM and MCP Gateway](https://git.ltms.dev/systems/vms/wiki/LLM-and-MCP-Gateway)
· **Upstream issue:** [systems/vms#31](https://git.ltms.dev/systems/vms/issues/31)
The gateway went live on 2026-08-15 and replaced Bifrost. This plan says what that means for a
**member definition** in `bridged.yaml`, because that is the part of this repo the change actually
touches.
---
## 1. What changed upstream
One front door for every LLM and MCP client: `https://llm.ltms.dev`, one token per consumer.
| Surface | URL |
|---|---|
| OpenAI chat | `https://llm.ltms.dev/v1/chat/completions` |
| OpenAI models | `https://llm.ltms.dev/v1/models` |
| **Anthropic messages** | `https://llm.ltms.dev/anthropic/v1/messages` |
| MCP, all servers multiplexed | `https://llm.ltms.dev/mcp` |
Anything outside that list returns **404 before any token is checked**, on purpose — the gateway must
never become a blanket proxy.
The model backend is unchanged: GX10 vLLM at `10.10.10.26:8000` (`gx00.gw`), model name exactly
`deepseek-v4-flash`. The direct LAN path stays open on purpose as an escape hatch.
---
## 2. Where claude-bridge sits today
We do **not** use the gateway. The `local` profile talks straight to the vLLM:
```yaml
local:
kind: claude-code
baseUrl: http://gx00.gw:8000 # direct vLLM — no auth, LAN only
model: deepseek-v4-flash
configDir: /Users/dai.ha/.ccs/instances/gx10
```
Three facts about our side that decide the shape of this work:
1. **`baseUrl` becomes `ANTHROPIC_BASE_URL`** in the member's environment, and `tokenEnv` becomes
`ANTHROPIC_AUTH_TOKEN` (the value is read from a host env var and never stored in config).
`local` sets no `tokenEnv` today, because a direct vLLM needs no token.
2. **`SubscriptionGuard` refuses any host not on an allowlist**, and that allowlist is
`guard.offSubscriptionHosts: [gx00.gw]`. It is built once in `Bridged.java:93` and handed to the
launcher, so **it is a restart-required key**, not a hot one. Changing `baseUrl` without changing
this makes every `local` spawn throw.
3. **The wiki names us as a blocker.** Under *Not done yet*: retiring the shared `legacy` token is
blocked because "kb, brain, **claude-bridge** and the workstation still share it. Each needs its
own consumer first."
Context7 is mounted twice today, both times straight at `https://ct7.ltms.dev/mcp` — once in
`.mcp.json` (the primary) and once in `opencode.json` (the `sol` and `terra` members).
```mermaid
flowchart LR
subgraph now["Today"]
M1["local member<br/>claude-code"] -->|"ANTHROPIC_BASE_URL"| V1["vLLM gx00.gw:8000<br/>no auth, LAN only"]
M2["sol / terra<br/>opencode"] --> CT1["ct7.ltms.dev/mcp"]
P1["primary"] --> CT1
end
subgraph after["Proposed"]
M3["local member"] -->|"ANTHROPIC_BASE_URL<br/>+ ANTHROPIC_AUTH_TOKEN"| G["llm.ltms.dev/anthropic<br/>consumer: claude-bridge"]
G --> V2["vLLM gx00.gw:8000"]
M4["local-direct<br/>weight 0, escape hatch"] --> V2
end
```
*The member definition is the only thing that moves. The model behind it does not.*
---
## 3. The member definition change
The gateway serves an Anthropic surface *and* an OpenAI surface, so **both member kinds can point at
it**. That is the main opportunity here, and it is bigger than the `local` profile alone.
### 3a. `local` — claude-code, on `/anthropic`
| Key | Today | After | Note |
|---|---|---|---|
| `baseUrl` | `http://gx00.gw:8000` | `https://llm.ltms.dev/anthropic` | see the schema warning below |
| `tokenEnv` | *(unset)* | `AI_GATEWAY_TOKEN` | new consumer token, `llmk-claude-bridge-<32 hex>` |
| `model` | `deepseek-v4-flash` | unchanged | must stay **exact**; a regex match returns an empty `/v1/models` while completions keep working |
| `guard.offSubscriptionHosts` | `[gx00.gw]` | `[gx00.gw, llm.ltms.dev]` | **restart required** |
### 3b. A new opencode profile on `/v1` — no code needed
`OpenCodeLauncher` already supports a pinned OpenAI-compatible endpoint (CB-508). Given `baseUrl` it
writes a custom provider block into the worker's opencode config:
- `baseUrl` → `options.baseURL`. `openAiBaseUrl` appends `/v1` to a bare host, and takes a URL that
already has a path **as-is** — so `https://llm.ltms.dev/v1` works unchanged.
- `tokenEnv` → `options.apiKey` (falls back to a placeholder when unset, since a local vLLM ignores it).
- `model:` **must** be `<provider>/<model>` when `baseUrl` is set — a bare name is rejected loudly
rather than silently falling back to opencode's default gateway.
So the profile is pure config:
```yaml
gx:
kind: opencode
baseUrl: https://llm.ltms.dev/v1
tokenEnv: AI_GATEWAY_TOKEN
model: gx/deepseek-v4-flash # provider id is ours to choose; the half after / is the model
argv: ["opencode"]
mcpUrl: http://127.0.0.1:8765/mcp
gitTokenEnv: WORKER_GITEA_TOKEN
weight: 100 # same tier as `local` — free
maxLoad: 2
# deliberately NO credentialId — this is our own box, not the shared OpenAI account
```
**Why this matters more than it looks.** Today every opencode member is `sol` or `terra`, and those
are two models on **one** OpenAI account sharing `credentialId: openai-shared` — so an exhaustion on
either locks out both, and half the fleet's opencode capacity dies at once. A gateway-backed opencode
profile is free, is not on that credential, and therefore is not in that quarantine pair. It removes
a single point of failure rather than just adding capacity.
**Note the asymmetry, it is deliberate:** `SubscriptionGuard` does not apply to opencode at all — the
guard exists to stop a *Claude* worker borrowing the operator's subscription, and opencode reads its
own provider credentials. So 3b needs **no allowlist change**; only 3a does.
**Both still need a restart, for a different reason.** `tokenEnv` is resolved by
`HerdrPeerLauncher.resolveEnv` → `env.apply(name)`, which reads the **daemon's own process
environment**. The running `bridged` inherited its environment when it started, so a variable added to
`secrets.sh` afterwards is simply not there — the launcher would inject an empty token and the
gateway would answer 401. This is the same failure as trap 1 in `scripts/redeploy-bridged.sh`
(`WORKER_GITEA_TOKEN`), and it has the same fix: **restart from a login shell**, and use
`scripts/redeploy-bridged.sh --check` to confirm the name resolves before restarting anything.
### 3c. What this does to `ccs`
Once a profile carries `baseUrl`, `tokenEnv` and `model` itself, the ccs instance stops being what
routes a member. Be precise about what is left, though: `configDir` still supplies **folder trust**
and `settings.json`, and dropping it is what produced the trust dialog and the wrong-model error
recorded in `bridged.yaml`. So ccs goes from *deciding where the tokens go* to *holding client-side
state*. Less load-bearing, not removable.
### Why `/anthropic` and never `/v1/chat/completions`
The gateway declares its Anthropic backend as `schema.name: Anthropic`, which means **no
translation** — streaming, tool use and thinking blocks pass through exactly as they do against vLLM
directly.
Declared as `OpenAI`, Envoy's translator looks for a `thinking_blocks` field that our vLLM does not
send (it sends `reasoning_content`), and **every thinking delta disappears silently**. Claude Code
speaks the Anthropic protocol, so `/anthropic` is both correct and the only safe choice.
This is the exact failure shape this repo keeps hitting: it compiles, it answers, it looks healthy,
and a capability is quietly off. Treat it as a `silent-default` risk, not a config preference.
**Open question for 3b — ANSWERED, 2026-08-15.** The worry was that the OpenAI surface might drop
reasoning the way the wiki documents for a mis-declared Anthropic backend. It does not. Checked at
the API before any profile was switched:
| surface | request | result |
|---|---|---|
| `/anthropic/v1/messages` | `deepseek-v4-flash`, 64 tokens | 200, response carries a real `"type":"thinking"` block |
| `/v1/chat/completions` | same | 200, message carries a populated `reasoning_content` (and a `reasoning` field) |
| `/v1/models` | — | 200, exactly `["deepseek-v4-flash"]` — the exact-name trap is clear |
| `/v1/models`, **no token** | — | **401** — Caddy is gating, as designed |
So reasoning survives on **both** surfaces, and the `/anthropic` choice for `local` is about protocol
correctness rather than a repair for a known loss. The last row matters on its own: the wiki warns
the gateway's own `SecurityPolicy` fails open, so it is worth knowing the proxy in front really does
refuse an unauthenticated request here.
---
## 4. Decisions
### D1 — switch, but keep the direct path as an explicit profile · **recommended**
Switching buys four things we do not have:
- **Free opencode capacity, off the shared credential.** The largest single win. See §3b — it retires
a real single point of failure, not just a cost line.
- **Per-consumer usage figures.** The cockpit counts requests per consumer. That is the first real
measurement of what the fleet consumes, and it feeds [CB-589](https://git.ltms.dev/lms/claude-bridge/issues/74) Gap 2 directly.
- **Our own revocable token.** One consumer to revoke if a worker ever leaks it, instead of a shared
`legacy` token used by four systems.
- **It works off-LAN.** `gx00.gw` resolves on the LAN only.
The cost is honest and worth stating: we add a TLS edge, an auth proxy and a gateway to the path of
every member spawn. The wiki keeps the direct route open precisely because "if the gateway breaks,
nothing that matters is blocked."
So keep it. Add a second profile `local-direct` pointing at `http://gx00.gw:8000` with **`weight: 0`**
— never auto-selected, still spawnable with an explicit `bridge_spawn{profile: "local-direct"}`.
That is exactly what CB-554 made `weight: 0` mean, and it turns the escape hatch into something the
lead can actually reach during an incident.
### D2 — do members also mount the gateway's `/mcp`? · **OPEN, operator's call**
Not a detail. `CLAUDE.md` states in two places that a member mounts **only** the bridge MCP, and a
worker's honesty rule leans on it ("never claim the result of a check you had no way to run").
- **Keep bridge-only.** The invariant stays true and simple. Workers stay cheap and narrow.
- **Add the gateway MCP.** Implementers get context7 documentation lookups, which is genuinely useful
for library work. But `mcpUrl` in `BridgedConfig.Profile` is a **single `String`**, so a
claude-code member can mount exactly one MCP — this needs a code change, not a config edit.
Note the invariant is **already inaccurate**: `opencode.json` gives `sol` and `terra` both context7
and gitea. So the choice is really "make the rule true" or "make the rule match reality". Either is
defensible; picking one is not mine to do.
### D3 — token scope
One consumer, `claude-bridge`, its token in `${SHARED_ENV}/tools/secrets.sh` as `AI_GATEWAY_TOKEN`,
referenced by name only. Never the literal value in `bridged.yaml` — `tokenEnv` exists for this.
---
## 5. Units of work
```mermaid
flowchart TB
U1["U1 · consumer token<br/>issue via cockpit, add to secrets.sh"]
U2["U2 · profile + guard<br/>bridged.yaml, restart"]
U3["U3 · verify live<br/>spawn, prove thinking survives"]
U4["U4 · context7 via gateway<br/>.mcp.json + opencode.json"]
U5["U5 · docs<br/>CLAUDE.md, wiki 11-Features"]
U1 --> U2 --> U3
U4 --> U5
U3 --> U5
```
| # | Scope | Who | Why |
|---|---|---|---|
| U1 | Issue the `claude-bridge` consumer at `auth.ltms.dev`; store as `AI_GATEWAY_TOKEN` | **operator** | touches secrets and a host we do not own |
| U2a | New `gx` opencode profile on `/v1` — pure config, no guard change | **lead** | `bridged.yaml` is gitignored, so a worker cannot see or edit it |
| U2b | `local` → `/anthropic`; add `local-direct` weight 0; add `llm.ltms.dev` to the guard allowlist | **lead** | same |
| U2c | One restart from a **login shell**, after U2a and U2b | **lead** | picks up `AI_GATEWAY_TOKEN` into the daemon env *and* the guard allowlist, in one stop |
| U3 | Live spawn on both new profiles; confirm reasoning survives on each surface | **lead** | needs real spawns and the running daemon |
| U4 | Point `.mcp.json` and `opencode.json` context7 at the gateway `/mcp`; rename pinned tools | delegatable | tracked files, self-contained |
| U5 | Fix the "members mount only the bridge" claim; add a `wiki/11-Features.md` entry | delegatable | writing, clear criteria |
U1 blocks U2a, U2b and U3. U4 and U5 do not depend on it.
**Write U2a and U2b, then restart once (U2c), then verify `gx` before `local`.** Since both profiles
need the same restart there is no reason to do two, but there is still a reason to *verify* in order:
`gx` exercises the token and the gateway with no guard involved, so if it fails the cause is upstream.
`local` adds the guard allowlist on top, so a failure there points at our config instead. Testing them
in that order separates the two causes instead of confusing them.
> **U1 status, 2026-08-15:** the operator issued the consumer and exported it as `AI_GATEWAY_TOKEN`
> (one key for every agent and MCP client behind `llm.ltms.dev`). Confirmed: it resolves in a login
> shell, is 48 characters and carries the documented `llmk-` prefix. The value was never printed.
---
## 6. Traps carried over from the wiki
Each of these cost someone real debugging time upstream. They apply to us.
1. **Rotating a token restarts the auth proxy, which drops in-flight streaming responses.** For us
that means rotating `AI_GATEWAY_TOKEN` kills every live member mid-turn, and an async ticket's
report goes with it. This is the same rule as a daemon redeploy: **drain the fleet first**
(`bridge_list` → `bridge_poll` anything wanted → `bridge_stop`), then rotate.
2. **The gateway's own `SecurityPolicy` fails open.** Standalone `aigw run` accepts it and silently
ignores it — an unauthenticated request returned **200**. Auth is the Caddy proxy in front, and
nothing else. Never reason as if the gateway authenticates.
3. **Exact model name.** A regex match routes fine but returns an **empty** `/v1/models` list while
completions keep working. A wrong name returns a bare 404 that reads exactly like a dead gateway.
4. **MCP tool names changed prefix separator.** Bifrost used one dash (`ct7-resolve-library-id`); the
gateway uses **two underscores** (`ct7__resolve-library-id`). Relevant only if U4 is done.
5. **`/v1/models` 404 vs empty list are different faults.** 404 means no route loaded at all; empty
means the model match is a regex. Do not conflate them when diagnosing.
---
## 7. Verification — what would prove this works
Merging config is not proving it. The checks, in order:
1. `bridge_spawn{profile: "gx"}` succeeds and the member completes a real turn ending in
`bridge_reply`. This is the first proof of the token, the URL and the model name, and it risks
nothing the fleet depends on.
2. `bridge_spawn{profile: "local"}` succeeds. If the guard allowlist was missed, this **throws** — a
loud, self-correcting failure, which is the good kind. If the restart was missed, it also throws,
for the same reason.
3. A `local` member completes a turn. That exercises streaming through two TLS edges, the auth proxy
and the gateway.
4. **Reasoning survives, checked separately on each surface.** For `local` on `/anthropic` this is
the check that catches the `/v1` versus `/anthropic` mistake, and it is the only one that does —
nothing else distinguishes a working passthrough from a translator quietly dropping thinking
deltas. For `gx` on `/v1`, this answers the open question in §3 rather than assuming it.
5. The cockpit at `auth.ltms.dev` shows requests counted against the `claude-bridge` consumer, not
`legacy`. That is the whole point of taking our own token.
6. `bridge_spawn{profile: "local-direct"}` still works, so the escape hatch is real rather than
theoretical.
7. `bridge_list` shows `gx` carrying no `credentialId`, so a `sol`/`terra` exhaustion cannot
quarantine it. This is the single-point-of-failure claim in §3b, checked rather than asserted.
---
## 7.1 What the live run actually found — 2026-08-15
U1–U2c were done, the daemon restarted onto them, and both new profiles were spawned for real. The
migration was then **reverted**. This section is the result, so none of it has to be re-derived.
### The blocker
`llm.ltms.dev` answers **HTTP 413 Request Entity Too Large** above **32 KiB (32768 bytes)**, on both
surfaces:
```
/v1 32695 bytes -> 200 /anthropic 32095 bytes -> 200
/v1 32795 bytes -> 413 /anthropic 32855 bytes -> 413
```
32 KiB is far below one real agent turn.
**Root cause — confirmed by the systems/vms side, 2026-08-15.** My guess that it was a Caddy
`request_body max_size` was **wrong**. It is Envoy, inside `aigw` on `llm.vm`. Envoy Gateway defaults
a listener's `per_connection_buffer_limit_bytes` to **32768**, and the AI Gateway buffers the *whole*
request body before it can route on the model name — so that default is not a network tuning knob
here, it is a hard ceiling on prompt size. Read out of the live Envoy `config_dump`:
```
listener default/llm/http per_connection_buffer_limit_bytes: 32768
```
Nobody chose 32 KiB; it was inherited from the default. Both TLS edges are innocent: the same
boundary reproduces on the LAN path and the internet path, and both 413s carry an `x-llm-consumer`
header their auth proxy sets only *after* authenticating — so the body cleared both edges and the
auth. Directly on `llm.vm`, `aigw` 413s at 39 KB while the vLLM backend accepts the same 39 KB and
answers 200.
**Do not plan around 32 KiB.** The intended ceiling is far higher. Their fix — a `ClientTrafficPolicy`
setting `bufferLimit: 8Mi` — is written but **not deployed** as of this note, pending their operator's
approval. I have not re-tested and will not until they confirm, so as not to measure a half-changed
system. Fixed in **systems/vms**, not here.
### The part worth remembering
Two members were spawned at the same moment with the same message:
| | `local` (claude-code, `/anthropic`) | `gx` (opencode, `/v1`) |
|---|---|---|
| READY → BUSY | 19:07:26 | 19:07:45 |
| BUSY → DONE | **19:08:51 (66s)** | **never — 10+ min, ticket FAILED** |
**`local` passed.** It passed only because the probe was three trivial questions in a fresh session,
so the request fit under 32 KiB. The profile looked healthy and was a landmine set to fire on the
first turn that reads a file.
So §7's checklist was not wrong, it was **too easy**. Any future run of it must use a task that reads
a real file. A liveness probe proves the token and the URL; it does not prove the path.
`gx` did not fail loudly either. Reproduced outside the bridge by running `opencode` by hand with the
launcher's own generated config:
```
Error: Request Entity Too Large
...compacts context, retries...
Error: Request Entity Too Large
```
opencode **catches the 413, compacts, and retries — indefinitely**. A member that fails loudly costs
one turn; this one costs the whole task and is indistinguishable from a slow worker.
> **Diagnosing a stuck opencode member.** Do not read its pane. The launcher writes its config to a
> temp dir and passes it as `OPENCODE_CONFIG` — find it with
> `ls -dt /var/folders/*/*/T/bridged-opencode-* | head -1`, check the provider block and the key's
> length and prefix (never its value), then reproduce with `opencode run --auto -m <provider>/<model>`
> using the same `OPENCODE_CONFIG`. That is what turned "it hangs" into a one-line error.
### What checked out, and needs no re-testing
- Token accepted on both surfaces. **Unauthenticated → 401**, so the Caddy proxy really does gate —
the wiki's "SecurityPolicy fails open" warning is about the gateway itself, not the edge.
- `/v1/models` returns exactly `["deepseek-v4-flash"]`, so trap 3 is clear.
- **Reasoning survives both surfaces** — see §3b above.
- The launcher's generated opencode provider block is correct, carrying a real 48-character `llmk-`
key rather than the `bridged-local-noauth` placeholder.
- `SubscriptionGuard` accepted `llm.ltms.dev` after the allowlist edit and the restart: `local`
spawned without throwing, which is the check that catches a missed restart.
### Resolution — both ceilings fixed, migration completed
systems/vms fixed both, and each was re-checked from this side rather than taken on trust:
| ceiling | was | now | our own check |
|---|---|---|---|
| listener buffer | 32 KiB | 32 Mi | 1.2 MB body → **200** (was 413) |
| LLM route timeout | 60s | 86400s | the request that truncated: **101s, `message_stop` present, 4000/4000** |
The timeout moved in two steps on 2026-08-15: 60s → 1800s, then 1800s → **86400s (24 hours)** after
the truncation risk below was discussed. They tried `request: 0s` first, which removes the
total-duration timer completely. It works, but on an `AIGatewayRoute` the **idle timeout is derived
from the request timeout**, so `0s` also removed any bound on a stalled connection. 86400s keeps a
reaper for dead connections while putting the truncation timer out of practical reach.
Neither was deliberate. The 32 KiB was Envoy Gateway's default `per_connection_buffer_limit_bytes`;
the 60s was Envoy AI Gateway's own documented default. The 60s bounded **generation** as well as
prompt size — a tiny prompt with a long answer returned 504 at 60.05s.
Two configuration facts worth keeping, from their bisection:
- **`ClientTrafficPolicy` is honoured in standalone `aigw run`; `BackendTrafficPolicy` is NOT.** A
`BackendTrafficPolicy` setting `requestTimeout` is accepted, logs nothing, and leaves the routes
unchanged (upstream `envoyproxy/gateway#9513`). What works is `timeouts: {request: …}` on each
`AIGatewayRoute` rule. Nothing from the outside distinguishes the two — the same silent-default
shape as their `SecurityPolicy` caveat.
- In that stack, "the config was accepted" proves nothing. Read the live `config_dump`.
## 7.2 The risk we accepted, and why we could not remove it
Raising the timeout made the failure **rare, not impossible**, and the residual failure is silent.
On a mid-response timeout over chunked HTTP/1.1, Envoy ends the chunked encoding *cleanly* instead of
resetting the connection, so the client receives what looks like a complete transfer
(`envoyproxy/envoy#17186` — acknowledged as a bug in 2021, closed by a stale bot, never fixed). The
December 2025 fix `envoyproxy/envoy#42269` changes locally-originated resets from `NO_ERROR` to
`INTERNAL_ERROR`, but it is **HTTP/2 only** and SSE clients here speak HTTP/1.1.
Measured on our side while the timeout was still 60s:
```
HTTP 200 61.07s 141992 bytes
message_stop 0 message_delta 0 error events 0
emitted 2473 of 4000, ending on a WELL-FORMED SSE frame
```
A syntactically valid stream that simply stops. Any timer firing mid-stream — route timeout, idle
timeout, `max_stream_duration` — fails this same way.
**The recommended defence does not transfer to us.** The right fix is to treat a stream with no
`message_stop` / `[DONE]` / `finish_reason` as failed. We cannot: our members are Claude Code and
opencode, third-party clients whose SSE parsing we do not own, and there is no seam to insert the
check. Whether either detects a missing terminator is unverified — and opencode's handling of the 413
(swallow, compact, retry forever, never surface an error) does not suggest it is strict.
So the honest statement of our position:
> Gateway traffic is acceptable at 86400s because a single request would have to run for 24 hours to
> trip the bug — **not** because we could detect it if it did.
At 86400s our **own** limit binds first, which is the ordering we want. `MessageService.ASYNC_TIMEOUT_MS`
caps a turn at 30 minutes, so a runaway request ends as a clean `FAILED` ticket that we raised, rather
than as a silently truncated `200` that we cannot see. While the gateway sat at 1800s the two numbers
were equal and did not nest, so a gateway-side stall could have been misread as a bug in our own ticket
handling. That ambiguity is now gone.
**If a member ever returns a confident but truncated answer, suspect this before anything in our own
code.** That is the whole reason this section exists.
---
## 8. Related
- [CB-589 / #74](https://git.ltms.dev/lms/claude-bridge/issues/74) — cost-first placement and a
gateway that reports live capacity. The per-consumer figures this migration unlocks are the first
input that ticket actually needs.
- `docs/CB-500-Multi-Tier-Coordination.md` §11 — the distributed-sandbox topology this gateway is
part of.
+1 -1
View File
@@ -26,7 +26,7 @@
],
"enabled": true,
"environment": {
"GITEA_ACCESS_TOKEN": "{env:GITEA_ACCESS_TOKEN}",
"GITEA_ACCESS_TOKEN": "{env:WORKER_GITEA_TOKEN}",
"GITEA_HOST": "{env:GITEA_HOST}"
}
}
+16 -2
View File
@@ -8,8 +8,10 @@
# This script exists to turn five remembered traps into one auditable command:
#
# 1. A piped `mvn` hides BUILD FAILURE behind a zero exit, so the build here is never piped.
# 2. The daemon must start from a LOGIN shell, or WORKER_GITEA_TOKEN is empty and workers cannot
# open a PR. Nothing in the daemon logs this, so the script checks it and says so out loud.
# 2. The daemon must start from a LOGIN shell, or the tokens it hands to members are empty:
# WORKER_GITEA_TOKEN (workers cannot open a PR) and AI_GATEWAY_TOKEN (401 at llm.ltms.dev).
# Both are read from the DAEMON's own environment at spawn time, so a value added to
# secrets.sh after startup is absent. Nothing logs this, so the script checks and says so.
# 3. An old daemon that never actually died looks identical from the outside, so the script waits
# for the process to exit and for the port to free before it starts a new one.
# 4. "It started" is not "it works": the script polls /healthz until it answers, and reports the
@@ -82,6 +84,18 @@ else
warn "Fix \${SHARED_ENV}/tools/secrets.sh before relying on worker checkpoints."
fi
# Same trap, second variable (CB-591). A profile's `tokenEnv:` is resolved from the DAEMON's own
# process environment by HerdrPeerLauncher.resolveEnv, so a token added to secrets.sh after the
# daemon started is simply absent. The launcher then injects an empty token and llm.ltms.dev answers
# 401 — long after the restart, and with nothing tying the two together.
if zsh -lc '[ -n "${AI_GATEWAY_TOKEN:-}" ]' 2>/dev/null; then
ok "AI_GATEWAY_TOKEN resolves in a login shell"
else
warn "AI_GATEWAY_TOKEN is EMPTY in a login shell."
warn "Any profile whose tokenEnv is AI_GATEWAY_TOKEN will get an empty token and 401 at the gateway."
warn "This only matters once a profile points at llm.ltms.dev — harmless before that."
fi
if [ "$CHECK_ONLY" = 1 ]; then
say "--check: nothing changed"
exit 0