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
159 changes: 159 additions & 0 deletions apps/daemon/src/__tests__/conversation-store.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -111,6 +111,16 @@ function openOperation(operationId: string, sessionId = 's-1'): ConversationOper
});
}

function openFork(operationId: string) {
return {
operationId: OperationIdSchema.parse(operationId),
sessionId: SessionIdSchema.parse('s-1'),
kind: 'session.fork' as const,
state: 'open' as const,
createdAt: 4,
};
}

async function seedIntent(store: ConversationStore): Promise<void> {
await store.persistTurnIntent({
turn: turn({
Expand Down Expand Up @@ -481,4 +491,153 @@ describe('SQLite conversation store', () => {
{ ...first, state: 'failed' },
]);
});

it('a turn-less operation respects the open-operation gate and a replayed id', async () => {
const { database } = await databaseWithSessions('s-1');
const store = createConversationStore(database.client);
const fork = openFork('op-fork');
await store.persistOperation(fork);

expect(await store.listOpenOperations(SessionIdSchema.parse('s-1'))).toEqual([fork]);
await expect(async () =>
store.persistTurnIntent({ turn: turn({ turnId: 't-1' }), operation: openOperation('op-1') }),
).rejects.toBeInstanceOf(ConversationSessionBusyError);
await expect(async () => store.persistOperation(openFork('op-fork-2'))).rejects.toBeInstanceOf(
ConversationSessionBusyError,
);
await store.resolveOperation({
...fork,
state: 'failed',
error: { code: 'unsupported', message: 'no checkpoint' },
resolvedAt: 5,
});
await expect(async () => store.persistOperation(fork)).rejects.toThrow('UNIQUE');
});

it('commitFork writes the child session, its runs, and its turns with the operation atomically', async () => {
const { path, database } = await databaseWithSessions('s-1');
const store = createConversationStore(database.client);
await seedIntent(store);
await store.resolveOperation({
...openOperation('op-1'),
state: 'succeeded',
turnId: TurnIdSchema.parse('t-prompted'),
resolvedAt: 5,
});
const fork = openFork('op-fork');
await store.persistOperation(fork);
const child = SessionRecordSchema.parse({
sessionId: 's-child',
kind: 'claude-code',
cwd: '/repo',
origin: { type: 'created' },
forkOrigin: { sourceSessionId: 's-1', sourceTurnId: 't-prompted', forkedAt: 6 },
createdAt: 6,
updatedAt: 6,
runs: [
{ runId: 'run-child', baseTurnId: 't-copied-2', historyId: 'native-child', startedAt: 6 },
],
graphRevision: 0,
eventEpoch: 0,
});
const copied = [
turn({
turnId: 't-copied-1',
sessionId: 's-child',
input: { type: 'prompt', promptId: PromptIdSchema.parse('p-1') },
runId: 'run-child',
state: 'completed',
createdAt: 2,
}),
turn({
turnId: 't-copied-2',
sessionId: 's-child',
parentTurnId: TurnIdSchema.parse('t-copied-1'),
runId: 'run-child',
state: 'completed',
createdAt: 3,
}),
];
const succeeded = {
...fork,
state: 'succeeded' as const,
turnId: TurnIdSchema.parse('t-copied-2'),
resolvedAt: 7,
};

expect(await store.commitFork({ child, turns: copied, operation: succeeded })).toBe(true);

closeDatabase(database);
const reopened = openDatabase(path);
const reopenedStore = createConversationStore(reopened.client);
expect(await createSessionStore(reopened.client).load()).toContainEqual(child);
expect(await reopenedStore.listTurns(child.sessionId)).toEqual(copied);
expect(await reopenedStore.getTurn(TurnIdSchema.parse('t-copied-2'))).toEqual(copied[1]);
expect(await reopenedStore.getOperation(fork.operationId)).toEqual(succeeded);
// The source's deletion keeps the prompt the child still references.
await reopenedStore.deleteSession(SessionIdSchema.parse('s-1'));
expect(await reopenedStore.getPrompt(PromptIdSchema.parse('p-1'))).toEqual(prompt('p-1'));
});

it('commitFork writes nothing when the operation already resolved or a row is refused', async () => {
const { database } = await databaseWithSessions('s-1');
const store = createConversationStore(database.client);
const sessionStore = createSessionStore(database.client);
const child = SessionRecordSchema.parse({
sessionId: 's-child',
kind: 'claude-code',
cwd: '/repo',
origin: { type: 'created' },
createdAt: 6,
updatedAt: 6,
runs: [],
});
const failed = openFork('op-fork');
await store.persistOperation(failed);
await store.resolveOperation({
...failed,
state: 'failed',
error: { code: 'timeout', message: 'too slow' },
resolvedAt: 5,
});
expect(
await store.commitFork({
child,
turns: [turn({ turnId: 't-late', sessionId: 's-child', state: 'completed' })],
operation: {
...failed,
state: 'succeeded',
turnId: TurnIdSchema.parse('t-late'),
resolvedAt: 7,
},
}),
).toBe(false);
expect(await sessionStore.load()).toHaveLength(1);

// A copied turn naming a prompt that no longer exists rolls the whole commit back.
const open = openFork('op-fork-2');
await store.persistOperation(open);
await expect(async () =>
store.commitFork({
child,
turns: [
turn({
turnId: 't-orphan',
sessionId: 's-child',
input: { type: 'prompt', promptId: PromptIdSchema.parse('p-gone') },
state: 'completed',
}),
],
operation: {
...open,
state: 'succeeded',
turnId: TurnIdSchema.parse('t-orphan'),
resolvedAt: 8,
},
}),
).rejects.toThrow('FOREIGN KEY');
expect(await sessionStore.load()).toHaveLength(1);
expect(await store.getOperation(open.operationId)).toEqual(open);
expect(await store.getTurn(TurnIdSchema.parse('t-orphan'))).toBeUndefined();
});
});
100 changes: 76 additions & 24 deletions apps/daemon/src/conversation-store.ts
Original file line number Diff line number Diff line change
@@ -1,4 +1,8 @@
import type { ConversationStore, ConversationTurnIntent } from '@linkcode/engine';
import type {
ConversationForkCommit,
ConversationStore,
ConversationTurnIntent,
} from '@linkcode/engine';
import { ConversationSessionBusyError } from '@linkcode/engine';
import type {
ConversationOperation,
Expand All @@ -24,8 +28,11 @@ import {
promptAttachmentRefs,
prompts,
providerTurnBindings,
sessionRuns,
sessions,
uploadLeases,
} from './db/schema';
import { toSessionRow, toSessionRunRows } from './session-store';

type TurnRow = typeof conversationTurns.$inferSelect;
type PromptRow = typeof prompts.$inferSelect;
Expand All @@ -47,6 +54,36 @@ export function createConversationStore(db: DaemonDatabaseClient): ConversationS
.run();
}

function hasOpenOperation(tx: DbOrTx, sessionId: SessionId): boolean {
const open = tx
.select({ operationId: conversationOperations.operationId })
.from(conversationOperations)
.where(
and(
eq(conversationOperations.sessionId, sessionId),
eq(conversationOperations.state, 'open'),
),
)
.get();
return open !== undefined;
}

/** The open→terminal transition every resolver races for; `changes === 0` means a first writer
* already stored a terminal result. */
function transitionOperation(tx: DbOrTx, operation: ConversationOperation): boolean {
const result = tx
.update(conversationOperations)
.set(toOperationRow(operation))
.where(
and(
eq(conversationOperations.operationId, operation.operationId),
eq(conversationOperations.state, 'open'),
),
)
.run();
return result.changes > 0;
}

return {
listTurns(sessionId: SessionId): Promise<ConversationTurn[]> {
const rows = db
Expand All @@ -58,6 +95,15 @@ export function createConversationStore(db: DaemonDatabaseClient): ConversationS
return Promise.resolve(rows.map(toTurn));
},

getTurn(turnId: TurnId): Promise<ConversationTurn | undefined> {
const row = db
.select()
.from(conversationTurns)
.where(eq(conversationTurns.turnId, turnId))
.get();
return Promise.resolve(row ? toTurn(row) : undefined);
},

saveTurn(turn: ConversationTurn): Promise<void> {
upsertTurn(db, turn);
return Promise.resolve();
Expand Down Expand Up @@ -113,17 +159,7 @@ export function createConversationStore(db: DaemonDatabaseClient): ConversationS
persistTurnIntent(intent: ConversationTurnIntent): Promise<ConversationTurn> {
const { parentTurnId, sessionId } = intent.turn;
const persisted = db.transaction((tx) => {
const open = tx
.select({ operationId: conversationOperations.operationId })
.from(conversationOperations)
.where(
and(
eq(conversationOperations.sessionId, sessionId),
eq(conversationOperations.state, 'open'),
),
)
.get();
if (open) throw new ConversationSessionBusyError(sessionId);
if (hasOpenOperation(tx, sessionId)) throw new ConversationSessionBusyError(sessionId);
const siblings = tx
.select({ value: count() })
.from(conversationTurns)
Expand Down Expand Up @@ -167,26 +203,42 @@ export function createConversationStore(db: DaemonDatabaseClient): ConversationS
return Promise.resolve(persisted);
},

persistOperation(operation: Extract<ConversationOperation, { state: 'open' }>): Promise<void> {
db.transaction((tx) => {
if (hasOpenOperation(tx, operation.sessionId)) {
throw new ConversationSessionBusyError(operation.sessionId);
}
// Plain insert: a replayed operationId must conflict here, never re-open a terminal row.
tx.insert(conversationOperations).values(toOperationRow(operation)).run();
});
return Promise.resolve();
},

resolveOperation(operation: ConversationOperation, turn?: ConversationTurn): Promise<boolean> {
const transitioned = db.transaction((tx) => {
const result = tx
.update(conversationOperations)
.set(toOperationRow(operation))
.where(
and(
eq(conversationOperations.operationId, operation.operationId),
eq(conversationOperations.state, 'open'),
),
)
.run();
// A concurrent resolver already stored a terminal result; the first writer stands.
if (result.changes === 0) return false;
if (!transitionOperation(tx, operation)) return false;
if (turn) upsertTurn(tx, turn);
return true;
});
return Promise.resolve(transitioned);
},

commitFork(commit: ConversationForkCommit): Promise<boolean> {
const transitioned = db.transaction((tx) => {
if (!transitionOperation(tx, commit.operation)) return false;
// Plain inserts throughout: the child id is fresh, and a parent row must precede its child
// (the caller hands the copied lineage root-first) for the self-referencing foreign key.
tx.insert(sessions).values(toSessionRow(commit.child)).run();
const runs = toSessionRunRows(commit.child);
if (runs.length > 0) tx.insert(sessionRuns).values(runs).run();
for (let i = 0, len = commit.turns.length; i < len; i++) {
tx.insert(conversationTurns).values(toTurnRow(commit.turns[i])).run();
}
return true;
});
return Promise.resolve(transitioned);
},

deleteSession(sessionId: SessionId): Promise<void> {
db.transaction((tx) => {
const rows = tx
Expand Down
45 changes: 23 additions & 22 deletions apps/daemon/src/session-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -43,27 +43,8 @@ export function createSessionStore(db: DaemonDatabaseClient): SessionStore {
.run();
// Runs are few per session; rewriting them keeps save() a whole-record upsert.
tx.delete(sessionRuns).where(eq(sessionRuns.sessionId, record.sessionId)).run();
if (record.runs.length > 0) {
tx.insert(sessionRuns)
.values(
record.runs.map((run, seq) => ({
sessionId: record.sessionId,
seq,
// runId is optional at the wire parse boundary only; every writer mints it, so a
// runId-less run here is a bug — minting one would drift the durable id per save.
runId: nullthrow(run.runId, `Session run without runId: ${record.sessionId}`),
baseTurnId: run.baseTurnId ?? null,
historyId: run.historyId ?? null,
accountId: run.accountId ?? null,
model: run.model ?? null,
effort: run.effort ?? null,
approvalPolicyId: run.approvalPolicyId ?? null,
startedAt: run.startedAt,
endedAt: run.endedAt ?? null,
})),
)
.run();
}
const runs = toSessionRunRows(record);
if (runs.length > 0) tx.insert(sessionRuns).values(runs).run();
});
return Promise.resolve();
},
Expand All @@ -76,7 +57,9 @@ export function createSessionStore(db: DaemonDatabaseClient): SessionStore {
};
}

function toSessionRow(record: SessionRecord): typeof sessions.$inferInsert {
/** Also used by the conversation store, whose fork commit inserts the child session row in the
* same transaction as the turns that reference it. */
export function toSessionRow(record: SessionRecord): typeof sessions.$inferInsert {
return {
sessionId: record.sessionId,
kind: record.kind,
Expand All @@ -99,6 +82,24 @@ function toSessionRow(record: SessionRecord): typeof sessions.$inferInsert {
};
}

export function toSessionRunRows(record: SessionRecord): Array<typeof sessionRuns.$inferInsert> {
return record.runs.map((run, seq) => ({
sessionId: record.sessionId,
seq,
// runId is optional at the wire parse boundary only; every writer mints it, so a runId-less
// run here is a bug — minting one would drift the durable id per save.
runId: nullthrow(run.runId, `Session run without runId: ${record.sessionId}`),
baseTurnId: run.baseTurnId ?? null,
historyId: run.historyId ?? null,
accountId: run.accountId ?? null,
model: run.model ?? null,
effort: run.effort ?? null,
approvalPolicyId: run.approvalPolicyId ?? null,
startedAt: run.startedAt,
endedAt: run.endedAt ?? null,
}));
}

function toRecord(row: SessionRow, runRows: RunRow[]): SessionRecord {
return SessionRecordSchema.parse({
sessionId: row.sessionId,
Expand Down
Loading
Loading