From b8262f9ad15eb2b5f0204d97562cd851ca5fbe38 Mon Sep 17 00:00:00 2001 From: Patrick Roza Date: Sat, 25 Jul 2026 11:11:40 +0200 Subject: [PATCH 1/2] fix(server): bound provider shutdown drains --- .../Layers/ProjectionSnapshotQuery.ts | 37 ------------ .../provider/Layers/ProviderService.test.ts | 58 ++++++++++++++++++- .../src/provider/Layers/ProviderService.ts | 38 +++++++++++- 3 files changed, 93 insertions(+), 40 deletions(-) diff --git a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts index bed3b1cb3ec..090b97abc92 100644 --- a/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts +++ b/apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts @@ -2331,43 +2331,6 @@ const makeProjectionSnapshotQuery = Effect.gen(function* () { ), ), ); - const getThreadActivitiesPage: ProjectionSnapshotQueryShape["getThreadActivitiesPage"] = ( - input, - ) => - Effect.gen(function* () { - const limit = Math.min( - Math.max(1, input.limit ?? THREAD_DETAIL_ACTIVITY_WINDOW), - THREAD_DETAIL_ACTIVITY_WINDOW, - ); - // Fetch one extra to detect whether older activities remain. - const rowsEffect = - "beforeSequence" in input - ? listThreadActivityRowsBeforeSequence({ - threadId: input.threadId, - beforeSequence: input.beforeSequence, - limit: limit + 1, - }) - : listUnsequencedThreadActivityRowsBeforeActivity({ - threadId: input.threadId, - beforeCreatedAt: input.beforeCreatedAt, - beforeActivityId: input.beforeActivityId, - limit: limit + 1, - }); - const rows = yield* rowsEffect.pipe( - Effect.mapError( - toPersistenceSqlOrDecodeError( - "ProjectionSnapshotQuery.getThreadActivitiesPage:query", - "ProjectionSnapshotQuery.getThreadActivitiesPage:decodeRows", - ), - ), - ); - const hasMore = rows.length > limit; - // Rows are newest-first; keep the page closest to the cursor, then reverse - // to ascending for display. - const page = (hasMore ? rows.slice(0, limit) : rows).map(mapThreadActivityRow).toReversed(); - return { activities: page, hasMore }; - }); - return { getCommandReadModel, getSnapshot, diff --git a/apps/server/src/provider/Layers/ProviderService.test.ts b/apps/server/src/provider/Layers/ProviderService.test.ts index 1fb1cd92c7a..7799b450ae9 100644 --- a/apps/server/src/provider/Layers/ProviderService.test.ts +++ b/apps/server/src/provider/Layers/ProviderService.test.ts @@ -367,6 +367,53 @@ it.effect("ProviderServiceLive catches stopAll failures during shutdown", () => }), ); +it.effect("ProviderServiceLive bounds a provider that wedges during shutdown", () => + Effect.gen(function* () { + const codex = makeFakeCodexAdapter(); + codex.stopAll.mockImplementation(() => Effect.never); + const registry = makeAdapterRegistryMock({ + [CODEX_DRIVER]: codex.adapter, + }); + const providerAdapterLayer = Layer.succeed( + ProviderAdapterRegistry.ProviderAdapterRegistry, + registry, + ); + const runtimeRepositoryLayer = ProviderSessionRuntime.layer.pipe( + Layer.provide(SqlitePersistenceMemory), + ); + const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); + const providerLayer = Layer.mergeAll( + makeProviderServiceLive({ shutdownGracePeriod: "50 millis" }).pipe( + Layer.provide(providerAdapterLayer), + Layer.provide(directoryLayer), + Layer.provide(defaultServerSettingsLayer), + Layer.provideMerge(AnalyticsService.layerTest), + Layer.provide( + Layer.succeed( + ProviderEventLoggers.ProviderEventLoggers, + ProviderEventLoggers.NoOpProviderEventLoggers, + ), + ), + ), + directoryLayer, + runtimeRepositoryLayer, + NodeServices.layer, + ); + const scope = yield* Scope.make(); + const runtimeServices = yield* Layer.build(providerLayer).pipe(Scope.provide(scope)); + + yield* ProviderService.ProviderService.pipe(Effect.provide(runtimeServices)); + const closeFiber = yield* Scope.close(scope, Exit.void).pipe( + Effect.forkChild({ startImmediately: true }), + ); + yield* advanceTestClock(50); + const closeExit = yield* Fiber.join(closeFiber).pipe(Effect.exit); + + assert.equal(Exit.isSuccess(closeExit), true); + assert.equal(codex.stopAll.mock.calls.length, 1); + }), +); + it.effect("graceful shutdown preserves recovery intent only for working sessions", () => Effect.gen(function* () { const tempDir = NodeFS.mkdtempSync(NodePath.join(NodeOS.tmpdir(), "t3-provider-recovery-")); @@ -377,7 +424,7 @@ it.effect("graceful shutdown preserves recovery intent only for working sessions ); const directoryLayer = ProviderSessionDirectoryLive.pipe(Layer.provide(runtimeRepositoryLayer)); const codex = makeFakeCodexAdapter(); - const providerLayer = makeProviderServiceLive().pipe( + const providerLayer = makeProviderServiceLive({ shutdownGracePeriod: "50 millis" }).pipe( Layer.provide( Layer.succeed( ProviderAdapterRegistry.ProviderAdapterRegistry, @@ -428,7 +475,14 @@ it.effect("graceful shutdown preserves recovery intent only for working sessions })); yield* provider.stopSession({ threadId: stoppedThreadId }); - yield* Scope.close(scope, Exit.void); + // Model a provider protocol drain that never completes. Recovery intent + // must already be durable when the global shutdown deadline interrupts it. + codex.stopAll.mockImplementation(() => Effect.never); + const closeFiber = yield* Scope.close(scope, Exit.void).pipe( + Effect.forkChild({ startImmediately: true }), + ); + yield* advanceTestClock(50); + yield* Fiber.join(closeFiber); const rows = yield* Effect.gen(function* () { const repository = yield* ProviderSessionRuntime.ProviderSessionRuntimeRepository; diff --git a/apps/server/src/provider/Layers/ProviderService.ts b/apps/server/src/provider/Layers/ProviderService.ts index b7870fd105e..678c9cb590f 100644 --- a/apps/server/src/provider/Layers/ProviderService.ts +++ b/apps/server/src/provider/Layers/ProviderService.ts @@ -26,6 +26,7 @@ import { } from "@t3tools/contracts"; import { causeErrorTag } from "@t3tools/shared/observability"; import * as DateTime from "effect/DateTime"; +import type * as Duration from "effect/Duration"; import * as Effect from "effect/Effect"; import * as Layer from "effect/Layer"; import * as Option from "effect/Option"; @@ -71,6 +72,12 @@ const isModelSelection = Schema.is(ModelSelection); */ export interface ProviderServiceLiveOptions { readonly canonicalEventLogger?: EventNdjsonLogger; + /** + * Maximum time the server gives all provider adapters, collectively, to + * stop during process shutdown. Recovery intent is persisted before this + * clock starts. + */ + readonly shutdownGracePeriod?: Duration.Input; } type ProviderServiceMethod = @@ -1182,7 +1189,36 @@ const makeProviderService = Effect.fn("makeProviderService")(function* ( }); }), ).pipe(Effect.asVoid); - yield* Effect.forEach(currentAdapters, ([, adapter]) => adapter.stopAll()).pipe(Effect.asVoid); + const adapterStops = yield* Effect.forEach( + currentAdapters, + ([instanceId, adapter]) => + adapter.stopAll().pipe( + Effect.exit, + Effect.map((exit) => ({ instanceId, exit })), + ), + { concurrency: "unbounded" }, + ).pipe( + // Scope finalizers are uninterruptible by default. Restore + // interruptibility here so the timeout can release a provider whose + // protocol drain never completes. + Effect.interruptible, + Effect.timeoutOption(options?.shutdownGracePeriod ?? "1 minute"), + ); + if (Option.isNone(adapterStops)) { + yield* Effect.logWarning("provider shutdown grace period elapsed", { + timeout: String(options?.shutdownGracePeriod ?? "1 minute"), + sessionCount: activeSessions.length, + }); + } else { + yield* Effect.forEach( + adapterStops.value, + ({ instanceId, exit }) => + exit._tag === "Failure" + ? Effect.logWarning("provider adapter failed during shutdown", { instanceId }) + : Effect.void, + { discard: true }, + ); + } yield* McpSessionRegistry.revokeAllActiveMcpCredentials(); McpProviderSession.clearAllMcpProviderSessions(); const bindings = yield* directory.listBindings().pipe(Effect.orElseSucceed(() => [])); From c8b59626d44cce749280a432e6515fa834b518e0 Mon Sep 17 00:00:00 2001 From: Patrick Roza Date: Sat, 25 Jul 2026 11:42:59 +0200 Subject: [PATCH 2/2] fix: repair memory candidate integration checks --- apps/mobile/src/features/threads/ThreadFeed.tsx | 2 -- apps/server/src/vcs/VcsStatusBroadcaster.test.ts | 8 ++++---- 2 files changed, 4 insertions(+), 6 deletions(-) diff --git a/apps/mobile/src/features/threads/ThreadFeed.tsx b/apps/mobile/src/features/threads/ThreadFeed.tsx index f9c763d55ff..7e991e0aaae 100644 --- a/apps/mobile/src/features/threads/ThreadFeed.tsx +++ b/apps/mobile/src/features/threads/ThreadFeed.tsx @@ -1836,8 +1836,6 @@ export const ThreadFeed = memo(function ThreadFeed(props: ThreadFeedProps) { // maintainVisibleContentPosition also keeps the viewport anchored // when older history prepends at the top. maintainVisibleContentPosition={maintainVisibleContentPosition} - onStartReached={onStartReachedOlderHistory} - onStartReachedThreshold={0.5} data={presentedFeed} extraData={listAppearanceData} renderItem={renderItem} diff --git a/apps/server/src/vcs/VcsStatusBroadcaster.test.ts b/apps/server/src/vcs/VcsStatusBroadcaster.test.ts index cc56d7c06c0..10b86894440 100644 --- a/apps/server/src/vcs/VcsStatusBroadcaster.test.ts +++ b/apps/server/src/vcs/VcsStatusBroadcaster.test.ts @@ -705,12 +705,12 @@ describe("VcsStatusBroadcaster", () => { event._tag === "localUpdated" ? Deferred.succeed(firstLocal, event).pipe(Effect.ignore) : Effect.void, - ).pipe(Effect.forkIn(firstScope)); + ).pipe(Effect.forkIn(firstScope, { startImmediately: true })); yield* Stream.runForEach(broadcaster.streamStatus({ cwd: "/repo" }), (event) => event._tag === "localUpdated" ? Deferred.succeed(secondLocal, event).pipe(Effect.ignore) : Effect.void, - ).pipe(Effect.forkIn(secondScope)); + ).pipe(Effect.forkIn(secondScope, { startImmediately: true })); yield* Deferred.await(firstLocal); yield* Deferred.await(secondLocal); @@ -729,10 +729,10 @@ describe("VcsStatusBroadcaster", () => { const nextSnapshot = yield* Deferred.make(); const nextScope = yield* Scope.make(); yield* Stream.runForEach(broadcaster.streamStatus({ cwd: "/repo" }), (event) => - event._tag === "snapshot" + event._tag === "localUpdated" ? Deferred.succeed(nextSnapshot, event).pipe(Effect.ignore) : Effect.void, - ).pipe(Effect.forkIn(nextScope)); + ).pipe(Effect.forkIn(nextScope, { startImmediately: true })); yield* Deferred.await(nextSnapshot); // Releasing the final poller also evicts its cwd cache entry, so a later