Skip to content
Open
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
Original file line numberDiff line numberDiff line change
Expand Up@@ -20,6 +20,7 @@ import {
ProviderService,
type ProviderServiceShape,
} from "../../provider/Services/ProviderService.ts";
import { ProviderSessionDirectory } from "../../provider/Services/ProviderSessionDirectory.ts";
import * as TerminalManager from "../../terminal/Manager.ts";
import {
OrchestrationEngineService,
Expand DownExpand Up@@ -108,8 +109,12 @@ describe("ThreadDeletionReactor drain", () => {
const terminalManager = {
close: () => Effect.void,
} as unknown as TerminalManager.TerminalManager["Service"];
const directory = {
remove: () => Effect.void,
} as unknown as ProviderSessionDirectory["Service"];
const layer = ThreadDeletionReactorLive.pipe(
Layer.provide(Layer.succeed(ProviderService, providerService)),
Layer.provide(Layer.succeed(ProviderSessionDirectory, directory)),
Layer.provide(Layer.succeed(TerminalManager.TerminalManager, terminalManager)),
Layer.provide(Layer.succeed(OrchestrationEngineService, engine)),
);
Expand All@@ -136,4 +141,42 @@ describe("ThreadDeletionReactor drain", () => {
).pipe(Effect.provide(layer));
}),
);

effectIt.effect("forgets the provider binding for a deleted thread", () =>
Effect.gen(function* () {
const removed: Array<ThreadId> = [];
const engine = {
latestSequence: Effect.succeed(0),
streamDomainEvents: Stream.make(deletedEvent(1)),
} as unknown as OrchestrationEngineShape;
const providerService = {
stopSession: () => Effect.void,
} as unknown as ProviderServiceShape;
const terminalManager = {
close: () => Effect.void,
} as unknown as TerminalManager.TerminalManager["Service"];
const directory = {
remove: (removedThreadId: ThreadId) =>
Effect.sync(() => {
removed.push(removedThreadId);
}),
} as unknown as ProviderSessionDirectory["Service"];
const layer = ThreadDeletionReactorLive.pipe(
Layer.provide(Layer.succeed(ProviderService, providerService)),
Layer.provide(Layer.succeed(ProviderSessionDirectory, directory)),
Layer.provide(Layer.succeed(TerminalManager.TerminalManager, terminalManager)),
Layer.provide(Layer.succeed(OrchestrationEngineService, engine)),
);

yield* Effect.scoped(
Effect.gen(function* () {
const reactor = yield* ThreadDeletionReactor;
yield* reactor.start();
yield* reactor.drainThrough(1);

expect(removed).toEqual([threadId]);
}),
).pipe(Effect.provide(layer));
}),
);
});
12 changes: 12 additions & 0 deletions apps/server/src/orchestration/Layers/ThreadDeletionReactor.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -7,6 +7,7 @@ import * as Stream from "effect/Stream";
import * as SubscriptionRef from "effect/SubscriptionRef";

import { ProviderService } from "../../provider/Services/ProviderService.ts";
import { ProviderSessionDirectory } from "../../provider/Services/ProviderSessionDirectory.ts";
import * as TerminalManager from "../../terminal/Manager.ts";
import { OrchestrationEngineService } from "../Services/OrchestrationEngine.ts";
import {
Expand DownExpand Up@@ -41,6 +42,7 @@ export const logCleanupCauseUnlessInterrupted = <R, E>({
const make = Effect.gen(function* () {
const orchestrationEngine = yield* OrchestrationEngineService;
const providerService = yield* ProviderService;
const providerSessionDirectory = yield* ProviderSessionDirectory;
const terminalManager = yield* TerminalManager.TerminalManager;

const stopProviderSession = (threadId: ThreadDeletedEvent["payload"]["threadId"]) =>
Expand All@@ -50,6 +52,15 @@ const make = Effect.gen(function* () {
threadId,
});

// Nothing else prunes the runtime table, so a deleted thread would keep its
// binding for the life of the install.
const forgetProviderBinding = (threadId: ThreadDeletedEvent["payload"]["threadId"]) =>
logCleanupCauseUnlessInterrupted({
effect: providerSessionDirectory.remove(threadId),
message: "thread deletion cleanup skipped provider binding removal",
threadId,
});

const closeThreadTerminals = (threadId: ThreadDeletedEvent["payload"]["threadId"]) =>
logCleanupCauseUnlessInterrupted({
effect: terminalManager.close({ threadId, deleteHistory: true }),
Expand All@@ -62,6 +73,7 @@ const make = Effect.gen(function* () {
) {
const { threadId } = event.payload;
yield* stopProviderSession(threadId);
yield* forgetProviderBinding(threadId);

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Retry can lose new provider binding

Medium Severity

Draft retries reuse a soft-deleted thread id. forgetProviderBinding always calls remove on the async deletion worker and does not check for a later thread.created. If the retry upserts a binding before that worker finishes, the new provider_session_runtime row is deleted and later routing cannot find a session.

Fix in CursorFix in Web

Reviewed by Cursor Bugbot for commit 4fe4855. Configure here.

yield* closeThreadTerminals(threadId);
});

Expand Down
1 change: 1 addition & 0 deletions apps/server/src/provider/Layers/CodexAdapter.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -225,6 +225,7 @@ const providerSessionDirectoryTestLayer = Layer.succeed(ProviderSessionDirectory
getProvider: () =>
Effect.die(new Error("ProviderSessionDirectory.getProvider is not used in test")),
getBinding: () => Effect.succeed(Option.none()),
remove: () => Effect.void,
listThreadIds: () => Effect.succeed([]),
listBindings: () => Effect.succeed([]),
});
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/provider/Layers/OpenCodeAdapter.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -424,6 +424,7 @@ const providerSessionDirectoryTestLayer = Layer.succeed(ProviderSessionDirectory
getProvider: () =>
Effect.die(new Error("ProviderSessionDirectory.getProvider is not used in test")),
getBinding: () => Effect.succeed(Option.none()),
remove: () => Effect.void,
listThreadIds: () => Effect.succeed([]),
listBindings: () => Effect.succeed([]),
});
Expand Down
1 change: 1 addition & 0 deletions apps/server/src/provider/Layers/ProviderService.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -2467,6 +2467,7 @@ const boundedListing = makeProviderServiceLayer({
upsert: () => Effect.void,
getProvider: () => Effect.die("ProviderService.listSessions does not use getProvider"),
getBinding,
remove: () => Effect.void,
listThreadIds,
listBindings: () => Effect.die("ProviderService.listSessions does not use listBindings"),
},
Expand Down
8 changes: 8 additions & 0 deletions apps/server/src/provider/Layers/ProviderSessionDirectory.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -164,6 +164,13 @@ const makeProviderSessionDirectory = Effect.gen(function* () {
),
);

const remove: ProviderSessionDirectoryShape["remove"] = (threadId) =>
repository
.deleteByThreadId({ threadId })
.pipe(
Effect.mapError(toPersistenceError("ProviderSessionDirectory.remove:deleteByThreadId")),
);

const listThreadIds: ProviderSessionDirectoryShape["listThreadIds"] = () =>
repository.list().pipe(
Effect.mapError(toPersistenceError("ProviderSessionDirectory.listThreadIds:list")),
Expand All@@ -186,6 +193,7 @@ const makeProviderSessionDirectory = Effect.gen(function* () {
upsert,
getProvider,
getBinding,
remove,
listThreadIds,
listBindings,
} satisfies ProviderSessionDirectoryShape;
Expand Down
Original file line numberDiff line numberDiff line change
Expand Up@@ -53,6 +53,11 @@ export interface ProviderSessionDirectoryShape {
threadId: ThreadId,
) => Effect.Effect<Option.Option<ProviderRuntimeBinding>, ProviderSessionDirectoryReadError>;

/** Forgets a thread's runtime binding, for example when the thread is deleted. */
readonly remove: (
threadId: ThreadId,
) => Effect.Effect<void, ProviderSessionDirectoryPersistenceError>;

readonly listThreadIds: () => Effect.Effect<
ReadonlyArray<ThreadId>,
ProviderSessionDirectoryPersistenceError
Expand Down
8 changes: 8 additions & 0 deletions apps/server/src/serverRuntimeStartup.reconcile.test.ts
Original file line numberDiff line numberDiff line change
Expand Up@@ -132,6 +132,7 @@ it.effect("marks active running sessions that have persisted resume state", () =
),
upsert: (binding) => Effect.sync(() => upserts.push(binding)),
getProvider: () => Effect.die("unused"),
remove: () => Effect.void,
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
}),
Expand DownExpand Up@@ -232,6 +233,7 @@ it.effect("continues marked sessions after activation with provider-specific inp
),
),
getProvider: () => Effect.die("unused"),
remove: () => Effect.void,
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
},
Expand DownExpand Up@@ -355,6 +357,7 @@ it.effect("does not continue archived or deleted marked sessions", () => {
},
upsert: () => Effect.void,
getProvider: () => Effect.die("unused"),
remove: () => Effect.void,
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
},
Expand DownExpand Up@@ -410,6 +413,7 @@ it.effect("retries continuation preparation before settling a persistent failure
),
upsert: () => Effect.void,
getProvider: () => Effect.die("unused"),
remove: () => Effect.void,
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
},
Expand DownExpand Up@@ -482,6 +486,7 @@ it.effect("reconciles multiple active and archived orphans but skips live sessio
upsert: (binding) => Effect.sync(() => upserts.push(binding)),
getProvider: () => Effect.die("unused"),
listThreadIds: () => Effect.die("unused"),
remove: () => Effect.void,
listBindings: () => Effect.die("unused"),
},
dispatch: (command) =>
Expand DownExpand Up@@ -560,6 +565,7 @@ it.effect(
upsert: () => Effect.fail(writeFailure),
getProvider: () => Effect.die("unused"),
listThreadIds: () => Effect.die("unused"),
remove: () => Effect.void,
listBindings: () => Effect.die("unused"),
},
dispatch: (command) =>
Expand DownExpand Up@@ -597,6 +603,7 @@ it.effect("retries failed projections and continues after a persistent failure",
upsert: () => Effect.void,
getProvider: () => Effect.die("unused"),
listThreadIds: () => Effect.die("unused"),
remove: () => Effect.void,
listBindings: () => Effect.die("unused"),
},
dispatch: (command) => {
Expand DownExpand Up@@ -645,6 +652,7 @@ it.effect("does not fail startup when the live provider session inventory cannot
upsert: () => Effect.die("unused"),
getProvider: () => Effect.die("unused"),
listThreadIds: () => Effect.die("unused"),
remove: () => Effect.void,
listBindings: () => Effect.die("unused"),
}),
Effect.provideService(OrchestrationEngine.OrchestrationEngineService, {
Expand Down
Loading