diff --git a/apps/server/src/orchestration/PullRequestSyncReactor.test.ts b/apps/server/src/orchestration/PullRequestSyncReactor.test.ts index 414f76f324..1406eae239 100644 --- a/apps/server/src/orchestration/PullRequestSyncReactor.test.ts +++ b/apps/server/src/orchestration/PullRequestSyncReactor.test.ts @@ -1,9 +1,11 @@ import { + EventId, ProjectId, ProviderInstanceId, PullRequestOperationError, ThreadId, type OrchestrationCommand, + type OrchestrationEvent, type OrchestrationProjectShell, type OrchestrationShellSnapshot, type OrchestrationThreadShell, @@ -18,6 +20,7 @@ import * as Crypto from "effect/Crypto"; import * as Deferred from "effect/Deferred"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; +import * as PubSub from "effect/PubSub"; import * as Queue from "effect/Queue"; import * as Ref from "effect/Ref"; import * as Stream from "effect/Stream"; @@ -31,6 +34,7 @@ import { } from "./Services/OrchestrationEngine.ts"; import { ProjectionSnapshotQuery } from "./Services/ProjectionSnapshotQuery.ts"; import * as PullRequestSyncReactor from "./PullRequestSyncReactor.ts"; +import { resolveAutoSettlementAt } from "./ThreadSettlementPolicy.ts"; const NOW = "2026-08-28T12:00:00.000Z"; const PROJECT_ID = ProjectId.make("sync-project"); @@ -152,6 +156,7 @@ function makeSummary( } interface HarnessOptions { + readonly onDispatch?: (command: SyncCommand | LinkCommand) => Effect.Effect; readonly invalidate?: PullRequestService["Service"]["invalidate"]; readonly snapshot: OrchestrationShellSnapshot; readonly summary?: ( @@ -165,6 +170,7 @@ interface HarnessOptions { const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options: HarnessOptions) { const activation = yield* Deferred.make(); const snapshots = yield* Ref.make(options.snapshot); + const events = yield* PubSub.unbounded(); const snapshotReads = yield* Queue.unbounded(); const syncCommands = yield* Ref.make>([]); const linkCommands = yield* Ref.make>([]); @@ -188,11 +194,13 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options: const dispatch: OrchestrationEngineShape["dispatch"] = (command) => { if (command.type === "thread.pull-request-link.sync") { return Ref.update(syncCommands, (recorded) => [...recorded, command]).pipe( + Effect.andThen(options.onDispatch?.(command) ?? Effect.void), Effect.as({ sequence: 1 }), ); } if (command.type === "thread.pull-request.link") { return Ref.update(linkCommands, (recorded) => [...recorded, command]).pipe( + Effect.andThen(options.onDispatch?.(command) ?? Effect.void), Effect.as({ sequence: 1 }), ); } @@ -213,6 +221,9 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options: readEvents: () => Stream.empty, dispatch, streamDomainEvents: Stream.empty, + subscribeDomainEvents: PubSub.subscribe(events).pipe( + Effect.map((subscription) => Stream.fromSubscription(subscription)), + ), latestSequence: Effect.succeed(0), }), Layer.succeed(ServerActivation, Deferred.await(activation)), @@ -220,6 +231,7 @@ const makeHarness = Effect.fn("makePullRequestSyncHarness")(function* (options: ); return { + events, activation, snapshots, snapshotReads, @@ -274,6 +286,42 @@ function applySync( } describe("PullRequestSyncReactor", () => { + it.effect("syncs a newly linked merged PR without waiting for the periodic sweep", () => + Effect.scoped( + Effect.gen(function* () { + yield* TestClock.setTime(Date.parse(NOW)); + const fixture = yield* makeHarness({ + snapshot: makeSnapshot([makeThread("one")]), + summary: (input) => + Effect.succeed(makeSummary(input, { state: "merged", mergedAt: NOW })), + }); + yield* Effect.gen(function* () { + const reactor = yield* startAndSweep(fixture); + const link = makeLink(42); + yield* Ref.set( + fixture.snapshots, + makeSnapshot([makeThread("one", { pullRequests: [link] })]), + ); + yield* PubSub.publish(fixture.events, { + type: "thread.pull-request-linked", + sequence: 2, + eventId: EventId.make("linked"), + aggregateKind: "thread", + aggregateId: ThreadId.make("one"), + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + payload: { threadId: ThreadId.make("one"), link, updatedAt: NOW }, + }); + yield* Queue.take(fixture.snapshotReads); + yield* reactor.drain; + assert.strictEqual((yield* Ref.get(fixture.syncCommands))[0]?.snapshot.state, "merged"); + }).pipe(Effect.provide(fixture.layer)); + }), + ), + ); it.effect("retries a failed stack read after the summary becomes terminal", () => Effect.scoped( Effect.gen(function* () { @@ -303,6 +351,7 @@ describe("PullRequestSyncReactor", () => { yield* Effect.gen(function* () { const reactor = yield* startAndSweep(fixture); const commands = yield* Ref.get(fixture.syncCommands); + assert.deepStrictEqual(commands, []); yield* Ref.update(fixture.snapshots, (snapshot) => applySync(snapshot, commands)); yield* sweepAgain(fixture, reactor); assert.strictEqual(attempts, 2); @@ -315,6 +364,85 @@ describe("PullRequestSyncReactor", () => { ), ); + it.effect("retries a failed sibling link before publishing a terminal snapshot", () => + Effect.scoped( + Effect.gen(function* () { + yield* TestClock.setTime(Date.parse(NOW)); + let failSibling = true; + const fixture = yield* makeHarness({ + snapshot: makeSnapshot([ + makeThread("one", { pullRequests: [makeLink(7, { state: "closed" })] }), + ]), + summary: (input) => + Effect.succeed(makeSummary(input, { state: "merged", mergedAt: NOW })), + stack: () => + Effect.succeed({ + id: "stack", + number: 7, + url: "https://github.com/owner/repository/stacks/7", + base: "main", + layers: [ + { number: 7, headBranch: "feature", state: "merged" }, + { number: 8, headBranch: "sibling", state: "open" }, + ], + }), + onDispatch: (command) => + command.type === "thread.pull-request.link" && failSibling + ? Effect.die("temporary link failure") + : Effect.void, + }); + yield* Effect.gen(function* () { + const reactor = yield* startAndSweep(fixture); + assert.deepStrictEqual(yield* Ref.get(fixture.syncCommands), []); + failSibling = false; + yield* sweepAgain(fixture, reactor); + assert.strictEqual((yield* Ref.get(fixture.syncCommands))[0]?.snapshot.state, "merged"); + assert.strictEqual((yield* Ref.get(fixture.linkCommands)).at(-1)?.number, 8); + }).pipe(Effect.provide(fixture.layer)); + }), + ), + ); + + it.effect("concurrent stack reads persist a shared sibling once before syncing both roots", () => + Effect.scoped( + Effect.gen(function* () { + const readsReady = yield* Deferred.make(); + let reads = 0; + let linked = false; + const fixture = yield* makeHarness({ + snapshot: makeSnapshot([makeThread("one", { pullRequests: [makeLink(7), makeLink(8)] })]), + stack: () => + Effect.gen(function* () { + if (++reads === 2) yield* Deferred.succeed(readsReady, undefined); + yield* Deferred.await(readsReady); + return { + id: "stack", + number: 7, + url: "https://github.com/owner/repository/stacks/7", + base: "main", + layers: [{ number: 9, headBranch: "sibling", state: "open" as const }], + }; + }), + onDispatch: (command) => + Effect.gen(function* () { + if (command.type === "thread.pull-request.link") { + yield* Effect.yieldNow; + assert.strictEqual(linked, false); + linked = true; + } else { + assert.strictEqual(linked, true); + } + }), + }); + yield* Effect.gen(function* () { + yield* startAndSweep(fixture); + assert.strictEqual((yield* Ref.get(fixture.linkCommands)).length, 1); + assert.strictEqual((yield* Ref.get(fixture.syncCommands)).length, 2); + }).pipe(Effect.provide(fixture.layer)); + }), + ), + ); + it.effect("explicit refresh reads a changed stack even when its PR summary is unchanged", () => Effect.scoped( Effect.gen(function* () { @@ -631,20 +759,38 @@ describe("PullRequestSyncReactor", () => { base: "main", layers: [ { number: 41, headBranch: "layer-1", state: "merged" }, - { number: 42, headBranch: "layer-2", state: "open" }, + { number: 42, headBranch: "layer-2", state: "merged" }, { number: 43, headBranch: "layer-3", state: "open" }, ], }; + let thread = makeThread("one", { + pullRequests: [ + makeLink(42, { state: "open" }), + makeLink(41, { state: "merged" }, { source: "stack-dismissed" }), + ], + }); const fixture = yield* makeHarness({ - snapshot: makeSnapshot([ - makeThread("one", { - pullRequests: [ - makeLink(42), - makeLink(41, { state: "merged" }, { source: "stack-dismissed" }), - ], - }), - ]), + snapshot: makeSnapshot([thread]), + summary: (input) => + Effect.succeed(makeSummary(input, { state: "merged", mergedAt: NOW })), stack: () => Effect.succeed(stack), + onDispatch: (command) => + Effect.sync(() => { + thread = + command.type === "thread.pull-request.link" + ? { ...thread, pullRequests: [...thread.pullRequests, makeLink(command.number)] } + : applySync(makeSnapshot([thread]), [command]).threads[0]!; + // Every projected event may wake settlement, including the terminal root update. + assert.isNull( + resolveAutoSettlementAt({ + thread, + pullRequest: null, + now: NOW, + autoSettleAfterDays: null, + autoSettleOnMerge: true, + }), + ); + }), }); yield* Effect.gen(function* () { diff --git a/apps/server/src/orchestration/PullRequestSyncReactor.ts b/apps/server/src/orchestration/PullRequestSyncReactor.ts index 0a39fa5d91..4afd74987f 100644 --- a/apps/server/src/orchestration/PullRequestSyncReactor.ts +++ b/apps/server/src/orchestration/PullRequestSyncReactor.ts @@ -22,6 +22,8 @@ import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Schedule from "effect/Schedule"; import type * as Scope from "effect/Scope"; +import * as Semaphore from "effect/Semaphore"; +import * as Stream from "effect/Stream"; import * as PullRequestService from "../pullRequest/PullRequestService.ts"; import { forkParked } from "../serverActivation.ts"; @@ -149,7 +151,7 @@ export const make = Effect.gen(function* () { (cause: Cause.Cause): Effect.Effect => Cause.hasInterruptsOnly(cause) ? Effect.failCause(cause) : Effect.logWarning(message, fields); - const sweep = Effect.fn("PullRequestSyncReactor.sweep")(function* () { + const sweep = Effect.fn("PullRequestSyncReactor.sweep")(function* (requestedKey?: string) { const snapshot = yield* snapshots.getShellSnapshot(); const now = yield* DateTime.now; const nowMs = DateTime.toEpochMillis(now); @@ -173,6 +175,7 @@ export const make = Effect.gen(function* () { // Layers auto-linked this sweep, so two links of one thread that share a // stack do not both try to add the same sibling. const linkedThisSweep = new Set(); + const persistence = yield* Semaphore.make(1); const syncEntry = Effect.fn("PullRequestSyncReactor.syncEntry")(function* ( entry: LinkEntry, @@ -185,21 +188,8 @@ export const make = Effect.gen(function* () { link.snapshot === null || !snapshotFieldsEqual(link.snapshot, fields) || !stacksEqual(link.stack, nextStack); - if (changed) { - const uuid = yield* crypto.randomUUIDv4; - yield* engine.dispatch({ - type: "thread.pull-request-link.sync", - commandId: CommandId.make(`server:pr-sync:${thread.id}:${uuid}`), - threadId: thread.id, - host: link.host, - repository: link.repository, - number: link.number, - snapshot: { ...fields, syncedAt: nowIso }, - stack: nextStack, - }); - } - if (fetchedStack === null || fetchedStack.stack === null) return; - for (const layer of fetchedStack.stack.layers) { + // Persist discovered siblings before a terminal snapshot can trigger settlement. + for (const layer of fetchedStack?.stack?.layers ?? []) { const layerKey = { host: link.host, repository: link.repository, number: layer.number }; const dedupeKey = `${thread.id}:${threadPullRequestKeyOf(layerKey)}`; if (linkedThisSweep.has(dedupeKey)) continue; @@ -211,25 +201,29 @@ export const make = Effect.gen(function* () { } const url = siblingPullRequestUrl(link.url, layer.number); if (url === null) continue; + const uuid = yield* crypto.randomUUIDv4; + yield* engine.dispatch({ + type: "thread.pull-request.link", + commandId: CommandId.make(`server:pr-stack-link:${thread.id}:${uuid}`), + threadId: thread.id, + ...layerKey, + url, + source: "stack", + }); linkedThisSweep.add(dedupeKey); + } + if (changed) { const uuid = yield* crypto.randomUUIDv4; - yield* engine - .dispatch({ - type: "thread.pull-request.link", - commandId: CommandId.make(`server:pr-stack-link:${thread.id}:${uuid}`), - threadId: thread.id, - ...layerKey, - url, - source: "stack", - }) - .pipe( - Effect.catchCause( - logSkipped("pull request stack layer link skipped", { - threadId: thread.id, - number: layer.number, - }), - ), - ); + yield* engine.dispatch({ + type: "thread.pull-request-link.sync", + commandId: CommandId.make(`server:pr-sync:${thread.id}:${uuid}`), + threadId: thread.id, + host: link.host, + repository: link.repository, + number: link.number, + snapshot: { ...fields, syncedAt: nowIso }, + stack: nextStack, + }); } }); @@ -270,8 +264,11 @@ export const make = Effect.gen(function* () { ) : null; if (needsStack) { - if (fetchedStack === null) retryStacks.add(key); - else retryStacks.delete(key); + if (fetchedStack === null) { + retryStacks.add(key); + return; + } + retryStacks.delete(key); } // The host answered, so the cadence clock ticks even if a dispatch below is rejected. lastSyncedAt.set(key, nowMs); @@ -281,9 +278,13 @@ export const make = Effect.gen(function* () { entries, (entry) => syncEntry(entry, fields, fetchedStack).pipe( - Effect.catchCause( - logSkipped("pull request sync skipped", { threadId: entry.thread.id, key }), - ), + persistence.withPermits(1), + Effect.catchCause((cause) => { + if (!Cause.hasInterruptsOnly(cause)) retryStacks.add(key); + return logSkipped("pull request sync skipped", { threadId: entry.thread.id, key })( + cause, + ); + }), ), { discard: true }, ); @@ -292,7 +293,7 @@ export const make = Effect.gen(function* () { yield* Effect.forEach( groups, ([key, entries]) => - isDue(key, entries, nowMs) + (requestedKey === undefined || requestedKey === key) && isDue(key, entries, nowMs) ? syncGroup(key, entries).pipe( Effect.catchCause(logSkipped("pull request sync skipped", { key })), ) @@ -301,13 +302,19 @@ export const make = Effect.gen(function* () { ); }); - const worker = yield* makeDrainableWorker(() => - sweep().pipe(Effect.catchCause(logSkipped("pull request sync sweep failed", {}))), + const worker = yield* makeDrainableWorker((key: string | undefined) => + sweep(key).pipe(Effect.catchCause(logSkipped("pull request sync sweep failed", {}))), ); const start: PullRequestSyncReactor["Service"]["start"] = Effect.fn( "PullRequestSyncReactor.start", )(function* () { + const events = yield* engine.subscribeDomainEvents; + yield* forkParked( + Stream.runForEach(events, (event) => + event.type === "thread.pull-request-linked" ? requestSync(event.payload.link) : Effect.void, + ), + ); yield* forkParked( Effect.gen(function* () { yield* worker.enqueue(undefined); @@ -318,8 +325,9 @@ export const make = Effect.gen(function* () { const requestSync: PullRequestSyncReactor["Service"]["requestSync"] = (key) => Effect.suspend(() => { - requested.set(threadPullRequestKeyOf(key), ++requestGeneration); - return worker.enqueue(undefined); + const syncKey = threadPullRequestKeyOf(key); + requested.set(syncKey, ++requestGeneration); + return worker.enqueue(syncKey); }); return { start, drain: worker.drain, requestSync } satisfies PullRequestSyncReactor["Service"]; diff --git a/apps/server/src/orchestration/ThreadSettlementReactor.test.ts b/apps/server/src/orchestration/ThreadSettlementReactor.test.ts index 455cb37a61..0e23ae9498 100644 --- a/apps/server/src/orchestration/ThreadSettlementReactor.test.ts +++ b/apps/server/src/orchestration/ThreadSettlementReactor.test.ts @@ -1,10 +1,12 @@ import { DEFAULT_SERVER_SETTINGS, + EventId, ProjectId, ProviderInstanceId, PullRequestOperationError, ThreadId, type OrchestrationCommand, + type OrchestrationEvent, type OrchestrationProjectShell, type OrchestrationShellSnapshot, type OrchestrationThreadShell, @@ -171,6 +173,7 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options: const settingsReads = yield* Queue.unbounded(); const settingsChanges = yield* PubSub.unbounded(); const mergedPullRequests = yield* PubSub.unbounded(); + const domainEvents = yield* PubSub.unbounded(); const commands = yield* Ref.make>([]); const branchCalls = yield* Ref.make< ReadonlyArray<{ readonly cwd: string; readonly branch: string }> @@ -261,6 +264,9 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options: readEvents: () => Stream.empty, dispatch, streamDomainEvents: Stream.empty, + subscribeDomainEvents: PubSub.subscribe(domainEvents).pipe( + Effect.map((subscription) => Stream.fromSubscription(subscription)), + ), latestSequence: Effect.succeed(0), }), Layer.succeed(ServerSettingsService, serverSettings), @@ -283,6 +289,7 @@ const makeHarness = Effect.fn("makeThreadSettlementHarness")(function* (options: summaryRecovery, invalidatedCwds, updateSettings, + publishEvent: (event: OrchestrationEvent) => PubSub.publish(domainEvents, event), publishMerge: PubSub.publish(mergedPullRequests, { projectId: PROJECT_ID, repository: "owner/repository", @@ -306,12 +313,12 @@ const startHarness = Effect.fn("startThreadSettlementHarness")(function* ( describe("ThreadSettlementReactor", () => { it.effect( - "settles all-terminal links from snapshots and keeps open or unsynced links active", + "settles synced terminal links immediately while open, unsynced, or running threads wait", () => Effect.scoped( Effect.gen(function* () { yield* TestClock.setTime(Date.parse(NOW)); - const link = (number: number, state: "open" | "merged" | null) => ({ + const link = (number: number, state: "open" | "closed" | "merged" | null) => ({ host: "example.test", repository: "owner/repository", number, @@ -331,14 +338,30 @@ describe("ThreadSettlementReactor", () => { updatedAt: NOW, syncedAt: NOW, mergedAt: state === "merged" ? NOW : null, + closedAt: state === "closed" ? NOW : null, }, }); + const runningSession = { + threadId: ThreadId.make("running"), + status: "running" as const, + providerName: "Codex", + runtimeMode: "full-access" as const, + activeTurnId: null, + lastError: null, + updatedAt: NOW, + }; + const threads = [ + makeThread("merged", { pullRequests: [link(1, "open"), link(2, "merged")] }), + makeThread("closed", { pullRequests: [link(1, "open")] }), + makeThread("open", { pullRequests: [link(1, "open"), link(2, "open")] }), + makeThread("unsynced", { pullRequests: [link(1, "open"), link(2, null)] }), + makeThread("running", { + pullRequests: [link(1, "open")], + session: runningSession, + }), + ]; const fixture = yield* makeHarness({ - snapshot: makeSnapshot([ - makeThread("merged", { pullRequests: [link(1, "merged"), link(2, "merged")] }), - makeThread("open", { pullRequests: [link(1, "merged"), link(2, "open")] }), - makeThread("unsynced", { pullRequests: [link(1, "merged"), link(2, null)] }), - ]), + snapshot: makeSnapshot(threads), settings: { ...DEFAULT_SERVER_SETTINGS, sidebarAutoSettleOnMerge: true }, branchPullRequest: () => Effect.die("linked threads must not query the branch"), pullRequestSummary: () => Effect.die("linked threads must use their snapshots"), @@ -346,9 +369,66 @@ describe("ThreadSettlementReactor", () => { yield* Effect.gen(function* () { const reactor = yield* ThreadSettlementReactor.ThreadSettlementReactor; yield* startHarness(reactor, fixture.activation, fixture.snapshotReads); + assert.deepStrictEqual(yield* Ref.get(fixture.commands), []); + const eventBase = { + sequence: 2, + eventId: EventId.make("pull-request-synced"), + aggregateKind: "thread" as const, + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + }; + for (const thread of threads) { + const terminalLink = link(1, thread.id === "closed" ? "closed" : "merged"); + yield* Ref.update(fixture.snapshots, (snapshot) => ({ + ...snapshot, + threads: snapshot.threads.map((current) => + current.id === thread.id + ? { ...current, pullRequests: [terminalLink, ...current.pullRequests.slice(1)] } + : current, + ), + })); + yield* fixture.publishEvent({ + ...eventBase, + type: "thread.pull-request-synced", + aggregateId: thread.id, + payload: { + threadId: thread.id, + host: terminalLink.host, + repository: terminalLink.repository, + number: terminalLink.number, + snapshot: terminalLink.snapshot!, + stack: null, + updatedAt: NOW, + }, + }); + yield* Queue.take(fixture.snapshotReads); + yield* reactor.drain; + } + assert.deepStrictEqual( + (yield* Ref.get(fixture.commands)).map(({ threadId }) => threadId), + [ThreadId.make("merged"), ThreadId.make("closed")], + ); + const readySession = { ...runningSession, status: "ready" as const }; + yield* Ref.update(fixture.snapshots, (snapshot) => ({ + ...snapshot, + threads: snapshot.threads.map((thread) => + thread.id === readySession.threadId ? { ...thread, session: readySession } : thread, + ), + })); + yield* fixture.publishEvent({ + ...eventBase, + type: "thread.session-set", + aggregateId: readySession.threadId, + payload: { threadId: readySession.threadId, session: readySession }, + }); + yield* Queue.take(fixture.snapshotReads); + yield* reactor.drain; assert.deepStrictEqual( (yield* Ref.get(fixture.commands)).map(({ threadId }) => threadId), - [ThreadId.make("merged")], + [ThreadId.make("merged"), ThreadId.make("closed"), ThreadId.make("running")], ); assert.deepStrictEqual(yield* Ref.get(fixture.branchCalls), []); assert.deepStrictEqual(yield* Ref.get(fixture.summaryCalls), []); @@ -357,6 +437,109 @@ describe("ThreadSettlementReactor", () => { ), ); + it.effect("rechecks only the affected thread after link and unlink events", () => + Effect.scoped( + Effect.gen(function* () { + yield* TestClock.setTime(Date.parse(NOW)); + const link = (number: number, state: "open" | "merged") => ({ + host: "example.test", + repository: "owner/repository", + number, + url: `https://example.test/owner/repository/pull/${number}`, + source: "manual" as const, + linkedAt: NOW, + stack: null, + snapshot: { + state, + title: "Review", + headBranch: "feature", + baseBranch: "main", + isDraft: false, + updatedAt: NOW, + syncedAt: NOW, + mergedAt: state === "merged" ? NOW : null, + closedAt: null, + }, + }); + const merged = link(1, "merged"); + const open = link(2, "open"); + const fixture = yield* makeHarness({ + snapshot: makeSnapshot([ + makeThread("new-link"), + makeThread("removed-link", { pullRequests: [open, merged] }), + ]), + settings: { + ...DEFAULT_SERVER_SETTINGS, + sidebarAutoSettleOnMerge: true, + sidebarAutoSettleAfterDays: null, + }, + branchPullRequest: () => Effect.die("linked event must use the projected snapshot"), + pullRequestSummary: () => Effect.die("linked event must not query a provider"), + }); + yield* Effect.gen(function* () { + const reactor = yield* ThreadSettlementReactor.ThreadSettlementReactor; + yield* startHarness(reactor, fixture.activation, fixture.snapshotReads); + assert.deepStrictEqual(yield* Ref.get(fixture.commands), []); + const eventBase = { + sequence: 2, + eventId: EventId.make("link-change"), + aggregateKind: "thread" as const, + occurredAt: NOW, + commandId: null, + causationEventId: null, + correlationId: null, + metadata: {}, + }; + + yield* Ref.update(fixture.snapshots, (snapshot) => ({ + ...snapshot, + threads: snapshot.threads.map((thread) => + thread.id === ThreadId.make("new-link") + ? { ...thread, pullRequests: [merged] } + : thread, + ), + })); + yield* fixture.publishEvent({ + ...eventBase, + type: "thread.pull-request-linked", + aggregateId: ThreadId.make("new-link"), + payload: { threadId: ThreadId.make("new-link"), link: merged, updatedAt: NOW }, + }); + yield* Queue.take(fixture.snapshotReads); + yield* reactor.drain; + + yield* Ref.update(fixture.snapshots, (snapshot) => ({ + ...snapshot, + threads: snapshot.threads.map((thread) => + thread.id === ThreadId.make("removed-link") + ? { ...thread, pullRequests: [merged] } + : thread, + ), + })); + yield* fixture.publishEvent({ + ...eventBase, + type: "thread.pull-request-unlinked", + aggregateId: ThreadId.make("removed-link"), + payload: { + threadId: ThreadId.make("removed-link"), + host: open.host, + repository: open.repository, + number: open.number, + updatedAt: NOW, + }, + }); + yield* Queue.take(fixture.snapshotReads); + yield* reactor.drain; + assert.deepStrictEqual( + (yield* Ref.get(fixture.commands)).map(({ threadId }) => threadId), + [ThreadId.make("new-link"), ThreadId.make("removed-link")], + ); + assert.deepStrictEqual(yield* Ref.get(fixture.summaryCalls), []); + }).pipe(Effect.provide(fixture.layer)); + }), + ), + ); + it("distinguishes a project that inherits the threshold from one that disables it", () => { const inherits = ThreadSettlementReactor.autoSettlementSettingsKey({ ...DEFAULT_SERVER_SETTINGS, diff --git a/apps/server/src/orchestration/ThreadSettlementReactor.ts b/apps/server/src/orchestration/ThreadSettlementReactor.ts index b9041d2976..0dec28a7df 100644 --- a/apps/server/src/orchestration/ThreadSettlementReactor.ts +++ b/apps/server/src/orchestration/ThreadSettlementReactor.ts @@ -1,4 +1,9 @@ -import { CommandId, type ServerSettings as ServerSettingsValue } from "@t3tools/contracts"; +import { + CommandId, + type OrchestrationEvent, + type ServerSettings as ServerSettingsValue, + type ThreadId, +} from "@t3tools/contracts"; import { resolveProjectSettings } from "@t3tools/shared/projectSettings"; import { makeDrainableWorker } from "@t3tools/shared/DrainableWorker"; import * as Cause from "effect/Cause"; @@ -84,6 +89,7 @@ export const make = Effect.gen(function* () { const sweep = Effect.fn("ThreadSettlementReactor.sweep")(function* ( mergedPullRequest: PullRequestService.PullRequestMergeEvent | null, + threadId?: ThreadId, ) { const settings = yield* settingsService.getSettings; if (!autoSettlementConfigured(settings)) { @@ -94,7 +100,11 @@ export const make = Effect.gen(function* () { const projects = new Map(snapshot.projects.map((project) => [project.id, project])); // A merge rechecks all candidates, including branches that discovery has // not linked yet. Those lookups can still have cached the PR as open. - const candidates = snapshot.threads.filter((thread) => isAutoSettlementCandidate(thread, now)); + const candidates = snapshot.threads.filter( + (thread) => + (threadId === undefined || thread.id === threadId) && + isAutoSettlementCandidate(thread, now), + ); // Return the thread when it still needs a pull request decision. A rejected // dispatch skips it for this snapshot instead of retrying through a lookup. @@ -279,8 +289,11 @@ export const make = Effect.gen(function* () { ); }); - const runSweep = (mergedPullRequest: PullRequestService.PullRequestMergeEvent | null) => - sweep(mergedPullRequest).pipe( + const runSweep = ( + mergedPullRequest: PullRequestService.PullRequestMergeEvent | null, + threadId?: ThreadId, + ) => + sweep(mergedPullRequest, threadId).pipe( Effect.catchCause((cause) => Cause.hasInterruptsOnly(cause) ? Effect.failCause(cause) @@ -289,13 +302,36 @@ export const make = Effect.gen(function* () { }), ), ); - const worker = yield* makeDrainableWorker(() => runSweep(null)); + const worker = yield* makeDrainableWorker((threadId: ThreadId | undefined) => + runSweep(null, threadId), + ); + + const processEvent = (event: OrchestrationEvent) => { + switch (event.type) { + case "thread.pull-request-linked": + case "thread.pull-request-synced": + case "thread.pull-request-unlinked": + // Merge notifications can arrive before the linked snapshot is projected. + // Recheck the persisted state so terminal links settle without the timer. + return worker.enqueue(event.payload.threadId); + case "thread.session-set": + if ( + event.payload.session.status !== "running" && + event.payload.session.status !== "starting" + ) { + return worker.enqueue(event.payload.threadId); + } + break; + } + return Effect.void; + }; const start: ThreadSettlementReactor["Service"]["start"] = Effect.fn( "ThreadSettlementReactor.start", )(function* () { const settingsChanges = yield* settingsService.subscribeChanges; const mergedPullRequests = yield* pullRequests.subscribeMerges; + const events = yield* engine.subscribeDomainEvents; const initialSettings = yield* settingsService.getSettings.pipe(Effect.orDie); let lastSettlementSettings = autoSettlementSettingsKey(initialSettings); yield* forkParked( @@ -314,7 +350,8 @@ export const make = Effect.gen(function* () { return worker.enqueue(undefined); }), ); - yield* forkParked(Stream.runForEach(mergedPullRequests, runSweep)); + yield* forkParked(Stream.runForEach(mergedPullRequests, (event) => runSweep(event))); + yield* forkParked(Stream.runForEach(events, processEvent)); }); return { start, drain: worker.drain } satisfies ThreadSettlementReactor["Service"]; diff --git a/apps/server/src/pullRequest/PullRequestService.test.ts b/apps/server/src/pullRequest/PullRequestService.test.ts index a4f866b59a..e84703f241 100644 --- a/apps/server/src/pullRequest/PullRequestService.test.ts +++ b/apps/server/src/pullRequest/PullRequestService.test.ts @@ -1081,8 +1081,13 @@ it.effect("publishes a merge for immediate settlement only after host confirmati // Queueing succeeds while the host still reports an open PR. yield* service.runAction({ ...reference, action: "merge" }); + const queuedRefresh = Option.getOrThrow(yield* Stream.runHead(service.subscribeRefreshes)); confirmationFails = true; yield* service.runAction({ ...reference, action: "merge" }); + assert.isAbove( + Option.getOrThrow(yield* Stream.runHead(service.subscribeRefreshes)), + queuedRefresh, + ); confirmationFails = false; state = "merged"; yield* TestClock.setTime(Date.parse(mergedAt)); @@ -1101,6 +1106,53 @@ it.effect("publishes a merge for immediate settlement only after host confirmati ), ); +it.effect("refreshes every reader before a queued merge confirmation finishes", () => + Effect.scoped( + Effect.gen(function* () { + const confirmationStarted = yield* Deferred.make(); + const confirm = yield* Deferred.make(); + const reference = { projectId: "p1" as ProjectId, repository: "acme/web", number: 1 }; + const service = yield* makeService({ + projects: [ + project({ id: "p1", title: "web", workspaceRoot: "/a", repository: "acme/web" }), + ], + providers: [ + fakeProvider("github", { + getChangeRequestSummary: () => + Effect.gen(function* () { + yield* Deferred.succeed(confirmationStarted, undefined); + yield* Deferred.await(confirm); + return changeRequest(1, "2026-09-16T00:00:00.000Z"); + }), + }), + ], + }); + const merges = yield* service.subscribeMerges; + const observedMerge = yield* Stream.runHead(merges).pipe( + Effect.forkChild({ startImmediately: true }), + ); + const readers = yield* Effect.forEach([0, 1], () => + Stream.runHead(service.subscribeRefreshes).pipe( + Effect.forkChild({ startImmediately: true }), + ), + ); + const action = yield* service + .runAction({ ...reference, action: "merge" }) + .pipe(Effect.forkChild({ startImmediately: true })); + yield* Deferred.await(confirmationStarted); + const revisions = yield* Effect.forEach(readers, (reader) => + Fiber.join(reader).pipe(Effect.map(Option.getOrThrow)), + ); + assert.isAbove(revisions[0]!, 0); + assert.strictEqual(revisions[0], revisions[1]); + assert.isUndefined(action.pollUnsafe()); + yield* Deferred.succeed(confirm, undefined); + yield* Fiber.join(action); + assert.isUndefined(observedMerge.pollUnsafe()); + }), + ), +); + it.effect("refuses an action this viewer may not take, and says what access it takes", () => Effect.gen(function* () { let ran: string | null = null; @@ -2254,10 +2306,16 @@ it.effect("invalidates the cached activity after reacting, like the other mutati ], }); + yield* service.refreshAfterTurn(reference.projectId); + const previousRefresh = Option.getOrThrow(yield* Stream.runHead(service.subscribeRefreshes)); yield* service.activity(reference); assert.strictEqual(activityCalls, 1); yield* service.setReaction({ ...reference, content: "heart", reacted: true }); + assert.isAbove( + Option.getOrThrow(yield* Stream.runHead(service.subscribeRefreshes)), + previousRefresh, + ); yield* service.activity(reference); assert.strictEqual(activityCalls, 2); @@ -3192,31 +3250,106 @@ it.effect("explicit and turn invalidations make the next listing ask the host ag }), ); -it.effect("a mutation makes the next listing ask the host again, with no client asking", () => - Effect.gen(function* () { - let hostCalls = 0; - const service = yield* makeService({ - projects: [project({ id: "p1", title: "web", workspaceRoot: "/a", repository: "acme/web" })], - providers: [ - fakeProvider("github", { - listChangeRequests: () => { - hostCalls += 1; - return Effect.succeed({ items: [], truncated: false, continues: false }); - }, - }), - ], - }); +it.effect("close and reopen notify subscribed readers after invalidating their cached state", () => + Effect.scoped( + Effect.gen(function* () { + let hostCalls = 0; + let state: "open" | "closed" = "open"; + const reference = { projectId: "p1" as ProjectId, repository: "acme/web", number: 1 }; + const service = yield* makeService({ + projects: [ + project({ id: "p1", title: "web", workspaceRoot: "/a", repository: "acme/web" }), + ], + providers: [ + fakeProvider("github", { + getChangeRequestSummary: () => + Effect.succeed({ ...changeRequest(1, "2026-09-16T00:00:00.000Z"), state }), + runAction: (input) => + Effect.sync(() => { + state = input.action === "close" ? "closed" : "open"; + }), + listChangeRequests: () => { + hostCalls += 1; + return Effect.succeed({ items: [], truncated: false, continues: false }); + }, + }), + ], + }); - yield* service.list({ state: "open" }); - yield* service.runAction({ - projectId: "p1" as ProjectId, - repository: "acme/web", - number: 1, - action: "close", - }); - yield* service.list({ state: "open" }); - assert.strictEqual(hostCalls, 2); - }), + yield* service.refreshAfterTurn(reference.projectId); + yield* service.list({ state: "open" }); + assert.strictEqual((yield* service.summary(reference)).state, "open"); + for (const action of ["close", "reopen"] as const) { + const refreshed = yield* service.subscribeRefreshes.pipe( + Stream.drop(1), + Stream.take(1), + Stream.mapEffect(() => + Effect.gen(function* () { + yield* service.list({ state: "open" }); + return yield* service.summary(reference); + }), + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + yield* service.runAction({ ...reference, action }); + assert.strictEqual( + Option.getOrThrow(yield* Fiber.join(refreshed)).state, + action === "close" ? "closed" : "open", + ); + } + assert.strictEqual(hostCalls, 3); + }), + ), +); + +it.effect("explicit invalidation refreshes origin readers after a routed host mutation", () => + Effect.scoped( + Effect.gen(function* () { + let state: "open" | "closed" = "open"; + const reference = { projectId: "p1" as ProjectId, repository: "acme/web", number: 1 }; + const service = yield* makeService({ + projects: [ + project({ id: "p1", title: "web", workspaceRoot: "/a", repository: "acme/web" }), + ], + providers: [ + fakeProvider("github", { + getChangeRequestSummary: () => + Effect.succeed({ ...changeRequest(1, "2026-09-16T00:00:00.000Z"), state }), + }), + ], + }); + yield* service.refreshAfterTurn(reference.projectId); + let revision = Option.getOrThrow(yield* Stream.runHead(service.subscribeRefreshes)); + // The sync reactor invalidates before reading; it must not notify itself again. + yield* service.invalidate({ reference }); + assert.strictEqual( + Option.getOrThrow(yield* Stream.runHead(service.subscribeRefreshes)), + revision, + ); + assert.strictEqual((yield* service.summary(reference)).state, "open"); + + for (const nextState of ["closed", "open"] as const) { + const refreshed = yield* service.subscribeRefreshes.pipe( + Stream.drop(1), + Stream.take(1), + Stream.mapEffect((nextRevision) => + service + .summary(reference) + .pipe(Effect.map((summary) => ({ revision: nextRevision, state: summary.state }))), + ), + Stream.runHead, + Effect.forkChild({ startImmediately: true }), + ); + state = nextState; + yield* service.invalidate({ reference }, { notifyReaders: true }); + const result = Option.getOrThrow(yield* Fiber.join(refreshed)); + assert.strictEqual(result.state, nextState); + assert.isAbove(result.revision, revision); + revision = result.revision; + } + }), + ), ); it.effect("does not cache a failed listing", () => diff --git a/apps/server/src/pullRequest/PullRequestService.ts b/apps/server/src/pullRequest/PullRequestService.ts index c61ce781a8..ff53821727 100644 --- a/apps/server/src/pullRequest/PullRequestService.ts +++ b/apps/server/src/pullRequest/PullRequestService.ts @@ -214,7 +214,10 @@ export class PullRequestService extends Context.Service< readonly setLabels: ( input: PullRequestLabelChangeInput, ) => Effect.Effect; - readonly invalidate: (input: PullRequestInvalidateInput) => Effect.Effect; + readonly invalidate: ( + input: PullRequestInvalidateInput, + options?: { readonly notifyReaders?: boolean }, + ) => Effect.Effect; } >()("t3/pullRequest/PullRequestService") {} @@ -2732,10 +2735,12 @@ export const make = Effect.gen(function* () { return { stats: [...held, ...result.stats] }; }); - const invalidate: PullRequestService["Service"]["invalidate"] = (input) => { + const invalidate: PullRequestService["Service"]["invalidate"] = Effect.fn( + "PullRequestService.invalidate", + )(function* (input, options) { const reference = input.reference; if (reference !== undefined) { - return canonicalRef(reference).pipe( + yield* canonicalRef(reference).pipe( Effect.flatMap((ref) => readCache .invalidate(refScope(ref)) @@ -2743,12 +2748,15 @@ export const make = Effect.gen(function* () { ), Effect.ignore, ); - } - return Effect.sync(() => { + } else { listingsEpoch = ++epochCounter; viewersByHost.clear(); - }).pipe(Effect.andThen(Cache.invalidateAll(viewerFlights))); - }; + yield* Cache.invalidateAll(viewerFlights); + } + if (options?.notifyReaders) { + yield* SubscriptionRef.set(pullRequestRefreshes, ++epochCounter); + } + }); const refreshAfterTurn: PullRequestService["Service"]["refreshAfterTurn"] = (projectId) => Effect.suspend(() => { @@ -2767,9 +2775,7 @@ export const make = Effect.gen(function* () { .pipe(Effect.andThen(SubscriptionRef.set(pullRequestRefreshes, listingsEpoch))); }); - // A mutation's own client re-reads right after it, and every other client's next read must - // see the action too — so a write forgets the change request it touched and the listings its - // state change reorders, for everyone, without any client asking. + // Invalidate before notifying every client so mounted readers immediately fetch the edit. const invalidatedByMutation = ( method: (input: I) => Effect.Effect, @@ -2787,6 +2793,7 @@ export const make = Effect.gen(function* () { }), ), ); + yield* SubscriptionRef.set(pullRequestRefreshes, listingsEpoch); }); const runActionAndInvalidate: PullRequestService["Service"]["runAction"] = Effect.fn( "PullRequestService.runActionAndInvalidate", @@ -2798,6 +2805,7 @@ export const make = Effect.gen(function* () { ); bumpRefEpoch({ ...ref, repository }); listingsEpoch = ++epochCounter; + yield* SubscriptionRef.set(pullRequestRefreshes, listingsEpoch); if (input.action === "merge") { // A successful merge action can merely enqueue the PR or enable auto-merge. const confirmed = yield* summaryUncached({ ...input, repository }).pipe( diff --git a/apps/server/src/ws.ts b/apps/server/src/ws.ts index 9d4328a8cc..6864bd1c21 100644 --- a/apps/server/src/ws.ts +++ b/apps/server/src/ws.ts @@ -3097,7 +3097,7 @@ const makeWsRpcLayer = ( [WS_METHODS.pullRequestsInvalidate]: (input) => observeRpcEffect( WS_METHODS.pullRequestsInvalidate, - pullRequests.invalidate(input).pipe( + pullRequests.invalidate(input, { notifyReaders: true }).pipe( // A reader asking for fresh host state also wants the thread badges it feeds to // catch up, including a merged link the sweep would otherwise never revisit. Effect.andThen( diff --git a/packages/client-runtime/src/state/pullRequests.test.ts b/packages/client-runtime/src/state/pullRequests.test.ts index 6187cf2e7b..fe6ea90a59 100644 --- a/packages/client-runtime/src/state/pullRequests.test.ts +++ b/packages/client-runtime/src/state/pullRequests.test.ts @@ -21,6 +21,7 @@ import * as EnvironmentSupervisor from "../connection/supervisor.ts"; import type { WsRpcProtocolClient } from "../rpc/protocol.ts"; import type { RpcSession } from "../rpc/session.ts"; import { + createLinkedPullRequestSummaryAtomFamily, createPullRequestEnvironmentAtoms, createPullRequestStackAtomFamily, } from "./pullRequests.ts"; @@ -134,6 +135,108 @@ it.effect("keeps concurrent diff file reads on different hosts separate", () => ), ); +it.effect("shares close, reopen, and merge with an untouched client's mounted PR readers", () => + Effect.scoped( + Effect.gen(function* () { + const revision = yield* SubscriptionRef.make(0); + let state: "open" | "closed" | "merged" = "open"; + const client = { + [WS_METHODS.pullRequestsSubscribeRefreshes]: () => SubscriptionRef.changes(revision), + [WS_METHODS.pullRequestsSummary]: () => Effect.sync(() => ({ state })), + [WS_METHODS.pullRequestsDetail]: () => Effect.sync(() => ({ state })), + [WS_METHODS.pullRequestsList]: () => + Effect.sync(() => ({ entries: [{ number: 1, state }] })), + [WS_METHODS.pullRequestsRunAction]: (input: { + readonly action: "close" | "reopen" | "merge"; + }) => + Effect.gen(function* () { + state = + input.action === "close" ? "closed" : input.action === "merge" ? "merged" : "open"; + yield* SubscriptionRef.update(revision, (value) => value + 1); + }), + } as unknown as WsRpcProtocolClient; + const writer = yield* makeTestRuntime(client); + const reader = yield* makeTestRuntime(client); + const target = { + environmentId: TARGET.environmentId, + input: { + projectId: ProjectId.make("project-1"), + host: "github.example.com", + repository: "acme/web", + number: 1, + }, + }; + const detail = reader.atoms.detail(target); + const summary = createLinkedPullRequestSummaryAtomFamily( + reader.runtime, + reader.atoms.refreshes, + )(target); + const list = reader.atoms.list({ + environmentId: TARGET.environmentId, + input: { state: "all" }, + }); + const unmountDetail = reader.registry.mount(detail); + const unmountSummary = reader.registry.mount(summary); + const unmountList = reader.registry.mount(list); + yield* Effect.addFinalizer(() => + Effect.sync(() => { + unmountDetail(); + unmountSummary(); + unmountList(); + }), + ); + expect((yield* AtomRegistry.getResult(reader.registry, detail)).state).toBe("open"); + expect((yield* AtomRegistry.getResult(reader.registry, summary)).state).toBe("open"); + expect((yield* AtomRegistry.getResult(reader.registry, list)).entries[0]?.state).toBe("open"); + + for (const [action, expected] of [ + ["close", "closed"], + ["reopen", "open"], + ["merge", "merged"], + ] as const) { + const detailChanged = Latch.makeUnsafe(); + const summaryChanged = Latch.makeUnsafe(); + const listChanged = Latch.makeUnsafe(); + const stops = [ + reader.registry.subscribe(detail, (result) => { + if (AsyncResult.isSuccess(result) && result.value.state === expected) { + detailChanged.openUnsafe(); + } + }), + reader.registry.subscribe(summary, (result) => { + if (AsyncResult.isSuccess(result) && result.value.state === expected) { + summaryChanged.openUnsafe(); + } + }), + reader.registry.subscribe(list, (result) => { + if (AsyncResult.isSuccess(result) && result.value.entries[0]?.state === expected) { + listChanged.openUnsafe(); + } + }), + ]; + yield* Effect.addFinalizer(() => Effect.sync(() => stops.forEach((stop) => stop()))); + const result = yield* Effect.promise(() => + writer.atoms.runAction.run(writer.registry, { + ...target, + input: { ...target.input, action }, + }), + ); + expect(AsyncResult.isSuccess(result)).toBe(true); + // The second client receives only the server push: no local refresh or timer tick. + yield* detailChanged.await; + yield* summaryChanged.await; + yield* listChanged.await; + expect((yield* AtomRegistry.getResult(reader.registry, detail)).state).toBe(expected); + expect((yield* AtomRegistry.getResult(reader.registry, summary)).state).toBe(expected); + expect((yield* AtomRegistry.getResult(reader.registry, list)).entries[0]?.state).toBe( + expected, + ); + stops.forEach((stop) => stop()); + } + }), + ), +); + it.effect("refreshes pull request activity after a comment is updated", () => Effect.scoped( Effect.gen(function* () {