From a8af2558db081f3d8bc396d72025d37d66d44ded Mon Sep 17 00:00:00 2001 From: Kristof Siket Date: Thu, 17 Sep 2026 15:29:47 +0200 Subject: [PATCH] fix(sql-runtime): verify contract marker before BEGIN A fresh runtime whose first operation was a transaction read the contract marker after BEGIN. On a single-connection driver that read was the first statement of the transaction: SET TRANSACTION failed with SQLSTATE 25001, REPEATABLE READ and SERIALIZABLE snapshots were taken at the marker read, and a failed marker read aborted the transaction. On a pooled driver the read waited for a second client, which never arrives when the database serves one connection at a time. The runtime now awaits the single-flight marker gate before it acquires a connection, in acquireRawConnection(), and connection() acquires through it. The gate is unchanged: one read per runtime, none when verifyMarker is false. Co-Authored-By: Claude Fable 5.1 Signed-off-by: Kristof Siket --- .../4. Runtime & Middleware Framework.md | 4 +- packages/2-sql/5-runtime/src/sql-runtime.ts | 22 +++-- .../test/marker-verification.test.ts | 99 +++++++++++++++++++ .../supabase/test/supabase-runtime.test.ts | 33 ++++++- .../framework/test/transaction-orm.test.ts | 6 -- test/e2e/framework/test/transaction.test.ts | 17 ++++ test/e2e/framework/test/utils.ts | 5 +- 7 files changed, 168 insertions(+), 18 deletions(-) diff --git a/docs/architecture docs/subsystems/4. Runtime & Middleware Framework.md b/docs/architecture docs/subsystems/4. Runtime & Middleware Framework.md index 482164f5e7e1..468db24a9468 100644 --- a/docs/architecture docs/subsystems/4. Runtime & Middleware Framework.md +++ b/docs/architecture docs/subsystems/4. Runtime & Middleware Framework.md @@ -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). diff --git a/packages/2-sql/5-runtime/src/sql-runtime.ts b/packages/2-sql/5-runtime/src/sql-runtime.ts index 954b464492d6..77eae03f2263 100644 --- a/packages/2-sql/5-runtime/src/sql-runtime.ts +++ b/packages/2-sql/5-runtime/src/sql-runtime.ts @@ -363,18 +363,28 @@ export abstract class SqlRuntimeBase = 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 { + protected async acquireRawConnection(): Promise { + await this.ensureMarkerVerified(); return this.driver.acquireConnection(); } - private async setupDriverExecution(exec: SqlExecutionPlan): Promise { - this.familyAdapter.validatePlan(exec, this.contract); - this._telemetry = null; + private ensureMarkerVerified(): Promise { if (this.verifyMarkerPromise === null) { this.verifyMarkerPromise = this.verifyMarker(); } - await this.verifyMarkerPromise; + return this.verifyMarkerPromise; + } + + private async setupDriverExecution(exec: SqlExecutionPlan): Promise { + this.familyAdapter.validatePlan(exec, this.contract); + this._telemetry = null; + await this.ensureMarkerVerified(); } protected getListDecoder(): ListDecoder { @@ -763,7 +773,7 @@ export abstract class SqlRuntimeBase = Co } async connection(): Promise { - const driverConn = await this.driver.acquireConnection(); + const driverConn = await this.acquireRawConnection(); const self = this; const wrappedConnection: RuntimeConnection & diff --git a/packages/2-sql/5-runtime/test/marker-verification.test.ts b/packages/2-sql/5-runtime/test/marker-verification.test.ts index 2d56a25b0b01..495aaa3e3185 100644 --- a/packages/2-sql/5-runtime/test/marker-verification.test.ts +++ b/packages/2-sql/5-runtime/test/marker-verification.test.ts @@ -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'; @@ -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; + }), + 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([]); + }); +}); diff --git a/packages/3-extensions/supabase/test/supabase-runtime.test.ts b/packages/3-extensions/supabase/test/supabase-runtime.test.ts index 86156f621182..68c320a03369 100644 --- a/packages/3-extensions/supabase/test/supabase-runtime.test.ts +++ b/packages/3-extensions/supabase/test/supabase-runtime.test.ts @@ -216,7 +216,7 @@ function createRecordingDriver( return driver; } -function createStubAdapter() { +function createStubAdapter(readMarker?: () => Promise<{ kind: 'absent' }>) { const codec: Codec = { id: 'pg/int4@1', targetTypes: ['int4'], @@ -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['lower']>[0]) { const params = [...new Set(ast.collectParamRefs())].map((ref) => @@ -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); @@ -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 ?? [], }; @@ -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(); diff --git a/test/e2e/framework/test/transaction-orm.test.ts b/test/e2e/framework/test/transaction-orm.test.ts index eabffef9bd1d..11003009e041 100644 --- a/test/e2e/framework/test/transaction-orm.test.ts +++ b/test/e2e/framework/test/transaction-orm.test.ts @@ -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(); diff --git a/test/e2e/framework/test/transaction.test.ts b/test/e2e/framework/test/transaction.test.ts index 0b3c47aacc2e..066804a596fd 100644 --- a/test/e2e/framework/test/transaction.test.ts +++ b/test/e2e/framework/test/transaction.test.ts @@ -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(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(contractJsonPath, async ({ db, runtime, client }) => { await expect( diff --git a/test/e2e/framework/test/utils.ts b/test/e2e/framework/test/utils.ts index 066bb2f9b92d..f6fb9203eecb 100644 --- a/test/e2e/framework/test/utils.ts +++ b/test/e2e/framework/test/utils.ts @@ -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'; @@ -213,6 +213,8 @@ export interface TestRuntimeContext> { readonly runtime: Runtime; /** The sql-builder proxy for building and executing queries */ readonly db: Db; + /** The raw SQL lane for statements the builder does not express */ + readonly raw: PostgresClient['raw']; /** The raw pg client for direct SQL queries */ readonly client: Client; /** The DDL SQL generated for the contract */ @@ -268,6 +270,7 @@ export async function withTestRuntime>( context: postgresClient.context, runtime, db: postgresClient.sql, + raw: postgresClient.raw, client, sql, });