Compare commits
12 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| cec48832be | |||
| 65c9deb4d1 | |||
| 28ae27b8e1 | |||
| 8a837a2830 | |||
| 0a2b3a4a56 | |||
| 613ece92dc | |||
| 8d4206c2b5 | |||
| 03286a589b | |||
| 88b9503c3b | |||
| c553d795d8 | |||
| 3ba6d6784c | |||
| 4fa6553db5 |
@@ -88,15 +88,28 @@ bind:
|
||||
# backoffMs: 60000
|
||||
# quietNudgeCap: 3
|
||||
|
||||
# Fleet health detection is dormant unless enabled. It reads one whole-fleet agent list per tick.
|
||||
# It can run without a webhook; bridge_list then reports healthCoverage: detection-only.
|
||||
# Fleet health detection is dormant unless enabled (CB-573). It reads one whole-fleet agent list
|
||||
# per tick.
|
||||
# intervalSeconds → how often a tick runs (default 30). ENFORCED floor of 15: the code computes
|
||||
# Math.max(15, intervalSeconds), so a lower value is silently raised, not
|
||||
# rejected.
|
||||
# workingSuspectAfterSeconds, paneProbeIntervalSeconds → accepted and parsed, but NOT YET READ by
|
||||
# anything — the dormant monitor only consumes intervalSeconds today (CB-573
|
||||
# shipped ahead of the evidence publishers these two knobs are for). Setting
|
||||
# them changes nothing right now, and no minimum is enforced on either, because
|
||||
# nothing reads them to enforce one. They exist so a later build can start
|
||||
# honouring them without another config-shape change.
|
||||
# notifications.mode → "webhook" flips what bridge_list REPORTS (healthCoverage: "full" instead
|
||||
# of "detection-only") — it does NOT make bridged send any webhook call; no
|
||||
# delivery mechanism is implemented yet. Any other value, or omitting the
|
||||
# block, reports "detection-only".
|
||||
# health:
|
||||
# enabled: true
|
||||
# intervalSeconds: 30 # minimum 15
|
||||
# workingSuspectAfterSeconds: 600 # minimum 300
|
||||
# paneProbeIntervalSeconds: 60 # minimum 60
|
||||
# intervalSeconds: 30
|
||||
# workingSuspectAfterSeconds: 600
|
||||
# paneProbeIntervalSeconds: 60
|
||||
# notifications:
|
||||
# mode: disabled # disabled (default) or webhook
|
||||
# mode: disabled
|
||||
|
||||
# herdr Unix socket. Omit to use the client default
|
||||
# (${HERDR_SOCKET_PATH:-~/.config/herdr/herdr.sock}).
|
||||
@@ -286,6 +299,11 @@ placement: weighted
|
||||
# / credentialId. Those are hot because the placement policy (and, for credentialId,
|
||||
# the CB-578 stage B quarantine check) reads them through a supplier — being config is
|
||||
# not by itself enough to make a key hot.
|
||||
# EXCEPT `fleet.leaders`: Bridged.main reads it once at startup to build the lead tab
|
||||
# scanner and launcher, and neither is rebuilt on reload. A changed/added/removed
|
||||
# `fleet.leaders` entry is silently accepted — the reload reports "config reloaded"
|
||||
# with nothing in the deferred list — but has NO effect until you restart. Treat it
|
||||
# as deferred in practice, even though today's reload output does not say so.
|
||||
# DEFERRED → accepted into the new config, but the wiring built at startup keeps the old value
|
||||
# until you restart: `lifecycle:`, `leadHeartbeat:`, `guard:`, `worktreeRoot:`,
|
||||
# `spawnReadyTimeoutMs` / `spawnReadyPollMs`, `quarantineCooldownSeconds` (CB-578
|
||||
@@ -368,6 +386,12 @@ fleet:
|
||||
# An auto-launched lead is NOT a member: it gets no worker reply charter, is never registered with
|
||||
# the session lifecycle (the idle reaper would kill your orchestrator), and stays on the
|
||||
# subscription — ANTHROPIC_BASE_URL/AUTH_TOKEN are stripped from its env whatever the profile says.
|
||||
#
|
||||
# GET THE `tab:` VALUE RIGHT. A pane that does not match any configured `tab:` (a typo, a renamed
|
||||
# tab, a pane no entry names at all) is not recognised as a lead — it resolves as an ordinary
|
||||
# WORKER instead, silently, and every orchestration call it makes (spawn/stop/send/drain) is
|
||||
# refused. There is no error at startup for this: an unmatched pane is simply not a lead. If your
|
||||
# primary suddenly can't spawn or send, check this section first.
|
||||
# leaders:
|
||||
# opus-5.0:
|
||||
# profile: opus # omit to never create this lead, only recognise it
|
||||
@@ -423,10 +447,15 @@ guard:
|
||||
# idleTtlSeconds → reap READY/DONE sessions idle longer than this (never BUSY/SPAWNING)
|
||||
# contextCap → force-release a session after this many delegated turns
|
||||
# drainTimeoutSeconds → seconds to wait for BUSY sessions on shutdown before forced teardown
|
||||
# clearAfterTurn → whether a reusable worker discards its conversation context after every
|
||||
# completed delegated turn (default false). Works for claude-code workers
|
||||
# only — any other peer kind (e.g. opencode) logs "context reset is
|
||||
# unsupported for peer kind …" once and the reset is a no-op.
|
||||
# lifecycle:
|
||||
# idleTtlSeconds: 300
|
||||
# contextCap: 10
|
||||
# drainTimeoutSeconds: 5
|
||||
# clearAfterTurn: false
|
||||
|
||||
# Durable reply delivery (CB-307 Stage 2). OMIT this block entirely to keep the default
|
||||
# in-memory, soft-state reply inbox (late worker replies are held only until a daemon bounce).
|
||||
@@ -434,10 +463,14 @@ guard:
|
||||
# on a durable per-target queue (agent.<target>.inbox) and survive a restart — the broker
|
||||
# redelivers anything the primary had not yet drained. Production default is LavinMQ; a stock
|
||||
# RabbitMQ speaks the same AMQP 0-9-1, so it is a URI-only swap.
|
||||
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
|
||||
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
|
||||
# uri → AMQP connection URI. No trailing slash ⇒ the default vhost "/"; an empty path ("/")
|
||||
# is vhost "" and will NOT connect. Encode a named vhost as .../%2Fmyvhost.
|
||||
# prefetch → CB-527: consumer basicQos, capping how many unacked messages the inbox holds
|
||||
# in-heap per owned target (the rest sits on the broker's durable queue instead of
|
||||
# growing the JVM heap). Default 32 when omitted.
|
||||
# broker:
|
||||
# uri: amqp://guest:guest@127.0.0.1:5672
|
||||
# prefetch: 32
|
||||
|
||||
# Active push-to-primary (CB-307 Stage 3). When a worker reply lands with no open bridge_send,
|
||||
# the ReplyPushLoop injects a *drain nudge* (never the payload) into the primary's own herdr
|
||||
|
||||
@@ -83,6 +83,10 @@ public final class Bridged {
|
||||
static void main(String[] args) {
|
||||
Path configPath = Path.of(args.length > 0 ? args[0] : "bridged.yaml");
|
||||
BridgedConfig cfg = BridgedConfig.load(configPath);
|
||||
// CB-594: report which secret env vars the config actually needs, by name, before anything
|
||||
// else can fail on a silently-empty one. A daemon started without a login shell (launchd)
|
||||
// boots fine either way — this is the only thing that says so out loud.
|
||||
reportRequiredSecrets(cfg);
|
||||
// CB-559: `cfg` stays the startup snapshot — every validation and every piece of one-time
|
||||
// wiring below reads it, and must, because those decisions cannot be unmade. `config` is the
|
||||
// live reference the hot paths read per use. Which keys can actually move is ConfigRef's
|
||||
@@ -348,8 +352,9 @@ public final class Bridged {
|
||||
// connection, so keep the reference to close it in the ordered shutdown hook.
|
||||
final ReplyInbox replyInbox;
|
||||
if (cfg.broker() != null && cfg.broker().isConfigured()) {
|
||||
replyInbox = AmqpReplyInbox.open(cfg.broker().uri());
|
||||
log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri());
|
||||
replyInbox = AmqpReplyInbox.open(cfg.broker().uri(), cfg.broker().prefetchOrDefault());
|
||||
log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})",
|
||||
cfg.broker().uri(), cfg.broker().prefetchOrDefault());
|
||||
} else {
|
||||
replyInbox = new InMemoryReplyInbox();
|
||||
log.info("reply inbox: in-memory (soft-state)");
|
||||
@@ -543,6 +548,66 @@ public final class Bridged {
|
||||
return target -> presence.isPresent(target) || leads.get().containsKey(target);
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-594: which env vars the loaded config actually needs, and why — every non-{@code
|
||||
* subscription} profile's {@code tokenEnv} (a subscription profile never reads one, see
|
||||
* {@link BridgedConfig.Profile#isSubscription()}), plus every profile's {@code gitTokenEnv}
|
||||
* where set (opt-in). Derived from the config, not hard-coded, so a new profile is covered for
|
||||
* free. A var required by more than one profile is one entry naming every profile that needs
|
||||
* it. Deliberately excludes {@code auth.tokenEnv}: that one is already enforced loudly, by a
|
||||
* startup throw in {@code main()} — about 370 lines <em>below</em> this method's call site
|
||||
* ({@link #reportRequiredSecrets(BridgedConfig)}), not a few lines above it. That throw only
|
||||
* fires when {@code auth.mode: token} is configured; under the default loopback-trust mode it
|
||||
* never runs, and {@code auth.tokenEnv} is simply not required.
|
||||
*
|
||||
* <p>Package-private and pure (no I/O, no logging) so the derivation is unit-testable without
|
||||
* capturing log output; {@link #reportRequiredSecrets(BridgedConfig)} is the logging caller.
|
||||
*/
|
||||
static Map<String, List<String>> requiredSecretEnvVars(BridgedConfig cfg) {
|
||||
Map<String, List<String>> requiredBy = new LinkedHashMap<>();
|
||||
cfg.profiles().forEach((name, profile) -> {
|
||||
if (!profile.isSubscription()) {
|
||||
requiredBy.computeIfAbsent(profile.tokenEnv(), _ -> new ArrayList<>())
|
||||
.add("profile '" + name + "' tokenEnv");
|
||||
}
|
||||
if (profile.hasGitToken()) {
|
||||
requiredBy.computeIfAbsent(profile.gitTokenEnv(), _ -> new ArrayList<>())
|
||||
.add("profile '" + name + "' gitTokenEnv");
|
||||
}
|
||||
});
|
||||
return requiredBy;
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-594: log, by name only, which required env vars (see {@link #requiredSecretEnvVars}) are
|
||||
* set in the daemon's own process environment — the environment every profile's {@code
|
||||
* tokenEnv}/{@code gitTokenEnv} is read from at spawn time (see
|
||||
* {@code HerdrPeerLauncher.resolveEnv}). Never logs a value, a prefix, or a length.
|
||||
*
|
||||
* <p>A missing entry only warns — it must never refuse to start. A daemon that boots and says
|
||||
* what is wrong is strictly more useful than one that will not boot at all.
|
||||
*/
|
||||
private static void reportRequiredSecrets(BridgedConfig cfg) {
|
||||
Map<String, List<String>> requiredBy = requiredSecretEnvVars(cfg);
|
||||
if (requiredBy.isEmpty()) {
|
||||
log.info("startup secrets: no profile references a token env var — nothing to check");
|
||||
return;
|
||||
}
|
||||
Map<String, String> env = System.getenv();
|
||||
requiredBy.forEach((varName, sources) -> {
|
||||
String value = env.get(varName);
|
||||
if (value != null && !value.isBlank()) {
|
||||
log.info("startup secret {}: set ({})", varName, String.join(", ", sources));
|
||||
} else {
|
||||
log.warn("startup secret {}: MISSING ({}) — the daemon will start anyway, and this "
|
||||
+ "failure stays invisible until a worker actually needs it. Fix "
|
||||
+ "${SHARED_ENV}/tools/secrets.sh and restart bridged from a LOGIN "
|
||||
+ "shell (see scripts/redeploy-bridged.sh).",
|
||||
varName, String.join(", ", sources));
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
/**
|
||||
* Poll herdr's {@code ping} until it answers or {@link #HERDR_WAIT_SECONDS} elapses (CB-504).
|
||||
*
|
||||
|
||||
@@ -5,6 +5,7 @@ import com.fasterxml.jackson.core.JsonParser;
|
||||
import com.fasterxml.jackson.core.JsonToken;
|
||||
import com.fasterxml.jackson.databind.ObjectMapper;
|
||||
import com.fasterxml.jackson.dataformat.yaml.YAMLFactory;
|
||||
import dev.ltms.bridged.msg.AmqpReplyInbox;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
@@ -532,16 +533,24 @@ public record BridgedConfig(
|
||||
* stays soft-state. Production default is LavinMQ; a stock RabbitMQ speaks the same AMQP 0-9-1
|
||||
* and is a URI-only swap.
|
||||
*
|
||||
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/{@code null}
|
||||
* ⇒ the broker block is treated as absent (in-memory adapter).
|
||||
* @param uri AMQP connection URI, e.g. {@code amqp://guest:guest@127.0.0.1:5672/}. Blank/
|
||||
* {@code null} ⇒ the broker block is treated as absent (in-memory adapter).
|
||||
* @param prefetch CB-527: the consumer's {@code basicQos} prefetch count, bounding how many
|
||||
* unacked messages the AMQP inbox holds in-heap per owned target. {@code null}/
|
||||
* non-positive ⇒ {@link AmqpReplyInbox#DEFAULT_PREFETCH}.
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Broker(String uri) {
|
||||
public record Broker(String uri, Integer prefetch) {
|
||||
|
||||
/** True when a usable broker URI is configured (an empty block does not enable AMQP). */
|
||||
public boolean isConfigured() {
|
||||
return uri != null && !uri.isBlank();
|
||||
}
|
||||
|
||||
/** The prefetch to use, defaulting to {@link AmqpReplyInbox#DEFAULT_PREFETCH} when unset. */
|
||||
public int prefetchOrDefault() {
|
||||
return (prefetch != null && prefetch > 0) ? prefetch : AmqpReplyInbox.DEFAULT_PREFETCH;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -15,6 +15,7 @@ import dev.ltms.bridged.msg.MessageService;
|
||||
import dev.ltms.bridged.msg.Rendezvous;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
import dev.ltms.bridged.placement.BackendQuarantine;
|
||||
import dev.ltms.bridged.placement.PlacementException;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.session.MemberSession;
|
||||
import dev.ltms.bridged.session.WorktreeRequest;
|
||||
@@ -692,6 +693,10 @@ public final class BridgeMcp {
|
||||
return text(json(memberView(member)));
|
||||
} catch (GuardException e) {
|
||||
return error("subscription boundary: " + e.getMessage());
|
||||
} catch (PlacementException e) {
|
||||
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — distinct
|
||||
// from "profile does not exist" below.
|
||||
return error("no capacity: " + e.getMessage());
|
||||
} catch (IllegalArgumentException e) {
|
||||
return error(e.getMessage()); // unknown / no-default profile, or a refused resumeSessionId
|
||||
} catch (PeerUnreachableException e) {
|
||||
|
||||
@@ -7,6 +7,7 @@ import com.rabbitmq.client.ConnectionFactory;
|
||||
import com.rabbitmq.client.DeliverCallback;
|
||||
import com.rabbitmq.client.Recoverable;
|
||||
import com.rabbitmq.client.RecoveryListener;
|
||||
import com.rabbitmq.client.Return;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
@@ -14,7 +15,13 @@ import java.io.IOException;
|
||||
import java.nio.charset.StandardCharsets;
|
||||
import java.util.LinkedHashMap;
|
||||
import java.util.List;
|
||||
import java.util.NavigableMap;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ConcurrentSkipListMap;
|
||||
import java.util.concurrent.ExecutionException;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
|
||||
/**
|
||||
* AMQP-backed {@link ReplyInbox} (CB-307 Stage 2): genuine cross-restart durability behind the same
|
||||
@@ -29,6 +36,27 @@ import java.util.concurrent.ConcurrentHashMap;
|
||||
* leaves them on the broker — it redelivers on reconnect. That is the durability the in-memory
|
||||
* adapter cannot give, with the port contract preserved.
|
||||
*
|
||||
* <p><strong>Prefetch bounds the held backlog (CB-527).</strong> The consumer channel calls
|
||||
* {@code basicQos} with a configurable prefetch count ({@link #DEFAULT_PREFETCH} unless the caller
|
||||
* passes another value to {@link #open(String, int)}) before starting any consumer. Without a bound,
|
||||
* the broker pushes its entire queue into {@link #held} the instant a target is {@link #own owned},
|
||||
* so an undrained primary grows the JVM heap without limit and any queue-level control
|
||||
* ({@code x-max-length}, per-message TTL) never fires because the queue never actually holds a
|
||||
* backlog. Prefetch keeps the backlog where it is visible — on the broker — until the owner drains it.
|
||||
*
|
||||
* <p><strong>Publishes require a confirmed, routable delivery (CB-528).</strong> {@link #publish}
|
||||
* runs on a channel separate from the consume/ack channel ({@link #channel}), so a slow or blocked
|
||||
* publish confirm can never hold {@link #channelLock} and stall an ack — the ack path never waits on
|
||||
* a publish confirm. That publish channel is in publisher-confirm mode and every publish sets the
|
||||
* {@code mandatory} flag, so an unroutable publish (queue not declared, e.g. the owner never called
|
||||
* {@link #own}) is returned by the broker instead of silently dropped. The broker sends the
|
||||
* <em>return</em> for an unroutable message before the <em>confirm</em> that covers it — the ack/nack
|
||||
* callback checks the returned-set at confirm time rather than assuming an ack means routed — so
|
||||
* "confirmed" here means "durably queued", not merely "accepted by the broker". A returned or nacked
|
||||
* (or un-confirmed within the timeout) publish surfaces as an {@link IllegalStateException} on the
|
||||
* caller's thread; the caller — {@link MessageService#reply} — must not report success for a
|
||||
* black-holed reply.
|
||||
*
|
||||
* <p><strong>Ownership is explicit.</strong> {@link #own} declares the queue and starts the consumer;
|
||||
* {@link #release} cancels it. {@link #publish} sends to the queue but does <em>not</em> imply ownership
|
||||
* and does not attach a consumer. This split is required by CB-308 federation, where one gateway may
|
||||
@@ -54,6 +82,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
private static final String QUEUE_PREFIX = "agent.";
|
||||
private static final String QUEUE_SUFFIX = ".inbox";
|
||||
|
||||
/** CB-527: the prefetch used when a caller does not pass an explicit value to {@link #open(String, int)}. */
|
||||
public static final int DEFAULT_PREFETCH = 32;
|
||||
|
||||
/** How long {@link #publish} waits for its publisher confirm before failing the call (CB-528). */
|
||||
private static final long CONFIRM_TIMEOUT_MS = 10_000L;
|
||||
|
||||
private final Connection connection;
|
||||
private final Channel channel;
|
||||
/** All channel operations (publish/declare/ack/cancel) serialize on this — a Channel is not thread-safe. */
|
||||
@@ -63,39 +97,87 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
|
||||
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
* CB-528: a dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume
|
||||
* + ack) so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
|
||||
*/
|
||||
private final Channel publishChannel;
|
||||
private final Object publishChannelLock = new Object();
|
||||
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
|
||||
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
|
||||
/**
|
||||
* The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery
|
||||
* tag. Assumes {@code msgId} is unique per in-flight publish: a second {@link #publish} for a
|
||||
* {@code msgId} still awaiting its confirm would overwrite this entry and misdirect
|
||||
* {@link #onReturn}'s lookup. Not reachable today — {@code MessageService.reply} generates a fresh
|
||||
* {@code UUID} per call — so no guard is added for it.
|
||||
*/
|
||||
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
|
||||
|
||||
/** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */
|
||||
private record Held(long deliveryTag, InboxMessage message) {}
|
||||
|
||||
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */
|
||||
/** A publish awaiting its confirm; {@link #returned} records whether the broker already returned it. */
|
||||
private static final class Pending {
|
||||
final String msgId;
|
||||
final CompletableFuture<Void> confirmed = new CompletableFuture<>();
|
||||
volatile boolean returned;
|
||||
|
||||
Pending(String msgId) {
|
||||
this.msgId = msgId;
|
||||
}
|
||||
}
|
||||
|
||||
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) with {@link #DEFAULT_PREFETCH}. */
|
||||
public static AmqpReplyInbox open(String uri) {
|
||||
return open(uri, DEFAULT_PREFETCH);
|
||||
}
|
||||
|
||||
/** As {@link #open(String)}, with an explicit consumer prefetch (CB-527: caps the held backlog per target). */
|
||||
public static AmqpReplyInbox open(String uri, int prefetch) {
|
||||
try {
|
||||
ConnectionFactory factory = new ConnectionFactory();
|
||||
factory.setUri(uri);
|
||||
// Self-heal transient blips; topology recovery re-declares queues and re-attaches consumers.
|
||||
factory.setAutomaticRecoveryEnabled(true);
|
||||
factory.setTopologyRecoveryEnabled(true);
|
||||
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"));
|
||||
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"), prefetch);
|
||||
} catch (Exception e) {
|
||||
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** Wrap an already-open connection (injection seam for the contract test). */
|
||||
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
|
||||
AmqpReplyInbox(Connection connection) {
|
||||
this(connection, DEFAULT_PREFETCH);
|
||||
}
|
||||
|
||||
/** As above, with an explicit prefetch (injection seam for the contract test). */
|
||||
AmqpReplyInbox(Connection connection, int prefetch) {
|
||||
this.connection = connection;
|
||||
try {
|
||||
this.channel = connection.createChannel();
|
||||
// CB-527: bound the held backlog per owned target — must be set before any own()/basicConsume.
|
||||
this.channel.basicQos(prefetch);
|
||||
this.publishChannel = connection.createChannel();
|
||||
this.publishChannel.confirmSelect();
|
||||
this.publishChannel.addReturnListener(this::onReturn);
|
||||
this.publishChannel.addConfirmListener(this::onAck, this::onNack);
|
||||
} catch (IOException e) {
|
||||
throw new IllegalStateException("cannot open AMQP channel", e);
|
||||
}
|
||||
// On automatic recovery the broker redelivers unacked messages with FRESH delivery-tags; the
|
||||
// tags we were holding are now stale. Drop the held snapshot so the re-attached consumer
|
||||
// repopulates it with valid tags (dedup by msgId still prevents any double-queue).
|
||||
// repopulates it with valid tags (dedup by msgId still prevents any double-queue). Any publish
|
||||
// confirm still in flight when the connection dropped is equally stale — its sequence number
|
||||
// meant nothing on the old channel and means nothing on the recovered one, so fail it now
|
||||
// rather than let it silently ride out CONFIRM_TIMEOUT_MS.
|
||||
if (connection instanceof Recoverable recoverable) {
|
||||
recoverable.addRecoveryListener(new RecoveryListener() {
|
||||
@Override
|
||||
public void handleRecovery(Recoverable recoverable) {
|
||||
held.clear();
|
||||
failPendingPublishesOnRecovery();
|
||||
log.info("AMQP connection recovered; cleared held replies for fresh redelivery");
|
||||
}
|
||||
|
||||
@@ -141,6 +223,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Publish {@code content} and block until the broker's publisher confirm for it lands (CB-528).
|
||||
* Throws {@link IllegalStateException} if the message is returned as unroutable, nacked, or not
|
||||
* confirmed within {@link #CONFIRM_TIMEOUT_MS} — the caller must treat that as a failed publish,
|
||||
* not a lost-and-forgotten one.
|
||||
*/
|
||||
@Override
|
||||
public void publish(String target, String msgId, String content) {
|
||||
AMQP.BasicProperties props = new AMQP.BasicProperties.Builder()
|
||||
@@ -148,12 +236,35 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
.deliveryMode(2) // persistent — survives a broker restart
|
||||
.contentType("text/plain")
|
||||
.build();
|
||||
try {
|
||||
synchronized (channelLock) {
|
||||
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
|
||||
Pending pending = new Pending(msgId);
|
||||
long seq;
|
||||
synchronized (publishChannelLock) {
|
||||
seq = publishChannel.getNextPublishSeqNo();
|
||||
pendingBySeq.put(seq, pending);
|
||||
pendingByMsgId.put(msgId, pending);
|
||||
try {
|
||||
publishChannel.basicPublish("", queueName(target), true, props,
|
||||
content.getBytes(StandardCharsets.UTF_8));
|
||||
} catch (IOException e) {
|
||||
pendingBySeq.remove(seq, pending);
|
||||
pendingByMsgId.remove(msgId, pending);
|
||||
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
|
||||
}
|
||||
} catch (IOException e) {
|
||||
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
|
||||
}
|
||||
try {
|
||||
pending.confirmed.get(CONFIRM_TIMEOUT_MS, TimeUnit.MILLISECONDS);
|
||||
} catch (ExecutionException e) {
|
||||
Throwable cause = e.getCause();
|
||||
throw cause instanceof RuntimeException re ? re : new IllegalStateException(cause);
|
||||
} catch (TimeoutException e) {
|
||||
throw new IllegalStateException("publish confirm for reply " + msgId + " to " + queueName(target)
|
||||
+ " timed out after " + CONFIRM_TIMEOUT_MS + "ms — broker may be unreachable or overloaded", e);
|
||||
} catch (InterruptedException e) {
|
||||
Thread.currentThread().interrupt();
|
||||
throw new IllegalStateException("interrupted awaiting publish confirm for " + msgId, e);
|
||||
} finally {
|
||||
pendingBySeq.remove(seq, pending);
|
||||
pendingByMsgId.remove(msgId, pending);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -222,17 +333,115 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
};
|
||||
}
|
||||
|
||||
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm (CB-528). */
|
||||
private void onReturn(Return r) {
|
||||
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
|
||||
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
|
||||
if (pending != null) {
|
||||
pending.returned = true;
|
||||
} else {
|
||||
log.warn("AMQP return for reply {} (routingKey={}, {} {}) with no matching in-flight publish"
|
||||
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
|
||||
r.getReplyText());
|
||||
}
|
||||
}
|
||||
|
||||
private void onAck(long seq, boolean multiple) {
|
||||
resolveConfirm(seq, multiple, true);
|
||||
}
|
||||
|
||||
private void onNack(long seq, boolean multiple) {
|
||||
resolveConfirm(seq, multiple, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
|
||||
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
|
||||
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
|
||||
* a return means "confirmed but never routed", not "durably queued".
|
||||
*/
|
||||
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
|
||||
NavigableMap<Long, Pending> covered = multiple
|
||||
? pendingBySeq.headMap(seq, true)
|
||||
: pendingBySeq.subMap(seq, true, seq, true);
|
||||
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
if (ack && !pending.returned) {
|
||||
pending.confirmed.complete(null);
|
||||
} else if (ack) {
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"reply " + pending.msgId + " was returned as unroutable (queue not declared/owned)"));
|
||||
} else {
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"broker nacked publish of reply " + pending.msgId));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery.
|
||||
* Guarded by {@link #publishChannelLock}, the same lock {@link #publish} holds while it takes its
|
||||
* sequence number and registers its {@link Pending}: without it, a {@link #publish} that starts
|
||||
* after the connection has already recovered (so it publishes — and will be confirmed — on the
|
||||
* <em>new</em> channel) can register between this sweep's iteration and its clear, and this sweep
|
||||
* then fails a publish that actually succeeded. {@link #publish} only holds the lock for the
|
||||
* seq/map-put/{@code basicPublish} — it awaits the confirm outside it — so this sweep can only ever
|
||||
* wait for an in-flight {@code basicPublish} call to return, never for a broker round trip. No
|
||||
* deadlock.
|
||||
*
|
||||
* <p>Package-private (rather than {@code private}) only so the unit test can drive it directly
|
||||
* against a concurrent {@link #publish} without a live broker reconnect.
|
||||
*/
|
||||
void failPendingPublishesOnRecovery() {
|
||||
synchronized (publishChannelLock) {
|
||||
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId
|
||||
+ " is unknown"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Fail every publish still awaiting its confirm with a clear, immediate error instead of leaving it
|
||||
* to time out after {@link #CONFIRM_TIMEOUT_MS} once the channels are closed underneath it. Guarded
|
||||
* by {@link #publishChannelLock} for the same reason as {@link #failPendingPublishesOnRecovery}.
|
||||
*/
|
||||
private void failPendingPublishesOnClose() {
|
||||
synchronized (publishChannelLock) {
|
||||
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"AMQP reply inbox closed while publish of reply " + pending.msgId
|
||||
+ " was still awaiting its confirm"));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private static String queueName(String target) {
|
||||
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
|
||||
}
|
||||
|
||||
@Override
|
||||
public void close() {
|
||||
failPendingPublishesOnClose();
|
||||
try {
|
||||
channel.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP channel close: {}", e.toString());
|
||||
}
|
||||
try {
|
||||
publishChannel.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP publish channel close: {}", e.toString());
|
||||
}
|
||||
try {
|
||||
connection.close();
|
||||
} catch (Exception e) {
|
||||
|
||||
@@ -37,10 +37,12 @@ import java.util.stream.Collectors;
|
||||
* <p>Each tick examines <em>everything</em> pending for that lead — reply targets whose inbox
|
||||
* still holds an unacked message ({@link #pendingReplies}) and tickets not yet collected
|
||||
* ({@link #pendingTickets}) — and sends at most one combined nudge per tick
|
||||
* ({@link #injectNudge(String, int)}). Work that arrives while the lead is busy is never lost:
|
||||
* it is re-read fresh on every tick until the lead is injectable or the shared reminder cap
|
||||
* ({@link #maxReminders}) is reached, whichever the durable inbox / pending-ticket set doesn't
|
||||
* already answer via {@code STOP}.
|
||||
* ({@link #injectNudge(String, int, int)}). Work that arrives while the lead is busy is never
|
||||
* lost: it is re-read fresh on every tick until the lead is injectable or its own reminder cap
|
||||
* ({@link #maxReminders}) is reached — reply and ticket work each spend from their own budget, so
|
||||
* one source exhausting its cap does not stop nudges about the other (post-CB-590 regression fix;
|
||||
* see {@link #decide}) — whichever the durable inbox / pending-ticket set doesn't already answer
|
||||
* via {@code STOP}.
|
||||
*/
|
||||
public final class ReplyPushLoop {
|
||||
|
||||
@@ -144,20 +146,32 @@ public final class ReplyPushLoop {
|
||||
* Pure decision function: examine everything pending for {@code lead} — reply targets and
|
||||
* tickets alike — and return what the loop should do.
|
||||
*
|
||||
* @param lead the lead terminal to nudge
|
||||
* @param reminderCount how many nudges have been sent so far for this lead (shared across
|
||||
* both reply and ticket work — CB-590 collapses the reminder cap onto
|
||||
* one counter per lead, so alternating sources cannot outrun the bound)
|
||||
* <p><strong>CB-590-fix: one schedule, two budgets.</strong> The single per-lead schedule
|
||||
* (CB-590) still ticks once for both sources, but each source is capped independently —
|
||||
* {@code replyReminderCount} against a reply target still pending, {@code ticketReminderCount}
|
||||
* against a ticket still pending. A busy reply stream that exhausts its own cap must not stop
|
||||
* the loop from nudging about a ticket that still has budget left, and vice versa: either
|
||||
* source being eligible (has pending work AND is under its own cap) is enough for
|
||||
* {@link Action#INJECT}. Only when neither source has eligible work does the loop
|
||||
* {@link Action#STOP}.
|
||||
*
|
||||
* @param lead the lead terminal to nudge
|
||||
* @param replyReminderCount how many nudges have covered pending reply work for this lead
|
||||
* @param ticketReminderCount how many nudges have covered pending ticket work for this lead
|
||||
* @return the action the caller should take
|
||||
*/
|
||||
Action decide(String lead, int reminderCount) {
|
||||
boolean anyPending = !pendingReplyTargetsFor(lead).isEmpty() || !pendingTicketIdsFor(lead).isEmpty();
|
||||
if (!anyPending) {
|
||||
Action decide(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
boolean hasReplyWork = !pendingReplyTargetsFor(lead).isEmpty();
|
||||
boolean hasTicketWork = !pendingTicketIdsFor(lead).isEmpty();
|
||||
if (!hasReplyWork && !hasTicketWork) {
|
||||
log.debug("push: nothing pending for lead {}, stopping reminder", lead);
|
||||
return Action.STOP;
|
||||
}
|
||||
if (reminderCount >= maxReminders) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {}, stopping", maxReminders, lead);
|
||||
boolean replyEligible = hasReplyWork && replyReminderCount < maxReminders;
|
||||
boolean ticketEligible = hasTicketWork && ticketReminderCount < maxReminders;
|
||||
if (!replyEligible && !ticketEligible) {
|
||||
log.debug("push: reminder cap ({}) reached for lead {} on every source with pending work, stopping",
|
||||
maxReminders, lead);
|
||||
countNudge("exhausted");
|
||||
return Action.STOP;
|
||||
}
|
||||
@@ -241,21 +255,28 @@ public final class ReplyPushLoop {
|
||||
return;
|
||||
}
|
||||
log.debug("push: starting reminder loop for lead {}", lead);
|
||||
scheduleNext(lead, 0);
|
||||
scheduleNext(lead, 0, 0);
|
||||
}
|
||||
|
||||
/** Execute one loop tick — called on the scheduler thread. */
|
||||
private void tick(String lead, int reminderCount) {
|
||||
private void tick(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
Set<String> repliesBefore = pendingReplyTargetsFor(lead);
|
||||
Set<String> ticketsBefore = pendingTicketIdsFor(lead);
|
||||
var action = decide(lead, reminderCount);
|
||||
var action = decide(lead, replyReminderCount, ticketReminderCount);
|
||||
switch (action) {
|
||||
case INJECT -> {
|
||||
injectNudge(lead, reminderCount);
|
||||
scheduleNext(lead, reminderCount + 1);
|
||||
injectNudge(lead, replyReminderCount, ticketReminderCount);
|
||||
// Only the source(s) actually eligible this tick spend a unit of their own budget —
|
||||
// an exhausted source riding along in the combined message (still pending, still
|
||||
// named) does not get charged again; its count stays put until it drains.
|
||||
boolean replyEligible = !repliesBefore.isEmpty() && replyReminderCount < maxReminders;
|
||||
boolean ticketEligible = !ticketsBefore.isEmpty() && ticketReminderCount < maxReminders;
|
||||
scheduleNext(lead,
|
||||
replyEligible ? replyReminderCount + 1 : replyReminderCount,
|
||||
ticketEligible ? ticketReminderCount + 1 : ticketReminderCount);
|
||||
}
|
||||
// Re-check after the configured backoff; the lead may become injectable soon.
|
||||
case WAIT_BUSY -> scheduleNext(lead, reminderCount);
|
||||
case WAIT_BUSY -> scheduleNext(lead, replyReminderCount, ticketReminderCount);
|
||||
case STOP -> stopOrRestart(lead, repliesBefore, ticketsBefore);
|
||||
}
|
||||
}
|
||||
@@ -296,14 +317,14 @@ public final class ReplyPushLoop {
|
||||
|| pendingTicketIdsFor(lead).stream().anyMatch(t -> !ticketsBefore.contains(t));
|
||||
if (racedIn && activeLeads.putIfAbsent(lead, Boolean.TRUE) == null) {
|
||||
log.debug("push: new work for lead {} raced the reminder loop's stop — restarting", lead);
|
||||
scheduleNext(lead, 0);
|
||||
scheduleNext(lead, 0, 0);
|
||||
return;
|
||||
}
|
||||
log.debug("push: reminder loop ended for lead {}", lead);
|
||||
}
|
||||
|
||||
/** Send one combined nudge covering everything currently pending for {@code lead}. */
|
||||
private void injectNudge(String lead, int reminderCount) {
|
||||
private void injectNudge(String lead, int replyReminderCount, int ticketReminderCount) {
|
||||
// Re-read rather than threading it down from decide(): a reply can drain, or a ticket be
|
||||
// collected (or another arrive), between the decision and the injection.
|
||||
Set<String> replyTargets = pendingReplyTargetsFor(lead);
|
||||
@@ -315,18 +336,20 @@ public final class ReplyPushLoop {
|
||||
String nudge = formatNudge(replyTargets, tickets);
|
||||
try {
|
||||
agents.send(lead, nudge);
|
||||
log.debug("push: nudge {}/{} sent to lead {} ({} reply target(s), {} ticket(s))",
|
||||
reminderCount + 1, maxReminders, lead, replyTargets.size(), tickets.size());
|
||||
log.debug("push: nudge sent to lead {} (reply {}/{}, ticket {}/{}; {} reply target(s), {} ticket(s))",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders,
|
||||
replyTargets.size(), tickets.size());
|
||||
countNudge("delivered");
|
||||
} catch (RuntimeException e) {
|
||||
log.warn("push: failed to nudge lead {} (reminder {}/{}): {}",
|
||||
lead, reminderCount + 1, maxReminders, e.toString());
|
||||
log.warn("push: failed to nudge lead {} (reply {}/{}, ticket {}/{}): {}",
|
||||
lead, replyReminderCount + 1, maxReminders, ticketReminderCount + 1, maxReminders, e.toString());
|
||||
}
|
||||
}
|
||||
|
||||
/** Schedule the next tick on the scheduler thread pool. */
|
||||
private void scheduleNext(String lead, int nextReminderCount) {
|
||||
scheduler.schedule(() -> tick(lead, nextReminderCount), backoffMs, TimeUnit.MILLISECONDS);
|
||||
private void scheduleNext(String lead, int nextReplyReminderCount, int nextTicketReminderCount) {
|
||||
scheduler.schedule(() -> tick(lead, nextReplyReminderCount, nextTicketReminderCount),
|
||||
backoffMs, TimeUnit.MILLISECONDS);
|
||||
}
|
||||
|
||||
// --- nudge formatting ------------------------------------------------------------------------
|
||||
|
||||
@@ -13,6 +13,7 @@ import dev.ltms.bridged.herdr.HerdrClient;
|
||||
import dev.ltms.bridged.herdr.HerdrException;
|
||||
import dev.ltms.bridged.inject.MemberPresence;
|
||||
import dev.ltms.bridged.peer.PeerUnreachableException;
|
||||
import dev.ltms.bridged.placement.PlacementException;
|
||||
import dev.ltms.bridged.msg.MessageService;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.peer.MemberRole;
|
||||
@@ -292,6 +293,11 @@ public final class BridgedApp {
|
||||
ctx.status(201).json(view(member));
|
||||
} catch (GuardException e) {
|
||||
ctx.status(403).json(Map.of("error", "subscription_boundary", "detail", e.getMessage()));
|
||||
} catch (PlacementException e) {
|
||||
// CB-599: no candidate had capacity (maxLoad, quarantine, or all-exhausted) — a benign,
|
||||
// likely-transient refusal, distinct from "profile does not exist" below. 503: the
|
||||
// request was valid and will likely succeed later.
|
||||
ctx.status(503).json(Map.of("error", "no_capacity", "detail", e.getMessage()));
|
||||
} catch (IllegalArgumentException e) {
|
||||
ctx.status(400).json(Map.of("error", "unknown_profile", "detail", e.getMessage()));
|
||||
} catch (PeerUnreachableException e) {
|
||||
|
||||
@@ -0,0 +1,114 @@
|
||||
package dev.ltms.bridged;
|
||||
|
||||
import dev.ltms.bridged.config.BridgedConfig;
|
||||
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.util.List;
|
||||
import java.util.Map;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertFalse;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-594: {@link Bridged#requiredSecretEnvVars(BridgedConfig)} is what decides what the startup
|
||||
* secret report checks — it must derive that set from the config, not a hand-written list, or a
|
||||
* new profile's token silently stops being reported.
|
||||
*/
|
||||
class RequiredSecretEnvVarsTest {
|
||||
|
||||
private static BridgedConfig load(Path dir, String yaml) throws Exception {
|
||||
Path f = dir.resolve("bridged.yaml");
|
||||
Files.writeString(f, yaml);
|
||||
return BridgedConfig.load(f);
|
||||
}
|
||||
|
||||
@Test
|
||||
void collectsATokenEnvPerNonSubscriptionProfile(@TempDir Path dir) throws Exception {
|
||||
BridgedConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
tokenEnv: AI_GATEWAY_TOKEN
|
||||
""");
|
||||
|
||||
Map<String, List<String>> required = Bridged.requiredSecretEnvVars(cfg);
|
||||
|
||||
assertTrue(required.containsKey("AI_GATEWAY_TOKEN"));
|
||||
assertEquals(List.of("profile 'local' tokenEnv"), required.get("AI_GATEWAY_TOKEN"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aSubscriptionProfileNeedsNoTokenEnv(@TempDir Path dir) throws Exception {
|
||||
BridgedConfig cfg = load(dir, """
|
||||
profiles:
|
||||
opus:
|
||||
subscription: true
|
||||
model: claude-opus-5
|
||||
""");
|
||||
|
||||
assertTrue(Bridged.requiredSecretEnvVars(cfg).isEmpty(),
|
||||
"subscription: true never reads ANTHROPIC_AUTH_TOKEN — see Profile#isSubscription");
|
||||
}
|
||||
|
||||
@Test
|
||||
void gitTokenEnvIsOptInAndCollectedWhenSet(@TempDir Path dir) throws Exception {
|
||||
BridgedConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
tokenEnv: AI_GATEWAY_TOKEN
|
||||
gitTokenEnv: WORKER_GITEA_TOKEN
|
||||
""");
|
||||
|
||||
Map<String, List<String>> required = Bridged.requiredSecretEnvVars(cfg);
|
||||
|
||||
assertTrue(required.containsKey("WORKER_GITEA_TOKEN"));
|
||||
assertEquals(List.of("profile 'local' gitTokenEnv"), required.get("WORKER_GITEA_TOKEN"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void noGitTokenEnvMeansNothingIsRequiredForIt(@TempDir Path dir) throws Exception {
|
||||
BridgedConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
tokenEnv: AI_GATEWAY_TOKEN
|
||||
""");
|
||||
|
||||
assertFalse(Bridged.requiredSecretEnvVars(cfg).containsKey("WORKER_GITEA_TOKEN"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void aVarSharedByTwoProfilesIsReportedOnceNamingBoth(@TempDir Path dir) throws Exception {
|
||||
BridgedConfig cfg = load(dir, """
|
||||
profiles:
|
||||
local:
|
||||
baseUrl: http://gx00.gw:8000
|
||||
tokenEnv: AI_GATEWAY_TOKEN
|
||||
gitTokenEnv: WORKER_GITEA_TOKEN
|
||||
gx:
|
||||
kind: opencode
|
||||
baseUrl: https://llm.ltms.dev/v1
|
||||
tokenEnv: AI_GATEWAY_TOKEN
|
||||
gitTokenEnv: WORKER_GITEA_TOKEN
|
||||
""");
|
||||
|
||||
Map<String, List<String>> required = Bridged.requiredSecretEnvVars(cfg);
|
||||
|
||||
assertEquals(List.of("profile 'local' tokenEnv", "profile 'gx' tokenEnv"),
|
||||
required.get("AI_GATEWAY_TOKEN"));
|
||||
assertEquals(List.of("profile 'local' gitTokenEnv", "profile 'gx' gitTokenEnv"),
|
||||
required.get("WORKER_GITEA_TOKEN"));
|
||||
}
|
||||
|
||||
@Test
|
||||
void noProfilesMeansNothingIsRequired(@TempDir Path dir) throws Exception {
|
||||
BridgedConfig cfg = load(dir, "bind:\n host: 127.0.0.1\n port: 8765\n");
|
||||
|
||||
assertTrue(Bridged.requiredSecretEnvVars(cfg).isEmpty());
|
||||
}
|
||||
}
|
||||
@@ -16,12 +16,15 @@ import dev.ltms.bridged.inject.MemberPresence;
|
||||
import dev.ltms.bridged.session.MemberSession;
|
||||
import dev.ltms.bridged.session.WorktreeRequest;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.member.CompositePeerLauncher;
|
||||
import dev.ltms.bridged.placement.BackendQuarantine;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import io.modelcontextprotocol.spec.McpSchema;
|
||||
import dev.ltms.bridged.msg.InMemoryReplyInbox;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.CompletableFuture;
|
||||
@@ -423,6 +426,35 @@ class BridgeMcpTest {
|
||||
assertTrue(textOf(res).contains("unknown worker profile"), textOf(res));
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-599: a profile at its {@code maxLoad} cap must surface a readable reason on the MCP
|
||||
* surface too, not merely flip {@code isError} with an opaque or absent message.
|
||||
*/
|
||||
@Test
|
||||
void spawnAtMaxLoadSurfacesTheCapacityReason() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
BridgedConfig.Profile wcfg = new BridgedConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
|
||||
"tab", "bridged-workers", "worker: {profile} #{n}", null,
|
||||
null, null, null, null, null, null, null, 0, null, null, null);
|
||||
Map<String, BridgedConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
|
||||
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
|
||||
new AgentControl(h), new WorkspaceControl(h), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
profiles, wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok" : null);
|
||||
CompositePeerLauncher composite = new CompositePeerLauncher(
|
||||
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), _ -> 0);
|
||||
SessionManager sm = new SessionManager(composite);
|
||||
|
||||
McpSchema.CallToolResult res = BridgeMcp.spawn(sm, "ltms-local");
|
||||
|
||||
assertTrue(res.isError());
|
||||
String text = textOf(res);
|
||||
assertTrue(text.contains("no capacity"), "surfaces a capacity reason, not a bare error: " + text);
|
||||
assertTrue(text.contains("ltms-local"), "names the profile: " + text);
|
||||
assertTrue(text.contains("maxLoad"), "explains the refusal: " + text);
|
||||
assertFalse(h.called("agent.start"), "at cap, the spawn is refused before any herdr call");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnPassesTheRequestedCwdToTheWorker() {
|
||||
FakeHerdr h = new FakeHerdr();
|
||||
|
||||
@@ -16,6 +16,7 @@ import java.util.List;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertEquals;
|
||||
import static org.junit.jupiter.api.Assertions.assertThrows;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
@@ -166,6 +167,88 @@ class AmqpReplyInboxContractTest {
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void prefetchBoundsTheHeldBacklog() throws Exception {
|
||||
String target = "worker-prefetch-" + System.nanoTime();
|
||||
int prefetch = 4;
|
||||
int published = 10;
|
||||
try (AmqpReplyInbox inbox = new AmqpReplyInbox(newConnection(), prefetch);
|
||||
Connection inspect = newConnection()) {
|
||||
inbox.own(target);
|
||||
for (int i = 0; i < published; i++) {
|
||||
inbox.publish(target, "m" + i, "payload " + i);
|
||||
}
|
||||
awaitHeldAtLeast(inbox, target, prefetch);
|
||||
|
||||
long depth;
|
||||
try (Channel ch = inspect.createChannel()) {
|
||||
depth = ch.queueDeclarePassive(queueName(target)).getMessageCount();
|
||||
}
|
||||
assertTrue(depth >= published - prefetch,
|
||||
"broker should still hold at least " + (published - prefetch)
|
||||
+ " undelivered messages behind a prefetch of " + prefetch + ", saw " + depth);
|
||||
|
||||
drainUntilEmpty(inbox, target, inspect);
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void unroutablePublishReportsFailureNotSilentSuccess() throws Exception {
|
||||
String target = "worker-unroutable-" + System.nanoTime();
|
||||
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
|
||||
// Deliberately never own(target): the queue is never declared, so the default-exchange
|
||||
// route to agent.<target>.inbox does not exist and the broker must return the publish.
|
||||
IllegalStateException ex = assertThrows(IllegalStateException.class,
|
||||
() -> inbox.publish(target, "m1", "nobody home"));
|
||||
assertTrue(ex.getMessage() != null && ex.getMessage().toLowerCase().contains("unroutable"),
|
||||
"expected an unroutable-publish failure, got: " + ex.getMessage());
|
||||
}
|
||||
}
|
||||
|
||||
@Test
|
||||
void confirmedPublishDeliversNormally() throws Exception {
|
||||
String target = "worker-confirm-" + System.nanoTime();
|
||||
try (AmqpReplyInbox inbox = AmqpReplyInbox.open(uri())) {
|
||||
inbox.own(target);
|
||||
inbox.publish(target, "m1", "confirmed delivery"); // must return normally: routed and confirmed
|
||||
|
||||
List<ReplyInbox.InboxMessage> got = awaitPeek(inbox, target);
|
||||
assertEquals(1, got.size());
|
||||
assertEquals("confirmed delivery", got.getFirst().content());
|
||||
inbox.ack(target, "m1");
|
||||
}
|
||||
}
|
||||
|
||||
/** Poll peek until at least {@code n} replies for {@code target} are held, or ~10s elapse. */
|
||||
@SuppressWarnings("BusyWait")
|
||||
private static void awaitHeldAtLeast(AmqpReplyInbox inbox, String target, int n) throws InterruptedException {
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
|
||||
while (inbox.peek(target).size() < n && System.nanoTime() < deadline) {
|
||||
Thread.sleep(50);
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
* Repeatedly ack whatever is currently held (freeing prefetch slots for the next delivery) until
|
||||
* both the local snapshot and the broker's own queue depth are empty, or ~10s elapse.
|
||||
*/
|
||||
@SuppressWarnings("BusyWait")
|
||||
private static void drainUntilEmpty(AmqpReplyInbox inbox, String target, Connection inspect) throws Exception {
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(10);
|
||||
while (System.nanoTime() < deadline) {
|
||||
for (ReplyInbox.InboxMessage msg : inbox.peek(target)) {
|
||||
inbox.ack(target, msg.msgId());
|
||||
}
|
||||
try (Channel ch = inspect.createChannel()) {
|
||||
if (ch.queueDeclarePassive(queueName(target)).getMessageCount() == 0 && inbox.peek(target).isEmpty()) {
|
||||
return;
|
||||
}
|
||||
}
|
||||
Thread.sleep(50);
|
||||
}
|
||||
throw new AssertionError("did not drain " + target + " to empty within the deadline");
|
||||
}
|
||||
|
||||
/** Poll peek (broker delivery is async) until a reply for {@code target} appears or ~10s elapse. */
|
||||
@SuppressWarnings("BusyWait") // deliberate poll for async broker delivery, bounded by the deadline
|
||||
private static List<ReplyInbox.InboxMessage> awaitPeek(AmqpReplyInbox inbox, String target)
|
||||
|
||||
@@ -0,0 +1,264 @@
|
||||
package dev.ltms.bridged.msg;
|
||||
|
||||
import com.rabbitmq.client.AMQP;
|
||||
import com.rabbitmq.client.Channel;
|
||||
import com.rabbitmq.client.ConfirmCallback;
|
||||
import com.rabbitmq.client.Connection;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.junit.jupiter.api.Timeout;
|
||||
|
||||
import java.lang.reflect.InvocationHandler;
|
||||
import java.lang.reflect.Proxy;
|
||||
import java.util.List;
|
||||
import java.util.concurrent.CopyOnWriteArrayList;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import java.util.concurrent.atomic.AtomicLong;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import static org.junit.jupiter.api.Assertions.assertNotNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertNull;
|
||||
import static org.junit.jupiter.api.Assertions.assertTrue;
|
||||
|
||||
/**
|
||||
* CB-528 follow-up: {@link AmqpReplyInbox#failPendingPublishesOnRecovery()} must not fail a publish
|
||||
* that registers concurrently with (and was not yet visible when) the recovery sweep began — that
|
||||
* would report a publish that actually succeeded as failed, and {@code MessageService.reply} retries
|
||||
* with a fresh {@code msgId}, so the reply is delivered twice. A live broker reconnect cannot be
|
||||
* forced reliably, so this drives {@link AmqpReplyInbox#failPendingPublishesOnRecovery} and a real
|
||||
* {@link AmqpReplyInbox#publish} against each other directly, against fake AMQP channels built with
|
||||
* {@link Proxy} (no mocking library is on the classpath).
|
||||
*
|
||||
* <p>Also covers Finding 2 (CB-528 follow-up): {@link AmqpReplyInbox#close()} must fail an in-flight
|
||||
* publish promptly instead of leaving it to idle out the 10s confirm timeout.
|
||||
*/
|
||||
class AmqpReplyInboxRecoveryRaceTest {
|
||||
|
||||
/** Large enough that the (unfixed) unsynchronized sweep's iteration is a real, observable window
|
||||
* a concurrently-started publish can land in — not just a best case, single-entry sprint. */
|
||||
private static final int STALE_PUBLISHES = 100_000;
|
||||
|
||||
@Test
|
||||
@Timeout(30)
|
||||
void recoverySweepDoesNotFailAPublishThatRegistersWhileItIsRunning() throws Exception {
|
||||
AtomicLong seqCounter = new AtomicLong();
|
||||
List<Long> seqOrder = new CopyOnWriteArrayList<>();
|
||||
List<String> msgIdOrder = new CopyOnWriteArrayList<>();
|
||||
AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
|
||||
AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
|
||||
|
||||
Channel publishChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
|
||||
Channel consumeChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
|
||||
Connection connection = fakeConnection(consumeChannel, publishChannel);
|
||||
|
||||
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
|
||||
|
||||
// STALE_PUBLISHES in-flight publishes that never get confirmed — they sit in pendingBySeq /
|
||||
// pendingByMsgId exactly like publishes whose confirm never arrived before a connection drop.
|
||||
// Virtual threads make this many concurrent blocking publish() calls cheap.
|
||||
CountDownLatch staleStarted = new CountDownLatch(STALE_PUBLISHES);
|
||||
for (int i = 0; i < STALE_PUBLISHES; i++) {
|
||||
String msgId = "stale-" + i;
|
||||
Thread.ofVirtual().start(() -> {
|
||||
staleStarted.countDown();
|
||||
try {
|
||||
inbox.publish("worker-stale", msgId, "x");
|
||||
} catch (IllegalStateException expected) {
|
||||
// resolved (failed by the sweep) — that is exactly what this thread is here for
|
||||
}
|
||||
});
|
||||
}
|
||||
staleStarted.await();
|
||||
// Let the registrations (the synchronized put into pendingBySeq/pendingByMsgId) actually land
|
||||
// for all of them before the sweep starts, so the sweep begins with a large, real backlog.
|
||||
Thread.sleep(300);
|
||||
|
||||
AtomicReference<Throwable> sweepError = new AtomicReference<>();
|
||||
Thread sweepThread = new Thread(() -> {
|
||||
try {
|
||||
inbox.failPendingPublishesOnRecovery();
|
||||
} catch (Throwable t) {
|
||||
sweepError.set(t);
|
||||
}
|
||||
}, "recovery-sweep");
|
||||
sweepThread.start();
|
||||
// A short, deliberate head start: with STALE_PUBLISHES this large, the (unfixed) sweep's own
|
||||
// iteration takes several milliseconds, so this guarantees the sweep has already begun —
|
||||
// and, once guarded, is already holding publishChannelLock — before "fresh" attempts to
|
||||
// register. Without this head start, "fresh" sometimes wins the race for the lock and
|
||||
// registers before the sweep even starts, which is the accepted "already in flight when
|
||||
// recovery fires" case (correctly failed either way) rather than the bug under test.
|
||||
Thread.sleep(5);
|
||||
|
||||
// This is the exact interleaving CB-528's follow-up describes: "the still-running recovery
|
||||
// sweep" racing a publish that registers while it is mid-flight.
|
||||
AtomicReference<Throwable> publishError = new AtomicReference<>();
|
||||
Thread freshThread = new Thread(() -> {
|
||||
try {
|
||||
inbox.publish("worker-fresh", "fresh", "hello");
|
||||
} catch (Throwable t) {
|
||||
publishError.set(t);
|
||||
}
|
||||
}, "fresh-publish");
|
||||
freshThread.start();
|
||||
|
||||
sweepThread.join(20_000);
|
||||
assertNull(sweepError.get(), "sweep threw: " + sweepError.get());
|
||||
|
||||
// Simulate the broker's real confirm for "fresh" now that the sweep is done, so a correct
|
||||
// implementation's publish() returns normally instead of idling out CONFIRM_TIMEOUT_MS. Poll
|
||||
// for the registration rather than checking once: freshThread may still be contending for
|
||||
// publishChannelLock (behind the 20,000 stale threads' own lock acquisitions) even though the
|
||||
// sweep itself has already finished.
|
||||
long deadline = System.nanoTime() + TimeUnit.SECONDS.toNanos(9);
|
||||
int idx = -1;
|
||||
while (idx < 0 && System.nanoTime() < deadline) {
|
||||
idx = msgIdOrder.indexOf("fresh");
|
||||
if (idx < 0) {
|
||||
Thread.sleep(20);
|
||||
}
|
||||
}
|
||||
if (idx >= 0 && ackCallback.get() != null) {
|
||||
ackCallback.get().handle(seqOrder.get(idx), false);
|
||||
}
|
||||
freshThread.join(15_000);
|
||||
|
||||
assertNull(publishError.get(),
|
||||
"a publish that registered while the recovery sweep was running must not be failed by "
|
||||
+ "it, but got: " + publishError.get());
|
||||
}
|
||||
|
||||
@Test
|
||||
@Timeout(15)
|
||||
void closeFailsInFlightPublishPromptlyInsteadOfWaitingOutTheConfirmTimeout() throws Exception {
|
||||
AtomicLong seqCounter = new AtomicLong();
|
||||
List<Long> seqOrder = new CopyOnWriteArrayList<>();
|
||||
List<String> msgIdOrder = new CopyOnWriteArrayList<>();
|
||||
AtomicReference<ConfirmCallback> ackCallback = new AtomicReference<>();
|
||||
AtomicReference<ConfirmCallback> nackCallback = new AtomicReference<>();
|
||||
|
||||
Channel publishChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
|
||||
Channel consumeChannel = fakeChannel(seqCounter, seqOrder, msgIdOrder, ackCallback, nackCallback);
|
||||
Connection connection = fakeConnection(consumeChannel, publishChannel);
|
||||
|
||||
AmqpReplyInbox inbox = new AmqpReplyInbox(connection, AmqpReplyInbox.DEFAULT_PREFETCH);
|
||||
|
||||
CountDownLatch publishReturned = new CountDownLatch(1);
|
||||
AtomicReference<Throwable> publishError = new AtomicReference<>();
|
||||
AtomicLong elapsedMillis = new AtomicLong();
|
||||
Thread publishThread = new Thread(() -> {
|
||||
long start = System.nanoTime();
|
||||
try {
|
||||
inbox.publish("worker-close", "never-confirmed", "x");
|
||||
} catch (Throwable t) {
|
||||
publishError.set(t);
|
||||
} finally {
|
||||
elapsedMillis.set((System.nanoTime() - start) / 1_000_000);
|
||||
publishReturned.countDown();
|
||||
}
|
||||
}, "publish-during-close");
|
||||
publishThread.start();
|
||||
|
||||
Thread.sleep(200); // let publish() register before close() runs
|
||||
inbox.close();
|
||||
|
||||
assertTrue(publishReturned.await(5, TimeUnit.SECONDS), "publish() did not return after close()");
|
||||
assertNotNull(publishError.get(), "a publish in flight when close() runs must fail, not hang");
|
||||
assertTrue(publishError.get().getMessage() != null
|
||||
&& publishError.get().getMessage().toLowerCase().contains("closed"),
|
||||
"expected a clear closed-inbox message, got: " + publishError.get());
|
||||
assertTrue(elapsedMillis.get() < 5_000,
|
||||
"close() should fail the in-flight publish promptly, not wait out the confirm timeout — took "
|
||||
+ elapsedMillis.get() + "ms");
|
||||
}
|
||||
|
||||
/** A {@link Proxy}-backed {@link Channel}: only the calls {@link AmqpReplyInbox} actually makes
|
||||
* are meaningfully implemented; everything else returns a harmless default. */
|
||||
private static Channel fakeChannel(AtomicLong seqCounter, List<Long> seqOrder, List<String> msgIdOrder,
|
||||
AtomicReference<ConfirmCallback> ackCallback,
|
||||
AtomicReference<ConfirmCallback> nackCallback) {
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
String name = method.getName();
|
||||
if (name.equals("getNextPublishSeqNo")) {
|
||||
long value = seqCounter.incrementAndGet();
|
||||
seqOrder.add(value);
|
||||
return value;
|
||||
}
|
||||
if (name.equals("basicPublish")) {
|
||||
AMQP.BasicProperties props = (AMQP.BasicProperties) args[3];
|
||||
msgIdOrder.add(props.getMessageId());
|
||||
return null;
|
||||
}
|
||||
if (name.equals("addConfirmListener")) {
|
||||
ackCallback.set((ConfirmCallback) args[0]);
|
||||
nackCallback.set((ConfirmCallback) args[1]);
|
||||
return null;
|
||||
}
|
||||
if (name.equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (name.equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
if (name.equals("toString")) {
|
||||
return "FakeChannel";
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Channel) Proxy.newProxyInstance(AmqpReplyInboxRecoveryRaceTest.class.getClassLoader(),
|
||||
new Class<?>[] {Channel.class}, handler);
|
||||
}
|
||||
|
||||
/** A {@link Proxy}-backed {@link Connection} handing out {@code first} then {@code second} from
|
||||
* successive {@code createChannel()} calls, matching {@link AmqpReplyInbox}'s constructor. */
|
||||
private static Connection fakeConnection(Channel first, Channel second) {
|
||||
AtomicInteger calls = new AtomicInteger();
|
||||
InvocationHandler handler = (proxy, method, args) -> {
|
||||
String name = method.getName();
|
||||
if (name.equals("createChannel") && (args == null || args.length == 0)) {
|
||||
return calls.getAndIncrement() == 0 ? first : second;
|
||||
}
|
||||
if (name.equals("equals")) {
|
||||
return proxy == args[0];
|
||||
}
|
||||
if (name.equals("hashCode")) {
|
||||
return System.identityHashCode(proxy);
|
||||
}
|
||||
if (name.equals("toString")) {
|
||||
return "FakeConnection";
|
||||
}
|
||||
return defaultValue(method.getReturnType());
|
||||
};
|
||||
return (Connection) Proxy.newProxyInstance(AmqpReplyInboxRecoveryRaceTest.class.getClassLoader(),
|
||||
new Class<?>[] {Connection.class}, handler);
|
||||
}
|
||||
|
||||
private static Object defaultValue(Class<?> type) {
|
||||
if (!type.isPrimitive() || type == void.class) {
|
||||
return null;
|
||||
}
|
||||
if (type == boolean.class) {
|
||||
return Boolean.FALSE;
|
||||
}
|
||||
if (type == long.class) {
|
||||
return 0L;
|
||||
}
|
||||
if (type == short.class) {
|
||||
return (short) 0;
|
||||
}
|
||||
if (type == byte.class) {
|
||||
return (byte) 0;
|
||||
}
|
||||
if (type == char.class) {
|
||||
return (char) 0;
|
||||
}
|
||||
if (type == double.class) {
|
||||
return 0.0d;
|
||||
}
|
||||
if (type == float.class) {
|
||||
return 0.0f;
|
||||
}
|
||||
return 0;
|
||||
}
|
||||
}
|
||||
@@ -81,7 +81,7 @@ class ReplyPushLoopTest {
|
||||
@Test
|
||||
void decideWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -90,7 +90,7 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(2, 100);
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -99,7 +99,7 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -108,7 +108,7 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
|
||||
"BLOCKED is injectable");
|
||||
}
|
||||
|
||||
@@ -118,7 +118,7 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0),
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0),
|
||||
"DONE is injectable");
|
||||
}
|
||||
|
||||
@@ -128,7 +128,7 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -137,7 +137,7 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -146,9 +146,9 @@ class ReplyPushLoopTest {
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop();
|
||||
loop.onReplyQueued(WORKER);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
|
||||
inbox.ack(WORKER, "m1");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
// --- onReplyQueued integration -------------------------------------------------------------
|
||||
@@ -240,7 +240,7 @@ class ReplyPushLoopTest {
|
||||
@Test
|
||||
void decideTicketsWithNothingPendingIsStop() {
|
||||
agents = agentWithStatus("idle");
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop().decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -248,7 +248,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(2, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -256,7 +256,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -264,7 +264,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("working");
|
||||
var loop = loop(5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.WAIT_BUSY, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -272,7 +272,7 @@ class ReplyPushLoopTest {
|
||||
agents = agentWithStatus("idle");
|
||||
var loop = new ReplyPushLoop(new PrimaryRegistry(null), agents, inbox, scheduler, 5, 100_000);
|
||||
loop.onTicketTerminal("task-1", WORKER, false); // no lead known -> never registered as pending
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 0));
|
||||
}
|
||||
|
||||
// --- CB-588: onTicketTerminal integration ---------------------------------------------------
|
||||
@@ -420,6 +420,49 @@ class ReplyPushLoopTest {
|
||||
+ "restart the loop — that would defeat the reminder cap");
|
||||
}
|
||||
|
||||
// --- mirror of the two stopOrRestart tests above, for the reply arm ------------------------
|
||||
//
|
||||
// Both tests above only ever passed Set.of() for repliesBefore, so racedIn's reply branch
|
||||
// (`pendingReplyTargetsFor(lead).stream().anyMatch(t -> !repliesBefore.contains(t))`) was
|
||||
// never exercised by anything other than an always-empty snapshot. The reviewer flagged this:
|
||||
// racedIn is symmetric in the code, and only half of it was pinned by a test.
|
||||
|
||||
@Test
|
||||
void aReplyStillPendingWhenTheLoopStopsIsNotStranded() {
|
||||
// Mirrors aTicketStillPendingWhenTheLoopStopsIsNotStranded: a reply target that raced in
|
||||
// during the decision-to-release window (absent from the "before" snapshot) must reclaim
|
||||
// the schedule slot rather than being stranded with no schedule left to nudge about it.
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
ReplyPushLoop loop = loop(1, 100_000); // long backoff — no natural tick fires during this test
|
||||
|
||||
loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY}
|
||||
|
||||
loop.stopOrRestart(PRIMARY, Set.of(), Set.of());
|
||||
|
||||
assertTrue(loop.isActive(), "a reply that raced the loop's stop must reclaim the schedule "
|
||||
+ "slot, not be stranded with no schedule left to ever nudge about it");
|
||||
}
|
||||
|
||||
@Test
|
||||
void aStaleUncollectedReplyAtCapDoesNotRestartTheLoop() {
|
||||
// Mirrors aStaleUncollectedTicketAtCapDoesNotRestartTheLoop: a reply target already
|
||||
// accounted for at decide time (present in repliesBefore) must not restart the loop —
|
||||
// that is the reminder cap doing its job, not a race.
|
||||
var rec = recordingClient();
|
||||
agents = new AgentControl(rec);
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
ReplyPushLoop loop = loop(1, 100_000);
|
||||
|
||||
loop.onReplyQueued(WORKER); // pendingReplies={term_worker}; activeLeads={PRIMARY}
|
||||
|
||||
loop.stopOrRestart(PRIMARY, Set.of(WORKER), Set.of());
|
||||
|
||||
assertFalse(loop.isActive(), "a stale reply target already accounted for at decide time "
|
||||
+ "must not restart the loop — that would defeat the reminder cap");
|
||||
}
|
||||
|
||||
@Test
|
||||
void ticketNudgeFormatIsCorrect() {
|
||||
String single = ReplyPushLoop.TICKET_NUDGE_FORMAT.formatted("task-1", "", "task-1");
|
||||
@@ -485,6 +528,57 @@ class ReplyPushLoopTest {
|
||||
assertTrue(nudge.contains("task-1"), "the ticket must not be dropped: " + nudge);
|
||||
}
|
||||
|
||||
// --- CB-590 follow-up: per-source reminder budgets — the regression this round exists for ---
|
||||
|
||||
@Test
|
||||
void oneExhaustedSourceDoesNotBlockANudgeForTheOtherSource() {
|
||||
// The live trace this ticket was filed from: an undrained reply target got nudged up to
|
||||
// its cap (5 reminders), then a ticket for the SAME lead went terminal shortly before the
|
||||
// next scheduled tick — so it coalesced onto the still-active schedule (arriving BEFORE
|
||||
// that tick's "before" snapshot, not during the decision-to-release race stopOrRestart
|
||||
// guards). With CB-590's single shared reminder counter, that tick's decide() saw
|
||||
// reminderCount already at the cap and returned STOP regardless of the ticket, and because
|
||||
// the ticket was already present in that tick's "before" snapshot, stopOrRestart's
|
||||
// racedIn check (proven correct on its own above) did not save it either — it is not a
|
||||
// race, it looks like ordinary stale backlog. The ticket was then stranded: pending
|
||||
// forever with no live schedule, never named in any nudge.
|
||||
//
|
||||
// Fixed by giving each source its own counter. Here the reply source is AT its cap (2/2)
|
||||
// and the ticket source has NEVER been nudged (0/2) — decide() must still return INJECT,
|
||||
// because the ticket is still eligible on its own budget.
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(2, 100_000); // huge backoff — this test drives decide()/isActive() directly
|
||||
loop.onReplyQueued(WORKER);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.INJECT, loop.decide(PRIMARY, 2, 0),
|
||||
"the reply source is exhausted (2/2), but the ticket source has never been "
|
||||
+ "nudged (0/2) — the lead must still be injected so the ticket is not "
|
||||
+ "lost, exactly the CB-590 follow-up regression");
|
||||
|
||||
// isActive() (criterion #5): a real tick that takes the INJECT branch above never calls
|
||||
// stopOrRestart, so the schedule started by onTicketTerminal above stays live — the
|
||||
// ticket is not left stranded with isActive()==false while it is still pending.
|
||||
assertTrue(loop.isActive(), "the schedule must stay active while the ticket source still "
|
||||
+ "has budget left, even though the reply source sharing it is exhausted");
|
||||
}
|
||||
|
||||
@Test
|
||||
void bothSourcesExhaustedIsStillStop() {
|
||||
// The flip side: per-source budgets must not turn into unbounded nudging. When BOTH
|
||||
// sources are at their cap, decide() must still STOP — a per-source budget is still a
|
||||
// budget.
|
||||
agents = agentWithStatus("idle");
|
||||
inbox.publish(WORKER, "m1", "hello");
|
||||
var loop = loop(2, 100_000);
|
||||
loop.onReplyQueued(WORKER);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 2),
|
||||
"both the reply and the ticket source are at their own cap — must still stop");
|
||||
}
|
||||
|
||||
// --- metrics (CB-512) ----------------------------------------------------------------------
|
||||
|
||||
@Test
|
||||
@@ -513,7 +607,7 @@ class ReplyPushLoopTest {
|
||||
var loop = loop(2, 100, metrics);
|
||||
loop.onReplyQueued(WORKER);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2, 0));
|
||||
|
||||
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the reminder cap must count as exhausted");
|
||||
@@ -541,7 +635,7 @@ class ReplyPushLoopTest {
|
||||
var loop = loop(2, 100_000, metrics);
|
||||
loop.onTicketTerminal("task-1", WORKER, false);
|
||||
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 2));
|
||||
assertEquals(ReplyPushLoop.Action.STOP, loop.decide(PRIMARY, 0, 2));
|
||||
|
||||
assertEquals(1, metrics.count(BridgedMetrics.PUSH_NUDGES, "outcome", "exhausted"),
|
||||
"hitting the ticket reminder cap must count as exhausted");
|
||||
|
||||
@@ -18,6 +18,8 @@ import dev.ltms.bridged.session.GitWorktrees;
|
||||
import dev.ltms.bridged.session.SessionManager;
|
||||
import dev.ltms.bridged.session.Worktrees;
|
||||
import dev.ltms.bridged.member.ClaudeCodeLauncher;
|
||||
import dev.ltms.bridged.member.CompositePeerLauncher;
|
||||
import dev.ltms.bridged.placement.PlacementPolicies;
|
||||
import io.javalin.Javalin;
|
||||
import org.junit.jupiter.api.AfterEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
@@ -240,6 +242,43 @@ class BridgedAppTest {
|
||||
assertFalse(herdr.called("agent.start"), "an unknown profile must not spawn anything");
|
||||
}
|
||||
|
||||
/**
|
||||
* CB-599: a profile at its {@code maxLoad} cap must not surface as a bare 500 — the caller
|
||||
* needs a structured, readable reason, distinct from "unknown_profile".
|
||||
*/
|
||||
@Test
|
||||
void spawnAtMaxLoadIs503WithTheCapacityReasonNotABare500() throws Exception {
|
||||
FakeHerdr herdr = new FakeHerdr();
|
||||
BridgedConfig.Profile wcfg = new BridgedConfig.Profile(
|
||||
"ltms-local", "http://gx00.gw:8000", "coder", null, "BRIDGED_WORKER_TOKEN", null,
|
||||
"tab", "bridged-workers", "worker: {profile} #{n}", null,
|
||||
null, null, null, null, null, null, null, 0, null, null, null);
|
||||
Map<String, BridgedConfig.Profile> profiles = Map.of(wcfg.profile(), wcfg);
|
||||
ClaudeCodeLauncher delegate = new ClaudeCodeLauncher(
|
||||
new AgentControl(herdr), new WorkspaceControl(herdr), new SubscriptionGuard(Set.of("gx00.gw")),
|
||||
profiles, wcfg.profile(), k -> "BRIDGED_WORKER_TOKEN".equals(k) ? "tok-abc" : null);
|
||||
CompositePeerLauncher workers = new CompositePeerLauncher(
|
||||
List.of(delegate), wcfg.profile(), profiles, PlacementPolicies.fixed(), _ -> 0);
|
||||
SessionManager sessions = new SessionManager(workers, new GitWorktrees());
|
||||
this.presence = sessions.asPresence();
|
||||
Injector injector = new Injector(new AgentControl(herdr));
|
||||
InMemoryReplyInbox inbox = new InMemoryReplyInbox();
|
||||
sessions.onAcquire(inbox::own);
|
||||
MessageService messages = new MessageService(new AgentControl(herdr), injector, new Rendezvous(), inbox);
|
||||
app = new BridgedApp(herdr, workers, sessions, messages, this.presence, null).build().start("127.0.0.1", 0);
|
||||
int port = app.port();
|
||||
|
||||
HttpResponse<String> res = req(port, "POST", "/members?profile=ltms-local");
|
||||
|
||||
assertEquals(503, res.statusCode(), res.body());
|
||||
JsonNode body = mapper.readTree(res.body());
|
||||
assertEquals("no_capacity", body.get("error").asText());
|
||||
String detail = body.get("detail").asText();
|
||||
assertTrue(detail.contains("ltms-local"), "detail names the profile: " + detail);
|
||||
assertTrue(detail.contains("maxLoad"), "detail explains the refusal: " + detail);
|
||||
assertFalse(herdr.called("agent.start"), "at cap, the spawn is refused before any herdr call");
|
||||
}
|
||||
|
||||
@Test
|
||||
void spawnWorkerReusesExistingWorkerSpace() throws Exception {
|
||||
// A space labelled "bridged-workers" already exists → no second workspace.create.
|
||||
|
||||
@@ -1,7 +1,7 @@
|
||||
<?xml version="1.0" encoding="UTF-8"?>
|
||||
<!DOCTYPE plist PUBLIC "-//Apple//DTD PLIST 1.0//EN" "http://www.apple.com/DTDs/PropertyList-1.0.dtd">
|
||||
<!--
|
||||
CB-504 — launchd agent for bridged (macOS).
|
||||
CB-504 / CB-594 — launchd agent for bridged (macOS).
|
||||
|
||||
This is the real supervision target today: the dogfooded daemon runs on macOS, where there is
|
||||
no systemd. A systemd unit ships alongside (deploy/bridged.service) for the Linux gateways
|
||||
@@ -9,14 +9,33 @@
|
||||
|
||||
Install:
|
||||
cp deploy/dev.ltms.bridged.plist ~/Library/LaunchAgents/
|
||||
# edit the paths + JAVA_HOME below to match this host, then:
|
||||
launchctl load -w ~/Library/LaunchAgents/dev.ltms.bridged.plist
|
||||
launchctl list | grep bridged
|
||||
|
||||
The paths below are already filled in for this host (resolved 2026-08-16 from
|
||||
`/usr/libexec/java_home`... except that reported the system Applet-plugin JVM, not the jenv-
|
||||
managed JDK 25 actually used to build/run bridged, so JAVA_HOME here is the real one:
|
||||
`JENV_VERSION=25.0.3 java -XshowSettings:properties -version 2>&1 | grep java.home`; `which mvn`;
|
||||
`echo $HOME`). If this file is copied to a different host, re-resolve all three paths and check
|
||||
no placeholder path is left behind; scripts/redeploy-bridged.sh's check mode does not (and
|
||||
cannot) check this file for you.
|
||||
|
||||
CB-594 — launchd cannot run a login shell (see the PATH comment on EnvironmentVariables below,
|
||||
and scripts/bridged-launchd-wrapper.sh for the fix): ProgramArguments below execs THAT wrapper,
|
||||
not java directly, so WORKER_GITEA_TOKEN and AI_GATEWAY_TOKEN still get sourced from
|
||||
${SHARED_ENV}/tools/secrets.sh even though launchd itself never sources anything.
|
||||
|
||||
Note on ordering: launchd has no "start after herdr" primitive for user agents, and neither
|
||||
does systemd in a way that survives a socket appearing late. bridged retries the herdr socket
|
||||
on startup instead, so an agent that comes up before herdr converges rather than dying — that
|
||||
retry is the actual fix; KeepAlive below is the backstop.
|
||||
|
||||
CB-594 — KeepAlive vs. scripts/redeploy-bridged.sh: a bare SIGTERM makes this JVM exit 143 even
|
||||
with its shutdown hook running to completion (measured, see the CB-594 report), which
|
||||
SuccessfulExit:false below reads as a crash and races to restart the OLD jar. The redeploy
|
||||
script now detects a loaded agent and uses `launchctl unload`/`load` instead of a raw kill, so
|
||||
only one supervisor ever touches the process at a time — read that script's own output on a
|
||||
redeploy for the confirmation.
|
||||
-->
|
||||
<plist version="1.0">
|
||||
<dict>
|
||||
@@ -25,22 +44,23 @@
|
||||
|
||||
<key>ProgramArguments</key>
|
||||
<array>
|
||||
<string>/Users/CHANGEME/Tool/jdk-25.0.2.jdk/Contents/Home/bin/java</string>
|
||||
<string>/Users/dai.ha/LTMS/claude-bridge/scripts/bridged-launchd-wrapper.sh</string>
|
||||
<string>/Users/dai.ha/Softwares/jdks/jdk-25.0.3.jdk/Contents/Home/bin/java</string>
|
||||
<string>-jar</string>
|
||||
<string>/Users/CHANGEME/src/claude-bridge/bridged/target/bridged.jar</string>
|
||||
<string>/Users/dai.ha/LTMS/claude-bridge/bridged/target/bridged.jar</string>
|
||||
<string>bridged.yaml</string>
|
||||
</array>
|
||||
|
||||
<!-- Config path in ProgramArguments is relative, so the working directory must be the module. -->
|
||||
<key>WorkingDirectory</key>
|
||||
<string>/Users/CHANGEME/src/claude-bridge/bridged</string>
|
||||
<string>/Users/dai.ha/LTMS/claude-bridge/bridged</string>
|
||||
|
||||
<key>EnvironmentVariables</key>
|
||||
<dict>
|
||||
<key>JAVA_HOME</key>
|
||||
<string>/Users/CHANGEME/Tool/jdk-25.0.2.jdk/Contents/Home</string>
|
||||
<string>/Users/dai.ha/Softwares/jdks/jdk-25.0.3.jdk/Contents/Home</string>
|
||||
<key>HERDR_SOCKET_PATH</key>
|
||||
<string>/Users/CHANGEME/.config/herdr/herdr.sock</string>
|
||||
<string>/Users/dai.ha/.config/herdr/herdr.sock</string>
|
||||
<!--
|
||||
PATH matters more than it looks (CB-511): bridged propagates its own PATH to every worker
|
||||
it spawns, so this line decides whether the fleet can run a build at all. launchd does NOT
|
||||
@@ -48,19 +68,39 @@
|
||||
bare /usr/bin:/bin and no JDK or Maven. Keep the toolchain entries first.
|
||||
-->
|
||||
<key>PATH</key>
|
||||
<string>/Users/CHANGEME/Tool/jdk-25.0.2.jdk/Contents/Home/bin:/Users/CHANGEME/Tool/apache-maven-3.9.16/bin:/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin</string>
|
||||
<string>/Users/dai.ha/Softwares/jdks/jdk-25.0.3.jdk/Contents/Home/bin:/Users/dai.ha/Softwares/apache-maven/bin:/opt/homebrew/bin:/usr/local/bin:/usr/bin:/bin:/usr/sbin:/sbin</string>
|
||||
<!--
|
||||
Worker/API tokens are NOT set here: this file is committed. Export them from a private
|
||||
launchd override or a wrapper script. bridged reads the API token from the env var named
|
||||
by auth.tokenEnv (default BRIDGED_API_TOKEN) and only in auth.mode: token.
|
||||
Worker/API tokens are NOT set here: this file is committed. CB-594 —
|
||||
scripts/bridged-launchd-wrapper.sh (named in ProgramArguments above) is what supplies
|
||||
them, by execing a login shell that sources ${SHARED_ENV}/tools/secrets.sh before the
|
||||
daemon itself starts. bridged also reads the API token from the env var named by
|
||||
auth.tokenEnv (default BRIDGED_API_TOKEN) and only in auth.mode: token — the wrapper
|
||||
covers that one too, since it is the same login shell.
|
||||
-->
|
||||
</dict>
|
||||
|
||||
<key>RunAtLoad</key>
|
||||
<true/>
|
||||
|
||||
<!-- Restart on crash, but not in a tight loop if the config is bad (bridged fails fast on a
|
||||
non-loopback bind without token auth — that is a config error, not a transient one). -->
|
||||
<!--
|
||||
CB-600 — read this before assuming ThrottleInterval bounds anything. It paces restarts to at
|
||||
most one per 10s; it does NOT cap how many times launchd retries. If bridged fails fast on
|
||||
every start — a bad bridged.yaml, for example auth.mode: token with the token env var unset,
|
||||
which throws in main() before the daemon ever binds a port — launchd restarts it forever,
|
||||
once every 10s, until a human intervenes. LaunchAgents have no "give up after N attempts"
|
||||
primitive, so this is not something a config change here can fix.
|
||||
|
||||
That loop stops only two ways: (1) `launchctl unload -w ~/Library/LaunchAgents/dev.ltms.bridged.plist`,
|
||||
or (2) the underlying cause gets fixed, so the process starts successfully and stays up (no
|
||||
more exits to restart). scripts/redeploy-bridged.sh does not add a third way — it does not
|
||||
make bridged self-disable on a config error, on purpose: a fail-fast exit path that
|
||||
sometimes decides "this is unrecoverable, stop trying" is one more thing that can misfire,
|
||||
and a wrongly self-disabled daemon needs the exact same manual `launchctl load -w` recovery
|
||||
this comment already names — so it buys nothing an operator watching for the crash loop
|
||||
doesn't already have, at the cost of a new way to be silently down. Watch for it with
|
||||
`launchctl list dev.ltms.bridged` (a high restart count) or by tailing bridged.out for the
|
||||
same startup error repeating every ~10s.
|
||||
-->
|
||||
<key>KeepAlive</key>
|
||||
<dict>
|
||||
<key>SuccessfulExit</key>
|
||||
@@ -69,10 +109,17 @@
|
||||
<key>ThrottleInterval</key>
|
||||
<integer>10</integer>
|
||||
|
||||
<!--
|
||||
CB-594 — same file scripts/redeploy-bridged.sh already tails ($BRIDGED/bridged.out), and both
|
||||
streams point at it, not two separate log files: the script's fresh-line / ERROR-count checks
|
||||
after a restart read this one path regardless of whether launchd or the script started the
|
||||
process, and a stdout/stderr split would make half of what happens during a launchd-driven
|
||||
restart invisible to it.
|
||||
-->
|
||||
<key>StandardOutPath</key>
|
||||
<string>/Users/CHANGEME/src/claude-bridge/bridged/logs/bridged.out.log</string>
|
||||
<string>/Users/dai.ha/LTMS/claude-bridge/bridged/bridged.out</string>
|
||||
<key>StandardErrorPath</key>
|
||||
<string>/Users/CHANGEME/src/claude-bridge/bridged/logs/bridged.err.log</string>
|
||||
<string>/Users/dai.ha/LTMS/claude-bridge/bridged/bridged.out</string>
|
||||
|
||||
<key>ProcessType</key>
|
||||
<string>Background</string>
|
||||
|
||||
Executable
+32
@@ -0,0 +1,32 @@
|
||||
#!/usr/bin/env bash
|
||||
#
|
||||
# CB-594 — the only reason this file exists: launchd does not run a login shell.
|
||||
#
|
||||
# WORKER_GITEA_TOKEN and AI_GATEWAY_TOKEN live in ${SHARED_ENV}/tools/secrets.sh, sourced only by a
|
||||
# LOGIN shell (.zprofile/.zshrc etc). launchd execs a job's ProgramArguments directly — no shell, no
|
||||
# profile, nothing sourced (the plist's own PATH comment documents the same gap one variable over).
|
||||
# A daemon started that way boots fine and looks healthy; the failure is invisible until a worker
|
||||
# tries to open a PR (WORKER_GITEA_TOKEN empty) or a gateway profile gets a 401 (AI_GATEWAY_TOKEN
|
||||
# empty) — hours later, with nothing tying the two together (CB-591, CLAUDE.md "Redeploying the
|
||||
# daemon"). Bridged now also logs which required secret names resolved at startup (see
|
||||
# Bridged.reportRequiredSecrets), but that log line can only tell the truth if the tokens had a
|
||||
# chance to be sourced in the first place — which is this script's entire job.
|
||||
#
|
||||
# So: launchd execs THIS script instead of java directly. This script execs a login shell
|
||||
# ('zsh -l'), which sources secrets.sh, and that shell execs the real command in its place — one
|
||||
# process throughout (exec, not a subshell fork), so launchd's PID tracking, KeepAlive, and
|
||||
# StandardOut/ErrorPath all still see the one process they expect.
|
||||
#
|
||||
# The plist passes the full command as THIS script's own arguments, e.g.:
|
||||
# ProgramArguments = [ .../bridged-launchd-wrapper.sh, /path/to/java, -jar, /path/to/bridged.jar,
|
||||
# bridged.yaml ]
|
||||
# so the wrapper stays generic and the actual command lives in exactly one place (the plist), not
|
||||
# duplicated here.
|
||||
set -euo pipefail
|
||||
|
||||
if [ "$#" -eq 0 ]; then
|
||||
echo "bridged-launchd-wrapper.sh: no command given — check the plist's ProgramArguments" >&2
|
||||
exit 2
|
||||
fi
|
||||
|
||||
exec /bin/zsh -lc 'exec "$@"' -- "$@"
|
||||
+137
-9
@@ -5,13 +5,14 @@
|
||||
# A merge is not a deployment: the running daemon holds the jar it was started with, so code merged
|
||||
# to main does nothing until this runs. See CLAUDE.md -> "Redeploying the daemon".
|
||||
#
|
||||
# This script exists to turn five remembered traps into one auditable command:
|
||||
# This script exists to turn six remembered traps into one auditable command:
|
||||
#
|
||||
# 1. A piped `mvn` hides BUILD FAILURE behind a zero exit, so the build here is never piped.
|
||||
# 2. The daemon must start from a LOGIN shell, or the tokens it hands to members are empty:
|
||||
# WORKER_GITEA_TOKEN (workers cannot open a PR) and AI_GATEWAY_TOKEN (401 at llm.ltms.dev).
|
||||
# Both are read from the DAEMON's own environment at spawn time, so a value added to
|
||||
# secrets.sh after startup is absent. Nothing logs this, so the script checks and says so.
|
||||
# secrets.sh after startup is absent. Nothing logs this here, so the script checks and says
|
||||
# so — and since CB-594, bridged's own startup log says so too, by env var name.
|
||||
# 3. An old daemon that never actually died looks identical from the outside, so the script waits
|
||||
# for the process to exit and for the port to free before it starts a new one.
|
||||
# 4. "It started" is not "it works": the script polls /healthz until it answers, and reports the
|
||||
@@ -19,6 +20,13 @@
|
||||
# mismatch.
|
||||
# 5. Restarting under live members drops their tickets, so the script refuses unless you confirm
|
||||
# the fleet is drained.
|
||||
# 6. CB-594 — the launchd agent (deploy/dev.ltms.bridged.plist), if installed and loaded, is a
|
||||
# SECOND supervisor: its KeepAlive.SuccessfulExit=false restarts the daemon on any nonzero
|
||||
# exit, and a bare SIGTERM makes this JVM exit 143 even with its shutdown hook running to
|
||||
# completion (measured — see the CB-594 report). A plain `kill` here would race launchd's own
|
||||
# restart of the OLD jar. So this script detects whether the agent is loaded and, only then,
|
||||
# swaps `kill` + manual `nohup` for `launchctl unload`/`load` — the one supervisor in control
|
||||
# at any moment is whichever one you asked to act, never both.
|
||||
#
|
||||
# Usage:
|
||||
# scripts/redeploy-bridged.sh # build, confirm, restart, verify
|
||||
@@ -43,13 +51,17 @@ HEALTH='http://127.0.0.1:8765/healthz'
|
||||
STOP_WAIT=30 # seconds to wait for a clean exit before reporting failure
|
||||
HEALTH_WAIT=60 # seconds to wait for /healthz to answer after start
|
||||
|
||||
# CB-594: the launchd agent this script must not fight with (see trap 6 above).
|
||||
LAUNCHD_LABEL='dev.ltms.bridged'
|
||||
LAUNCHD_PLIST="$HOME/Library/LaunchAgents/$LAUNCHD_LABEL.plist"
|
||||
|
||||
DO_BUILD=1; ASSUME_YES=0; CHECK_ONLY=0
|
||||
for arg in "$@"; do
|
||||
case "$arg" in
|
||||
--yes|-y) ASSUME_YES=1 ;;
|
||||
--no-build) DO_BUILD=0 ;;
|
||||
--check) CHECK_ONLY=1 ;;
|
||||
-h|--help) sed -n '3,30p' "${BASH_SOURCE[0]}"; exit 0 ;;
|
||||
-h|--help) sed -n '3,37p' "${BASH_SOURCE[0]}"; exit 0 ;;
|
||||
*) echo "unknown option: $arg (try --help)" >&2; exit 2 ;;
|
||||
esac
|
||||
done
|
||||
@@ -61,6 +73,58 @@ die() { printf '\n FAIL %s\n\n' "$*" >&2; exit 1; }
|
||||
|
||||
jar_id() { [ -f "$JAR" ] && shasum -a 256 "$JAR" | cut -c1-12 || echo "absent"; }
|
||||
running_pid() { pgrep -f "$PATTERN" || true; }
|
||||
# `launchctl list <label>` exits 0 iff the label is loaded (registered with launchd) — true whether
|
||||
# or not it is currently running, which is exactly "supervision is active" for our purposes. Read-
|
||||
# only: neither helper below changes anything, so both are also safe under --check.
|
||||
launchd_installed() { [ -f "$LAUNCHD_PLIST" ]; }
|
||||
launchd_loaded() { launchctl list "$LAUNCHD_LABEL" >/dev/null 2>&1; }
|
||||
|
||||
# CB-600: the script computes its own log path from where it sits on disk (REPO, above); the
|
||||
# plist hard-codes an absolute StandardOutPath. Nothing forced the two to agree — if this script
|
||||
# were ever run from a checkout other than the one the loaded plist names, launchd would start and
|
||||
# log the daemon correctly, while every check below (the fresh "bridged listening" line, the
|
||||
# ERROR-count scan) would read a different, empty or stale file and the script would report a
|
||||
# clean restart while the daemon crash-loops. Pure and side-effect-free besides `die`/`ok` — reads
|
||||
# the two paths, resolves them, compares — so it never touches launchd or the daemon and can be
|
||||
# exercised by sourcing this script (see the SOURCED guard below) without installing the agent.
|
||||
check_log_path_matches_plist() {
|
||||
local script_out="$1" plist_path="$2"
|
||||
local plist_out resolved_out resolved_plist_out
|
||||
# Checked by exit status, not by emptiness: on a missing file/key PlistBuddy exits nonzero but
|
||||
# still writes a message ("File Doesn't Exist, Will Create: ...") that command substitution
|
||||
# would happily capture as if it were the real value — testing only `-z` missed that case.
|
||||
if ! plist_out="$(/usr/libexec/PlistBuddy -c 'Print :StandardOutPath' "$plist_path" 2>/dev/null)" \
|
||||
|| [ -z "$plist_out" ]; then
|
||||
die "launchd agent is loaded but PlistBuddy could not read StandardOutPath from
|
||||
$plist_path
|
||||
— cannot verify the daemon logs where this script is about to look. Fix the plist before
|
||||
redeploying supervised."
|
||||
fi
|
||||
resolved_out="$(cd "$(dirname "$script_out")" 2>/dev/null && pwd -P)/$(basename "$script_out")" || true
|
||||
resolved_plist_out="$(cd "$(dirname "$plist_out")" 2>/dev/null && pwd -P)/$(basename "$plist_out")" || true
|
||||
if [ -z "$resolved_out" ] || [ -z "$resolved_plist_out" ] || [ "$resolved_out" != "$resolved_plist_out" ]; then
|
||||
die "log path mismatch — this script reads
|
||||
$script_out (resolved: ${resolved_out:-<directory does not exist>})
|
||||
but the loaded plist's StandardOutPath is
|
||||
$plist_out (resolved: ${resolved_plist_out:-<directory does not exist>})
|
||||
Under supervision the daemon writes to the PLIST's path, not necessarily this script's — every
|
||||
post-restart check below (the fresh 'bridged listening' line, the ERROR-count scan) would read
|
||||
the wrong file and could report a clean restart while the daemon crash-loops. Fix the mismatch
|
||||
(move this checkout to match the plist, or edit the plist's StandardOutPath/StandardErrorPath)
|
||||
before redeploying supervised."
|
||||
fi
|
||||
ok "log path check: script and plist agree ($resolved_out)"
|
||||
}
|
||||
|
||||
# CB-600: sourceable for testing. When this file is SOURCED (not executed) it stops here — nothing
|
||||
# below runs — so a test harness can `source` it to call check_log_path_matches_plist (or the
|
||||
# other pure helpers above) against a throwaway plist fixture without ever reaching the mutating
|
||||
# flow (build/stop/start) or touching the real daemon or launchd. On a normal `./redeploy-bridged.sh`
|
||||
# invocation `(return 0 2>/dev/null)` fails (return is illegal at top level of an executed script),
|
||||
# so this whole block is a no-op and every line below still runs exactly as before.
|
||||
if (return 0 2>/dev/null); then
|
||||
return 0
|
||||
fi
|
||||
|
||||
# ---------------------------------------------------------------- report state
|
||||
|
||||
@@ -74,6 +138,25 @@ fi
|
||||
ok "jar on disk: $(jar_id) ($([ -f "$JAR" ] && date -r "$JAR" '+%Y-%m-%d %H:%M:%S' || echo 'none'))"
|
||||
ok "HEAD: $(git -C "$REPO" log --oneline -1)"
|
||||
|
||||
# CB-594: supervision state. Installed and loaded are different facts — a copied-but-never-loaded
|
||||
# plist supervises nothing, and a loaded label with no file backing it (rare, but possible after an
|
||||
# edited/moved plist) is still what launchd will act on.
|
||||
if launchd_installed; then
|
||||
ok "launchd agent installed: $LAUNCHD_PLIST"
|
||||
else
|
||||
warn "launchd agent NOT installed (no supervision — a crash will not restart the daemon)."
|
||||
fi
|
||||
SUPERVISED=0
|
||||
if launchd_loaded; then
|
||||
SUPERVISED=1
|
||||
ok "launchd agent loaded ($LAUNCHD_LABEL) — launchd supervises this daemon"
|
||||
# CB-600: fail loudly here, before ANY other check runs, if this script and the loaded plist
|
||||
# would read different log files — every check after this point is worthless otherwise.
|
||||
check_log_path_matches_plist "$OUT" "$LAUNCHD_PLIST"
|
||||
else
|
||||
warn "launchd agent not loaded — this script is the only thing that will restart the daemon."
|
||||
fi
|
||||
|
||||
# The trap with no log line. Checked in a LOGIN shell, because that is how the daemon is started
|
||||
# below. Never prints the value — only whether it resolved.
|
||||
if zsh -lc '[ -n "${WORKER_GITEA_TOKEN:-}" ]' 2>/dev/null; then
|
||||
@@ -137,11 +220,26 @@ if [ -n "$OLD_PID" ] && [ "$ASSUME_YES" = 0 ]; then
|
||||
fi
|
||||
|
||||
# ------------------------------------------------------------------ stop
|
||||
#
|
||||
# CB-594: when SUPERVISED, launchd owns the stop — never a raw `kill` here. A bare SIGTERM makes
|
||||
# this JVM exit 143 even with its shutdown hook running to completion (verified separately: a
|
||||
# throwaway Java process with an equivalent shutdown hook, sent SIGTERM from a login shell that
|
||||
# could `wait` on it directly, reported exit code 143 every time — never 0). launchd's
|
||||
# KeepAlive.SuccessfulExit=false treats any nonzero exit as a crash and restarts the OLD jar,
|
||||
# which would race this script's own restart of the NEW one. `launchctl unload` avoids that race
|
||||
# by deregistering the job first, so no KeepAlive is left armed when the process actually stops.
|
||||
|
||||
if [ -n "$OLD_PID" ]; then
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)" # verify a FRESH line appears later
|
||||
kill "$OLD_PID"
|
||||
if [ "$SUPERVISED" = 1 ]; then
|
||||
echo " supervision is ON: using 'launchctl unload' (not kill) so launchd's own KeepAlive"
|
||||
echo " cannot restart the OLD jar out from under this script — see the CB-594 comment above."
|
||||
launchctl unload -w "$LAUNCHD_PLIST" \
|
||||
|| die "launchctl unload failed — the daemon may still be under supervision; investigate before retrying"
|
||||
else
|
||||
kill "$OLD_PID"
|
||||
fi
|
||||
for _ in $(seq "$STOP_WAIT"); do
|
||||
[ -z "$(running_pid)" ] && break
|
||||
sleep 1
|
||||
@@ -152,18 +250,48 @@ if [ -n "$OLD_PID" ]; then
|
||||
leave worktrees and panes behind. Investigate, then kill -9 by hand if you accept that."
|
||||
fi
|
||||
ok "pid $OLD_PID exited"
|
||||
elif [ "$SUPERVISED" = 1 ]; then
|
||||
# Loaded but not currently running (e.g. throttled after a crash loop). Unload it anyway so the
|
||||
# start step below does a clean load, never a load stacked on an already-loaded label.
|
||||
say "stop"
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
launchctl unload -w "$LAUNCHD_PLIST" 2>/dev/null || true
|
||||
ok "launchd agent unloaded (was already not running)"
|
||||
else
|
||||
RESTART_MARK="$(wc -l < "$OUT" 2>/dev/null || echo 0)"
|
||||
fi
|
||||
|
||||
# ------------------------------------------------------------------ start
|
||||
# Login shell (zsh -l) is what puts the secrets on the daemon's environment. cwd must be bridged/
|
||||
# because the daemon resolves bridged.yaml, logs/ and target/ relative to it.
|
||||
# Unsupervised: login shell (zsh -l) is what puts the secrets on the daemon's environment, and cwd
|
||||
# must be bridged/ because the daemon resolves bridged.yaml, logs/ and target/ relative to it.
|
||||
# Supervised: launchd does both — deploy/dev.ltms.bridged.plist points ProgramArguments at
|
||||
# scripts/bridged-launchd-wrapper.sh (CB-594), which is what execs the login shell in launchd's
|
||||
# place, and WorkingDirectory in the plist already pins bridged/.
|
||||
|
||||
say "start"
|
||||
# Absolute jar path so `ps` names which checkout is running; cwd still bridged/ because the daemon
|
||||
# resolves bridged.yaml, logs/ and target/ relative to it.
|
||||
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )
|
||||
if [ "$SUPERVISED" = 1 ]; then
|
||||
echo " supervision is ON: using 'launchctl load' so launchd starts and keeps supervising this"
|
||||
echo " process, instead of a manual nohup that launchd would know nothing about."
|
||||
# CB-600: 'launchctl unload -w' above already persisted Disabled=true for this label. A load -w
|
||||
# that succeeds clears it; a load -w that FAILS leaves the agent both stopped and disabled — worse
|
||||
# than before this script ran, because a later reboot or login will not bring it back either. One
|
||||
# retry covers a transient race (e.g. launchd not yet fully done deregistering); if it still fails,
|
||||
# die with the exact recovery command rather than a bare "failed".
|
||||
if ! launchctl load -w "$LAUNCHD_PLIST" 2>/dev/null; then
|
||||
warn "launchctl load failed on the first attempt — retrying once after a short pause"
|
||||
sleep 2
|
||||
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed twice.
|
||||
The agent is now STOPPED and DISABLED — it will NOT come back on its own, not even after a
|
||||
reboot or login, because 'launchctl unload -w' above persisted Disabled=true and load -w
|
||||
never got the chance to clear it. Recover with:
|
||||
launchctl load -w \"$LAUNCHD_PLIST\"
|
||||
If that still fails, check 'launchctl list $LAUNCHD_LABEL', validate the plist with
|
||||
'plutil -lint \"$LAUNCHD_PLIST\"', and check $OUT before assuming a retry will succeed."
|
||||
fi
|
||||
else
|
||||
# Absolute jar path so `ps` names which checkout is running.
|
||||
( cd "$BRIDGED" && zsh -lc "nohup java -jar '$JAR' >> bridged.out 2>&1 &" )
|
||||
fi
|
||||
|
||||
for _ in $(seq 10); do
|
||||
NEW_PID="$(running_pid)"
|
||||
|
||||
Reference in New Issue
Block a user