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
Original file line number Diff line number Diff line change
Expand Up @@ -529,7 +529,9 @@ The callback receives a `TransactionContext` (`tx`) with three slots, all bound

### Lifecycle and connection scoping

The runtime acquires a connection (auto-connecting lazily, like `db.runtime()`), issues `BEGIN`, runs the callback, and either commits and resolves with the callback's return value or rolls back and re-throws the original error. The connection is released in both branches — including when the callback hangs or the process crashes — via `finally` semantics. If `COMMIT` itself fails, the promise rejects with the commit error. If `ROLLBACK` fails after a callback throw, the rollback error wraps the original.
The runtime resolves the contract marker gate if this is its first operation, acquires a connection (auto-connecting lazily, like `db.runtime()`), issues `BEGIN`, runs the callback, and either commits and resolves with the callback's return value or rolls back and re-throws the original error. The connection is released in both branches — including when the callback hangs or the process crashes — via `finally` semantics. If `COMMIT` itself fails, the promise rejects with the commit error. If `ROLLBACK` fails after a callback throw, the rollback error wraps the original.

The marker read always finishes before the connection is acquired, so it never runs inside the transaction. The runtime reads the marker through the driver. A single-connection driver serves that read on the socket the transaction uses, so a read after `BEGIN` would be the transaction's first statement: PostgreSQL would then reject a following `SET TRANSACTION` with SQLSTATE 25001, a REPEATABLE READ or SERIALIZABLE transaction would take its snapshot at the marker read, and a failed marker read would abort the transaction. A pooled driver serves the read on a second client, which never arrives when the pool's only client is the one the transaction holds. `connection()` and the protected `acquireRawConnection()` await the same single-flight gate statements use, so the marker is still read at most once per runtime and never when `verifyMarker` is `false`.

ORM nested mutations (`withMutationScope` in the SQL ORM mutation executor) detect that they are running inside a transaction context and reuse the transaction's `RuntimeQueryable` instead of acquiring a new connection. This is why `RuntimeConnection` and `RuntimeTransaction` extend the same `RuntimeScope` slice that `RuntimeQueryable` does — both compose against one canonical surface (see [ORM Client Integration](#orm-client-integration) above).

Expand Down
22 changes: 16 additions & 6 deletions packages/2-sql/5-runtime/src/sql-runtime.ts
Original file line number Diff line number Diff line change
Expand Up @@ -363,18 +363,28 @@ export abstract class SqlRuntimeBase<TContract extends Contract<SqlStorage> = Co
* issued on it runs below the middleware/codec/telemetry pipeline. It carries
* its own lifecycle (`release`/`destroy`/`beginTransaction`); the caller owns
* disposal.
*
* The contract marker gate resolves before the connection is acquired: the marker is read
* through the driver, which is this connection's socket on a single-connection driver and a
* second client on a pool. The read therefore never runs inside a transaction begun on the
* connection and never waits on a pool whose only client is already held.
*/
protected acquireRawConnection(): Promise<SqlConnection> {
protected async acquireRawConnection(): Promise<SqlConnection> {
await this.ensureMarkerVerified();
return this.driver.acquireConnection();
}

private async setupDriverExecution(exec: SqlExecutionPlan): Promise<void> {
this.familyAdapter.validatePlan(exec, this.contract);
this._telemetry = null;
private ensureMarkerVerified(): Promise<void> {
if (this.verifyMarkerPromise === null) {
this.verifyMarkerPromise = this.verifyMarker();
}
await this.verifyMarkerPromise;
return this.verifyMarkerPromise;
}

private async setupDriverExecution(exec: SqlExecutionPlan): Promise<void> {
this.familyAdapter.validatePlan(exec, this.contract);
this._telemetry = null;
await this.ensureMarkerVerified();
}

protected getListDecoder(): ListDecoder {
Expand Down Expand Up @@ -763,7 +773,7 @@ export abstract class SqlRuntimeBase<TContract extends Contract<SqlStorage> = Co
}

async connection(): Promise<RuntimeConnection> {
const driverConn = await this.driver.acquireConnection();
const driverConn = await this.acquireRawConnection();
const self = this;

const wrappedConnection: RuntimeConnection &
Expand Down
99 changes: 99 additions & 0 deletions packages/2-sql/5-runtime/test/marker-verification.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,7 @@ import type {
SqlRuntimeTargetDescriptor,
} from '../src/sql-context';
import { createExecutionContext, createSqlExecutionStack } from '../src/sql-context';
import { withTransaction } from '../src/sql-runtime';
import { defineTestCodec } from './test-codec';
import { createTestRuntime as createRuntime, descriptorsFromCodecs } from './utils';

Expand Down Expand Up @@ -376,3 +377,101 @@ describe('verifyMarker', () => {
expect(log.warn).toHaveBeenCalledTimes(1);
});
});

describe('verifyMarker and transactions', () => {
function createTransactionalDriver(events: string[]): SqlDriver {
const transaction = {
execute: vi.fn().mockResolvedValue({ affectedRows: 0 }),
query: vi.fn().mockImplementation(async function* (_request: SqlExecuteRequest) {
events.push('statement');
yield {} as Record<string, unknown>;
}),
commit: vi.fn().mockResolvedValue(undefined),
rollback: vi.fn().mockResolvedValue(undefined),
};
const connection = {
execute: vi.fn().mockResolvedValue({ affectedRows: 0 }),
query: vi.fn(),
release: vi.fn().mockImplementation(async () => {
events.push('release');
}),
destroy: vi.fn().mockResolvedValue(undefined),
beginTransaction: vi.fn().mockImplementation(async () => {
events.push('begin');
return transaction;
}),
};
return {
...createDriver(),
acquireConnection: vi.fn().mockImplementation(async () => {
events.push('acquire');
return connection;
}),
};
}

it('reads the marker before acquiring the connection when a transaction is the first operation', async () => {
const events: string[] = [];
const readMarkerSpy = vi.fn().mockImplementation(async () => {
events.push('marker');
return { kind: 'absent' };
});
const runtime = buildRuntime({
markerResult: { kind: 'absent' },
driver: createTransactionalDriver(events),
readMarkerSpy,
});

await withTransaction(runtime, (tx) => tx.query(createPlan()).toArray());

expect(events).toEqual(['marker', 'acquire', 'begin', 'statement', 'release']);
});

it('reads the marker once across transactions and plain queries', async () => {
const events: string[] = [];
const readMarkerSpy = vi.fn().mockResolvedValue({ kind: 'absent' });
const runtime = buildRuntime({
markerResult: { kind: 'absent' },
driver: createTransactionalDriver(events),
readMarkerSpy,
});

await withTransaction(runtime, (tx) => tx.query(createPlan()).toArray());
await withTransaction(runtime, (tx) => tx.query(createPlan()).toArray());
await runtime.query(createPlan()).toArray();

expect(readMarkerSpy).toHaveBeenCalledTimes(1);
});

it('begins the transaction without a marker read when verifyMarker is false', async () => {
const events: string[] = [];
const readMarkerSpy = vi.fn().mockResolvedValue({ kind: 'absent' });
const runtime = buildRuntime({
markerResult: { kind: 'absent' },
verifyMarker: false,
driver: createTransactionalDriver(events),
readMarkerSpy,
});

await withTransaction(runtime, (tx) => tx.query(createPlan()).toArray());

expect(readMarkerSpy).not.toHaveBeenCalled();
expect(events).toEqual(['acquire', 'begin', 'statement', 'release']);
});

it('acquires no connection when the marker read fails', async () => {
const events: string[] = [];
const readMarkerSpy = vi.fn().mockRejectedValue(new Error('marker read failed'));
const runtime = buildRuntime({
markerResult: { kind: 'absent' },
driver: createTransactionalDriver(events),
readMarkerSpy,
});

await expect(
withTransaction(runtime, (tx) => tx.query(createPlan()).toArray()),
).rejects.toThrow('marker read failed');

expect(events).toEqual([]);
});
});
33 changes: 29 additions & 4 deletions packages/3-extensions/supabase/test/supabase-runtime.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -216,7 +216,7 @@ function createRecordingDriver(
return driver;
}

function createStubAdapter() {
function createStubAdapter(readMarker?: () => Promise<{ kind: 'absent' }>) {
const codec: Codec<string> = {
id: 'pg/int4@1',
targetTypes: ['int4'],
Expand All @@ -233,7 +233,7 @@ function createStubAdapter() {
id: 'test-profile',
target: 'postgres',
capabilities: {},
readMarker: async () => ({ kind: 'absent' as const }),
readMarker: readMarker ?? (async () => ({ kind: 'absent' as const })),
},
lower(ast: Parameters<SqlRuntimeAdapterInstance<'postgres'>['lower']>[0]) {
const params = [...new Set(ast.collectParamRefs())].map((ref) =>
Expand Down Expand Up @@ -284,8 +284,9 @@ function createTestTargetDescriptor(): SqlRuntimeTargetDescriptor<'postgres'> {
function createTestSetup(options?: {
middleware?: readonly SqlMiddleware[];
affectedRows?: number;
readMarker?: () => Promise<{ kind: 'absent' }>;
}) {
const adapter = createStubAdapter();
const adapter = createStubAdapter(options?.readMarker);
const driver = createRecordingDriver(undefined, options?.affectedRows);
const targetDescriptor = createTestTargetDescriptor();
const adapterDescriptor = createTestAdapterDescriptor(adapter);
Expand Down Expand Up @@ -314,7 +315,7 @@ function createTestSetup(options?: {
context,
adapter: stackInstance.adapter,
driver: driver as unknown as SqlDriver,
verifyMarker: false,
verifyMarker: options?.readMarker === undefined ? false : 'onFirstUse',
middleware: options?.middleware ?? [],
};

Expand Down Expand Up @@ -516,6 +517,30 @@ describe('SupabaseRuntimeImpl', () => {
});
});

describe('openRoleSession — marker verification', () => {
it('reads the marker before the role session transaction begins', async () => {
const events: string[] = [];
const { runtime, driver } = createTestSetup({
readMarker: async () => {
events.push('marker');
return { kind: 'absent' };
},
});
driver.connection.beginTransactionSpy.mockImplementation(async () => {
events.push('begin');
return driver.connection.transaction;
});
const session = await runtime.openRoleSession({ role: 'authenticated' });

const tx = await session.transaction();
await tx.query(stubPlan()).toArray();
await tx.commit();
await session.release();

expect(events).toEqual(['marker', 'begin']);
});
});

describe('openRoleSession — release', () => {
it('release() sends RESET ALL then releases the connection to the pool', async () => {
const { runtime, driver } = createTestSetup();
Expand Down
6 changes: 0 additions & 6 deletions test/e2e/framework/test/transaction-orm.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -38,12 +38,6 @@ async function withPostgresClient(
try {
runtime = await db.connect();

// Warm up the runtime so that contract verification (which acquires its
// own connection) runs before the first transaction. PGlite only allows
// one concurrent connection, so verification inside a transaction would
// deadlock.
await db.orm.public.User.first();

await callback(db);
} finally {
await runtime?.close();
Expand Down
17 changes: 17 additions & 0 deletions test/e2e/framework/test/transaction.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,23 @@ describe('transaction E2E', { timeout: 30000 }, () => {
});
});

it('runs no statement between BEGIN and the first user statement on a fresh runtime', async () => {
await withTestRuntime<Contract>(contractJsonPath, async ({ raw, runtime }) => {
const level = await withTransaction(runtime, async (tx) => {
await tx.execute(
raw.sql`SET TRANSACTION ISOLATION LEVEL SERIALIZABLE`.affectedCount().build(),
);
return tx.query(
raw.sql`SHOW transaction_isolation`
.returnsRow({ transaction_isolation: 'pg/text@1' })
.build(),
);
});

expect(level).toEqual([{ transaction_isolation: 'serializable' }]);
});
});

it('rolls back all writes on error', async () => {
await withTestRuntime<Contract>(contractJsonPath, async ({ db, runtime, client }) => {
await expect(
Expand Down
5 changes: 4 additions & 1 deletion test/e2e/framework/test/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -14,7 +14,7 @@ import type { SqlStorage } from '@prisma/orm-postgres/family-contract/types';
import type { ExecutionContext, Runtime } from '@prisma/orm-postgres/family-runtime';
import { materialiseMigrationPackage } from '@prisma/orm-postgres/migration-tools/io';
import { emitContractSpaceArtifacts } from '@prisma/orm-postgres/migration-tools/spaces';
import postgres from '@prisma/orm-postgres/runtime';
import postgres, { type PostgresClient } from '@prisma/orm-postgres/runtime';
import { PostgresContractSerializer } from '@prisma/orm-postgres/target/runtime';
import { withClient, withDevDatabase } from '@repo/test-utils';
import type { Client } from 'pg';
Expand Down Expand Up @@ -213,6 +213,8 @@ export interface TestRuntimeContext<TContract extends Contract<SqlStorage>> {
readonly runtime: Runtime;
/** The sql-builder proxy for building and executing queries */
readonly db: Db<TContract>;
/** The raw SQL lane for statements the builder does not express */
readonly raw: PostgresClient<TContract>['raw'];
/** The raw pg client for direct SQL queries */
readonly client: Client;
/** The DDL SQL generated for the contract */
Expand Down Expand Up @@ -268,6 +270,7 @@ export async function withTestRuntime<TContract extends Contract<SqlStorage>>(
context: postgresClient.context,
runtime,
db: postgresClient.sql,
raw: postgresClient.raw,
client,
sql,
});
Expand Down
Loading