Skip to content

fix(server): restore queued turn identity after restart - #10611

Closed
saphid wants to merge 9 commits into
pingdotgg:mainfrom
saphid:repair/d6096318-restart-projection-stack
Closed

saphid wants to merge 9 commits into
pingdotgg:mainfrom
saphid:repair/d6096318-restart-projection-stack

Conversation

@saphid

@saphid saphid commented Sep 7, 2026 •

Copy link
Copy Markdown
Contributor

Rebased onto the canonical #10608 head 8cc4ce35ff (squashed dependency). Latest head ed4f1c9b8b adds request-sequence generation bounds on top of the dependency’s durable pending/submitted turn-start rows.

What changed

  • Snapshot reconstruction hydrates the newest persisted pending turn-start message ID into the full snapshot, command read model, and per-thread detail, so a queued user message survives restart with its original identity.
  • thread.turn.start.acknowledge carries an optional requestSequence (alongside the dependency’s sourceProposedPlan); migration 053 persists it on placeholder rows.
  • Turn-start clears are bounded by request sequence: provider failures, compaction cancels, session-stop cleanup, and reactor-startup sweeps pass a throughRequestSequence bound so a stale answer cannot retire a newer re-requested generation of the same message.
  • A re-requested send reuses its pending placeholder: the row is re-stamped to the newest request sequence, and a submitted row no longer blocks a fresh generation.
  • Session stop now fails both pending and submitted turn starts that never surfaced a turn, reported once per request.
  • packages/client-runtime thread reducer mirrors the same acknowledgement generation matching so local and remote projections agree.

Why

A turn may be accepted and persisted immediately before the server restarts while still lacking a running provider lifecycle. The previous snapshot query omitted that pending row from restart hydration, losing the queued message or manufacturing a duplicate provider turn. The dependency dedupes pending placeholders by message id; without generation bounds, a stale acknowledgement or failure for an older request could claim or clear a newer re-requested send.

Verification

Focused tests at head ed4f1c9b8b on dependency 8cc4ce35ff:

  • ProjectionPipeline.test.ts — 45 pass (persisted-SQL coverage for generation-bounded clears, dedupe re-stamp, submitted-row adoption across a running turn, unsequenced-ack fallback)
  • ProviderCommandReactor.test.ts — 74 pass (startup sweep of pre-subscription pending starts, stop-after-send cleanup, acknowledgement retry until durable)
  • projector.test.ts, ProjectionSnapshotQuery.test.ts, serverRuntimeStartup(.reconcile).test.ts, CheckpointReactor.test.ts, GitVcsDriver.test.ts, ProviderRuntimeIngestion.test.ts, decider.settled.test.ts — 355 pass
  • threadReducer.test.ts — 67 pass
  • tsc --noEmit clean for apps/server, packages/client-runtime, packages/contracts; scoped lint 0 errors; git diff --check clean.

Independent cross-provider review was skipped: Codex weekly headroom reported 0% remaining, below the policy threshold.

Dependency

Stacked on #10608 at head 8cc4ce35ff (repair/d6096318-checkpoint-revert), which owns durable submitted-start persistence, the compaction serialization, and checkpoint-revert plumbing this builds on. Draft until #10608 merges.

Checklist

  • Nonvisual persistence/restart change — screenshots and interaction video are not applicable
  • One concern: queued turn-start identity and its request-generation bounds

Model: SWE-2 Max via T3 Code

Coordination trace: T3 thread 6c0b4244-aecc-49fe-9393-c37a86e87972

@github-actions github-actions Bot added vouch:trusted PR author is trusted by repo permissions or the VOUCHED list. size:L 100-499 changed lines (additions + deletions). labels Sep 7, 2026
Comment thread apps/server/src/orchestration/Layers/CheckpointReactor.ts Outdated
Comment thread apps/server/src/orchestration/Layers/ProviderCommandReactor.ts Outdated
Comment thread apps/server/src/server.ts Outdated
@macroscopeapp

macroscopeapp Bot commented Sep 8, 2026

Copy link
Copy Markdown
Contributor

Effect service convention findings were posted inline.

Posted via Macroscope — Effect Service Conventions

@saphid
saphid force-pushed the repair/d6096318-restart-projection-stack branch from decb3cc to 9dbecbd Compare September 8, 2026 05:54
@github-actions github-actions Bot added size:XXL 1,000+ changed lines (additions + deletions). and removed size:L 100-499 changed lines (additions + deletions). labels Sep 8, 2026
Comment thread apps/server/src/orchestration/Layers/ProjectionPipeline.ts
Comment thread apps/server/src/serverRuntimeStartup.ts Outdated
Comment thread apps/server/src/orchestration/decider.ts Outdated
@saphid
saphid force-pushed the repair/d6096318-restart-projection-stack branch from aec4382 to 1bd48eb Compare September 8, 2026 07:19
Comment thread apps/server/src/orchestration/Layers/ProviderCommandReactor.ts Outdated
Comment thread apps/server/src/persistence/Layers/ProjectionTurns.ts
@maxibotstef

Copy link
Copy Markdown

Adding a focused reproduction for the accepted-before-reactor-subscription window related to this stack, rather than opening a duplicate issue.

Tested: official tag v0.0.41-nightly.20260911.1520, commit c52b8d96e4b34201f19b5e5bb12c6b2a77bfaa9a, using the existing ProviderCommandReactor.test.ts harness and its fake provider. No production code was changed. This has not been tested against this PR or its #10608 dependency.

Reproduction

  1. Create the harness project and thread with a fixed model selection.
  2. Before reactor.start(), dispatch one thread.turn.start with stable command and message IDs. Record its accepted receipt.
  3. Start the reactor and replay the identical command.
  4. Drain the reactor and check the receipt and fake provider calls.
Assertion Observed
First dispatch accepted Yes
Replay returns the original receipt sequence Yes
Provider continuations after replay 0, expected 1 for recovery
Existing post-subscription positive control Pass

The failing assertion is expect(harness.sendTurn).toHaveBeenCalledTimes(1): expected 1, observed 0. The two focused tests were independently rerun together with the same result: one expected failure, one passing control, 61 skipped.

Scope of the evidence

This is an in-process pre-subscription test, modeling a persisted request that predates a new reactor subscription. It does not kill/restart an app or provider process, exercise the entire server startup/reconciliation path, or test the “Continue threads after restarts” setting. In particular, it does not show that normal public requests are accepted before server startup completes, or that this PR fails.

It isolates a narrower boundary: replaying an accepted command cannot itself recover a provider wake that the reactor never observed. The receipt path returns the existing sequence without republishing the event (engine); the reactor subscribes to a hot event stream (reactor).

This matters for an external client that retries the same command after a disconnect: a successful command receipt is distinct from recovery of the provider continuation. Would this be a useful additional regression case for the combined #10608/#10611 startup path, alongside the persisted-SQL tests already described here?

Reproduce on the released tag

Apply the test-only patch below in an isolated checkout of the pinned tag, install dependencies using the repository's documented setup, and run this from apps/server:

./node_modules/.bin/vp test run src/orchestration/Layers/ProviderCommandReactor.test.ts -t 'recovers one provider continuation when an accepted turn is replayed after reactor subscription starts|reacts to thread.turn.start by ensuring session and sending provider turn'

Environment: macOS arm64, Node 26.5.0, Vitest 4.1.11, fake provider only. The runner warned that Node was outside the preferred 24.13.1 version; both cases ran under the same environment. This is a focused seam reproduction, not release qualification. Fixture identifiers, message text and the test title in the patch are normalized for publication; test logic is unchanged.

Test-only patch against the pinned official tag
--- a/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts
+++ b/apps/server/src/orchestration/Layers/ProviderCommandReactor.test.ts
@@ -173,6 +173,7 @@
     readonly unreadableHistory?: boolean;
     readonly titleRegenerationCompletionDispatchFailures?: number;
     readonly titleRegenerationBeforeStart?: "one" | "two";
+    readonly dispatchTurnBeforeReactorStart?: boolean;
     readonly serverActivation?: Effect.Effect<void>;
     readonly beforeReadySessionDispatch?: () => Effect.Effect<void>;
     readonly compactThreadEffect?: () => Effect.Effect<void, ProviderAdapterRequestError>;
@@ -566,6 +567,25 @@
       );
     }
 
+    const preReactorTurnStartCommand = {
+      type: "thread.turn.start" as const,
+      commandId: CommandId.make("cmd-restart-window"),
+      threadId: ThreadId.make("thread-1"),
+      message: {
+        messageId: MessageId.make("message-restart-window"),
+        role: "user" as const,
+        text: "[AUTOMATION] Evaluate candidate candidate-1.",
+        attachments: [],
+      },
+      modelSelection,
+      interactionMode: DEFAULT_PROVIDER_INTERACTION_MODE,
+      runtimeMode: "approval-required" as const,
+      createdAt: now,
+    };
+    const preReactorTurnStartReceipt = input?.dispatchTurnBeforeReactorStart
+      ? await Effect.runPromise(engine.dispatch(preReactorTurnStartCommand))
+      : null;
+
     scope = await Effect.runPromise(Scope.make("sequential"));
     await Effect.runPromise(
       reactor
@@ -608,6 +628,8 @@
       generateThreadTitle,
       runtimeSessions,
       stateDir,
+      preReactorTurnStartCommand,
+      preReactorTurnStartReceipt,
       drain,
       runEffect,
       get titleRegenerationCompletionDispatchAttempts() {
@@ -615,6 +637,31 @@
       },
     };
   }
+
+  effectIt.effect(
+    "recovers one provider continuation when an accepted turn is replayed after reactor subscription starts",
+    () =>
+      Effect.gen(function* () {
+        const harness = yield* Effect.promise(() =>
+          createHarness({ dispatchTurnBeforeReactorStart: true }),
+        );
+
+        expect(harness.preReactorTurnStartReceipt).not.toBeNull();
+        expect(harness.sendTurn).toHaveBeenCalledTimes(0);
+
+        const replayReceipt = yield* harness.engine.dispatch(harness.preReactorTurnStartCommand);
+        yield* Effect.promise(() => harness.drain());
+
+        expect(replayReceipt.sequence).toBe(harness.preReactorTurnStartReceipt?.sequence);
+        expect(harness.sendTurn).toHaveBeenCalledTimes(1);
+        expect(harness.sendTurn).toHaveBeenCalledWith(
+          expect.objectContaining({
+            threadId: ThreadId.make("thread-1"),
+            input: "[AUTOMATION] Evaluate candidate candidate-1.",
+          }),
+        );
+      }),
+  );
 
   effectIt.effect.each(["new", "ready", "stopped"] as const)(
     "handles sign-out for a %s thread before worktree repair, text helpers, or startup",

Prepared with Codex on behalf of @maxibotstef.

@saphid
saphid force-pushed the repair/d6096318-restart-projection-stack branch from 67e2ffc to ed4f1c9 Compare September 16, 2026 03:31
@saphid

saphid commented Sep 16, 2026

Copy link
Copy Markdown
Contributor Author

Thanks for the careful reproduction — the seam you isolated is real, and the combined stack now has dedicated coverage for it.

In the rebased stack (ed4f1c9b8b on #10608 8cc4ce35ff), ProviderCommandReactor.start() captures latestSequence as the handoff bound, lists durable pending turn starts at that bound via listPendingTurnStarts(handoffSequence), and clears each with a provider.turn.start.failed activity carrying throughRequestSequence: handoffSequence before processing live events (sequence > handoffSequence). So a start accepted before the subscription is resolved as a bounded, user-visible failure rather than silently dropped — and the replayed command still returns the original receipt through commandId dedupe, exactly as you observed.

Your proposed case is already exercised in ProviderCommandReactor.test.ts:

  • marks every turn start left pending across reactor startup as failed dispatches thread.turn.start before reactor.start() (the turnStartBeforeStart harness option — same pre-subscription window as your patch), then asserts sendTurn is never called and both durable rows get provider.turn.start.failed activities.
  • continues clearing startup turn starts after one failure dispatch fails covers a mid-sweep dispatch failure.
  • processes a turn committed during the startup snapshot handoff covers the adjacent boundary: an event committed between subscription and snapshot is still processed live.

The one place behavior differs from your expected assertion is deliberate: the stack fails the pre-subscription start (sendTurn stays 0) instead of resuming it (sendTurn = 1). A hot-stream gap cannot distinguish "reactor not yet subscribed" from "provider may have already accepted this send" — auto-resending risks duplicating a turn, so the row is resolved and the user is told to resend. Resuming queued sends after restart is a product decision on the "Continue threads after restarts" surface, outside this PR's identity/durability scope.

Given the existing harness coverage of the same window, I'm not adding the variant as a separate regression case here — but flagging it was useful confirmation that the pre-subscription boundary needed to be explicit. Happy to revisit if you see the covered assertions diverging from the behavior you measured.

saphid and others added 9 commits September 16, 2026 14:45
Turn-start rows could outlive their queue and survive a session stop —
for example an in-flight /compact that finished compacting but was still
blocked resuming, or a start the provider acknowledged but never began.
The rows then stayed durable forever, guarded settle/snooze, and would
be resurrected as queued turns after restart.

On thread.session-stop-requested, report every remaining pending or
submitted turn start by its own message id so the projection deletes
exactly those rows.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Extends session-stop cleanup beyond pending turn starts: a start the
provider already accepted (state = 'submitted') but never began would
otherwise survive as a durable row, blocking settle, snooze, and
checkpoint revert until a restart reconciled it.

Both pending and submitted rows are now queried through the stop
event's own sequence, so a turn-start request decided after the stop is
never canceled by it, and a row acknowledged between the two reads is
reported only once via its message-correlated failure activity.

Also aligns GitVcsDriver's restore test with the preserve contract the
stack established: restores never clean files outside the checkpoint
diff, so untracked files survive while index-tracked staged files are
rewritten with the index.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Clearing a turn start by message id alone could retire a newer request:
a same-message-id re-request landing between the failure's cause and its
projection would be deleted too, stranding a live queued start until
restart reconciliation.

Stamp each projection row with the sequence of the request event that
created it (migration 053), and carry a request-sequence bound through
provider turn-start failure, sign-out, and compaction-cancel activities.
Bounded clears — in the projection delete, the in-memory projector, and
the client reducer — then retire only the request generation they answer;
a re-requested start survives. Startup reconciliation bounds each orphan
report by the row's own request sequence. Rows written before the column
existed carry NULL and count as older than every bound, preserving legacy
unbounded-clear behavior.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
The pending-start insert dedupes by message id so compaction re-dispatch
stays idempotent. With request-sequence bounds, a re-requested send must
move that placeholder forward: the row now answers the newest request
generation, a submitted row no longer blocks a fresh generation, and a
bounded failure cannot retire a re-requested start by its older sequence.

Generated with [Devin](https://devin.ai)

Co-Authored-By: Devin <158243242+devin-ai-integration[bot]@users.noreply.github.com>
Recovery dispatches thread.title.refine and bounded turn-start failures
whose follow-up events land on the domain PubSub. The subscription was
captured before the handoff sequence, but the consumer only forked after
recovery completed, so drain() could return before recovery-emitted
events reached the workers. Fork the filtered consumer first — the
subscription buffers events from subscribe-time, and the handoff filter
still excludes pre-start sequences.
@saphid
saphid force-pushed the repair/d6096318-restart-projection-stack branch from ed4f1c9 to 9a76af5 Compare September 16, 2026 05:10
@juliusmarminge

Copy link
Copy Markdown
Member

Thanks for the PR. We're not taking changes to the orchestration and provider layers right now: that part of the server is being rewritten for V2, and merging into the current code would either conflict with or be thrown away by that work.

Closing for now. If this is still an issue once V2 lands, please reopen (or open a fresh PR against the new code) and we'll take a proper look.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

size:XXL 1,000+ changed lines (additions + deletions). vouch:trusted PR author is trusted by repo permissions or the VOUCHED list.

Projects

None yet

Development

Successfully merging this pull request may close these issues.

3 participants