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
86 changes: 86 additions & 0 deletions packages/local-runtime-v2/src/infra/db/write-transaction.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,86 @@
import { setTimeout as delay } from 'node:timers/promises';
import { sql } from 'drizzle-orm';
import type { AppDb } from './client.js';

const WRITE_LOCK_BUDGET_MS = 10_000;
const WRITE_LOCK_ATTEMPT_MS = 50;

/** Cancellation before the mutation callback starts; no write needs to be replayed. */
export class WriteLockWaitAbortedError extends Error {
override readonly name = 'WriteLockWaitAbortedError';

constructor(readonly signal: AbortSignal) {
super('SQLite write lock wait was cancelled', { cause: signal.reason });
}
}

/**
* Retry only transaction admission: a callback that has started is never replayed.
* The signal cancels contention waits, not an immediately available write. This
* lets post-cancellation tool completion and cleanup messages remain durable.
*/
export async function runWithWriteLock<T>(
db: AppDb,
mutation: (tx: AppDb) => T,
options: { readonly signal?: AbortSignal; readonly timeoutMs?: number } = {},
): Promise<T> {
const budget = options.timeoutMs ?? WRITE_LOCK_BUDGET_MS;
if (!Number.isFinite(budget) || budget <= 0)
throw new RangeError('Invalid write lock budget');
const deadline = performance.now() + budget;
let attempt = 0;
let hasContended = false;
let lastBusy: unknown = new Error('SQLite write lock wait exceeded its deadline');
for (;;) {
if (hasContended) throwIfWaitAborted(options.signal);
const remaining = deadline - performance.now();
if (remaining <= 0) throw lastBusy;
const previous = db.get<{ timeout: number }>(sql`PRAGMA busy_timeout`).timeout;
let entered = false;
try {
const nativeWaitMs = options.signal?.aborted
? 0
: Math.ceil(Math.min(WRITE_LOCK_ATTEMPT_MS, remaining));
db.run(sql.raw(`PRAGMA busy_timeout = ${nativeWaitMs}`));
return db.transaction(
(tx) => {
entered = true;
// Only lock acquisition gets a short timeout. Restore the connection's
// policy before callbacks (including nested transactions) can use it.
db.run(sql.raw(`PRAGMA busy_timeout = ${previous}`));
return mutation(tx);
},
{ behavior: 'immediate' },
);
} catch (error) {
if (entered || !isBusy(error)) throw error;
lastBusy = error;
} finally {
// No await occurs while the shared connection has a temporary timeout.
db.run(sql.raw(`PRAGMA busy_timeout = ${previous}`));
}
hasContended = true;
throwIfWaitAborted(options.signal);
const wait = Math.min(
deadline - performance.now(),
25 * 2 ** Math.min(attempt++, 3) + Math.random() * 25,
);
if (wait <= 0) throw lastBusy;
try {
await delay(wait, undefined, { signal: options.signal });
} catch (error) {
if (options.signal?.aborted && error instanceof Error && error.name === 'AbortError') {
throw new WriteLockWaitAbortedError(options.signal);
}
throw error;
}
}
}

function throwIfWaitAborted(signal: AbortSignal | undefined): void {
if (signal?.aborted) throw new WriteLockWaitAbortedError(signal);
}

function isBusy(error: unknown): boolean {
return error instanceof Error && Reflect.get(error, 'code') === 'SQLITE_BUSY';
}
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,7 @@ export interface SessionAgentEventContext {
export interface SessionAgentProjectionInput {
readonly context: SessionAgentEventContext;
readonly event: RuntimeEvent;
readonly signal?: AbortSignal;
}

export interface SessionSystemAgentProjectionOptions {
Expand Down Expand Up @@ -76,19 +77,25 @@ export function createSessionSystemAgentProjection(options: SessionSystemAgentPr
const message = displayMessage(input.event);
if (message && !isTerminalAssistantDisplayError(message)) {
const queryKey = await queryKeyForTurn(options, input.context);
await options.messages.upsert({
sessionId: input.context.sessionId,
turnId: input.context.turnId,
message: {
...message,
turn_id: input.context.turnId,
...(queryKey ? { query_key: queryKey } : {}),
await options.messages.upsert(
{
sessionId: input.context.sessionId,
turnId: input.context.turnId,
message: {
...message,
turn_id: input.context.turnId,
...(queryKey ? { query_key: queryKey } : {}),
},
source: input.context.provenance?.source ?? 'agent',
...(input.context.provenance?.sourceContext
? { sourceContext: input.context.provenance.sourceContext }
: {}),
},
source: input.context.provenance?.source ?? 'agent',
...(input.context.provenance?.sourceContext
? { sourceContext: input.context.provenance.sourceContext }
: {}),
});
// Complete tool messages describe work already executed, including
// abort cleanup. Persist these facts within the write-lock budget
// even when the lease is cancelled; text-only waits may stop early.
message.tool_calls?.length ? undefined : { signal: input.signal },
);
}
await projectQueryCollapse(() => options.queryCollapse?.projectRuntimeEvent(input));
},
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,10 @@ export interface MessageUpsertInput {
readonly source?: string;
readonly sourceContext?: Record<string, unknown>;
}
export interface MessageWriteOptions {
/** Cancels lock contention waits; immediately available cleanup writes still commit. */
readonly signal?: AbortSignal;
}
export interface UserMessageCommitInput extends MessageUpsertInput {
readonly unconsumedFromTurnIds?: readonly string[];
/** Trusted Queue startup lineage; the initial Host never read these rows. */
Expand Down Expand Up @@ -97,8 +101,11 @@ export interface MessageRepository {
},
): Promise<DisplayMessageRecord[]>;
commitUserMessage(input: UserMessageCommitInput): Promise<UserMessageCommitResult>;
upsert(input: MessageUpsertInput): Promise<NormalizedDisplayMessage>;
upsertMany(inputs: readonly MessageUpsertInput[]): Promise<readonly NormalizedDisplayMessage[]>;
upsert(input: MessageUpsertInput, options?: MessageWriteOptions): Promise<NormalizedDisplayMessage>;
upsertMany(
inputs: readonly MessageUpsertInput[],
options?: MessageWriteOptions,
): Promise<readonly NormalizedDisplayMessage[]>;
replace(input: MessageReplaceInput): Promise<void>;
replaceStream(input: MessageReplaceStreamInput): Promise<void>;
rewindInclusive(input: MessageRewindInclusiveInput): Promise<MessageRewindInclusiveResult>;
Expand Down
Original file line number Diff line number Diff line change
@@ -1,6 +1,7 @@
import { and, asc, desc, eq, gt, gte, inArray, lt, lte, sql } from 'drizzle-orm';

import type { AppDb } from '../../../../infra/db/client.js';
import { runWithWriteLock } from '../../../../infra/db/write-transaction.js';
import {
legacyMessages,
messageRowMigrations,
Expand All @@ -22,6 +23,7 @@ import type {
MessageRewindInclusiveInput,
MessageRewindInclusiveResult,
MessageUpsertInput,
MessageWriteOptions,
DisplayMessageRecord,
NormalizedDisplayMessage,
UserMessageCommitInput,
Expand Down Expand Up @@ -205,14 +207,18 @@ class DrizzleMessageRepository implements MessageRepository {
);
}

async upsert(input: MessageUpsertInput): Promise<NormalizedDisplayMessage> {
const [result] = await this.upsertMany([input]);
async upsert(
input: MessageUpsertInput,
options?: MessageWriteOptions,
): Promise<NormalizedDisplayMessage> {
const [result] = await this.upsertMany([input], options);
if (!result) throw new Error('Message upsert returned no row');
return result;
}

async upsertMany(
inputs: readonly MessageUpsertInput[],
options?: MessageWriteOptions,
): Promise<readonly NormalizedDisplayMessage[]> {
const sessionId = inputs[0]?.sessionId;
if (!sessionId) return [];
Expand All @@ -229,11 +235,16 @@ class DrizzleMessageRepository implements MessageRepository {
}),
replaceProvenance: hasMessageProvenance(input),
}));
this.mutationTransaction(sessionId, (tx) => {
normalized.forEach(({ message, replaceProvenance }) =>
this.write(tx, sessionId, message, replaceProvenance),
);
});
await runWithWriteLock(
this.options.db,
(tx) => {
this.ensureReadyInTransaction(tx, sessionId);
normalized.forEach(({ message, replaceProvenance }) =>
this.write(tx, sessionId, message, replaceProvenance),
);
},
options,
);
return normalized.map(({ message }) => message);
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -30,6 +30,10 @@ export type AgentEventResult =
};

export interface AgentEventDelivery {
handleRuntimeEvent(context: AgentEventContext, event: RuntimeEvent): Promise<AgentEventResult>;
handleRuntimeEvent(
context: AgentEventContext,
event: RuntimeEvent,
signal?: AbortSignal,
): Promise<AgentEventResult>;
handleHistoryCommitted(context: AgentEventContext, change: CommittedHistoryChange): Promise<void>;
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,8 @@ type ObservationStage = 'runtime-event' | 'history-committed' | 'history-failure
interface RuntimeProjectionInput {
readonly context: AgentEventContext;
readonly event: RuntimeEvent;
/** Process-local control; excluded from semantic snapshots and replay identities. */
readonly signal?: AbortSignal;
}

interface HistoryProjectionInput {
Expand Down Expand Up @@ -156,7 +158,11 @@ export class RequiredAgentEventDelivery implements AgentEventDelivery, AgentHost
this.historyReplays = new SemanticReplayRegistry(maximum);
}

handleRuntimeEvent(context: AgentEventContext, event: RuntimeEvent): Promise<AgentEventResult> {
handleRuntimeEvent(
context: AgentEventContext,
event: RuntimeEvent,
signal?: AbortSignal,
): Promise<AgentEventResult> {
try {
const snapshot = captureSemanticSnapshot({ context, event });
validateRuntimeInput(snapshot.value.context, snapshot.value.event);
Expand All @@ -167,7 +173,7 @@ export class RequiredAgentEventDelivery implements AgentEventDelivery, AgentHost
conflict: () => new AgentEventIdentityConflictError('runtime-event', identity),
execute: () =>
this.lane.run(snapshot.value.context.sessionId, () =>
this.projectRuntime(snapshot.value.context, snapshot.value.event),
this.projectRuntime(snapshot.value.context, snapshot.value.event, signal),
),
});
} catch (error) {
Expand Down Expand Up @@ -223,6 +229,7 @@ export class RequiredAgentEventDelivery implements AgentEventDelivery, AgentHost
private async projectRuntime(
context: AgentEventContext,
event: RuntimeEvent,
signal?: AbortSignal,
): Promise<AgentEventResult> {
const runtimeSequence = this.validateSequence(context, event);
const authoritative = await this.options.projectors.session.projectRuntimeEvent({
Expand All @@ -233,7 +240,11 @@ export class RequiredAgentEventDelivery implements AgentEventDelivery, AgentHost
throw new AgentEventAcknowledgementError(terminalOutcome(event) ?? 'non-terminal', 'missing');
}
validateAcknowledgement(event, authoritative);
await this.options.projectors.messages.projectRuntimeEvent({ context, event });
await this.options.projectors.messages.projectRuntimeEvent({
context,
event,
...(signal ? { signal } : {}),
});
await this.options.projectors.stream.projectRuntimeEvent({ context, event });
await this.options.projectors.turnFacts.projectRuntimeEvent({ context, event });
this.commitSequence(context, runtimeSequence);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,7 @@ import { TurnCommittedHistoryState } from '../history/turn-committed-history-sta
import { CanonicalUserMessageIdentityLane } from '../history/canonical-user-message-identities.js';
import { AgentEventAssociationError } from '../preparation/turn-preflight.js';
import type { UserMessageId } from '../../../session-system/index.js';
import { WriteLockWaitAbortedError } from '../../../../infra/db/write-transaction.js';

export class AgentTerminalConfirmationError extends Error {
override readonly name = 'AgentTerminalConfirmationError';
Expand Down Expand Up @@ -147,11 +148,26 @@ export class TurnCommitPipeline {
}
await this.lane.enqueue(async () => {
validateRuntimeAssociation(this.dependencies.context, event);
const result = await this.dependencies.events.handleRuntimeEvent(
this.dependencies.context,
event,
);
new TerminalConfirmation().observe(event, result);
const signal = this.dependencies.lease.signal;
try {
const result = await this.dependencies.events.handleRuntimeEvent(
this.dependencies.context,
event,
signal,
);
new TerminalConfirmation().observe(event, result);
} catch (error) {
if (
error instanceof WriteLockWaitAbortedError &&
error.signal === signal &&
signal.aborted
) {
// This lease cancelled a projection before its write began. Keep the
// lane available for the runner's abort reconciliation and terminal.
return;
}
throw error;
}
});
};

Expand Down
2 changes: 2 additions & 0 deletions release/public-source.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

Loading
Loading