Compare commits
10 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 83cac07f6e | |||
| f472e0f782 | |||
| 91c9f981c5 | |||
| dd906526c0 | |||
| ec3001796a | |||
| 86cf4c285a | |||
| cd18887b69 | |||
| 4637c68295 | |||
| 0114bd1fa7 | |||
| 6dee84ca71 |
@@ -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;
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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");
|
||||
}
|
||||
}
|
||||
@@ -1,6 +1,8 @@
|
||||
# CB-500 — Multi-Tier Coordination (Stage 6)
|
||||
|
||||
**Status:** design note (proposal — ticket split deferred)
|
||||
**Status:** design note. Developments A/B remain proposals; Development C (§6 and Figures 7–8) is
|
||||
**SUPERSEDED** by the advisory-architect design in Gitea issue #16 and the `architects:` configuration
|
||||
block (CB-548).
|
||||
**Depends on:** CB-401/402 (Peer Launcher SPI + composite router — placement-neutral spawn),
|
||||
CB-308 (per-agent broker channels + global id + federated roster — the addressing substrate),
|
||||
CB-307 (durable inbox + push loop), CB-301/303 (session FSM + context-cap/idle-ttl), CB-304
|
||||
@@ -16,8 +18,9 @@ workers) into a **multi-tier** one, along three axes the lead has asked for:
|
||||
1. **Sandboxed workers** — each worker runs in a **separated, peer-owned sandbox** carrying its own
|
||||
toolchain (Claude routed via `ANTHROPIC_BASE_URL`, a headless IDE, git, MCP, dev-tools), with
|
||||
**per-role** sandboxes (a backend-agent image, a frontend-agent image).
|
||||
2. **Main-agent pairs** — the "main" tier becomes a **pair** (on-subscription Opus + one cloud
|
||||
module) collaborating, instead of a lone primary.
|
||||
2. **Main-agent pairs** — **SUPERSEDED.** The considered model made the "main" tier a pair
|
||||
(on-subscription Opus + one cloud module). The actual fleet is one human-driven lead plus two
|
||||
short-lived advisory architects on different model families.
|
||||
3. **An orchestrator tier** — a supervisor **above** the mains that owns their **session identity**
|
||||
(naming, resume) and **curates context**, so every main→worker delegation carries the *exact*
|
||||
slice of context it needs and nothing else.
|
||||
@@ -63,6 +66,10 @@ Four concrete bake-ins assume a single tier:
|
||||
|
||||
## 3. Target multi-tier architecture
|
||||
|
||||
> **SUPERSEDED fleet sketch.** Figure 2 records the former two-main model. The actual fleet is one
|
||||
> lead, two independent advisory architects, and N workers; architects are sideways peers, not leads
|
||||
> and not a tier above the lead. See Gitea issue #16 and the `architects:` block.
|
||||
|
||||
```mermaid
|
||||
flowchart TB
|
||||
human["human"]
|
||||
@@ -93,10 +100,9 @@ flowchart TB
|
||||
chan --- roster
|
||||
```
|
||||
|
||||
*Figure 2 — three tiers. Tier 0 owns the mains' session identity + context scope; Tier 1 is a
|
||||
collaborating pair, each an MCP client with its own pull inbox; Tier 2 is peer-owned sandboxes the
|
||||
bus launches into. The middle is CB-308's per-agent-channel + federated-roster substrate, now
|
||||
carrying tier-to-tier traffic, not just host-to-host.*
|
||||
*Figure 2 — **SUPERSEDED historical fleet sketch.** It proposed a collaborating pair of managed mains.
|
||||
The actual fleet keeps one human-driven lead and uses two independent, short-lived advisory architects
|
||||
on different model families, so agreement is evidence rather than correlated echo.*
|
||||
|
||||
The recursion is the key idea: **`orchestrator : mains :: main : workers`** — the same
|
||||
spawn/name/resume/scope verbs at two levels.
|
||||
@@ -166,6 +172,11 @@ gateway, because herdr keystroke-injection needs a locally-owned PTY.**
|
||||
|
||||
## 5. Development B — Main-agent pairs
|
||||
|
||||
> **SUPERSEDED — do not implement this model.** The two-main fleet was replaced by one human-driven
|
||||
> lead and two independent advisory architects. They are deliberately different model families (Claude
|
||||
> Sonnet 5 and GPT-5.6 through opencode), receive the same brief, and work independently so agreement
|
||||
> is evidence rather than correlated echo. See Gitea issue #16 and `architects:`.
|
||||
|
||||
Both mains are MCP **clients**, so **neither can be called into** — each needs a **pull-based
|
||||
per-agent inbox**, which is precisely CB-308 item #1 (per-agent AMQP channels). The primary machinery
|
||||
that is singular today (single-slot `PrimaryRegistry`, a push-loop aimed at one terminal, "these
|
||||
@@ -221,6 +232,19 @@ push-loop fan-out; relax "orchestration tools only the primary calls" to "any re
|
||||
|
||||
## 6. Development C — Orchestrator tier
|
||||
|
||||
> **SUPERSEDED — do not implement this model.** The operator rejected a supervisor above the lead.
|
||||
> The human continues to drive the pre-existing lead directly; bridged neither spawns nor resumes that
|
||||
> lead. What replaced this proposal is **one lead, two short-lived advisory architects, and N workers**:
|
||||
> the lead engages architects sideways for a strong-model assessment, then discards them. Architect
|
||||
> slots are declared in `architects:` (see Gitea issue #16), rather than making leads managed sessions.
|
||||
> The two architects deliberately use different model families — Claude Sonnet 5 and GPT-5.6 through
|
||||
> opencode — and receive the same brief independently. Agreement is evidence, not correlated echo
|
||||
> from one provider or one conversation.
|
||||
|
||||
> **Historical alternative retained.** The text and figures below record the considered model and why it
|
||||
> was rejected: it re-rooted the human-facing session above the lead, violating the still-true premise
|
||||
> that configured leaders pre-exist, are recognised, and cannot be resumed by bridged.
|
||||
|
||||
The orchestrator is **`SessionManager` recursed one tier up**: today it spawns/names/reaps *worker*
|
||||
sessions; the orchestrator does the same for *main* sessions, and adds **context scoping**.
|
||||
|
||||
@@ -249,11 +273,10 @@ flowchart TB
|
||||
m2 -->|"scoped delegation"| w
|
||||
```
|
||||
|
||||
*Figure 7 — the recursion. Tiers 1 and 2 run the identical spawn/name/resume machinery; the
|
||||
orchestrator merely operates it one level higher. **Re-rooting caveat:** today the primary IS the
|
||||
human's live session; here the human drives the orchestrator, and the mains become managed,
|
||||
resumable sessions. That moves the human-facing top up a tier — an intentional re-root, not an
|
||||
add-on.*
|
||||
*Figure 7 — **SUPERSEDED historical alternative.** The recursion re-rooted the human-facing session:
|
||||
the human drove an orchestrator and the mains became managed, resumable sessions. The operator rejected
|
||||
that re-root. The replacement keeps the human-driven, pre-existing lead and engages architects sideways
|
||||
as short-lived advisory peers; see Gitea issue #16 and `architects:`.*
|
||||
|
||||
```mermaid
|
||||
sequenceDiagram
|
||||
@@ -270,10 +293,9 @@ sequenceDiagram
|
||||
O->>O: fold into orchestrator context, pick next main/turn
|
||||
```
|
||||
|
||||
*Figure 8 — context focus. The orchestrator holds the global context and hands each main only the
|
||||
slice a given delegation needs, so the main→worker conversation stays on-point. Context *scoping* is
|
||||
coordination (the bus already owns session/turn lifecycle) — it stays inside the identity boundary
|
||||
(§7), unlike toolchain ownership which does not.*
|
||||
*Figure 8 — **SUPERSEDED historical alternative.** This proposed an orchestrator holding global context
|
||||
and slicing it for managed mains. The replacement has the human-driven lead send the same advisory brief
|
||||
issue #16 and `architects:`.*
|
||||
|
||||
**Deltas:** a second, higher `SessionManager` instance whose "peers" are mains; the orchestrator
|
||||
becomes the top MCP client; context-slice selection (new) layered on CB-303's `context_cap` +
|
||||
@@ -315,7 +337,7 @@ flowchart LR
|
||||
cb402["CB-401/402<br/>Peer Launcher SPI + composite<br/>(DONE / in-flight)"]
|
||||
A["A · SandboxLauncher<br/>(placement-neutral, independent)"]
|
||||
cb308["CB-308 substrate<br/>per-agent channels + global id<br/>+ federated roster"]
|
||||
B["B · main-agent pair<br/>(multi-slot PrimaryRegistry)"]
|
||||
B["B · main-agent pair (SUPERSEDED)<br/>(multi-slot PrimaryRegistry)"]
|
||||
C["C · orchestrator tier<br/>(SessionManager recursed up)"]
|
||||
cb402 --> A
|
||||
cb402 --> cb308
|
||||
@@ -333,7 +355,7 @@ flowchart LR
|
||||
2. **A · SandboxLauncher** — independent; a second proof of the SPI (placement-neutral). Ships anytime.
|
||||
3. **CB-308 substrate** — per-agent channels + global id + federated roster (the multi-host work,
|
||||
promoted from host-to-host to tier-to-tier).
|
||||
4. **B · main-agent pair** — multi-slot `PrimaryRegistry` + per-main inbox routing, on the substrate.
|
||||
4. **B · main-agent pair** — **SUPERSEDED** by lead + two advisory architects.
|
||||
5. **C · orchestrator tier** — the capstone; the recursive session manager + context scoping.
|
||||
|
||||
## 9. Open questions (to resolve at ticket-split)
|
||||
@@ -341,8 +363,8 @@ flowchart LR
|
||||
- **Sandbox mechanism:** container (`docker exec`) vs devcontainer — how the role→image mapping is
|
||||
expressed on the profile. *(Topology **resolved** in §11: distributed = gateway-per-host × local
|
||||
sandboxes; the remaining choice is only the local launch mechanism, not the shape.)*
|
||||
- **Pair semantics:** are the two mains fully symmetric peers, or is one a co-primary that may also
|
||||
delegate? Affects how `PrimaryRegistry` and the "orchestration tools" identity relax.
|
||||
- **Pair semantics:** **SUPERSEDED.** The two-main question is replaced by the architect role's
|
||||
least-privilege boundary: advisory architects can send/reply/ask/read but cannot spawn/stop/drain.
|
||||
- **Orchestrator drivenness:** the mains become programmatically spawned/resumed — does the human
|
||||
still ever type directly into a main, or only into the orchestrator? (The re-root caveat, Fig 7.)
|
||||
- **Context-slice selection:** who decides the slice — orchestrator heuristics, explicit tool args,
|
||||
|
||||
+1
-1
Submodule wiki updated: 8c63db5da6...e5424f4665
Reference in New Issue
Block a user