Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
2 changes: 0 additions & 2 deletions apps/mobile/src/features/threads/ThreadFeed.tsx
Original file line number Diff line number Diff line change
Expand Up @@ -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}
Expand Down
37 changes: 0 additions & 37 deletions apps/server/src/orchestration/Layers/ProjectionSnapshotQuery.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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,
Expand Down
58 changes: 56 additions & 2 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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-"));
Expand All @@ -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,
Expand Down Expand Up @@ -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;
Expand Down
38 changes: 37 additions & 1 deletion apps/server/src/provider/Layers/ProviderService.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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";
Expand Down Expand Up @@ -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<Name extends keyof ProviderService.ProviderService["Service"]> =
Expand Down Expand Up @@ -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(() => []));
Expand Down
8 changes: 4 additions & 4 deletions apps/server/src/vcs/VcsStatusBroadcaster.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -729,10 +729,10 @@ describe("VcsStatusBroadcaster", () => {
const nextSnapshot = yield* Deferred.make<VcsStatusStreamEvent>();
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
Expand Down
Loading