diff --git a/deploy/agent/codex/Dockerfile b/deploy/agent/codex/Dockerfile index b9b2b5c9..ca27fe48 100644 --- a/deploy/agent/codex/Dockerfile +++ b/deploy/agent/codex/Dockerfile @@ -17,8 +17,11 @@ RUN apt-get update -qq && apt-get install -y -qq --no-install-recommends \ ca-certificates git openssh-client \ && rm -rf /var/lib/apt/lists/* -# Pinned to the version the adapter's event parser was measured against (RUN-TOPOLOGY §10). -RUN npm i -g @openai/codex@0.146.0 && npm cache clean --force +# Pinned. The event parser was first measured on 0.146.0 (RUN-TOPOLOGY §10); raised to 0.156.1 on +# 2026-09-23 because the model list is read from this CLI, and 0.146.0 predates the gpt-6 models. Every +# flag the adapter uses was re-checked on 0.156.1; the --json event shape is re-proved by the next paid +# run (UNVERIFIED section B). +RUN npm i -g @openai/codex@0.156.1 && npm cache clean --force # The mount points exist, owned by the agent user, so a fresh named volume mounted there inherits # that ownership; without this the volume is root's and the agent cannot write its own workspace. diff --git a/docs/UNVERIFIED.md b/docs/UNVERIFIED.md index 49d0459e..5d858da2 100644 --- a/docs/UNVERIFIED.md +++ b/docs/UNVERIFIED.md @@ -413,13 +413,16 @@ so in the outage this exists for it fails too. What the operator sees depends on WHEN the outage began, and the two are not the same: - **After the prompt was stored** — the row is PROMPTED and carries an expiry, so the screen counts down - and then shows a code that has plainly run out. -- **Before it** — the prompt is the only thing that ever writes {@code expires_at}, and the worker does - not check whether that send arrived either. The row stays PENDING with **no countdown at all**: the - screen says "starting the sign-in", indefinitely. + and then shows a code that has plainly run out. Two minutes past that expiry the orchestrator closes the + row itself as expired (`HarnessSignIns.resendUnclaimed`, 2026-09-23). +- **Before it** — the row stays PENDING. Since 2026-09-23 (review of PR #168) the orchestrator re-sends + an unanswered start every 30 seconds, and closes the row as + `sign_in_not_started` if no worker has shown a code within six minutes (a four-minute start window, carried in the start itself, plus two for the code to cross the bus). A worker that cannot + start a unit reports nothing, because another worker may hold it. So the screen no longer says "starting the sign-in" for ever; it ends within + the wait. -Neither resolves itself. Cancelling is the way out, and cancelling publishes before it changes the row, -so it too needs the broker back; until then the screen shows only the cancel's own error. +Both now end on their own, without the broker. What is still lost is the credential itself: a finished +sign-in whose result cannot be delivered is discarded, and the operator signs in again. **Evidence needed.** None — this is a deliberate trade, not a suspicion. Closing it means a durable terminal record the screen can read without the broker, which is the same transactional-outbox treatment @@ -468,6 +471,7 @@ Each has a runbook mode. None has been run by an operator. | The whole M1 lifecycle against a real forge | **Mode Q** | Cancel, steer, the watchdog, the push gate and the charge ledger have only ever met a WireMock LLM and a local origin | | Corporate-only bundle → the failure it produces | Mode R §5 | The documented trap (internal forge works, model API fails) is asserted nowhere; it is the mistake an operator will actually make | | A private-registry pull | Mode S §4 | Nothing pulls from a private registry in any test. `authFor` and the attachment are unit-tested; the *pull* is not | +| **Codex CLI 0.156.1 in the agent image** (2026-09-23) | none yet | Raised from 0.146.0 so the model list includes the gpt-6 models. Re-checked on 0.156.1: every flag the adapter passes, the API-key login, and the shape of the file it writes (`auth_mode=apikey`). NOT re-checked: the `--json` event stream the usage parser reads, and the device sign-in output. Both need a paid run or a real sign-in, and the first of each proves or breaks them | | **Runs pinned to the image their model list came from** (M3.5 part M, 2026-09-23) | none yet | Choosing the pin is unit-tested against given daemon answers, and one real-daemon test pins a LOCAL build by its image id. No test pulls a registry image and pins it by its registry digest, and none runs two workers. So "two workers holding different images under one tag run the same one" is argued from the code, not watched. A local-only image is pinned by an id that exists on one daemon only: on a second worker such a run fails to pull — by design, but unobserved | | **OIDC sessions actually renew instead of re-authenticating** | **Mode J check 11** (2026-09-10) | The bug it fixes needs a real browser, a real Keycloak and **fifteen elapsed minutes**. No suite here has any of the three: there are zero WebSocket client tests, and nothing observes a token reaching its `exp`. `OidcSessionsAreRenewedTest` asserts the four `application.yml` files *say* renewal is on — it cannot assert Quarkus *does* it | diff --git a/spire-contract/src/main/java/dev/codespire/contract/command/HarnessImageCommand.java b/spire-contract/src/main/java/dev/codespire/contract/command/HarnessImageCommand.java index ee4aa576..73f17f53 100644 --- a/spire-contract/src/main/java/dev/codespire/contract/command/HarnessImageCommand.java +++ b/spire-contract/src/main/java/dev/codespire/contract/command/HarnessImageCommand.java @@ -26,8 +26,18 @@ public sealed interface HarnessImageCommand { * @param harness the name the orchestrator dispatches under, echoed back so the answer files itself * @param image the exact reference the orchestrator will run — the answer describes THIS image, and a * different tag of the same repository may carry a different CLI and a different list + * @param askedAt when the orchestrator asked. The channels replay from their oldest record, so a + * question can arrive long after it was asked: the worker skips one too old to matter, and the + * answer carries this back so an older answer cannot replace a newer one (review of PR #168). + * Null in a question sent before this existed. */ - record Describe(String requestId, String harness, String image) implements HarnessImageCommand { + record Describe(String requestId, String harness, String image, java.time.Instant askedAt) + implements HarnessImageCommand { + + public Describe(String requestId, String harness, String image) { + this(requestId, harness, image, null); + } + public Describe { if (requestId == null || requestId.isBlank()) throw new IllegalArgumentException("A request id is required"); if (harness == null || harness.isBlank()) throw new IllegalArgumentException("A harness is required"); diff --git a/spire-contract/src/main/java/dev/codespire/contract/command/HarnessSignInCommand.java b/spire-contract/src/main/java/dev/codespire/contract/command/HarnessSignInCommand.java index d74931ea..869c811e 100644 --- a/spire-contract/src/main/java/dev/codespire/contract/command/HarnessSignInCommand.java +++ b/spire-contract/src/main/java/dev/codespire/contract/command/HarnessSignInCommand.java @@ -41,14 +41,57 @@ public sealed interface HarnessSignInCommand { * @param maxWaitSeconds how long the unit may wait for the operator. The vendor states a 15-minute * code lifetime; this is the deployment's own ceiling, so a unit cannot outlive the code it is * waiting on and sit holding a container for ever. + * @param requestedAt when the operator pressed start. The wait is counted from HERE, not from delivery: + * the channel replays from its oldest record and the orchestrator re-sends a start nobody picked + * up, so a start can arrive late or twice. One whose wait has already run out opens no unit, and + * a late one gets only the time left (review of PR #168). Null in a start sent before this existed. + * @param startWithinSeconds how long after the press a unit may still be OPENED, as distinct from how + * long a person may take to approve ({@code maxWaitSeconds}). The orchestrator ends a sign-in that + * shows no code a little after this, so a worker must not open one it could only prompt for once + * that has happened. Carried rather than configured twice, so the two sides cannot drift. Zero in + * a start sent before this existed, which then opens nothing. */ - record Start(String signInId, String harness, String image, long maxWaitSeconds) + record Start(String signInId, String harness, String image, long maxWaitSeconds, java.time.Instant requestedAt, + long startWithinSeconds) implements HarnessSignInCommand { + public Start(String signInId, String harness, String image, long maxWaitSeconds) { + this(signInId, harness, image, maxWaitSeconds, null, 0); + } + + /** When the code must be on screen by; after this the orchestrator ends the row. Null if unknown. */ + public java.time.Instant promptDeadline() { + return requestedAt == null ? null : requestedAt.plusSeconds(startWithinSeconds); + } + + /** When the person's time to approve runs out, counted from the press. Null if unknown. */ + public java.time.Instant approvalDeadline() { + return requestedAt == null ? null : requestedAt.plusSeconds(maxWaitSeconds); + } + + /** + * Whether a unit opened now could still show its code in time: the code must be on screen by the + * end of the start window, and the unit may take {@code toPrompt} to print it. + */ + public boolean mayOpenAt(java.time.Instant now, java.time.Duration toPrompt) { + if (requestedAt == null || startWithinSeconds <= 0) return false; + return !now.plus(toPrompt).isAfter(requestedAt.plusSeconds(startWithinSeconds)); + } + + /** How long the unit may still wait, counted from the request; the full wait when that is unknown. */ + public java.time.Duration remainingWait(java.time.Instant now) { + java.time.Duration full = java.time.Duration.ofSeconds(maxWaitSeconds); + if (requestedAt == null) return full; + java.time.Duration left = java.time.Duration.between(now, requestedAt.plus(full)); + return left.isNegative() ? java.time.Duration.ZERO : (left.compareTo(full) > 0 ? full : left); + } + public Start { if (signInId == null || signInId.isBlank()) throw new IllegalArgumentException("A sign-in id is required"); if (harness == null || harness.isBlank()) throw new IllegalArgumentException("A harness name is required"); if (image == null || image.isBlank()) throw new IllegalArgumentException("An agent image is required"); + if (startWithinSeconds < 0 || startWithinSeconds > maxWaitSeconds) throw new IllegalArgumentException( + "A start window of " + startWithinSeconds + " seconds must lie within the wait of " + maxWaitSeconds); if (maxWaitSeconds <= 0) throw new IllegalArgumentException( "A sign-in wait of " + maxWaitSeconds + " seconds would end the unit before the operator" + " could read the code; the wait is a ceiling, not a switch"); diff --git a/spire-contract/src/main/java/dev/codespire/contract/event/HarnessImageResult.java b/spire-contract/src/main/java/dev/codespire/contract/event/HarnessImageResult.java index 3ff54b70..b867ecf2 100644 --- a/spire-contract/src/main/java/dev/codespire/contract/event/HarnessImageResult.java +++ b/spire-contract/src/main/java/dev/codespire/contract/event/HarnessImageResult.java @@ -68,13 +68,21 @@ enum Status { * when it could not be reached. A run of this harness uses it rather than the tag, * so it runs the image these models were read from (review of PR #167). Null in an * answer sent before pins existed. + * @param askedAt when the question this answers was asked, echoed back. The cache keeps the answer to + * the NEWEST question, so a replayed older answer cannot roll it back. Null in an answer + * sent before this existed. */ record Described(String requestId, String harness, String image, Status status, List models, - String pinnedImage) + String pinnedImage, java.time.Instant askedAt) implements HarnessImageResult { public Described(String requestId, String harness, String image, Status status, List models) { - this(requestId, harness, image, status, models, null); + this(requestId, harness, image, status, models, null, null); + } + + public Described(String requestId, String harness, String image, Status status, List models, + String pinnedImage) { + this(requestId, harness, image, status, models, pinnedImage, null); } public Described { diff --git a/spire-contract/src/main/java/dev/codespire/contract/event/HarnessSignInResult.java b/spire-contract/src/main/java/dev/codespire/contract/event/HarnessSignInResult.java index 7d4b64d8..f97b5dad 100644 --- a/spire-contract/src/main/java/dev/codespire/contract/event/HarnessSignInResult.java +++ b/spire-contract/src/main/java/dev/codespire/contract/event/HarnessSignInResult.java @@ -105,6 +105,8 @@ record Failed(String signInId, String cause, String detail) implements HarnessSi public static final String UNIT_FAILED = "sign_in_unit_failed"; /** The CLI signed in, but as an API key rather than a subscription. */ public static final String WRONG_MODE = "sign_in_wrong_mode"; + /** No worker picked the sign-in up before its wait ran out, re-sends included. */ + public static final String NOT_STARTED = "sign_in_not_started"; public Failed { if (signInId == null || signInId.isBlank()) throw new IllegalArgumentException("A sign-in id is required"); diff --git a/spire-contract/src/test/java/dev/codespire/contract/command/HarnessSignInStartTest.java b/spire-contract/src/test/java/dev/codespire/contract/command/HarnessSignInStartTest.java new file mode 100644 index 00000000..708eeca9 --- /dev/null +++ b/spire-contract/src/test/java/dev/codespire/contract/command/HarnessSignInStartTest.java @@ -0,0 +1,64 @@ +package dev.codespire.contract.command; + +import org.junit.jupiter.api.Test; + +import java.time.Duration; +import java.time.Instant; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * A start's wait is counted from the operator's press, not from delivery, because a start can be + * replayed or re-sent (review of PR #168). + */ +class HarnessSignInStartTest { + + private static final Instant PRESSED = Instant.parse("2026-09-23T07:00:00Z"); + + private static HarnessSignInCommand.Start start(Instant requestedAt) { + return new HarnessSignInCommand.Start("TEST-sign-in", "codex", "TEST-image", 840, requestedAt, 240); + } + + /** A unit may be opened only if it can print its code before the start window closes. */ + @Test + void aUnitMayBeOpenedOnlyWhileItsCodeCanStillBeShownInTime() { + Duration toPrompt = Duration.ofSeconds(60); + assertEquals(true, start(PRESSED).mayOpenAt(PRESSED.plusSeconds(180), toPrompt)); + assertEquals(false, start(PRESSED).mayOpenAt(PRESSED.plusSeconds(181), toPrompt), + "the wait is far from over, but the code would land after the window"); + } + + @Test + void theDeadlinesAreCountedFromThePress() { + assertEquals(PRESSED.plusSeconds(240), start(PRESSED).promptDeadline()); + assertEquals(PRESSED.plusSeconds(840), start(PRESSED).approvalDeadline()); + } + + @Test + void aStartWithNoWindowOrNoPressOpensNothing() { + assertEquals(false, new HarnessSignInCommand.Start("TEST-sign-in", "codex", "TEST-image", 840, PRESSED, 0) + .mayOpenAt(PRESSED, Duration.ZERO)); + assertEquals(false, start(null).mayOpenAt(PRESSED, Duration.ZERO)); + } + + @Test + void aLateStartGetsOnlyWhatIsLeft() { + assertEquals(Duration.ofSeconds(240), start(PRESSED).remainingWait(PRESSED.plusSeconds(600))); + } + + @Test + void aStartWhoseWaitIsGoneGetsNothing() { + assertEquals(Duration.ZERO, start(PRESSED).remainingWait(PRESSED.plusSeconds(900))); + } + + /** A clock behind the orchestrator's must not grant more than the full wait. */ + @Test + void neverMoreThanTheFullWait() { + assertEquals(Duration.ofSeconds(840), start(PRESSED).remainingWait(PRESSED.minusSeconds(60))); + } + + @Test + void aStartFromBeforeTimesExistedGetsTheFullWait() { + assertEquals(Duration.ofSeconds(840), start(null).remainingWait(PRESSED)); + } +} diff --git a/spire-orchestrator/src/main/java/dev/codespire/orchestrator/dlq/DlqTopics.java b/spire-orchestrator/src/main/java/dev/codespire/orchestrator/dlq/DlqTopics.java index f66ed544..e5226ce6 100644 --- a/spire-orchestrator/src/main/java/dev/codespire/orchestrator/dlq/DlqTopics.java +++ b/spire-orchestrator/src/main/java/dev/codespire/orchestrator/dlq/DlqTopics.java @@ -49,6 +49,19 @@ final class DlqTopics { private static final Set RUN_RESULT_TYPES = Set.of("RunStarted", "RunFinished", "RunFailed"); + /** + * The harness side channels (M3.5 parts F and M). Without these a dead-lettered model-list or + * sign-in record was replayed onto {@code cs.commands}, where nothing reads it (review of PR #168). + */ + static final String HARNESS_IMAGE_COMMANDS = "cs.harness-image-commands"; + static final String HARNESS_IMAGE_RESULTS = "cs.harness-image-results"; + static final String HARNESS_SIGN_IN_COMMANDS = "cs.harness-sign-in-commands"; + static final String HARNESS_SIGN_IN_RESULTS = "cs.harness-sign-in-results"; + + private static final Set HARNESS_SIGN_IN_COMMAND_TYPES = Set.of("Start", "Cancel"); + + private static final Set HARNESS_SIGN_IN_RESULT_TYPES = Set.of("Prompted", "Completed", "Failed"); + private DlqTopics() { } @@ -74,6 +87,10 @@ static String forType(String type) { if (RUN_RESULT_TYPES.contains(type)) { return RUN_RESULTS; } + if ("Describe".equals(type)) return HARNESS_IMAGE_COMMANDS; + if ("Described".equals(type)) return HARNESS_IMAGE_RESULTS; + if (HARNESS_SIGN_IN_COMMAND_TYPES.contains(type)) return HARNESS_SIGN_IN_COMMANDS; + if (HARNESS_SIGN_IN_RESULT_TYPES.contains(type)) return HARNESS_SIGN_IN_RESULTS; return COMMANDS; } } diff --git a/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessCatalogues.java b/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessCatalogues.java index 5de18faa..0d1062ff 100644 --- a/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessCatalogues.java +++ b/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessCatalogues.java @@ -81,13 +81,16 @@ void onStart(@Observes StartupEvent event) { * because a worker that was down at startup would otherwise leave the cache empty until the next * orchestrator restart. The answer is cheap: a label read, and a pull only the first time. */ - @Scheduled(every = "${spire.harness-catalogue-interval:10m}", + // Delayed: the startup ask already covers boot, and a first tick at startup ran before the Kafka + // emitter was connected, logging an injection ERROR on every start (dev stack, 2026-09-23). + @Scheduled(every = "${spire.harness-catalogue-interval:10m}", delayed = "1m", concurrentExecution = Scheduled.ConcurrentExecution.SKIP) void refresh() { for (Map.Entry harness : config.agentImage().entrySet()) { try { KafkaSends.sendAndAwait(commands, harness.getKey(), - new HarnessImageCommand.Describe(UUID.randomUUID().toString(), harness.getKey(), harness.getValue()), + new HarnessImageCommand.Describe(UUID.randomUUID().toString(), harness.getKey(), harness.getValue(), + Instant.now()), "describe the image for " + harness.getKey()); } catch (RuntimeException undelivered) { // The next interval asks again; one lost question is not worth failing startup over. @@ -110,17 +113,27 @@ public void record(HarnessImageResult.Described answer) { return; } try (Connection c = dataSource.getConnection(); PreparedStatement ps = c.prepareStatement(""" - INSERT INTO harness_catalogue (harness, image, status, models, observed_at, pinned_image) - VALUES (?, ?, ?, ?::jsonb, now(), ?) + INSERT INTO harness_catalogue (harness, image, status, models, observed_at, pinned_image, asked_at) + VALUES (?, ?, ?, ?::jsonb, now(), ?, ?) ON CONFLICT (harness) DO UPDATE SET image=excluded.image, status=excluded.status, models=excluded.models, observed_at=now(), - pinned_image=excluded.pinned_image + pinned_image=excluded.pinned_image, asked_at=excluded.asked_at + -- The answer to the NEWEST question wins, not the last one to arrive: the channel + -- replays from its oldest record, and an older answer must not roll the list and its + -- pin back (review of PR #168). An answer with no time predates times, so it can only + -- be a replay and replaces nothing. A row about an image no longer configured is stale + -- whatever its time, so any timed current answer replaces it. + WHERE excluded.asked_at IS NOT NULL + AND (harness_catalogue.image <> excluded.image + OR harness_catalogue.asked_at IS NULL + OR excluded.asked_at >= harness_catalogue.asked_at) """)) { ps.setString(1, answer.harness()); ps.setString(2, answer.image()); ps.setString(3, answer.status().name()); ps.setString(4, mapper.writeValueAsString(answer.models())); ps.setString(5, answer.pinnedImage()); + ps.setTimestamp(6, answer.askedAt() == null ? null : java.sql.Timestamp.from(answer.askedAt())); ps.executeUpdate(); } catch (SQLException | JsonProcessingException failure) { throw new IllegalStateException("The model catalogue for " + answer.harness() + " could not be stored", failure); diff --git a/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessSignIns.java b/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessSignIns.java index f4794c94..9f1b3104 100644 --- a/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessSignIns.java +++ b/spire-orchestrator/src/main/java/dev/codespire/orchestrator/factory/HarnessSignIns.java @@ -74,6 +74,9 @@ public record Started(View view, String refusal) {} */ public Started start(String label, String harness, String actor) { UUID id = UUID.randomUUID(); + // Stored as the row's own start time and sent with every start, so the wait is counted from the + // operator's press however late, or however often, the start is delivered. + Instant requested = Instant.now(); String image = config.agentImage().get(harness); if (image == null) return new Started(null, "harness_unconfigured"); @@ -87,10 +90,11 @@ public Started start(String label, String harness, String actor) { if (inProgress(c, harness).isPresent()) return new Started(null, "sign_in_already_running"); if (pool.hasLabel(c, label)) return new Started(null, "harness_credential_label_taken"); try (PreparedStatement ps = c.prepareStatement(""" - INSERT INTO harness_sign_in (id, label, harness, state, started_by) - VALUES (?, ?, ?, 'PENDING', ?) + INSERT INTO harness_sign_in (id, label, harness, state, started_by, created_at) + VALUES (?, ?, ?, 'PENDING', ?, ?) """)) { ps.setObject(1, id); ps.setString(2, label); ps.setString(3, harness); ps.setString(4, actor); + ps.setTimestamp(5, Timestamp.from(requested)); ps.executeUpdate(); } return new Started(read(c, id).orElseThrow(), null); @@ -110,7 +114,8 @@ INSERT INTO harness_sign_in (id, label, harness, state, started_by) // screen shows and the operator cancels — visible, rather than a unit nobody knows about. try { KafkaSends.sendAndAwait(commands, id.toString(), - new HarnessSignInCommand.Start(id.toString(), harness, image, MAX_WAIT.toSeconds()), + new HarnessSignInCommand.Start(id.toString(), harness, image, MAX_WAIT.toSeconds(), requested, + START_WITHIN.toSeconds()), "harness sign-in start for " + id); } catch (RuntimeException undelivered) { LOG.errorf(undelivered, "sign-in %s could not be asked for", id); @@ -121,6 +126,108 @@ INSERT INTO harness_sign_in (id, label, harness, state, started_by) return opened; } + /** + * How long a start may go unanswered before it is sent again. + * + *

Longer than a consumer group takes to hand a partition to a restarted worker, which is where + * a start used to vanish: the worker read only new records, and one sent while it was rejoining was + * never read, leaving the sign-in PENDING for ever (review of PR #168). + */ + static final Duration RESEND_AFTER = Duration.ofSeconds(45); + + /** + * How long after the press a worker may still OPEN a unit. Sent with every start, so the worker and + * this class use one number (review of PR #168). + * + *

Much shorter than {@link #MAX_WAIT}, which is how long a PERSON may take. It has to be short + * because a worker that cannot start a unit says nothing — another worker may hold it — so the + * deadline below is the only way the operator of a broken worker hears anything. + */ + static final Duration START_WITHIN = Duration.ofMinutes(4); + + /** + * When a sign-in that still shows no code is ended as not started: the start window, plus room for + * the code to cross the bus. A worker opens a unit only if it can print the code inside the window. + */ + static final Duration UNCLAIMED_DEADLINE = START_WITHIN.plusMinutes(2); + + /** + * Re-sends every start nobody has picked up, and fails the ones whose wait has run out. + * + *

Safe to repeat: the start carries the original request time, so the worker gives it only the + * time left and opens nothing once that is gone; and a worker already running the unit ignores the + * repeat. The two exits here are the only ways a PENDING row ends without a worker — which is the + * point: before this, nothing ended one. + */ + @io.quarkus.scheduler.Scheduled(every = "${spire.harness-sign-in-retry-interval:30s}", delayed = "30s", + concurrentExecution = io.quarkus.scheduler.Scheduled.ConcurrentExecution.SKIP) + void resendUnclaimed() { + // A prompt whose code has run out, and whose worker's own "expired" never arrived, would stay on + // screen for ever. The grace covers the worker reporting it the ordinary way first. + update(""" + UPDATE harness_sign_in SET state='FAILED', reason=?, updated_at=now() + WHERE state='PROMPTED' AND expires_at < ? + """, ps -> { + ps.setString(1, HarnessSignInResult.Failed.EXPIRED); + ps.setTimestamp(2, Timestamp.from(Instant.now().minus(EXPIRY_GRACE))); + }); + for (Unclaimed row : unclaimed(Instant.now().minus(RESEND_AFTER))) { + String image = config.agentImage().get(row.harness()); + boolean windowClosed = !row.requestedAt().plus(START_WITHIN).isAfter(Instant.now()); + if (row.requestedAt().plus(UNCLAIMED_DEADLINE).isBefore(Instant.now()) || image == null) { + failUnclaimed(row.id(), image == null ? "harness_unconfigured" : HarnessSignInResult.Failed.NOT_STARTED); + continue; + } + // Past the start window a worker drops the start unopened, so sending it is only noise. The row + // stays PENDING until the deadline above, which leaves room for a code already on its way. + if (windowClosed) continue; + try { + KafkaSends.sendAndAwait(commands, row.id().toString(), new HarnessSignInCommand.Start( + row.id().toString(), row.harness(), image, MAX_WAIT.toSeconds(), row.requestedAt(), + START_WITHIN.toSeconds()), + "harness sign-in re-send for " + row.id()); + } catch (RuntimeException undelivered) { + // The next pass tries again; the row keeps its own deadline either way. + LOG.warnf("sign-in %s could not be re-sent (%s)", row.id(), undelivered.getClass().getSimpleName()); + } + } + } + + /** How long past a code's expiry the worker has to say so itself before the row is closed here. */ + static final Duration EXPIRY_GRACE = Duration.ofMinutes(2); + + private record Unclaimed(UUID id, String harness, Instant requestedAt) {} + + /** + * Fails a row only if it is STILL waiting for a worker when the write runs. + * + *

The list of unclaimed rows was read a moment earlier; a prompt that landed since then has + * given the row an expiry and its own two-minute grace, and ordinary {@link #fail} accepts a + * PROMPTED row too — so it would have ended a sign-in the operator was about to approve + * (review of PR #168). + */ + private void failUnclaimed(UUID id, String reason) { + update(""" + UPDATE harness_sign_in SET state='FAILED', reason=?, updated_at=now() + WHERE id=? AND state='PENDING' + """, ps -> { ps.setString(1, reason); ps.setObject(2, id); }); + } + + private java.util.List unclaimed(Instant before) { + try (Connection c = dataSource.getConnection(); PreparedStatement ps = c.prepareStatement( + "SELECT id, harness, created_at FROM harness_sign_in WHERE state='PENDING' AND created_at < ?")) { + ps.setTimestamp(1, Timestamp.from(before)); + java.util.List rows = new java.util.ArrayList<>(); + try (ResultSet rs = ps.executeQuery()) { + while (rs.next()) rows.add(new Unclaimed(rs.getObject("id", UUID.class), rs.getString("harness"), + rs.getTimestamp("created_at").toInstant())); + } + return rows; + } catch (SQLException failure) { + throw new IllegalStateException("The waiting sign-ins could not be read", failure); + } + } + /** What the operator must do. Written straight through: the screen is polling for exactly this. */ public void prompted(HarnessSignInResult.Prompted result) { update(""" diff --git a/spire-orchestrator/src/main/resources/application.yml b/spire-orchestrator/src/main/resources/application.yml index 6d0f7f49..21b5a222 100644 --- a/spire-orchestrator/src/main/resources/application.yml +++ b/spire-orchestrator/src/main/resources/application.yml @@ -278,9 +278,14 @@ mp: id: spire-orchestrator-harness-image value: deserializer: dev.codespire.orchestrator.factory.HarnessImageResultDeserializer + # Earliest, not latest. A new group, or one rejoining while a dead member's session runs + # out after a restart, is only assigned its partition after a delay; with latest, a message + # sent in that window is skipped, and the model list stayed unknown until the next refresh. + # Measured on the dev stack on 2026-09-23. Replaying is harmless: answers are idempotent, + # handled in order, and an answer about an image no longer configured is dropped. auto: offset: - reset: latest + reset: earliest failure-strategy: dead-letter-queue dead-letter-queue: topic: cs.dlq @@ -291,10 +296,14 @@ mp: id: spire-orchestrator-harness-sign-in value: deserializer: dev.codespire.orchestrator.factory.HarnessSignInResultDeserializer - # Latest: a restart must not replay prompts whose codes expired days ago onto a screen. + # Earliest (was latest). Under latest a result sent while this group was being assigned its + # partition was never read: a prompt the screen was waiting for, or a finished sign-in, was + # lost (review of PR #168). Replay only moves a sign-in forward -- a prompt applies to a PENDING + # row, a completion or failure to an open one -- and an open row past its code's expiry is + # closed by HarnessSignIns.resendUnclaimed, so an old prompt cannot keep one on screen. auto: offset: - reset: latest + reset: earliest failure-strategy: dead-letter-queue dead-letter-queue: topic: cs.dlq @@ -476,6 +485,7 @@ spire: work-gate-expiry-interval: "off" work-preparation-interval: "off" harness-catalogue-interval: "off" + harness-sign-in-retry-interval: "off" work-run-interval: "off" work-delivery-interval: "off" work-activity-interval: "off" diff --git a/spire-orchestrator/src/main/resources/db/migration/V82__freshness_of_harness_answers.sql b/spire-orchestrator/src/main/resources/db/migration/V82__freshness_of_harness_answers.sql new file mode 100644 index 00000000..5a10b34d --- /dev/null +++ b/spire-orchestrator/src/main/resources/db/migration/V82__freshness_of_harness_answers.sql @@ -0,0 +1,4 @@ +-- Review of PR #168: the harness channels now replay from their oldest record, so an answer can +-- arrive after a newer one. The catalogue keeps the answer to the NEWEST question, compared by when the +-- question was asked. NULL for a row written before this, which any timed answer replaces. +ALTER TABLE harness_catalogue ADD COLUMN asked_at TIMESTAMPTZ; diff --git a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/dlq/DlqTopicsTest.java b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/dlq/DlqTopicsTest.java index 90527c4a..63d645a7 100644 --- a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/dlq/DlqTopicsTest.java +++ b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/dlq/DlqTopicsTest.java @@ -11,6 +11,17 @@ void repositoryDeliveriesReplayWithTheirProvenance() { assertEquals("cs.repository-integration", DlqTopics.forType("RepositoryDelivery")); } + /** Before these were mapped, each fell through to cs.commands, where nothing reads it (review of PR #168). */ + @Test void harnessRecordsReplayOntoTheirOwnTopics() { + assertEquals("cs.harness-image-commands", DlqTopics.forType("Describe")); + assertEquals("cs.harness-image-results", DlqTopics.forType("Described")); + assertEquals("cs.harness-sign-in-commands", DlqTopics.forType("Start")); + assertEquals("cs.harness-sign-in-commands", DlqTopics.forType("Cancel")); + assertEquals("cs.harness-sign-in-results", DlqTopics.forType("Prompted")); + assertEquals("cs.harness-sign-in-results", DlqTopics.forType("Completed")); + assertEquals("cs.harness-sign-in-results", DlqTopics.forType("Failed")); + } + @Test void registrationReplaysOntoItsRegistryTopic() { assertEquals("cs.registry-integration", DlqTopics.forType("RepositoryRegistration")); } diff --git a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessCataloguesTest.java b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessCataloguesTest.java index 8148bd7d..e06414fd 100644 --- a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessCataloguesTest.java +++ b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessCataloguesTest.java @@ -40,6 +40,53 @@ private String image() { return config.agentImage().get(HARNESS); } + /** + * The channel replays from its oldest record, so an answer can arrive after a newer one. The answer to + * the NEWEST question is kept, whatever order they arrive in (review of PR #168). + */ + @Test + void anOlderAnswerArrivingLateDoesNotRollTheListBack() { + java.time.Instant earlier = java.time.Instant.parse("2026-09-23T07:00:00Z"); + java.time.Instant later = earlier.plusSeconds(600); + catalogues.record(new HarnessImageResult.Described("TEST-new", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-current", true, 1)), "TEST-pin-new", later)); + + catalogues.record(new HarnessImageResult.Described("TEST-old", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-previous", true, 1)), "TEST-pin-old", earlier)); + // An answer from before times were sent is the oldest of all. + catalogues.record(new HarnessImageResult.Described("TEST-untimed", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-untimed", true, 1)), "TEST-pin-untimed")); + + var kept = catalogues.get(HARNESS).orElseThrow(); + assertEquals("TEST-pin-new", kept.pinnedImage()); + assertTrue(kept.find("TEST-current").isPresent()); + } + + /** An answer with no time predates times: it may fill an empty cache, never replace a row. */ + @Test + void anAnswerWithNoTimeReplacesNothing() { + catalogues.record(new HarnessImageResult.Described("TEST-first", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-first", true, 1)), "TEST-pin-first")); + + catalogues.record(new HarnessImageResult.Described("TEST-second", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-second", true, 1)), "TEST-pin-second")); + + assertEquals("TEST-pin-first", catalogues.get(HARNESS).orElseThrow().pinnedImage()); + } + + @Test + void aNewerAnswerReplacesAnOlderOne() { + java.time.Instant earlier = java.time.Instant.parse("2026-09-23T07:00:00Z"); + catalogues.record(new HarnessImageResult.Described("TEST-old", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-previous", true, 1)), "TEST-pin-old", earlier)); + + catalogues.record(new HarnessImageResult.Described("TEST-new", HARNESS, image(), + HarnessImageResult.Status.OK, List.of(model("TEST-current", true, 1)), "TEST-pin-new", + earlier.plusSeconds(600))); + + assertEquals("TEST-pin-new", catalogues.get(HARNESS).orElseThrow().pinnedImage()); + } + /** * A run uses the exact image the list was read from; with nothing read, the tag, as before part M. * And an answer about an image the harness has since left pins nothing (review of PR #167). @@ -64,7 +111,9 @@ private static HarnessImageResult.Model model(String slug, boolean visible, int private HarnessImageResult.Described answer(String image, HarnessImageResult.Status status, List models) { - return new HarnessImageResult.Described("TEST-request", HARNESS, image, status, models); + // Timed, as every answer is now: an untimed one predates times and replaces nothing. + return new HarnessImageResult.Described("TEST-request", HARNESS, image, status, models, null, + java.time.Instant.now()); } /** "Not asked yet" and "the image declares nothing" send an operator to different places. */ diff --git a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessImageResultsOrderTest.java b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessImageResultsOrderTest.java index f121d399..8f635eee 100644 --- a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessImageResultsOrderTest.java +++ b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessImageResultsOrderTest.java @@ -13,6 +13,42 @@ */ class HarnessImageResultsOrderTest { + /** + * An answer sent while the orchestrator's group was still being assigned, after a restart, must + * not be skipped (measured on the dev stack, 2026-09-23). + */ + @Test + void anAnswerSentBeforeTheOrchestratorJoinedIsStillRecorded() throws java.io.IOException { + assertTrue(channel("harness-image-results-in").contains("reset: earliest")); + } + + /** A sign-in result lost the same way left the screen waiting for a code already printed. */ + @Test + void aSignInResultSentBeforeTheOrchestratorJoinedIsStillRecorded() throws java.io.IOException { + assertTrue(channel("harness-sign-in-results-in").contains("reset: earliest")); + } + + /** One channel's settings, up to its failure strategy. */ + private static String channel(String name) throws java.io.IOException { + String yaml; + try (var in = HarnessImageResults.class.getResourceAsStream("/application.yml")) { + yaml = new String(java.util.Objects.requireNonNull(in, "application.yml").readAllBytes(), + java.nio.charset.StandardCharsets.UTF_8); + } + int start = yaml.indexOf(" " + name + ":"); + int end = yaml.indexOf("failure-strategy", start); + assertTrue(start >= 0 && end > start, name + " is not declared"); + return yaml.substring(start, end); + } + + /** The first tick waits: at startup it ran before the Kafka emitter was connected. */ + @Test + void theRefreshTimerDoesNotFireDuringStartup() throws NoSuchMethodException { + var scheduled = HarnessCatalogues.class.getDeclaredMethod("refresh") + .getAnnotation(io.quarkus.scheduler.Scheduled.class); + assertTrue(!scheduled.delayed().isBlank(), "the first refresh must be delayed"); + } + @Test void theAnswersAreRecordedOneAtATime() throws NoSuchMethodException { Blocking blocking = HarnessImageResults.class.getMethod("onResult", Message.class).getAnnotation(Blocking.class); diff --git a/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessSignInRetriesTest.java b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessSignInRetriesTest.java new file mode 100644 index 00000000..095ac223 --- /dev/null +++ b/spire-orchestrator/src/test/java/dev/codespire/orchestrator/factory/HarnessSignInRetriesTest.java @@ -0,0 +1,175 @@ +package dev.codespire.orchestrator.factory; + +import com.fasterxml.jackson.databind.JsonNode; +import com.fasterxml.jackson.databind.ObjectMapper; +import dev.codespire.contract.event.HarnessSignInResult; +import io.quarkus.test.junit.QuarkusTest; +import jakarta.inject.Inject; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.clients.consumer.KafkaConsumer; +import org.apache.kafka.common.serialization.StringDeserializer; +import org.eclipse.microprofile.config.inject.ConfigProperty; +import org.junit.jupiter.api.AfterEach; +import org.junit.jupiter.api.Test; + +import javax.sql.DataSource; +import java.sql.Connection; +import java.sql.PreparedStatement; +import java.sql.SQLException; +import java.sql.Statement; +import java.sql.Timestamp; +import java.time.Duration; +import java.time.Instant; +import java.time.temporal.ChronoUnit; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.UUID; + +import static org.junit.jupiter.api.Assertions.*; + +/** + * A sign-in nobody picked up (review of PR #168). + * + *

A start sent while the worker's group was still being assigned its partition was never read, and + * nothing ended the row: it stayed PENDING, and the screen said "starting" for ever. These are the two + * ways out — send it again while its wait lasts, fail it once the wait is gone — and the matching exit + * for a prompt whose code ran out without the worker saying so. + */ +@QuarkusTest +class HarnessSignInRetriesTest { + + @Inject HarnessSignIns signIns; + @Inject DataSource dataSource; + @Inject ObjectMapper mapper; + @ConfigProperty(name = "kafka.bootstrap.servers") String bootstrap; + + private static final String OWNED = "TEST-retry-%"; + + @AfterEach + void clean() throws SQLException { + try (Connection c = dataSource.getConnection(); Statement s = c.createStatement()) { + s.executeUpdate("DELETE FROM harness_sign_in WHERE label LIKE '" + OWNED + "'"); + } + } + + private UUID start(String label) { + HarnessSignIns.Started started = signIns.start(label, "codex", "TEST-operator"); + assertNull(started.refusal(), started.refusal()); + return started.view().id(); + } + + /** Moves the operator's press into the past, as if the start had gone unanswered that long. */ + private Instant pressedAgo(UUID id, Duration ago) throws SQLException { + Instant pressed = Instant.now().minus(ago).truncatedTo(ChronoUnit.MILLIS); + try (Connection c = dataSource.getConnection(); + PreparedStatement ps = c.prepareStatement("UPDATE harness_sign_in SET created_at=? WHERE id=?")) { + ps.setTimestamp(1, Timestamp.from(pressed)); ps.setObject(2, id); ps.executeUpdate(); + } + return pressed; + } + + /** Every start sent for this sign-in, oldest first. */ + private List startsFor(UUID id, int atLeast) throws Exception { + List own = new ArrayList<>(); + try (KafkaConsumer consumer = new KafkaConsumer<>(Map.of( + ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, bootstrap, + ConsumerConfig.GROUP_ID_CONFIG, "TEST-sign-in-retries-" + UUID.randomUUID(), + ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class, + ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class, + ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"))) { + consumer.subscribe(List.of("cs.harness-sign-in-commands")); + long deadline = System.nanoTime() + Duration.ofSeconds(15).toNanos(); + while (System.nanoTime() < deadline && own.size() < atLeast) collect(consumer, id, own, Duration.ofMillis(500)); + collect(consumer, id, own, Duration.ofSeconds(1)); + } + return own; + } + + private void collect(KafkaConsumer consumer, UUID id, List own, Duration wait) throws Exception { + for (ConsumerRecord record : consumer.poll(wait)) { + if (!id.toString().equals(record.key())) continue; + JsonNode value = mapper.readTree(record.value()); + if ("Start".equals(value.path("type").asText())) own.add(value); + } + } + + @Test + void aStartNobodyPickedUpIsSentAgainCountedFromTheOriginalPress() throws Exception { + UUID id = start("TEST-retry-resend"); + Instant pressed = pressedAgo(id, Duration.ofSeconds(60)); + + signIns.resendUnclaimed(); + + List starts = startsFor(id, 2); + assertEquals(2, starts.size(), "the original start and one re-send"); + assertFalse(starts.get(0).path("requestedAt").isNull() || starts.get(0).path("requestedAt").isMissingNode(), + "the first start carries the press time too"); + // The wait is anchored to the press, so a worker receiving this late gives it only what is left. + assertEquals(pressed, mapper.convertValue(starts.get(1).path("requestedAt"), Instant.class)); + // The window the worker must open a unit in is the one this side ends unclaimed rows by. + assertEquals(HarnessSignIns.START_WITHIN.toSeconds(), starts.get(0).path("startWithinSeconds").asLong()); + assertEquals(HarnessSignIns.START_WITHIN.toSeconds(), starts.get(1).path("startWithinSeconds").asLong()); + assertEquals("PENDING", signIns.get(id).orElseThrow().state()); + } + + @Test + void aStartStillUnansweredWhenItsWaitRanOutFailsAndSaysWhy() throws Exception { + UUID id = start("TEST-retry-expired"); + pressedAgo(id, HarnessSignIns.UNCLAIMED_DEADLINE.plusSeconds(30)); + + signIns.resendUnclaimed(); + + HarnessSignIns.View view = signIns.get(id).orElseThrow(); + assertEquals("FAILED", view.state()); + assertEquals(HarnessSignInResult.Failed.NOT_STARTED, view.reason()); + assertEquals(1, startsFor(id, 1).size(), "a sign-in whose wait is gone is not sent again"); + } + + /** Past the start window a worker opens nothing, so nothing is sent; the row waits for its deadline. */ + @Test + void aStartPastItsWindowIsNotSentAgainButNotYetEnded() throws Exception { + UUID id = start("TEST-retry-window-closed"); + pressedAgo(id, HarnessSignIns.START_WITHIN.plusSeconds(30)); + + signIns.resendUnclaimed(); + + assertEquals(1, startsFor(id, 1).size()); + assertEquals("PENDING", signIns.get(id).orElseThrow().state(), "a code may still be on its way"); + } + + @Test + void aFreshStartIsNotSentAgainYet() throws Exception { + UUID id = start("TEST-retry-fresh"); + + signIns.resendUnclaimed(); + + assertEquals(1, startsFor(id, 1).size()); + } + + /** The worker normally reports an expired code itself; the row must close even when that is lost. */ + @Test + void aPromptWhoseCodeRanOutIsClosedOnceTheGraceHasPassed() { + UUID expired = start("TEST-retry-prompt-old"); + signIns.prompted(new HarnessSignInResult.Prompted(expired.toString(), "https://auth.example.test/device", + "ABCD-12345", Instant.now().minus(HarnessSignIns.EXPIRY_GRACE).minusSeconds(30))); + + signIns.resendUnclaimed(); + + HarnessSignIns.View view = signIns.get(expired).orElseThrow(); + assertEquals("FAILED", view.state()); + assertEquals(HarnessSignInResult.Failed.EXPIRED, view.reason()); + } + + @Test + void aPromptJustPastItsExpiryWaitsForTheWorkerToSaySo() { + UUID recent = start("TEST-retry-prompt-new"); + signIns.prompted(new HarnessSignInResult.Prompted(recent.toString(), "https://auth.example.test/device", + "ABCD-12345", Instant.now().minusSeconds(30))); + + signIns.resendUnclaimed(); + + assertEquals("PROMPTED", signIns.get(recent).orElseThrow().state()); + } +} diff --git a/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessImageWorker.java b/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessImageWorker.java index a1dc99d9..a77a4462 100644 --- a/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessImageWorker.java +++ b/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessImageWorker.java @@ -51,11 +51,35 @@ public class HarnessImageWorker { @Blocking public CompletionStage onCommand(Message message) { if (message.getPayload() instanceof HarnessImageCommand.Describe describe) { + if (isStale(describe, java.time.Instant.now())) { + // Replayed from before a restart. The orchestrator asks again on its own schedule, and a + // pull for an image nobody runs any more could hold up the current question for minutes. + LOG.infof("skipping a question about %s asked at %s; a newer one follows", describe.harness(), + describe.askedAt()); + return message.ack(); + } publish(describe(describe)); } return message.ack(); } + /** + * How old a question may be and still be answered. + * + *

Longer than the orchestrator's refresh interval (ten minutes by default), so an ordinary + * question is never skipped; short enough that a backlog replayed after an outage is dropped instead + * of pulled one image at a time ahead of the question that matters (review of PR #168). + */ + static final java.time.Duration STALE_AFTER = java.time.Duration.ofMinutes(15); + + /** + * A question with no time was sent before times existed, so it can only be a replay: skipped, like + * any other old one. The orchestrator asks again, with a time, at start-up and on its schedule. + */ + static boolean isStale(HarnessImageCommand.Describe question, java.time.Instant now) { + return question.askedAt() == null || question.askedAt().plus(STALE_AFTER).isBefore(now); + } + HarnessImageResult.Described describe(HarnessImageCommand.Describe command) { dev.codespire.runtime.ImageDescription image; try { @@ -66,11 +90,11 @@ HarnessImageResult.Described describe(HarnessImageCommand.Describe command) { LOG.warnf("the image for harness %s could not be read (%s)", command.harness(), unreachable.getClass().getSimpleName()); return new HarnessImageResult.Described(command.requestId(), command.harness(), command.image(), - HarnessImageResult.Status.IMAGE_UNAVAILABLE, List.of()); + HarnessImageResult.Status.IMAGE_UNAVAILABLE, List.of(), null, command.askedAt()); } ModelCatalogueLabel.Read read = ModelCatalogueLabel.of(image.labels(), mapper); return new HarnessImageResult.Described(command.requestId(), command.harness(), command.image(), - read.status(), read.models(), image.pinned()); + read.status(), read.models(), image.pinned(), command.askedAt()); } private static String keyOf(HarnessImageResult result) { diff --git a/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessSignInWorker.java b/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessSignInWorker.java index d6733115..4bbf4bf3 100644 --- a/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessSignInWorker.java +++ b/spire-run-worker/src/main/java/dev/codespire/runworker/HarnessSignInWorker.java @@ -54,6 +54,9 @@ public class HarnessSignInWorker { */ private static final Duration ABANDONED_AFTER = Duration.ofMinutes(30); + /** The time the deadlines are read against. Replaceable so a test can make a unit start slowly. */ + java.time.Clock clock = java.time.Clock.systemUTC(); + @Inject SignInRuntime runtime; @Inject EncryptionService encryption; @Inject HarnessRegistry harnesses; @@ -112,6 +115,22 @@ private void start(HarnessSignInCommand.Start command) { LOG.infof("sign-in %s is already running here; ignoring a repeated start", command.signInId()); return; } + // Counted from the operator's press, not from delivery: a start can be replayed or re-sent. One + // that could not show its code before its start window closes — the window the orchestrator ends + // unclaimed sign-ins by — opens nothing, and neither does one with no press time (sent before + // times existed, so only a replay). Its code would reach a row already ended (review of PR #168). + // + // And it SAYS nothing. "Too late to start another unit" is not "this sign-in expired": with two + // workers, a late copy reaching one of them would otherwise fail a sign-in the other is running + // and the operator can still approve. Ending a sign-in nobody runs is the orchestrator's job, and + // it does it from the row's own clock (HarnessSignIns.resendUnclaimed). + if (!command.mayOpenAt(Instant.now(clock), PROMPT_TIMEOUT)) { + LOG.infof("not starting sign-in %s: requested at %s, and its start window has closed", + command.signInId(), command.requestedAt()); + running.remove(command.signInId()); + return; + } + Duration wait = command.remainingWait(Instant.now(clock)); if (cancelled.remove(command.signInId()) != null) { // The cancel got here first. Creating the unit now would mean the operator's cancel did // nothing and a container waited out its ceiling on a code nobody would type. @@ -133,7 +152,7 @@ private void start(HarnessSignInCommand.Start command) { SignInPrompt prompt = new SignInPrompt(flow.get().verificationHost()); SignInUnitSpec spec = new SignInUnitSpec(command.signInId(), command.image(), flow.get().command(), flow.get().resultPath(), EnterpriseEnvironment.NONE, 512 * MEGABYTE, 1_000_000_000L, - 64 * MEGABYTE, Duration.ofSeconds(command.maxWaitSeconds())); + 64 * MEGABYTE, wait); SignInRuntime.Handle handle; try { handle = runtime.start(spec, prompt::accept); } @@ -142,7 +161,16 @@ private void start(HarnessSignInCommand.Start command) { adopt(command.signInId(), taken); return; } - catch (RuntimeException failed) { running.remove(command.signInId()); throw failed; } + catch (RuntimeException failed) { + // Before this worker owns a unit, a failure says nothing about the sign-in: another worker may + // hold it, and a timeout here cannot tell. Reporting it ended a sign-in the other worker was + // running (review of PR #168). So it is logged and dropped; the orchestrator re-sends the start + // and, if no worker ever shows a code, ends the row itself as not started. + running.remove(command.signInId()); + LOG.warnf("sign-in %s: no unit could be started here (%s); leaving it to a retry", + command.signInId(), failed.getClass().getSimpleName()); + return; + } running.put(command.signInId(), handle); if (cancelled.remove(command.signInId()) != null) { // It arrived while the container was being created — the window the claim alone cannot @@ -150,7 +178,12 @@ private void start(HarnessSignInCommand.Start command) { runtime.cancel(handle); } try { - if (!awaitPrompt(prompt)) { + // Absolute deadlines, taken from the press and checked where they matter, not durations taken + // before the unit started: a slow start or a paused process used to stretch both — a code + // published after the row had been ended, and a wait longer than the person's (review of PR #168). + Instant promptBy = command.promptDeadline(); + Instant approveBy = command.approvalDeadline(); + if (!awaitPrompt(prompt, promptBy)) { // The CLI said nothing this build can read, or said two different things. Never invent a // link or a code, and never pick between two: an operator sent to a guessed address is // worse than one told the sign-in did not start, because they will type their account @@ -161,15 +194,26 @@ private void start(HarnessSignInCommand.Start command) { : "the sign-in tool printed no link and code this build could read")); return; } - // The EFFECTIVE deadline: the sooner of what the vendor promised and what this worker was - // given. The screen counted down the vendor's figure while the worker waited on its own, so - // it could show a minute remaining on a unit that had already been destroyed. - Duration budget = Duration.ofSeconds(command.maxWaitSeconds()); - Duration life = prompt.expiresIn().filter(vendor -> vendor.compareTo(budget) < 0).orElse(budget); - emit(new HarnessSignInResult.Prompted(command.signInId(), prompt.link(), prompt.code(), - Instant.now().plus(life))); + if (Instant.now(clock).isAfter(promptBy)) { + // Printed, but too late: the orchestrator has ended, or is about to end, this row as not + // started. Publishing now would put a code on screen for a sign-in nobody can finish. The + // unit is destroyed below; nothing is reported, because the row already has its answer. + LOG.infof("sign-in %s printed its code after its window closed; discarding the unit", + command.signInId()); + return; + } + // The EFFECTIVE deadline: the sooner of what the vendor promised and what is left of the + // person's time, measured NOW. The screen counted down the vendor's figure while the worker + // waited on its own, so it could show a minute remaining on a unit already destroyed. + // ONE reading of the clock for the expiry: a duration from one reading added to a later one + // put the countdown past the real deadline by however long passed between them — visible + // after a pause (review of PR #168). + Instant shown = Instant.now(clock); + Instant expires = prompt.expiresIn().map(shown::plus).filter(vendor -> vendor.isBefore(approveBy)) + .orElse(approveBy); + emit(new HarnessSignInResult.Prompted(command.signInId(), prompt.link(), prompt.code(), expires)); - SignInRuntime.Exit exit = runtime.awaitExit(handle, Duration.ofSeconds(command.maxWaitSeconds())); + SignInRuntime.Exit exit = runtime.awaitExit(handle, nonNegative(Duration.between(Instant.now(clock), approveBy))); if (!(exit instanceof SignInRuntime.Exit.Observed observed)) { runtime.cancel(handle); boolean fault = exit instanceof SignInRuntime.Exit.Unobservable; @@ -189,10 +233,9 @@ private void start(HarnessSignInCommand.Start command) { // A second send is attempted, and it is NOT claimed to arrive. It travels the same // broker path that has just failed, so in the outage this exists for it fails too. // - // What the operator sees depends on when the outage began: a row that reached PROMPTED - // counts down and expires visibly, while one that never did stays PENDING with no - // countdown, because the prompt is the only thing that writes an expiry. Both are - // cleared by cancelling, which also needs the broker back. UNVERIFIED.md A5. + // What the operator sees: a row that reached PROMPTED counts down and is closed after its + // expiry; one that never did is ended as not started. Both by the orchestrator on its own + // clock, without this broker (HarnessSignIns.resendUnclaimed). UNVERIFIED.md A5. LOG.errorf("sign-in %s completed but could not be delivered; the credential is discarded" + " and the row stays open until the operator cancels it", command.signInId()); emit(new HarnessSignInResult.Failed(command.signInId(), HarnessSignInResult.Failed.UNIT_FAILED, @@ -280,9 +323,15 @@ private void cancel(HarnessSignInCommand.Cancel command) { emit(new HarnessSignInResult.Failed(command.signInId(), HarnessSignInResult.Failed.CANCELLED, command.reason())); } - private boolean awaitPrompt(SignInPrompt prompt) { - Instant deadline = Instant.now().plus(PROMPT_TIMEOUT); - while (Instant.now().isBefore(deadline)) { + private static Duration nonNegative(Duration duration) { + return duration.isNegative() ? Duration.ZERO : duration; + } + + /** Waits for the code, but never past the window it must be on screen by. */ + private boolean awaitPrompt(SignInPrompt prompt, Instant promptBy) { + Instant patience = Instant.now(clock).plus(PROMPT_TIMEOUT); + Instant deadline = patience.isBefore(promptBy) ? patience : promptBy; + while (Instant.now(clock).isBefore(deadline)) { // Ambiguity does not improve by waiting, and waiting out the whole timeout for a decision // already made just keeps a container and an operator hanging. if (prompt.ambiguous()) return false; diff --git a/spire-run-worker/src/main/resources/application.yml b/spire-run-worker/src/main/resources/application.yml index dce2c9f6..ab0cdc17 100644 --- a/spire-run-worker/src/main/resources/application.yml +++ b/spire-run-worker/src/main/resources/application.yml @@ -225,11 +225,14 @@ mp: id: spire-run-worker-sign-in value: deserializer: dev.codespire.runworker.HarnessSignInCommandDeserializer - # Latest: a restart must not re-open sign-ins whose codes expired days ago and whose - # operator has long gone. An abandoned one is re-started from the screen, by a person. + # Earliest (was latest). Under latest a start sent while this group was being assigned its + # partition was never read, and the sign-in sat PENDING for ever (review of PR #168). What + # latest protected against is now held by the start itself: it carries the operator's press + # time, and one whose wait has run out opens no unit. The orchestrator also re-sends a start + # nobody picked up, which the worker ignores while that unit is already running. auto: offset: - reset: latest + reset: earliest failure-strategy: dead-letter-queue dead-letter-queue: topic: cs.dlq @@ -242,9 +245,14 @@ mp: id: spire-run-worker-harness-image value: deserializer: dev.codespire.runworker.HarnessImageCommandDeserializer + # Earliest, not latest. A new group, or one rejoining while a dead member's session runs + # out after a restart, is only assigned its partition after a delay; with latest, a message + # sent in that window is skipped, and the model list stayed unknown until the next refresh. + # Measured on the dev stack on 2026-09-23. Replaying is harmless: answers are idempotent, + # handled in order, and an answer about an image no longer configured is dropped. auto: offset: - reset: latest + reset: earliest failure-strategy: dead-letter-queue dead-letter-queue: topic: cs.dlq diff --git a/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessImageWorkerTest.java b/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessImageWorkerTest.java index 54b42209..764388cd 100644 --- a/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessImageWorkerTest.java +++ b/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessImageWorkerTest.java @@ -38,7 +38,8 @@ void anAnswerIsKeyedByItsHarnessSoOneHarnessKeepsOnePartition() { worker.ackSeconds = 5; worker.results = new Capturing(sent); - worker.onCommand(Message.of(new HarnessImageCommand.Describe("TEST-request-1", "codex", "TEST-image"))); + worker.onCommand(Message.of(new HarnessImageCommand.Describe("TEST-request-1", "codex", "TEST-image", + java.time.Instant.now()))); assertEquals(1, sent.size()); assertEquals("codex", sent.getFirst().key()); @@ -47,6 +48,56 @@ void anAnswerIsKeyedByItsHarnessSoOneHarnessKeepsOnePartition() { ((HarnessImageResult.Described) sent.getFirst().value()).pinnedImage()); } + /** A replayed backlog is dropped instead of pulled image by image ahead of the current question. */ + @Test + void aQuestionTooOldToMatterIsSkippedWithoutReadingTheImage() { + List> sent = new ArrayList<>(); + HarnessImageWorker worker = new HarnessImageWorker(); + worker.runtime = (RunRuntime) Proxy.newProxyInstance(RunRuntime.class.getClassLoader(), new Class[] { RunRuntime.class }, + (proxy, method, args) -> { throw new AssertionError("a stale question must not reach the image"); }); + worker.mapper = new ObjectMapper(); + worker.ackSeconds = 5; + worker.results = new Capturing(sent); + java.time.Instant longAgo = java.time.Instant.now().minus(HarnessImageWorker.STALE_AFTER).minusSeconds(60); + + worker.onCommand(Message.of(new HarnessImageCommand.Describe("TEST-request-old", "codex", "TEST-image", longAgo))); + + assertEquals(0, sent.size()); + } + + /** A question with no time predates times, so it can only be a replay (review of PR #168). */ + @Test + void aQuestionWithNoTimeIsSkippedWithoutReadingTheImage() { + List> sent = new ArrayList<>(); + HarnessImageWorker worker = new HarnessImageWorker(); + worker.runtime = (RunRuntime) Proxy.newProxyInstance(RunRuntime.class.getClassLoader(), new Class[] { RunRuntime.class }, + (proxy, method, args) -> { throw new AssertionError("an undated question must not reach the image"); }); + worker.mapper = new ObjectMapper(); + worker.ackSeconds = 5; + worker.results = new Capturing(sent); + + worker.onCommand(Message.of(new HarnessImageCommand.Describe("TEST-request-undated", "codex", "TEST-image"))); + + assertEquals(0, sent.size()); + } + + @Test + void theAnswerCarriesTheTimeOfTheQuestion() { + List> sent = new ArrayList<>(); + HarnessImageWorker worker = new HarnessImageWorker(); + worker.runtime = (RunRuntime) Proxy.newProxyInstance(RunRuntime.class.getClassLoader(), new Class[] { RunRuntime.class }, + (proxy, method, args) -> method.getName().equals("describeImage") + ? new dev.codespire.runtime.ImageDescription("TEST-pin", Map.of()) : null); + worker.mapper = new ObjectMapper(); + worker.ackSeconds = 5; + worker.results = new Capturing(sent); + java.time.Instant asked = java.time.Instant.now().minusSeconds(5); + + worker.onCommand(Message.of(new HarnessImageCommand.Describe("TEST-request-new", "codex", "TEST-image", asked))); + + assertEquals(asked, ((HarnessImageResult.Described) sent.getFirst().value()).askedAt()); + } + /** One partition keeps order only if the worker does not answer two of its messages at once. */ @Test void theQuestionsAreAnsweredOneAtATime() throws NoSuchMethodException { diff --git a/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessSignInStartTimeTest.java b/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessSignInStartTimeTest.java new file mode 100644 index 00000000..67aaa264 --- /dev/null +++ b/spire-run-worker/src/test/java/dev/codespire/runworker/HarnessSignInStartTimeTest.java @@ -0,0 +1,257 @@ +package dev.codespire.runworker; + +import dev.codespire.contract.command.HarnessSignInCommand; +import dev.codespire.contract.event.HarnessSignInResult; +import dev.codespire.runtime.SignInRuntime; +import io.smallrye.reactive.messaging.kafka.Record; +import org.eclipse.microprofile.reactive.messaging.Emitter; +import org.eclipse.microprofile.reactive.messaging.Message; +import org.junit.jupiter.api.Test; + +import java.lang.reflect.Proxy; +import java.time.Instant; +import java.util.ArrayList; +import java.util.List; +import java.util.Map; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionStage; + +import static org.junit.jupiter.api.Assertions.assertEquals; + +/** + * A start is replayed after a restart and re-sent when nobody picked it up, so its wait is counted from + * the operator's press (review of PR #168). Plain unit tests: the paths under test end before any + * container, and the runtime fails the test if it is reached. + */ +class HarnessSignInStartTimeTest { + + private final List sent = new ArrayList<>(); + + private HarnessSignInWorker worker() { + HarnessSignInWorker worker = new HarnessSignInWorker(); + worker.runtime = (SignInRuntime) Proxy.newProxyInstance(SignInRuntime.class.getClassLoader(), + new Class[] { SignInRuntime.class }, + (proxy, method, args) -> { throw new AssertionError("no unit may be started: " + method.getName()); }); + worker.ackSeconds = 5; + worker.results = new Capturing(sent); + return worker; + } + + private static HarnessSignInCommand.Start pressed(Instant at) { + return new HarnessSignInCommand.Start("TEST-sign-in", "codex", "TEST-image", 840, at, 240); + } + + /** + * Opens no unit, and says nothing: with two workers, a late copy reporting "expired" would end a + * sign-in the other worker is running (review of PR #168). The orchestrator ends unclaimed rows. + */ + @Test + void aStartWhoseWaitHasRunOutOpensNoUnitAndReportsNothing() { + worker().onCommand(Message.of(pressed(Instant.now().minusSeconds(900)))); + + assertEquals(0, sent.size()); + } + + /** + * Most of the person's wait is left, but a code printed now would land after the orchestrator ended + * the row (review of PR #168). The start window, not the wait, decides. + */ + @Test + void aStartPastItsWindowOpensNoUnitEvenWithWaitLeft() { + worker().onCommand(Message.of(pressed(Instant.now().minusSeconds(200)))); + + assertEquals(0, sent.size()); + } + + /** A start with no press time predates this change, so it can only be a replay. */ + @Test + void aStartWithNoPressTimeOpensNoUnit() { + worker().onCommand(Message.of(new HarnessSignInCommand.Start("TEST-sign-in", "codex", "TEST-image", 840))); + + assertEquals(0, sent.size()); + } + + /** + * A failure before this worker owns a unit — a timeout, a refused create — says nothing about the + * sign-in: another worker may be running it. So it reports nothing (review of PR #168). + */ + @Test + void aFailureBeforeOwnershipReportsNothing() { + HarnessSignInWorker worker = worker(); + worker.harnesses = new HarnessRegistry(); + worker.runtime = (SignInRuntime) Proxy.newProxyInstance(SignInRuntime.class.getClassLoader(), + new Class[] { SignInRuntime.class }, (proxy, method, args) -> { + if (method.getName().equals("start")) throw new IllegalStateException("TEST connection reset"); + throw new AssertionError("nothing else may be reached: " + method.getName()); + }); + + worker.onCommand(Message.of(pressed(Instant.now()))); + + assertEquals(0, sent.size()); + } + + /** + * A unit admitted in time that then starts slowly — or whose process is paused — prints its code + * after the window has closed. The row has been ended by then, so nothing may be published; the unit + * is destroyed (review of PR #168). + */ + @Test + void aCodePrintedAfterTheWindowIsNotPublished() { + HarnessSignInWorker worker = worker(); + worker.harnesses = new HarnessRegistry(); + Instant pressed = Instant.parse("2026-09-24T12:00:00Z"); + java.util.concurrent.atomic.AtomicReference now = new java.util.concurrent.atomic.AtomicReference<>(pressed); + worker.clock = new java.time.Clock() { + @Override public java.time.ZoneId getZone() { return java.time.ZoneOffset.UTC; } + @Override public java.time.Clock withZone(java.time.ZoneId zone) { return this; } + @Override public Instant instant() { return now.get(); } + }; + List destroyed = new ArrayList<>(); + worker.runtime = (SignInRuntime) Proxy.newProxyInstance(SignInRuntime.class.getClassLoader(), + new Class[] { SignInRuntime.class }, (proxy, method, args) -> switch (method.getName()) { + case "start" -> { + now.set(pressed.plusSeconds(600)); + @SuppressWarnings("unchecked") + java.util.function.Consumer lines = (java.util.function.Consumer) args[1]; + lines.accept("1. Open this link in your browser and sign in to your account"); + lines.accept(" https://auth.openai.com/codex/device"); + lines.accept("2. Enter this one-time code (expires in 15 minutes)"); + lines.accept(" ABCD-12345"); + yield new SignInRuntime.Handle("TEST-sign-in", "TEST-unit"); + } + case "destroy" -> { destroyed.add("TEST-unit"); yield true; } + default -> throw new AssertionError("a late unit must not be waited on: " + method.getName()); + }); + + worker.onCommand(Message.of(pressed(pressed))); + + assertEquals(0, sent.size(), "no code for a row already ended"); + assertEquals(List.of("TEST-unit"), destroyed); + } + + /** + * A slow start inside the window: the code is shown, but the time left for the person is measured + * after the start, not before it, so the fourteen-minute limit is not stretched (review of PR #168). + */ + @Test + void theTimeLeftIsMeasuredAfterASlowStart() { + HarnessSignInWorker worker = worker(); + worker.harnesses = new HarnessRegistry(); + Instant pressed = Instant.parse("2026-09-24T12:00:00Z"); + java.util.concurrent.atomic.AtomicReference now = new java.util.concurrent.atomic.AtomicReference<>(pressed); + worker.clock = new java.time.Clock() { + @Override public java.time.ZoneId getZone() { return java.time.ZoneOffset.UTC; } + @Override public java.time.Clock withZone(java.time.ZoneId zone) { return this; } + @Override public Instant instant() { return now.get(); } + }; + List waited = new ArrayList<>(); + worker.runtime = (SignInRuntime) Proxy.newProxyInstance(SignInRuntime.class.getClassLoader(), + new Class[] { SignInRuntime.class }, (proxy, method, args) -> switch (method.getName()) { + case "start" -> { + now.set(pressed.plusSeconds(150)); + @SuppressWarnings("unchecked") + java.util.function.Consumer lines = (java.util.function.Consumer) args[1]; + lines.accept(" https://auth.openai.com/codex/device"); + lines.accept(" ABCD-12345"); + yield new SignInRuntime.Handle("TEST-sign-in", "TEST-unit"); + } + case "awaitExit" -> { waited.add((java.time.Duration) args[1]); yield new SignInRuntime.Exit.StillRunning(); } + case "destroy" -> true; + case "cancel" -> null; + default -> throw new AssertionError("unexpected: " + method.getName()); + }); + + worker.onCommand(Message.of(pressed(pressed))); + + var prompted = (HarnessSignInResult.Prompted) sent.getFirst(); + assertEquals(pressed.plusSeconds(840), prompted.expiresAt(), "the approval time ends 840s after the press"); + assertEquals(List.of(java.time.Duration.ofSeconds(690)), waited); + } + + /** + * Time passes between any two readings of the clock. The countdown shown must still end exactly at + * the approval deadline, never after it (review of PR #168). + */ + @Test + void theShownExpiryNeverPassesTheApprovalDeadline() { + HarnessSignInWorker worker = worker(); + worker.harnesses = new HarnessRegistry(); + Instant pressed = Instant.parse("2026-09-24T12:00:00Z"); + java.util.concurrent.atomic.AtomicReference now = new java.util.concurrent.atomic.AtomicReference<>(pressed); + java.util.concurrent.atomic.AtomicBoolean ticking = new java.util.concurrent.atomic.AtomicBoolean(); + worker.clock = new java.time.Clock() { + @Override public java.time.ZoneId getZone() { return java.time.ZoneOffset.UTC; } + @Override public java.time.Clock withZone(java.time.ZoneId zone) { return this; } + // After the unit starts, every reading is a second later than the one before it. + @Override public Instant instant() { return ticking.get() ? now.updateAndGet(at -> at.plusSeconds(1)) : now.get(); } + }; + worker.runtime = (SignInRuntime) Proxy.newProxyInstance(SignInRuntime.class.getClassLoader(), + new Class[] { SignInRuntime.class }, (proxy, method, args) -> switch (method.getName()) { + case "start" -> { + now.set(pressed.plusSeconds(150)); + ticking.set(true); + @SuppressWarnings("unchecked") + java.util.function.Consumer lines = (java.util.function.Consumer) args[1]; + lines.accept(" https://auth.openai.com/codex/device"); + lines.accept(" ABCD-12345"); + yield new SignInRuntime.Handle("TEST-sign-in", "TEST-unit"); + } + case "awaitExit" -> new SignInRuntime.Exit.StillRunning(); + case "destroy" -> true; + case "cancel" -> null; + default -> throw new AssertionError("unexpected: " + method.getName()); + }); + + worker.onCommand(Message.of(pressed(pressed))); + + var prompted = (HarnessSignInResult.Prompted) sent.getFirst(); + assertEquals(pressed.plusSeconds(840), prompted.expiresAt()); + } + + /** A replay of the start for a unit still running here must not fail that unit. */ + @Test + void aReplayedStartForAUnitStillRunningIsIgnored() throws Exception { + HarnessSignInWorker worker = worker(); + var running = HarnessSignInWorker.class.getDeclaredField("running"); + running.setAccessible(true); + @SuppressWarnings("unchecked") + Map units = (Map) running.get(worker); + units.put("TEST-sign-in", new SignInRuntime.Handle("TEST-unit", "TEST-unit")); + + worker.onCommand(Message.of(pressed(Instant.now().minusSeconds(900)))); + + assertEquals(0, sent.size()); + } + + private record Capturing(List sent) + implements Emitter> { + @Override + public CompletionStage send(Record record) { + sent.add(record.value()); + return CompletableFuture.completedFuture(null); + } + + @Override + public >> void send(M message) { + throw new UnsupportedOperationException("the worker sends records, not messages"); + } + + @Override + public void complete() { + } + + @Override + public void error(Exception e) { + } + + @Override + public boolean isCancelled() { + return false; + } + + @Override + public boolean hasRequests() { + return true; + } + } +} diff --git a/spire-run-worker/src/test/java/dev/codespire/runworker/MessagingChannelsAreDeclaredTest.java b/spire-run-worker/src/test/java/dev/codespire/runworker/MessagingChannelsAreDeclaredTest.java index 52444045..58d8c61f 100644 --- a/spire-run-worker/src/test/java/dev/codespire/runworker/MessagingChannelsAreDeclaredTest.java +++ b/spire-run-worker/src/test/java/dev/codespire/runworker/MessagingChannelsAreDeclaredTest.java @@ -103,6 +103,18 @@ void executionStaysOnTheOrderedWorkChannel() { assertEquals(WORK_CHANNEL, incoming.value()); } + /** A question sent while the worker's group was still being assigned must not be skipped. */ + @Test + void anImageQuestionSentBeforeTheWorkerJoinedIsStillAnswered() { + assertTrue(channelBlock(applicationYaml(), "harness-image-commands-in").contains("reset: earliest")); + } + + /** A start lost the same way left the sign-in PENDING for ever (review of PR #168). */ + @Test + void aSignInStartSentBeforeTheWorkerJoinedIsStillRead() { + assertTrue(channelBlock(applicationYaml(), "harness-sign-in-commands-in").contains("reset: earliest")); + } + private static Method onControl() { return method(RunControlListener.class, "onControl", dev.codespire.contract.command.RunCommand.class); } diff --git a/spire-runtime-docker/src/main/java/dev/codespire/runtime/docker/DockerSignInRuntime.java b/spire-runtime-docker/src/main/java/dev/codespire/runtime/docker/DockerSignInRuntime.java index 29c912d5..d99dc5ba 100644 --- a/spire-runtime-docker/src/main/java/dev/codespire/runtime/docker/DockerSignInRuntime.java +++ b/spire-runtime-docker/src/main/java/dev/codespire/runtime/docker/DockerSignInRuntime.java @@ -111,8 +111,7 @@ public Handle start(SignInUnitSpec spec, Consumer lines) { // By IDENTITY, never through the age-fenced sweep query: that one hides anything created in // the current second, which is precisely the container whose name just caused this // conflict. Missing it made the loser rethrow and fail a live owner's sign-in. - Handle existing = find(spec.unitId()).orElseThrow(() -> taken); - throw new SignInRuntime.AlreadyClaimed(existing, stillRunning(existing)); + throw claimedBy(spec.unitId(), () -> find(spec.unitId()), this::stillRunning); } // Everything after creation is guarded. A failure here used to leave a container with no // handle in anyone's hands — the one orphan that can hold a credential and that no finally @@ -224,6 +223,27 @@ private static boolean interrupted(Throwable failure) { } /** Whether the daemon still reports this container as running. False for gone, stopped or unknown. */ + /** + * What a name conflict means: the unit is HELD, whatever else can be learned about it. + * + *

The conflict itself is the proof; the lookup only adds detail. When the lookup fails, or finds + * nothing because the owner has just removed its container, the answer is still "held by someone + * else" — never an error. An error became UNIT_FAILED, and with two workers that ended a sign-in + * the other one was running and the operator could still approve (review of PR #168). A holder + * that no longer exists is resolved by the orchestrator's own deadlines, not here. + */ + static SignInRuntime.AlreadyClaimed claimedBy(String unitId, java.util.function.Supplier> lookup, + java.util.function.Predicate running) { + Handle existing; + try { + existing = lookup.get().orElse(null); + } catch (RuntimeException unreadable) { + existing = null; + } + if (existing == null) return new SignInRuntime.AlreadyClaimed(new Handle(unitId, "unknown"), false); + return new SignInRuntime.AlreadyClaimed(existing, running.test(existing)); + } + private boolean stillRunning(Handle handle) { try { return Boolean.TRUE.equals(client.inspectContainerCmd(handle.reference()).exec().getState().getRunning()); diff --git a/spire-runtime-docker/src/test/java/dev/codespire/runtime/docker/DockerSignInClaimTest.java b/spire-runtime-docker/src/test/java/dev/codespire/runtime/docker/DockerSignInClaimTest.java new file mode 100644 index 00000000..a548558f --- /dev/null +++ b/spire-runtime-docker/src/test/java/dev/codespire/runtime/docker/DockerSignInClaimTest.java @@ -0,0 +1,44 @@ +package dev.codespire.runtime.docker; + +import dev.codespire.runtime.SignInRuntime; +import org.junit.jupiter.api.Test; + +import java.util.Optional; + +import static org.junit.jupiter.api.Assertions.assertEquals; +import static org.junit.jupiter.api.Assertions.assertFalse; +import static org.junit.jupiter.api.Assertions.assertTrue; + +/** + * A name conflict proves another worker holds the unit; failing to look it up must not turn that into + * a failure that ends the owner's sign-in (review of PR #168). + */ +class DockerSignInClaimTest { + + private static final SignInRuntime.Handle OWNER = new SignInRuntime.Handle("TEST-sign-in", "TEST-container"); + + @Test + void aLookupThatFailsStillMeansHeld() { + SignInRuntime.AlreadyClaimed claimed = DockerSignInRuntime.claimedBy("TEST-sign-in", + () -> { throw new IllegalStateException("TEST daemon hiccup"); }, handle -> true); + + assertEquals("TEST-sign-in", claimed.existing().unitId()); + assertFalse(claimed.existingIsRunning(), "nothing known about it, so not claimed to be running"); + } + + /** The owner removed its container between the conflict and the lookup. */ + @Test + void aHolderThatJustVanishedStillMeansHeld() { + SignInRuntime.AlreadyClaimed claimed = DockerSignInRuntime.claimedBy("TEST-sign-in", Optional::empty, handle -> true); + + assertEquals("TEST-sign-in", claimed.existing().unitId()); + } + + @Test + void aHolderThatIsFoundIsNamedWithItsState() { + SignInRuntime.AlreadyClaimed claimed = DockerSignInRuntime.claimedBy("TEST-sign-in", () -> Optional.of(OWNER), handle -> true); + + assertEquals("TEST-container", claimed.existing().reference()); + assertTrue(claimed.existingIsRunning()); + } +} diff --git a/spire-ui/src/components/HarnessSubscriptionSignIn.tsx b/spire-ui/src/components/HarnessSubscriptionSignIn.tsx index e3099f62..94bcc2a7 100644 --- a/spire-ui/src/components/HarnessSubscriptionSignIn.tsx +++ b/spire-ui/src/components/HarnessSubscriptionSignIn.tsx @@ -14,6 +14,7 @@ const REASONS: Record = { sign_in_cancelled: 'The sign-in was cancelled.', sign_in_unit_failed: 'The sign-in tool did not start, or printed something this version cannot read. Nothing was stored.', sign_in_wrong_mode: 'That signed in as an API key, not a subscription. Add it as an API key instead.', + sign_in_not_started: 'No run worker showed a code within six minutes. Check that the run worker is running and has the agent image, then start again.', }; const sentence = (reason: string | null) => diff --git a/spire-ui/src/components/repositories/factory/BuildModelFields.test.tsx b/spire-ui/src/components/repositories/factory/BuildModelFields.test.tsx index 443c9b3e..a5b0e0c2 100644 --- a/spire-ui/src/components/repositories/factory/BuildModelFields.test.tsx +++ b/spire-ui/src/components/repositories/factory/BuildModelFields.test.tsx @@ -43,6 +43,14 @@ it('shows a model the harness runs but that has no price, and says what is missi expect(choices.find(choice => choice.value === 'TEST-fast')?.blocked).toBeNull(); }); +// Every option disabled makes the select ignore clicks and keys, which reads as broken (operator, +// 2026-09-23). The screen names the reason; see BuildStep. +it('marks every model blocked when none has a price', () => { + const choices = modelChoices('codex', KNOWN, [], REPORTED); + + expect(choices.every(choice => choice.blocked !== null)).toBe(true); +}); + // When the image did not say, the price list is all there is — what this screen offered before — and it // is not dressed up as the harness's own list. it('falls back to the price list, and says why, when the harness list is unknown', () => { diff --git a/spire-ui/src/components/repositories/factory/BuildStep.tsx b/spire-ui/src/components/repositories/factory/BuildStep.tsx index 06b74314..7ec90ea4 100644 --- a/spire-ui/src/components/repositories/factory/BuildStep.tsx +++ b/spire-ui/src/components/repositories/factory/BuildStep.tsx @@ -105,6 +105,10 @@ export default function BuildStep({ repositoryId, defaults, open, setOpen, chang setForm(previous => ({ ...previous, model }))} setEffort={effort => setForm(previous => ({ ...previous, effort }))} /> + {/* A select whose every option is disabled ignores clicks and keys alike, and looked broken (operator, + 2026-09-23). Say why nothing can be chosen, and where to fix it. */} + {form.harness && offered.length > 0 && offered.every(choice => choice.blocked !== null) &&

+ No model can be picked yet: each one is missing a price. Add the rates for one in Settings → LLM.

} {picked?.blocked &&

{form.model} has {picked.blocked}. {form.harness} reports those token types, and an API-key run needs a rate for each — or a mark in Settings → LLM that the vendor does not bill it.

} diff --git a/spire-ui/src/components/repositories/factory/RepositoryFactory.test.tsx b/spire-ui/src/components/repositories/factory/RepositoryFactory.test.tsx index e7a0ce6c..a19a134f 100644 --- a/spire-ui/src/components/repositories/factory/RepositoryFactory.test.tsx +++ b/spire-ui/src/components/repositories/factory/RepositoryFactory.test.tsx @@ -243,6 +243,19 @@ describe('build setup', () => { }); /** A model with no output price is refused at dispatch, so it cannot be chosen here either. */ + // Every option disabled: the select ignores clicks and keys, which the operator read as broken. + it('says why no model can be picked when every one lacks a price', async () => { + const types = ['INPUT', 'CACHED_INPUT', 'CACHE_WRITE', 'OUTPUT', 'REASONING']; + vi.mocked(build.buildOptions).mockResolvedValue({ harnesses: ['codex'], reportedTypes: { codex: types }, + models: { codex: { status: 'OK', offered: [{ slug: 'TEST-unpriced-only', displayName: 'TEST unpriced only', + defaultEffort: 'medium', efforts: ['medium'], visible: true, priority: 1 }] } } }); + renderFactory(); + await open(); + fireEvent.change(await screen.findByLabelText('Harness', field), { target: { value: 'codex' } }); + + expect(await screen.findByText(/No model can be picked yet/)).toBeInTheDocument(); + }); + it('offers an unpriced model as unselectable and says why', async () => { renderFactory(); await open();