Compare commits
29 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 7d41ccccee | |||
| bf616e192a | |||
| f8bd5d0c51 | |||
| ac40de1d30 | |||
| 7930a31b94 | |||
| 850fb12807 | |||
| b14b66ab03 | |||
| 15ff6bcde5 | |||
| 7822772905 | |||
| 0efe1567c0 | |||
| 837fed7690 | |||
| aa4ee64a34 | |||
| bb750cdba3 | |||
| d56c77b368 | |||
| 83e2ff06cf | |||
| f5deaafd06 | |||
| fe46311266 | |||
| 08968bb1b7 | |||
| 27bbd11f06 | |||
| 48d7841fbf | |||
| cc1df11f69 | |||
| fdfd4ac491 | |||
| d5dd5639ae | |||
| 7a120b3256 | |||
| 863d477966 | |||
| a36b7ccd7c | |||
| 32bf324a1e | |||
| 16de9df000 | |||
| 81a0cf4710 |
@@ -78,9 +78,11 @@ below are the procedure — run them in order, every task, not only the big ones
|
|||||||
what makes them reliable. Where the project ships no such skill, spell the procedure out in the
|
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
|
brief instead. The brief is self-contained — the worker sees your message and the repo, nothing
|
||||||
of your context, your plan, or your screen.
|
of your context, your plan, or your screen.
|
||||||
5. **Collect** — `bridge_poll{ticket}` → `bridge_ack{ticket, msgId}`. Answer a worker's `bridge_ask`
|
5. **Collect** — `bridge_poll{ticket}` → `bridge_ack{target, msgId}`. Answer a worker's `bridge_ask`
|
||||||
with `bridge_send{turnId, content}` — **not** `sessionId`. A worker gone quiet is diagnosed with
|
with `bridge_send{turnId, content}` — **not** `sessionId`. A worker gone quiet is diagnosed with
|
||||||
`bridge_status`, never by reading its terminal.
|
`bridge_status`, never by reading its terminal; it also reports an open question and the `turnId`
|
||||||
|
that answers it. **A worker's ask waits ~55 seconds, and no nudge makes that longer** — so never
|
||||||
|
brief a worker to "ask me". Decide before you delegate, or give it an explicit default.
|
||||||
6. **Verify yourself.** Re-run the build and the checks. A worker cannot run your IDE tooling, any
|
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
|
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.
|
(`… | tail`) hides failures behind a zero exit — never promote a worker's "clean" to a fact.
|
||||||
@@ -105,8 +107,8 @@ the merge — and merging on a reviewer's word is delegating it by proxy.
|
|||||||
|---|---|
|
|---|---|
|
||||||
| Confirm your own role | `bridge_whoami` |
|
| Confirm your own role | `bridge_whoami` |
|
||||||
| See backends available | `bridge_profiles` |
|
| See backends available | `bridge_profiles` |
|
||||||
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?}` → `sessionId` + `paneId` |
|
| Start a member | `bridge_spawn{role?, profile?, cwd?, worktree?, ticket?, sessionName?, resumeSessionId?}` → `sessionId` + `paneId` |
|
||||||
| See the fleet | `bridge_list` → `leads` (your peers) + `members` · one peer's state: `bridge_status{sessionId}` |
|
| See the fleet | `bridge_list` → `leads` (your peers) + `members` (each carries `agentSessionId` when its backend knows one) · one peer's state: `bridge_status{sessionId}` |
|
||||||
| Delegate (blocking) | `bridge_send{sessionId, content}` |
|
| Delegate (blocking) | `bridge_send{sessionId, content}` |
|
||||||
| Delegate (long task) | `bridge_send{sessionId, content, wait:false}` → ticket → `bridge_poll{ticket}` |
|
| Delegate (long task) | `bridge_send{sessionId, content, wait:false}` → ticket → `bridge_poll{ticket}` |
|
||||||
| Answer a member's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
|
| Answer a member's `bridge_ask` | `bridge_send{turnId, content}` — **not** `sessionId` |
|
||||||
@@ -190,8 +192,10 @@ charter, not here.
|
|||||||
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
|
- **Never commit** `.mcp.json` (the primary's local copy, flagged `--skip-worktree`) or `wiki/`
|
||||||
(a submodule with its own remote).
|
(a submodule with its own remote).
|
||||||
- **Flows and the error model** — rendezvous, `bridge_ask`, detached delivery, turn-done fallback —
|
- **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
|
are diagrammed in `docs/MCP-Contract.md` **§6 only**. The rest of that page is a pre-build design
|
||||||
session's context.
|
doc whose tool names, parameter names and REST paths never caught up with the code, so do not use
|
||||||
|
it as the tool reference (CB-609). Section 6 is kept out of this file because this file loads into
|
||||||
|
every session's context.
|
||||||
|
|
||||||
### Redeploying the daemon — the lead may do this (primary only)
|
### Redeploying the daemon — the lead may do this (primary only)
|
||||||
|
|
||||||
|
|||||||
@@ -18,8 +18,9 @@ one unified Claude setup and the **sole communication gateway** (REST/SSE stays
|
|||||||
clients; any broker is `bridged`-internal, below the gateway).
|
clients; any broker is `bridged`-internal, below the gateway).
|
||||||
herdr owns the PTYs, multiplexing, persistence, and **agent-status events**; `bridged` owns
|
herdr owns the PTYs, multiplexing, persistence, and **agent-status events**; `bridged` owns
|
||||||
policy (subscription boundary, session lifecycle, status-gated delivery) and the client
|
policy (subscription boundary, session lifecycle, status-gated delivery) and the client
|
||||||
contract. The worker `claude` launches with `ANTHROPIC_BASE_URL=https://ollama.ltms.dev` + a
|
contract. A Claude member launches with `ANTHROPIC_BASE_URL` pointed at the gateway,
|
||||||
bearer token; the primary Opus stays env-clean and calls `bridged`'s MCP tools.
|
`https://llm.ltms.dev/anthropic`, plus a bearer token; the lead stays env-clean and calls
|
||||||
|
`bridged`'s MCP tools. See the wiki's **[13 User Guide](wiki/13-User-Guide.md)** to run it.
|
||||||
|
|
||||||
```mermaid
|
```mermaid
|
||||||
flowchart LR
|
flowchart LR
|
||||||
@@ -31,7 +32,7 @@ flowchart LR
|
|||||||
end
|
end
|
||||||
HERDR["herdr<br/>panes · agent-status"]
|
HERDR["herdr<br/>panes · agent-status"]
|
||||||
W["worker claude pane<br/>ANTHROPIC_BASE_URL set<br/>MCP client"]
|
W["worker claude pane<br/>ANTHROPIC_BASE_URL set<br/>MCP client"]
|
||||||
M["ollama.ltms.dev<br/>(worker model)"]
|
M["llm.ltms.dev<br/>(the one gateway)"]
|
||||||
|
|
||||||
OPUS -->|"MCP bridge_send (blocks)"| SRV
|
OPUS -->|"MCP bridge_send (blocks)"| SRV
|
||||||
W -.->|"MCP bridge_reply"| SRV
|
W -.->|"MCP bridge_reply"| SRV
|
||||||
|
|||||||
@@ -208,6 +208,29 @@ profiles:
|
|||||||
# `bridge_spawn{profile:"gx10"}` against it is refused too; a cap holds even when the profile
|
# `bridge_spawn{profile:"gx10"}` against it is refused too; a cap holds even when the profile
|
||||||
# is named directly. Negative is refused at config load — there is no sane meaning for it.
|
# is named directly. Negative is refused at config load — there is no sane meaning for it.
|
||||||
maxLoad: 2
|
maxLoad: 2
|
||||||
|
# subscription: true
|
||||||
|
# THE KNOB THAT DECIDES WHO PAYS (CB-539). Default false. When true, this profile's members
|
||||||
|
# run on the OPERATOR'S OWN Claude subscription instead of a metered endpoint — every spawn
|
||||||
|
# bills your plan and eats your usage limit. Off-subscription is the whole point of this
|
||||||
|
# daemon, so treat `true` as a deliberate exception, not a convenience.
|
||||||
|
#
|
||||||
|
# What changes when it is set (ClaudeCodeLauncher):
|
||||||
|
# - no ANTHROPIC_BASE_URL and no ANTHROPIC_AUTH_TOKEN are injected — the member inherits
|
||||||
|
# the operator's own Claude Code auth, which is exactly why it bills the plan;
|
||||||
|
# - SubscriptionGuard never vets it, because there is no baseUrl to vet;
|
||||||
|
# - no token is required, so `tokenEnv` is irrelevant here.
|
||||||
|
#
|
||||||
|
# MUTUALLY EXCLUSIVE with `baseUrl` — setting both is refused at config load (CB-542). On the
|
||||||
|
# subscription path no guard would vet the URL, so allowing both would be a way around the
|
||||||
|
# guard rather than a configuration.
|
||||||
|
#
|
||||||
|
# GOTCHA 1 — it is invisible to the startup secret check. `Bridged.reportRequiredSecrets`
|
||||||
|
# skips subscription profiles on purpose (they need no token), so a boot log that reports
|
||||||
|
# every secret as fine says nothing about these profiles.
|
||||||
|
#
|
||||||
|
# GOTCHA 2 — `maxLoad` is the ONLY throttle you have here. There is no metering, no budget
|
||||||
|
# and no refusal on cost; the cap on live members is the single thing standing between a
|
||||||
|
# fan-out and your monthly limit. Set it deliberately and keep it small.
|
||||||
# gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302)
|
# gitTokenEnv: GITEA_TOKEN # opt-in: let this profile's workers open their own PR (CB-302)
|
||||||
# gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv
|
# gitHostEnv: GITEA_HOST # defaults to GITEA_HOST; injected only with gitTokenEnv
|
||||||
# exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578)
|
# exhaustedPattern: "usage limit has been reached" # opt-in: classify a usage-limit refusal (CB-578)
|
||||||
@@ -275,6 +298,26 @@ profiles:
|
|||||||
# argv: ["opencode"]
|
# argv: ["opencode"]
|
||||||
# How an unqualified spawn chooses a profile: fixed (default, reproduces pre-CB-518 behaviour),
|
# How an unqualified spawn chooses a profile: fixed (default, reproduces pre-CB-518 behaviour),
|
||||||
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
|
# round-robin, or weighted. Omitting this key is a strict no-op for existing configs.
|
||||||
|
#
|
||||||
|
# `weighted` IS NOT "cheapest first" — read this before you set weights (CB-589).
|
||||||
|
# It is smooth weighted round-robin: it spreads spawns across EVERY profile that has a free slot,
|
||||||
|
# in weight ratio. It has no idea which profile costs money. So with local:10 / paid:2 you do not
|
||||||
|
# get "use local, overflow to paid" — you get roughly one spawn in six going to the paid profile
|
||||||
|
# while the local box still has a free slot.
|
||||||
|
#
|
||||||
|
# There is a sharper second effect. The policy's running score map lives for the daemon's whole
|
||||||
|
# life. While a profile is at maxLoad it is filtered out and its score FREEZES, so the paid
|
||||||
|
# profiles keep accumulating against it. When the local slot frees up it returns with a stale
|
||||||
|
# score and can LOSE the next pick — a paid spawn while the free box sits idle.
|
||||||
|
#
|
||||||
|
# Until a real cost-first policy exists, the workaround is to make the ratio decisive rather than
|
||||||
|
# proportional: give the free profile a weight so large that it wins every pick it is eligible
|
||||||
|
# for, and paid profiles only ever take genuine overflow. On this host that is local weight 100
|
||||||
|
# against paid weights of ~1.
|
||||||
|
#
|
||||||
|
# The gotcha with that workaround: it expresses a PREFERENCE ORDER through a RATIO knob. Add a
|
||||||
|
# future profile at weight 150 and it silently outranks the free box, with nothing to warn you.
|
||||||
|
# Re-check the weights whenever you add a profile.
|
||||||
placement: weighted
|
placement: weighted
|
||||||
|
|
||||||
# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in
|
# How long a credential sits out after a BACKEND_EXHAUSTED classification (CB-578 stage B), in
|
||||||
@@ -430,6 +473,91 @@ guard:
|
|||||||
- gx00.gw
|
- gx00.gw
|
||||||
- gx01.gw
|
- gx01.gw
|
||||||
|
|
||||||
|
# Member credential policy (CB-596, gitea issue #82). A herdr pane runs a LOGIN shell, and that
|
||||||
|
# shell re-sources the operator's own secret store — so a spawned member inherits every credential
|
||||||
|
# the operator's shell holds, not just the ones bridged means to give it. Measured on this host:
|
||||||
|
# 31 credential names, all set, with only ONE (GITEA_ACCESS_TOKEN) blocked before this — and that
|
||||||
|
# block was a single name hardcoded in HerdrPeerLauncher.java, not driven by this file. This block
|
||||||
|
# replaces that hardcoded shadow with a config-driven list of names.
|
||||||
|
#
|
||||||
|
# ROUND-2 CORRECTION, measured live: the pane-creation env overlay below (applied at tab.create /
|
||||||
|
# pane.split, BEFORE the pane's login shell runs) does NOT survive that login shell for any name
|
||||||
|
# secrets.sh actually exports — the shell re-exports it afterwards and overwrites the sentinel.
|
||||||
|
# Proof: GITEA_ACCESS_TOKEN comes back blocked only because secrets.sh itself carries a guarded
|
||||||
|
# export (`[ -n "${BRIDGED_MEMBER:-}" ] || export GITEA_ACCESS_TOKEN=...`) — that guard, not this
|
||||||
|
# file, is what wins. No other name in `known` below has a matching guard in secrets.sh yet (1
|
||||||
|
# guard measured against 33 export lines there). So today this block's overlay is REAL protection
|
||||||
|
# only for a name secrets.sh does not export, or a peer kind whose pane never runs a login shell —
|
||||||
|
# for everything secrets.sh exports and guards, the guard in secrets.sh (out of scope for this
|
||||||
|
# ticket) is what actually blocks it, not this list. An exec-time fix (winning after the login
|
||||||
|
# shell finishes, before the agent process starts) was attempted and found to have no seam in the
|
||||||
|
# current herdr protocol — AgentControl.start takes a fixed `kind` (herdr resolves the executable)
|
||||||
|
# plus trailing CLI args for that binary, not an arbitrary argv or an env map; only tab.create /
|
||||||
|
# pane.split accept `env`, and that is this same pane-creation overlay. See gitea #82 for the open
|
||||||
|
# design question this leaves.
|
||||||
|
#
|
||||||
|
# DENY-BY-DEFAULT, NOT A DENY-LIST. A deny-list (name the bad ones, let everything else through) is
|
||||||
|
# silently wrong the moment the operator's store gains a new secret — nothing would ever report it.
|
||||||
|
# Deny-by-default inverts that: `known` bounds the blast radius to names actually enumerated below,
|
||||||
|
# and EVERY one of them is blocked UNLESS it is also in `allow`. Omitting this block entirely (the
|
||||||
|
# shipped default) blocks NOTHING — unlike most optional blocks in this file, absence here is a real
|
||||||
|
# gap, not a safe "feature off". A name that is neither `known` nor `allow`-ed is not silently let
|
||||||
|
# through either: the daemon logs a WARN naming any credential-shaped env var it finds on neither
|
||||||
|
# list (never its value), so a secret added to the store later does not go unnoticed forever.
|
||||||
|
#
|
||||||
|
# policy → only "deny-by-default" exists today (an operator-authored deny-list was deliberately
|
||||||
|
# rejected — see above). An unrecognized value refuses to start, naming it.
|
||||||
|
# allow → credential names a member legitimately needs. Left OUT of the pane's env overlay
|
||||||
|
# entirely, so the value the pane's own (login) shell exports passes through untouched.
|
||||||
|
# known → every credential name the operator's store is known to export. Every name here NOT
|
||||||
|
# also in `allow` is overlaid with a non-secret sentinel value before the pane's login
|
||||||
|
# shell runs — real protection only for names the login shell does not itself re-export
|
||||||
|
# (see the ROUND-2 CORRECTION note above for the ones it does).
|
||||||
|
#
|
||||||
|
# HOT-RELOADABLE the same way `fleet:` is (CB-559): read fresh on every spawn, so editing this list
|
||||||
|
# and reloading config (or restarting) changes what the NEXT spawn inherits; already-running members
|
||||||
|
# are unaffected either way.
|
||||||
|
# memberCredentials:
|
||||||
|
# policy: deny-by-default
|
||||||
|
# allow:
|
||||||
|
# - AI_GATEWAY_TOKEN # named in a profile's tokenEnv (local/gx) — a member reaching the
|
||||||
|
# # gateway is by design, not a leak
|
||||||
|
# - WORKER_GITEA_TOKEN # the repo-scoped forge token a member needs to open its own PR (CB-302)
|
||||||
|
# - CONTEXT7_TOKEN # already decided as allowed by CB-593
|
||||||
|
# - GITEA_HOST # not a credential — a hostname, paired with the forge token above
|
||||||
|
# known:
|
||||||
|
# - AI_GATEWAY_TOKEN
|
||||||
|
# - BESZEL_ADMIN_EMAIL
|
||||||
|
# - BESZEL_ADMIN_PASSWORD
|
||||||
|
# - BESZEL_HUB_URL
|
||||||
|
# - BESZEL_KEY
|
||||||
|
# - BESZEL_UNIVERSAL_TOKEN
|
||||||
|
# - BRAIN_MCP_TOKEN
|
||||||
|
# - CF_ACCOUNT_ID
|
||||||
|
# - CF_API_TOKEN
|
||||||
|
# - CF_USER_TOKEN
|
||||||
|
# - CONFLUENCE_API_TOKEN
|
||||||
|
# - CONFLUENCE_USERNAME
|
||||||
|
# - CONTEXT7_TOKEN
|
||||||
|
# - GITEA_ACCESS_TOKEN
|
||||||
|
# - GITEA_HOST
|
||||||
|
# - GITLAB_OAUTH_CLIENT_SECRET
|
||||||
|
# - GITLAB_PERSONAL_ACCESS_TOKEN
|
||||||
|
# - GRAFANA_ADMIN_PASSWORD
|
||||||
|
# - GRAFANA_ADMIN_USER
|
||||||
|
# - HASS_TOKEN
|
||||||
|
# - HW_PASSWORD
|
||||||
|
# - HW_USER
|
||||||
|
# - LTMS_API_KEY
|
||||||
|
# - MEMORY_MCP_TOKEN
|
||||||
|
# - METRICS_PUSH_TOKEN
|
||||||
|
# - OPENCODE_AUTOMODE_MODEL
|
||||||
|
# - TELEGRAM_BOT_TOKEN
|
||||||
|
# - TELEGRAM_CHAT_ID
|
||||||
|
# - TS_API_KEY
|
||||||
|
# - TS_AUTHKEY
|
||||||
|
# - WORKER_GITEA_TOKEN
|
||||||
|
|
||||||
# Spawn-readiness gate (CB-306). The launcher blocks until the worker's herdr status is
|
# Spawn-readiness gate (CB-306). The launcher blocks until the worker's herdr status is
|
||||||
# injectable (IDLE/BLOCKED/DONE) or the timeout elapses. 0 disables the gate.
|
# injectable (IDLE/BLOCKED/DONE) or the timeout elapses. 0 disables the gate.
|
||||||
# NOTE: keys are camelCase — config is bound by plain Jackson with no naming strategy and
|
# NOTE: keys are camelCase — config is bound by plain Jackson with no naming strategy and
|
||||||
|
|||||||
@@ -87,6 +87,11 @@ public final class Bridged {
|
|||||||
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
|
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
|
||||||
// boots fine either way — this is the only thing that says so out loud.
|
// boots fine either way — this is the only thing that says so out loud.
|
||||||
reportRequiredSecrets(cfg);
|
reportRequiredSecrets(cfg);
|
||||||
|
// CB-596: an absent (or empty) memberCredentials: block blocks NOTHING — no credential
|
||||||
|
// name is hardcoded any more to fall back on. Say so loudly, the same way a missing
|
||||||
|
// secret is reported above, so upgrading past this commit never silently drops CB-592's
|
||||||
|
// protection.
|
||||||
|
reportMemberCredentialsGap(cfg);
|
||||||
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
|
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
|
||||||
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
|
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
|
||||||
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
|
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
|
||||||
@@ -139,13 +144,15 @@ public final class Bridged {
|
|||||||
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
|
adapters.add(new ClaudeCodeLauncher(agents, spaces, guard,
|
||||||
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
claudeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||||
() -> config.get().fleet()));
|
() -> config.get().fleet(),
|
||||||
|
() -> config.get().memberCredentials()));
|
||||||
}
|
}
|
||||||
if (!opencodeProfiles.isEmpty()) {
|
if (!opencodeProfiles.isEmpty()) {
|
||||||
adapters.add(new OpenCodeLauncher(agents, spaces,
|
adapters.add(new OpenCodeLauncher(agents, spaces,
|
||||||
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
opencodeProfiles, cfg.effectiveDefaultProfile(), System::getenv,
|
||||||
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
cfg.spawnReadyTimeoutMs(), cfg.spawnReadyPollMs(),
|
||||||
() -> config.get().fleet()));
|
() -> config.get().fleet(),
|
||||||
|
() -> config.get().memberCredentials()));
|
||||||
}
|
}
|
||||||
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
AtomicReference<Function<String, Integer>> liveCountRef = new AtomicReference<>(_ -> 0);
|
||||||
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
// CB-578 stage B: one quarantine tracker for the whole daemon, shared between the launcher
|
||||||
@@ -441,6 +448,11 @@ public final class Bridged {
|
|||||||
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
reason += "; worktree=" + detail.worktreePath() + " branch=" + detail.branch()
|
||||||
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
+ " snapshot=" + (detail.snapshotRef() != null ? detail.snapshotRef() : "none");
|
||||||
}
|
}
|
||||||
|
// CB-584 (issue #65 criterion 5): also name the agent session, so a lead can resume the
|
||||||
|
// member's conversation instead of only re-dispatching a fresh one onto the same files.
|
||||||
|
if (detail.agentSessionId() != null) {
|
||||||
|
reason += " agentSessionId=" + detail.agentSessionId();
|
||||||
|
}
|
||||||
messages.abandon(detail.terminalId(), reason);
|
messages.abandon(detail.terminalId(), reason);
|
||||||
replyInbox.release(detail.terminalId());
|
replyInbox.release(detail.terminalId());
|
||||||
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
primaryRegistry.forgetDelegation(detail.terminalId()); // CB-532: don't leak the lead binding
|
||||||
@@ -608,6 +620,31 @@ public final class Bridged {
|
|||||||
});
|
});
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: {@code known:} empty (block absent entirely, or present but empty) means {@link
|
||||||
|
* BridgedConfig.MemberCredentials#blockedSet()} is empty too — every member pane inherits the
|
||||||
|
* operator's whole secret store, unblocked, exactly the defect this ticket fixes. Unlike a
|
||||||
|
* missing token ({@link #reportRequiredSecrets}), there is no name to point at: the point is
|
||||||
|
* that the block itself is missing. Warn once at startup and say what to add; never refuse to
|
||||||
|
* start over it — see {@link #reportRequiredSecrets} for why a daemon that boots and says
|
||||||
|
* what is wrong beats one that will not boot at all.
|
||||||
|
*
|
||||||
|
* <p>Package-private so the test can capture the log directly, the same way {@link
|
||||||
|
* #requiredSecretEnvVars} is exposed for {@link #reportRequiredSecrets}'s own test.
|
||||||
|
*/
|
||||||
|
static void reportMemberCredentialsGap(BridgedConfig cfg) {
|
||||||
|
BridgedConfig.MemberCredentials creds = cfg.memberCredentials();
|
||||||
|
if (creds != null && !creds.known().isEmpty()) {
|
||||||
|
log.info("memberCredentials: {} known name(s), {} allowed — blocking {} on every spawn",
|
||||||
|
creds.known().size(), creds.allow().size(), creds.blockedSet().size());
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
log.warn("memberCredentials: absent or empty — the daemon will start anyway, and every "
|
||||||
|
+ "member pane inherits the operator's WHOLE secret store, unblocked (CB-592's "
|
||||||
|
+ "protection is lost). Add a memberCredentials: block (policy/allow/known) to "
|
||||||
|
+ "bridged.yaml — see bridged.example.yaml — and restart.");
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import com.fasterxml.jackson.databind.ObjectMapper;
|
|||||||
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
||||||
import dev.ltms.bridged.msg.AmqpReplyInbox;
|
import dev.ltms.bridged.msg.AmqpReplyInbox;
|
||||||
import dev.ltms.bridged.peer.MemberRole;
|
import dev.ltms.bridged.peer.MemberRole;
|
||||||
|
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||||
import org.slf4j.Logger;
|
import org.slf4j.Logger;
|
||||||
import org.slf4j.LoggerFactory;
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
@@ -65,6 +66,9 @@ import java.util.Set;
|
|||||||
* {@code BackendQuarantine} built at startup, so it is DEFERRED: changing it
|
* {@code BackendQuarantine} built at startup, so it is DEFERRED: changing it
|
||||||
* needs a restart, and a quarantine already running keeps whatever cooldown was
|
* needs a restart, and a quarantine already running keeps whatever cooldown was
|
||||||
* live when it started.
|
* live when it started.
|
||||||
|
* @param memberCredentials deny-by-default policy (CB-596) for which of the operator's own host
|
||||||
|
* credentials a spawned member's pane inherits. {@code null} (the block
|
||||||
|
* omitted) blocks nothing — see {@link MemberCredentials}.
|
||||||
*/
|
*/
|
||||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||||
public record BridgedConfig(
|
public record BridgedConfig(
|
||||||
@@ -84,7 +88,19 @@ public record BridgedConfig(
|
|||||||
String placement,
|
String placement,
|
||||||
Auth auth,
|
Auth auth,
|
||||||
ConfigReload configReload,
|
ConfigReload configReload,
|
||||||
Integer quarantineCooldownSeconds) {
|
Integer quarantineCooldownSeconds,
|
||||||
|
MemberCredentials memberCredentials) {
|
||||||
|
|
||||||
|
/** Back-compat form before the CB-596 {@code memberCredentials:} block was added. */
|
||||||
|
public BridgedConfig(Bind bind, String herdrSocket, Map<String, Profile> profiles, Guard guard,
|
||||||
|
String worktreeRoot, Lifecycle lifecycle, Integer spawnReadyTimeoutMs,
|
||||||
|
Integer spawnReadyPollMs, Broker broker, Primary primary, Fleet fleet,
|
||||||
|
LeadHeartbeat leadHeartbeat, Health health, String placement, Auth auth,
|
||||||
|
ConfigReload configReload, Integer quarantineCooldownSeconds) {
|
||||||
|
this(bind, herdrSocket, profiles, guard, worktreeRoot, lifecycle, spawnReadyTimeoutMs,
|
||||||
|
spawnReadyPollMs, broker, primary, fleet, leadHeartbeat, health, placement, auth,
|
||||||
|
configReload, quarantineCooldownSeconds, null);
|
||||||
|
}
|
||||||
|
|
||||||
/** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */
|
/** Default cooldown (CB-578 stage B) when {@code quarantineCooldownSeconds} is absent/non-positive. */
|
||||||
public static final int DEFAULT_QUARANTINE_COOLDOWN_SECONDS = 1800;
|
public static final int DEFAULT_QUARANTINE_COOLDOWN_SECONDS = 1800;
|
||||||
@@ -925,6 +941,58 @@ public record BridgedConfig(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: which of the operator's own host credentials a spawned member's pane may inherit.
|
||||||
|
*
|
||||||
|
* <p>A herdr pane runs a login shell that re-sources the operator's own secret store, so a
|
||||||
|
* member inherits every credential the operator's shell holds — measured at 31 names on this
|
||||||
|
* host, of which only one ({@code GITEA_ACCESS_TOKEN}) used to be blocked, and that block was a
|
||||||
|
* single name hardcoded in {@link HerdrPeerLauncher} rather than driven by config (gitea issue
|
||||||
|
* #82). This record replaces that hardcoded shadow with a config-driven one.
|
||||||
|
*
|
||||||
|
* <p><b>deny-by-default, not a deny-list.</b> A deny-list (block these specific names, let
|
||||||
|
* everything else through) is silently wrong the moment a new secret is added to the operator's
|
||||||
|
* store — nothing would ever report it. Deny-by-default inverts that: {@link #known} bounds the
|
||||||
|
* blast radius to names the operator has actually enumerated, and every one of them is blocked
|
||||||
|
* UNLESS it is also in {@link #allow}. A name that shows up in neither list is not silently
|
||||||
|
* allowed — see {@code HerdrPeerLauncher}'s gap detector, which logs it.
|
||||||
|
*
|
||||||
|
* @param policy how the block is computed. Only {@link #POLICY_DENY_BY_DEFAULT} is understood
|
||||||
|
* today; {@code null}/blank defaults to it. An operator's own deny-list is
|
||||||
|
* deliberately not supported — see above.
|
||||||
|
* @param allow credential names a member legitimately needs (e.g. the gateway token it reaches
|
||||||
|
* the LLM through, the repo-scoped forge token it opens its own PR with). Every
|
||||||
|
* name here is left unmentioned in the pane's env overlay, so the value the pane's
|
||||||
|
* own (login) shell exports passes through untouched.
|
||||||
|
* @param known every credential name the operator's store is known to export. Every name here
|
||||||
|
* that is NOT also in {@link #allow} is overlaid with a non-secret sentinel value,
|
||||||
|
* shadowing whatever the pane's login shell would otherwise export for it.
|
||||||
|
*/
|
||||||
|
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||||
|
public record MemberCredentials(String policy, List<String> allow, List<String> known) {
|
||||||
|
|
||||||
|
/** The only policy this build understands: block every {@code known} name not in {@code allow}. */
|
||||||
|
public static final String POLICY_DENY_BY_DEFAULT = "deny-by-default";
|
||||||
|
|
||||||
|
public MemberCredentials {
|
||||||
|
policy = (policy == null || policy.isBlank()) ? POLICY_DENY_BY_DEFAULT : policy.toLowerCase();
|
||||||
|
allow = allow == null ? List.of() : List.copyOf(allow);
|
||||||
|
known = known == null ? List.of() : List.copyOf(known);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** {@link #allow} as a set, for membership checks. */
|
||||||
|
public Set<String> allowSet() {
|
||||||
|
return Set.copyOf(allow);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** {@link #known} minus {@link #allow} — the names a spawn must shadow. */
|
||||||
|
public Set<String> blockedSet() {
|
||||||
|
Set<String> blocked = new java.util.LinkedHashSet<>(known);
|
||||||
|
blocked.removeAll(allowSet());
|
||||||
|
return blocked;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The candidate profiles an unqualified spawn of {@code role} chooses between, in definition
|
* The candidate profiles an unqualified spawn of {@code role} chooses between, in definition
|
||||||
* order (CB-557).
|
* order (CB-557).
|
||||||
@@ -967,11 +1035,16 @@ public record BridgedConfig(
|
|||||||
/**
|
/**
|
||||||
* Top-level keys this version understands. Used only to warn about the rest — see
|
* Top-level keys this version understands. Used only to warn about the rest — see
|
||||||
* {@link #warnUnknownTopLevelKeys}. Keep in step with the record components.
|
* {@link #warnUnknownTopLevelKeys}. Keep in step with the record components.
|
||||||
|
*
|
||||||
|
* <p>Package-private (not {@code private}) so a test can assert every key here is documented in
|
||||||
|
* {@code bridged.example.yaml} — the only committed description of the config schema, since
|
||||||
|
* {@code bridged.yaml} itself is gitignored.
|
||||||
*/
|
*/
|
||||||
private static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
|
static final Set<String> KNOWN_TOP_LEVEL_KEYS = Set.of(
|
||||||
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
|
"bind", "herdrSocket", "profiles", "guard", "worktreeRoot",
|
||||||
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
"lifecycle", "spawnReadyTimeoutMs", "spawnReadyPollMs", "broker", "primary", "fleet",
|
||||||
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds");
|
"leadHeartbeat", "health", "placement", "auth", "configReload", "quarantineCooldownSeconds",
|
||||||
|
"memberCredentials");
|
||||||
|
|
||||||
/** Load and validate config from {@code path}. */
|
/** Load and validate config from {@code path}. */
|
||||||
public static BridgedConfig load(Path path) {
|
public static BridgedConfig load(Path path) {
|
||||||
@@ -982,7 +1055,15 @@ public record BridgedConfig(
|
|||||||
warnUnknownTopLevelKeys(yaml, path);
|
warnUnknownTopLevelKeys(yaml, path);
|
||||||
rejectDuplicateMemberSlots(yaml);
|
rejectDuplicateMemberSlots(yaml);
|
||||||
rejectNegativeMaxLoad(yaml);
|
rejectNegativeMaxLoad(yaml);
|
||||||
|
rejectUnknownKind(yaml);
|
||||||
|
rejectUnknownAuthMode(yaml);
|
||||||
|
rejectUnknownPlacement(yaml);
|
||||||
|
rejectUnknownMemberCredentialsPolicy(yaml);
|
||||||
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
|
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
|
||||||
|
// CB-606: validated here, eagerly, using PlacementPolicies.fromName as the single source
|
||||||
|
// of truth — not lazily at first spawn (see CompositePeerLauncher's placementPolicy
|
||||||
|
// Supplier), where a bad name would still start a daemon that looks healthy.
|
||||||
|
rejectUnknownPlacementPolicy(cfg.placement());
|
||||||
return cfg.withDefaults();
|
return cfg.withDefaults();
|
||||||
} catch (IOException e) {
|
} catch (IOException e) {
|
||||||
throw new UncheckedIOException("cannot read bridged config at " + path, e);
|
throw new UncheckedIOException("cannot read bridged config at " + path, e);
|
||||||
@@ -1276,6 +1357,201 @@ public record BridgedConfig(
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** The peer kinds this build has an adapter for — {@link Profile#kind()}'s only valid values. */
|
||||||
|
private static final Set<String> KNOWN_KINDS = Set.of(Profile.KIND_CLAUDE_CODE, Profile.KIND_OPENCODE);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reject a profile whose {@code kind:} is not one of {@link #KNOWN_KINDS} (CB-604), naming the
|
||||||
|
* profile, the value it set, and the accepted set.
|
||||||
|
*
|
||||||
|
* <p>{@link Profile}'s compact constructor only lower-cases {@code kind} and compares it against
|
||||||
|
* {@code KIND_OPENCODE} — anything else, including a typo like {@code opencod}, silently falls
|
||||||
|
* into the claude-code bucket ({@link dev.ltms.bridged.member.CompositePeerLauncher} routes by
|
||||||
|
* exact adapter claim, not by membership in a known set). With {@code argv:} also unset, the argv
|
||||||
|
* default special-cases only the exact string {@code "claude-code"}, so the launch command falls
|
||||||
|
* back to {@code List.of(kind)} — the daemon then tries to run a program literally named after the
|
||||||
|
* typo. {@code CompositePeerLauncher}'s constructor already treats a profile claimed by two
|
||||||
|
* adapters as fatal (CB-402); an unrecognized kind is the same class of adapter-routing mistake
|
||||||
|
* and gets the same treatment here, at config load, rather than surfacing later as a failed spawn.
|
||||||
|
*
|
||||||
|
* @param yaml the raw config text
|
||||||
|
* @throws IllegalStateException when any profile's {@code kind} is a non-blank value not in
|
||||||
|
* {@link #KNOWN_KINDS} (case-insensitive)
|
||||||
|
*/
|
||||||
|
static void rejectUnknownKind(String yaml) {
|
||||||
|
Map<?, ?> raw;
|
||||||
|
try {
|
||||||
|
raw = YAML.readValue(yaml, Map.class);
|
||||||
|
} catch (IOException | IllegalArgumentException e) {
|
||||||
|
return; // a malformed file is reported by the real parse, not here
|
||||||
|
}
|
||||||
|
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
List<String> bad = profiles.entrySet().stream()
|
||||||
|
.filter(e -> e.getValue() instanceof Map<?, ?> p
|
||||||
|
&& p.get("kind") instanceof String k && !k.isBlank()
|
||||||
|
&& !KNOWN_KINDS.contains(k.toLowerCase()))
|
||||||
|
.map(e -> String.valueOf(e.getKey()) + "=" + ((Map<?, ?>) e.getValue()).get("kind"))
|
||||||
|
.sorted()
|
||||||
|
.toList();
|
||||||
|
if (!bad.isEmpty()) {
|
||||||
|
throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad)
|
||||||
|
+ "] set an unrecognized kind — accepted values are "
|
||||||
|
+ String.join(", ", KNOWN_KINDS.stream().sorted().toList())
|
||||||
|
+ " (case-insensitive); an unrecognized kind would otherwise fall back to the"
|
||||||
|
+ " claude-code adapter and try to launch a program named after the typo.");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The auth modes this build understands — {@link Auth#mode()}'s only valid values. */
|
||||||
|
private static final Set<String> KNOWN_AUTH_MODES = Set.of(Auth.MODE_LOOPBACK_TRUST, Auth.MODE_TOKEN);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reject an {@code auth.mode} that is not one of {@link #KNOWN_AUTH_MODES} (CB-606), naming the
|
||||||
|
* value and the accepted set.
|
||||||
|
*
|
||||||
|
* <p>{@link Auth}'s compact constructor only lower-cases {@code mode}, and
|
||||||
|
* {@link Auth#tokenMode()} only compares the result against {@code MODE_TOKEN} — anything else,
|
||||||
|
* including a typo like {@code toekn}, silently behaves as {@code loopback-trust}. That fallback
|
||||||
|
* is otherwise checked only by {@link #validateAuthExposure()}, and only when the bind is
|
||||||
|
* non-loopback: on a loopback bind (the common case) the typo is invisible end to end — the
|
||||||
|
* daemon starts cleanly and authenticates nobody while the operator believes {@code token} mode
|
||||||
|
* is active. Refuse it here, unconditionally, at config load, rather than let it hide behind the
|
||||||
|
* bind check.
|
||||||
|
*
|
||||||
|
* @param yaml the raw config text
|
||||||
|
* @throws IllegalStateException when {@code auth.mode} is a non-blank value not in
|
||||||
|
* {@link #KNOWN_AUTH_MODES} (case-insensitive)
|
||||||
|
*/
|
||||||
|
static void rejectUnknownAuthMode(String yaml) {
|
||||||
|
Map<?, ?> raw;
|
||||||
|
try {
|
||||||
|
raw = YAML.readValue(yaml, Map.class);
|
||||||
|
} catch (IOException | IllegalArgumentException e) {
|
||||||
|
return; // a malformed file is reported by the real parse, not here
|
||||||
|
}
|
||||||
|
if (raw == null || !(raw.get("auth") instanceof Map<?, ?> auth)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (!(auth.get("mode") instanceof String mode) || mode.isBlank()
|
||||||
|
|| KNOWN_AUTH_MODES.contains(mode.toLowerCase())) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
throw new IllegalStateException("refusing to start: auth.mode=" + mode
|
||||||
|
+ " is not recognized — accepted values are "
|
||||||
|
+ String.join(", ", KNOWN_AUTH_MODES.stream().sorted().toList())
|
||||||
|
+ " (case-insensitive); an unrecognized mode would otherwise silently fall back to"
|
||||||
|
+ " loopback-trust, which authenticates nobody.");
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The per-profile placements this build understands — {@link Profile#placement()}'s only valid values. */
|
||||||
|
private static final Set<String> KNOWN_PLACEMENTS = Set.of("tab", "pane");
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reject a profile whose {@code placement:} is not one of {@link #KNOWN_PLACEMENTS} (CB-606),
|
||||||
|
* naming the profile, the value it set, and the accepted set.
|
||||||
|
*
|
||||||
|
* <p>{@link Profile}'s compact constructor only lower-cases {@code placement}, and
|
||||||
|
* {@link Profile#tabPlacement()} only compares the result against {@code "tab"} — anything else,
|
||||||
|
* including a typo like {@code tabb}, silently falls back to the legacy pane placement with no
|
||||||
|
* signal anywhere.
|
||||||
|
*
|
||||||
|
* @param yaml the raw config text
|
||||||
|
* @throws IllegalStateException when any profile's {@code placement} is a non-blank value not in
|
||||||
|
* {@link #KNOWN_PLACEMENTS} (case-insensitive)
|
||||||
|
*/
|
||||||
|
static void rejectUnknownPlacement(String yaml) {
|
||||||
|
Map<?, ?> raw;
|
||||||
|
try {
|
||||||
|
raw = YAML.readValue(yaml, Map.class);
|
||||||
|
} catch (IOException | IllegalArgumentException e) {
|
||||||
|
return; // a malformed file is reported by the real parse, not here
|
||||||
|
}
|
||||||
|
if (raw == null || !(raw.get("profiles") instanceof Map<?, ?> profiles)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
List<String> bad = profiles.entrySet().stream()
|
||||||
|
.filter(e -> e.getValue() instanceof Map<?, ?> p
|
||||||
|
&& p.get("placement") instanceof String pl && !pl.isBlank()
|
||||||
|
&& !KNOWN_PLACEMENTS.contains(pl.toLowerCase()))
|
||||||
|
.map(e -> String.valueOf(e.getKey()) + "=" + ((Map<?, ?>) e.getValue()).get("placement"))
|
||||||
|
.sorted()
|
||||||
|
.toList();
|
||||||
|
if (!bad.isEmpty()) {
|
||||||
|
throw new IllegalStateException("refusing to start: profile(s) [" + String.join(", ", bad)
|
||||||
|
+ "] set an unrecognized placement — accepted values are "
|
||||||
|
+ String.join(", ", KNOWN_PLACEMENTS.stream().sorted().toList())
|
||||||
|
+ " (case-insensitive); an unrecognized placement would otherwise fall back to"
|
||||||
|
+ " legacy pane placement with no signal anywhere.");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** The member-credential policies this build understands — {@link MemberCredentials#policy()}'s only valid value. */
|
||||||
|
private static final Set<String> KNOWN_MEMBER_CREDENTIALS_POLICIES =
|
||||||
|
Set.of(MemberCredentials.POLICY_DENY_BY_DEFAULT);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reject a {@code memberCredentials.policy} that is not {@link #KNOWN_MEMBER_CREDENTIALS_POLICIES}
|
||||||
|
* (CB-596), naming the value and the accepted set.
|
||||||
|
*
|
||||||
|
* <p>{@link MemberCredentials}'s compact constructor only lower-cases {@code policy} and defaults
|
||||||
|
* a blank one to {@link MemberCredentials#POLICY_DENY_BY_DEFAULT} — nothing rejects an actual
|
||||||
|
* typo like {@code deny-by-defualt}. There is only one policy today, so such a typo would
|
||||||
|
* currently behave identically to the real value by accident; the day a second policy exists
|
||||||
|
* that accident becomes a silent behavior change. Refuse it now, at config load, following the
|
||||||
|
* same pattern as {@link #rejectUnknownAuthMode} and {@link #rejectUnknownPlacement}.
|
||||||
|
*
|
||||||
|
* @param yaml the raw config text
|
||||||
|
* @throws IllegalStateException when {@code memberCredentials.policy} is a non-blank value not in
|
||||||
|
* {@link #KNOWN_MEMBER_CREDENTIALS_POLICIES} (case-insensitive)
|
||||||
|
*/
|
||||||
|
static void rejectUnknownMemberCredentialsPolicy(String yaml) {
|
||||||
|
Map<?, ?> raw;
|
||||||
|
try {
|
||||||
|
raw = YAML.readValue(yaml, Map.class);
|
||||||
|
} catch (IOException | IllegalArgumentException e) {
|
||||||
|
return; // a malformed file is reported by the real parse, not here
|
||||||
|
}
|
||||||
|
if (raw == null || !(raw.get("memberCredentials") instanceof Map<?, ?> mc)) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (!(mc.get("policy") instanceof String policy) || policy.isBlank()
|
||||||
|
|| KNOWN_MEMBER_CREDENTIALS_POLICIES.contains(policy.toLowerCase())) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
throw new IllegalStateException("refusing to start: memberCredentials.policy=" + policy
|
||||||
|
+ " is not recognized — accepted values are "
|
||||||
|
+ String.join(", ", KNOWN_MEMBER_CREDENTIALS_POLICIES.stream().sorted().toList())
|
||||||
|
+ " (case-insensitive).");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Reject a top-level {@code placement:} policy name {@link PlacementPolicies#fromName} does not
|
||||||
|
* recognize (CB-606), at config load rather than lazily at first spawn.
|
||||||
|
*
|
||||||
|
* <p>{@code CompositePeerLauncher} only calls {@link PlacementPolicies#fromName} per spawn,
|
||||||
|
* through a {@code Supplier} that re-reads live config (CB-559, so a hot-reloaded placement
|
||||||
|
* policy takes effect without a restart) — so a bad name still starts a daemon that looks
|
||||||
|
* healthy and fails only the first time something spawns without naming a profile. Every other
|
||||||
|
* field this class validates fails here, at load; this one gets the same treatment, calling
|
||||||
|
* {@link PlacementPolicies#fromName} itself as the single source of truth for what is valid
|
||||||
|
* rather than duplicating its accepted set.
|
||||||
|
*
|
||||||
|
* @param placement the raw, possibly null/blank {@code placement} value as parsed (before
|
||||||
|
* {@link #withDefaults()} runs); {@code fromName} itself treats null/blank as
|
||||||
|
* {@code fixed}, so this call changes no default
|
||||||
|
* @throws IllegalStateException when {@code placement} is a name {@link PlacementPolicies} does
|
||||||
|
* not recognize
|
||||||
|
*/
|
||||||
|
private static void rejectUnknownPlacementPolicy(String placement) {
|
||||||
|
try {
|
||||||
|
PlacementPolicies.fromName(placement);
|
||||||
|
} catch (IllegalArgumentException e) {
|
||||||
|
throw new IllegalStateException("refusing to start: " + e.getMessage(), e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
static List<String> unknownTopLevelKeys(String yaml) {
|
static List<String> unknownTopLevelKeys(String yaml) {
|
||||||
Map<?, ?> raw;
|
Map<?, ?> raw;
|
||||||
try {
|
try {
|
||||||
@@ -1320,9 +1596,17 @@ public record BridgedConfig(
|
|||||||
// config that never mentions it should still get a sane cooldown rather than a null one.
|
// config that never mentions it should still get a sane cooldown rather than a null one.
|
||||||
Integer quarantineCooldown = (quarantineCooldownSeconds != null && quarantineCooldownSeconds > 0)
|
Integer quarantineCooldown = (quarantineCooldownSeconds != null && quarantineCooldownSeconds > 0)
|
||||||
? quarantineCooldownSeconds : DEFAULT_QUARANTINE_COOLDOWN_SECONDS;
|
? quarantineCooldownSeconds : DEFAULT_QUARANTINE_COOLDOWN_SECONDS;
|
||||||
|
// memberCredentials IS defaulted, like guard/lifecycle/auth above, so no reader ever sees a
|
||||||
|
// null. CB-596: an empty MemberCredentials (empty known, empty allow) blocks NOTHING — unlike
|
||||||
|
// guard/lifecycle, an absent block is not a safe "feature off" default here, it is a gap. It
|
||||||
|
// is deliberately not pre-populated with a Java-side name list (that would just reintroduce
|
||||||
|
// the hardcoded-list defect this record replaces); the block must be configured in
|
||||||
|
// bridged.yaml to protect anything. See bridged.example.yaml's memberCredentials: comment.
|
||||||
|
MemberCredentials mc = memberCredentials != null ? memberCredentials
|
||||||
|
: new MemberCredentials(null, List.of(), List.of());
|
||||||
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
return new BridgedConfig(b, herdrSocket, profiles, g, worktreeRoot, l, timeout, pollMs,
|
||||||
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
broker, primary, f, leadHeartbeat, health, placementOrDefault, a, configReload,
|
||||||
quarantineCooldown);
|
quarantineCooldown, mc);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -577,13 +577,25 @@ public final class BridgeMcp {
|
|||||||
return text("acknowledged " + msgId);
|
return text("acknowledged " + msgId);
|
||||||
}
|
}
|
||||||
|
|
||||||
/** {@code bridge_status}: the live lifecycle status of a worker session. */
|
/**
|
||||||
|
* {@code bridge_status}: the live lifecycle status of a worker session, plus — when the worker
|
||||||
|
* is paused mid-turn in an async {@code bridge_ask} (CB-582) — the open question and how to
|
||||||
|
* answer it, so a lead on its normal poll cadence does not need the ticket to notice.
|
||||||
|
*/
|
||||||
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
|
static McpSchema.CallToolResult status(MessageService messages, String sessionId) {
|
||||||
if (isBlank(sessionId)) {
|
if (isBlank(sessionId)) {
|
||||||
return error("sessionId is required");
|
return error("sessionId is required");
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
return text(messages.status(sessionId).name().toLowerCase());
|
String base = messages.status(sessionId).name().toLowerCase();
|
||||||
|
MessageService.PendingAsk ask = messages.pendingAsk(sessionId);
|
||||||
|
if (ask == null) {
|
||||||
|
return text(base);
|
||||||
|
}
|
||||||
|
return text(base + "\n\n[question — worker is waiting for your answer]\n" + ask.question()
|
||||||
|
+ "\n\nAnswer it by calling bridge_send again with turnId=\"" + ask.turnId()
|
||||||
|
+ "\" and content set to your answer; the worker resumes the same turn."
|
||||||
|
+ " (ticket " + ask.ticket() + ")");
|
||||||
} catch (HerdrException e) {
|
} catch (HerdrException e) {
|
||||||
return error("herdr error for session " + sessionId + ": " + e.getMessage());
|
return error("herdr error for session " + sessionId + ": " + e.getMessage());
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -83,6 +83,21 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
|||||||
fleet);
|
fleet);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Production constructor, plus the CB-596 {@code memberCredentials} policy supplier.
|
||||||
|
*/
|
||||||
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
|
||||||
|
this(agents, spaces, guard, profiles, defaultProfile, env,
|
||||||
|
spawnReadyTimeoutMs,
|
||||||
|
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||||
|
fleet, memberCredentials);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Full testability constructor. Every injectable collaborator is explicit so unit tests supply
|
* Full testability constructor. Every injectable collaborator is explicit so unit tests supply
|
||||||
* fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}). The
|
* fakes for the clock ({@code nowMillis}) and poll-loop wait ({@code sleeper}). The
|
||||||
@@ -125,6 +140,39 @@ public final class ClaudeCodeLauncher extends HerdrPeerLauncher {
|
|||||||
this.guard = guard;
|
this.guard = guard;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Full testability constructor, plus the CB-596 {@code memberCredentials} policy supplier.
|
||||||
|
*/
|
||||||
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs,
|
||||||
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
|
||||||
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
|
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||||
|
this.guard = guard;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Full testability constructor, plus an injectable host-env-names source for the CB-596
|
||||||
|
* criterion-4 gap detector. Test seam only — every production call site leaves this at the
|
||||||
|
* default (the real {@code System.getenv()} key set) via the constructor above.
|
||||||
|
*/
|
||||||
|
public ClaudeCodeLauncher(AgentControl agents, WorkspaceControl spaces, SubscriptionGuard guard,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs,
|
||||||
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials,
|
||||||
|
Supplier<Set<String>> hostEnvNames) {
|
||||||
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
|
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials, hostEnvNames);
|
||||||
|
this.guard = guard;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* {@inheritDoc}
|
* {@inheritDoc}
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -20,6 +20,7 @@ import org.slf4j.LoggerFactory;
|
|||||||
import java.security.SecureRandom;
|
import java.security.SecureRandom;
|
||||||
import java.util.ArrayList;
|
import java.util.ArrayList;
|
||||||
import java.util.Collection;
|
import java.util.Collection;
|
||||||
|
import java.util.HashSet;
|
||||||
import java.util.LinkedHashMap;
|
import java.util.LinkedHashMap;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
@@ -90,6 +91,27 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
*/
|
*/
|
||||||
private final Supplier<BridgedConfig.Fleet> fleet;
|
private final Supplier<BridgedConfig.Fleet> fleet;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: the live {@code memberCredentials:} policy, read once per spawn (same hot-reload shape
|
||||||
|
* as {@link #fleet}). {@code null} — either the supplier itself, or what it returns — means no
|
||||||
|
* policy is configured and {@link #applyMemberCredentialPolicy} shadows nothing.
|
||||||
|
*/
|
||||||
|
private final Supplier<BridgedConfig.MemberCredentials> memberCredentials;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Enumerates the daemon's own process environment variable NAMES ONLY, never values — the CB-596
|
||||||
|
* criterion-4 gap detector's data source (see {@link #logCredentialGap}). Injectable for tests;
|
||||||
|
* production always resolves to the real {@code System.getenv()} key set.
|
||||||
|
*
|
||||||
|
* <p>Deliberately the daemon's own environment, not the spawned pane's: nothing in the herdr
|
||||||
|
* client surface lets the daemon read back an arbitrary command's output from a pane before the
|
||||||
|
* peer starts in it, so there is no channel to inspect the pane's environment directly. The
|
||||||
|
* daemon's own process is started the same way (a login shell sourcing the same secret store —
|
||||||
|
* see CB-592's investigation of {@code secrets.sh}), so on a single-host deployment its env
|
||||||
|
* mirrors what the pane's login shell is about to export.
|
||||||
|
*/
|
||||||
|
private final Supplier<Set<String>> hostEnvNames;
|
||||||
|
|
||||||
/** The final instruction always requires a bridge reply when the bridge MCP is mounted. */
|
/** The final instruction always requires a bridge reply when the bridge MCP is mounted. */
|
||||||
protected static final String REPLY_CHARTER =
|
protected static final String REPLY_CHARTER =
|
||||||
"You are a spawned member in the claude-bridge fleet. Every message you receive arrives "
|
"You are a spawned member in the claude-bridge fleet. Every message you receive arrives "
|
||||||
@@ -166,6 +188,41 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
long spawnReadyTimeoutMs,
|
long spawnReadyTimeoutMs,
|
||||||
LongSupplier nowMillis, Runnable sleeper,
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
Supplier<BridgedConfig.Fleet> fleet) {
|
Supplier<BridgedConfig.Fleet> fleet) {
|
||||||
|
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||||
|
nowMillis, sleeper, fleet, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* As above, plus the live {@code memberCredentials} policy (CB-596).
|
||||||
|
*
|
||||||
|
* @param memberCredentials live member-credential policy, read once per spawn; {@code null} ⇒
|
||||||
|
* no policy configured, so a spawn shadows nothing. A separate
|
||||||
|
* constructor rather than a new parameter on the one above, so every
|
||||||
|
* existing call site keeps the pre-CB-596 default without an edit.
|
||||||
|
*/
|
||||||
|
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs,
|
||||||
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
|
||||||
|
this(namePrefix, agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||||
|
nowMillis, sleeper, fleet, memberCredentials, null);
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* As above, plus an injectable {@link #hostEnvNames} source for the CB-596 gap detector. Test
|
||||||
|
* seam only — every production call site leaves this {@code null} and gets the real host env.
|
||||||
|
*/
|
||||||
|
protected HerdrPeerLauncher(String namePrefix, AgentControl agents, WorkspaceControl spaces,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs,
|
||||||
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials,
|
||||||
|
Supplier<Set<String>> hostEnvNames) {
|
||||||
this.fleet = fleet;
|
this.fleet = fleet;
|
||||||
this.namePrefix = namePrefix;
|
this.namePrefix = namePrefix;
|
||||||
this.agents = agents;
|
this.agents = agents;
|
||||||
@@ -176,6 +233,8 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
this.spawnReadyTimeoutMs = spawnReadyTimeoutMs;
|
this.spawnReadyTimeoutMs = spawnReadyTimeoutMs;
|
||||||
this.nowMillis = nowMillis;
|
this.nowMillis = nowMillis;
|
||||||
this.sleeper = sleeper;
|
this.sleeper = sleeper;
|
||||||
|
this.memberCredentials = memberCredentials;
|
||||||
|
this.hostEnvNames = hostEnvNames != null ? hostEnvNames : () -> System.getenv().keySet();
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- adapter seams -------------------------------------------------------------------------
|
// --- adapter seams -------------------------------------------------------------------------
|
||||||
@@ -767,12 +826,13 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* CB-592: overlay value that shadows the admin {@code GITEA_ACCESS_TOKEN} a herdr pane
|
* CB-596: overlay value that shadows any host credential a herdr pane otherwise inherits from
|
||||||
* otherwise inherits from herdr's own login-shell process environment (gitea issue #77).
|
* herdr's own login-shell process environment (gitea issue #82, superseding CB-592's single
|
||||||
* herdr spawns a pane from its <em>own</em> process environment and layers our map on top —
|
* hardcoded {@code GITEA_ACCESS_TOKEN} name — see {@link #applyMemberCredentialPolicy}). 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
|
* {@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
|
* the keys we put in that map, so any key we never mention passes straight through from
|
||||||
* herdr's own shell, admin token included.
|
* herdr's own shell, admin credentials included.
|
||||||
*
|
*
|
||||||
* <p>Deliberately a non-blank sentinel, not {@code ""}. Whether an empty-string overlay value
|
* <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
|
* overrides an inherited variable or is skipped as blank could not be settled by reading this
|
||||||
@@ -782,21 +842,21 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
* javadoc), and that is only demonstrated for a non-blank value, so this reuses the same,
|
* javadoc), and that is only demonstrated for a non-blank value, so this reuses the same,
|
||||||
* proven-reliable shape rather than the unverified one.
|
* 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
|
* <p><b>MEASURED ON A LIVE PANE, 2026-08-15 (CB-592): this sentinel alone does NOT hold.</b> The
|
||||||
* itself works — {@code GITEA_TOKEN} is injected here, is exported by no shell file, and does
|
* overlay itself works — {@code GITEA_TOKEN} is injected here, is exported by no shell file, and
|
||||||
* reach the pane. The sentinel loses one step later. A herdr pane runs a <em>login</em> shell,
|
* does reach the pane. The sentinel loses one step later. A herdr pane runs a <em>login</em>
|
||||||
* {@code ~/.zprofile} sources {@code ${SHARED_ENV}/tools/secrets.sh}, and that file does a plain
|
* shell, {@code ~/.zprofile} sources {@code ${SHARED_ENV}/tools/secrets.sh}, and that file does a
|
||||||
* unconditional {@code export GITEA_ACCESS_TOKEN=...}. A login shell overwrites a value already
|
* plain unconditional {@code export GITEA_ACCESS_TOKEN=...}. A login shell overwrites a value
|
||||||
* in the environment, so the real admin token is put back over this sentinel before the member
|
* already in the environment, so the real admin token is put back over this sentinel before the
|
||||||
* process ever starts. That defeat applies to <em>every</em> name {@code secrets.sh} exports,
|
* member process ever starts. That defeat applies to <em>every</em> name {@code secrets.sh}
|
||||||
* and no launcher-side overlay can win against it.
|
* 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.
|
* <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
|
* 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.
|
* 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 =
|
private static final String BLOCKED_CREDENTIAL_SENTINEL =
|
||||||
"blocked-by-bridged-cb592-see-gitea-issue-77";
|
"blocked-by-bridged-cb596-see-gitea-issue-82";
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* CB-592: marks a pane as a bridged member so a shell startup file can decline to export
|
* CB-592: marks a pane as a bridged member so a shell startup file can decline to export
|
||||||
@@ -835,11 +895,11 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
* dev.ltms.bridged.guard.SubscriptionGuard}, which is checked against the profile's
|
* dev.ltms.bridged.guard.SubscriptionGuard}, which is checked against the profile's
|
||||||
* {@code baseUrl} and nothing else.
|
* {@code baseUrl} and nothing else.
|
||||||
*
|
*
|
||||||
* <p>The CB-592 shadow and marker are put in <em>last</em>, after the profile's own
|
* <p>The CB-596 credential shadow and the CB-592 marker are put in <em>last</em>, after the
|
||||||
* {@code env:}, so no profile — present or future — can restore the admin token, or hide that
|
* profile's own {@code env:}, so no profile — present or future — can restore a blocked
|
||||||
* the pane is a member, by naming either in config. This is the one place both are applied:
|
* credential, or hide that the pane is a member, by naming either in config. This is the one
|
||||||
* every {@code buildLaunch} in every adapter calls this first, so a new profile, and a peer
|
* place both are applied: every {@code buildLaunch} in every adapter calls this first, so a new
|
||||||
* kind not yet written, gets them for free.
|
* profile, and a peer kind not yet written, gets them for free.
|
||||||
*/
|
*/
|
||||||
protected Map<String, String> baseEnv(BridgedConfig.Profile cfg) {
|
protected Map<String, String> baseEnv(BridgedConfig.Profile cfg) {
|
||||||
Map<String, String> workerEnv = new LinkedHashMap<>();
|
Map<String, String> workerEnv = new LinkedHashMap<>();
|
||||||
@@ -850,11 +910,70 @@ public abstract class HerdrPeerLauncher implements PeerLauncher {
|
|||||||
if (cfg != null && cfg.env() != null) {
|
if (cfg != null && cfg.env() != null) {
|
||||||
workerEnv.putAll(cfg.env());
|
workerEnv.putAll(cfg.env());
|
||||||
}
|
}
|
||||||
workerEnv.put("GITEA_ACCESS_TOKEN", BLOCKED_GITEA_ACCESS_TOKEN);
|
applyMemberCredentialPolicy(workerEnv);
|
||||||
workerEnv.put(MEMBER_MARKER, "1");
|
workerEnv.put(MEMBER_MARKER, "1");
|
||||||
return workerEnv;
|
return workerEnv;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: shadow every configured {@code memberCredentials.known} name that is not also
|
||||||
|
* {@code allow}-ed, replacing CB-592's single hardcoded {@code GITEA_ACCESS_TOKEN} name (gitea
|
||||||
|
* issue #82). An allow-listed name is deliberately left unmentioned here — see {@link
|
||||||
|
* #BLOCKED_CREDENTIAL_SENTINEL}'s javadoc for why an overlay entry is the only way to shadow an
|
||||||
|
* inherited value, which is exactly why an allowed name must get NO entry: any entry at all,
|
||||||
|
* blank or not, risks overriding the real value the pane needs.
|
||||||
|
*
|
||||||
|
* <p>No {@code memberCredentials} configured — the supplier is {@code null}, or it resolves to
|
||||||
|
* one whose {@code known} list is empty — shadows nothing. This is a real, config-driven gap
|
||||||
|
* (see {@link BridgedConfig.MemberCredentials}'s javadoc), not a safe default: deny-by-default
|
||||||
|
* only defends names the operator has actually enumerated in {@code known}.
|
||||||
|
*/
|
||||||
|
private void applyMemberCredentialPolicy(Map<String, String> workerEnv) {
|
||||||
|
BridgedConfig.MemberCredentials creds = memberCredentials == null ? null : memberCredentials.get();
|
||||||
|
if (creds == null) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
for (String name : creds.blockedSet()) {
|
||||||
|
workerEnv.put(name, BLOCKED_CREDENTIAL_SENTINEL);
|
||||||
|
}
|
||||||
|
logCredentialGap(creds);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Credential-shaped env var name heuristic for {@link #logCredentialGap} — case-insensitive. */
|
||||||
|
private static final Pattern CREDENTIAL_SHAPED_NAME =
|
||||||
|
Pattern.compile("(?i).*(TOKEN|SECRET|_KEY|APIKEY|PASSWORD|CREDENTIAL|AUTH).*");
|
||||||
|
|
||||||
|
/** Guards {@link #logCredentialGap} to one WARN per launcher instance, not one per spawn. */
|
||||||
|
private final AtomicBoolean credentialGapLogged = new AtomicBoolean();
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
|
||||||
|
* {@code allow} is not silently allowed — it is reported. {@link #hostEnvNames} enumerates the
|
||||||
|
* daemon's own environment (see that field's javadoc for why the daemon's env is read rather
|
||||||
|
* than the spawned pane's, which the daemon has no channel to inspect at spawn time); this logs
|
||||||
|
* every such NAME, at WARN, at most once per launcher instance — never a value, a prefix of a
|
||||||
|
* value, or a hash of a value, so the log itself cannot leak anything.
|
||||||
|
*/
|
||||||
|
private void logCredentialGap(BridgedConfig.MemberCredentials creds) {
|
||||||
|
Set<String> covered = new HashSet<>(creds.known());
|
||||||
|
covered.addAll(creds.allow());
|
||||||
|
List<String> gap = hostEnvNames.get().stream()
|
||||||
|
.filter(name -> CREDENTIAL_SHAPED_NAME.matcher(name).matches())
|
||||||
|
.filter(name -> !covered.contains(name))
|
||||||
|
.sorted()
|
||||||
|
.toList();
|
||||||
|
if (gap.isEmpty()) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
if (credentialGapLogged.compareAndSet(false, true)) {
|
||||||
|
log.warn("memberCredentials gap: {} credential-shaped env var name(s) are on neither "
|
||||||
|
+ "known: nor allow: — every member pane inherits them UNBLOCKED — {}. "
|
||||||
|
+ "Add each to memberCredentials.known (blocked by default) or .allow "
|
||||||
|
+ "(if a member legitimately needs it).",
|
||||||
|
gap.size(), gap);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/** Defensive copy of {@code argv} plus room to append launch flags. */
|
/** Defensive copy of {@code argv} plus room to append launch flags. */
|
||||||
protected static List<String> mutableArgv(List<String> argv) {
|
protected static List<String> mutableArgv(List<String> argv) {
|
||||||
return new ArrayList<>(argv);
|
return new ArrayList<>(argv);
|
||||||
|
|||||||
@@ -104,6 +104,20 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
|||||||
defaultConfigRoot(), defaultDiscoveryRoot(), fleet);
|
defaultConfigRoot(), defaultDiscoveryRoot(), fleet);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Production constructor, plus the CB-596 {@code memberCredentials} policy supplier.
|
||||||
|
*/
|
||||||
|
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs, long spawnReadyPollMs,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
|
||||||
|
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||||
|
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||||
|
defaultConfigRoot(), defaultDiscoveryRoot(), fleet, memberCredentials);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Full testability constructor. Every injectable collaborator is explicit so unit tests supply a
|
* Full testability constructor. Every injectable collaborator is explicit so unit tests supply a
|
||||||
* fake clock ({@code nowMillis}), poll-loop wait ({@code sleeper}), and a temp {@code configRoot}
|
* fake clock ({@code nowMillis}), poll-loop wait ({@code sleeper}), and a temp {@code configRoot}
|
||||||
@@ -151,6 +165,23 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
|||||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Full testability constructor, plus the CB-596 {@code memberCredentials} policy supplier.
|
||||||
|
*/
|
||||||
|
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||||
|
Map<String, BridgedConfig.Profile> profiles, String defaultProfile,
|
||||||
|
Function<String, String> env,
|
||||||
|
long spawnReadyTimeoutMs,
|
||||||
|
LongSupplier nowMillis, Runnable sleeper,
|
||||||
|
Path configRoot, Path discoveryRoot,
|
||||||
|
Supplier<BridgedConfig.Fleet> fleet,
|
||||||
|
Supplier<BridgedConfig.MemberCredentials> memberCredentials) {
|
||||||
|
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||||
|
spawnReadyTimeoutMs, nowMillis, sleeper, fleet, memberCredentials);
|
||||||
|
this.configRoot = configRoot;
|
||||||
|
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||||
|
}
|
||||||
|
|
||||||
private static Path defaultConfigRoot() {
|
private static Path defaultConfigRoot() {
|
||||||
return Path.of(System.getProperty("java.io.tmpdir"));
|
return Path.of(System.getProperty("java.io.tmpdir"));
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -164,18 +164,30 @@ public final class MessageService {
|
|||||||
|
|
||||||
/** An in-flight or finished async delegation, keyed by its ticket. */
|
/** An in-flight or finished async delegation, keyed by its ticket. */
|
||||||
private static final class Task {
|
private static final class Task {
|
||||||
|
private final String ticket;
|
||||||
private final String target;
|
private final String target;
|
||||||
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
private final CompletableFuture<Reply> future = new CompletableFuture<>();
|
||||||
private final long createdNanos;
|
private final long createdNanos;
|
||||||
private volatile Reply question;
|
private volatile Reply question;
|
||||||
private volatile String turnId;
|
private volatile String turnId;
|
||||||
|
|
||||||
private Task(String target, long createdNanos) {
|
private Task(String ticket, String target, long createdNanos) {
|
||||||
|
this.ticket = ticket;
|
||||||
this.target = target;
|
this.target = target;
|
||||||
this.createdNanos = createdNanos;
|
this.createdNanos = createdNanos;
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A worker session's currently-open {@code bridge_ask} question, surfaced so {@code bridge_status}
|
||||||
|
* can show it without the caller needing the ticket first (CB-582). Only covers async
|
||||||
|
* (fire-and-poll) delegations, which track the question on their {@link Task}; a blocking
|
||||||
|
* ({@code wait:true}) send already hands the question straight back to its own caller, so there is
|
||||||
|
* nothing hidden left for {@code bridge_status} to surface in that case.
|
||||||
|
*/
|
||||||
|
public record PendingAsk(String ticket, String question, String turnId) {
|
||||||
|
}
|
||||||
|
|
||||||
private final AgentControl agents;
|
private final AgentControl agents;
|
||||||
private final Injector injector;
|
private final Injector injector;
|
||||||
private final Rendezvous rendezvous;
|
private final Rendezvous rendezvous;
|
||||||
@@ -200,9 +212,11 @@ public final class MessageService {
|
|||||||
* Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}.
|
* Create with an explicit {@link ReplyInbox} and optional {@link ReplyPushLoop}.
|
||||||
*
|
*
|
||||||
* @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply
|
* @param pushLoop nullable — when non-null, the push loop is notified on the no-waiter reply
|
||||||
* branch ({@link #reply}) so it can nudge the primary to drain the inbox, and
|
* branch ({@link #reply}) so it can nudge the primary to drain the inbox,
|
||||||
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
|
* (CB-588) whenever an async ticket started by {@link #sendAsync} reaches a
|
||||||
* terminal phase, and whenever {@link #poll} hands a terminal ticket to its caller
|
* terminal phase, whenever {@link #poll} hands a terminal ticket to its caller,
|
||||||
|
* and (CB-582) whenever an async ticket's worker pauses mid-turn in
|
||||||
|
* {@code bridge_ask} or that pause ends (answered or lapsed)
|
||||||
*/
|
*/
|
||||||
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
|
public MessageService(AgentControl agents, Injector injector, Rendezvous rendezvous,
|
||||||
ReplyInbox inbox, ReplyPushLoop pushLoop) {
|
ReplyInbox inbox, ReplyPushLoop pushLoop) {
|
||||||
@@ -495,6 +509,15 @@ public final class MessageService {
|
|||||||
rendezvous.closeAsk(ticket.turnId());
|
rendezvous.closeAsk(ticket.turnId());
|
||||||
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
|
return new AskResult(AskOutcome.NO_WAITER, null); // no primary is blocked on this worker
|
||||||
}
|
}
|
||||||
|
// CB-582: the question just became visible via bridge_poll (Phase.ASKING) for an async
|
||||||
|
// (wait:false) delegation — nudge the lead's own pane the same way a terminal ticket does
|
||||||
|
// (CB-588), since the lead's normal poll cadence is minutes away and the reverse-rendezvous
|
||||||
|
// window (~55s, see BridgeMcp/BridgedApp) is far shorter. A blocking (wait:true) send has
|
||||||
|
// no Task and gets the question directly in its own reply, so task == null there — nothing
|
||||||
|
// to nudge.
|
||||||
|
if (task != null && pushLoop != null) {
|
||||||
|
pushLoop.onQuestionOpened(task.ticket, workerSession, ticket.turnId(), question);
|
||||||
|
}
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
|
String answer = ticket.answer().get(timeoutMillis, TimeUnit.MILLISECONDS);
|
||||||
@@ -513,6 +536,14 @@ public final class MessageService {
|
|||||||
// Only the fresh owner tears down the shared turn; a duplicate must leave it open.
|
// Only the fresh owner tears down the shared turn; a duplicate must leave it open.
|
||||||
if (ticket.fresh()) {
|
if (ticket.fresh()) {
|
||||||
rendezvous.closeAsk(ticket.turnId());
|
rendezvous.closeAsk(ticket.turnId());
|
||||||
|
// CB-582: tear the push loop's copy down at the same point, not only on the three
|
||||||
|
// paths that call clearAsyncQuestion. The answer future can complete exceptionally
|
||||||
|
// (ExecutionException) or the thread be interrupted, and both leave this method by
|
||||||
|
// throwing — the question would stay pending forever, keep being named in nudges
|
||||||
|
// until its own cap, and never be removed from the map. Already-closed is a no-op.
|
||||||
|
if (pushLoop != null) {
|
||||||
|
pushLoop.questionClosed(ticket.turnId());
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -589,7 +620,7 @@ public final class MessageService {
|
|||||||
*/
|
*/
|
||||||
public String sendAsync(String target, String content, Runnable onAccepted) {
|
public String sendAsync(String target, String content, Runnable onAccepted) {
|
||||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||||
Task task = new Task(target, nowNanos.getAsLong());
|
Task task = new Task(ticket, target, nowNanos.getAsLong());
|
||||||
tasks.put(ticket, task);
|
tasks.put(ticket, task);
|
||||||
if (pushLoop != null) {
|
if (pushLoop != null) {
|
||||||
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
// CB-588: task.future only ever completes on a terminal phase (DONE or a failure) — a
|
||||||
@@ -715,6 +746,12 @@ public final class MessageService {
|
|||||||
|
|
||||||
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
|
/** Clear an answered or lapsed question, but only when it matches the ticket's current turn. */
|
||||||
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
|
private void clearAsyncQuestion(String turnId, boolean forgetTurn) {
|
||||||
|
// CB-582: tell the push loop first — like ticketCollected, a removal for a turnId it never
|
||||||
|
// nudged about (or already dropped) is a harmless no-op, so this is safe to call unconditionally
|
||||||
|
// rather than threading the guard below through it.
|
||||||
|
if (pushLoop != null) {
|
||||||
|
pushLoop.questionClosed(turnId);
|
||||||
|
}
|
||||||
Task task = asyncTasksByTurn.get(turnId);
|
Task task = asyncTasksByTurn.get(turnId);
|
||||||
if (task != null && turnId.equals(task.turnId)) {
|
if (task != null && turnId.equals(task.turnId)) {
|
||||||
task.question = null;
|
task.question = null;
|
||||||
@@ -746,6 +783,23 @@ public final class MessageService {
|
|||||||
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
return asyncTasksByTurn.values().stream().anyMatch(task -> target.equals(task.target));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The question {@code workerSession} is currently paused on via {@code bridge_ask}, if any
|
||||||
|
* (CB-582) — {@code bridge_status} uses this to show a pending question without the caller
|
||||||
|
* needing the ticket. {@code null} when the session has no open async question (including a
|
||||||
|
* session mid a <em>blocking</em> {@code bridge_ask}, which has no {@link Task} to look up — see
|
||||||
|
* {@link PendingAsk}).
|
||||||
|
*/
|
||||||
|
public PendingAsk pendingAsk(String workerSession) {
|
||||||
|
for (Task task : tasks.values()) {
|
||||||
|
Reply q = task.question;
|
||||||
|
if (q != null && workerSession.equals(task.target)) {
|
||||||
|
return new PendingAsk(task.ticket, q.text(), q.turnId());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return null;
|
||||||
|
}
|
||||||
|
|
||||||
/** Release the async executor. */
|
/** Release the async executor. */
|
||||||
public void close() {
|
public void close() {
|
||||||
asyncExecutor.shutdown();
|
asyncExecutor.shutdown();
|
||||||
|
|||||||
@@ -19,30 +19,33 @@ import java.util.stream.Collectors;
|
|||||||
|
|
||||||
/**
|
/**
|
||||||
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
|
* A status-gated push loop that nudges a lead's own herdr pane when it has uncollected work
|
||||||
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), or an
|
* waiting: a worker reply queued with no live {@code bridge_send} to resolve it (CB-307), an
|
||||||
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
|
* async delegation ticket ({@code bridge_send(wait:false)}) that reached a terminal phase
|
||||||
* (CB-588).
|
* (CB-588), or an async ticket's worker pausing mid-turn in {@code bridge_ask} to await an answer
|
||||||
|
* (CB-582).
|
||||||
*
|
*
|
||||||
* <p><strong>CB-590: one schedule per lead.</strong> Both kinds of work are triggered through
|
* <p><strong>CB-590: one schedule per lead.</strong> All three kinds of work are triggered
|
||||||
* their own entry point — {@link #onReplyQueued(String)} and
|
* through their own entry point — {@link #onReplyQueued(String)},
|
||||||
* {@link #onTicketTerminal(String, String, boolean)} — but both resolve the lead that should be
|
* {@link #onTicketTerminal(String, String, boolean)}, and
|
||||||
* nudged and coalesce onto a single per-lead reminder schedule, tracked in {@link #activeLeads}.
|
* {@link #onQuestionOpened(String, String, String, String)} — but each resolves the lead that
|
||||||
* Earlier this was two independent schedules (one keyed by worker target for replies, one keyed
|
* should be nudged and coalesces onto a single per-lead reminder schedule, tracked in
|
||||||
* by lead for tickets) that could both decide to inject into the same pane in the same window —
|
* {@link #activeLeads}. Earlier this was two independent schedules (one keyed by worker target
|
||||||
* a race, not routine behaviour, but the expensive kind: it interrupts the lead's live turn
|
* for replies, one keyed by lead for tickets) that could both decide to inject into the same pane
|
||||||
* twice. Collapsing to one schedule per lead makes that structurally impossible: at most one
|
* in the same window — a race, not routine behaviour, but the expensive kind: it interrupts the
|
||||||
* scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
|
* lead's live turn twice. Collapsing to one schedule per lead makes that structurally impossible:
|
||||||
|
* at most one scheduled tick chain is ever live for a given lead (guarded by {@link #activeLeads}'
|
||||||
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
|
* compare-and-set), so at most one {@code agents.send} to that lead's pane is ever in flight.
|
||||||
*
|
*
|
||||||
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
|
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
|
||||||
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
|
* still holds an unacked message ({@link #pendingReplies}), tickets not yet collected
|
||||||
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
|
* ({@link #pendingTickets}), and open questions not yet answered or lapsed
|
||||||
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
|
* ({@link #pendingQuestions}) — and sends at most one combined nudge per tick
|
||||||
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
|
* ({@link #injectNudge(String, int, int, int)}). Work that arrives while the lead is busy is
|
||||||
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
|
* never lost: it is re-read fresh on every tick until the lead is injectable or its own reminder
|
||||||
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
|
* cap ({@link #maxReminders}) is reached — each source spends from its own budget, so one source
|
||||||
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
|
* exhausting its cap does not stop nudges about the others (post-CB-590 regression fix; see
|
||||||
* via {@code STOP}.
|
* {@link #decide}) — whichever the durable inbox / pending set doesn't already answer via
|
||||||
|
* {@code STOP}.
|
||||||
*/
|
*/
|
||||||
public final class ReplyPushLoop {
|
public final class ReplyPushLoop {
|
||||||
|
|
||||||
@@ -57,6 +60,14 @@ public final class ReplyPushLoop {
|
|||||||
/** Coalesced form, several uncollected tickets for the same lead. */
|
/** Coalesced form, several uncollected tickets for the same lead. */
|
||||||
static final String TICKETS_NUDGE_FORMAT =
|
static final String TICKETS_NUDGE_FORMAT =
|
||||||
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
|
"%d tickets finished%s — run bridge_poll(ticket=...) for each to collect them: %s";
|
||||||
|
/** Singular form, one worker paused mid-turn in bridge_ask (CB-582) — names the answer call directly. */
|
||||||
|
static final String QUESTION_NUDGE_FORMAT =
|
||||||
|
"Worker %s asked a question (ticket %s) — answer it with bridge_send(turnId=\"%s\", "
|
||||||
|
+ "content=...) to resume its turn:\n%s";
|
||||||
|
/** Coalesced form, several open questions for the same lead. */
|
||||||
|
static final String QUESTIONS_NUDGE_FORMAT =
|
||||||
|
"%d workers are paused on a question — run bridge_poll(ticket=...) for each, then answer "
|
||||||
|
+ "with bridge_send(turnId=..., content=...): %s";
|
||||||
|
|
||||||
private final PrimaryRegistry primaryRegistry;
|
private final PrimaryRegistry primaryRegistry;
|
||||||
private final AgentControl agents;
|
private final AgentControl agents;
|
||||||
@@ -66,11 +77,23 @@ public final class ReplyPushLoop {
|
|||||||
private final long backoffMs;
|
private final long backoffMs;
|
||||||
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
|
private final Metrics metrics; // CB-512: nullable — no registry in unit tests
|
||||||
|
|
||||||
/** Worker targets with a reply queued, and the lead to nudge about it, keyed by target. */
|
/**
|
||||||
private final ConcurrentHashMap<String, String> pendingReplies = new ConcurrentHashMap<>();
|
* Worker targets with a reply queued, keyed by target. Each entry carries its own nudge
|
||||||
|
* count (CB-598) rather than sharing one counter per lead per source: a target's count only
|
||||||
|
* ever reflects nudges that actually named that target, so a target that joins while the
|
||||||
|
* schedule is already deep into another target's reminders still reads as fresh.
|
||||||
|
*/
|
||||||
|
private final ConcurrentHashMap<String, ReplyEntry> pendingReplies = new ConcurrentHashMap<>();
|
||||||
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
/** Tickets that have gone terminal but not yet been polled, keyed by ticket. */
|
||||||
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
|
private final ConcurrentHashMap<String, PendingTicket> pendingTickets = new ConcurrentHashMap<>();
|
||||||
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets). */
|
/**
|
||||||
|
* Open {@code bridge_ask} questions not yet answered or lapsed, keyed by {@code turnId}
|
||||||
|
* (CB-582). A question's own nudge count is tracked the same per-item way as
|
||||||
|
* {@link #pendingTickets} (CB-598): a fresh question keeps its source eligible regardless of
|
||||||
|
* how depleted an older, still-open question's count is.
|
||||||
|
*/
|
||||||
|
private final ConcurrentHashMap<String, PendingQuestion> pendingQuestions = new ConcurrentHashMap<>();
|
||||||
|
/** CB-590: leads with an active combined reminder schedule (replies and/or tickets and/or questions). */
|
||||||
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
|
private final ConcurrentHashMap<String, Boolean> activeLeads = new ConcurrentHashMap<>();
|
||||||
|
|
||||||
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
public ReplyPushLoop(PrimaryRegistry primaryRegistry, AgentControl agents, ReplyInbox inbox,
|
||||||
@@ -116,10 +139,10 @@ public final class ReplyPushLoop {
|
|||||||
Set<String> result = new HashSet<>();
|
Set<String> result = new HashSet<>();
|
||||||
for (var entry : pendingReplies.entrySet()) {
|
for (var entry : pendingReplies.entrySet()) {
|
||||||
String target = entry.getKey();
|
String target = entry.getKey();
|
||||||
String owningLead = entry.getValue();
|
ReplyEntry owning = entry.getValue();
|
||||||
if (!lead.equals(owningLead)) continue;
|
if (!lead.equals(owning.lead())) continue;
|
||||||
if (inbox.peek(target).isEmpty()) {
|
if (inbox.peek(target).isEmpty()) {
|
||||||
pendingReplies.remove(target, owningLead);
|
pendingReplies.remove(target, owning);
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
result.add(target);
|
result.add(target);
|
||||||
@@ -127,8 +150,15 @@ public final class ReplyPushLoop {
|
|||||||
return result;
|
return result;
|
||||||
}
|
}
|
||||||
|
|
||||||
/** A ticket awaiting collection: which lead to nudge, and whether it ended in failure. */
|
/** A pending reply target: which lead to nudge, and how many nudges have named it so far. */
|
||||||
private record PendingTicket(String ticket, String lead, boolean failed) {
|
private record ReplyEntry(String lead, int nudgeCount) {
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A ticket awaiting collection: which lead to nudge, whether it ended in failure, and how
|
||||||
|
* many nudges have named it so far (CB-598 — tracked per ticket, not per lead per source).
|
||||||
|
*/
|
||||||
|
private record PendingTicket(String ticket, String lead, boolean failed, int nudgeCount) {
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
|
/** Tickets still pending for {@code lead}, snapshotted fresh for one tick. */
|
||||||
@@ -142,6 +172,70 @@ public final class ReplyPushLoop {
|
|||||||
.collect(Collectors.toUnmodifiableSet());
|
.collect(Collectors.toUnmodifiableSet());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* An open question awaiting the lead's answer: which ticket it belongs to, which worker asked,
|
||||||
|
* which lead to nudge, the question text, and how many nudges have named it so far (CB-598 —
|
||||||
|
* tracked per question, not per lead per source).
|
||||||
|
*/
|
||||||
|
private record PendingQuestion(String turnId, String ticket, String target, String lead,
|
||||||
|
String question, int nudgeCount) {
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Questions still open for {@code lead}, snapshotted fresh for one tick. */
|
||||||
|
private List<PendingQuestion> pendingQuestionsFor(String lead) {
|
||||||
|
return pendingQuestions.values().stream().filter(q -> lead.equals(q.lead())).toList();
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Question turnIds still open for {@code lead} — a plain snapshot for race comparison. */
|
||||||
|
private Set<String> pendingQuestionTurnIdsFor(String lead) {
|
||||||
|
return pendingQuestionsFor(lead).stream().map(PendingQuestion::turnId)
|
||||||
|
.collect(Collectors.toUnmodifiableSet());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The reply-source reminder count {@link #decide} should see for {@code lead} on this tick:
|
||||||
|
* the <em>minimum</em> nudge count among the reply targets currently pending for it (CB-598).
|
||||||
|
*
|
||||||
|
* <p>Before this, the count passed to {@code decide} was a single counter carried forward
|
||||||
|
* across scheduled ticks ({@code scheduleNext(lead, count + 1, ...)}), incremented whenever
|
||||||
|
* the source had <em>any</em> pending work — not tied to which target that work was. A target
|
||||||
|
* that joined while an older target's count was already near the cap inherited that count on
|
||||||
|
* its very next tick, even though no nudge had ever named it. Taking the minimum over what is
|
||||||
|
* actually pending now means a fresh target (count 0) keeps the source eligible regardless of
|
||||||
|
* how many times an older, still-undrained target has already been nudged; that older target
|
||||||
|
* keeps riding along in the combined nudge text without spending any more of its own budget
|
||||||
|
* (see {@link #bumpNudgeCounts}). Returns 0 when nothing is pending — {@link #decide} never
|
||||||
|
* consults the count in that case, since {@code hasReplyWork} is false.
|
||||||
|
*/
|
||||||
|
private int minReplyNudgeCountFor(String lead) {
|
||||||
|
int min = Integer.MAX_VALUE;
|
||||||
|
for (String target : pendingReplyTargetsFor(lead)) {
|
||||||
|
ReplyEntry entry = pendingReplies.get(target);
|
||||||
|
if (entry != null) {
|
||||||
|
min = Math.min(min, entry.nudgeCount());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return min == Integer.MAX_VALUE ? 0 : min;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** As {@link #minReplyNudgeCountFor}, for the ticket source. */
|
||||||
|
private int minTicketNudgeCountFor(String lead) {
|
||||||
|
int min = Integer.MAX_VALUE;
|
||||||
|
for (PendingTicket ticket : pendingTicketsFor(lead)) {
|
||||||
|
min = Math.min(min, ticket.nudgeCount());
|
||||||
|
}
|
||||||
|
return min == Integer.MAX_VALUE ? 0 : min;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** As {@link #minReplyNudgeCountFor}, for the question source (CB-582). */
|
||||||
|
private int minQuestionNudgeCountFor(String lead) {
|
||||||
|
int min = Integer.MAX_VALUE;
|
||||||
|
for (PendingQuestion q : pendingQuestionsFor(lead)) {
|
||||||
|
min = Math.min(min, q.nudgeCount());
|
||||||
|
}
|
||||||
|
return min == Integer.MAX_VALUE ? 0 : min;
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
||||||
* tickets alike — and return what the loop should do.
|
* tickets alike — and return what the loop should do.
|
||||||
@@ -155,21 +249,40 @@ public final class ReplyPushLoop {
|
|||||||
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
|
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
|
||||||
* {@link Action#STOP}.
|
* {@link Action#STOP}.
|
||||||
*
|
*
|
||||||
|
* <p><strong>CB-598: the counts are per-item, not per-tick.</strong> {@link #tick} no longer
|
||||||
|
* carries these counts forward across scheduled calls — it recomputes them fresh every tick via
|
||||||
|
* {@link #minReplyNudgeCountFor} / {@link #minTicketNudgeCountFor}, so this function itself did
|
||||||
|
* not need to change; only what its caller feeds it did.
|
||||||
|
*
|
||||||
* @param lead the lead terminal to nudge
|
* @param lead the lead terminal to nudge
|
||||||
* @param replyReminderCount how many nudges have covered pending reply work for this lead
|
* @param replyReminderCount the lowest nudge count among reply targets pending for this lead
|
||||||
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
|
* @param ticketReminderCount the lowest nudge count among tickets pending for this lead
|
||||||
* @return the action the caller should take
|
* @return the action the caller should take
|
||||||
*/
|
*/
|
||||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||||
|
return decide(lead, replyReminderCount, ticketReminderCount, minQuestionNudgeCountFor(lead));
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* As {@link #decide(String, int, int)}, with the question source (CB-582) folded in on the
|
||||||
|
* same footing as replies and tickets: its own eligibility (has open questions AND under its
|
||||||
|
* own {@link #maxReminders} budget) is enough on its own to {@link Action#INJECT}, exactly like
|
||||||
|
* the other two.
|
||||||
|
*
|
||||||
|
* @param questionReminderCount the lowest nudge count among questions open for this lead
|
||||||
|
*/
|
||||||
|
Action decide(String lead, int replyReminderCount, int ticketReminderCount, int questionReminderCount) {
|
||||||
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
||||||
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
||||||
if (!hasReplyWork && !hasTicketWork) {
|
boolean hasQuestionWork = !pendingQuestionTurnIdsFor(lead).isEmpty();
|
||||||
|
if (!hasReplyWork && !hasTicketWork && !hasQuestionWork) {
|
||||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||||
return Action.STOP;
|
return Action.STOP;
|
||||||
}
|
}
|
||||||
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
||||||
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
||||||
if (!replyEligible && !ticketEligible) {
|
boolean questionEligible = hasQuestionWork && questionReminderCount < maxReminders;
|
||||||
|
if (!replyEligible && !ticketEligible && !questionEligible) {
|
||||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||||
maxReminders, lead);
|
maxReminders, lead);
|
||||||
countNudge("exhausted");
|
countNudge("exhausted");
|
||||||
@@ -204,7 +317,8 @@ public final class ReplyPushLoop {
|
|||||||
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
|
log.debug("push: no lead is known to be waiting on {}, skipping reminder", target);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
pendingReplies.put(target, lead.get());
|
pendingReplies.compute(target, (t, existing) ->
|
||||||
|
new ReplyEntry(lead.get(), existing == null ? 0 : existing.nudgeCount()));
|
||||||
startOrCoalesce(lead.get());
|
startOrCoalesce(lead.get());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -232,7 +346,8 @@ public final class ReplyPushLoop {
|
|||||||
ticket, target);
|
ticket, target);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
pendingTickets.put(ticket, new PendingTicket(ticket, lead.get(), failed));
|
pendingTickets.compute(ticket, (id, existing) ->
|
||||||
|
new PendingTicket(ticket, lead.get(), failed, existing == null ? 0 : existing.nudgeCount()));
|
||||||
startOrCoalesce(lead.get());
|
startOrCoalesce(lead.get());
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -246,6 +361,39 @@ public final class ReplyPushLoop {
|
|||||||
pendingTickets.remove(ticket);
|
pendingTickets.remove(ticket);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Called when an async ticket's worker pauses mid-turn in {@code bridge_ask} (CB-582): the
|
||||||
|
* question is now visible via {@code bridge_poll} (Phase.ASKING), but the reverse-rendezvous
|
||||||
|
* window it opened with (~55s default, see {@code BridgeMcp}/{@code BridgedApp}) is far shorter
|
||||||
|
* than a lead's normal minutes-long poll cadence — exactly the gap this closes. Resolves the
|
||||||
|
* delegating lead the same way {@link #onTicketTerminal} does and coalesces onto the same
|
||||||
|
* per-lead schedule (CB-590).
|
||||||
|
*
|
||||||
|
* @param ticket the async ticket the question belongs to (for {@code bridge_poll})
|
||||||
|
* @param target the worker session that asked
|
||||||
|
* @param turnId correlation id the lead answers with ({@code bridge_send turnId=...})
|
||||||
|
* @param question the question text
|
||||||
|
*/
|
||||||
|
public void onQuestionOpened(String ticket, String target, String turnId, String question) {
|
||||||
|
var lead = primaryRegistry.nudgeTargetFor(target);
|
||||||
|
if (lead.isEmpty()) {
|
||||||
|
log.debug("push: no lead is known to be waiting on {}'s question (turnId {}), skipping nudge",
|
||||||
|
target, turnId);
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
pendingQuestions.put(turnId, new PendingQuestion(turnId, ticket, target, lead.get(), question, 0));
|
||||||
|
startOrCoalesce(lead.get());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Called when a worker's {@code bridge_ask} resolves — answered or lapsed unanswered — so a
|
||||||
|
* scheduled tick never nudges about a question the lead already handled. A {@code turnId} that
|
||||||
|
* was never pending (never nudged, or already closed) is a no-op.
|
||||||
|
*/
|
||||||
|
public void questionClosed(String turnId) {
|
||||||
|
pendingQuestions.remove(turnId);
|
||||||
|
}
|
||||||
|
|
||||||
// --- the schedule ----------------------------------------------------------------------------
|
// --- the schedule ----------------------------------------------------------------------------
|
||||||
|
|
||||||
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
/** Start a reminder schedule for {@code lead}, or join the one already running. */
|
||||||
@@ -255,29 +403,40 @@ public final class ReplyPushLoop {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
log.debug("push: starting reminder loop for lead {}", lead);
|
log.debug("push: starting reminder loop for lead {}", lead);
|
||||||
scheduleNext(lead, 0, 0);
|
scheduleNext(lead);
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Execute one loop tick — called on the scheduler thread. */
|
/**
|
||||||
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
|
* Execute one loop tick — called on the scheduler thread (or directly by a test; package-private
|
||||||
|
* for the same reason as {@link #stopOrRestart}).
|
||||||
|
*
|
||||||
|
* <p><strong>CB-598.</strong> The reminder counts fed into {@link #decide} are recomputed fresh
|
||||||
|
* every tick from what is actually pending right now ({@link #minReplyNudgeCountFor} /
|
||||||
|
* {@link #minTicketNudgeCountFor}), rather than carried forward as running counters across
|
||||||
|
* scheduled calls. A counter carried forward has no memory of which item it was counting for:
|
||||||
|
* a target or ticket that joined mid-backoff — after the previous tick fired but before this one
|
||||||
|
* did — is already sitting in {@code repliesBefore} / {@code ticketsBefore} below by the time this
|
||||||
|
* tick takes its snapshot, indistinguishable at that point from backlog the cap is meant to
|
||||||
|
* silence. Recomputing from the per-item counts fixes that: a newly-joined item's own count is
|
||||||
|
* still 0, so it keeps its source eligible regardless of how depleted an older, still-undrained
|
||||||
|
* item's count is.
|
||||||
|
*/
|
||||||
|
void tick(String lead) {
|
||||||
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
||||||
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
||||||
var action = decide(lead, replyReminderCount, ticketReminderCount);
|
Set<String> questionsBefore = pendingQuestionTurnIdsFor(lead);
|
||||||
|
int replyReminderCount = minReplyNudgeCountFor(lead);
|
||||||
|
int ticketReminderCount = minTicketNudgeCountFor(lead);
|
||||||
|
int questionReminderCount = minQuestionNudgeCountFor(lead);
|
||||||
|
var action = decide(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||||
switch (action) {
|
switch (action) {
|
||||||
case INJECT -> {
|
case INJECT -> {
|
||||||
injectNudge(lead, replyReminderCount, ticketReminderCount);
|
injectNudge(lead, replyReminderCount, ticketReminderCount, questionReminderCount);
|
||||||
// Only the source(s) actually eligible this tick spend a unit of their own budget —
|
scheduleNext(lead);
|
||||||
// an exhausted source riding along in the combined message (still pending, still
|
|
||||||
// named) does not get charged again; its count stays put until it drains.
|
|
||||||
boolean replyEligible = !repliesBefore.isEmpty() && replyReminderCount < maxReminders;
|
|
||||||
boolean ticketEligible = !ticketsBefore.isEmpty() && ticketReminderCount < maxReminders;
|
|
||||||
scheduleNext(lead,
|
|
||||||
replyEligible ? replyReminderCount + 1 : replyReminderCount,
|
|
||||||
ticketEligible ? ticketReminderCount + 1 : ticketReminderCount);
|
|
||||||
}
|
}
|
||||||
// Re-check after the configured backoff; the lead may become injectable soon.
|
// Re-check after the configured backoff; the lead may become injectable soon.
|
||||||
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
|
case WAIT_BUSY -> scheduleNext(lead);
|
||||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore, questionsBefore);
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -312,50 +471,92 @@ public final class ReplyPushLoop {
|
|||||||
* side always wins; neither can miss the other, so this never loops on its own account.
|
* side always wins; neither can miss the other, so this never loops on its own account.
|
||||||
*/
|
*/
|
||||||
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
|
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore) {
|
||||||
|
stopOrRestart(lead, repliesBefore, ticketsBefore, pendingQuestionTurnIdsFor(lead));
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* As {@link #stopOrRestart(String, Set, Set)}, with the question source's (CB-582) own "before"
|
||||||
|
* snapshot folded into the same race check: a question that raced in during the
|
||||||
|
* decision-to-release window reclaims the schedule slot exactly like a raced-in reply or ticket.
|
||||||
|
*/
|
||||||
|
void stopOrRestart(String lead, Set<String> repliesBefore, Set<String> ticketsBefore,
|
||||||
|
Set<String> questionsBefore) {
|
||||||
activeLeads.remove(lead);
|
activeLeads.remove(lead);
|
||||||
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|
boolean racedIn = pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))
|
||||||
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
|
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t))
|
||||||
|
|| pendingQuestionTurnIdsFor(lead).stream().anyMatch(t -> !questionsBefore.contains(t));
|
||||||
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
||||||
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
||||||
scheduleNext(lead, 0, 0);
|
scheduleNext(lead);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
log.debug("push: reminder loop ended for lead {}", lead);
|
log.debug("push: reminder loop ended for lead {}", lead);
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||||
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
|
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount,
|
||||||
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
|
int questionReminderCount) {
|
||||||
// collected (or another arrive), between the decision and the injection.
|
// Re-read rather than threading it down from decide(): a reply can drain, a ticket be
|
||||||
|
// collected, or a question be answered (or another arrive), between the decision and the
|
||||||
|
// injection.
|
||||||
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
||||||
List<PendingTicket> tickets = pendingTicketsFor(lead);
|
List<PendingTicket> tickets = pendingTicketsFor(lead);
|
||||||
if (replyTargets.isEmpty() && tickets.isEmpty()) {
|
List<PendingQuestion> questions = pendingQuestionsFor(lead);
|
||||||
|
if (replyTargets.isEmpty() && tickets.isEmpty() && questions.isEmpty()) {
|
||||||
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
log.debug("push: pending work for lead {} drained before the nudge could be sent", lead);
|
||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
String nudge = formatNudge(replyTargets, tickets);
|
String nudge = formatNudge(replyTargets, tickets, questions);
|
||||||
try {
|
try {
|
||||||
agents.send(lead, nudge);
|
agents.send(lead, nudge);
|
||||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
|
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}, question {}/{}; "
|
||||||
|
+ "{} reply target(s), {} ticket(s), {} question(s))",
|
||||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||||
replyTargets.size(), tickets.size());
|
questionReminderCount + 1, maxReminders,
|
||||||
|
replyTargets.size(), tickets.size(), questions.size());
|
||||||
countNudge("delivered");
|
countNudge("delivered");
|
||||||
} catch (RuntimeException e) {
|
} catch (RuntimeException e) {
|
||||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
|
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}, question {}/{}): {}",
|
||||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
|
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||||
|
questionReminderCount + 1, maxReminders, e.toString());
|
||||||
|
}
|
||||||
|
// Bump every item actually named in this nudge, not just whatever the shared source-level
|
||||||
|
// eligibility used to gate (CB-598) — each item's own count is what the next tick's
|
||||||
|
// minReplyNudgeCountFor / minTicketNudgeCountFor / minQuestionNudgeCountFor will read. An
|
||||||
|
// item already at or over the cap keeps riding along in the text (still pending, still
|
||||||
|
// named) but its extra bumps here are inert: decide() already treats it as ineligible once
|
||||||
|
// its count reaches maxReminders.
|
||||||
|
bumpNudgeCounts(replyTargets, tickets, questions);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Record that every one of these items was just named in a sent (or attempted) nudge. */
|
||||||
|
private void bumpNudgeCounts(Set<String> replyTargets, List<PendingTicket> tickets,
|
||||||
|
List<PendingQuestion> questions) {
|
||||||
|
for (String target : replyTargets) {
|
||||||
|
pendingReplies.computeIfPresent(target, (t, e) -> new ReplyEntry(e.lead(), e.nudgeCount() + 1));
|
||||||
|
}
|
||||||
|
for (PendingTicket ticket : tickets) {
|
||||||
|
pendingTickets.computeIfPresent(ticket.ticket(),
|
||||||
|
(id, e) -> new PendingTicket(e.ticket(), e.lead(), e.failed(), e.nudgeCount() + 1));
|
||||||
|
}
|
||||||
|
for (PendingQuestion question : questions) {
|
||||||
|
pendingQuestions.computeIfPresent(question.turnId(), (id, e) ->
|
||||||
|
new PendingQuestion(e.turnId(), e.ticket(), e.target(), e.lead(), e.question(),
|
||||||
|
e.nudgeCount() + 1));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
/** Schedule the next tick on the scheduler thread pool. */
|
/** Schedule the next tick on the scheduler thread pool. */
|
||||||
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
|
private void scheduleNext(String lead) {
|
||||||
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
|
scheduler.schedule(() -> tick(lead),
|
||||||
backoffMs, TimeUnit.MILLISECONDS);
|
backoffMs, TimeUnit.MILLISECONDS);
|
||||||
}
|
}
|
||||||
|
|
||||||
// --- nudge formatting ------------------------------------------------------------------------
|
// --- nudge formatting ------------------------------------------------------------------------
|
||||||
|
|
||||||
/** Render everything pending for one lead as a single nudge line. */
|
/** Render everything pending for one lead as a single nudge line. */
|
||||||
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets) {
|
private static String formatNudge(Set<String> replyTargets, List<PendingTicket> tickets,
|
||||||
|
List<PendingQuestion> questions) {
|
||||||
List<String> parts = new ArrayList<>();
|
List<String> parts = new ArrayList<>();
|
||||||
if (!replyTargets.isEmpty()) {
|
if (!replyTargets.isEmpty()) {
|
||||||
parts.add(formatRepliesNudge(replyTargets));
|
parts.add(formatRepliesNudge(replyTargets));
|
||||||
@@ -363,6 +564,9 @@ public final class ReplyPushLoop {
|
|||||||
if (!tickets.isEmpty()) {
|
if (!tickets.isEmpty()) {
|
||||||
parts.add(formatTicketsNudge(tickets));
|
parts.add(formatTicketsNudge(tickets));
|
||||||
}
|
}
|
||||||
|
if (!questions.isEmpty()) {
|
||||||
|
parts.add(formatQuestionsNudge(questions));
|
||||||
|
}
|
||||||
return String.join(" | ", parts);
|
return String.join(" | ", parts);
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -390,6 +594,18 @@ public final class ReplyPushLoop {
|
|||||||
return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids);
|
return TICKETS_NUDGE_FORMAT.formatted(pending.size(), failedNote, ids);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Render one or several open questions (CB-582). */
|
||||||
|
private static String formatQuestionsNudge(List<PendingQuestion> pending) {
|
||||||
|
if (pending.size() == 1) {
|
||||||
|
PendingQuestion q = pending.get(0);
|
||||||
|
return QUESTION_NUDGE_FORMAT.formatted(q.target(), q.ticket(), q.turnId(), q.question());
|
||||||
|
}
|
||||||
|
String ids = pending.stream()
|
||||||
|
.map(q -> q.ticket() + " (turnId=" + q.turnId() + ")")
|
||||||
|
.collect(Collectors.joining(", "));
|
||||||
|
return QUESTIONS_NUDGE_FORMAT.formatted(pending.size(), ids);
|
||||||
|
}
|
||||||
|
|
||||||
// --- lifecycle -----------------------------------------------------------------------------
|
// --- lifecycle -----------------------------------------------------------------------------
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -397,8 +613,8 @@ public final class ReplyPushLoop {
|
|||||||
* uses this to stand aside: while the push loop is actively nudging a lead, a concurrent
|
* uses this to stand aside: while the push loop is actively nudging a lead, a concurrent
|
||||||
* heartbeat injection would start a second competing turn in the same pane — racing loops
|
* heartbeat injection would start a second competing turn in the same pane — racing loops
|
||||||
* multiply turns and context burn. "Active" means a schedule exists in {@link #activeLeads},
|
* multiply turns and context burn. "Active" means a schedule exists in {@link #activeLeads},
|
||||||
* which now covers both reply-queued (CB-307) and ticket-terminal (CB-588) work (CB-590) —
|
* which now covers reply-queued (CB-307), ticket-terminal (CB-588), and question-open (CB-582)
|
||||||
* bounded by what has been triggered, not by any persistent state.
|
* work (CB-590) — bounded by what has been triggered, not by any persistent state.
|
||||||
*/
|
*/
|
||||||
public boolean isActive() {
|
public boolean isActive() {
|
||||||
return !activeLeads.isEmpty();
|
return !activeLeads.isEmpty();
|
||||||
@@ -410,6 +626,7 @@ public final class ReplyPushLoop {
|
|||||||
activeLeads.clear();
|
activeLeads.clear();
|
||||||
pendingReplies.clear();
|
pendingReplies.clear();
|
||||||
pendingTickets.clear();
|
pendingTickets.clear();
|
||||||
|
pendingQuestions.clear();
|
||||||
}
|
}
|
||||||
|
|
||||||
/** @see #stop() */
|
/** @see #stop() */
|
||||||
|
|||||||
@@ -235,7 +235,14 @@ public final class BridgedApp {
|
|||||||
List<Map<String, Object>> out = sessions.roster().stream()
|
List<Map<String, Object>> out = sessions.roster().stream()
|
||||||
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
.map(s -> SessionManager.rosterView(s, live.get(s.terminalId())))
|
||||||
.toList();
|
.toList();
|
||||||
ctx.status(200).json(Map.of("workers", out));
|
Map<String, Object> body = new LinkedHashMap<>();
|
||||||
|
body.put("workers", out);
|
||||||
|
// CB-586: operator visibility for the refs/wip snapshot store without shelling into the
|
||||||
|
// repo — how many snapshot refs exist and roughly what they cost. Present only once a
|
||||||
|
// worktree session has established the repo, so a never-snapshotted fleet reports nothing.
|
||||||
|
sessions.wipRefs().ifPresent(st -> body.put("wipRefs",
|
||||||
|
Map.of("count", st.count(), "costBytes", st.costBytes())));
|
||||||
|
ctx.status(200).json(body);
|
||||||
}
|
}
|
||||||
|
|
||||||
/** The configured worker profiles and which one a no-argument spawn uses. */
|
/** The configured worker profiles and which one a no-argument spawn uses. */
|
||||||
@@ -509,10 +516,20 @@ public final class BridgedApp {
|
|||||||
return;
|
return;
|
||||||
}
|
}
|
||||||
try {
|
try {
|
||||||
ctx.status(200).json(Map.of(
|
Map<String, Object> body = new LinkedHashMap<>();
|
||||||
"sessionId", id,
|
body.put("sessionId", id);
|
||||||
"status", messages.status(id).name().toLowerCase(),
|
body.put("status", messages.status(id).name().toLowerCase());
|
||||||
"ready", presence.isPresent(id)));
|
body.put("ready", presence.isPresent(id));
|
||||||
|
// CB-582: a worker paused mid-turn in an async bridge_ask is otherwise invisible to a
|
||||||
|
// status poll — surface the open question and how to answer it, same as bridge_poll's
|
||||||
|
// Phase.ASKING view.
|
||||||
|
MessageService.PendingAsk ask = messages.pendingAsk(id);
|
||||||
|
if (ask != null) {
|
||||||
|
body.put("question", ask.question());
|
||||||
|
body.put("turnId", ask.turnId());
|
||||||
|
body.put("ticket", ask.ticket());
|
||||||
|
}
|
||||||
|
ctx.status(200).json(body);
|
||||||
} catch (HerdrException e) {
|
} catch (HerdrException e) {
|
||||||
herdrError(ctx, e);
|
herdrError(ctx, e);
|
||||||
}
|
}
|
||||||
@@ -538,6 +555,11 @@ public final class BridgedApp {
|
|||||||
if (v.detail() != null) {
|
if (v.detail() != null) {
|
||||||
body.put("detail", v.detail());
|
body.put("detail", v.detail());
|
||||||
}
|
}
|
||||||
|
// CB-582: Phase.ASKING carries the question in v.reply() (handled above) and its answer-
|
||||||
|
// correlation id here — a REST caller polling this ticket otherwise has no way to answer it.
|
||||||
|
if (v.turnId() != null) {
|
||||||
|
body.put("turnId", v.turnId());
|
||||||
|
}
|
||||||
ctx.status(200).json(body);
|
ctx.status(200).json(body);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -12,9 +12,12 @@ import java.nio.file.Files;
|
|||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
import java.nio.file.StandardCopyOption;
|
import java.nio.file.StandardCopyOption;
|
||||||
import java.security.SecureRandom;
|
import java.security.SecureRandom;
|
||||||
|
import java.util.ArrayList;
|
||||||
|
import java.util.HashSet;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Optional;
|
import java.util.Optional;
|
||||||
|
import java.util.Set;
|
||||||
import java.util.concurrent.TimeUnit;
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.atomic.AtomicLong;
|
import java.util.concurrent.atomic.AtomicLong;
|
||||||
import java.util.stream.Collectors;
|
import java.util.stream.Collectors;
|
||||||
@@ -299,6 +302,129 @@ public final class GitWorktrees implements Worktrees {
|
|||||||
return index;
|
return index;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* One {@code refs/wip/<branch>} snapshot ref as read by {@link #listWipRefs}: its full ref name,
|
||||||
|
* the snapshot commit's sha, and that commit's committer time in unix millis (the age of the
|
||||||
|
* snapshot — a snapshot is written once and never rewritten, so the commit date is the ref's).
|
||||||
|
*/
|
||||||
|
private record WipRef(String refName, String sha, long committerMillis) {
|
||||||
|
String branch() {
|
||||||
|
return refName.substring("refs/wip/".length());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public WipRefStats wipRefs(String repoRoot) {
|
||||||
|
List<WipRef> refs = listWipRefs(repoRoot);
|
||||||
|
long costBytes = 0;
|
||||||
|
for (WipRef ref : refs) {
|
||||||
|
costBytes += treeSize(repoRoot, ref.sha());
|
||||||
|
}
|
||||||
|
return new WipRefStats(refs.size(), costBytes);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||||
|
// The rule is documented on Worktrees#pruneWipRefs: delete only a snapshot whose tree
|
||||||
|
// content is already reachable from main AND that is older than minAgeMillis. Reachability
|
||||||
|
// is the floor that keeps a worker's last copy; the age floor keeps a just-written snapshot
|
||||||
|
// from being swept while a lead may still be looking at it.
|
||||||
|
List<WipRef> refs = listWipRefs(repoRoot);
|
||||||
|
if (refs.isEmpty()) {
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
long nowMillis = System.currentTimeMillis();
|
||||||
|
// Resolve what main carries once per sweep, not once per ref.
|
||||||
|
Set<String> mainObjects = reachableObjectsFromMain(repoRoot);
|
||||||
|
int deleted = 0;
|
||||||
|
for (WipRef ref : refs) {
|
||||||
|
long ageMillis = nowMillis - ref.committerMillis();
|
||||||
|
if (ageMillis <= minAgeMillis) {
|
||||||
|
continue; // too recent — never swept, even if it looks recoverable (CB-586)
|
||||||
|
}
|
||||||
|
String tree = exec("git", "-C", repoRoot, "rev-parse", ref.sha() + "^{tree}").trim();
|
||||||
|
if (!mainObjects.contains(tree)) {
|
||||||
|
// Last copy of the snapshot's content — the worker's work exists nowhere else.
|
||||||
|
// Never delete automatically (CB-586 criterion 2).
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
exec("git", "-C", repoRoot, "update-ref", "-d", ref.refName());
|
||||||
|
deleted++;
|
||||||
|
log.info("pruned snapshot ref refs/wip/{} commit={} (age {}h): its tree is already "
|
||||||
|
+ "reachable from main, so the work is preserved; recover from reflog via "
|
||||||
|
+ "git update-ref refs/wip/{} {}",
|
||||||
|
ref.branch(), ref.sha(), TimeUnit.MILLISECONDS.toHours(ageMillis),
|
||||||
|
ref.branch(), ref.sha());
|
||||||
|
}
|
||||||
|
return deleted;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Every {@code refs/wip/*} ref (see {@link WipRef}). The committer date is read as a unix
|
||||||
|
* count of seconds and converted to millis. {@code %00} (NUL) separates the fields because a
|
||||||
|
* branch name may contain spaces.
|
||||||
|
*/
|
||||||
|
private List<WipRef> listWipRefs(String repoRoot) {
|
||||||
|
String out = exec("git", "-C", repoRoot, "for-each-ref",
|
||||||
|
"--format=%(refname)%00%(objectname)%00%(committerdate:unix)", "refs/wip/");
|
||||||
|
List<WipRef> refs = new ArrayList<>();
|
||||||
|
for (String line : out.split("\\R")) {
|
||||||
|
if (line.isBlank()) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
String[] parts = line.split("\u0000", -1);
|
||||||
|
if (parts.length == 3 && !parts[1].isBlank()) {
|
||||||
|
refs.add(new WipRef(parts[0], parts[1], Long.parseLong(parts[2]) * 1000L));
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return refs;
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* The set of object shas reachable from {@code main}, or an empty set when {@code main} cannot
|
||||||
|
* be resolved. An empty set is the safe direction: the retention sweep then concludes nothing
|
||||||
|
* is recoverable, so it deletes nothing — a repo with no {@code main} must never cause a
|
||||||
|
* worker's last copy of a snapshot to be dropped on a reachability misreading.
|
||||||
|
*/
|
||||||
|
private Set<String> reachableObjectsFromMain(String repoRoot) {
|
||||||
|
if (exitCode("git", "-C", repoRoot, "rev-parse", "--verify", "main") != 0) {
|
||||||
|
log.debug("refs/wip retention: no 'main' ref in {} — treating nothing as reachable", repoRoot);
|
||||||
|
return Set.of();
|
||||||
|
}
|
||||||
|
String out = exec("git", "-C", repoRoot, "rev-list", "--objects", "main");
|
||||||
|
Set<String> objects = new HashSet<>();
|
||||||
|
for (String line : out.split("\\R")) {
|
||||||
|
if (line.isBlank()) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
int sp = line.indexOf(' ');
|
||||||
|
objects.add(sp < 0 ? line : line.substring(0, sp));
|
||||||
|
}
|
||||||
|
return objects;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Approximate cost of a snapshot: the sum of every blob's size in its committed tree. */
|
||||||
|
private long treeSize(String repoRoot, String sha) {
|
||||||
|
String out = exec("git", "-C", repoRoot, "ls-tree", "-r", "-l", sha);
|
||||||
|
long total = 0;
|
||||||
|
for (String line : out.split("\\R")) {
|
||||||
|
if (line.isBlank()) {
|
||||||
|
continue;
|
||||||
|
}
|
||||||
|
// ls-tree -l row: "<mode> <type> <object> <size>\t<path>"; the size is only numeric for
|
||||||
|
// blobs (trees read "-"), so gate on the type token and take the 4th whitespace field.
|
||||||
|
String[] parts = line.split("\\s+");
|
||||||
|
if (parts.length >= 4 && "blob".equals(parts[1])) {
|
||||||
|
try {
|
||||||
|
total += Long.parseLong(parts[3]);
|
||||||
|
} catch (NumberFormatException ignored) {
|
||||||
|
// a '-' size (or any anomaly) contributes nothing to the rough figure
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
return total;
|
||||||
|
}
|
||||||
|
|
||||||
/** Resolve the directory that will hold per-session worktree checkouts. */
|
/** Resolve the directory that will hold per-session worktree checkouts. */
|
||||||
private Path resolveRoot(String repoRoot) {
|
private Path resolveRoot(String repoRoot) {
|
||||||
if (configuredRoot != null && !configuredRoot.isBlank()) {
|
if (configuredRoot != null && !configuredRoot.isBlank()) {
|
||||||
|
|||||||
@@ -53,6 +53,14 @@ public final class SessionManager implements TurnListener {
|
|||||||
private final int contextCap;
|
private final int contextCap;
|
||||||
private final boolean clearAfterTurn;
|
private final boolean clearAfterTurn;
|
||||||
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
private volatile MemberLifecycle memberLifecycle = MemberLifecycle.NONE;
|
||||||
|
/**
|
||||||
|
* CB-586: the repo root the fleet actually works in, remembered the first time a worktree
|
||||||
|
* session is spawned (worktrees are checkouts of it). {@code refs/wip/*} live there, and this
|
||||||
|
* single cached value is what the snapshot retention sweep and the operator-visible census run
|
||||||
|
* against. The daemon is bridged into one project at a time, so "the first worktree's repo" is
|
||||||
|
* the repo; {@code null} until any worktree is spawned, meaning nothing to sweep or measure.
|
||||||
|
*/
|
||||||
|
private volatile String fleetRepoRoot;
|
||||||
|
|
||||||
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
|
/** CB-520: notified with a terminalId on every acquire; no-op until wired. */
|
||||||
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
private final List<Consumer<String>> acquireListeners = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||||
@@ -293,8 +301,10 @@ public final class SessionManager implements TurnListener {
|
|||||||
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
|
// blocked caller fails fast with a real reason instead of sitting on a rendezvous
|
||||||
// nothing will ever resolve. CB-578 stage C: carry the worktree/branch/snapshot ref
|
// nothing will ever resolve. CB-578 stage C: carry the worktree/branch/snapshot ref
|
||||||
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
|
// too, so a failed ticket's detail can point a lead at the same tree to re-dispatch.
|
||||||
|
// CB-584 (issue #65 criterion 5): carry agentSessionId alongside them, so a lead can
|
||||||
|
// also resume the member's conversation, not just re-dispatch onto its files.
|
||||||
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
|
notifyReleased(new ReleaseDetail(removed.terminalId(), removed.worktree(),
|
||||||
removed.branch(), snapshotRef));
|
removed.branch(), snapshotRef, removed.agentSessionId()));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
// CB-581: the pane must always stop, even if the dirty check above threw. A session removed
|
||||||
@@ -348,7 +358,8 @@ public final class SessionManager implements TurnListener {
|
|||||||
* are {@code null} for a shared-tree session; {@code snapshotRef} is {@code null} unless this
|
* are {@code null} for a shared-tree session; {@code snapshotRef} is {@code null} unless this
|
||||||
* release snapshotted a dirty worktree into {@code refs/wip/<branch>}.
|
* release snapshotted a dirty worktree into {@code refs/wip/<branch>}.
|
||||||
*/
|
*/
|
||||||
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef) {
|
public record ReleaseDetail(String terminalId, String worktreePath, String branch, String snapshotRef,
|
||||||
|
String agentSessionId) {
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -450,6 +461,11 @@ public final class SessionManager implements TurnListener {
|
|||||||
// The non-worktree path always used this chain; only this branch was missed.
|
// The non-worktree path always used this chain; only this branch was missed.
|
||||||
String repoRoot = worktrees.repoRoot(
|
String repoRoot = worktrees.repoRoot(
|
||||||
launcher.effectiveCwd(new SpawnRequest(preResolvedProfile, requestedCwd, callerCwd)));
|
launcher.effectiveCwd(new SpawnRequest(preResolvedProfile, requestedCwd, callerCwd)));
|
||||||
|
if (fleetRepoRoot == null) {
|
||||||
|
// CB-586: remember the repo whose worktrees the fleet spawns — its refs/wip/* are the
|
||||||
|
// snapshot store the retention sweep and the operator census operate on.
|
||||||
|
fleetRepoRoot = repoRoot;
|
||||||
|
}
|
||||||
String branch = "worker/" + slug(wt.ticketSlug()) + "-" + nonce();
|
String branch = "worker/" + slug(wt.ticketSlug()) + "-" + nonce();
|
||||||
String path = null;
|
String path = null;
|
||||||
PeerHandle handle;
|
PeerHandle handle;
|
||||||
@@ -740,6 +756,26 @@ public final class SessionManager implements TurnListener {
|
|||||||
return registry.size();
|
return registry.size();
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586: the operator-visible census of {@code refs/wip/*} in the repo the fleet works in —
|
||||||
|
* how many snapshot refs exist and roughly what they cost. Empty (no repo known) until at
|
||||||
|
* least one worktree session has been spawned, exactly so a fleet that has never snapshotted
|
||||||
|
* anything surfaces nothing new, as it did before CB-586.
|
||||||
|
*/
|
||||||
|
public Optional<Worktrees.WipRefStats> wipRefs() {
|
||||||
|
String repo = fleetRepoRoot;
|
||||||
|
return repo == null ? Optional.empty() : Optional.of(worktrees.wipRefs(repo));
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586: run the snapshot retention sweep in the fleet's repo (a no-op until a worktree has
|
||||||
|
* been spawned, which establishes the repo). Returns how many {@code refs/wip/*} it deleted.
|
||||||
|
*/
|
||||||
|
public int sweepWipRefs(long minAgeMillis) {
|
||||||
|
String repo = fleetRepoRoot;
|
||||||
|
return repo == null ? 0 : worktrees.pruneWipRefs(repo, minAgeMillis);
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The registered session owning {@code terminalId}, or {@code null} if none does.
|
* The registered session owning {@code terminalId}, or {@code null} if none does.
|
||||||
*
|
*
|
||||||
|
|||||||
@@ -14,12 +14,29 @@ public final class SessionReaper {
|
|||||||
|
|
||||||
private static final Logger log = LoggerFactory.getLogger(SessionReaper.class);
|
private static final Logger log = LoggerFactory.getLogger(SessionReaper.class);
|
||||||
private static final long DEFAULT_INTERVAL_MILLIS = 5000;
|
private static final long DEFAULT_INTERVAL_MILLIS = 5000;
|
||||||
|
/** CB-586: the refs/wip age floor — never sweep a snapshot younger than 24h (the CB-586 rule). */
|
||||||
|
private static final long WIP_MIN_AGE_MILLIS = TimeUnit.HOURS.toMillis(24);
|
||||||
|
/**
|
||||||
|
* CB-586: how often the retention sweep runs. Given the 24h age floor, running it every few
|
||||||
|
* hours means a ref is dropped within hours of becoming eligible, never within minutes.
|
||||||
|
*/
|
||||||
|
private static final long WIP_SWEEP_INTERVAL_NANOS = TimeUnit.HOURS.toNanos(6);
|
||||||
|
|
||||||
private final SessionManager sessions;
|
private final SessionManager sessions;
|
||||||
private final long idleTtlNanos;
|
private final long idleTtlNanos;
|
||||||
private final long intervalMillis;
|
private final long intervalMillis;
|
||||||
private volatile boolean running;
|
private volatile boolean running;
|
||||||
private Thread thread;
|
private Thread thread;
|
||||||
|
/**
|
||||||
|
* When the retention sweep last ran, and whether it ever has. The flag is not a convenience:
|
||||||
|
* a "never yet" sentinel value cannot be compared by subtraction. {@code Long.MIN_VALUE} was
|
||||||
|
* the obvious choice and it silently overflows — {@code System.nanoTime()} is positive on this
|
||||||
|
* platform, so {@code now - Long.MIN_VALUE} wraps to a large negative number, the interval gate
|
||||||
|
* reads it as "swept moments ago", and it returns before ever assigning the field. The sweep
|
||||||
|
* then never runs at all, for the life of the process, with nothing in the log to say so.
|
||||||
|
*/
|
||||||
|
private volatile boolean sweptOnce;
|
||||||
|
private volatile long lastWipSweepNanos;
|
||||||
|
|
||||||
/** Construct a reaper with the default 5-second polling interval. */
|
/** Construct a reaper with the default 5-second polling interval. */
|
||||||
public SessionReaper(SessionManager sessions, long idleTtlSeconds) {
|
public SessionReaper(SessionManager sessions, long idleTtlSeconds) {
|
||||||
@@ -49,10 +66,39 @@ public final class SessionReaper {
|
|||||||
} catch (RuntimeException e) {
|
} catch (RuntimeException e) {
|
||||||
log.warn("session reaper iteration failed; continuing", e);
|
log.warn("session reaper iteration failed; continuing", e);
|
||||||
}
|
}
|
||||||
|
maybeSweepWipRefs();
|
||||||
sleep();
|
sleep();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586: run the refs/wip retention sweep on a slow cadence (hours, not the per-iteration
|
||||||
|
* millisecond loop). Best-effort — a failure must never take the idle-reap loop down with it.
|
||||||
|
*/
|
||||||
|
private void maybeSweepWipRefs() {
|
||||||
|
long now = System.nanoTime();
|
||||||
|
// The first pass always sweeps: a restart is a fine moment to sweep, the 24h age floor
|
||||||
|
// makes it safe, and it means the feature is observable right after a redeploy instead of
|
||||||
|
// six hours later. Only after that does the interval gate apply, and by then both operands
|
||||||
|
// come from nanoTime, so the subtraction is well-defined.
|
||||||
|
if (sweptOnce && now - lastWipSweepNanos < WIP_SWEEP_INTERVAL_NANOS) {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
try {
|
||||||
|
int deleted = sessions.sweepWipRefs(WIP_MIN_AGE_MILLIS);
|
||||||
|
if (deleted > 0) {
|
||||||
|
log.info("refs/wip retention sweep deleted {} snapshot ref(s) older than 24h whose "
|
||||||
|
+ "content was already reachable from main", deleted);
|
||||||
|
}
|
||||||
|
} catch (RuntimeException e) {
|
||||||
|
log.warn("refs/wip retention sweep failed; continuing", e);
|
||||||
|
}
|
||||||
|
// Set even when the sweep threw, so a broken repo is retried on the slow cadence rather
|
||||||
|
// than hammering git on every 5-second iteration.
|
||||||
|
lastWipSweepNanos = now;
|
||||||
|
sweptOnce = true;
|
||||||
|
}
|
||||||
|
|
||||||
private void sleep() {
|
private void sleep() {
|
||||||
try {
|
try {
|
||||||
Thread.sleep(intervalMillis);
|
Thread.sleep(intervalMillis);
|
||||||
|
|||||||
@@ -54,4 +54,48 @@ public interface Worktrees {
|
|||||||
* tolerance — a worktree that is gone holds nothing to snapshot)
|
* tolerance — a worktree that is gone holds nothing to snapshot)
|
||||||
*/
|
*/
|
||||||
Optional<String> snapshot(String worktreePath, String branch, String message);
|
Optional<String> snapshot(String worktreePath, String branch, String message);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586: how many {@code refs/wip/*} snapshot refs exist in {@code repoRoot} and roughly what
|
||||||
|
* they cost. This is the operator-visible surface for the snapshot growth CB-578 stage C left
|
||||||
|
* behind — counts of refs alone hide that each one pins a whole tree for {@code git gc}.
|
||||||
|
*
|
||||||
|
* @param repoRoot the repository to scan
|
||||||
|
* @return count of snapshot refs, and {@code costBytes} = the approximate total working-tree
|
||||||
|
* size of every snapshot's committed content (summed per ref, so shared objects are
|
||||||
|
* counted once per ref that carries them)
|
||||||
|
*/
|
||||||
|
WipRefStats wipRefs(String repoRoot);
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586: run the {@code refs/wip/*} retention sweep and return how many refs it deleted.
|
||||||
|
*
|
||||||
|
* <p>The retention rule is <em>reachability plus an age floor</em>. A snapshot ref is deleted
|
||||||
|
* only when <strong>both</strong> hold:
|
||||||
|
* <ol>
|
||||||
|
* <li>its commit's <em>tree content</em> is already reachable from {@code main} — the work
|
||||||
|
* the snapshot preserved has been recovered, so dropping the ref loses nothing; and</li>
|
||||||
|
* <li>the ref is older than {@code minAgeMillis} — a very recent snapshot is never swept
|
||||||
|
* while a lead may still be looking at it.</li>
|
||||||
|
* </ol>
|
||||||
|
*
|
||||||
|
* <p>Reachability is the safety property. A snapshot exists precisely because the work was not
|
||||||
|
* committed anywhere else, so a snapshot whose content is <em>not</em> reachable from
|
||||||
|
* {@code main} is the <strong>last copy</strong> of a worker's work and must never be deleted
|
||||||
|
* automatically — that is the failure CB-576 and CB-578 stage C were built to stop. Age alone
|
||||||
|
* must never drive a deletion, because age-based sweeping is exactly how the last copy gets
|
||||||
|
* destroyed. (Both numbers and the rule are CB-586's decision; this method only implements it.)
|
||||||
|
*
|
||||||
|
* <p>Every deletion logs the ref name and the commit sha, so an operator who finds they lost
|
||||||
|
* the wrong thing can still recover it from git's reflog.
|
||||||
|
*
|
||||||
|
* @param repoRoot the repository whose {@code refs/wip/*} to sweep
|
||||||
|
* @param minAgeMillis the age floor; a ref younger than this is never touched
|
||||||
|
* @return the number of snapshot refs deleted
|
||||||
|
*/
|
||||||
|
int pruneWipRefs(String repoRoot, long minAgeMillis);
|
||||||
|
|
||||||
|
/** CB-586: the operator-visible census of {@code refs/wip/*} in one repository. */
|
||||||
|
record WipRefStats(int count, long costBytes) {
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -0,0 +1,108 @@
|
|||||||
|
package dev.ltms.bridged;
|
||||||
|
|
||||||
|
import ch.qos.logback.classic.Logger;
|
||||||
|
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||||
|
import ch.qos.logback.core.read.ListAppender;
|
||||||
|
import dev.ltms.bridged.config.BridgedConfig;
|
||||||
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.junit.jupiter.api.io.TempDir;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
|
import java.nio.file.Files;
|
||||||
|
import java.nio.file.Path;
|
||||||
|
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: an absent (or empty) {@code memberCredentials:} block blocks nothing — no name is
|
||||||
|
* hardcoded any more to fall back on, so the daemon must say so out loud at startup rather than
|
||||||
|
* silently dropping CB-592's protection. Mirrors {@link RequiredSecretEnvVarsTest}'s pattern for
|
||||||
|
* the CB-594 startup-secrets report, capturing the real log via a {@link ListAppender}.
|
||||||
|
*/
|
||||||
|
class MemberCredentialsGapReportTest {
|
||||||
|
|
||||||
|
private static BridgedConfig load(Path dir, String yaml) throws Exception {
|
||||||
|
Path f = dir.resolve("bridged.yaml");
|
||||||
|
Files.writeString(f, yaml);
|
||||||
|
return BridgedConfig.load(f);
|
||||||
|
}
|
||||||
|
|
||||||
|
private static ListAppender<ILoggingEvent> attach() {
|
||||||
|
Logger logger = (Logger) LoggerFactory.getLogger(Bridged.class);
|
||||||
|
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||||
|
appender.start();
|
||||||
|
logger.addAppender(appender);
|
||||||
|
return appender;
|
||||||
|
}
|
||||||
|
|
||||||
|
private static void detach(ListAppender<ILoggingEvent> appender) {
|
||||||
|
((Logger) LoggerFactory.getLogger(Bridged.class)).detachAppender(appender);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void anAbsentBlockWarnsThatEveryMemberInheritsTheWholeStore(@TempDir Path dir) throws Exception {
|
||||||
|
BridgedConfig cfg = load(dir, "bind:\n host: 127.0.0.1\n port: 8765\n");
|
||||||
|
|
||||||
|
ListAppender<ILoggingEvent> appender = attach();
|
||||||
|
try {
|
||||||
|
Bridged.reportMemberCredentialsGap(cfg);
|
||||||
|
} finally {
|
||||||
|
detach(appender);
|
||||||
|
}
|
||||||
|
|
||||||
|
assertTrue(appender.list.stream().anyMatch(e ->
|
||||||
|
e.getLevel() == ch.qos.logback.classic.Level.WARN
|
||||||
|
&& e.getFormattedMessage().contains("memberCredentials")
|
||||||
|
&& e.getFormattedMessage().contains("WHOLE secret store")),
|
||||||
|
"an absent block must WARN that protection is lost, not stay silent");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void anEmptyKnownListWarnsTheSameAsAbsent(@TempDir Path dir) throws Exception {
|
||||||
|
BridgedConfig cfg = load(dir, """
|
||||||
|
memberCredentials:
|
||||||
|
policy: deny-by-default
|
||||||
|
""");
|
||||||
|
|
||||||
|
ListAppender<ILoggingEvent> appender = attach();
|
||||||
|
try {
|
||||||
|
Bridged.reportMemberCredentialsGap(cfg);
|
||||||
|
} finally {
|
||||||
|
detach(appender);
|
||||||
|
}
|
||||||
|
|
||||||
|
assertTrue(appender.list.stream().anyMatch(e ->
|
||||||
|
e.getLevel() == ch.qos.logback.classic.Level.WARN
|
||||||
|
&& e.getFormattedMessage().contains("memberCredentials")),
|
||||||
|
"policy: with no known: names still blocks nothing and must warn the same way");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void aPopulatedKnownListLogsInfoNotWarn(@TempDir Path dir) throws Exception {
|
||||||
|
BridgedConfig cfg = load(dir, """
|
||||||
|
memberCredentials:
|
||||||
|
policy: deny-by-default
|
||||||
|
allow: [AI_GATEWAY_TOKEN]
|
||||||
|
known: [AI_GATEWAY_TOKEN, GITEA_ACCESS_TOKEN]
|
||||||
|
""");
|
||||||
|
|
||||||
|
// logback-test.xml pins dev.ltms.bridged to WARN (see its own comment); raise it here so
|
||||||
|
// the INFO line this test asserts on actually reaches the appender, and restore after.
|
||||||
|
Logger logger = (Logger) LoggerFactory.getLogger(Bridged.class);
|
||||||
|
ch.qos.logback.classic.Level original = logger.getLevel();
|
||||||
|
logger.setLevel(ch.qos.logback.classic.Level.INFO);
|
||||||
|
ListAppender<ILoggingEvent> appender = attach();
|
||||||
|
try {
|
||||||
|
Bridged.reportMemberCredentialsGap(cfg);
|
||||||
|
} finally {
|
||||||
|
detach(appender);
|
||||||
|
logger.setLevel(original);
|
||||||
|
}
|
||||||
|
|
||||||
|
assertFalse(appender.list.stream().anyMatch(e -> e.getLevel() == ch.qos.logback.classic.Level.WARN),
|
||||||
|
"a configured, non-empty known: list must not warn — the block is doing its job");
|
||||||
|
assertTrue(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("blocking 1")),
|
||||||
|
"the INFO line should say how many names are actually blocked (known minus allow)");
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -10,6 +10,7 @@ import java.nio.file.Path;
|
|||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
|
import java.util.regex.Pattern;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.*;
|
import static org.junit.jupiter.api.Assertions.*;
|
||||||
|
|
||||||
@@ -1039,6 +1040,44 @@ class BridgedConfigTest {
|
|||||||
"an opencode worker with no argv defaults to the opencode binary, never claude");
|
"an opencode worker with no argv defaults to the opencode binary, never claude");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-604: an unrecognized {@code kind:} used to silently fall into the claude-code bucket — not
|
||||||
|
* matching {@code "opencode"} was the only check. With {@code argv:} also unset that meant the
|
||||||
|
* daemon tried to launch a program literally named after the typo.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void unknownKindIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("kind-typo.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
profiles:
|
||||||
|
gemini:
|
||||||
|
kind: opencod
|
||||||
|
model: google/gemini-2.5-pro
|
||||||
|
""");
|
||||||
|
|
||||||
|
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||||
|
assertTrue(e.getMessage().contains("gemini"), "error names the profile: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("opencod"), "error names the bad value: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("claude-code") && e.getMessage().contains("opencode"),
|
||||||
|
"error names the accepted set: " + e.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void blankKindStillDefaultsToClaudeCode(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("kind-blank.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
profiles:
|
||||||
|
gx10:
|
||||||
|
baseUrl: http://gx10.gw:8000
|
||||||
|
kind: ""
|
||||||
|
argv: ["ccs", "gx10"]
|
||||||
|
""");
|
||||||
|
|
||||||
|
BridgedConfig cfg = BridgedConfig.load(f);
|
||||||
|
assertEquals(BridgedConfig.Profile.KIND_CLAUDE_CODE, cfg.profiles().get("gx10").kind(),
|
||||||
|
"a blank kind: is documented to behave exactly like an absent one");
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void authDefaultsToLoopbackTrustSoExistingConfigsBehaveAsBefore(@TempDir Path dir) throws Exception {
|
void authDefaultsToLoopbackTrustSoExistingConfigsBehaveAsBefore(@TempDir Path dir) throws Exception {
|
||||||
Path f = dir.resolve("no-auth-block.yaml");
|
Path f = dir.resolve("no-auth-block.yaml");
|
||||||
@@ -1157,6 +1196,10 @@ class BridgedConfigTest {
|
|||||||
idleAfterSeconds: 600
|
idleAfterSeconds: 600
|
||||||
backoffMs: 45000
|
backoffMs: 45000
|
||||||
quietNudgeCap: 5
|
quietNudgeCap: 5
|
||||||
|
memberCredentials:
|
||||||
|
policy: deny-by-default
|
||||||
|
allow: [AI_GATEWAY_TOKEN]
|
||||||
|
known: [AI_GATEWAY_TOKEN, GITEA_ACCESS_TOKEN]
|
||||||
""");
|
""");
|
||||||
|
|
||||||
BridgedConfig cfg = BridgedConfig.load(f);
|
BridgedConfig cfg = BridgedConfig.load(f);
|
||||||
@@ -1187,6 +1230,84 @@ class BridgedConfigTest {
|
|||||||
assertEquals(600, cfg.leadHeartbeat().idleAfterSeconds(), "leadHeartbeat binds at the top level");
|
assertEquals(600, cfg.leadHeartbeat().idleAfterSeconds(), "leadHeartbeat binds at the top level");
|
||||||
assertEquals(45_000L, cfg.leadHeartbeat().backoffMs());
|
assertEquals(45_000L, cfg.leadHeartbeat().backoffMs());
|
||||||
assertEquals(5, cfg.leadHeartbeat().quietNudgeCap());
|
assertEquals(5, cfg.leadHeartbeat().quietNudgeCap());
|
||||||
|
|
||||||
|
assertEquals(BridgedConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT, cfg.memberCredentials().policy(),
|
||||||
|
"memberCredentials binds at the top level");
|
||||||
|
assertEquals(Set.of("AI_GATEWAY_TOKEN"), cfg.memberCredentials().allowSet());
|
||||||
|
assertEquals(Set.of("GITEA_ACCESS_TOKEN"), cfg.memberCredentials().blockedSet());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* A top-level key {@code BridgedConfig} reads but that appears nowhere in
|
||||||
|
* {@code bridged.example.yaml} — live or commented — is invisible drift: {@code bridged.yaml}
|
||||||
|
* is gitignored, so the example is the ONLY committed description of the config schema, and
|
||||||
|
* neither {@link #shippedExampleConfigParses} (example → code: does the example still parse)
|
||||||
|
* nor {@link #everyOptionalKnobDocumentedInTheExampleBinds} (a hand-maintained list of keys
|
||||||
|
* that must bind) can catch a brand-new key nobody added to either.
|
||||||
|
*
|
||||||
|
* <p>This test compares the OTHER direction: every key in {@link BridgedConfig#KNOWN_TOP_LEVEL_KEYS}
|
||||||
|
* (the parser's own accepted set, which backs the unknown-key WARN) must appear as a top-level
|
||||||
|
* key in the example text, live or commented-out — see {@link #topLevelKeyDocumented}.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void everyKnownTopLevelKeyIsDocumentedInTheExample() throws Exception {
|
||||||
|
Path example = Path.of("bridged.example.yaml");
|
||||||
|
assertTrue(Files.exists(example), "bridged.example.yaml must ship next to the pom");
|
||||||
|
String text = Files.readString(example);
|
||||||
|
|
||||||
|
List<String> undocumented = BridgedConfig.KNOWN_TOP_LEVEL_KEYS.stream()
|
||||||
|
.filter(key -> !topLevelKeyDocumented(text, key))
|
||||||
|
.sorted()
|
||||||
|
.toList();
|
||||||
|
|
||||||
|
assertTrue(undocumented.isEmpty(), () -> "key(s) " + undocumented
|
||||||
|
+ " are read by BridgedConfig but appear nowhere in bridged.example.yaml — "
|
||||||
|
+ "document each one there, commented out if optional. bridged.yaml is "
|
||||||
|
+ "gitignored, so this file is the only committed description of the config "
|
||||||
|
+ "schema an operator or a worker can see.");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Most of {@code bridged.example.yaml} is deliberately commented out — optional sections are
|
||||||
|
* documented as commented blocks so the shipped file stays a working minimal config. A key
|
||||||
|
* documented ONLY as a comment must still count as documented; parsing the file as YAML and
|
||||||
|
* reading its live key set (as an earlier attempt at this guard did) gets this wrong, because
|
||||||
|
* every commented section then looks entirely absent.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void commentedOnlyTopLevelKeyCountsAsDocumented() {
|
||||||
|
String yaml = """
|
||||||
|
bind:
|
||||||
|
port: 8765
|
||||||
|
# broker:
|
||||||
|
# uri: amqp://guest:guest@127.0.0.1:5672
|
||||||
|
""";
|
||||||
|
assertTrue(topLevelKeyDocumented(yaml, "broker"),
|
||||||
|
"a key documented only inside a commented-out block must still count as documented");
|
||||||
|
}
|
||||||
|
|
||||||
|
/** A key that appears in neither a live nor a commented top-level line must NOT count. */
|
||||||
|
@Test
|
||||||
|
void absentTopLevelKeyIsNotDocumented() {
|
||||||
|
String yaml = """
|
||||||
|
bind:
|
||||||
|
port: 8765
|
||||||
|
""";
|
||||||
|
assertFalse(topLevelKeyDocumented(yaml, "broker"),
|
||||||
|
"a key mentioned nowhere in the example must not be reported as documented");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* True when {@code key} appears as a top-level YAML key in {@code yaml} — either live
|
||||||
|
* ({@code key:} at column 0) or commented out ({@code # key:}, also at column 0, with only
|
||||||
|
* whitespace between the {@code #} and the key). Anchoring on column 0 is what keeps this a
|
||||||
|
* top-level check: an indented occurrence (a nested field, or prose inside a comment that
|
||||||
|
* happens to end in a colon) never matches, because {@code ^} requires the key's own first
|
||||||
|
* character — or the sole leading {@code #} — to sit at the very start of the line.
|
||||||
|
*/
|
||||||
|
private static boolean topLevelKeyDocumented(String yaml, String key) {
|
||||||
|
Pattern p = Pattern.compile("(?m)^(?:#\\s*)?" + Pattern.quote(key) + ":");
|
||||||
|
return p.matcher(yaml).find();
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@@ -1295,6 +1416,158 @@ class BridgedConfigTest {
|
|||||||
assertTrue(e.getMessage().contains("maxLoad"), "error names the key: " + e.getMessage());
|
assertTrue(e.getMessage().contains("maxLoad"), "error names the key: " + e.getMessage());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-606: an unrecognized {@code auth.mode} used to silently fall back to
|
||||||
|
* {@code loopback-trust} — {@link BridgedConfig.Auth#tokenMode()} only checked equality
|
||||||
|
* against {@code "token"}. On a loopback bind {@link BridgedConfig#validateAuthExposure()}
|
||||||
|
* never runs (it only fires for a non-loopback bind), so the typo was completely invisible:
|
||||||
|
* the daemon started cleanly and authenticated nobody while the operator believed token mode
|
||||||
|
* was active.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void unknownAuthModeIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("auth-mode-typo.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
bind:
|
||||||
|
host: 127.0.0.1
|
||||||
|
port: 8765
|
||||||
|
auth:
|
||||||
|
mode: toekn
|
||||||
|
""");
|
||||||
|
|
||||||
|
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||||
|
assertTrue(e.getMessage().contains("toekn"), "error names the bad value: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("loopback-trust") && e.getMessage().contains("token"),
|
||||||
|
"error names the accepted set: " + e.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-606: an unrecognized per-profile {@code placement:} used to silently fall back to legacy
|
||||||
|
* pane placement — {@link BridgedConfig.Profile#tabPlacement()} only checked equality against
|
||||||
|
* {@code "tab"}.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void unknownProfilePlacementIsRefusedAtLoadNamingTheProfileAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("placement-typo.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
profiles:
|
||||||
|
gx10:
|
||||||
|
baseUrl: http://gx10.gw:8000
|
||||||
|
placement: tabb
|
||||||
|
""");
|
||||||
|
|
||||||
|
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||||
|
assertTrue(e.getMessage().contains("gx10"), "error names the profile: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("tabb"), "error names the bad value: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("tab") && e.getMessage().contains("pane"),
|
||||||
|
"error names the accepted set: " + e.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: {@code memberCredentials.policy} is validated the same way {@code auth.mode} and
|
||||||
|
* per-profile {@code placement} are (CB-606's pattern) — a typo must not silently behave as the
|
||||||
|
* one real policy, because the day a second policy exists that silent fallback becomes a real
|
||||||
|
* behavior change instead of a happy accident.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void unknownMemberCredentialsPolicyIsRefusedAtLoadNamingTheValueAndTheAcceptedSet(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("member-credentials-policy-typo.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
bind:
|
||||||
|
host: 127.0.0.1
|
||||||
|
port: 8765
|
||||||
|
memberCredentials:
|
||||||
|
policy: deny-by-defualt
|
||||||
|
known:
|
||||||
|
- GITEA_ACCESS_TOKEN
|
||||||
|
""");
|
||||||
|
|
||||||
|
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||||
|
assertTrue(e.getMessage().contains("deny-by-defualt"), "error names the bad value: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("deny-by-default"), "error names the accepted set: " + e.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
|
/** {@code allow}/{@code known} bind and {@link BridgedConfig.MemberCredentials#blockedSet()} is known minus allow. */
|
||||||
|
@Test
|
||||||
|
void memberCredentialsBindsAllowAndKnownAndComputesBlockedSet(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("member-credentials.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
bind:
|
||||||
|
port: 8080
|
||||||
|
memberCredentials:
|
||||||
|
policy: deny-by-default
|
||||||
|
allow:
|
||||||
|
- AI_GATEWAY_TOKEN
|
||||||
|
- WORKER_GITEA_TOKEN
|
||||||
|
known:
|
||||||
|
- AI_GATEWAY_TOKEN
|
||||||
|
- WORKER_GITEA_TOKEN
|
||||||
|
- GITEA_ACCESS_TOKEN
|
||||||
|
- GITLAB_PERSONAL_ACCESS_TOKEN
|
||||||
|
""");
|
||||||
|
|
||||||
|
BridgedConfig.MemberCredentials mc = BridgedConfig.load(f).memberCredentials();
|
||||||
|
assertEquals(BridgedConfig.MemberCredentials.POLICY_DENY_BY_DEFAULT, mc.policy());
|
||||||
|
assertEquals(Set.of("AI_GATEWAY_TOKEN", "WORKER_GITEA_TOKEN"), mc.allowSet());
|
||||||
|
assertEquals(Set.of("GITEA_ACCESS_TOKEN", "GITLAB_PERSONAL_ACCESS_TOKEN"), mc.blockedSet(),
|
||||||
|
"blockedSet is known minus allow");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596: omitting {@code memberCredentials:} entirely must NOT crash a reader that assumes a
|
||||||
|
* non-null block (the same "fill in nested defaults" contract every other structural field
|
||||||
|
* gets — see {@link BridgedConfig#withDefaults()}), but it also must not pretend anything is
|
||||||
|
* blocked: an empty {@code known} list blocks nothing, and that is a real gap the operator must
|
||||||
|
* close by configuring this block, not a safe default.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void absentMemberCredentialsDefaultsToAnEmptyNonNullBlock(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("member-credentials-absent.yaml");
|
||||||
|
Files.writeString(f, "bind:\n port: 8080\n");
|
||||||
|
|
||||||
|
BridgedConfig.MemberCredentials mc = BridgedConfig.load(f).memberCredentials();
|
||||||
|
assertNotNull(mc, "withDefaults() must never leave this null");
|
||||||
|
assertTrue(mc.known().isEmpty(), "no known list configured — nothing is blocked");
|
||||||
|
assertTrue(mc.allowSet().isEmpty());
|
||||||
|
assertTrue(mc.blockedSet().isEmpty());
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void absentProfilePlacementDefaultsToTab(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("placement-absent.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
profiles:
|
||||||
|
gx10:
|
||||||
|
baseUrl: http://gx10.gw:8000
|
||||||
|
""");
|
||||||
|
|
||||||
|
BridgedConfig cfg = BridgedConfig.load(f);
|
||||||
|
BridgedConfig.Profile w = cfg.profiles().get("gx10");
|
||||||
|
assertTrue(w.tabPlacement(), "an absent placement must keep defaulting to tab");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-606: the top-level {@code placement:} policy name WAS validated, but only lazily, by
|
||||||
|
* {@code PlacementPolicies.fromName} through {@code CompositePeerLauncher}'s per-spawn
|
||||||
|
* {@code Supplier} — so a bad name still started a daemon that looked healthy and failed only
|
||||||
|
* the first time something spawned without naming a profile. This must now fail at load.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void unknownTopLevelPlacementPolicyIsRefusedAtLoadNotLazilyAtFirstSpawn(@TempDir Path dir) throws Exception {
|
||||||
|
Path f = dir.resolve("placement-policy-typo.yaml");
|
||||||
|
Files.writeString(f, """
|
||||||
|
bind:
|
||||||
|
port: 8080
|
||||||
|
placement: weightd
|
||||||
|
""");
|
||||||
|
|
||||||
|
IllegalStateException e = assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||||
|
assertTrue(e.getMessage().contains("weightd"), "error names the bad value: " + e.getMessage());
|
||||||
|
assertTrue(e.getMessage().contains("fixed") && e.getMessage().contains("round-robin")
|
||||||
|
&& e.getMessage().contains("weighted"),
|
||||||
|
"error names the accepted set: " + e.getMessage());
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void subscriptionFlagBindsAndDefaultsFalse(@TempDir Path dir) throws Exception {
|
void subscriptionFlagBindsAndDefaultsFalse(@TempDir Path dir) throws Exception {
|
||||||
Path f = dir.resolve("subscription.yaml");
|
Path f = dir.resolve("subscription.yaml");
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ import java.util.ArrayList;
|
|||||||
import java.util.LinkedHashMap;
|
import java.util.LinkedHashMap;
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
|
import java.util.concurrent.CopyOnWriteArrayList;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* Recording fake {@link HerdrClient} for unit/acceptance tests. Returns canned frames
|
* Recording fake {@link HerdrClient} for unit/acceptance tests. Returns canned frames
|
||||||
@@ -22,7 +23,13 @@ public final class FakeHerdr implements HerdrClient {
|
|||||||
public static final long WORKER_PID = 4242;
|
public static final long WORKER_PID = 4242;
|
||||||
|
|
||||||
private final ObjectMapper mapper = new ObjectMapper();
|
private final ObjectMapper mapper = new ObjectMapper();
|
||||||
public final List<Call> calls = new ArrayList<>();
|
/**
|
||||||
|
* Thread-safe on purpose. Background loops — {@link dev.ltms.bridged.msg.ReplyPushLoop} and the
|
||||||
|
* lead heartbeat — call this fake from their own scheduler threads while a test polls
|
||||||
|
* {@link #called} from the test thread. A plain {@code ArrayList} threw
|
||||||
|
* {@code ConcurrentModificationException} out of {@code called()} when a nudge landed mid-stream.
|
||||||
|
*/
|
||||||
|
public final List<Call> calls = new CopyOnWriteArrayList<>();
|
||||||
private boolean healthy = true;
|
private boolean healthy = true;
|
||||||
private final List<String> extraWorkspaces = new ArrayList<>();
|
private final List<String> extraWorkspaces = new ArrayList<>();
|
||||||
private final List<String> extraAgents = new ArrayList<>();
|
private final List<String> extraAgents = new ArrayList<>();
|
||||||
|
|||||||
@@ -773,6 +773,52 @@ class BridgeMcpTest {
|
|||||||
assertEquals("blocked", textOf(res));
|
assertEquals("blocked", textOf(res));
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-582: a lead polling {@code bridge_status} on its normal cadence — not {@code bridge_poll}
|
||||||
|
* — must also see a worker's open async {@code bridge_ask} question, since the reverse-rendezvous
|
||||||
|
* window it opened with is far shorter than that cadence.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void statusReportsAnOpenQuestionWhenTheWorkerIsMidAsk() throws Exception {
|
||||||
|
String ticket = messages.sendAsync(T, "task that asks");
|
||||||
|
long deadline = System.currentTimeMillis() + 3000;
|
||||||
|
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||||
|
Thread.sleep(5);
|
||||||
|
}
|
||||||
|
assertTrue(rendezvous.isWaiting(T), "sendAsync should have opened its rendezvous waiter");
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.AskResult> ask =
|
||||||
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||||
|
|
||||||
|
MessageService.TaskView asking;
|
||||||
|
deadline = System.currentTimeMillis() + 3000;
|
||||||
|
do {
|
||||||
|
asking = messages.poll(ticket);
|
||||||
|
Thread.sleep(5);
|
||||||
|
} while (asking.phase() != MessageService.Phase.ASKING && System.currentTimeMillis() < deadline);
|
||||||
|
assertEquals(MessageService.Phase.ASKING, asking.phase());
|
||||||
|
|
||||||
|
McpSchema.CallToolResult res = BridgeMcp.status(messages, T);
|
||||||
|
assertNotEquals(Boolean.TRUE, res.isError());
|
||||||
|
String out = textOf(res);
|
||||||
|
assertTrue(out.startsWith("idle"), "the live status must still lead the text: " + out);
|
||||||
|
assertTrue(out.contains("which config file?"), "the question text must be shown: " + out);
|
||||||
|
assertTrue(out.contains("turnId=\"" + asking.turnId() + "\""), "the turnId must be shown: " + out);
|
||||||
|
assertTrue(out.contains("(ticket " + ticket + ")"), "the ticket must be shown: " + out);
|
||||||
|
|
||||||
|
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||||
|
String turnId = asking.turnId();
|
||||||
|
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||||
|
() -> messages.answer(turnId, "config.yaml", 5000));
|
||||||
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
|
deadline = System.currentTimeMillis() + 3000;
|
||||||
|
while (!rendezvous.isWaiting(T) && System.currentTimeMillis() < deadline) {
|
||||||
|
Thread.sleep(5);
|
||||||
|
}
|
||||||
|
assertTrue(rendezvous.resolve(T, "done"));
|
||||||
|
answer.get(5, TimeUnit.SECONDS);
|
||||||
|
}
|
||||||
|
|
||||||
// --- bridge_whoami: the caller's own identity, so an agent never has to guess its role -------
|
// --- bridge_whoami: the caller's own identity, so an agent never has to guess its role -------
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
@@ -1,5 +1,8 @@
|
|||||||
package dev.ltms.bridged.member;
|
package dev.ltms.bridged.member;
|
||||||
|
|
||||||
|
import ch.qos.logback.classic.Logger;
|
||||||
|
import ch.qos.logback.classic.spi.ILoggingEvent;
|
||||||
|
import ch.qos.logback.core.read.ListAppender;
|
||||||
import dev.ltms.bridged.config.BridgedConfig;
|
import dev.ltms.bridged.config.BridgedConfig;
|
||||||
import dev.ltms.bridged.guard.GuardException;
|
import dev.ltms.bridged.guard.GuardException;
|
||||||
import dev.ltms.bridged.guard.SubscriptionGuard;
|
import dev.ltms.bridged.guard.SubscriptionGuard;
|
||||||
@@ -12,6 +15,7 @@ import dev.ltms.bridged.peer.PeerHandle;
|
|||||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||||
import dev.ltms.bridged.peer.SpawnRequest;
|
import dev.ltms.bridged.peer.SpawnRequest;
|
||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
|
import org.slf4j.LoggerFactory;
|
||||||
|
|
||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
@@ -676,39 +680,90 @@ class ClaudeCodeLauncherTest {
|
|||||||
"the guard-checked baseUrl must win over any env: entry, or the boundary is bypassable");
|
"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 -----------------------------
|
// --- CB-596: config-driven member-credential policy (replaces CB-592's hardcoded single name) --
|
||||||
|
|
||||||
/**
|
/** A representative {@code memberCredentials} — 4 allowed, 4 blocked, matching the real ticket shape. */
|
||||||
* herdr's env map is an overlay onto its own (login-shell) process environment, so a worker
|
private static final BridgedConfig.MemberCredentials TEST_MEMBER_CREDENTIALS = new BridgedConfig.MemberCredentials(
|
||||||
* inherits whatever the daemon's shell carries — including the admin GITEA_ACCESS_TOKEN — for
|
null,
|
||||||
* every key baseEnv does not explicitly shadow. This pins that the launcher DOES send an
|
List.of("AI_GATEWAY_TOKEN", "WORKER_GITEA_TOKEN", "CONTEXT7_TOKEN", "GITEA_HOST"),
|
||||||
* explicit (non-blank) GITEA_ACCESS_TOKEN to herdr on every spawn, whatever the profile is, so
|
List.of("AI_GATEWAY_TOKEN", "WORKER_GITEA_TOKEN", "CONTEXT7_TOKEN", "GITEA_HOST",
|
||||||
* a future baseEnv refactor cannot silently drop it and reopen the leak. Asserted against what
|
"GITEA_ACCESS_TOKEN", "GITLAB_PERSONAL_ACCESS_TOKEN", "TS_AUTHKEY", "HASS_TOKEN"));
|
||||||
* 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");
|
private ClaudeCodeLauncher serviceWithCredentials(FakeHerdr herdr, BridgedConfig.MemberCredentials creds) {
|
||||||
assertNotNull(shadowed, "GITEA_ACCESS_TOKEN must be explicitly overlaid, not left unmentioned");
|
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||||
assertFalse(shadowed.isBlank(), "a blank overlay value's override behaviour is unverified — must be non-blank");
|
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||||
|
List.of("claude"), "tab", "bridged-workers", "worker: {profile} #{n}", null, null, null);
|
||||||
|
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
|
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||||
|
_ -> null, 0, System::currentTimeMillis, () -> {}, null, () -> creds);
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* No profile — present or future — may restore the admin token by naming it in {@code env:}.
|
* 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 operator's own credentials — for
|
||||||
|
* every key {@code baseEnv} does not explicitly shadow. This pins that every {@code known} name
|
||||||
|
* NOT also {@code allow}-ed gets an explicit (non-blank) sentinel overlay, whatever the profile
|
||||||
|
* is. Asserted against what tab.create's params actually carry, not an internal map built in the
|
||||||
|
* test (gitea #82).
|
||||||
|
*
|
||||||
|
* <p>Scope, measured on a live pane 2026-08-15 (CB-592): this pins what the launcher SENDS, and
|
||||||
|
* that is all a unit test can pin. It does not prove the value survives, and for names the
|
||||||
|
* operator's secrets.sh also exports it does not: the pane runs a login shell that puts the real
|
||||||
|
* value back over this sentinel unless the export is guarded on BRIDGED_MEMBER — see
|
||||||
|
* everySpawnMarksThePaneAsAMember below.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void everyKnownNameNotAllowedIsShadowedWithTheSentinel() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
serviceWithCredentials(herdr, TEST_MEMBER_CREDENTIALS).spawn();
|
||||||
|
|
||||||
|
Map<String, String> env = startEnv(herdr);
|
||||||
|
for (String blocked : List.of("GITEA_ACCESS_TOKEN", "GITLAB_PERSONAL_ACCESS_TOKEN", "TS_AUTHKEY", "HASS_TOKEN")) {
|
||||||
|
String shadowed = env.get(blocked);
|
||||||
|
assertNotNull(shadowed, blocked + " must be explicitly overlaid, not left unmentioned");
|
||||||
|
assertFalse(shadowed.isBlank(), blocked + "'s overlay value must be non-blank");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* An {@code allow}-ed name must get NO overlay entry at all — any entry, blank or not, risks
|
||||||
|
* overriding the real value the pane needs, and the whole point of {@code allow} is that the
|
||||||
|
* pane's own inherited value passes through untouched.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void everyAllowedNameGetsNoOverlayEntrySoTheRealValuePassesThrough() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
serviceWithCredentials(herdr, TEST_MEMBER_CREDENTIALS).spawn();
|
||||||
|
|
||||||
|
Map<String, String> env = startEnv(herdr);
|
||||||
|
for (String allowed : List.of("AI_GATEWAY_TOKEN", "WORKER_GITEA_TOKEN", "CONTEXT7_TOKEN", "GITEA_HOST")) {
|
||||||
|
assertFalse(env.containsKey(allowed),
|
||||||
|
allowed + " is allow-listed — the launcher must not mention it at all");
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* No {@code memberCredentials} configured (the pre-CB-596 constructor overloads still used
|
||||||
|
* throughout this file, and the shape a fresh {@code bridged.yaml} with no memberCredentials:
|
||||||
|
* block resolves to) blocks NOTHING. This documents the transitional gap rather than hiding it —
|
||||||
|
* see {@code HerdrPeerLauncher#applyMemberCredentialPolicy}'s javadoc.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void noMemberCredentialsConfiguredBlocksNothing() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
service(herdr, List.of("claude"), null).spawn();
|
||||||
|
|
||||||
|
assertNull(startEnv(herdr).get("GITEA_ACCESS_TOKEN"),
|
||||||
|
"with no memberCredentials configured, nothing is shadowed — config must supply the policy");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* No profile — present or future — may restore a blocked name by naming it in {@code env:}.
|
||||||
* The shadow is applied after the profile's own env in {@link HerdrPeerLauncher#baseEnv}
|
* 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.
|
* precisely so this can never happen; this test pins that ordering.
|
||||||
*/
|
*/
|
||||||
@Test
|
@Test
|
||||||
void aProfileEnvEntryCannotRestoreTheAdminGiteaAccessToken() {
|
void aProfileEnvEntryCannotRestoreABlockedName() {
|
||||||
FakeHerdr herdr = new FakeHerdr();
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||||
@@ -716,10 +771,48 @@ class ClaudeCodeLauncherTest {
|
|||||||
null, Map.of("GITEA_ACCESS_TOKEN", "admin-secret-from-profile-config"), null, null);
|
null, Map.of("GITEA_ACCESS_TOKEN", "admin-secret-from-profile-config"), null, null);
|
||||||
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||||
_ -> null).spawn();
|
_ -> null, 0, System::currentTimeMillis, () -> {}, null, () -> TEST_MEMBER_CREDENTIALS)
|
||||||
|
.spawn(cfg.profile(), null, null);
|
||||||
|
|
||||||
assertNotEquals("admin-secret-from-profile-config", startEnv(herdr).get("GITEA_ACCESS_TOKEN"),
|
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");
|
"a profile's own env: must not be able to smuggle a blocked name back in");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-596 criterion 4: a credential-shaped host env var name on neither {@code known} nor
|
||||||
|
* {@code allow} is not silently allowed — it must be reported (never its value). This pins the
|
||||||
|
* WARN naming the gap, using an injected host-env-names source rather than the real
|
||||||
|
* {@code System.getenv()} so the test is deterministic.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void aCredentialShapedNameOnNeitherListIsLoggedAsAGap() {
|
||||||
|
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", "worker: {profile} #{n}", null, null, null);
|
||||||
|
ClaudeCodeLauncher svc = new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
|
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(),
|
||||||
|
_ -> null, 0, System::currentTimeMillis, () -> {}, null, () -> TEST_MEMBER_CREDENTIALS,
|
||||||
|
() -> Set.of("PATH", "HOME", "AI_GATEWAY_TOKEN", "A_BRAND_NEW_SECRET_TOKEN"));
|
||||||
|
|
||||||
|
Logger logger = (Logger) LoggerFactory.getLogger(HerdrPeerLauncher.class);
|
||||||
|
ListAppender<ILoggingEvent> appender = new ListAppender<>();
|
||||||
|
appender.start();
|
||||||
|
logger.addAppender(appender);
|
||||||
|
try {
|
||||||
|
svc.spawn();
|
||||||
|
} finally {
|
||||||
|
logger.detachAppender(appender);
|
||||||
|
}
|
||||||
|
|
||||||
|
assertTrue(appender.list.stream().anyMatch(e ->
|
||||||
|
e.getFormattedMessage().contains("memberCredentials gap")
|
||||||
|
&& e.getFormattedMessage().contains("A_BRAND_NEW_SECRET_TOKEN")),
|
||||||
|
"the gap must name the unrecognized credential-shaped var, never a value");
|
||||||
|
assertFalse(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("PATH")),
|
||||||
|
"PATH/HOME are not credential-shaped and must not be reported as a gap");
|
||||||
|
assertFalse(appender.list.stream().anyMatch(e -> e.getFormattedMessage().contains("AI_GATEWAY_TOKEN")),
|
||||||
|
"a name already on allow: is covered, not a gap");
|
||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
|
|||||||
@@ -52,6 +52,14 @@ class OpenCodeLauncherTest {
|
|||||||
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, fleet);
|
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, fleet);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
private static OpenCodeLauncher serviceWithCredentials(FakeHerdr herdr, Path configRoot,
|
||||||
|
BridgedConfig.Profile cfg,
|
||||||
|
BridgedConfig.MemberCredentials creds) {
|
||||||
|
return new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
|
Map.of(cfg.profile(), cfg), cfg.profile(), k -> "GITEA_ACCESS_TOKEN".equals(k) ? "tok" : null,
|
||||||
|
0, System::currentTimeMillis, () -> { }, configRoot, configRoot, null, () -> creds);
|
||||||
|
}
|
||||||
|
|
||||||
@SuppressWarnings("unchecked")
|
@SuppressWarnings("unchecked")
|
||||||
private static Map<String, Object> lastStart(FakeHerdr herdr) {
|
private static Map<String, Object> lastStart(FakeHerdr herdr) {
|
||||||
return (Map<String, Object>) herdr.lastCall("agent.start").params();
|
return (Map<String, Object>) herdr.lastCall("agent.start").params();
|
||||||
@@ -209,18 +217,33 @@ class OpenCodeLauncherTest {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* CB-592: the shadow lives in {@link HerdrPeerLauncher#baseEnv}, shared by every adapter — this
|
* CB-596: the config-driven shadow lives in {@link HerdrPeerLauncher#baseEnv}, shared by every
|
||||||
* pins that the opencode path gets it too, not just Claude's. See the matching test in
|
* adapter — this pins that the opencode path gets it too, not just Claude's. See the matching
|
||||||
* {@code ClaudeCodeLauncherTest} for the full rationale (gitea #77).
|
* tests in {@code ClaudeCodeLauncherTest} for the full rationale (gitea #82, superseding CB-592's
|
||||||
|
* single hardcoded name).
|
||||||
*/
|
*/
|
||||||
@Test
|
@Test
|
||||||
void everySpawnShadowsTheAdminGiteaAccessToken(@TempDir Path root) {
|
void aKnownNameNotAllowedIsShadowedWithTheSentinel(@TempDir Path root) {
|
||||||
FakeHerdr herdr = new FakeHerdr();
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
service(herdr, root, opencodeCfg(null, null, null)).spawn();
|
BridgedConfig.MemberCredentials creds = new BridgedConfig.MemberCredentials(
|
||||||
|
null, List.of("AI_GATEWAY_TOKEN"), List.of("AI_GATEWAY_TOKEN", "GITEA_ACCESS_TOKEN"));
|
||||||
|
serviceWithCredentials(herdr, root, opencodeCfg(null, null, null), creds).spawn();
|
||||||
|
|
||||||
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
|
String shadowed = startEnv(herdr).get("GITEA_ACCESS_TOKEN");
|
||||||
assertNotNull(shadowed, "GITEA_ACCESS_TOKEN must be explicitly overlaid, not left unmentioned");
|
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");
|
assertFalse(shadowed.isBlank(), "a blank overlay value's override behaviour is unverified — must be non-blank");
|
||||||
|
assertFalse(startEnv(herdr).containsKey("AI_GATEWAY_TOKEN"),
|
||||||
|
"an allow-listed name must get no overlay entry at all");
|
||||||
|
}
|
||||||
|
|
||||||
|
/** No {@code memberCredentials} configured (the pre-CB-596 constructor overloads) blocks nothing. */
|
||||||
|
@Test
|
||||||
|
void noMemberCredentialsConfiguredBlocksNothing(@TempDir Path root) {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
service(herdr, root, opencodeCfg(null, null, null)).spawn();
|
||||||
|
|
||||||
|
assertNull(startEnv(herdr).get("GITEA_ACCESS_TOKEN"),
|
||||||
|
"with no memberCredentials configured, nothing is shadowed — config must supply the policy");
|
||||||
}
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
@@ -35,9 +35,12 @@ import static org.junit.jupiter.api.Assertions.assertTrue;
|
|||||||
*/
|
*/
|
||||||
class AmqpReplyInboxRecoveryRaceTest {
|
class AmqpReplyInboxRecoveryRaceTest {
|
||||||
|
|
||||||
/** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window
|
/** Large enough that thousands of entries are still unprocessed by the time the very first one
|
||||||
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
|
* is observed as failed (see {@code sweepIsHoldingTheLock} below) — that gap is what makes the
|
||||||
private static final int STALE_PUBLISHES = 100_000;
|
* head start deterministic instead of a coin flip. 100,000 gave the same guarantee but made the
|
||||||
|
* test far more expensive than the guarantee needs; the ordering no longer depends on a timing
|
||||||
|
* window sized to the full backlog; just to the tail of it. */
|
||||||
|
private static final int STALE_PUBLISHES = 2_000;
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
@Timeout(30)
|
@Timeout(30)
|
||||||
@@ -58,6 +61,12 @@ class AmqpReplyInboxRecoveryRaceTest {
|
|||||||
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
|
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
|
||||||
// Virtual threads make this many concurrent blocking publish() calls cheap.
|
// Virtual threads make this many concurrent blocking publish() calls cheap.
|
||||||
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
|
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
|
||||||
|
// Counted down by the FIRST stale publish thread to observe its own failure. That can only
|
||||||
|
// happen from inside failPendingPublishesOnRecovery() — nothing else in this test ever
|
||||||
|
// completes a stale Pending exceptionally (no nack/return is simulated for any "stale-*"
|
||||||
|
// msgId) — so seeing it fire is direct, observable proof the sweep is inside its loop, not a
|
||||||
|
// timing guess. It replaces the old fixed Thread.sleep(5) head start.
|
||||||
|
CountDownLatch sweepIsHoldingTheLock = new CountDownLatch(1);
|
||||||
for (int i = 0; i < STALE_PUBLISHES; i++) {
|
for (int i = 0; i < STALE_PUBLISHES; i++) {
|
||||||
String msgId = "stale-" + i;
|
String msgId = "stale-" + i;
|
||||||
Thread.ofVirtual().start(() -> {
|
Thread.ofVirtual().start(() -> {
|
||||||
@@ -65,7 +74,7 @@ class AmqpReplyInboxRecoveryRaceTest {
|
|||||||
try {
|
try {
|
||||||
inbox.publish("worker-stale", msgId, "x");
|
inbox.publish("worker-stale", msgId, "x");
|
||||||
} catch (IllegalStateException expected) {
|
} catch (IllegalStateException expected) {
|
||||||
// resolved (failed by the sweep) — that is exactly what this thread is here for
|
sweepIsHoldingTheLock.countDown();
|
||||||
}
|
}
|
||||||
});
|
});
|
||||||
}
|
}
|
||||||
@@ -83,13 +92,28 @@ class AmqpReplyInboxRecoveryRaceTest {
|
|||||||
}
|
}
|
||||||
}, "recovery-sweep");
|
}, "recovery-sweep");
|
||||||
sweepThread.start();
|
sweepThread.start();
|
||||||
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
|
|
||||||
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
|
// Deterministic head start: block until the sweep has actually failed one of the stale
|
||||||
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
|
// publishes. failPendingPublishesOnRecovery() (once guarded, as it is on main) holds
|
||||||
// register. Without this head start, "fresh" sometimes wins the race for the lock and
|
// publishChannelLock for its ENTIRE loop, not just per entry — so this failure proves the
|
||||||
// registers before the sweep even starts, which is the accepted "already in flight when
|
// sweep is, at this instant, still holding that lock. With STALE_PUBLISHES this large, the
|
||||||
// recovery fires" case (correctly failed either way) rather than the bug under test.
|
// remaining ~1,999 entries give an enormous margin between "first failure observed" and "sweep
|
||||||
Thread.sleep(5);
|
// releases the lock": there is no window left for "fresh" to slip in before the sweep starts,
|
||||||
|
// or to win the lock ahead of it — see the case-2 note below. This also means Case 1 (the sweep
|
||||||
|
// is already inside its loop, holding the lock, when "fresh" tries to register) is now
|
||||||
|
// guaranteed by construction rather than merely likely under a fixed sleep.
|
||||||
|
assertTrue(sweepIsHoldingTheLock.await(20, TimeUnit.SECONDS),
|
||||||
|
"the sweep never failed a single stale publish — it may not have started");
|
||||||
|
|
||||||
|
// Case 2 ("fresh" wins publishChannelLock before the sweep even starts, so it genuinely
|
||||||
|
// published on the stale channel and the sweep correctly fails it) is impossible by
|
||||||
|
// construction in this test: freshThread.start() below is reached only after
|
||||||
|
// sweepIsHoldingTheLock has counted down, which can only happen once
|
||||||
|
// failPendingPublishesOnRecovery() is already running and has already failed a stale entry.
|
||||||
|
// There is no code path that lets "fresh" start before the sweep starts. That case is real
|
||||||
|
// and correct production behaviour (see AmqpReplyInbox#failPendingPublishesOnRecovery's
|
||||||
|
// javadoc), it is just not reachable from this deterministic ordering, so it does not need a
|
||||||
|
// separate assertion here.
|
||||||
|
|
||||||
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
|
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
|
||||||
// sweep" racing a publish that registers while it is mid-flight.
|
// sweep" racing a publish that registers while it is mid-flight.
|
||||||
@@ -109,8 +133,8 @@ class AmqpReplyInboxRecoveryRaceTest {
|
|||||||
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
|
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
|
||||||
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
|
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
|
||||||
// for the registration rather than checking once: freshThread may still be contending for
|
// for the registration rather than checking once: freshThread may still be contending for
|
||||||
// publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the
|
// publishChannelLock (behind the sweep's own hold on it, and possibly other stale threads
|
||||||
// sweep itself has already finished.
|
// still unwinding) even though the sweep itself has already finished.
|
||||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
|
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
|
||||||
int idx = -1;
|
int idx = -1;
|
||||||
while (idx < 0 && System.nanoTime() < deadline) {
|
while (idx < 0 && System.nanoTime() < deadline) {
|
||||||
|
|||||||
@@ -714,6 +714,47 @@ class MessageServiceTest {
|
|||||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- CB-582: bridge_status pendingAsk() ------------------------------------------------------
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void pendingAskReturnsNullWhenNoQuestionIsOpen() throws Exception {
|
||||||
|
assertNull(messages.pendingAsk(T), "no async ticket at all -> no pending ask");
|
||||||
|
|
||||||
|
String ticket = messages.sendAsync(T, "long task");
|
||||||
|
awaitWaiting();
|
||||||
|
assertNull(messages.pendingAsk(T), "a plain pending delegation is not a question");
|
||||||
|
|
||||||
|
injectDelivery();
|
||||||
|
assertTrue(rendezvous.resolve(T, "done"));
|
||||||
|
awaitTicketPhase(ticket, MessageService.Phase.DONE);
|
||||||
|
assertNull(messages.pendingAsk(T), "a finished ticket carries no open question either");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void pendingAskReturnsTheOpenQuestionForAnAsyncTicket() throws Exception {
|
||||||
|
String ticket = messages.sendAsync(T, "task that asks");
|
||||||
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.AskResult> ask =
|
||||||
|
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config file?", 5000));
|
||||||
|
MessageService.TaskView asking = awaitTicketPhase(ticket, MessageService.Phase.ASKING);
|
||||||
|
|
||||||
|
MessageService.PendingAsk pending = messages.pendingAsk(T);
|
||||||
|
assertNotNull(pending, "bridge_status should see the open question");
|
||||||
|
assertEquals(ticket, pending.ticket());
|
||||||
|
assertEquals("which config file?", pending.question());
|
||||||
|
assertEquals(asking.turnId(), pending.turnId());
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||||
|
() -> messages.answer(asking.turnId(), "config.yaml", 5000));
|
||||||
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
|
assertNull(messages.pendingAsk(T), "an answered question is no longer pending");
|
||||||
|
awaitWaiting();
|
||||||
|
assertTrue(rendezvous.resolve(T, "done"));
|
||||||
|
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||||
|
}
|
||||||
|
|
||||||
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
|
// --- CB-588: async ticket terminal nudges ---------------------------------------------------
|
||||||
//
|
//
|
||||||
// MessageService.reply's rendezvous fast path is exactly what an async ticket always takes
|
// MessageService.reply's rendezvous fast path is exactly what an async ticket always takes
|
||||||
@@ -833,6 +874,110 @@ class MessageServiceTest {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- CB-582: bridge_ask question-open nudges --------------------------------------------------
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void anAsyncTicketThatPausesOnAQuestionNudgesTheLeadWithNoPriorPollCall() throws Exception {
|
||||||
|
try (var wiring = wireWithPushLoop(5, 50)) {
|
||||||
|
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||||
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||||
|
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||||
|
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
||||||
|
|
||||||
|
awaitNudge(wiring.leadHerdr());
|
||||||
|
String nudge = wiring.leadHerdr().lastCall("agent.prompt").params().toString();
|
||||||
|
assertTrue(nudge.contains(ticket), "the nudge should name the ticket: " + nudge);
|
||||||
|
assertTrue(nudge.contains(asking.turnId()), "the nudge should name the turnId: " + nudge);
|
||||||
|
assertTrue(nudge.contains("bridge_send(turnId="),
|
||||||
|
"the nudge should name the exact answer call: " + nudge);
|
||||||
|
assertTrue(nudge.contains("which config file?"), "the nudge should include the question: " + nudge);
|
||||||
|
|
||||||
|
// Clean up the still-open ask so the background thread does not linger past the test.
|
||||||
|
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||||
|
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||||
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
|
awaitWaiting();
|
||||||
|
assertTrue(rendezvous.resolve(T, "done"));
|
||||||
|
answer.get(5, TimeUnit.SECONDS);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void answeringAQuestionStopsFurtherNudgesAboutIt() throws Exception {
|
||||||
|
try (var wiring = wireWithPushLoop(5, 50)) {
|
||||||
|
String ticket = wiring.service().sendAsync(T, "task that asks");
|
||||||
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.AskResult> ask = CompletableFuture.supplyAsync(
|
||||||
|
() -> wiring.service().ask(T, "which config file?", 5000));
|
||||||
|
MessageService.TaskView asking = awaitTicketPhaseOn(wiring.service(), ticket, MessageService.Phase.ASKING);
|
||||||
|
|
||||||
|
awaitNudge(wiring.leadHerdr());
|
||||||
|
long callsBeforeAnswer = wiring.leadHerdr().calls.stream()
|
||||||
|
.filter(c -> c.method().equals("agent.prompt")).count();
|
||||||
|
|
||||||
|
CompletableFuture<MessageService.Reply> answer = CompletableFuture.supplyAsync(
|
||||||
|
() -> wiring.service().answer(asking.turnId(), "config.yaml", 5000));
|
||||||
|
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||||
|
awaitWaiting();
|
||||||
|
assertTrue(rendezvous.resolve(T, "done"));
|
||||||
|
answer.get(5, TimeUnit.SECONDS);
|
||||||
|
|
||||||
|
// Let several more ticks (and the ticket's own now-legitimate terminal nudge) fire —
|
||||||
|
// none of them may still name the question's turnId, which is closed.
|
||||||
|
Thread.sleep(300);
|
||||||
|
boolean anyNamesClosedQuestion = wiring.leadHerdr().calls.stream()
|
||||||
|
.filter(c -> c.method().equals("agent.prompt"))
|
||||||
|
.skip(callsBeforeAnswer)
|
||||||
|
.anyMatch(c -> c.params().toString().contains(asking.turnId()));
|
||||||
|
assertFalse(anyNamesClosedQuestion,
|
||||||
|
"no nudge sent after the answer may still name the now-closed turnId " + asking.turnId());
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void anAskThatLeavesByThrowingStillClosesItsQuestion() throws Exception {
|
||||||
|
// CB-582 follow-up. ask() calls clearAsyncQuestion on three paths — no-waiter, timed out,
|
||||||
|
// and (from answer()) answered — but it can also leave by *throwing*: an interrupt while
|
||||||
|
// blocked on the answer, or an ExecutionException from the answer future. Those paths run
|
||||||
|
// only the finally block, so before the fix the push loop kept the question pending for
|
||||||
|
// good: named in every nudge until its own cap, then never removed from the map at all.
|
||||||
|
PrimaryRegistry registry = new PrimaryRegistry(null);
|
||||||
|
registry.recordDelegation(T, LEAD);
|
||||||
|
var scheduler = java.util.concurrent.Executors.newSingleThreadScheduledExecutor();
|
||||||
|
// A backoff far longer than the test: the schedule is started but no tick ever fires, so
|
||||||
|
// decide() is read directly and nothing here depends on timing.
|
||||||
|
ReplyPushLoop pushLoop = new ReplyPushLoop(registry, new AgentControl(new FakeHerdr()), inbox,
|
||||||
|
scheduler, 5, 60_000);
|
||||||
|
MessageService service = new MessageService(agents, injector, rendezvous, inbox, pushLoop);
|
||||||
|
try {
|
||||||
|
String ticket = service.sendAsync(T, "task that asks");
|
||||||
|
awaitWaiting();
|
||||||
|
injectDelivery();
|
||||||
|
|
||||||
|
Thread asker = new Thread(() -> assertThrows(IllegalStateException.class,
|
||||||
|
() -> service.ask(T, "which config file?", 30_000)));
|
||||||
|
asker.start();
|
||||||
|
awaitTicketPhaseOn(service, ticket, MessageService.Phase.ASKING);
|
||||||
|
assertEquals(ReplyPushLoop.Action.INJECT, pushLoop.decide(LEAD, 0, 0, 0),
|
||||||
|
"the open question should be the one thing keeping this lead's schedule alive");
|
||||||
|
|
||||||
|
asker.interrupt();
|
||||||
|
asker.join(5000);
|
||||||
|
assertFalse(asker.isAlive(), "the interrupted ask should have left ask() by throwing");
|
||||||
|
|
||||||
|
assertEquals(ReplyPushLoop.Action.STOP, pushLoop.decide(LEAD, 0, 0, 0),
|
||||||
|
"an ask that threw must still close its question, or the loop nudges about it for good");
|
||||||
|
} finally {
|
||||||
|
service.close();
|
||||||
|
scheduler.shutdownNow();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception {
|
void aFleetWithNoPushLoopConfiguredBehavesExactlyAsToday() throws Exception {
|
||||||
// `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw,
|
// `messages` (the shared field) uses the no-pushLoop constructor — poll() must not throw,
|
||||||
|
|||||||
@@ -42,6 +42,7 @@ class ReplyPushLoopTest {
|
|||||||
|
|
||||||
private static final String PRIMARY = "term_primary";
|
private static final String PRIMARY = "term_primary";
|
||||||
private static final String WORKER = "term_worker";
|
private static final String WORKER = "term_worker";
|
||||||
|
private static final String WORKER2 = "term_worker2";
|
||||||
private static final ObjectMapper MAPPER = new ObjectMapper();
|
private static final ObjectMapper MAPPER = new ObjectMapper();
|
||||||
|
|
||||||
private PrimaryRegistry registry;
|
private PrimaryRegistry registry;
|
||||||
@@ -579,6 +580,209 @@ class ReplyPushLoopTest {
|
|||||||
"both the reply and the ticket source are at their own cap — must still stop");
|
"both the reply and the ticket source are at their own cap — must still stop");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// --- CB-598: work arriving during a backoff must not read as stale backlog -------------------
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void aTargetArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() {
|
||||||
|
// The bug: reminder counts used to be a single counter per lead per source, carried
|
||||||
|
// forward across scheduled ticks (scheduleNext(lead, count + 1, ...)) rather than tracked
|
||||||
|
// per pending item. WORKER gets nudged once here, which — with cap=1 — exhausts the
|
||||||
|
// shared reply-source counter for this lead. WORKER2 then queues a reply for the SAME
|
||||||
|
// lead "during the backoff": while the schedule from WORKER's tick is still active, before
|
||||||
|
// the next tick's own start-of-tick snapshot runs. At that next tick, the OLD code passed
|
||||||
|
// the already-exhausted shared counter into decide() regardless of WORKER2 never having
|
||||||
|
// been named in any nudge, and — because WORKER2 was already present in that tick's
|
||||||
|
// "before" snapshot — stopOrRestart's race check (proven correct on its own elsewhere in
|
||||||
|
// this file) does not save it either: it looks like ordinary stale backlog, not a race.
|
||||||
|
// WORKER2 was then stranded forever with no live schedule and no nudge ever naming it.
|
||||||
|
//
|
||||||
|
// tick() is driven directly (package-private, same reasoning as stopOrRestart being
|
||||||
|
// directly testable) so the exact interleaving is deterministic instead of racing the
|
||||||
|
// scheduler thread over a real ~15s backoff.
|
||||||
|
//
|
||||||
|
// Before the fix, this test fails on the second assertEquals: rec.sendCount() stays at 1
|
||||||
|
// (decide() returns STOP on the second tick(), so injectNudge is never called a second
|
||||||
|
// time) and the "must still get one" assertion never even runs.
|
||||||
|
int cap = 1;
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
inbox.own(WORKER2);
|
||||||
|
inbox.publish(WORKER, "m1", "hello");
|
||||||
|
var loop = loop(cap, 100_000); // huge backoff — nothing fires on its own; we drive tick()
|
||||||
|
|
||||||
|
loop.onReplyQueued(WORKER);
|
||||||
|
loop.tick(PRIMARY); // first tick: nudges WORKER alone; WORKER's own count reaches the cap
|
||||||
|
assertEquals(1, rec.sendCount(), "the first tick should nudge about WORKER");
|
||||||
|
|
||||||
|
// WORKER2 "arrives during the backoff": queued for the same lead while the schedule from
|
||||||
|
// the tick above is still active (activeLeads still holds PRIMARY), before the next tick
|
||||||
|
// (simulated below) takes its own start-of-tick snapshot.
|
||||||
|
inbox.publish(WORKER2, "m2", "hello2");
|
||||||
|
loop.onReplyQueued(WORKER2);
|
||||||
|
|
||||||
|
loop.tick(PRIMARY); // the tick that would fire once that backoff elapsed
|
||||||
|
|
||||||
|
assertEquals(2, rec.sendCount(),
|
||||||
|
"WORKER2 was never named in any nudge yet and must still get one, even though "
|
||||||
|
+ "WORKER's own reminder count is already at the cap");
|
||||||
|
String secondNudge = rec.sentParams().get(1).getValue().toString();
|
||||||
|
assertTrue(secondNudge.contains(WORKER2), "the never-named target must be named: " + secondNudge);
|
||||||
|
|
||||||
|
// Criterion #3: isActive() must reflect that this lead still had a live nudge to give —
|
||||||
|
// the second tick took the INJECT branch, so the schedule stayed live rather than being
|
||||||
|
// torn down under WORKER2.
|
||||||
|
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh target");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void aTicketArrivingDuringTheBackoffGetsNudgedDespiteAnAlreadyCappedSibling() {
|
||||||
|
// Mirrors the reply-side test above for the ticket source.
|
||||||
|
int cap = 1;
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
var loop = loop(cap, 100_000);
|
||||||
|
|
||||||
|
loop.onTicketTerminal("task-1", WORKER, false);
|
||||||
|
loop.tick(PRIMARY); // first tick: nudges task-1 alone; its count reaches the cap
|
||||||
|
assertEquals(1, rec.sendCount(), "the first tick should nudge about task-1");
|
||||||
|
|
||||||
|
loop.onTicketTerminal("task-2", WORKER, false); // arrives during the backoff, same lead
|
||||||
|
loop.tick(PRIMARY);
|
||||||
|
|
||||||
|
assertEquals(2, rec.sendCount(),
|
||||||
|
"task-2 was never named in any nudge yet and must still get one, even though "
|
||||||
|
+ "task-1's reminder count is already at the cap");
|
||||||
|
String secondNudge = rec.sentParams().get(1).getValue().toString();
|
||||||
|
assertTrue(secondNudge.contains("task-2"), "the never-named ticket must be named: " + secondNudge);
|
||||||
|
assertTrue(loop.isActive(), "the schedule must stay active after nudging the fresh ticket");
|
||||||
|
}
|
||||||
|
|
||||||
|
// --- CB-582: bridge_ask question-open nudges -------------------------------------------------
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void onQuestionOpenedWithNoKnownLeadNeverStartsASchedule() throws Exception {
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 50);
|
||||||
|
|
||||||
|
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||||
|
|
||||||
|
assertFalse(loop.isActive(), "no lead known means nothing to nudge yet");
|
||||||
|
Thread.sleep(150);
|
||||||
|
assertEquals(0, rec.sendCount(), "must not nudge when no lead is known to be waiting");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void decideQuestionsWithNothingPendingIsStop() {
|
||||||
|
agents = agentWithStatus("idle");
|
||||||
|
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0, 0));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void decideQuestionsAtCapIsStop() {
|
||||||
|
agents = agentWithStatus("idle");
|
||||||
|
var loop = loop(2, 100_000);
|
||||||
|
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||||
|
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0, 2));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void decideQuestionsUnderCapWithInjectableLeadIsInject() {
|
||||||
|
agents = agentWithStatus("idle");
|
||||||
|
var loop = loop(5, 100_000);
|
||||||
|
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||||
|
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0, 0));
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void onQuestionOpenedCausesExactlyOneNudgeNamingTheTicketAndTurnId() throws Exception {
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
|
||||||
|
loop(1, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config file?");
|
||||||
|
|
||||||
|
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one question nudge should have been sent");
|
||||||
|
assertEquals(1, rec.sendCount());
|
||||||
|
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||||
|
assertTrue(nudge.contains("task-1"), "nudge should name the ticket: " + nudge);
|
||||||
|
assertTrue(nudge.contains("term_worker#1"), "nudge should name the turnId: " + nudge);
|
||||||
|
assertTrue(nudge.contains("bridge_send(turnId="), "nudge should name the exact answer call: " + nudge);
|
||||||
|
assertTrue(nudge.contains("which config file?"), "nudge should include the question text: " + nudge);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void questionClosedPreventsFurtherNudging() throws Exception {
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
var loop = loop(1, 100);
|
||||||
|
|
||||||
|
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||||
|
loop.questionClosed("term_worker#1"); // answered/lapsed before the first tick fired
|
||||||
|
|
||||||
|
Thread.sleep(300); // let the scheduled tick run
|
||||||
|
assertEquals(0, rec.sendCount(), "an already-closed question must never be nudged");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void questionNudgesSendUpToCapThenStop() throws Exception {
|
||||||
|
int cap = 2;
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
rec.sendLatch = new CountDownLatch(cap);
|
||||||
|
|
||||||
|
loop(cap, 50).onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||||
|
|
||||||
|
assertTrue(rec.sendLatch.await(5, TimeUnit.SECONDS), cap + " question nudges should have fired");
|
||||||
|
Thread.sleep(300);
|
||||||
|
assertEquals(cap, rec.sendCount(), "exactly " + cap + " question nudges (cap=" + cap + ")");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void questionAndTicketForTheSameLeadCoalesceIntoOneSend() throws Exception {
|
||||||
|
var rec = recordingClient();
|
||||||
|
agents = new AgentControl(rec);
|
||||||
|
var loop = loop(1, 300); // backoff wide enough that both entry points land before the first tick
|
||||||
|
|
||||||
|
loop.onTicketTerminal("task-1", WORKER, false);
|
||||||
|
loop.onQuestionOpened("task-2", WORKER, "term_worker#1", "which config?");
|
||||||
|
|
||||||
|
assertTrue(rec.sendLatch.await(3, TimeUnit.SECONDS), "one combined nudge should have been sent");
|
||||||
|
Thread.sleep(300);
|
||||||
|
assertEquals(1, rec.sendCount(),
|
||||||
|
"a ticket and a question for the same lead must coalesce onto ONE schedule");
|
||||||
|
String nudge = rec.sentParams().getFirst().getValue().toString();
|
||||||
|
assertTrue(nudge.contains("task-1"), "the combined nudge must still mention the ticket: " + nudge);
|
||||||
|
assertTrue(nudge.contains("term_worker#1"), "the combined nudge must still mention the question: " + nudge);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void oneExhaustedQuestionSourceDoesNotBlockANudgeForTheOtherSources() {
|
||||||
|
// Mirrors oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource for the question source:
|
||||||
|
// the question source is at its cap (2/2), but the ticket source has never been nudged
|
||||||
|
// (0/2) — decide() must still INJECT so the ticket is not stranded.
|
||||||
|
agents = agentWithStatus("idle");
|
||||||
|
var loop = loop(2, 100_000);
|
||||||
|
loop.onQuestionOpened("task-1", WORKER, "term_worker#1", "which config?");
|
||||||
|
loop.onTicketTerminal("task-2", WORKER, false);
|
||||||
|
|
||||||
|
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0, 2),
|
||||||
|
"the question source is exhausted (2/2), but the ticket source has never been "
|
||||||
|
+ "nudged (0/2) — the lead must still be injected so the ticket is not lost");
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void questionNudgeFormatIsCorrect() {
|
||||||
|
String single = ReplyPushLoop.QUESTION_NUDGE_FORMAT.formatted(
|
||||||
|
WORKER, "task-1", "term_worker#1", "which config?");
|
||||||
|
assertTrue(single.contains("Worker term_worker"));
|
||||||
|
assertTrue(single.contains("bridge_send(turnId=\"term_worker#1\""));
|
||||||
|
assertTrue(single.contains("which config?"));
|
||||||
|
|
||||||
|
String multi = ReplyPushLoop.QUESTIONS_NUDGE_FORMAT.formatted(2, "task-1 (turnId=t1), task-2 (turnId=t2)");
|
||||||
|
assertTrue(multi.contains("2 workers"));
|
||||||
|
assertTrue(multi.contains("bridge_poll(ticket=...)"));
|
||||||
|
}
|
||||||
|
|
||||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
|
|||||||
@@ -484,6 +484,67 @@ class BridgedAppTest {
|
|||||||
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
assertTrue(mapper.readTree(req(port, "GET", "/sessions/term_a/status").body()).get("ready").asBoolean());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-582: a lead polling {@code GET /sessions/{id}/status} on its normal cadence — not the
|
||||||
|
* ticket-scoped {@code /tasks/{ticket}} — must also see a worker's open async {@code bridge_ask}
|
||||||
|
* question, since the reverse-rendezvous window it opened with is far shorter than that cadence.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void sessionStatusReportsAnOpenQuestionWhenTheWorkerIsMidAsk() throws Exception {
|
||||||
|
FakeHerdr herdr = new FakeHerdr().agentStatus("idle"); // poller delivers the injection
|
||||||
|
int port = start(herdr, "http://gx00.gw:8000", Set.of("gx00.gw"));
|
||||||
|
|
||||||
|
HttpResponse<String> accepted = postMessage(port, "{\"content\":\"do it\",\"wait\":false}");
|
||||||
|
assertEquals(202, accepted.statusCode());
|
||||||
|
String ticket = mapper.readTree(accepted.body()).get("ticket").asText();
|
||||||
|
|
||||||
|
Thread.sleep(200); // let the background async send open its rendezvous waiter
|
||||||
|
|
||||||
|
var ask = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
|
||||||
|
try {
|
||||||
|
return postJson(port, "/sessions/term_a/ask",
|
||||||
|
"{\"question\":\"which config file?\",\"timeoutMs\":5000}");
|
||||||
|
} catch (Exception e) {
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
|
||||||
|
JsonNode task;
|
||||||
|
long deadline = System.currentTimeMillis() + 3000;
|
||||||
|
do {
|
||||||
|
task = mapper.readTree(req(port, "GET", "/tasks/" + ticket).body());
|
||||||
|
if ("asking".equals(task.path("phase").asText())) break;
|
||||||
|
//noinspection BusyWait
|
||||||
|
Thread.sleep(10);
|
||||||
|
} while (System.currentTimeMillis() < deadline);
|
||||||
|
assertEquals("asking", task.get("phase").asText());
|
||||||
|
String turnId = task.get("turnId").asText();
|
||||||
|
|
||||||
|
HttpResponse<String> status = req(port, "GET", "/sessions/term_a/status");
|
||||||
|
assertEquals(200, status.statusCode());
|
||||||
|
JsonNode body = mapper.readTree(status.body());
|
||||||
|
assertEquals("idle", body.get("status").asText(), "the live status must still be reported");
|
||||||
|
assertEquals("which config file?", body.get("question").asText());
|
||||||
|
assertEquals(turnId, body.get("turnId").asText());
|
||||||
|
assertEquals(ticket, body.get("ticket").asText());
|
||||||
|
|
||||||
|
// Answer it — via the same /message route bridge_send uses, keyed by turnId — so the
|
||||||
|
// background ask thread does not linger past the test.
|
||||||
|
var answer = java.util.concurrent.CompletableFuture.supplyAsync(() -> {
|
||||||
|
try {
|
||||||
|
return postJson(port, "/sessions/term_a/message",
|
||||||
|
"{\"content\":\"config.yaml\",\"turnId\":\"" + turnId + "\",\"timeoutMs\":4000}");
|
||||||
|
} catch (Exception e) {
|
||||||
|
throw new RuntimeException(e);
|
||||||
|
}
|
||||||
|
});
|
||||||
|
HttpResponse<String> askResponse = ask.get(6, java.util.concurrent.TimeUnit.SECONDS);
|
||||||
|
assertEquals(200, askResponse.statusCode());
|
||||||
|
assertEquals("config.yaml", mapper.readTree(askResponse.body()).get("answer").asText());
|
||||||
|
postJson(port, "/sessions/term_a/reply", "{\"content\":\"done\"}");
|
||||||
|
answer.get(6, java.util.concurrent.TimeUnit.SECONDS);
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
void stopWorkerInPanePlacementClosesOnlyThePane() throws Exception {
|
||||||
FakeHerdr herdr = new FakeHerdr();
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
|||||||
@@ -27,11 +27,15 @@ public final class FakeWorktrees implements Worktrees {
|
|||||||
public record SnapshotCall(String worktreePath, String branch, String message) {
|
public record SnapshotCall(String worktreePath, String branch, String message) {
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public record PruneCall(String repoRoot, long minAgeMillis) {
|
||||||
|
}
|
||||||
|
|
||||||
private final List<AddCall> addCalls = new CopyOnWriteArrayList<>();
|
private final List<AddCall> addCalls = new CopyOnWriteArrayList<>();
|
||||||
private final List<RemoveCall> removeCalls = new CopyOnWriteArrayList<>();
|
private final List<RemoveCall> removeCalls = new CopyOnWriteArrayList<>();
|
||||||
private final List<OverlayCall> overlayCalls = new CopyOnWriteArrayList<>();
|
private final List<OverlayCall> overlayCalls = new CopyOnWriteArrayList<>();
|
||||||
private final List<RepoRootCall> repoRootCalls = new CopyOnWriteArrayList<>();
|
private final List<RepoRootCall> repoRootCalls = new CopyOnWriteArrayList<>();
|
||||||
private final List<SnapshotCall> snapshotCalls = new CopyOnWriteArrayList<>();
|
private final List<SnapshotCall> snapshotCalls = new CopyOnWriteArrayList<>();
|
||||||
|
private final List<PruneCall> pruneCalls = new CopyOnWriteArrayList<>();
|
||||||
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
|
private final Set<String> existingPaths = ConcurrentHashMap.newKeySet();
|
||||||
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
|
private final Set<String> trackedPaths = ConcurrentHashMap.newKeySet();
|
||||||
private final AtomicLong snapshotSeq = new AtomicLong();
|
private final AtomicLong snapshotSeq = new AtomicLong();
|
||||||
@@ -40,6 +44,8 @@ public final class FakeWorktrees implements Worktrees {
|
|||||||
private volatile boolean dirty = false;
|
private volatile boolean dirty = false;
|
||||||
private volatile String repoRoot = "/repo";
|
private volatile String repoRoot = "/repo";
|
||||||
private volatile String prefix = "/worktrees";
|
private volatile String prefix = "/worktrees";
|
||||||
|
private volatile WipRefStats wipRefs = new WipRefStats(0, 0L);
|
||||||
|
private volatile int pruneResult = 0;
|
||||||
|
|
||||||
public FakeWorktrees withRepoRoot(String root) {
|
public FakeWorktrees withRepoRoot(String root) {
|
||||||
this.repoRoot = root;
|
this.repoRoot = root;
|
||||||
@@ -82,6 +88,18 @@ public final class FakeWorktrees implements Worktrees {
|
|||||||
return this;
|
return this;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Configure the value returned by {@link #wipRefs}. */
|
||||||
|
public FakeWorktrees withWipRefs(WipRefStats stats) {
|
||||||
|
this.wipRefs = stats;
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Configure the value returned by {@link #pruneWipRefs}. */
|
||||||
|
public FakeWorktrees withPruneResult(int deleted) {
|
||||||
|
this.pruneResult = deleted;
|
||||||
|
return this;
|
||||||
|
}
|
||||||
|
|
||||||
@Override
|
@Override
|
||||||
public String add(String repoRoot, String branch, String baseRef) {
|
public String add(String repoRoot, String branch, String baseRef) {
|
||||||
addCalls.add(new AddCall(repoRoot, branch, baseRef));
|
addCalls.add(new AddCall(repoRoot, branch, baseRef));
|
||||||
@@ -135,6 +153,17 @@ public final class FakeWorktrees implements Worktrees {
|
|||||||
return Optional.of("wip" + snapshotSeq.incrementAndGet());
|
return Optional.of("wip" + snapshotSeq.incrementAndGet());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public WipRefStats wipRefs(String repoRoot) {
|
||||||
|
return wipRefs;
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||||
|
pruneCalls.add(new PruneCall(repoRoot, minAgeMillis));
|
||||||
|
return pruneResult;
|
||||||
|
}
|
||||||
|
|
||||||
public List<AddCall> addCalls() {
|
public List<AddCall> addCalls() {
|
||||||
return List.copyOf(addCalls);
|
return List.copyOf(addCalls);
|
||||||
}
|
}
|
||||||
@@ -167,6 +196,10 @@ public final class FakeWorktrees implements Worktrees {
|
|||||||
return List.copyOf(snapshotCalls);
|
return List.copyOf(snapshotCalls);
|
||||||
}
|
}
|
||||||
|
|
||||||
|
public List<PruneCall> pruneCalls() {
|
||||||
|
return List.copyOf(pruneCalls);
|
||||||
|
}
|
||||||
|
|
||||||
public SnapshotCall lastSnapshot() {
|
public SnapshotCall lastSnapshot() {
|
||||||
return snapshotCalls.isEmpty() ? null : snapshotCalls.getLast();
|
return snapshotCalls.isEmpty() ? null : snapshotCalls.getLast();
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -3,6 +3,7 @@ package dev.ltms.bridged.session;
|
|||||||
import org.junit.jupiter.api.Test;
|
import org.junit.jupiter.api.Test;
|
||||||
import org.junit.jupiter.api.io.TempDir;
|
import org.junit.jupiter.api.io.TempDir;
|
||||||
|
|
||||||
|
import java.nio.charset.StandardCharsets;
|
||||||
import java.nio.file.Files;
|
import java.nio.file.Files;
|
||||||
import java.nio.file.Path;
|
import java.nio.file.Path;
|
||||||
import java.util.HashSet;
|
import java.util.HashSet;
|
||||||
@@ -138,6 +139,53 @@ class GitWorktreesTest {
|
|||||||
return out;
|
return out;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/** Write {@code content} as a blob into the object database; returns its sha. */
|
||||||
|
private static String blobOf(Path cwd, String content) throws Exception {
|
||||||
|
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "hash-object", "-w", "--stdin")
|
||||||
|
.redirectErrorStream(true).start();
|
||||||
|
p.getOutputStream().write(content.getBytes(StandardCharsets.UTF_8));
|
||||||
|
p.getOutputStream().close();
|
||||||
|
String out = new String(p.getInputStream().readAllBytes()).trim();
|
||||||
|
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git hash-object timed out");
|
||||||
|
assertEquals(0, p.exitValue(), "git hash-object failed:\n" + out);
|
||||||
|
return out;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Build a single-file tree object from {@code blob}; returns the tree's sha. */
|
||||||
|
private static String treeOf(Path cwd, String path, String blob) throws Exception {
|
||||||
|
Process p = new ProcessBuilder("git", "-C", cwd.toString(), "mktree")
|
||||||
|
.redirectErrorStream(true).start();
|
||||||
|
p.getOutputStream().write(("100644 blob " + blob + "\t" + path + "\n").getBytes(StandardCharsets.UTF_8));
|
||||||
|
p.getOutputStream().close();
|
||||||
|
String out = new String(p.getInputStream().readAllBytes()).trim();
|
||||||
|
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git mktree timed out");
|
||||||
|
assertEquals(0, p.exitValue(), "git mktree failed:\n" + out);
|
||||||
|
return out;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** {@code git commit-tree} rooted at {@code tree} with a chosen committer date; returns the sha. */
|
||||||
|
private static String commitTree(Path cwd, String tree, String parent, String committerDate,
|
||||||
|
String message) throws Exception {
|
||||||
|
ProcessBuilder pb = new ProcessBuilder("git", "-C", cwd.toString(), "commit-tree",
|
||||||
|
tree, "-p", parent, "-m", message);
|
||||||
|
pb.environment().put("GIT_COMMITTER_DATE", committerDate);
|
||||||
|
Process p = pb.redirectErrorStream(true).start();
|
||||||
|
String out = new String(p.getInputStream().readAllBytes()).trim();
|
||||||
|
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git commit-tree timed out");
|
||||||
|
assertEquals(0, p.exitValue(), "git commit-tree failed:\n" + out);
|
||||||
|
return out;
|
||||||
|
}
|
||||||
|
|
||||||
|
/** {@code git update-ref <ref> <sha>} — create the snapshot ref directly. */
|
||||||
|
private static void updateRef(Path cwd, String ref, String sha) throws Exception {
|
||||||
|
git(cwd, "update-ref", ref, sha);
|
||||||
|
}
|
||||||
|
|
||||||
|
/** True when {@code ref} exists in the repo (for-each-ref on a missing ref is empty, not an error). */
|
||||||
|
private static boolean refExists(Path cwd, String ref) throws Exception {
|
||||||
|
return !forEachRef(cwd, ref).trim().isEmpty();
|
||||||
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* The heart of CB-525: a provisioned worktree must not inherit the primary's MCP servers. Without
|
* The heart of CB-525: a provisioned worktree must not inherit the primary's MCP servers. Without
|
||||||
* the isolation step the checked-out {@code .mcp.json} carries them in, and a worker navigating
|
* the isolation step the checked-out {@code .mcp.json} carries them in, and a worker navigating
|
||||||
@@ -474,4 +522,104 @@ class GitWorktreesTest {
|
|||||||
assertThrows(WorktreeException.class, () -> gitWorktrees.snapshot(wt, branch, "test snapshot"),
|
assertThrows(WorktreeException.class, () -> gitWorktrees.snapshot(wt, branch, "test snapshot"),
|
||||||
"an unresolvable real index must fail loudly, not silently snapshot from an empty index");
|
"an unresolvable real index must fail loudly, not silently snapshot from an empty index");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586, criterion 2. A snapshot whose content is NOT reachable from {@code main} is the last
|
||||||
|
* copy of a worker's work, and must never be deleted automatically — even when it is old and
|
||||||
|
* even when the caller passes a zero age floor. Uses the real snapshot path on a dirty worktree,
|
||||||
|
* so the unreachable tree is exactly the shape CB-576/CB-578 stage C exist to protect.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void anUnreachableSnapshotIsNeverPruned(@TempDir Path tmp) throws Exception {
|
||||||
|
Path repo = initRepo(tmp.resolve("repo"));
|
||||||
|
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||||
|
String branch = "cb-586-unreachable";
|
||||||
|
String wt = gitWorktrees.add(repo.toString(), branch, "HEAD");
|
||||||
|
Files.writeString(Path.of(wt).resolve("worker-draft.txt"), "work that exists nowhere else\n");
|
||||||
|
|
||||||
|
Optional<String> ref = gitWorktrees.snapshot(wt, branch, "snapshot with unreachable content");
|
||||||
|
assertTrue(ref.isPresent());
|
||||||
|
|
||||||
|
// Age floor 0 makes age a non-issue: only reachability can save it — and it must.
|
||||||
|
assertEquals(0, gitWorktrees.pruneWipRefs(repo.toString(), 0),
|
||||||
|
"the unreachable snapshot is the last copy and must not be pruned");
|
||||||
|
assertTrue(refExists(repo, "refs/wip/" + branch),
|
||||||
|
"an unreachable snapshot must survive the sweep");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586, criterion 1 (the reachable half). A snapshot whose tree content IS already reachable
|
||||||
|
* from {@code main} and which is older than the age floor is pure duplication — the work is
|
||||||
|
* recovered — so it must be pruned.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void aReachableSnapshotOlderThanTheFloorIsPruned(@TempDir Path tmp) throws Exception {
|
||||||
|
Path repo = initRepo(tmp.resolve("repo"));
|
||||||
|
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||||
|
|
||||||
|
// A snapshot whose tree is exactly main's current tree: fully reachable from main.
|
||||||
|
String mainTree = revParse(repo, "main^{tree}");
|
||||||
|
String old = commitTree(repo, mainTree, revParse(repo, "HEAD"), "2020-01-01T00:00:00", "snapshot");
|
||||||
|
updateRef(repo, "refs/wip/recovered", old);
|
||||||
|
|
||||||
|
assertEquals(1, gitWorktrees.pruneWipRefs(repo.toString(), TimeUnit.HOURS.toMillis(24)),
|
||||||
|
"an old, main-reachable snapshot must be pruned");
|
||||||
|
assertFalse(refExists(repo, "refs/wip/recovered"),
|
||||||
|
"the reachable snapshot's ref must be gone after the sweep");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586, the age floor. A snapshot whose content IS reachable from {@code main} but which is
|
||||||
|
* younger than the age floor must not be swept — a lead may still be looking at it.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void aReachableButRecentSnapshotIsNotPruned(@TempDir Path tmp) throws Exception {
|
||||||
|
Path repo = initRepo(tmp.resolve("repo"));
|
||||||
|
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||||
|
|
||||||
|
// Reachable from main, but committed "now" — a fresh snapshot. The 24h floor must protect it.
|
||||||
|
String mainTree = revParse(repo, "main^{tree}");
|
||||||
|
String fresh = commitTree(repo, mainTree, revParse(repo, "HEAD"),
|
||||||
|
"2038-01-01T00:00:00", "snapshot just taken");
|
||||||
|
updateRef(repo, "refs/wip/fresh", fresh);
|
||||||
|
|
||||||
|
assertEquals(0, gitWorktrees.pruneWipRefs(repo.toString(), TimeUnit.HOURS.toMillis(24)),
|
||||||
|
"a recent snapshot must be kept even when reachable");
|
||||||
|
assertTrue(refExists(repo, "refs/wip/fresh"),
|
||||||
|
"the recent reachable snapshot must survive the sweep");
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586, criterion 5. A fleet that has never snapshotted anything has no {@code refs/wip/*},
|
||||||
|
* so a sweep is a no-op and the census reports none — identical to before CB-586 existed.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void aFleetWithNoSnapshotsPrunesNothingAndReportsNothing(@TempDir Path tmp) throws Exception {
|
||||||
|
Path repo = initRepo(tmp.resolve("repo"));
|
||||||
|
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||||
|
|
||||||
|
assertEquals(0, gitWorktrees.pruneWipRefs(repo.toString(), 0),
|
||||||
|
"no snapshot refs means nothing to prune");
|
||||||
|
Worktrees.WipRefStats stats = gitWorktrees.wipRefs(repo.toString());
|
||||||
|
assertEquals(0, stats.count(), "a never-snapshotted fleet has zero refs/wip refs");
|
||||||
|
assertEquals(0L, stats.costBytes(), "a never-snapshotted fleet costs zero bytes");
|
||||||
|
}
|
||||||
|
|
||||||
|
/** CB-586, criterion 4: the census reports how many refs exist and roughly what they cost. */
|
||||||
|
@Test
|
||||||
|
void wipRefsReportsCountAndCost(@TempDir Path tmp) throws Exception {
|
||||||
|
Path repo = initRepo(tmp.resolve("repo"));
|
||||||
|
GitWorktrees gitWorktrees = new GitWorktrees(tmp.resolve("wts").toString());
|
||||||
|
|
||||||
|
String blob = blobOf(repo, "a recoverable snapshot's worth of content");
|
||||||
|
String tree = treeOf(repo, "snapshot.txt", blob);
|
||||||
|
updateRef(repo, "refs/wip/one", commitTree(repo, tree, revParse(repo, "HEAD"),
|
||||||
|
"2020-01-01T00:00:00", "snapshot"));
|
||||||
|
updateRef(repo, "refs/wip/two", commitTree(repo, tree, revParse(repo, "HEAD"),
|
||||||
|
"2020-01-02T00:00:00", "snapshot"));
|
||||||
|
|
||||||
|
Worktrees.WipRefStats stats = gitWorktrees.wipRefs(repo.toString());
|
||||||
|
assertEquals(2, stats.count(), "two snapshot refs are reported");
|
||||||
|
assertTrue(stats.costBytes() > 0, "the cost of the snapshots is a positive byte count");
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -138,6 +138,16 @@ class SessionManagerTest {
|
|||||||
return java.util.Optional.of("wip" + snapshotSeq.incrementAndGet());
|
return java.util.Optional.of("wip" + snapshotSeq.incrementAndGet());
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public WipRefStats wipRefs(String repoRoot) {
|
||||||
|
return new WipRefStats(0, 0L);
|
||||||
|
}
|
||||||
|
|
||||||
|
@Override
|
||||||
|
public int pruneWipRefs(String repoRoot, long minAgeMillis) {
|
||||||
|
return 0;
|
||||||
|
}
|
||||||
|
|
||||||
List<String> removeCalls() {
|
List<String> removeCalls() {
|
||||||
return List.copyOf(removeCalls);
|
return List.copyOf(removeCalls);
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -11,10 +11,12 @@ import org.junit.jupiter.api.Test;
|
|||||||
import java.util.List;
|
import java.util.List;
|
||||||
import java.util.Map;
|
import java.util.Map;
|
||||||
import java.util.Set;
|
import java.util.Set;
|
||||||
|
import java.util.concurrent.TimeUnit;
|
||||||
import java.util.concurrent.atomic.AtomicLong;
|
import java.util.concurrent.atomic.AtomicLong;
|
||||||
|
|
||||||
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
|
||||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||||
|
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -90,6 +92,59 @@ class SessionReaperTest {
|
|||||||
return ticks.get() >= n;
|
return ticks.get() >= n;
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586: the retention sweep must actually reach the seam when the loop runs.
|
||||||
|
*
|
||||||
|
* <p>The sweep's own tests call {@code SessionManager.sweepWipRefs} directly, which walks
|
||||||
|
* around the reaper's interval gate — and the gate is where it broke. {@code lastWipSweepNanos}
|
||||||
|
* started at {@code Long.MIN_VALUE}, so {@code now - lastWipSweepNanos} overflowed to a large
|
||||||
|
* negative number, the gate read that as "swept moments ago", and it returned <em>before</em>
|
||||||
|
* the assignment that would have fixed the field. The sweep never ran once, for the life of
|
||||||
|
* the process, and every direct-call test still passed.
|
||||||
|
*
|
||||||
|
* <p>So this asserts through the loop: spawn a worktree session (which is what tells the
|
||||||
|
* manager which repo holds {@code refs/wip/*}), start the reaper, and require a real prune call
|
||||||
|
* to arrive at the fake.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void theLoopActuallyRunsTheWipRetentionSweep() throws InterruptedException {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt");
|
||||||
|
SessionManager sessions = new SessionManager(launcher(herdr), worktrees);
|
||||||
|
// Without a worktree session the repo is unknown and the sweep is a legitimate no-op, so
|
||||||
|
// this step is what makes the assertion below meaningful rather than vacuous.
|
||||||
|
sessions.acquire("ltms-local", null, "/caller/proj", null, new WorktreeRequest("cb-586", null));
|
||||||
|
assertTrue(worktrees.pruneCalls().isEmpty(), "nothing has swept before the reaper starts");
|
||||||
|
|
||||||
|
SessionReaper reaper = new SessionReaper(sessions, IDLE_TTL_SECONDS, SHORT_INTERVAL_MILLIS);
|
||||||
|
reaper.start();
|
||||||
|
try {
|
||||||
|
long deadline = System.currentTimeMillis() + 3000;
|
||||||
|
while (worktrees.pruneCalls().isEmpty() && System.currentTimeMillis() < deadline) {
|
||||||
|
Thread.sleep(10);
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
reaper.stop();
|
||||||
|
}
|
||||||
|
|
||||||
|
assertFalse(worktrees.pruneCalls().isEmpty(),
|
||||||
|
"the reaper loop must run the refs/wip retention sweep; it never reached the seam");
|
||||||
|
assertEquals("/repo", worktrees.pruneCalls().getFirst().repoRoot(),
|
||||||
|
"the sweep must target the repo the fleet's worktrees came from");
|
||||||
|
assertEquals(TimeUnit.HOURS.toMillis(24), worktrees.pruneCalls().getFirst().minAgeMillis(),
|
||||||
|
"the 24h age floor is the safety rule and must reach the seam intact");
|
||||||
|
}
|
||||||
|
|
||||||
|
/** As {@link #launcher()}, on a caller-supplied herdr so the test can inspect it. */
|
||||||
|
private static ClaudeCodeLauncher launcher(FakeHerdr herdr) {
|
||||||
|
BridgedConfig.Profile cfg = new BridgedConfig.Profile(
|
||||||
|
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN",
|
||||||
|
List.of("ccs", "ltms-local"), "tab", "bridged-workers",
|
||||||
|
"worker: {profile} #{n}", null, null, null);
|
||||||
|
return new ClaudeCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||||
|
new SubscriptionGuard(Set.of("gx00.gw")), Map.of(cfg.profile(), cfg), cfg.profile(), _ -> null);
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void stopIsIdempotentAndSafeBeforeStart() {
|
void stopIsIdempotentAndSafeBeforeStart() {
|
||||||
SessionReaper reaper = reaper();
|
SessionReaper reaper = reaper();
|
||||||
|
|||||||
@@ -422,6 +422,28 @@ class WorktreeSessionManagerTest {
|
|||||||
"acceptance criterion 6: a failed ticket's detail must carry the snapshot ref");
|
"acceptance criterion 6: a failed ticket's detail must carry the snapshot ref");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
void releaseNotifiesTheListenerWithTheAgentSessionId() {
|
||||||
|
// CB-584 (issue #65 criterion 5): a failed ticket's detail must also name the agent
|
||||||
|
// session, alongside worktree/branch/snapshot, so a lead can resume the conversation
|
||||||
|
// rather than only re-dispatch a fresh member onto the same files.
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||||
|
.withDirty(true);
|
||||||
|
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||||
|
java.util.List<SessionManager.ReleaseDetail> released = new java.util.concurrent.CopyOnWriteArrayList<>();
|
||||||
|
sessions.onRelease(released::add);
|
||||||
|
MemberSession s = sessions.acquire("ltms-local", MemberRole.DEV, null, "/caller/proj", null,
|
||||||
|
new WorktreeRequest("cb-584-e", null), "cb-584-session", null);
|
||||||
|
|
||||||
|
assertNotNull(s.agentSessionId(), "a named session must mint an agent session id to assert on");
|
||||||
|
sessions.release(s.paneId());
|
||||||
|
|
||||||
|
assertEquals(1, released.size());
|
||||||
|
SessionManager.ReleaseDetail detail = released.getFirst();
|
||||||
|
assertEquals(s.agentSessionId(), detail.agentSessionId());
|
||||||
|
}
|
||||||
|
|
||||||
@Test
|
@Test
|
||||||
void aFailingSnapshotStillPreservesTheWorktreeStopsThePaneAndNotifies() {
|
void aFailingSnapshotStillPreservesTheWorktreeStopsThePaneAndNotifies() {
|
||||||
FakeHerdr herdr = new FakeHerdr();
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
@@ -444,4 +466,37 @@ class WorktreeSessionManagerTest {
|
|||||||
"a failed snapshot leaves no ref to report");
|
"a failed snapshot leaves no ref to report");
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* CB-586, criteria 4 and 1. Once a worktree session is spawned the repo is known, so the
|
||||||
|
* operator census and the retention sweep delegate to that repo's {@code refs/wip/*}.
|
||||||
|
*/
|
||||||
|
@Test
|
||||||
|
void wipRefsAndSweepDelegateToTheFleetRepoOnceKnown() {
|
||||||
|
FakeHerdr herdr = new FakeHerdr();
|
||||||
|
FakeWorktrees worktrees = new FakeWorktrees().withRepoRoot("/repo").withPrefix("/wt")
|
||||||
|
.withWipRefs(new Worktrees.WipRefStats(3, 42L)).withPruneResult(2);
|
||||||
|
SessionManager sessions = new SessionManager(workerService(herdr), worktrees);
|
||||||
|
|
||||||
|
assertTrue(sessions.wipRefs().isEmpty(),
|
||||||
|
"no worktree spawned yet means no repo is known and nothing to report");
|
||||||
|
assertEquals(0, sessions.sweepWipRefs(TimeUnit.HOURS.toMillis(24)),
|
||||||
|
"no worktree spawned yet means the sweep is a no-op");
|
||||||
|
|
||||||
|
sessions.acquire("ltms-local", null, "/caller/proj", null,
|
||||||
|
new WorktreeRequest("cb-586", null));
|
||||||
|
|
||||||
|
Worktrees.WipRefStats stats = sessions.wipRefs().orElseThrow();
|
||||||
|
assertEquals(3, stats.count(), "the census comes from the fleet repo");
|
||||||
|
assertEquals(42L, stats.costBytes(), "the cost comes from the fleet repo");
|
||||||
|
assertEquals(2, sessions.sweepWipRefs(TimeUnit.HOURS.toMillis(24)),
|
||||||
|
"the sweep runs against the fleet repo");
|
||||||
|
|
||||||
|
List<FakeWorktrees.PruneCall> prunes = worktrees.pruneCalls();
|
||||||
|
assertEquals(1, prunes.size(), "the no-op short-circuits before reaching the seam, so only "
|
||||||
|
+ "the repo-known sweep issues a call");
|
||||||
|
assertEquals("/repo", prunes.getFirst().repoRoot(), "the sweep targets the fleet repo");
|
||||||
|
assertEquals(TimeUnit.HOURS.toMillis(24), prunes.getFirst().minAgeMillis(),
|
||||||
|
"the caller's age floor is passed through");
|
||||||
|
}
|
||||||
|
|
||||||
}
|
}
|
||||||
|
|||||||
+19
-4
@@ -1,9 +1,24 @@
|
|||||||
# MCP Contract — `bridged`'s unified gateway
|
# MCP Contract — `bridged`'s unified gateway
|
||||||
|
|
||||||
> **Status:** 🟡 Design (2026-07-14). Greenfield — no MCP code exists yet; the pom carries
|
> **Status: 🔴 HISTORICAL DESIGN — do NOT use as the tool reference.** Written 2026-07-14, before
|
||||||
> only Javalin/Jackson. This page defines the tool surface that CB-104 and its followers
|
> any MCP code existed. The system shipped and this page never caught up, so **its tool names,
|
||||||
> implement. It supersedes nothing; it fills the "MCP server face" left open by the
|
> parameter names and REST paths are wrong today**. Audited 2026-08-17; the specific drift:
|
||||||
> [Architecture](1-Architecture) page.
|
>
|
||||||
|
> - **Tools it names that do not exist:** `bridge_read`, `bridge_cancel`.
|
||||||
|
> - **Shipped tools it omits:** `bridge_poll`, `bridge_ack`, `bridge_profiles`, `bridge_whoami`.
|
||||||
|
> - **Parameter names are wrong nearly everywhere** — it says `message`/`target`/`timeout_seconds`/
|
||||||
|
> `block` where the code takes `content`/`sessionId`/`timeoutMs`/`wait`; `text` where
|
||||||
|
> `bridge_reply` takes `content`; `target` where `bridge_stop` takes `paneId`.
|
||||||
|
> - **REST paths are wrong:** it says `POST /workers` and `DELETE /workers/{paneId}`; the daemon
|
||||||
|
> serves `POST /members` and `DELETE /members/{paneId}`.
|
||||||
|
>
|
||||||
|
> **The authoritative tool surface is the live MCP schema** (each tool's own description and
|
||||||
|
> parameters, as mounted), with the intent→tool table in `CLAUDE.md` as the short form. Both were
|
||||||
|
> checked against `mcp/BridgeMcp.java` on 2026-08-17 and are accurate.
|
||||||
|
>
|
||||||
|
> What is still worth reading here is **§6 — the flows and the error model** (rendezvous,
|
||||||
|
> `bridge_ask`, detached delivery, the turn-done fallback). The shapes it describes are the ones
|
||||||
|
> that shipped; only the names around them drifted. Rewriting this page is tracked as **CB-609**.
|
||||||
|
|
||||||
`bridged` is the **sole communication gateway** for every Claude session in the bridge. Both
|
`bridged` is the **sole communication gateway** for every Claude session in the bridge. Both
|
||||||
the **primary** (Opus, on subscription) and every **worker** (off-subscription Claude Code)
|
the **primary** (Opus, on subscription) and every **worker** (off-subscription Claude Code)
|
||||||
|
|||||||
+2
-2
@@ -17,7 +17,7 @@ who the workers are, how the lead picks one, and how it runs many at once.
|
|||||||
- **Workers** — a herd of `claude` panes in herdr, each an addressable `bridged` session
|
- **Workers** — a herd of `claude` panes in herdr, each an addressable `bridged` session
|
||||||
with its **own model/env**:
|
with its **own model/env**:
|
||||||
- **Claude workers** (clean env, e.g. Sonnet) — reasoning-heavy or high-accuracy subtasks.
|
- **Claude workers** (clean env, e.g. Sonnet) — reasoning-heavy or high-accuracy subtasks.
|
||||||
- **Local workers** (`ANTHROPIC_BASE_URL=https://ollama.ltms.dev`) — bulk, cheap, or
|
- **Local workers** (`ANTHROPIC_BASE_URL=https://llm.ltms.dev/anthropic`) — bulk, cheap, or
|
||||||
embarrassingly parallel subtasks.
|
embarrassingly parallel subtasks.
|
||||||
|
|
||||||
Every worker is still a *real Claude Code process* (inherits `CLAUDE.md`, hooks, skills,
|
Every worker is still a *real Claude Code process* (inherits `CLAUDE.md`, hooks, skills,
|
||||||
@@ -35,7 +35,7 @@ flowchart TB
|
|||||||
WL1["w-local-1<br/>ANTHROPIC_BASE_URL set"]
|
WL1["w-local-1<br/>ANTHROPIC_BASE_URL set"]
|
||||||
WL2["w-local-2<br/>ANTHROPIC_BASE_URL set"]
|
WL2["w-local-2<br/>ANTHROPIC_BASE_URL set"]
|
||||||
ANT["api.anthropic.com<br/>(Pro/Max)"]
|
ANT["api.anthropic.com<br/>(Pro/Max)"]
|
||||||
OLL["ollama.ltms.dev<br/>(local model)"]
|
OLL["llm.ltms.dev<br/>(gateway to the local model)"]
|
||||||
|
|
||||||
LEAD -->|"blocking POST /message (target role)"| BD
|
LEAD -->|"blocking POST /message (target role)"| BD
|
||||||
BD -->|"Unix socket · send_text · events.subscribe"| HERDR
|
BD -->|"Unix socket · send_text · events.subscribe"| HERDR
|
||||||
|
|||||||
Executable
+135
@@ -0,0 +1,135 @@
|
|||||||
|
#!/usr/bin/env bash
|
||||||
|
#
|
||||||
|
# CB-596 step 1: measure which credentials a live member actually holds.
|
||||||
|
#
|
||||||
|
# The claim under test is that a member's herdr pane starts a LOGIN shell, that shell sources
|
||||||
|
# ${SHARED_ENV}/tools/secrets.sh, and so a member inherits every name that file exports — while
|
||||||
|
# CB-592 blocks exactly one of them (GITEA_ACCESS_TOKEN). That is an inference from the code, not a
|
||||||
|
# measurement, and issue #82 says plainly: do not build a fix on the inference. This is the
|
||||||
|
# measurement.
|
||||||
|
#
|
||||||
|
# WHY THIS IS A SCRIPT AND NOT A COMMAND SOMEONE TYPES
|
||||||
|
#
|
||||||
|
# Enumerating credential names inside a member is exactly the action that should need the operator's
|
||||||
|
# explicit approval, and the command classifier refuses it. That refusal is correct. This script is
|
||||||
|
# the seam: it is one auditable file the operator can read once, top to bottom, and then run — rather
|
||||||
|
# than approving an ad-hoc shell pipeline whose behaviour they have to take on trust.
|
||||||
|
#
|
||||||
|
# WHAT IT WILL NOT DO
|
||||||
|
#
|
||||||
|
# * It never prints a credential value, and never any prefix or suffix of one. Not one character.
|
||||||
|
# Issue #82's criterion 1 asked for a 6-character prefix; this prints a truncated SHA-256 instead.
|
||||||
|
# A prefix of a short secret is most of the secret, and it would end up pasted into a ticket. The
|
||||||
|
# hash answers every question the prefix was for — is it set, is it the same value as over there,
|
||||||
|
# is it the CB-592 sentinel — and answers none of the ones it should not.
|
||||||
|
# * It never writes anywhere, never contacts the network, and never touches secrets.sh, which is
|
||||||
|
# the operator's file.
|
||||||
|
#
|
||||||
|
# HOW TO RUN IT
|
||||||
|
#
|
||||||
|
# 1. As the operator, in a member's pane (a spawned worker's terminal):
|
||||||
|
# bash scripts/probe-member-credentials.sh
|
||||||
|
# 2. For the comparison row, in your OWN shell — a lead, not a member:
|
||||||
|
# bash scripts/probe-member-credentials.sh --allow-outside-member
|
||||||
|
#
|
||||||
|
# The two outputs side by side are the finding: any name whose hash matches between them is a
|
||||||
|
# credential the member holds in full.
|
||||||
|
#
|
||||||
|
set -uo pipefail
|
||||||
|
|
||||||
|
# The names ${SHARED_ENV}/tools/secrets.sh exports, recorded on 2026-08-16 (issue #82). Names only —
|
||||||
|
# this list contains no values and never should. If secrets.sh gains a name, this list goes stale and
|
||||||
|
# the probe silently stops asking about it; that staleness is itself part of what #82's criterion 4
|
||||||
|
# has to solve, so it is called out in the summary rather than hidden.
|
||||||
|
NAMES=(
|
||||||
|
AI_GATEWAY_TOKEN BESZEL_ADMIN_EMAIL BESZEL_ADMIN_PASSWORD
|
||||||
|
BESZEL_HUB_URL BESZEL_KEY BESZEL_UNIVERSAL_TOKEN
|
||||||
|
BRAIN_MCP_TOKEN CF_ACCOUNT_ID CF_API_TOKEN
|
||||||
|
CF_USER_TOKEN CONFLUENCE_API_TOKEN CONFLUENCE_USERNAME
|
||||||
|
CONTEXT7_TOKEN GITEA_HOST GITLAB_OAUTH_CLIENT_SECRET
|
||||||
|
GITLAB_PERSONAL_ACCESS_TOKEN GRAFANA_ADMIN_PASSWORD GRAFANA_ADMIN_USER
|
||||||
|
HASS_TOKEN HW_PASSWORD HW_USER
|
||||||
|
LTMS_API_KEY MEMORY_MCP_TOKEN METRICS_PUSH_TOKEN
|
||||||
|
OPENCODE_AUTOMODE_MODEL TELEGRAM_BOT_TOKEN TELEGRAM_CHAT_ID
|
||||||
|
TS_API_KEY TS_AUTHKEY WORKER_GITEA_TOKEN
|
||||||
|
GITEA_ACCESS_TOKEN
|
||||||
|
)
|
||||||
|
|
||||||
|
allow_outside=0
|
||||||
|
for arg in "$@"; do
|
||||||
|
case "$arg" in
|
||||||
|
--allow-outside-member) allow_outside=1 ;;
|
||||||
|
-h|--help) sed -n '2,40p' "$0"; exit 0 ;;
|
||||||
|
*) echo "unknown argument: $arg" >&2; exit 2 ;;
|
||||||
|
esac
|
||||||
|
done
|
||||||
|
|
||||||
|
if [ "${BRIDGED_MEMBER:-}" != "1" ] && [ "$allow_outside" -eq 0 ]; then
|
||||||
|
cat >&2 <<'EOF'
|
||||||
|
refusing to run: BRIDGED_MEMBER is not 1, so this is not a member's shell.
|
||||||
|
|
||||||
|
The finding this probe exists for is what a MEMBER holds. Run it in a spawned worker's pane. If you
|
||||||
|
meant to take the comparison reading from your own shell, pass --allow-outside-member and the output
|
||||||
|
will be labelled as such.
|
||||||
|
EOF
|
||||||
|
exit 1
|
||||||
|
fi
|
||||||
|
|
||||||
|
# Prefer sha256sum (Linux), fall back to shasum (macOS). If neither exists, report presence and
|
||||||
|
# length only — degraded, but never a value.
|
||||||
|
hasher=""
|
||||||
|
if command -v sha256sum >/dev/null 2>&1; then
|
||||||
|
hasher="sha256sum"
|
||||||
|
elif command -v shasum >/dev/null 2>&1; then
|
||||||
|
hasher="shasum -a 256"
|
||||||
|
fi
|
||||||
|
|
||||||
|
digest() { # value -> first 12 hex chars of its sha256, or "-" when no hasher is available
|
||||||
|
[ -z "$hasher" ] && { printf '%s' "-"; return; }
|
||||||
|
printf '%s' "$1" | $hasher | cut -c1-12
|
||||||
|
}
|
||||||
|
|
||||||
|
if [ "${BRIDGED_MEMBER:-}" = "1" ]; then
|
||||||
|
where="MEMBER (BRIDGED_MEMBER=1)"
|
||||||
|
else
|
||||||
|
where="NOT a member — comparison reading only"
|
||||||
|
fi
|
||||||
|
|
||||||
|
echo "CB-596 credential probe"
|
||||||
|
echo "reading from : $where"
|
||||||
|
echo "shell : ${SHELL:-unknown}"
|
||||||
|
echo "hash : ${hasher:-none available — lengths only}"
|
||||||
|
# Only printed so the two readings can be told apart when they are pasted side by side.
|
||||||
|
echo "host : $(hostname 2>/dev/null || echo unknown)"
|
||||||
|
echo
|
||||||
|
printf '%-30s %-7s %6s %s\n' "NAME" "STATE" "LEN" "SHA256-12"
|
||||||
|
printf '%-30s %-7s %6s %s\n' "------------------------------" "-------" "------" "------------"
|
||||||
|
|
||||||
|
set_count=0
|
||||||
|
for name in "${NAMES[@]}"; do
|
||||||
|
value="${!name:-}"
|
||||||
|
if [ -z "$value" ]; then
|
||||||
|
printf '%-30s %-7s %6s %s\n' "$name" "unset" "-" "-"
|
||||||
|
else
|
||||||
|
set_count=$((set_count + 1))
|
||||||
|
printf '%-30s %-7s %6s %s\n' "$name" "SET" "${#value}" "$(digest "$value")"
|
||||||
|
fi
|
||||||
|
done
|
||||||
|
|
||||||
|
echo
|
||||||
|
echo "$set_count of ${#NAMES[@]} names are set in this shell."
|
||||||
|
echo
|
||||||
|
cat <<'EOF'
|
||||||
|
How to read this:
|
||||||
|
|
||||||
|
* Take the MEMBER reading and the comparison reading, and line them up. A name whose SHA256-12
|
||||||
|
matches on both sides is a credential the member holds in full. That is the finding.
|
||||||
|
* GITEA_ACCESS_TOKEN is the control. CB-592 replaces it with a blocked sentinel, so its hash
|
||||||
|
should DIFFER between the two readings. If it matches, CB-592 is not working and that is the
|
||||||
|
most urgent thing on this page.
|
||||||
|
* AI_GATEWAY_TOKEN matching is expected and correct, not a leak: bridged.yaml names it in
|
||||||
|
`tokenEnv:` for the local and gx profiles, so a member reaching the gateway is by design.
|
||||||
|
* A name that is set here but is NOT in the list above will not appear at all. The list was
|
||||||
|
recorded on 2026-08-16 and does not update itself. Anything added to secrets.sh since then is
|
||||||
|
invisible to this probe — which is the same gap issue #82 criterion 4 asks to close properly.
|
||||||
|
EOF
|
||||||
Reference in New Issue
Block a user