Compare commits
9 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| dd906526c0 | |||
| ec3001796a | |||
| 86cf4c285a | |||
| cd18887b69 | |||
| 4637c68295 | |||
| 0114bd1fa7 | |||
| 6dee84ca71 | |||
| f004a0c654 | |||
| 6123576c68 |
@@ -203,17 +203,17 @@ public final class Bridged {
|
||||
leads = () -> leadTerminals;
|
||||
}
|
||||
|
||||
// CB-548: config-declared architect slots. Slots live in config (name → strong-model
|
||||
// profile); the terminal → slot binding is the live half, sourced from the slots' declared
|
||||
// terminals today and swapped for a live binding by the later spawn lifecycle. The registry
|
||||
// is what CallerResolver resolves against and what that lifecycle will read profiles from;
|
||||
// CB-548: config-declared architect slots. Config supplies only the stable name → profile
|
||||
// map; the terminal → slot binding is owned by the registry and is empty at startup, so no
|
||||
// pane resolves to an architect until the later spawn lifecycle binds one. The registry is
|
||||
// what CallerResolver resolves against and what that lifecycle will read profiles from;
|
||||
// nothing here spawns a slot.
|
||||
ArchitectRegistry architects = new ArchitectRegistry(
|
||||
cfg.architects() == null ? Map.of() : cfg.architects(),
|
||||
() -> cfg.architectTerminals());
|
||||
cfg.architects() == null ? Map.of() : cfg.architects());
|
||||
if (!architects.slots().isEmpty()) {
|
||||
log.info("architect slots: {} configured {}, terminals {}", architects.slots().size(),
|
||||
architects.slots().keySet(), cfg.architectTerminals().keySet());
|
||||
log.info("architect slots: {} configured {} — none bound yet (a slot is idle until the "
|
||||
+ "spawn lifecycle binds a live terminal to it)",
|
||||
architects.slots().size(), architects.slots().keySet());
|
||||
}
|
||||
|
||||
// Status-gated injector (CB-103): the single writer into workers, fed by a poller.
|
||||
@@ -326,12 +326,12 @@ public final class Bridged {
|
||||
+ " is unset or empty — export it before starting bridged");
|
||||
}
|
||||
callers = CallerResolver.withLeadsAndArchitects(identity, true, token, leads,
|
||||
architects::terminalBindings);
|
||||
architects::snapshot);
|
||||
log.info("auth: token mode (bearer required for non-worker callers, env {})",
|
||||
cfg.auth().tokenEnv());
|
||||
} else {
|
||||
callers = CallerResolver.withLeadsAndArchitects(identity, false, null, leads,
|
||||
architects::terminalBindings);
|
||||
architects::snapshot);
|
||||
log.info("auth: loopback-trust (any loopback non-worker caller is the primary)");
|
||||
}
|
||||
|
||||
|
||||
@@ -2,37 +2,38 @@ package dev.ltms.bridged.auth;
|
||||
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.function.Supplier;
|
||||
|
||||
/**
|
||||
* The architect-slot registry (CB-548): every gateway-local architect name and the strong-model
|
||||
* profile it points at, plus the live binding from a live architect's herdr terminal to its slot.
|
||||
* profile it points at, plus the <em>live</em> bindings from a live architect's herdr terminal to
|
||||
* its slot.
|
||||
*
|
||||
* <p>Two halves, split by who owns each:
|
||||
* <ul>
|
||||
* <li><b>slots</b> — configured once, keyed by the gateway-local unique name; each carries the
|
||||
* {@code profile} reference the <em>future</em> spawn lifecycle will read when it stands the
|
||||
* slot up. A read-only snapshot taken at construction.</li>
|
||||
* <li><b>terminal bindings</b> — a {@link Supplier} consulted on every read, so a binding
|
||||
* injected <em>after</em> startup (an operator pin, or the later lifecycle once it spawns a
|
||||
* session) takes effect without a restart. {@link CallerResolver} reads this to turn a pane
|
||||
* into an {@link Role#ARCHITECT}.</li>
|
||||
* {@code profile} reference the spawn lifecycle reads when it stands the slot up. A read-only
|
||||
* snapshot taken at construction.</li>
|
||||
* <li><b>terminal bindings</b> — owned by this registry and initially <em>empty</em>. Config
|
||||
* declares no architect terminal, so at startup every slot is idle and nothing resolves to an
|
||||
* architect; a session only becomes one when the spawn lifecycle {@linkplain #bind(String,
|
||||
* String) binds} its terminal to a slot. {@link CallerResolver} reads this through
|
||||
* {@link #snapshot()} to turn a pane into an {@link Role#ARCHITECT}.</li>
|
||||
* </ul>
|
||||
*
|
||||
* <p>Spawning/lifecycle is deliberately a separate unit: this class only exposes the map the
|
||||
* resolver resolves against and the profile lookup that lifecycle will call. Nothing here
|
||||
* creates or manages an architect session.
|
||||
* <p>Spawning/lifecycle is deliberately a separate unit: this class only owns the bindings and
|
||||
* exposes the map the resolver resolves against plus the profile lookup lifecycle will call.
|
||||
* Nothing here creates or manages an architect session.
|
||||
*/
|
||||
public final class ArchitectRegistry {
|
||||
|
||||
private final Map<String, BridgedConfig.Architect> slots;
|
||||
private final Supplier<Map<String, String>> terminalBindings;
|
||||
/** Live {@code terminal_id → slot name}; guarded by {@code this}. */
|
||||
private final Map<String, String> terminalToSlot = new HashMap<>();
|
||||
|
||||
public ArchitectRegistry(Map<String, BridgedConfig.Architect> slots,
|
||||
Supplier<Map<String, String>> terminalBindings) {
|
||||
public ArchitectRegistry(Map<String, BridgedConfig.Architect> slots) {
|
||||
this.slots = slots == null ? Map.of() : Map.copyOf(slots);
|
||||
this.terminalBindings = terminalBindings == null ? Map::of : terminalBindings;
|
||||
}
|
||||
|
||||
/** The configured slots, keyed by gateway-local unique name. Unmodifiable snapshot. */
|
||||
@@ -41,22 +42,30 @@ public final class ArchitectRegistry {
|
||||
}
|
||||
|
||||
/**
|
||||
* The live {@code terminal_id → slot name} bindings, re-read on every call.
|
||||
* An immutable copy of the live {@code terminal_id → slot name} bindings.
|
||||
*
|
||||
* <p>Passed to {@link CallerResolver} as the source of architect identity, and what
|
||||
* {@code bridge_whoami}/the roster will read to say which slot a pane hosts.
|
||||
* {@code bridge_whoami}/the roster will read to say which slot a pane hosts. Empty until the
|
||||
* spawn lifecycle binds a slot.
|
||||
*/
|
||||
public Map<String, String> terminalBindings() {
|
||||
return terminalBindings.get();
|
||||
public Map<String, String> snapshot() {
|
||||
synchronized (terminalToSlot) {
|
||||
return Map.copyOf(terminalToSlot);
|
||||
}
|
||||
}
|
||||
|
||||
/** The slot a live terminal is bound to, or {@code null} if it is no architect slot. */
|
||||
/** The slot a live terminal is bound to, or {@code null} if it is not an architect slot. */
|
||||
public String slotForTerminal(String terminal) {
|
||||
return terminal == null ? null : terminalBindings.get().get(terminal);
|
||||
if (terminal == null) {
|
||||
return null;
|
||||
}
|
||||
synchronized (terminalToSlot) {
|
||||
return terminalToSlot.get(terminal);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* The strong-model profile a slot runs under — what the future spawn lifecycle reads.
|
||||
* The strong-model profile a slot runs under — what the spawn lifecycle reads.
|
||||
*
|
||||
* @return the slot's configured {@code profile}, or {@code null} if the slot is unknown or
|
||||
* declares none
|
||||
@@ -70,4 +79,64 @@ public final class ArchitectRegistry {
|
||||
public boolean isSlot(String slotName) {
|
||||
return slots.containsKey(slotName);
|
||||
}
|
||||
|
||||
/**
|
||||
* Bind {@code terminal} to {@code slot} (CB-548).
|
||||
*
|
||||
* <p>The spawn lifecycle calls this when it stands a slot up. The bind is atomic and preserves
|
||||
* the two cardinality invariants: a terminal may occupy at most one slot, and a slot may host at
|
||||
* most one terminal. Binding the same terminal to the same slot again is a harmless no-op.
|
||||
*
|
||||
* @param slot a configured slot name, or the bind is refused
|
||||
* @param terminal the pane that will act as this architect
|
||||
* @return {@code true} if the binding is now {@code terminal → slot}; {@code false} if it was
|
||||
* refused — an unknown slot, a terminal already bound to a different slot, or a slot
|
||||
* already hosting a different terminal
|
||||
*/
|
||||
public boolean bind(String slot, String terminal) {
|
||||
if (slot == null || terminal == null || terminal.isBlank()) {
|
||||
return false;
|
||||
}
|
||||
synchronized (terminalToSlot) {
|
||||
if (!isSlot(slot)) {
|
||||
return false; // unknown slot — nothing to bind to
|
||||
}
|
||||
String existingSlot = terminalToSlot.get(terminal);
|
||||
if (existingSlot != null) {
|
||||
return slot.equals(existingSlot); // already this slot (idempotent) or a different one
|
||||
}
|
||||
if (terminalToSlot.containsValue(slot)) {
|
||||
return false; // slot already hosts a terminal — no second one
|
||||
}
|
||||
terminalToSlot.put(terminal, slot);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Compare-safe unbind of {@code expectedTerminal} from {@code slot} (CB-548).
|
||||
*
|
||||
* <p>The spawn lifecycle calls this when it tears a slot down. Only the exact binding
|
||||
* {@code expectedTerminal → slot} is removed; if that terminal was since rebound to a different
|
||||
* slot (or the slot to a different terminal), the call is a no-op returning {@code false} — a
|
||||
* stale unbind must never remove a replacement.
|
||||
*
|
||||
* @param slot the slot the caller believes the terminal is bound to
|
||||
* @param expectedTerminal the terminal it expects to be bound there
|
||||
* @return {@code true} if {@code expectedTerminal → slot} was removed; {@code false} if nothing
|
||||
* was (no such binding, or the binding had already moved)
|
||||
*/
|
||||
public boolean unbind(String slot, String expectedTerminal) {
|
||||
if (slot == null || expectedTerminal == null) {
|
||||
return false;
|
||||
}
|
||||
synchronized (terminalToSlot) {
|
||||
String current = terminalToSlot.get(expectedTerminal);
|
||||
if (current == null || !slot.equals(current)) {
|
||||
return false; // absent, or a replacement/moved binding — leave it in place
|
||||
}
|
||||
terminalToSlot.remove(expectedTerminal);
|
||||
return true;
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,6 +1,8 @@
|
||||
package dev.ltms.bridged.config;
|
||||
|
||||
import com.fasterxml.jackson.annotation.JsonIgnoreProperties;
|
||||
import com.fasterxml.jackson.core.JsonParser;
|
||||
import com.fasterxml.jackson.core.JsonToken;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
||||
import org.slf4j.Logger;
|
||||
@@ -11,6 +13,7 @@ import java.io.UncheckedIOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.Collections;
|
||||
import java.util.HashSet;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
@@ -44,10 +47,10 @@ import java.util.Set;
|
||||
* lead name; supersedes the singular {@code primary} pin, which stays honoured.
|
||||
* See {@link #leaderTerminals()} for how the two merge
|
||||
* @param architects CB-548 architect slots, keyed by gateway-local unique slot name; each points
|
||||
* at a strong-model profile, and the identity a live session is matched by is
|
||||
* its {@code terminal} binding (see {@link #architectTerminals()}). A slot is
|
||||
* the hook the future spawn lifecycle reads a profile back from — nothing here
|
||||
* spawns it.
|
||||
* at a strong-model profile the future spawn lifecycle reads back. An architect
|
||||
* is <em>not</em> recognised like a lead: config declares the slots only, and a
|
||||
* live session becomes an architect when the spawn lifecycle binds its terminal
|
||||
* to a slot. Nothing here spawns a slot.
|
||||
* @param leadScan opt-in discovery of leads by tab label (CB-531); {@code null} ⇒ no scanning,
|
||||
* and only {@code leaders:}/{@code primary:} name a lead
|
||||
* @param placement how to choose a worker profile for an unqualified spawn:
|
||||
@@ -389,26 +392,23 @@ public record BridgedConfig(
|
||||
* One entry of the CB-548 {@code architects:} registry — a gateway-local named slot that points
|
||||
* at a strong-model profile.
|
||||
*
|
||||
* <p>A lead and an architect differ in <em>authority</em>, not in how identity is established:
|
||||
* both are recognised by configuration rather than spawned. A lead resolves to
|
||||
* {@link dev.ltms.bridged.auth.Role#PRIMARY} and owns the whole lifecycle (spawn/stop/drain);
|
||||
* an architect resolves to {@link dev.ltms.bridged.auth.Role#ARCHITECT}, which delegates turns
|
||||
* ({@code SEND}) and replies/asks as its own pane but cannot stand up or tear down workers —
|
||||
* lifecycle stays in one pair of hands.
|
||||
* <p>A slot is <em>declared</em>, not recognised: config names the slot and the profile it runs,
|
||||
* and nothing else. Unlike a lead (which config pins by herdr {@code terminal_id} and is
|
||||
* recognised at startup), an architect slot is idle at boot — config supplies no terminal, so no
|
||||
* session resolves to one until the spawn lifecycle binds a live terminal to the slot. The
|
||||
* stable name + profile pair is the only config-time identity; live identity is defined purely
|
||||
* by the runtime {@link dev.ltms.bridged.auth.ArchitectRegistry} binding.
|
||||
*
|
||||
* <p>Why a {@code profile} reference: an architect is meant to run a strong model, and the slot
|
||||
* records which {@code workers:} profile that is — the value the future spawn lifecycle reads.
|
||||
* It must name a configured profile, enforced by {@link #validateArchitects()} (a stale or
|
||||
* typo'd reference fails at startup rather than silently spawning the wrong backend later).
|
||||
*
|
||||
* @param terminal the architect's herdr {@code terminal_id}; the field identity is matched by,
|
||||
* via the live terminal→slot binding. Optional at config time — binding may be
|
||||
* injected live — but a slot with no binding matches nothing yet.
|
||||
* @param profile the name of the strong-model {@code workers:} profile this slot runs;
|
||||
* required and validated against {@link #workerProfiles()}
|
||||
* @param profile the name of the strong-model {@code workers:} profile this slot runs;
|
||||
* required and validated against {@link #workerProfiles()}
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Architect(String terminal, String profile) {
|
||||
public record Architect(String profile) {
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -468,30 +468,6 @@ public record BridgedConfig(
|
||||
return Collections.unmodifiableMap(byTerminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* The terminal → architect-slot-name map that {@link dev.ltms.bridged.auth.CallerResolver}
|
||||
* resolves against (CB-548), derived from the {@code architects:} registry.
|
||||
*
|
||||
* <p>Keyed by terminal because a live session is matched by its pane; the value is the
|
||||
* gateway-local slot name. Slot names are inherently unique (a map key); a duplicate terminal
|
||||
* across two slots is last-wins here (the later entry overrides), which {@code leadership} has
|
||||
* always tolerated rather than refused. This is consumed as the <em>initial</em> live binding —
|
||||
* the supplier that feeds the resolver may be swapped for a live one by the future lifecycle.
|
||||
*
|
||||
* @return an unmodifiable map, empty when no architect slot is configured
|
||||
*/
|
||||
public Map<String, String> architectTerminals() {
|
||||
Map<String, String> byTerminal = new LinkedHashMap<>();
|
||||
if (architects != null) {
|
||||
architects.forEach((name, arch) -> {
|
||||
if (arch != null && arch.terminal() != null && !arch.terminal().isBlank()) {
|
||||
byTerminal.put(arch.terminal(), name);
|
||||
}
|
||||
});
|
||||
}
|
||||
return Collections.unmodifiableMap(byTerminal);
|
||||
}
|
||||
|
||||
/**
|
||||
* API authentication (CB-501). Governs how a caller that is <em>not</em> an on-host worker
|
||||
* pane proves it is the primary.
|
||||
@@ -602,6 +578,7 @@ public record BridgedConfig(
|
||||
try {
|
||||
String yaml = Files.readString(path);
|
||||
warnUnknownTopLevelKeys(yaml, path);
|
||||
rejectDuplicateArchitectSlots(yaml);
|
||||
BridgedConfig cfg = YAML.readValue(yaml, BridgedConfig.class);
|
||||
return cfg.withDefaults();
|
||||
} catch (IOException e) {
|
||||
@@ -609,6 +586,101 @@ public record BridgedConfig(
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject an {@code architects:} registry whose slot names repeat (CB-548).
|
||||
*
|
||||
* <p>The registry is a {@code Map} keyed by slot name, so by the time it is read duplicate keys
|
||||
* have already collapsed last-wins — a duplicated slot name would silently drop one slot and the
|
||||
* daemon would never know. Jackson's YAML parser does not fail on duplicate mapping keys by
|
||||
* default, so duplicates are caught here, at parse time, before the map is built. Only the
|
||||
* <em>top-level</em> {@code architects:} block is considered, and only its direct child keys (the
|
||||
* slot names) — a nested field elsewhere, even one also named {@code architects:}, is ignored, so
|
||||
* parsing of the rest of the config is unaffected.
|
||||
*
|
||||
* @throws IllegalStateException when two {@code architects:} entries share a slot name, naming it
|
||||
*/
|
||||
static void rejectDuplicateArchitectSlots(String yaml) {
|
||||
try (JsonParser p = YAML.createParser(yaml)) {
|
||||
if (p.nextToken() != JsonToken.START_OBJECT) {
|
||||
return; // not a mapping at top level — readValue reports the malformed file
|
||||
}
|
||||
// Scan the TOP-LEVEL mapping only. Every other field's value (however deep, including
|
||||
// any nested field also literally named "architects") is consumed whole by skipValue, so
|
||||
// the loop below can only ever see the top-level field names — a nested `architects:` can
|
||||
// neither suppress the real block nor be misread as one.
|
||||
JsonToken t;
|
||||
while ((t = p.nextToken()) != null && t != JsonToken.END_OBJECT) {
|
||||
if (t == JsonToken.FIELD_NAME) {
|
||||
String name = p.getCurrentName();
|
||||
JsonToken value = p.nextToken();
|
||||
if ("architects".equals(name)) {
|
||||
if (value == JsonToken.START_OBJECT) {
|
||||
rejectDuplicateChildSlotKeys(p);
|
||||
}
|
||||
return; // the single top-level architects block is handled; nothing more to check
|
||||
}
|
||||
skipValue(p, value);
|
||||
}
|
||||
}
|
||||
} catch (IOException e) {
|
||||
// Not a duplicate-name condition — let readValue report the malformed file itself.
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Reject a duplicated <em>direct child</em> key of the (already-positioned) {@code architects:}
|
||||
* mapping — i.e. a duplicated {@code slot name}.
|
||||
*
|
||||
* <p>Each slot's value is consumed whole by {@link #skipValue}, so a duplicated field <em>inside</em>
|
||||
* a slot (e.g. two {@code profile:} keys, or a duplicate nested {@code architects:}) is never seen
|
||||
* here and cannot masquerade as a duplicated slot name.
|
||||
*
|
||||
* @throws IllegalStateException when two {@code architects:} entries share a slot name, naming it
|
||||
*/
|
||||
private static void rejectDuplicateChildSlotKeys(JsonParser p) throws IOException {
|
||||
Set<String> seen = new HashSet<>();
|
||||
JsonToken t;
|
||||
while ((t = p.nextToken()) != null && t != JsonToken.END_OBJECT) {
|
||||
if (t == JsonToken.FIELD_NAME) {
|
||||
if (!seen.add(p.getCurrentName())) {
|
||||
throw new IllegalStateException("refusing to start: duplicate architect slot name '"
|
||||
+ p.getCurrentName() + "' — slot names must be unique; a later entry would "
|
||||
+ "silently overwrite the earlier one");
|
||||
}
|
||||
skipValue(p, p.nextToken()); // the slot's entire value
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Consume the whole value that starts at {@code start}, including every nested structure, and
|
||||
* leave the parser positioned just past it. Used so depth is handled structurally rather than by
|
||||
* a heuristic — a nested field is never interpreted as a top-level {@code architects:}.
|
||||
*/
|
||||
private static void skipValue(JsonParser p, JsonToken start) throws IOException {
|
||||
switch (start) {
|
||||
case START_OBJECT: {
|
||||
JsonToken t;
|
||||
while ((t = p.nextToken()) != null && t != JsonToken.END_OBJECT) {
|
||||
if (t == JsonToken.FIELD_NAME) {
|
||||
skipValue(p, p.nextToken());
|
||||
}
|
||||
}
|
||||
return;
|
||||
}
|
||||
case START_ARRAY: {
|
||||
JsonToken t;
|
||||
while ((t = p.nextToken()) != null && t != JsonToken.END_ARRAY) {
|
||||
skipValue(p, t);
|
||||
}
|
||||
return;
|
||||
}
|
||||
default:
|
||||
// A scalar (VALUE_* / VALUE_NULL) is already fully consumed by the nextToken that
|
||||
// returned it — nothing further to skip.
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Log a WARN naming any top-level key this version does not understand (CB-530).
|
||||
*
|
||||
@@ -790,7 +862,8 @@ public record BridgedConfig(
|
||||
* spawn that quietly has no backend to use.
|
||||
*
|
||||
* <p>Slot-name uniqueness needs no check here: the registry is a {@code Map} keyed by name, so
|
||||
* duplicates are unrepresentable by construction.
|
||||
* duplicates are unrepresentable by construction once loaded — and {@link #load(Path)} already
|
||||
* rejects a duplicated slot name at parse time, before the map collapses.
|
||||
*
|
||||
* @throws IllegalStateException when any architect slot is missing or names an unknown profile,
|
||||
* naming the slot and the offending reference
|
||||
|
||||
@@ -127,21 +127,31 @@ public final class BridgeMcp {
|
||||
str(req.arguments(), "sessionId"));
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
if (caller != null) primaryRegistry.record(caller);
|
||||
// CB-548: only a PRIMARY caller may claim the legacy singleton "primary" fallback.
|
||||
// An architect delegates as its own pane but must never become the fallback that
|
||||
// no-delegation inbox nudges target as if it were the primary (the per-target
|
||||
// delegation map does not cure the singleton).
|
||||
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
|
||||
Map<String, Object> a = req.arguments();
|
||||
// CB-532: remember WHICH lead is waiting on this worker, so its reply nudge goes
|
||||
// back to that lead rather than to whichever one happened to send first.
|
||||
primaryRegistry.recordDelegation(str(a, "sessionId"), caller);
|
||||
String target = str(a, "sessionId");
|
||||
String content = str(a, "content");
|
||||
String turnId = str(a, "turnId");
|
||||
if (turnId != null && !turnId.isBlank()) {
|
||||
// Answering a worker's bridge_ask (CB-205): resolve its blocked question and
|
||||
// block for the worker's reply as it resumes the same turn.
|
||||
return answer(messages, turnId, str(a, "content"), timeoutMs(a));
|
||||
// block for the worker's reply as it resumes the same turn. This is the same
|
||||
// delegation, so ownership is left untouched (CB-548) — never re-recorded.
|
||||
return answer(messages, turnId, content, timeoutMs(a));
|
||||
}
|
||||
// CB-548: delegator ownership (which lead's reply nudge this worker routes to,
|
||||
// CB-532) is recorded only once the send is ACCEPTED — MessageService has won the
|
||||
// session lock and queued delivery — via the accepted-delivery callback, never at
|
||||
// request time. A concurrent sender that times out BUSY therefore cannot steal a
|
||||
// live turn's reply routing without ever owning the turn.
|
||||
Runnable onAccepted = () -> primaryRegistry.recordDelegation(target, caller);
|
||||
// wait defaults to true (block for the reply); wait:false is fire-and-poll.
|
||||
return Boolean.FALSE.equals(a.get("wait"))
|
||||
? sendAsync(messages, str(a, "sessionId"), str(a, "content"))
|
||||
: send(messages, str(a, "sessionId"), str(a, "content"), timeoutMs(a));
|
||||
? sendAsync(messages, target, content, onAccepted)
|
||||
: send(messages, target, content, timeoutMs(a), onAccepted);
|
||||
})
|
||||
// bridge_reply's identity is the CONNECTION, never an argument — so the authz check
|
||||
// is "is this caller a worker at all", and it can only ever reply as itself.
|
||||
@@ -182,7 +192,9 @@ public final class BridgeMcp {
|
||||
McpSchema.CallToolResult denied = deny(exchange, Authz.Action.SPAWN, null);
|
||||
if (denied != null) return denied;
|
||||
String caller = callerTerminal(exchange);
|
||||
if (caller != null) primaryRegistry.record(caller);
|
||||
// SPAWN is already auth-gated to PRIMARY (architects can never call it), but
|
||||
// enforce the same invariant here: only a PRIMARY may claim the legacy singleton.
|
||||
recordPrimarySingleton(primaryRegistry, caller, principal(exchange));
|
||||
Map<String, Object> a = req.arguments();
|
||||
// CB-112: worker inherits the primary's cwd unless the call pins one.
|
||||
// CB-301: carry the caller's identity as the session owner (null for the primary).
|
||||
@@ -301,6 +313,24 @@ public final class BridgeMcp {
|
||||
return error(reason + ": " + caller.describe() + " may not " + action);
|
||||
}
|
||||
|
||||
/**
|
||||
* Update the legacy singleton "primary" fallback used for no-delegation inbox nudges (CB-548).
|
||||
*
|
||||
* <p>Only {@link Role#PRIMARY} callers — the unnamed primary and named leads alike — may claim
|
||||
* it. An architect delegates as its own pane but must never become the fallback: the per-target
|
||||
* delegation map ({@code PrimaryRegistry#recordDelegation}) does not cure the singleton, so an
|
||||
* architect left here would draw nudges that belong to a primary. The decision uses the resolved
|
||||
* role, never name/kind sniffing. A null {@code caller} (legacy/no-auth path) records nothing.
|
||||
*
|
||||
* <p>Split out of the tool handlers so the guard is unit-testable without fabricating an SDK
|
||||
* {@code McpSyncServerExchange} (same pattern as {@link #denyFor}/{@link #principalFrom}).
|
||||
*/
|
||||
static void recordPrimarySingleton(PrimaryRegistry registry, String callerTerminal, Principal caller) {
|
||||
if (caller != null && caller.isPrimary()) {
|
||||
registry.record(callerTerminal);
|
||||
}
|
||||
}
|
||||
|
||||
/** The worker identity resolved from this call's connection, or {@code null} if the primary. */
|
||||
private static String callerTerminal(McpSyncServerExchange exchange) {
|
||||
Object v = exchange.transportContext().get(CALLER_TERMINAL);
|
||||
@@ -343,12 +373,22 @@ public final class BridgeMcp {
|
||||
|
||||
/** {@code bridge_send}: delegate {@code content} to a worker session and block for its reply. */
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content, Long timeoutMs) {
|
||||
return send(messages, sessionId, content, timeoutMs, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #send(MessageService, String, String, Long)}, wiring an accepted-delivery hook
|
||||
* (CB-548): {@code onAccepted} records delegator ownership the instant the send is accepted, so
|
||||
* a BUSY interloper never claims a turn it did not win. {@code null} disables recording.
|
||||
*/
|
||||
static McpSchema.CallToolResult send(MessageService messages, String sessionId, String content,
|
||||
Long timeoutMs, Runnable onAccepted) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
long timeout = clamp(timeoutMs == null ? DEFAULT_TIMEOUT_MS : timeoutMs);
|
||||
try {
|
||||
return formatReply(messages.send(sessionId, content, timeout), timeout);
|
||||
return formatReply(messages.send(sessionId, content, timeout, onAccepted), timeout);
|
||||
} catch (HerdrException e) {
|
||||
return error("herdr error contacting session " + sessionId + ": " + e.getMessage());
|
||||
}
|
||||
@@ -417,10 +457,19 @@ public final class BridgeMcp {
|
||||
* immediately (fire-and-poll), so a long task isn't cut off by the caller's MCP call timeout.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content) {
|
||||
return sendAsync(messages, sessionId, content, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #sendAsync(MessageService, String, String)}, wiring the accepted-delivery hook
|
||||
* (CB-548) so an async flooding send records delegator ownership exactly once it is accepted.
|
||||
*/
|
||||
static McpSchema.CallToolResult sendAsync(MessageService messages, String sessionId, String content,
|
||||
Runnable onAccepted) {
|
||||
if (isBlank(sessionId) || isBlank(content)) {
|
||||
return error("sessionId and content are required");
|
||||
}
|
||||
String ticket = messages.sendAsync(sessionId, content);
|
||||
String ticket = messages.sendAsync(sessionId, content, onAccepted);
|
||||
return text("accepted — task delegated. Poll bridge_poll with ticket=" + ticket);
|
||||
}
|
||||
|
||||
|
||||
@@ -67,12 +67,14 @@ public final class PrimaryRegistry {
|
||||
}
|
||||
|
||||
/**
|
||||
* Record that {@code leadTerminal} delegated to worker {@code target} (CB-532).
|
||||
* Record that {@code leadTerminal} owns the accepted delegation of worker {@code target} (CB-532).
|
||||
*
|
||||
* <p>Called at {@code bridge_send} time, where both halves are known: the target is the tool's
|
||||
* argument and the lead is resolved from the connection. Last writer wins — if a second lead
|
||||
* takes over a worker, replies follow the lead that most recently delegated to it, which is the
|
||||
* one waiting.
|
||||
* <p>Called from the {@code MessageService} accepted-delivery hook — only after a send has won
|
||||
* the session's send lock and queued delivery — where both halves are known (CB-548). It is
|
||||
* deliberately <em>not</em> called at {@code bridge_send} request time: a concurrent sender that
|
||||
* times out {@code BUSY} must not steal a live delegation's reply routing without ever owning
|
||||
* the turn. Last writer wins — if a second lead's later send is accepted, replies follow the
|
||||
* lead that most recently delegated to it, which is the one waiting.
|
||||
*/
|
||||
public void recordDelegation(String target, String leadTerminal) {
|
||||
if (target == null || target.isBlank() || leadTerminal == null || leadTerminal.isBlank()) {
|
||||
|
||||
@@ -310,6 +310,22 @@ public final class MessageService {
|
||||
* worker replies via {@link Rendezvous} or {@code timeoutMillis} elapses.
|
||||
*/
|
||||
public Reply send(String target, String content, long timeoutMillis) {
|
||||
return send(target, content, timeoutMillis, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #send(String, String, long)}, but with an accepted-delivery hook.
|
||||
*
|
||||
* <p>{@code onAccepted} is invoked exactly once, once this send has won {@code target}'s send
|
||||
* lock and so become the <em>accepted target turn</em> — it runs <em>before</em> delivery is
|
||||
* queued, so a throwing hook fails the send cleanly (the waiter it already opened is closed and
|
||||
* nothing is left queued). It is <em>not</em> invoked when the send is {@link Outcome#BUSY}
|
||||
* (lock never taken). A caller uses this to record that <em>it</em> now owns the delegation's
|
||||
* reply routing (CB-548: {@code PrimaryRegistry} delegator ownership) — recording only on
|
||||
* acceptance means a concurrent sender that times out {@code BUSY} can never steal ownership it
|
||||
* never earned. {@code null} disables the hook.
|
||||
*/
|
||||
public Reply send(String target, String content, long timeoutMillis, Runnable onAccepted) {
|
||||
long deadlineNanos = System.nanoTime() + timeoutMillis * 1_000_000L;
|
||||
ReentrantLock lock = sessionLocks.computeIfAbsent(target, _ -> new ReentrantLock());
|
||||
|
||||
@@ -317,22 +333,36 @@ public final class MessageService {
|
||||
return new Reply(Outcome.BUSY, null); // another send held the session the whole window
|
||||
}
|
||||
try {
|
||||
CompletableFuture<Void> delivered = injector.enqueue(target, content);
|
||||
// Open the waiter BEFORE queueing delivery (CB-548). A fast reply — the worker already
|
||||
// injectable the instant we enqueue — otherwise arrives before the waiter is registered
|
||||
// and orphans into the inbox while this send blocks to the timeout (the enqueue-before-
|
||||
// open race). Opening first also means a throwing onAccepted (fired before enqueue) or an
|
||||
// enqueue failure is safely closed by the finally below: nothing is left queued, and the
|
||||
// failed send leaves no stale waiter behind.
|
||||
CompletableFuture<Rendezvous.Resolution> reply = rendezvous.open(target);
|
||||
try {
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
return recorded(new Reply(
|
||||
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("interrupted awaiting reply from " + target, e);
|
||||
// The send has won the lock; the accepted-delivery hook records delegator ownership
|
||||
// here (CB-548). It runs BEFORE enqueue so a throwing hook — onAccepted is now a
|
||||
// public callback — fails the send without queuing a message that would orphan.
|
||||
if (onAccepted != null) {
|
||||
onAccepted.run();
|
||||
}
|
||||
CompletableFuture<Void> delivered = injector.enqueue(target, content);
|
||||
try {
|
||||
Rendezvous.Resolution r = reply.get(remainingMillis(deadlineNanos), TimeUnit.MILLISECONDS);
|
||||
return recorded(new Reply(outcomeOf(r.kind()), r.text(), r.turnId()));
|
||||
} catch (TimeoutException e) {
|
||||
boolean wasDelivered = delivered.isDone() && !delivered.isCompletedExceptionally();
|
||||
log.debug("send to {} timed out (delivered={})", target, wasDelivered);
|
||||
return recorded(new Reply(
|
||||
wasDelivered ? Outcome.TIMED_OUT_WORKING : Outcome.TIMED_OUT_QUEUED, null));
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("interrupted awaiting reply from " + target, e);
|
||||
}
|
||||
} finally {
|
||||
rendezvous.close(target, reply);
|
||||
}
|
||||
@@ -440,9 +470,21 @@ public final class MessageService {
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content) {
|
||||
return sendAsync(target, content, null);
|
||||
}
|
||||
|
||||
/**
|
||||
* As {@link #sendAsync(String, String)}, with the accepted-delivery hook of
|
||||
* {@link #send(String, String, long, Runnable)} — the running {@code send} invokes {@code onAccepted}
|
||||
* the moment it becomes the accepted target turn, so async flooding records delegator ownership
|
||||
* exactly as the blocking path does (CB-548).
|
||||
*
|
||||
* @return the ticket to poll for the eventual result
|
||||
*/
|
||||
public String sendAsync(String target, String content, Runnable onAccepted) {
|
||||
String ticket = "task-" + ticketSeq.incrementAndGet();
|
||||
CompletableFuture<Reply> future =
|
||||
CompletableFuture.supplyAsync(() -> send(target, content, ASYNC_TIMEOUT_MS), asyncExecutor);
|
||||
CompletableFuture<Reply> future = CompletableFuture.supplyAsync(
|
||||
() -> send(target, content, ASYNC_TIMEOUT_MS, onAccepted), asyncExecutor);
|
||||
tasks.put(ticket, new Task(target, future, System.nanoTime()));
|
||||
pruneTerminalTickets();
|
||||
log.debug("async send {} -> {}", ticket, target);
|
||||
|
||||
@@ -74,15 +74,30 @@ public final class Rendezvous {
|
||||
/**
|
||||
* Register a waiter for {@code session} — the await side of the public {@code resolve*} methods.
|
||||
* The caller must hold that session's send lock.
|
||||
*
|
||||
* <p>Atomic fail-if-present (CB-548): if a waiter is already registered for {@code session}, an
|
||||
* {@link IllegalStateException} is thrown rather than replacing the first — so any future
|
||||
* invariant violation fails loudly instead of silently swapping the waiter another send is
|
||||
* blocked on. {@code MessageService} serializes sends per session (the send lock), so in correct
|
||||
* code a double open is impossible; this is a tripwire for the day that no longer holds.
|
||||
*/
|
||||
public CompletableFuture<Resolution> open(String session) {
|
||||
CompletableFuture<Resolution> waiter = new CompletableFuture<>();
|
||||
waiters.put(session, waiter);
|
||||
CompletableFuture<Resolution> existing = waiters.putIfAbsent(session, waiter);
|
||||
if (existing != null) {
|
||||
throw new IllegalStateException(
|
||||
"rendezvous double-open for session " + session + " — a waiter is already registered");
|
||||
}
|
||||
return waiter;
|
||||
}
|
||||
|
||||
/** Remove {@code waiter} for {@code session} (only if it is still the registered one). */
|
||||
void close(String session, CompletableFuture<Resolution> waiter) {
|
||||
/**
|
||||
* Remove {@code waiter} for {@code session}, only if it is still the registered one. The
|
||||
* symmetric complement of {@link #open}: a terminal send deregisters its waiter so the next
|
||||
* send on the session may {@link #open} a fresh one (CB-548 makes double-open an error, so a
|
||||
* successful {@code open} after a finished turn requires this close to have happened first).
|
||||
*/
|
||||
public void close(String session, CompletableFuture<Resolution> waiter) {
|
||||
waiters.remove(session, waiter);
|
||||
}
|
||||
|
||||
|
||||
@@ -27,12 +27,47 @@ public final class GitWorktrees implements Worktrees {
|
||||
|
||||
private static final Logger log = LoggerFactory.getLogger(GitWorktrees.class);
|
||||
|
||||
/** Project-level MCP config. Present in the repo, so every worktree checks the primary's out. */
|
||||
/** Project-level MCP config. Present in the repo, so every worktree would otherwise inherit the
|
||||
* primary's IDE server mounts (a CB-523 worker edited the primary checkout; see the isolation
|
||||
* javadoc). Neutralized unconditionally. */
|
||||
private static final String MCP_CONFIG = ".mcp.json";
|
||||
|
||||
/** What {@link #isolateToolSurface} writes: a valid, explicitly empty server map. */
|
||||
/** What {@link #isolateToolSurface} writes for {@code .mcp.json}: a valid, explicitly empty server map. */
|
||||
private static final String NEUTRAL_MCP_CONFIG = "{\n \"mcpServers\": {}\n}\n";
|
||||
|
||||
/** OpenCode's repo-level config. Tracked here, so it lands in every worktree; it carries
|
||||
* {@code {file:.secrets/...}} references to gitignored secrets that never reach a worktree, and
|
||||
* opencode refuses to start on a dangling reference — so it is neutralized and the worker gets
|
||||
* only the config its launcher writes via {@code OPENCODE_CONFIG}. */
|
||||
private static final String OPENCODE_CONFIG = "opencode.json";
|
||||
|
||||
/** What {@link #isolateToolSurface} writes for {@code opencode.json}: a valid, empty JSON object. */
|
||||
private static final String NEUTRAL_OPENCODE_CONFIG = "{}\n";
|
||||
|
||||
/** Autoenv's repo-level config. Not tracked today, but re-landing it must stay safe: autoenv
|
||||
* authorizes by path, so a fresh worktree path is always unauthorized and its interactive prompt
|
||||
* would block every spawn — neutralize it so it can never be committed. */
|
||||
private static final String AUTOENV_CONFIG = ".autoenv";
|
||||
|
||||
/** What {@link #isolateToolSurface} writes for {@code .autoenv}: a valid, empty env file. */
|
||||
private static final String NEUTRAL_AUTOENV_CONFIG = "";
|
||||
|
||||
/**
|
||||
* A tracked project config that is hostile in a provisioned worktree, and what to replace it
|
||||
* with. {@link #file} is the repo-relative path; {@link #stub} is a neutral but VALID payload for
|
||||
* that file's format — a malformed stub would only trade one crash for another;
|
||||
* {@link #createIfAbsent} keeps {@code .mcp.json}'s long-standing behaviour of writing its stub
|
||||
* even when the repo carries no such file, whereas the others are only touched when present.
|
||||
*/
|
||||
private record WorktreeHostileConfig(String file, String stub, boolean createIfAbsent) {}
|
||||
|
||||
/** The worktree-hostile configs neutralized in every provisioned worktree, in order. */
|
||||
private static final List<WorktreeHostileConfig> WORKTREE_HOSTILE_CONFIGS = List.of(
|
||||
new WorktreeHostileConfig(MCP_CONFIG, NEUTRAL_MCP_CONFIG, true),
|
||||
new WorktreeHostileConfig(OPENCODE_CONFIG, NEUTRAL_OPENCODE_CONFIG, false),
|
||||
new WorktreeHostileConfig(AUTOENV_CONFIG, NEUTRAL_AUTOENV_CONFIG, false)
|
||||
);
|
||||
|
||||
private final String configuredRoot;
|
||||
private final SecureRandom random = new SecureRandom();
|
||||
private final AtomicLong seq = new AtomicLong();
|
||||
@@ -66,8 +101,9 @@ public final class GitWorktrees implements Worktrees {
|
||||
}
|
||||
|
||||
/**
|
||||
* Neutralize the worktree's project MCP config so a worker inherits only the tools its launcher
|
||||
* mounts (the bridge, via {@code --mcp-config}) — never the primary's.
|
||||
* Neutralize the worktree's worktree-hostile project configs so a worker inherits only the tools
|
||||
* and environment its launcher mounts (the bridge via {@code --mcp-config}, the opencode config
|
||||
* via {@code OPENCODE_CONFIG}) — never the primary's.
|
||||
*
|
||||
* <p>This is unconditional, and it is not the same job as the parity overlay. The repo's own
|
||||
* committed {@code .mcp.json} declares the primary's IDE servers, so a fresh checkout mounts them
|
||||
@@ -75,25 +111,42 @@ public final class GitWorktrees implements Worktrees {
|
||||
* through tools bound to the <em>primary's</em> IntelliJ project, which silently hands it absolute
|
||||
* paths outside its own worktree. That is not hypothetical: a CB-523 worker made all 59 of its
|
||||
* edits in the primary checkout while compiling its worktree, so every build it ran was of code
|
||||
* that did not contain its changes.
|
||||
* that did not contain its changes. {@code opencode.json} is the same trap one tool over — tracked,
|
||||
* so it lands in every worktree, referencing gitignored {@code .secrets/} files that never do, and
|
||||
* opencode refuses to start on the dangling reference. {@code .autoenv} extends the principle to a
|
||||
* config that is not tracked today: autoenv authorizes by path, so a fresh worktree path is always
|
||||
* unauthorized and its interactive prompt would block every spawn, so re-landing one must be safe.
|
||||
*
|
||||
* <p>Writing an empty server map (rather than deleting the file) keeps a project-level
|
||||
* {@code .mcp.json} present and explicit, and the {@code --skip-worktree} bit keeps the
|
||||
* neutralized copy from ever showing up as a local modification the worker might commit.
|
||||
* <p>Where a config exists it is replaced by a valid neutral stub (an explicitly empty
|
||||
* map/object, or an empty env file — never a deletion, which would still let a later
|
||||
* {@code git checkout} restore the hostile copy). The {@code --skip-worktree} bit keeps the
|
||||
* neutralized copy from ever showing up as a local modification the worker might commit. A config
|
||||
* the repo does not carry is skipped silently — no stub is invented for a file the repo does not
|
||||
* have, and one missing file must never fail provisioning.
|
||||
*/
|
||||
private void isolateToolSurface(String worktreePath) {
|
||||
Path root = Path.of(worktreePath).toAbsolutePath().normalize();
|
||||
Path mcp = root.resolve(MCP_CONFIG);
|
||||
for (WorktreeHostileConfig cfg : WORKTREE_HOSTILE_CONFIGS) {
|
||||
neutralize(root, worktreePath, cfg);
|
||||
}
|
||||
}
|
||||
|
||||
private void neutralize(Path root, String worktreePath, WorktreeHostileConfig cfg) {
|
||||
Path target = root.resolve(cfg.file());
|
||||
if (!Files.exists(target) && !cfg.createIfAbsent()) {
|
||||
log.debug("{} absent in the worktree — skipping (repo does not carry it)", cfg.file());
|
||||
return;
|
||||
}
|
||||
try {
|
||||
Files.writeString(mcp, NEUTRAL_MCP_CONFIG);
|
||||
Files.writeString(target, cfg.stub());
|
||||
} catch (IOException e) {
|
||||
throw new WorktreeException("cannot neutralize " + MCP_CONFIG + " in the worktree: "
|
||||
throw new WorktreeException("cannot neutralize " + cfg.file() + " in the worktree: "
|
||||
+ e.getMessage(), e);
|
||||
}
|
||||
if (isTracked(root, MCP_CONFIG)) {
|
||||
exec("git", "-C", worktreePath, "update-index", "--skip-worktree", MCP_CONFIG);
|
||||
if (isTracked(root, cfg.file())) {
|
||||
exec("git", "-C", worktreePath, "update-index", "--skip-worktree", cfg.file());
|
||||
}
|
||||
log.debug("neutralized {} — worker tool surface is launcher-mounted only", MCP_CONFIG);
|
||||
log.debug("neutralized {} — worker tool surface is launcher-mounted only", cfg.file());
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -7,6 +7,8 @@ import dev.ltms.bridged.herdr.Agent;
|
||||
import dev.ltms.bridged.herdr.AgentControl;
|
||||
import dev.ltms.bridged.herdr.WorkspaceControl;
|
||||
import dev.ltms.bridged.peer.Capability;
|
||||
import dev.ltms.bridged.peer.PeerHandle;
|
||||
import dev.ltms.bridged.peer.SpawnRequest;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.io.UncheckedIOException;
|
||||
@@ -72,6 +74,24 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
/** Root under which per-spawn opencode config dirs are created (injectable for tests). */
|
||||
private final Path configRoot;
|
||||
|
||||
/**
|
||||
* The current spawn's resume-target session id, threaded from {@link #spawn(SpawnRequest)} to
|
||||
* {@link #buildLaunch} across the base's {@code spawn -> spawnInternal -> buildLaunch} chain,
|
||||
* which carries no request. A plain field would race under concurrent spawns (the base supports
|
||||
* them), so it is thread-local: each spawn captures its own request's id on its own thread, and
|
||||
* {@code buildLaunch}, synchronous and same-thread, reads exactly that one. Set only around the
|
||||
* {@code super.spawn} call and cleared in {@code finally}, so a paused/leftover value can never
|
||||
* bleed into the next spawn.
|
||||
*/
|
||||
private final ThreadLocal<String> resumeSessionId = new ThreadLocal<>();
|
||||
|
||||
/**
|
||||
* Session discovery against opencode's on-disk storage ({@link OpenCodeSessionDiscovery}) —
|
||||
* the one seam that knows opencode's private session-file layout. Its root is injectable for
|
||||
* tests so they never touch the operator's real {@code ~/.local/share/opencode}.
|
||||
*/
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
|
||||
/**
|
||||
* Production constructor — disables the spawn-ready gate ({@code spawnReadyTimeoutMs == 0}) so it
|
||||
* matches the legacy non-blocking spawn semantics. Config dirs are created under the JVM temp dir.
|
||||
@@ -81,7 +101,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
Function<String, String> env) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, 0,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(300),
|
||||
defaultConfigRoot());
|
||||
defaultConfigRoot(), defaultDiscoveryRoot());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -94,7 +114,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
long spawnReadyTimeoutMs, long spawnReadyPollMs) {
|
||||
this(agents, spaces, profiles, defaultProfile, env, spawnReadyTimeoutMs,
|
||||
System::currentTimeMillis, () -> sleepUninterruptibly(spawnReadyPollMs),
|
||||
defaultConfigRoot());
|
||||
defaultConfigRoot(), defaultDiscoveryRoot());
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -112,21 +132,31 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
* @param sleeper sleep/wait hook (encodes the poll interval; never called when the
|
||||
* gate is disabled)
|
||||
* @param configRoot existing directory under which per-spawn config dirs are created
|
||||
* @param discoveryRoot opencode's on-disk storage root to scan for session records
|
||||
* (injectable for tests; opencode's layout is matched at
|
||||
* {@link OpenCodeSessionDiscovery})
|
||||
*/
|
||||
public OpenCodeLauncher(AgentControl agents, WorkspaceControl spaces,
|
||||
Map<String, BridgedConfig.Worker> profiles, String defaultProfile,
|
||||
Function<String, String> env,
|
||||
long spawnReadyTimeoutMs,
|
||||
LongSupplier nowMillis, Runnable sleeper, Path configRoot) {
|
||||
LongSupplier nowMillis, Runnable sleeper,
|
||||
Path configRoot, Path discoveryRoot) {
|
||||
super(NAME_PREFIX, agents, spaces, profiles, defaultProfile, env,
|
||||
spawnReadyTimeoutMs, nowMillis, sleeper);
|
||||
this.configRoot = configRoot;
|
||||
this.discovery = new OpenCodeSessionDiscovery(discoveryRoot);
|
||||
}
|
||||
|
||||
private static Path defaultConfigRoot() {
|
||||
return Path.of(System.getProperty("java.io.tmpdir"));
|
||||
}
|
||||
|
||||
/** The default opencode storage root: {@code ~/.local/share/opencode} (the XDG data dir). */
|
||||
private static Path defaultDiscoveryRoot() {
|
||||
return Path.of(System.getProperty("user.home"), ".local", "share", "opencode");
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
*
|
||||
@@ -144,7 +174,7 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
workerEnv.put("OPENCODE_CONFIG", writeConfig(cfg).toString());
|
||||
}
|
||||
applyGitToken(workerEnv, cfg);
|
||||
return new Launch(workerEnv, argvWithModel(argvWithAuto(cfg), cfg));
|
||||
return new Launch(workerEnv, argvWithResume(argvWithModel(argvWithAuto(cfg), cfg)));
|
||||
}
|
||||
|
||||
/**
|
||||
@@ -177,6 +207,24 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
return argv;
|
||||
}
|
||||
|
||||
/**
|
||||
* The launch argv plus, on a resumed spawn, opencode's {@code -s <id>} flag to continue a prior
|
||||
* conversation by its session id. {@code -s, --session <id>} resumes an existing session; on a
|
||||
* fresh spawn (no resume target) no flag is added, letting opencode start a brand-new session.
|
||||
* The id comes from the current spawn request's {@code resumeSessionId}, threaded per-thread by
|
||||
* {@link #spawn(SpawnRequest)}.
|
||||
*/
|
||||
private List<String> argvWithResume(List<String> argv) {
|
||||
String id = resumeSessionId.get();
|
||||
if (id == null || id.isBlank()) {
|
||||
return argv;
|
||||
}
|
||||
List<String> withResume = mutableArgv(argv);
|
||||
withResume.add("-s");
|
||||
withResume.add(id);
|
||||
return withResume;
|
||||
}
|
||||
|
||||
/** The launch argv plus, when a model is configured, the opencode {@code -m provider/model} flag. */
|
||||
private List<String> argvWithModel(List<String> argv, BridgedConfig.Worker cfg) {
|
||||
if (cfg.model() != null && !cfg.model().isBlank()) {
|
||||
@@ -285,6 +333,82 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
return afterScheme.contains("/") ? trimmed : trimmed + "/v1";
|
||||
}
|
||||
|
||||
/**
|
||||
* {@inheritDoc}
|
||||
*
|
||||
* <p>adds this adapter's session-identity work around the base's spawn — as opencode cannot be
|
||||
* told its session id at spawn (see {@link Capability#SESSION_RESUME} vs
|
||||
* {@link Capability#SESSION_NAME}), identity is only ever adopted after the fact:
|
||||
* <ul>
|
||||
* <li>the request's {@code resumeSessionId} is remembered for {@link #buildLaunch} to turn
|
||||
* into {@code -s <id>}; and</li>
|
||||
* <li>the returned handle is wrapped so its
|
||||
* {@link dev.ltms.bridged.peer.PeerHandle#agentSessionId()} performs lazy session
|
||||
* discovery against opencode's storage (see {@link OpenCodeSessionDiscovery}) — always
|
||||
* non-blocking, {@code null} until opencode has persisted the session record.</li>
|
||||
* </ul>
|
||||
*/
|
||||
@Override
|
||||
public PeerHandle spawn(SpawnRequest req) {
|
||||
resumeSessionId.set(req.resumeSessionId());
|
||||
try {
|
||||
PeerHandle inner = super.spawn(req);
|
||||
return new SessionAwareHandle(inner, discovery, effectiveCwd(req));
|
||||
} finally {
|
||||
// Never let a paused/leftover resume id bleed into the next spawn on this thread.
|
||||
resumeSessionId.remove();
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* A {@link PeerHandle} that delegates everything to the base's worker handle but resolves
|
||||
* {@link #agentSessionId()} lazily through opencode session discovery. Delegate-only, so the
|
||||
* base's id/terminalId/profile semantics (CB-519's host-unique routing key, herdr coordinates)
|
||||
* are untouched — only the opencode-specific identity answer is added. {@code sessionName()}
|
||||
* stays null: opencode has no display-name seam, so the logical name lives only in the bridge's
|
||||
* roster (see the SESSION_NAME capability).
|
||||
*/
|
||||
private static final class SessionAwareHandle implements PeerHandle {
|
||||
private final PeerHandle delegate;
|
||||
private final OpenCodeSessionDiscovery discovery;
|
||||
private final String cwd;
|
||||
|
||||
SessionAwareHandle(PeerHandle delegate, OpenCodeSessionDiscovery discovery, String cwd) {
|
||||
this.delegate = delegate;
|
||||
this.discovery = discovery;
|
||||
this.cwd = cwd;
|
||||
}
|
||||
|
||||
@Override
|
||||
public String id() {
|
||||
return delegate.id();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String terminalId() {
|
||||
return delegate.terminalId();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String profile() {
|
||||
return delegate.profile();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String sessionName() {
|
||||
return delegate.sessionName();
|
||||
}
|
||||
|
||||
@Override
|
||||
public String agentSessionId() {
|
||||
// Lazy + retried, never a spawn-time blocker: opencode writes the session record only
|
||||
// when the session is first persisted, so null here is the correct interim answer and
|
||||
// the caller re-calls later (each call re-scans, picking up a record that has since
|
||||
// appeared).
|
||||
return discovery.sessionIdForDirectory(cwd);
|
||||
}
|
||||
}
|
||||
|
||||
// --- Agent-returning convenience spawns (used by callers/tests that want the herdr Agent) ---
|
||||
|
||||
/** Spawn a worker for the default profile in the resolved default cwd. */
|
||||
@@ -306,10 +430,15 @@ public final class OpenCodeLauncher extends HerdrPeerLauncher {
|
||||
|
||||
@Override
|
||||
public Set<Capability> capabilities() {
|
||||
Set<Capability> caps = EnumSet.of(Capability.MID_TURN_ASK, Capability.WORKTREE, Capability.ORPHAN_REAP);
|
||||
Set<Capability> caps = EnumSet.of(Capability.MID_TURN_ASK, Capability.WORKTREE,
|
||||
Capability.ORPHAN_REAP, Capability.SESSION_RESUME);
|
||||
if (hasGitTokenProfile()) {
|
||||
caps.add(Capability.SELF_PR);
|
||||
}
|
||||
// Deliberately NOT SESSION_NAME: opencode has no display-name flag, so the bridge's logical
|
||||
// name can't surface in the peer's own UI — declaring the capability would hide that
|
||||
// asymmetry rather than make it honest. For opencode the name lives only in the bridge's
|
||||
// roster (see PeerHandle.sessionName() returning null).
|
||||
return Set.copyOf(caps);
|
||||
}
|
||||
|
||||
|
||||
@@ -0,0 +1,122 @@
|
||||
package dev.ltms.bridged.worker;
|
||||
|
||||
import com.fasterxml.jackson.databind.JsonNode;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
|
||||
import java.io.IOException;
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.util.stream.Stream;
|
||||
|
||||
/**
|
||||
* Resolves the opencode session id for a bridged worker from opencode's on-disk storage — the
|
||||
* only place this adapter touches opencode's private layout, and deliberately the <em>only</em>
|
||||
* class that does.
|
||||
*
|
||||
* <p><strong>Why this is isolated behind one seam.</strong> The layout is version-coupled and not a
|
||||
* stable contract: opencode writes one JSON file per session under
|
||||
* {@code <storageRoot>/session/<projectID>/<ses_*.json>}, and each record carries a
|
||||
* {@code "version"} field (e.g. {@code "1.1.31"}), so the exact directory shape, file naming, and
|
||||
* field names can move between opencode releases. opencode also ships a headless HTTP server that
|
||||
* may supersede file scanning entirely. Everything this adapter knows about that private storage —
|
||||
* its shape, naming, and field names — lives here, so a layout change, or a switch to the HTTP
|
||||
* server, changes exactly one class and nothing in {@link OpenCodeLauncher}.
|
||||
*
|
||||
* <p>The determinism that makes this useful is structural, not a guess: every bridged worker runs
|
||||
* in its own unique git worktree, so the record's {@code directory} (its project root) equals the
|
||||
* worker's cwd identifies <em>its</em> session unambiguously. We match on {@code directory} rather
|
||||
* than diffing {@code opencode session list} before/after — that races under concurrent spawns, and
|
||||
* the CLI listing does not even show the directory.
|
||||
*
|
||||
* <p>All reads are best-effort and never throw: a missing or unreadable storage root, a record that
|
||||
* fails to parse, or a directory with no record yet all yield {@code null}, and the caller (the
|
||||
* session handle) treats that as "identity not resolved yet" and retries later.
|
||||
*/
|
||||
final class OpenCodeSessionDiscovery {
|
||||
|
||||
private final Path storageRoot; // e.g. ~/.local/share/opencode (injectable for tests)
|
||||
private final ObjectMapper json;
|
||||
|
||||
OpenCodeSessionDiscovery(Path storageRoot) {
|
||||
this.storageRoot = storageRoot;
|
||||
this.json = new ObjectMapper();
|
||||
}
|
||||
|
||||
/**
|
||||
* The opencode session id whose record references {@code directory} (the worker's cwd), or
|
||||
* {@code null} when no record matches yet. When several records share the directory — e.g.
|
||||
* repeated spawns into the same worktree — the <em>most recently modified</em> one wins: it is
|
||||
* the session the pane most likely corresponds to.
|
||||
*
|
||||
* <p>Never throws: a missing {@code storageRoot}, an unreadable/malformed record, or a
|
||||
* directory that has not been persisted yet all resolve to {@code null} rather than failing a
|
||||
* spawn. A bridged worker's session record is written lazily (when the session is first
|
||||
* persisted), so {@code null} here is the normal answer right after the pane is ready, and the
|
||||
* caller retries later.
|
||||
*
|
||||
* @param directory the worker's cwd, as resolved for this spawn
|
||||
* @return the matching session id, or {@code null} if none is known yet
|
||||
*/
|
||||
String sessionIdForDirectory(String directory) {
|
||||
if (directory == null || directory.isBlank()) {
|
||||
return null;
|
||||
}
|
||||
Path sessionRoot = storageRoot.resolve("session");
|
||||
if (!Files.isDirectory(sessionRoot)) {
|
||||
return null;
|
||||
}
|
||||
String best = null;
|
||||
long bestMtime = Long.MIN_VALUE;
|
||||
try (Stream<Path> projectDirs = Files.list(sessionRoot)) {
|
||||
for (Path projectDir : projectDirs.filter(Files::isDirectory).toList()) {
|
||||
try (Stream<Path> records = Files.list(projectDir)) {
|
||||
for (Path record : records.toList()) {
|
||||
String id = matchId(record, directory);
|
||||
if (id == null) {
|
||||
continue;
|
||||
}
|
||||
long mtime = lastModifiedEpochMillis(record);
|
||||
if (mtime > bestMtime) {
|
||||
bestMtime = mtime;
|
||||
best = id;
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
// one project dir unreadable — skip it; another may still match
|
||||
}
|
||||
}
|
||||
} catch (IOException ignored) {
|
||||
// storage root vanished or became unreadable — "no session known yet"
|
||||
return null;
|
||||
}
|
||||
return best;
|
||||
}
|
||||
|
||||
/**
|
||||
* The record's session id when it references {@code directory}, else {@code null}. A record
|
||||
* that is not JSON, lacks {@code id}/{@code directory}, or points at a different directory is
|
||||
* simply not our session; a malformed one is skipped, never fatal.
|
||||
*/
|
||||
private String matchId(Path record, String directory) {
|
||||
try {
|
||||
JsonNode node = json.readTree(record.toFile());
|
||||
JsonNode id = node == null ? null : node.get("id");
|
||||
JsonNode dir = node == null ? null : node.get("directory");
|
||||
if (id == null || dir == null || !directory.equals(dir.asText())) {
|
||||
return null;
|
||||
}
|
||||
return id.asText();
|
||||
} catch (IOException e) {
|
||||
return null;
|
||||
}
|
||||
}
|
||||
|
||||
/** The record's last-modified epoch ms, or {@code Long.MIN_VALUE} if unreadable (never wins). */
|
||||
private static long lastModifiedEpochMillis(Path record) {
|
||||
try {
|
||||
return Files.getLastModifiedTime(record).toMillis();
|
||||
} catch (IOException e) {
|
||||
return Long.MIN_VALUE;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -3,24 +3,30 @@ package dev.ltms.bridged.auth;
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.ArrayList;
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.concurrent.Future;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* CB-548 — the architect-slot registry: the config snapshot of slot → profile, and the live
|
||||
* terminal → slot binding the resolver reads. The role the binding produces is asserted in
|
||||
* {@link CallerResolverTest}; this pins the registry object itself.
|
||||
* terminal → slot bindings it owns. The role a binding produces is asserted in
|
||||
* {@link CallerResolverTest}; this pins the registry object itself — its invariants and their
|
||||
* thread-safety.
|
||||
*/
|
||||
class ArchitectRegistryTest {
|
||||
|
||||
private static final Map<String, BridgedConfig.Architect> SLOTS = Map.of(
|
||||
"lead-designer", new BridgedConfig.Architect("term_design", "sonnet"),
|
||||
"reviewer", new BridgedConfig.Architect(null, "gx10"));
|
||||
"lead-designer", new BridgedConfig.Architect("sonnet"),
|
||||
"reviewer", new BridgedConfig.Architect("gx10"));
|
||||
|
||||
private final ArchitectRegistry registry =
|
||||
new ArchitectRegistry(SLOTS, () -> Map.of("term_design", "lead-designer"));
|
||||
private final ArchitectRegistry registry = new ArchitectRegistry(SLOTS);
|
||||
|
||||
@Test
|
||||
void exposesTheConfiguredSlots() {
|
||||
@@ -37,30 +43,195 @@ class ArchitectRegistryTest {
|
||||
}
|
||||
|
||||
@Test
|
||||
void resolvesTheSlotOfALiveTerminal() {
|
||||
assertEquals("lead-designer", registry.slotForTerminal("term_design"));
|
||||
assertNull(registry.slotForTerminal("term_unbound"));
|
||||
void startsEmptySoNoTerminalResolvesToAnArchitect() {
|
||||
assertTrue(registry.snapshot().isEmpty());
|
||||
assertNull(registry.slotForTerminal("term_design"),
|
||||
"config declares no architect terminal — nothing is recognised until a bind");
|
||||
assertNull(registry.slotForTerminal(null), "no terminal ⇒ no slot");
|
||||
}
|
||||
|
||||
// ── bind ──────────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void theBindingIsLiveReReadPerCall() {
|
||||
Map<String, String> live = new HashMap<>();
|
||||
ArchitectRegistry r = new ArchitectRegistry(SLOTS, () -> live);
|
||||
|
||||
assertNull(r.slotForTerminal("term_design"));
|
||||
|
||||
live.put("term_design", "lead-designer"); // injected after construction
|
||||
|
||||
assertEquals("lead-designer", r.slotForTerminal("term_design"));
|
||||
void bindResolvesTheTerminalToTheSlot() {
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
assertEquals("lead-designer", registry.slotForTerminal("term_design"));
|
||||
assertEquals(Map.of("term_design", "lead-designer"), registry.snapshot());
|
||||
}
|
||||
|
||||
@Test
|
||||
void bindRefusesAnUnknownSlot() {
|
||||
assertFalse(registry.bind("nope", "term_x"),
|
||||
"a slot that is not configured must be refused — bind is not a way to invent one");
|
||||
assertNull(registry.slotForTerminal("term_x"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void bindRefusesATerminalInTwoSlots() {
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
assertFalse(registry.bind("reviewer", "term_design"),
|
||||
"a terminal may occupy at most one slot");
|
||||
assertEquals("lead-designer", registry.slotForTerminal("term_design"),
|
||||
"the first binding survives the refused second");
|
||||
}
|
||||
|
||||
@Test
|
||||
void bindRefusesASlotWithTwoTerminals() {
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
assertFalse(registry.bind("lead-designer", "term_other"),
|
||||
"a slot may host at most one terminal");
|
||||
assertEquals("lead-designer", registry.slotForTerminal("term_design"),
|
||||
"the first binding survives the refused second");
|
||||
assertNull(registry.slotForTerminal("term_other"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void rebindingTheSamePairIsAnIdempotentNoOp() {
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
assertTrue(registry.bind("lead-designer", "term_design"),
|
||||
"the same terminal → slot is harmless to repeat");
|
||||
assertEquals(1, registry.snapshot().size());
|
||||
}
|
||||
|
||||
// ── unbind ────────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void unbindRemovesTheExactBinding() {
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
assertTrue(registry.unbind("lead-designer", "term_design"));
|
||||
assertNull(registry.slotForTerminal("term_design"));
|
||||
assertTrue(registry.snapshot().isEmpty());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aStaleUnbindDoesNotRemoveAReplacement() {
|
||||
// Bind, tear down, and stand the slot back up with a NEW terminal.
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
registry.unbind("lead-designer", "term_design");
|
||||
assertTrue(registry.bind("lead-designer", "term_new"));
|
||||
|
||||
// A late unbind naming the OLD terminal must not remove the replacement binding.
|
||||
assertFalse(registry.unbind("lead-designer", "term_design"));
|
||||
assertEquals("lead-designer", registry.slotForTerminal("term_new"),
|
||||
"the replacement terminal stays bound");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aStaleUnbindForATerminalThatMovedSlotsDoesNothing() {
|
||||
// term_design starts in lead-designer, is torn down, and stands back up in a FREE slot.
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
registry.unbind("lead-designer", "term_design");
|
||||
assertTrue(registry.bind("reviewer", "term_design"));
|
||||
|
||||
// Unbinding against the slot it no longer occupies is refused; the new binding is intact.
|
||||
assertFalse(registry.unbind("lead-designer", "term_design"),
|
||||
"the old slot must not unbind a terminal that moved elsewhere");
|
||||
assertEquals("reviewer", registry.slotForTerminal("term_design"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void unbindOfNothingIsAFalseNoOp() {
|
||||
assertFalse(registry.unbind("lead-designer", "term_design"),
|
||||
"nothing was bound, so nothing is removed");
|
||||
}
|
||||
|
||||
// ── snapshot ─────────────────────────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void theSnapshotIsAnImmutableCopyNotAliveState() {
|
||||
assertTrue(registry.bind("lead-designer", "term_design"));
|
||||
Map<String, String> snap = registry.snapshot();
|
||||
|
||||
assertThrows(UnsupportedOperationException.class, () -> snap.put("x", "y"),
|
||||
"a handed-out snapshot cannot be mutated in place");
|
||||
|
||||
// Later binds must not leak into an earlier snapshot.
|
||||
assertTrue(registry.bind("reviewer", "term_review"));
|
||||
assertFalse(snap.containsKey("term_review"),
|
||||
"a snapshot is a point-in-time copy, not a live view");
|
||||
}
|
||||
|
||||
// ── concurrency (CB-548 invariants hold under contention) ─────────────────────────────────
|
||||
|
||||
@Test
|
||||
void concurrentBindsNeverGiveASlotTwoTerminals() throws Exception {
|
||||
int n = 16;
|
||||
ExecutorService pool = Executors.newFixedThreadPool(n);
|
||||
try {
|
||||
CountDownLatch go = new CountDownLatch(1);
|
||||
List<Future<Boolean>> results = new ArrayList<>();
|
||||
for (int i = 0; i < n; i++) {
|
||||
final String term = "term_" + i; // every thread races for the SAME slot
|
||||
results.add(pool.submit(() -> {
|
||||
go.await();
|
||||
return registry.bind("lead-designer", term);
|
||||
}));
|
||||
}
|
||||
go.countDown();
|
||||
|
||||
int won = 0;
|
||||
for (Future<Boolean> r : results) {
|
||||
if (r.get()) {
|
||||
won++;
|
||||
}
|
||||
}
|
||||
assertEquals(1, won, "exactly one terminal may win the sole slot, got " + won);
|
||||
assertEquals(1, registry.snapshot().size(),
|
||||
"the slot hosts at most one terminal after the race");
|
||||
} finally {
|
||||
pool.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void concurrentBindsNeverPutOneTerminalInTwoSlots() throws Exception {
|
||||
int n = 16;
|
||||
ExecutorService pool = Executors.newFixedThreadPool(n);
|
||||
try {
|
||||
CountDownLatch go = new CountDownLatch(1);
|
||||
List<Future<String>> results = new ArrayList<>();
|
||||
for (int i = 0; i < n; i++) {
|
||||
final String slot = (i % 2 == 0) ? "lead-designer" : "reviewer"; // all race for ONE terminal
|
||||
results.add(pool.submit(() -> {
|
||||
go.await();
|
||||
return registry.bind(slot, "shared_term")
|
||||
? registry.slotForTerminal("shared_term") : null;
|
||||
}));
|
||||
}
|
||||
go.countDown();
|
||||
|
||||
// Rebinding the same terminal to the same slot is a harmless idempotent true, so count
|
||||
// winners is not the assertion — agreement is: every thread that reported success must
|
||||
// have seen the terminal in the SAME slot, never in two at once.
|
||||
String bound = null;
|
||||
boolean conflict = false;
|
||||
for (Future<String> r : results) {
|
||||
String s = r.get();
|
||||
if (s != null) {
|
||||
if (bound == null) {
|
||||
bound = s;
|
||||
} else if (!bound.equals(s)) {
|
||||
conflict = true;
|
||||
}
|
||||
}
|
||||
}
|
||||
assertFalse(conflict, "a terminal was observed in two slots at once");
|
||||
assertNotNull(bound, "at least one thread bound the terminal");
|
||||
assertEquals(1, registry.snapshot().size(),
|
||||
"the terminal occupies exactly one slot in the final snapshot");
|
||||
assertEquals(bound, registry.slotForTerminal("shared_term"));
|
||||
} finally {
|
||||
pool.shutdownNow();
|
||||
}
|
||||
}
|
||||
|
||||
/** A handed-over slot map is snapshotted at construction, not offered as live state. */
|
||||
@Test
|
||||
void theSlotSnapshotIsFixedByConstruction() {
|
||||
Map<String, BridgedConfig.Architect> mutable = new HashMap<>(SLOTS);
|
||||
ArchitectRegistry r = new ArchitectRegistry(mutable, Map::of);
|
||||
ArchitectRegistry r = new ArchitectRegistry(mutable);
|
||||
|
||||
mutable.put("hijack", new BridgedConfig.Architect("t", "gx10"));
|
||||
mutable.put("hijack", new BridgedConfig.Architect("gx10"));
|
||||
|
||||
assertFalse(r.isSlot("hijack"), "a handed-over map is not offered as live state");
|
||||
}
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
package dev.ltms.bridged.config;
|
||||
|
||||
import dev.ltms.bridged.auth.ArchitectRegistry;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
@@ -330,7 +331,7 @@ class BridgedConfigTest {
|
||||
// ── CB-548: the architects registry ────────────────────────────────────────────────────────
|
||||
|
||||
@Test
|
||||
void architectsBlockBindsSlotsByGatewayLocalName(@TempDir Path dir) throws Exception {
|
||||
void architectsBlockDeclaresSlotsByNameAndProfileOnly(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("architects.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
@@ -340,7 +341,6 @@ class BridgedConfigTest {
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
terminal: term_design
|
||||
profile: sonnet
|
||||
reviewer:
|
||||
profile: sonnet
|
||||
@@ -351,15 +351,15 @@ class BridgedConfigTest {
|
||||
"slot names are the keys — gateway-local unique by construction");
|
||||
assertEquals("sonnet", cfg.architects().get("lead-designer").profile(),
|
||||
"each slot carries its strong-model profile reference");
|
||||
assertEquals("term_design", cfg.architects().get("lead-designer").terminal());
|
||||
// A slot with no terminal binds nothing yet — the live binding may supply it later.
|
||||
assertTrue(cfg.architects().get("reviewer").terminal() == null
|
||||
|| cfg.architects().get("reviewer").terminal().isBlank());
|
||||
assertEquals("sonnet", cfg.architects().get("reviewer").profile());
|
||||
}
|
||||
|
||||
@Test
|
||||
void architectTerminalsMapsEachBoundSlotByItsPane(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("arch-terminals.yaml");
|
||||
void anArchitectCarriesNoConfigTerminalSoNothingIsRecognisedYet(@TempDir Path dir) throws Exception {
|
||||
// The corrected CB-548 premise: config declares slots (name + profile) only. A `terminal:`
|
||||
// key left over from the earlier premise is ignored — an architect is NOT recognised from
|
||||
// config the way a lead is, so it binds nothing at startup and resolves no architect.
|
||||
Path f = dir.resolve("arch-stale-terminal.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
@@ -370,26 +370,24 @@ class BridgedConfigTest {
|
||||
lead-designer:
|
||||
terminal: term_design
|
||||
profile: sonnet
|
||||
reviewer:
|
||||
terminal: term_review
|
||||
profile: sonnet
|
||||
unbound:
|
||||
profile: sonnet
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
assertEquals("sonnet", cfg.architects().get("lead-designer").profile(),
|
||||
"the profile is still read even when a stray terminal is ignored");
|
||||
|
||||
assertEquals(Map.of("term_design", "lead-designer", "term_review", "reviewer"),
|
||||
BridgedConfig.load(f).architectTerminals(),
|
||||
"a slot with no terminal registers no binding; the value is the slot name");
|
||||
// The registry built from this config owns no bindings: the slot is idle at startup.
|
||||
ArchitectRegistry r = new ArchitectRegistry(cfg.architects());
|
||||
assertTrue(r.snapshot().isEmpty());
|
||||
assertNull(r.slotForTerminal("term_design"),
|
||||
"a config terminal must not resolve an architect — slots start idle");
|
||||
}
|
||||
|
||||
@Test
|
||||
void noArchitectsBlockLeavesNothingBound(@TempDir Path dir) throws Exception {
|
||||
void noArchitectsBlockLeavesNothingConfigured(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-arch.yaml");
|
||||
Files.writeString(f, "bind:\n port: 8080\n");
|
||||
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
assertNull(cfg.architects());
|
||||
assertTrue(cfg.architectTerminals().isEmpty(),
|
||||
assertNull(BridgedConfig.load(f).architects(),
|
||||
"no architects: block ⇒ no architect identity, exactly as before CB-548");
|
||||
}
|
||||
|
||||
@@ -406,7 +404,6 @@ class BridgedConfigTest {
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
terminal: term_design
|
||||
profile: ltms-local
|
||||
""");
|
||||
|
||||
@@ -426,7 +423,6 @@ class BridgedConfigTest {
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
terminal: term_design
|
||||
profile: sonnet
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
@@ -447,7 +443,6 @@ class BridgedConfigTest {
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
terminal: term_design
|
||||
profile: ""
|
||||
""");
|
||||
BridgedConfig cfg = BridgedConfig.load(f);
|
||||
@@ -469,7 +464,6 @@ class BridgedConfigTest {
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
terminal: term_design
|
||||
profile: sonnet
|
||||
reviewer:
|
||||
profile: gx10
|
||||
@@ -486,6 +480,105 @@ class BridgedConfigTest {
|
||||
assertDoesNotThrow(() -> BridgedConfig.load(f).validateArchitects());
|
||||
}
|
||||
|
||||
@Test
|
||||
void duplicateArchitectSlotNamesAreRejectedAtParseTime(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("arch-dup.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
workers:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
""");
|
||||
|
||||
IllegalStateException e =
|
||||
assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("lead-designer"),
|
||||
"the refusal names the duplicated slot, was: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("duplicate architect"),
|
||||
"the refusal says the slot name is duplicated");
|
||||
}
|
||||
|
||||
@Test
|
||||
void duplicateKeysOutsideArchitectsAreUnaffected(@TempDir Path dir) throws Exception {
|
||||
// The duplicate check is scoped to the architects block — a duplicate elsewhere is not this
|
||||
// guard's concern and must not change parsing of the rest of the config.
|
||||
Path f = dir.resolve("dup-other.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
workers:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
""");
|
||||
// Last-wins for a non-architect duplicate is untouched: only the architects block is walked.
|
||||
assertEquals(Set.of("sonnet"), BridgedConfig.load(f).workerProfiles().keySet());
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNestedArchitectsFieldDoesNotSuppressRealDuplicateDetection(@TempDir Path dir) throws Exception {
|
||||
// A field ALSO named `architects` nested under another block carries its own duplicate and
|
||||
// sits BEFORE the real top-level block. Only the top-level block is ever inspected: the
|
||||
// refusal must name the real slot (lead-designer), not the nested one (nested-slot).
|
||||
Path f = dir.resolve("nested-arch.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
architects:
|
||||
nested-slot:
|
||||
k: v
|
||||
nested-slot:
|
||||
k: v
|
||||
architects:
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
""");
|
||||
|
||||
IllegalStateException e =
|
||||
assertThrows(IllegalStateException.class, () -> BridgedConfig.load(f));
|
||||
assertTrue(e.getMessage().contains("lead-designer"),
|
||||
"the real top-level duplicate must be reported, was: " + e.getMessage());
|
||||
assertTrue(e.getMessage().contains("duplicate architect"),
|
||||
"the refusal says the slot name is duplicated");
|
||||
}
|
||||
|
||||
@Test
|
||||
void nestedDuplicateFieldsInsideASlotAreNotDuplicateSlotNames(@TempDir Path dir) throws Exception {
|
||||
// A duplicated field nested inside one slot's own value (here inside an ignored `extra:`
|
||||
// sub-block) is not a duplicate SLOT name — it must not be rejected as one. Only the direct
|
||||
// child keys of the architects mapping are slot names; whatever is deeper is the slot's
|
||||
// business and must not masquerade as a duplicate slot.
|
||||
Path f = dir.resolve("nested-dup-inside-slot.yaml");
|
||||
Files.writeString(f, """
|
||||
bind:
|
||||
port: 8080
|
||||
workers:
|
||||
sonnet:
|
||||
baseUrl: http://gx10.gw:8000
|
||||
architects:
|
||||
lead-designer:
|
||||
profile: sonnet
|
||||
extra:
|
||||
a: 1
|
||||
a: 1
|
||||
""");
|
||||
|
||||
BridgedConfig cfg = assertDoesNotThrow(() -> BridgedConfig.load(f));
|
||||
assertDoesNotThrow(cfg::validateArchitects,
|
||||
"a nested duplicate inside a slot is not a duplicate slot and must not refuse startup");
|
||||
assertEquals(Set.of("lead-designer"), cfg.architects().keySet());
|
||||
assertEquals("sonnet", cfg.architects().get("lead-designer").profile());
|
||||
}
|
||||
|
||||
@Test
|
||||
void absentBrokerBlockLeavesInboxSoftState(@TempDir Path dir) throws Exception {
|
||||
Path f = dir.resolve("no-broker.yaml");
|
||||
|
||||
@@ -286,10 +286,12 @@ class CompletionResolverTest {
|
||||
// The turn as the injector captured it at delivery (waiter + pre-turn baseline).
|
||||
var turnN = new CompletionResolver.InFlight(waiterN, "an earlier answer");
|
||||
|
||||
// Turn N is resolved by the worker's explicit reply.
|
||||
// Turn N is resolved by the worker's explicit reply, and its send deregisters the waiter.
|
||||
assertTrue(rendezvous.resolve("term_a", "N replied"));
|
||||
rendezvous.close("term_a", waiterN); // the sender's finally, before the next turn opens
|
||||
|
||||
// Turn N+1's send opens its own waiter on the same session (replacing the registered one).
|
||||
// Turn N+1's send opens its own waiter on the same session (CB-548: open fails if the
|
||||
// previous waiter is still registered, so a clean turn deregisters it first as above).
|
||||
var waiterN1 = rendezvous.open("term_a");
|
||||
|
||||
resolver.resolve("term_a", turnN); // turn N's completion fallback finally fires
|
||||
|
||||
@@ -484,4 +484,37 @@ class BridgeMcpTest {
|
||||
assertTrue(out.contains("\"architect\":\"lead-designer\""), out);
|
||||
assertTrue(out.contains("\"sessionId\":\"term_design\""), out);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-548: an architect SEND delegates as its own pane (recording the per-target delegation) but
|
||||
* must NEVER become the legacy singleton "primary" fallback — the per-target map does not cure
|
||||
* the singleton, so an architect left there would draw no-delegation inbox nudges meant for a
|
||||
* primary. Only PRIMARY callers (the unnamed primary and named leads alike) may claim it, and
|
||||
* the decision keys on the resolved role, not name/kind sniffing.
|
||||
*/
|
||||
@Test
|
||||
void architectSendDoesNotClaimThePrimarySingletonButALeadSendStillCan() {
|
||||
// Architect SEND: does not change the legacy primary fallback.
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
BridgeMcp.recordPrimarySingleton(reg, "term_design", Principal.architect("design", "term_design", 400));
|
||||
assertTrue(reg.primaryTerminal().isEmpty(),
|
||||
"an architect must never become the legacy primary fallback");
|
||||
|
||||
// Lead SEND (a named PRIMARY) still claims it — preserved from CB-530/CB-532.
|
||||
PrimaryRegistry leadReg = new PrimaryRegistry(null);
|
||||
BridgeMcp.recordPrimarySingleton(leadReg, "term_lead_opus", Principal.leader("opus", "term_lead_opus", 100));
|
||||
assertEquals("term_lead_opus", leadReg.primaryTerminal().orElseThrow(),
|
||||
"a named lead is a primary and may claim the fallback");
|
||||
|
||||
// Unnamed primary likewise.
|
||||
PrimaryRegistry primaryReg = new PrimaryRegistry(null);
|
||||
BridgeMcp.recordPrimarySingleton(primaryReg, "term_p", Principal.primary(50));
|
||||
assertEquals("term_p", primaryReg.primaryTerminal().orElseThrow(),
|
||||
"an unnamed primary may claim the fallback");
|
||||
|
||||
// A null caller (legacy/no-auth path) records nothing.
|
||||
PrimaryRegistry legacy = new PrimaryRegistry(null);
|
||||
BridgeMcp.recordPrimarySingleton(legacy, "term_x", null);
|
||||
assertTrue(legacy.primaryTerminal().isEmpty(), "no caller means nothing is recorded");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -5,6 +5,7 @@ import dev.ltms.bridged.herdr.AgentStatus;
|
||||
import dev.ltms.bridged.herdr.FakeHerdr;
|
||||
import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.inject.CompletionResolver;
|
||||
import dev.ltms.bridged.mcp.PrimaryRegistry;
|
||||
import dev.ltms.bridged.inject.Injector;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -16,6 +17,7 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
@@ -321,6 +323,141 @@ class MessageServiceTest {
|
||||
assertEquals("first done", firstReply.text());
|
||||
}
|
||||
|
||||
// --- CB-548: delegator ownership is recorded only on an ACCEPTED send ----------------------
|
||||
|
||||
private static final String LEAD_L = "term_lead_l";
|
||||
private static final String LEAD_A = "term_lead_a";
|
||||
|
||||
/**
|
||||
* The bug CB-548 fixes: L holds worker W, then architect A attempts W and times out BUSY. With
|
||||
* delegator ownership recorded at {@code bridge_send} <em>request</em> time, A's rejected call
|
||||
* would overwrite L — and W's late no-waiter reply would be pushed to A, who never owned the
|
||||
* turn. The accepted-delivery hook must not fire for a BUSY send, so L stays the delegator.
|
||||
*/
|
||||
@Test
|
||||
void busySenderDoesNotBecomeTheDelegatingOwner() throws Exception {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
// L accepts a delegation to W: the send wins the lock and queues delivery → L is recorded.
|
||||
CompletableFuture<MessageService.Reply> first = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
|
||||
awaitWaiting();
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"an accepted send owns the delegation");
|
||||
|
||||
// A attempts W while L holds it → BUSY (lock never taken) → its hook never fires.
|
||||
MessageService.Reply busy = messages.send(T, "second", 100, () -> reg.recordDelegation(T, LEAD_A));
|
||||
assertEquals(MessageService.Outcome.BUSY, busy.outcome());
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"a BUSY send must not steal the delegator ownership it never earned");
|
||||
|
||||
// L completes so the test thread is not left pinned.
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "first done"));
|
||||
MessageService.Reply firstReply = first.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.REPLIED, firstReply.outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* Once L's accepted delegation is fully done, a later <em>accepted</em> send from A may
|
||||
* legitimately become the new delegator — ownership follows the turn, not the first caller.
|
||||
*/
|
||||
@Test
|
||||
void anAcceptedSendAfterThePriorOwnerFinishesBecomesTheNewOwner() throws Exception {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
CompletableFuture<MessageService.Reply> first = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "first", 5000, () -> reg.recordDelegation(T, LEAD_L)));
|
||||
awaitWaiting();
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "first done"));
|
||||
MessageService.Reply firstReply = first.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.REPLIED, firstReply.outcome());
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owned the first turn");
|
||||
|
||||
// L finished; A's later accepted send takes the delegation over.
|
||||
CompletableFuture<MessageService.Reply> second = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "second", 5000, () -> reg.recordDelegation(T, LEAD_A)));
|
||||
awaitWaiting();
|
||||
assertEquals(LEAD_A, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"an accepted send after the owner finished becomes the new delegator");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
injector.onStatus(T, AgentStatus.WORKING);
|
||||
assertTrue(rendezvous.resolve(T, "second done"));
|
||||
MessageService.Reply secondReply = second.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.REPLIED, secondReply.outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-548 requirement: answering an existing {@code bridge_ask} is the SAME delegation, so it must
|
||||
* not rewrite ownership. L accepted the send (owned), the worker paused to ask, and L answers via
|
||||
* turnId — ownership stays L throughout; the answer path never touches the registry.
|
||||
*/
|
||||
@Test
|
||||
void answeringAnAskDoesNotRewriteDelegatorOwnership() throws Exception {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
CompletableFuture<MessageService.Reply> send = CompletableFuture.supplyAsync(
|
||||
() -> messages.send(T, "do X", 5000, () -> reg.recordDelegation(T, LEAD_L)));
|
||||
awaitWaiting();
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(), "L owns the delegation");
|
||||
|
||||
injector.onStatus(T, AgentStatus.IDLE); // deliver
|
||||
injector.onStatus(T, AgentStatus.WORKING); // worker picks it up, then pauses to ask
|
||||
|
||||
CompletableFuture<MessageService.AskResult> ask =
|
||||
CompletableFuture.supplyAsync(() -> messages.ask(T, "which config?", 5000));
|
||||
MessageService.Reply q = send.get(5, TimeUnit.SECONDS);
|
||||
assertEquals(MessageService.Outcome.QUESTION, q.outcome());
|
||||
assertNotNull(q.turnId());
|
||||
|
||||
// L answers the ask on the same turn; the answer path must not touch ownership.
|
||||
CompletableFuture<MessageService.Reply> answer =
|
||||
CompletableFuture.supplyAsync(() -> messages.answer(q.turnId(), "config.yaml", 5000));
|
||||
assertEquals("config.yaml", ask.get(5, TimeUnit.SECONDS).answer());
|
||||
awaitWaiting(); // the answering send reopened its forward waiter
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"answering an ask keeps L as the delegator — ownership is not rewritten");
|
||||
|
||||
assertTrue(rendezvous.resolve(T, "done"));
|
||||
assertEquals(MessageService.Outcome.REPLIED, answer.get(5, TimeUnit.SECONDS).outcome());
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-548: {@code onAccepted} is a public callback, so a throwing one must not orphan the turn.
|
||||
* The waiter is opened first, then the hook runs BEFORE delivery is queued — so a throw fails
|
||||
* the send loudly, closes its waiter, and never enqueues a message the worker would pick up and
|
||||
* reply into the void.
|
||||
*/
|
||||
@Test
|
||||
void aThrowingAcceptedHookLeavesNoStaleWaiterOrQueuedOrphan() {
|
||||
assertThrows(IllegalStateException.class,
|
||||
() -> messages.send(T, "doomed", 500,
|
||||
() -> { throw new IllegalStateException("ownership hook failed"); }),
|
||||
"a throwing ownership hook fails the send loudly");
|
||||
|
||||
assertFalse(rendezvous.isWaiting(T), "the failed send must not leave a stale rendezvous waiter");
|
||||
// Give the injector a delivery window: with nothing enqueued, nothing may reach the worker.
|
||||
injector.onStatus(T, AgentStatus.IDLE);
|
||||
boolean doomedQueued = herdr.calls.stream()
|
||||
.anyMatch(c -> c.method().equals("agent.prompt")
|
||||
&& String.valueOf(c.params()).contains("doomed"));
|
||||
assertFalse(doomedQueued, "a throwing ownership hook must not leave a queued, orphanable message");
|
||||
}
|
||||
|
||||
/**
|
||||
* The async (fire-and-poll) path runs the same {@code send} on a background thread, so the
|
||||
* accepted-delivery hook must thread through it — ownership is recorded exactly as blocking sends.
|
||||
*/
|
||||
@Test
|
||||
void asyncSendRecordsOwnershipOnAcceptance() throws Exception {
|
||||
PrimaryRegistry reg = new PrimaryRegistry(null);
|
||||
messages.sendAsync(T, "async task", () -> reg.recordDelegation(T, LEAD_L));
|
||||
awaitWaiting(); // the background send won the lock, queued, and opened its waiter
|
||||
assertEquals(LEAD_L, reg.nudgeTargetFor(T).orElseThrow(),
|
||||
"the async path records delegator ownership on acceptance, like the blocking path");
|
||||
}
|
||||
|
||||
// --- CB-307 reply inbox ----------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
|
||||
@@ -8,6 +8,8 @@ import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertNotEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertSame;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
@@ -20,6 +22,26 @@ class RendezvousTest {
|
||||
|
||||
private final Rendezvous rendezvous = new Rendezvous();
|
||||
|
||||
/**
|
||||
* CB-548: an {@code open} is atomic fail-if-present so a double open can never replace the first
|
||||
* waiter another send is blocked on. {@code MessageService} serializes sends per session, so in
|
||||
* correct code this cannot happen — the rejection is a loud tripwire for an invariant violation,
|
||||
* and the first waiter must survive it and stay resolvable.
|
||||
*/
|
||||
@Test
|
||||
void openRejectsADoubleOpenAndKeepsTheFirstWaiterRegisteredAndResolvable() {
|
||||
CompletableFuture<Rendezvous.Resolution> first = rendezvous.open(W);
|
||||
|
||||
assertThrows(IllegalStateException.class, () -> rendezvous.open(W),
|
||||
"a second open while one is registered is rejected loudly, not a silent replace");
|
||||
|
||||
assertSame(first, rendezvous.currentWaiter(W), "the first waiter remains the registered one");
|
||||
assertTrue(rendezvous.resolve(W, "first wins"), "the first waiter is still resolvable");
|
||||
Rendezvous.Resolution r = first.getNow(null);
|
||||
assertEquals(Rendezvous.Kind.REPLY, r.kind());
|
||||
assertEquals("first wins", r.text(), "the resolution lands on the first waiter, not the rejected one");
|
||||
}
|
||||
|
||||
@Test
|
||||
void openAskMintsAUniqueTurnScopedToItsSessionAndCoalescesDuplicates() {
|
||||
Rendezvous.AskTicket t1 = rendezvous.openAsk(W);
|
||||
|
||||
@@ -26,6 +26,18 @@ class GitWorktreesTest {
|
||||
}
|
||||
""";
|
||||
|
||||
/** An opencode config carrying a {@code {file:.secrets/...}} reference — CB-543's crash repro. */
|
||||
private static final String OPENCODE_WITH_FILE_REF = """
|
||||
{
|
||||
"env": {
|
||||
"CONTEXT7_TOKEN": "{file:.secrets/context7-token}"
|
||||
}
|
||||
}
|
||||
""";
|
||||
|
||||
/** A non-empty autoenv file — the form that would prompt for authorization in a worktree. */
|
||||
private static final String AUTOENV_WITH_DIRECTIVE = "export HELLO=world\n";
|
||||
|
||||
private static Path initRepo(Path dir) throws Exception {
|
||||
Files.createDirectories(dir);
|
||||
git(dir, "init", "-q", "-b", "main");
|
||||
@@ -47,9 +59,9 @@ class GitWorktreesTest {
|
||||
assertEquals(0, p.exitValue(), "git " + String.join(" ", args) + " failed:\n" + out);
|
||||
}
|
||||
|
||||
/** Pending changes to {@code .mcp.json} in {@code cwd}, empty when git considers it unmodified. */
|
||||
private static String mcpStatus(Path cwd) throws Exception {
|
||||
Process p = new ProcessBuilder("git", "status", "--porcelain", "--", ".mcp.json")
|
||||
/** Pending changes to {@code file} in {@code cwd}, empty when git considers it unmodified. */
|
||||
private static String status(Path cwd, String file) throws Exception {
|
||||
Process p = new ProcessBuilder("git", "status", "--porcelain", "--", file)
|
||||
.directory(cwd.toFile()).redirectErrorStream(true).start();
|
||||
String out = new String(p.getInputStream().readAllBytes());
|
||||
assertTrue(p.waitFor(30, TimeUnit.SECONDS), "git status timed out");
|
||||
@@ -82,7 +94,7 @@ class GitWorktreesTest {
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "cb-525-b", "HEAD");
|
||||
|
||||
assertEquals("", mcpStatus(Path.of(wt)),
|
||||
assertEquals("", status(Path.of(wt), ".mcp.json"),
|
||||
"the neutralized .mcp.json shows as modified — --skip-worktree did not take");
|
||||
}
|
||||
|
||||
@@ -116,4 +128,81 @@ class GitWorktreesTest {
|
||||
String body = Files.readString(Path.of(wt).resolve(".mcp.json"));
|
||||
assertTrue(body.replaceAll("\\s+", "").contains("\"mcpServers\":{}"), body);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-543's repro: a tracked {@code opencode.json} carries a {@code {file:.secrets/...}} reference
|
||||
* to a gitignored secret that never reaches a worktree, and opencode refuses to start on it. The
|
||||
* worktree's copy must be neutralized and hidden like {@code .mcp.json}.
|
||||
*/
|
||||
@Test
|
||||
void aTrackedOpencodeConfigIsNeutralizedAndHidden(@TempDir Path tmp) throws Exception {
|
||||
Path repo = tmp.resolve("repo");
|
||||
Files.createDirectories(repo);
|
||||
git(repo, "init", "-q", "-b", "main");
|
||||
git(repo, "config", "user.email", "test@example.invalid");
|
||||
git(repo, "config", "user.name", "Test");
|
||||
Files.writeString(repo.resolve(".mcp.json"), WITH_SERVERS);
|
||||
Files.writeString(repo.resolve("opencode.json"), OPENCODE_WITH_FILE_REF);
|
||||
Files.writeString(repo.resolve("README.md"), "seed\n");
|
||||
git(repo, "add", ".mcp.json", "opencode.json", "README.md");
|
||||
git(repo, "commit", "-q", "-m", "seed");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "cb-543-a", "HEAD");
|
||||
|
||||
String body = Files.readString(Path.of(wt).resolve("opencode.json"));
|
||||
assertFalse(body.contains(".secrets"),
|
||||
"worktree kept a dangling {file:...} secret reference:\n" + body);
|
||||
assertEquals("{}", body.replaceAll("\\s+", ""),
|
||||
"expected an empty JSON object stub, got:\n" + body);
|
||||
assertEquals("", status(Path.of(wt), "opencode.json"),
|
||||
"the neutralized opencode.json shows as modified — --skip-worktree did not take");
|
||||
}
|
||||
|
||||
/** A config the repo does not carry must be skipped — no stub invented, provisioning still succeeds. */
|
||||
@Test
|
||||
void anAbsentConfigIsSkippedWithoutError(@TempDir Path tmp) throws Exception {
|
||||
Path repo = initRepo(tmp.resolve("repo")); // only .mcp.json + README are committed
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "cb-543-b", "HEAD");
|
||||
|
||||
assertFalse(Files.exists(Path.of(wt).resolve("opencode.json")),
|
||||
"a stub was invented for a config the repo does not carry");
|
||||
assertFalse(Files.exists(Path.of(wt).resolve(".autoenv")),
|
||||
"a stub was invented for a config the repo does not carry");
|
||||
// .mcp.json's long-standing create-always behaviour must be unchanged.
|
||||
assertTrue(Files.exists(Path.of(wt).resolve(".mcp.json")), ".mcp.json stub was dropped");
|
||||
}
|
||||
|
||||
/** All three protected configs are covered: each one present in a worktree is neutralized and hidden. */
|
||||
@Test
|
||||
void allThreeConfigsAreNeutralizedWhenPresent(@TempDir Path tmp) throws Exception {
|
||||
Path repo = tmp.resolve("repo");
|
||||
Files.createDirectories(repo);
|
||||
git(repo, "init", "-q", "-b", "main");
|
||||
git(repo, "config", "user.email", "test@example.invalid");
|
||||
git(repo, "config", "user.name", "Test");
|
||||
Files.writeString(repo.resolve(".mcp.json"), WITH_SERVERS);
|
||||
Files.writeString(repo.resolve("opencode.json"), OPENCODE_WITH_FILE_REF);
|
||||
Files.writeString(repo.resolve(".autoenv"), AUTOENV_WITH_DIRECTIVE);
|
||||
Files.writeString(repo.resolve("README.md"), "seed\n");
|
||||
git(repo, "add", ".mcp.json", "opencode.json", ".autoenv", "README.md");
|
||||
git(repo, "commit", "-q", "-m", "seed");
|
||||
|
||||
String wt = new GitWorktrees(tmp.resolve("wts").toString())
|
||||
.add(repo.toString(), "cb-543-c", "HEAD");
|
||||
|
||||
assertTrue(Files.readString(Path.of(wt).resolve(".mcp.json"))
|
||||
.replaceAll("\\s+", "").contains("\"mcpServers\":{}"),
|
||||
".mcp.json was not neutralized");
|
||||
assertEquals("{}", Files.readString(Path.of(wt).resolve("opencode.json")).replaceAll("\\s+", ""),
|
||||
"opencode.json was not neutralized");
|
||||
assertEquals("", Files.readString(Path.of(wt).resolve(".autoenv")),
|
||||
".autoenv was not neutralized");
|
||||
|
||||
assertEquals("", status(Path.of(wt), ".mcp.json"), ".mcp.json still shows as modified");
|
||||
assertEquals("", status(Path.of(wt), "opencode.json"), "opencode.json still shows as modified");
|
||||
assertEquals("", status(Path.of(wt), ".autoenv"), ".autoenv still shows as modified");
|
||||
}
|
||||
}
|
||||
|
||||
@@ -37,7 +37,7 @@ class OpenCodeLauncherTest {
|
||||
private OpenCodeLauncher service(FakeHerdr herdr, Path configRoot, BridgedConfig.Worker cfg) {
|
||||
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);
|
||||
0, System::currentTimeMillis, () -> { }, configRoot, configRoot);
|
||||
}
|
||||
|
||||
@SuppressWarnings("unchecked")
|
||||
@@ -130,14 +130,65 @@ class OpenCodeLauncherTest {
|
||||
@Test
|
||||
void capabilitiesDeclareOrphanReapAndMcpAskAndConditionalSelfPr(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
assertEquals(java.util.Set.of(Capability.MID_TURN_ASK, Capability.WORKTREE, Capability.ORPHAN_REAP),
|
||||
assertEquals(java.util.Set.of(Capability.MID_TURN_ASK, Capability.WORKTREE, Capability.ORPHAN_REAP,
|
||||
Capability.SESSION_RESUME),
|
||||
service(herdr, root, opencodeCfg(null, null, null)).capabilities(),
|
||||
"no git token → no SELF_PR");
|
||||
"opencode can be resumed by its own session id, so SESSION_RESUME is always declared");
|
||||
assertFalse(service(herdr, root, opencodeCfg(null, null, null))
|
||||
.capabilities().contains(Capability.SESSION_NAME),
|
||||
"opencode has no display-name flag, so SESSION_NAME must NOT be declared");
|
||||
assertTrue(service(herdr, root, opencodeCfg(null, null, "GITEA_ACCESS_TOKEN"))
|
||||
.capabilities().contains(Capability.SELF_PR),
|
||||
"a git-token profile adds SELF_PR");
|
||||
}
|
||||
|
||||
// --- CB-547: resume + post-hoc session discovery --------------------------------------------
|
||||
|
||||
@Test
|
||||
void aResumeSpawnPassesTheSessionIdAsDashS(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
service(herdr, root, opencodeCfg("google/gemini-2.5-pro", null, null))
|
||||
.spawn(new SpawnRequest(null, null, null, null, "ses_41b79fc90ffeI9E8uZv6VprUn2"));
|
||||
|
||||
List<String> args = startArgs(herdr);
|
||||
int s = args.indexOf("-s");
|
||||
assertTrue(s >= 0, "a resumed spawn carries opencode's -s flag");
|
||||
assertEquals("ses_41b79fc90ffeI9E8uZv6VprUn2", args.get(s + 1),
|
||||
"the resume target id follows -s");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aFreshSpawnCarriesNoSessionFlag(@TempDir Path root) {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
service(herdr, root, opencodeCfg(null, null, null))
|
||||
.spawn(new SpawnRequest(null, null, null, null, null));
|
||||
|
||||
assertFalse(startArgs(herdr).contains("-s"),
|
||||
"no resume target → a fresh session with no -s flag");
|
||||
}
|
||||
|
||||
@Test
|
||||
void theHandleDiscoversTheSessionIdForTheWorkersCwdOnlyAfterItAppears(@TempDir Path root,
|
||||
@TempDir Path discRoot)
|
||||
throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
OpenCodeLauncher launcher = new OpenCodeLauncher(new AgentControl(herdr),
|
||||
new WorkspaceControl(herdr), Map.of("gemini", opencodeCfg(null, null, null)),
|
||||
"gemini", _ -> null, 0, System::currentTimeMillis, () -> { }, root, discRoot);
|
||||
|
||||
PeerHandle handle = launcher.spawn(new SpawnRequest(null, "/work/dir", null));
|
||||
|
||||
// opencode writes the record only when the session is first persisted — the instant the
|
||||
// pane is ready it does not exist, so agentSessionId() is null (never a spawn failure).
|
||||
assertNull(handle.agentSessionId(), "no record yet → null, not a spawn-time block");
|
||||
// Once the record appears (here: same cwd), lazy discovery resolves it — the handle's
|
||||
// session id matches its own worktree, not another's.
|
||||
OpenCodeSessionDiscoveryTest.writeRecord(discRoot, "p1", "ses_a.json",
|
||||
"ses_resolved", "/work/dir", 1000L);
|
||||
assertEquals("ses_resolved", handle.agentSessionId(),
|
||||
"agentSessionId() re-scans and picks up a record that has since been written");
|
||||
}
|
||||
|
||||
@Test
|
||||
void foreignWorkerMatchesOpencodePrefixButNotClaude() {
|
||||
String nonce = "abc123";
|
||||
@@ -169,7 +220,7 @@ class OpenCodeLauncherTest {
|
||||
long[] clock = {0};
|
||||
OpenCodeLauncher svc = new OpenCodeLauncher(new AgentControl(herdr), new WorkspaceControl(herdr),
|
||||
Map.of("gemini", opencodeCfg(null, null, null)), "gemini", _ -> null,
|
||||
1000, () -> clock[0], () -> clock[0] += 50, root);
|
||||
1000, () -> clock[0], () -> clock[0] += 50, root, root);
|
||||
|
||||
PeerUnreachableException ex = assertThrows(PeerUnreachableException.class,
|
||||
() -> svc.spawn(new SpawnRequest(null, null, null)));
|
||||
|
||||
@@ -0,0 +1,90 @@
|
||||
package dev.ltms.bridged.worker;
|
||||
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.io.TempDir;
|
||||
|
||||
import java.nio.file.Files;
|
||||
import java.nio.file.Path;
|
||||
import java.nio.file.attribute.FileTime;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.*;
|
||||
|
||||
/**
|
||||
* {@link OpenCodeSessionDiscovery} matches an opencode session record by the worker's cwd (its
|
||||
* {@code directory}) against opencode's on-disk storage. These tests populate a TEMP storage root
|
||||
* themselves — never the operator's real {@code ~/.local/share/opencode}.
|
||||
*/
|
||||
class OpenCodeSessionDiscoveryTest {
|
||||
|
||||
/**
|
||||
* Write a session record {@code {"id":..., "directory":...}} under
|
||||
* {@code <root>/session/<projectID>/<fileName>} and stamp it with a known last-modified time,
|
||||
* so "most recently modified wins" is deterministic. Static so the launcher test can reuse it.
|
||||
*/
|
||||
static void writeRecord(Path root, String projectId, String fileName, String id,
|
||||
String directory, long lastModifiedEpochMillis) throws Exception {
|
||||
Path dir = root.resolve("session").resolve(projectId);
|
||||
Files.createDirectories(dir);
|
||||
Path file = dir.resolve(fileName);
|
||||
Files.writeString(file, "{\"id\":\"" + id + "\",\"directory\":\"" + directory
|
||||
+ "\",\"projectID\":\"" + projectId + "\",\"version\":\"1.1.31\"}");
|
||||
Files.setLastModifiedTime(file, FileTime.fromMillis(lastModifiedEpochMillis));
|
||||
}
|
||||
|
||||
@Test
|
||||
void findsTheRecordWhoseDirectoryEqualsTheCwd(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
||||
writeRecord(root, "p2", "ses_b.json", "ses_bbb", "/w/b", 2000L);
|
||||
|
||||
assertEquals("ses_bbb", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/b"),
|
||||
"the record whose directory equals the cwd is the one found");
|
||||
assertEquals("ses_aaa", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aNonMatchingDirectoryYieldsNullRatherThanAMismatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "ses_a.json", "ses_aaa", "/w/a", 1000L);
|
||||
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/other"),
|
||||
"no record for this cwd yet → null, not a wrong session");
|
||||
}
|
||||
|
||||
@Test
|
||||
void prefersTheMostRecentlyModifiedRecordWhenSeveralMatch(@TempDir Path root) throws Exception {
|
||||
writeRecord(root, "p1", "old.json", "ses_old", "/w/a", 1000L);
|
||||
writeRecord(root, "p2", "new.json", "ses_new", "/w/a", 5000L);
|
||||
|
||||
assertEquals("ses_new", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"the freshest record for the cwd wins");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMissingOrEmptyStorageRootYieldsNullWithoutThrowing(@TempDir Path root) throws Exception {
|
||||
// Missing: no session dir at all under the root.
|
||||
assertNull(new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"));
|
||||
|
||||
// Present but empty: a session dir with nothing in it produces no match, not a throw.
|
||||
Path emptyRoot = root.resolve("empty");
|
||||
Files.createDirectories(emptyRoot.resolve("session"));
|
||||
assertNull(new OpenCodeSessionDiscovery(emptyRoot).sessionIdForDirectory("/w/a"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aBlankOrNullDirectoryYieldsNull(@TempDir Path root) {
|
||||
OpenCodeSessionDiscovery discovery = new OpenCodeSessionDiscovery(root);
|
||||
assertNull(discovery.sessionIdForDirectory(null));
|
||||
assertNull(discovery.sessionIdForDirectory(" "));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aMalformedRecordIsSkippedRatherThanFatal(@TempDir Path root) throws Exception {
|
||||
// A record that fails to parse must not abort the scan of its siblings.
|
||||
Path dir = root.resolve("session").resolve("p1");
|
||||
Files.createDirectories(dir);
|
||||
Files.writeString(dir.resolve("broken.json"), "{not valid json");
|
||||
writeRecord(root, "p1", "good.json", "ses_good", "/w/a", 1000L);
|
||||
|
||||
assertEquals("ses_good", new OpenCodeSessionDiscovery(root).sessionIdForDirectory("/w/a"),
|
||||
"an unreadable record is skipped; a later valid one still matches");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user