Skip to content
Closed
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
35 changes: 29 additions & 6 deletions apps/server/src/persistence/ProviderSessionRuntime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ export type RecordImportedTranscriptInput = typeof RecordImportedTranscriptInput

export interface ProviderSessionRuntimeUpsertOptions {
readonly onConflict?: "update" | "ignore";
readonly unlessNativeSessionId?: string;
}

/**
Expand All @@ -84,7 +85,7 @@ export class ProviderSessionRuntimeRepository extends Context.Service<
readonly upsert: (
runtime: ProviderSessionRuntime,
options?: ProviderSessionRuntimeUpsertOptions,
) => Effect.Effect<void, ProviderSessionRuntimeRepositoryError>;
) => Effect.Effect<boolean, ProviderSessionRuntimeRepositoryError>;

/** Record one source file without replacing the current session state. */
readonly recordImportedTranscript: (
Expand Down Expand Up @@ -237,8 +238,12 @@ export const make = Effect.gen(function* () {
`,
});

const insertRuntimeRow = SqlSchema.void({
Request: ProviderSessionRuntimeDbRowSchema,
// A no-op conflict update lets RETURNING report allowed reuse without changing the binding.
const insertRuntimeRow = SqlSchema.findOneOption({
Request: ProviderSessionRuntimeDbRowSchema.mapFields(
Struct.assign({ unlessNativeSessionId: Schema.NullOr(Schema.String) }),
),
Result: Schema.Struct({ threadId: ThreadId }),
execute: (runtime) =>
sql`
INSERT INTO provider_session_runtime (
Expand All @@ -252,7 +257,7 @@ export const make = Effect.gen(function* () {
resume_cursor_json,
runtime_payload_json
)
VALUES (
SELECT
${runtime.threadId},
${runtime.providerName},
${runtime.providerInstanceId},
Expand All @@ -266,8 +271,20 @@ export const make = Effect.gen(function* () {
THEN json_remove(${runtime.runtimePayload}, '$.importedTranscripts')
ELSE ${runtime.runtimePayload}
END
WHERE ${runtime.unlessNativeSessionId} IS NULL OR NOT EXISTS (
SELECT 1 FROM provider_session_runtime
WHERE thread_id NOT LIKE 'import:%'
AND provider_name = ${runtime.providerName}
AND COALESCE(provider_instance_id, provider_name) = ${runtime.providerInstanceId}
AND CASE WHEN json_valid(resume_cursor_json) THEN
json_extract(resume_cursor_json, CASE provider_name
WHEN 'claudeAgent' THEN '$.resume'
WHEN 'codex' THEN '$.threadId'
END)
END = ${runtime.unlessNativeSessionId}
)
ON CONFLICT (thread_id) DO NOTHING
ON CONFLICT (thread_id) DO UPDATE SET thread_id = provider_session_runtime.thread_id
RETURNING thread_id AS "threadId"
`,
});

Expand Down Expand Up @@ -369,7 +386,13 @@ export const make = Effect.gen(function* () {
});

const upsert: ProviderSessionRuntimeRepository["Service"]["upsert"] = (runtime, options) =>
(options?.onConflict === "ignore" ? insertRuntimeRow(runtime) : upsertRuntimeRow(runtime)).pipe(
(options?.onConflict === "ignore"
? insertRuntimeRow({
...runtime,
unlessNativeSessionId: options.unlessNativeSessionId ?? null,
}).pipe(Effect.map(Option.isSome))
: upsertRuntimeRow(runtime).pipe(Effect.as(true))
).pipe(
Effect.mapError(
toPersistenceSqlOrDecodeError(
"ProviderSessionRuntimeRepository.upsert:query",
Expand Down
237 changes: 227 additions & 10 deletions apps/server/src/project/AgentSessionImporter.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -200,7 +200,7 @@ const runImport = (input: {

it.layer(NodeServices.layer)("AgentSessionImporter", (it) => {
describe("importRecentAgentThreads", () => {
it.effect("uses the project root and stores provider-specific resume cursors", () =>
it.effect("imports sessions owned by another provider instance", () =>
Effect.gen(function* () {
const commands: Array<OrchestrationCommand> = [];
const bindings: Array<ProviderSessionDirectory.ProviderRuntimeBinding> = [];
Expand Down Expand Up @@ -230,12 +230,24 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => {
latestSequence: Effect.succeed(0),
});
const directory = ProviderSessionDirectory.ProviderSessionDirectory.of({
upsert: (binding) => Effect.sync(() => void bindings.push(binding)),
upsert: (binding) => Effect.sync(() => bindings.push(binding)).pipe(Effect.as(true)),
getProvider: () => Effect.die("unused"),
recordImportedTranscript: () => Effect.void,
getBinding: () => Effect.succeedNone,
getBinding: (threadId) =>
Effect.succeed(
Option.fromUndefinedOr(bindings.find((binding) => binding.threadId === threadId)),
),
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
listBindings: () =>
Effect.succeed([
{
threadId: ThreadId.make("native-other-instance"),
provider: ProviderDriverKind.make("codex"),
providerInstanceId: ProviderInstanceId.make("codex-other"),
resumeCursor: { threadId: "codex-session" },
lastSeenAt: "2026-08-24T10:00:00.000Z",
},
]),
});

const result = yield* runImport({
Expand Down Expand Up @@ -340,7 +352,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => {
recordImportedTranscript: () => Effect.die("unused"),
getBinding: () => Effect.die("must not read a scanner skip binding"),
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
listBindings: () => Effect.succeed([]),
});

const result = yield* runImport({
Expand Down Expand Up @@ -410,15 +422,15 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => {
}),
);
}
bindings.push(binding);
return Effect.void;
if (bindings.length === 0) bindings.push(binding);
return Effect.succeed(true);
},
getProvider: () => Effect.die("unused"),
recordImportedTranscript: () => Effect.void,
getBinding: () =>
Effect.succeed(bindings[0] === undefined ? Option.none() : Option.some(bindings[0])),
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
listBindings: () => Effect.succeed([]),
});
const snapshots = makeSnapshotsLayer({
project: makeProject(),
Expand Down Expand Up @@ -459,7 +471,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => {
recordImportedTranscript: () => Effect.void,
getBinding: () => Effect.succeedSome(runningBinding),
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
listBindings: () => Effect.succeed([]),
});
const engine = OrchestrationEngine.OrchestrationEngineService.of({
dispatch: () => Effect.die("must not replay history or settle active work"),
Expand Down Expand Up @@ -514,7 +526,7 @@ it.layer(NodeServices.layer)("AgentSessionImporter", (it) => {
recordImportedTranscript: () => Effect.die("unused"),
getBinding: () => Effect.succeedNone,
listThreadIds: () => Effect.die("unused"),
listBindings: () => Effect.die("unused"),
listBindings: () => Effect.succeed([]),
});

const result = yield* runImport({
Expand Down Expand Up @@ -582,6 +594,211 @@ const integrationLayer = Layer.mergeAll(
);

it.layer(integrationLayer)("AgentSessionImporter integration", (it) => {
for (const source of ["codex", "claudeAgent"] as const) {
for (const timing of ["before scan", "before reservation"] as const) {
it.effect(`skips native ${source} sessions bound ${timing}`, () =>
Effect.gen(function* () {
const engine = yield* OrchestrationEngine.OrchestrationEngineService;
const snapshots = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery;
const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory;
const projectId = ProjectId.make(`native-import-${source}-${timing}`);
const threadId = ThreadId.make(`native-${source}-${timing}`);
const thread = {
...makeThread(source),
providerSessionId:
source === "codex"
? `native-codex-session-${timing}`
: timing === "before scan"
? "123e4567-e89b-42d3-a456-426614174001"
: "123e4567-e89b-42d3-a456-426614174002",
};
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make(`create-${projectId}`),
projectId,
title: "Native project",
workspaceRoot: `/tmp/${projectId}`,
defaultModelSelection: null,
createdAt: thread.createdAt,
});
yield* engine.dispatch({
type: "thread.create",
commandId: CommandId.make(`create-${threadId}`),
threadId,
projectId,
title: "Native conversation",
modelSelection: { instanceId: thread.providerInstanceId, model: "default" },
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
createdAt: thread.createdAt,
});
const bindNative = directory.upsert({
threadId,
provider: ProviderDriverKind.make(source),
providerInstanceId: thread.providerInstanceId,
status: "stopped",
resumeCursor:
source === "codex"
? { threadId: thread.providerSessionId }
: { threadId, resume: thread.providerSessionId },
});
if (timing === "before scan") yield* bindNative;
const before = yield* snapshots.getThreadDetailById(threadId);
const result = yield* importRecentAgentThreads({ projectId }).pipe(
Effect.provideService(ProviderSessionDirectory.ProviderSessionDirectory, {
...directory,
upsert: (binding, options) =>
bindNative.pipe(Effect.andThen(directory.upsert(binding, options))),
}),
Effect.provideService(AgentSessionScanner.AgentSessionScanner, {
scan: Effect.die("unused"),
recentThreads: () => Stream.succeed(makeThreadOutcome(thread)),
}),
);
const importedId = ThreadId.make(`import:${source}:${thread.providerSessionId}`);
expect(result).toEqual({ importedCount: 0, skippedCount: 1 });
expect(yield* snapshots.getThreadDetailById(importedId)).toEqual(Option.none());
expect(yield* directory.getBinding(importedId)).toEqual(Option.none());
expect(yield* snapshots.getThreadDetailById(threadId)).toEqual(before);
const otherInstance = ProviderInstanceId.make(`${source}-other`);
const otherThreadId = ThreadId.make(
`import:${otherInstance}:${thread.providerSessionId}`,
);
yield* directory.upsert(
{
threadId: otherThreadId,
provider: ProviderDriverKind.make(source),
providerInstanceId: otherInstance,
resumeCursor:
source === "codex"
? { threadId: thread.providerSessionId }
: { resume: thread.providerSessionId },
},
{ onConflict: "ignore", unlessNativeSessionId: thread.providerSessionId },
);
expect(Option.isSome(yield* directory.getBinding(otherThreadId))).toBe(true);
}),
);
}
}

for (const source of ["codex", "claudeAgent"] as const) {
for (const owner of ["same instance", "other instance", "none"] as const) {
it.effect(`rechecks a reserved ${source} import with native owner in ${owner}`, () =>
Effect.gen(function* () {
const engine = yield* OrchestrationEngine.OrchestrationEngineService;
const snapshots = yield* ProjectionSnapshotQuery.ProjectionSnapshotQuery;
const directory = yield* ProviderSessionDirectory.ProviderSessionDirectory;
const projectId = ProjectId.make(`retry-native-${source}-${owner}`);
const nativeId = ThreadId.make(`native-${projectId}`);
const thread = {
...makeThread(source),
providerSessionId:
source === "codex"
? `retry-native-${owner}`
: owner === "same instance"
? "123e4567-e89b-42d3-a456-426614174011"
: owner === "other instance"
? "123e4567-e89b-42d3-a456-426614174012"
: "123e4567-e89b-42d3-a456-426614174013",
};
const importedId = ThreadId.make(`import:${source}:${thread.providerSessionId}`);
yield* engine.dispatch({
type: "project.create",
commandId: CommandId.make(`create-${projectId}`),
projectId,
title: "Retry native ownership",
workspaceRoot: `/tmp/${projectId}`,
defaultModelSelection: null,
createdAt: thread.createdAt,
});
yield* engine.dispatch({
type: "thread.create",
commandId: CommandId.make(`create-${nativeId}`),
threadId: nativeId,
projectId,
title: "Native conversation",
modelSelection: { instanceId: thread.providerInstanceId, model: "default" },
runtimeMode: "full-access",
interactionMode: "default",
branch: null,
worktreePath: null,
createdAt: thread.createdAt,
});
const scanner = AgentSessionScanner.AgentSessionScanner.of({
scan: Effect.die("unused"),
recentThreads: () => Stream.succeed(makeThreadOutcome(thread)),
});
const failed = yield* importRecentAgentThreads({ projectId }).pipe(
Effect.provideService(AgentSessionScanner.AgentSessionScanner, scanner),
Effect.provideService(OrchestrationEngine.OrchestrationEngineService, {
...engine,
dispatch: (command) =>
command.type === "thread.create"
? Effect.fail(
new OrchestrationCommandInvariantError({
commandType: command.type,
detail: "Injected thread creation failure after reservation.",
}),
)
: engine.dispatch(command),
}),
);
expect(failed).toEqual({ importedCount: 0, skippedCount: 1 });
expect(yield* snapshots.getThreadDetailById(importedId)).toEqual(Option.none());
const reservation = yield* directory.getBinding(importedId);
expect(Option.getOrThrow(reservation).status).toBe("stopped");
const nativeBefore = yield* snapshots.getThreadDetailById(nativeId);

const result = yield* importRecentAgentThreads({ projectId }).pipe(
Effect.provideService(AgentSessionScanner.AgentSessionScanner, {
...scanner,
recentThreads: () =>
Stream.fromEffect(
Effect.gen(function* () {
if (owner !== "none") {
yield* directory.upsert({
threadId: nativeId,
provider: ProviderDriverKind.make(source),
providerInstanceId:
owner === "same instance"
? thread.providerInstanceId
: ProviderInstanceId.make(`${source}-other`),
status: "running",
resumeCursor:
source === "codex"
? { threadId: thread.providerSessionId }
: { threadId: nativeId, resume: thread.providerSessionId },
});
}
return makeThreadOutcome(thread);
}).pipe(Effect.orDie),
),
}),
);
expect(yield* snapshots.getThreadDetailById(nativeId)).toEqual(nativeBefore);
if (owner === "same instance") {
expect(result).toEqual({ importedCount: 0, skippedCount: 1 });
expect(yield* snapshots.getThreadDetailById(importedId)).toEqual(Option.none());
expect(yield* directory.getBinding(importedId)).toEqual(reservation);
} else {
expect(result).toEqual({ importedCount: 1, skippedCount: 0 });
expect(
Option.getOrThrow(yield* snapshots.getThreadDetailById(importedId)).messages.map(
(message) => message.text,
),
).toEqual(thread.messages.map((message) => message.text));
expect(Option.getOrThrow(yield* directory.getBinding(importedId)).resumeCursor).toEqual(
Option.getOrThrow(reservation).resumeCursor,
);
}
}),
);
}
}

it.effect("imports once after the real engine persists an old rejected receipt", () =>
Effect.gen(function* () {
const engine = yield* OrchestrationEngine.OrchestrationEngineService;
Expand Down
Loading
Loading