Compare commits
1 Commits
| Author | SHA1 | Date | |
|---|---|---|---|
| 3ba6d6784c |
@@ -434,14 +434,10 @@ 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.
|
||||
# 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.
|
||||
# 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.
|
||||
# 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,9 +352,8 @@ 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(), cfg.broker().prefetchOrDefault());
|
||||
log.info("reply inbox: AMQP broker (durable) at {} (prefetch={})",
|
||||
cfg.broker().uri(), cfg.broker().prefetchOrDefault());
|
||||
replyInbox = AmqpReplyInbox.open(cfg.broker().uri());
|
||||
log.info("reply inbox: AMQP broker (durable) at {}", cfg.broker().uri());
|
||||
} else {
|
||||
replyInbox = new InMemoryReplyInbox();
|
||||
log.info("reply inbox: in-memory (soft-state)");
|
||||
@@ -544,6 +547,63 @@ 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, a few lines above this method's call site.
|
||||
*
|
||||
* <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,7 +5,6 @@ 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;
|
||||
@@ -533,24 +532,16 @@ 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 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}.
|
||||
* @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).
|
||||
*/
|
||||
@JsonIgnoreProperties(ignoreUnknown = true)
|
||||
public record Broker(String uri, Integer prefetch) {
|
||||
public record Broker(String uri) {
|
||||
|
||||
/** 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;
|
||||
}
|
||||
}
|
||||
|
||||
/**
|
||||
|
||||
@@ -7,7 +7,6 @@ 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;
|
||||
|
||||
@@ -15,13 +14,7 @@ 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
|
||||
@@ -36,27 +29,6 @@ import java.util.concurrent.TimeoutException;
|
||||
* 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
|
||||
@@ -82,12 +54,6 @@ 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. */
|
||||
@@ -97,81 +63,39 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
/** Targets whose queue is declared and consumer is running, mapped to their broker consumer tag. */
|
||||
private final ConcurrentHashMap<String, String> consumerTags = new ConcurrentHashMap<>();
|
||||
|
||||
/**
|
||||
* CB-528: a dedicated channel for {@link #publish}, kept separate from {@link #channel} (consume
|
||||
* + ack) so a publish confirm round trip never blocks under {@link #channelLock} and stalls an ack.
|
||||
*/
|
||||
private final Channel publishChannel;
|
||||
private final Object publishChannelLock = new Object();
|
||||
/** In-flight publishes awaiting their confirm, keyed by the publish channel's sequence number. */
|
||||
private final ConcurrentSkipListMap<Long, Pending> pendingBySeq = new ConcurrentSkipListMap<>();
|
||||
/** The same in-flight publishes, keyed by {@code msgId} — a broker {@code Return} carries no delivery tag. */
|
||||
private final ConcurrentHashMap<String, Pending> pendingByMsgId = new ConcurrentHashMap<>();
|
||||
|
||||
/** A message pulled off the broker but not yet acked: its delivery-tag plus the port payload. */
|
||||
private record Held(long deliveryTag, InboxMessage message) {}
|
||||
|
||||
/** 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}. */
|
||||
/** Connect to {@code uri} (e.g. {@code amqp://guest:guest@127.0.0.1:5672/}) and open the inbox. */
|
||||
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"), prefetch);
|
||||
return new AmqpReplyInbox(factory.newConnection("bridged-reply-inbox"));
|
||||
} catch (Exception e) {
|
||||
throw new IllegalStateException("cannot connect to AMQP broker at " + uri, e);
|
||||
}
|
||||
}
|
||||
|
||||
/** Wrap an already-open connection with {@link #DEFAULT_PREFETCH} (injection seam for the contract test). */
|
||||
/** Wrap an already-open connection (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). 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.
|
||||
// repopulates it with valid tags (dedup by msgId still prevents any double-queue).
|
||||
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");
|
||||
}
|
||||
|
||||
@@ -217,12 +141,6 @@ 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()
|
||||
@@ -230,35 +148,12 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
.deliveryMode(2) // persistent — survives a broker restart
|
||||
.contentType("text/plain")
|
||||
.build();
|
||||
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);
|
||||
}
|
||||
}
|
||||
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);
|
||||
synchronized (channelLock) {
|
||||
channel.basicPublish("", queueName(target), props, content.getBytes(StandardCharsets.UTF_8));
|
||||
}
|
||||
} catch (IOException e) {
|
||||
throw new IllegalStateException("cannot publish reply to " + queueName(target), e);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -327,65 +222,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
};
|
||||
}
|
||||
|
||||
/** Broker return for an unroutable {@code mandatory} publish — arrives BEFORE its confirm (CB-528). */
|
||||
private void onReturn(Return r) {
|
||||
String msgId = r.getProperties() == null ? null : r.getProperties().getMessageId();
|
||||
Pending pending = msgId == null ? null : pendingByMsgId.get(msgId);
|
||||
if (pending != null) {
|
||||
pending.returned = true;
|
||||
} else {
|
||||
log.warn("AMQP return for reply {} (routingKey={}, {} {}) with no matching in-flight publish"
|
||||
+ " — already resolved by a prior confirm", msgId, r.getRoutingKey(), r.getReplyCode(),
|
||||
r.getReplyText());
|
||||
}
|
||||
}
|
||||
|
||||
private void onAck(long seq, boolean multiple) {
|
||||
resolveConfirm(seq, multiple, true);
|
||||
}
|
||||
|
||||
private void onNack(long seq, boolean multiple) {
|
||||
resolveConfirm(seq, multiple, false);
|
||||
}
|
||||
|
||||
/**
|
||||
* Resolve every pending publish covered by this confirm (a single seq, or — {@code multiple} —
|
||||
* every seq up to and including it). Checks {@link Pending#returned} at confirm time: since the
|
||||
* broker's return for an unroutable message always precedes its confirm, an ack that arrives after
|
||||
* a return means "confirmed but never routed", not "durably queued".
|
||||
*/
|
||||
private void resolveConfirm(long seq, boolean multiple, boolean ack) {
|
||||
NavigableMap<Long, Pending> covered = multiple
|
||||
? pendingBySeq.headMap(seq, true)
|
||||
: pendingBySeq.subMap(seq, true, seq, true);
|
||||
for (var it = covered.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
if (ack && !pending.returned) {
|
||||
pending.confirmed.complete(null);
|
||||
} else if (ack) {
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"reply " + pending.msgId + " was returned as unroutable (queue not declared/owned)"));
|
||||
} else {
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"broker nacked publish of reply " + pending.msgId));
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Fail every publish still awaiting its confirm — their sequence numbers are stale after recovery. */
|
||||
private void failPendingPublishesOnRecovery() {
|
||||
for (var it = pendingBySeq.entrySet().iterator(); it.hasNext(); ) {
|
||||
Pending pending = it.next().getValue();
|
||||
it.remove();
|
||||
pendingByMsgId.remove(pending.msgId, pending);
|
||||
pending.confirmed.completeExceptionally(new IllegalStateException(
|
||||
"AMQP connection recovered mid-publish; confirm status of reply " + pending.msgId
|
||||
+ " is unknown"));
|
||||
}
|
||||
}
|
||||
|
||||
private static String queueName(String target) {
|
||||
return QUEUE_PREFIX + target + QUEUE_SUFFIX;
|
||||
}
|
||||
@@ -397,11 +233,6 @@ public final class AmqpReplyInbox implements ReplyInbox, AutoCloseable {
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP channel close: {}", e.toString());
|
||||
}
|
||||
try {
|
||||
publishChannel.close();
|
||||
} catch (Exception e) {
|
||||
log.debug("AMQP publish channel close: {}", e.toString());
|
||||
}
|
||||
try {
|
||||
connection.close();
|
||||
} catch (Exception e) {
|
||||
|
||||
@@ -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,7 +16,6 @@ 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;
|
||||
|
||||
/**
|
||||
@@ -167,88 +166,6 @@ 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)
|
||||
|
||||
@@ -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,11 +68,14 @@
|
||||
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>
|
||||
|
||||
@@ -69,10 +92,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 "$@"' -- "$@"
|
||||
@@ -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,11 @@ 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; }
|
||||
|
||||
# ---------------------------------------------------------------- report state
|
||||
|
||||
@@ -74,6 +91,22 @@ 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"
|
||||
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 +170,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 +200,33 @@ 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."
|
||||
launchctl load -w "$LAUNCHD_PLIST" || die "launchctl load failed"
|
||||
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