diff --git a/packages/entity-database-adapter-knex-testing-utils/src/StubPostgresDatabaseAdapter.ts b/packages/entity-database-adapter-knex-testing-utils/src/StubPostgresDatabaseAdapter.ts index 70d1c3293c..f38b92c058 100644 --- a/packages/entity-database-adapter-knex-testing-utils/src/StubPostgresDatabaseAdapter.ts +++ b/packages/entity-database-adapter-knex-testing-utils/src/StubPostgresDatabaseAdapter.ts @@ -212,7 +212,7 @@ export class StubPostgresDatabaseAdapter< tableName, tableFieldSingleValueEqualityOperands, tableFieldMultiValueEqualityOperands, - { orderBy: undefined, offset: undefined, limit: undefined }, + { orderBy: undefined, offset: undefined, limit: undefined, forUpdate: undefined }, ); return results.length; } diff --git a/packages/entity-database-adapter-knex/src/AuthorizationResultBasedKnexEntityLoader.ts b/packages/entity-database-adapter-knex/src/AuthorizationResultBasedKnexEntityLoader.ts index 12a1b0c124..a38cf4a3e7 100644 --- a/packages/entity-database-adapter-knex/src/AuthorizationResultBasedKnexEntityLoader.ts +++ b/packages/entity-database-adapter-knex/src/AuthorizationResultBasedKnexEntityLoader.ts @@ -85,6 +85,12 @@ export interface EntityLoaderQuerySelectionModifiers< * Limit the number of entities returned. */ limit?: number; + + /** + * Lock the selected rows for update using `SELECT ... FOR UPDATE`. The lock is held until the + * end of the transaction, so the query context must be transactional. + */ + forUpdate?: boolean; } /** diff --git a/packages/entity-database-adapter-knex/src/BasePostgresEntityDatabaseAdapter.ts b/packages/entity-database-adapter-knex/src/BasePostgresEntityDatabaseAdapter.ts index 76645de364..c8a2f469d2 100644 --- a/packages/entity-database-adapter-knex/src/BasePostgresEntityDatabaseAdapter.ts +++ b/packages/entity-database-adapter-knex/src/BasePostgresEntityDatabaseAdapter.ts @@ -134,6 +134,11 @@ export interface PostgresQuerySelectionModifiers> = @@ -152,6 +157,7 @@ export interface TableQuerySelectionModifiers[] | undefined; offset: number | undefined; limit: number | undefined; + forUpdate: boolean | undefined; } export abstract class BasePostgresEntityDatabaseAdapter< @@ -343,6 +349,7 @@ export abstract class BasePostgresEntityDatabaseAdapter< : undefined, offset: querySelectionModifiers.offset, limit: querySelectionModifiers.limit, + forUpdate: querySelectionModifiers.forUpdate, }; } } diff --git a/packages/entity-database-adapter-knex/src/BaseSQLQueryBuilder.ts b/packages/entity-database-adapter-knex/src/BaseSQLQueryBuilder.ts index 260c666244..551a6303c6 100644 --- a/packages/entity-database-adapter-knex/src/BaseSQLQueryBuilder.ts +++ b/packages/entity-database-adapter-knex/src/BaseSQLQueryBuilder.ts @@ -22,6 +22,7 @@ export abstract class BaseSQLQueryBuilder< limit?: number; offset?: number; orderBy?: readonly EntityLoaderOrderByClause[]; + forUpdate?: boolean; }, ) {} @@ -41,6 +42,15 @@ export abstract class BaseSQLQueryBuilder< return this; } + /** + * Lock the selected rows for update using `SELECT ... FOR UPDATE`. The lock is held until the + * end of the transaction, so the query must be executed in a transactional query context. + */ + forUpdate(): this { + this.modifiers.forUpdate = true; + return this; + } + /** * Order by a field. Can be called multiple times to add multiple order bys. */ diff --git a/packages/entity-database-adapter-knex/src/PostgresEntityDatabaseAdapter.ts b/packages/entity-database-adapter-knex/src/PostgresEntityDatabaseAdapter.ts index 2b21aea46b..5843b3a556 100644 --- a/packages/entity-database-adapter-knex/src/PostgresEntityDatabaseAdapter.ts +++ b/packages/entity-database-adapter-knex/src/PostgresEntityDatabaseAdapter.ts @@ -116,7 +116,7 @@ export class PostgresEntityDatabaseAdapter< tableValue: tableTuple[index], })), [], - { limit: 1, orderBy: undefined, offset: undefined }, + { limit: 1, orderBy: undefined, offset: undefined, forUpdate: undefined }, ); return results[0] ?? null; } @@ -125,7 +125,7 @@ export class PostgresEntityDatabaseAdapter< query: Knex.QueryBuilder, querySelectionModifiers: TableQuerySelectionModifiers, ): Knex.QueryBuilder { - const { orderBy, offset, limit } = querySelectionModifiers; + const { orderBy, offset, limit, forUpdate } = querySelectionModifiers; let ret = query; @@ -161,6 +161,10 @@ export class PostgresEntityDatabaseAdapter< ret = ret.limit(limit); } + if (forUpdate) { + ret = ret.forUpdate(); + } + return ret; } diff --git a/packages/entity-database-adapter-knex/src/__integration-tests__/PostgresEntityIntegration-test.ts b/packages/entity-database-adapter-knex/src/__integration-tests__/PostgresEntityIntegration-test.ts index 84a98fc749..a6661118fc 100644 --- a/packages/entity-database-adapter-knex/src/__integration-tests__/PostgresEntityIntegration-test.ts +++ b/packages/entity-database-adapter-knex/src/__integration-tests__/PostgresEntityIntegration-test.ts @@ -211,6 +211,166 @@ describe('postgres entity integration', () => { ); }); + describe('forUpdate', () => { + // Attempts to lock the row from a separate connection without waiting. Postgres raises + // lock_not_available (55P03) when another transaction already holds a FOR UPDATE lock on the row. + const tryLockRowFromOtherConnectionAsync = async (id: string): Promise => { + try { + await knexInstance.raw( + 'SELECT * FROM postgres_test_entities WHERE id = ? FOR UPDATE NOWAIT', + [id], + ); + return null; + } catch (e) { + return (e as any).code ?? null; + } + }; + + it('locks rows loaded with loadManyByFieldEqualityConjunctionAsync until the transaction ends', async () => { + const vc1 = new ViewerContext(createKnexIntegrationTestEntityCompanionProvider(knexInstance)); + const entity = await enforceAsyncResult( + PostgresTestEntity.creatorWithAuthorizationResults(vc1) + .setField('name', 'locked') + .createAsync(), + ); + + // not locked before the transaction + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBeNull(); + + await vc1.runInTransactionForDatabaseAdapterFlavorAsync('postgres', async (queryContext) => { + const results = await PostgresTestEntity.knexLoader( + vc1, + queryContext, + ).loadManyByFieldEqualityConjunctionAsync([{ fieldName: 'name', fieldValue: 'locked' }], { + forUpdate: true, + }); + expect(results).toHaveLength(1); + expect(results[0]!.getID()).toBe(entity.getID()); + + // locked while the transaction is open + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBe('55P03'); + }); + + // unlocked after the transaction commits + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBeNull(); + }); + + it('locks rows loaded with loadFirstByFieldEqualityConjunctionAsync', async () => { + const vc1 = new ViewerContext(createKnexIntegrationTestEntityCompanionProvider(knexInstance)); + const entity = await enforceAsyncResult( + PostgresTestEntity.creatorWithAuthorizationResults(vc1) + .setField('name', 'locked-first') + .createAsync(), + ); + + await vc1.runInTransactionForDatabaseAdapterFlavorAsync('postgres', async (queryContext) => { + const result = await PostgresTestEntity.knexLoader( + vc1, + queryContext, + ).loadFirstByFieldEqualityConjunctionAsync( + [{ fieldName: 'name', fieldValue: 'locked-first' }], + { + orderBy: [{ fieldName: 'name', order: OrderByOrdering.ASCENDING }], + forUpdate: true, + }, + ); + expect(result?.getID()).toBe(entity.getID()); + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBe('55P03'); + }); + + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBeNull(); + }); + + it('locks rows loaded with the loadManyBySQL query builder', async () => { + const vc1 = new ViewerContext(createKnexIntegrationTestEntityCompanionProvider(knexInstance)); + const entity = await enforceAsyncResult( + PostgresTestEntity.creatorWithAuthorizationResults(vc1) + .setField('name', 'locked-sql') + .createAsync(), + ); + + await vc1.runInTransactionForDatabaseAdapterFlavorAsync('postgres', async (queryContext) => { + // via the fluent builder method + const results = await PostgresTestEntity.knexLoader(vc1, queryContext) + .loadManyBySQL(sql`name = ${'locked-sql'}`) + .forUpdate() + .executeAsync(); + expect(results).toHaveLength(1); + expect(results[0]!.getID()).toBe(entity.getID()); + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBe('55P03'); + + // via the modifiers argument on the authorization-result-based loader + const authorizationResults = await PostgresTestEntity.knexLoaderWithAuthorizationResults( + vc1, + queryContext, + ) + .loadManyBySQL(sql`name = ${'locked-sql'}`, { forUpdate: true }) + .executeAsync(); + expect(authorizationResults).toHaveLength(1); + expect(authorizationResults[0]!.enforceValue().getID()).toBe(entity.getID()); + }); + + expect(await tryLockRowFromOtherConnectionAsync(entity.getID())).toBeNull(); + }); + + it('blocks a concurrent forUpdate load until the first transaction completes', async () => { + const vc1 = new ViewerContext(createKnexIntegrationTestEntityCompanionProvider(knexInstance)); + const entity = await enforceAsyncResult( + PostgresTestEntity.creatorWithAuthorizationResults(vc1) + .setField('name', 'counter') + .setField('hasADog', false) + .createAsync(), + ); + + const events: string[] = []; + const lockAndUpdateAsync = async (label: string, holdDuration: number): Promise => { + await vc1.runInTransactionForDatabaseAdapterFlavorAsync( + 'postgres', + async (queryContext) => { + const locked = await PostgresTestEntity.knexLoader(vc1, queryContext) + .loadManyBySQL(sql`id = ${entity.getID()}`) + .forUpdate() + .executeAsync(); + events.push(`${label}:locked`); + await setTimeout(holdDuration); + await PostgresTestEntity.updater(locked[0]!, queryContext) + .setField('name', locked[0]!.getField('name') + ',' + label) + .updateAsync(); + events.push(`${label}:updated`); + }, + ); + }; + + await Promise.all([ + lockAndUpdateAsync('a', 200), + setTimeout(50).then(() => lockAndUpdateAsync('b', 0)), + ]); + + // b cannot acquire the lock until a commits, so a's update always lands before b locks + expect(events).toEqual(['a:locked', 'a:updated', 'b:locked', 'b:updated']); + + const reloaded = await PostgresTestEntity.loader(vc1).loadByIDAsync(entity.getID()); + expect(reloaded.getField('name')).toBe('counter,a,b'); + }); + + it('throws when forUpdate is used outside of a transaction', async () => { + const vc1 = new ViewerContext(createKnexIntegrationTestEntityCompanionProvider(knexInstance)); + + await expect( + PostgresTestEntity.knexLoader(vc1).loadManyByFieldEqualityConjunctionAsync([], { + forUpdate: true, + }), + ).rejects.toThrow('forUpdate requires a transactional query context'); + + await expect( + PostgresTestEntity.knexLoader(vc1) + .loadManyBySQL(sql`TRUE`) + .forUpdate() + .executeAsync(), + ).rejects.toThrow('forUpdate requires a transactional query context'); + }); + }); + describe('JSON fields', () => { it('supports both types of array fields', async () => { const vc1 = new ViewerContext(createKnexIntegrationTestEntityCompanionProvider(knexInstance)); diff --git a/packages/entity-database-adapter-knex/src/__tests__/BasePostgresEntityDatabaseAdapter-test.ts b/packages/entity-database-adapter-knex/src/__tests__/BasePostgresEntityDatabaseAdapter-test.ts index ccc4effb01..479aed8832 100644 --- a/packages/entity-database-adapter-knex/src/__tests__/BasePostgresEntityDatabaseAdapter-test.ts +++ b/packages/entity-database-adapter-knex/src/__tests__/BasePostgresEntityDatabaseAdapter-test.ts @@ -6,8 +6,10 @@ import { instance, mock } from 'ts-mockito'; import type { TableFieldMultiValueEqualityCondition, TableFieldSingleValueEqualityCondition, + TableQuerySelectionModifiers, } from '../BasePostgresEntityDatabaseAdapter.ts'; import { BasePostgresEntityDatabaseAdapter } from '../BasePostgresEntityDatabaseAdapter.ts'; +import { sql } from '../SQLOperator.ts'; import type { TestFields } from './fixtures/TestEntity.ts'; import { testEntityConfiguration } from './fixtures/TestEntity.ts'; @@ -22,6 +24,12 @@ class TestEntityDatabaseAdapter extends BasePostgresEntityDatabaseAdapter< private readonly fetchEqualityConditionResults: object[]; private readonly fetchSQLFragmentResults: object[]; private readonly deleteCount: number; + public lastEqualityConditionQuerySelectionModifiers: + | TableQuerySelectionModifiers + | undefined; + public lastSQLFragmentQuerySelectionModifiers: + | TableQuerySelectionModifiers + | undefined; constructor({ fetchResults = [], @@ -76,7 +84,9 @@ class TestEntityDatabaseAdapter extends BasePostgresEntityDatabaseAdapter< _queryInterface: any, _tableName: string, _sqlFragment: any, + querySelectionModifiers: TableQuerySelectionModifiers, ): Promise { + this.lastSQLFragmentQuerySelectionModifiers = querySelectionModifiers; return this.fetchSQLFragmentResults; } @@ -85,7 +95,9 @@ class TestEntityDatabaseAdapter extends BasePostgresEntityDatabaseAdapter< _tableName: string, _tableFieldSingleValueEqualityOperands: TableFieldSingleValueEqualityCondition[], _tableFieldMultiValueEqualityOperands: TableFieldMultiValueEqualityCondition[], + querySelectionModifiers: TableQuerySelectionModifiers, ): Promise { + this.lastEqualityConditionQuerySelectionModifiers = querySelectionModifiers; return this.fetchEqualityConditionResults; } @@ -151,5 +163,43 @@ describe(BasePostgresEntityDatabaseAdapter, () => { const results = await adapter.fetchManyByFieldEqualityConjunctionAsync(queryContext, [], {}); expect(results).toEqual([{ stringField: 'hello' }]); }); + + it('converts query selection modifiers including forUpdate', async () => { + const queryContext = instance(mock(EntityQueryContext)); + const adapter = new TestEntityDatabaseAdapter({}); + await adapter.fetchManyByFieldEqualityConjunctionAsync(queryContext, [], { + limit: 2, + offset: 1, + forUpdate: true, + }); + expect(adapter.lastEqualityConditionQuerySelectionModifiers).toEqual({ + orderBy: undefined, + limit: 2, + offset: 1, + forUpdate: true, + }); + + await adapter.fetchManyByFieldEqualityConjunctionAsync(queryContext, [], {}); + expect(adapter.lastEqualityConditionQuerySelectionModifiers).toEqual({ + orderBy: undefined, + limit: undefined, + offset: undefined, + forUpdate: undefined, + }); + }); + }); + + describe('fetchManyBySQLFragmentAsync', () => { + it('converts query selection modifiers including forUpdate', async () => { + const queryContext = instance(mock(EntityQueryContext)); + const adapter = new TestEntityDatabaseAdapter({}); + await adapter.fetchManyBySQLFragmentAsync(queryContext, sql`TRUE`, { forUpdate: true }); + expect(adapter.lastSQLFragmentQuerySelectionModifiers).toEqual({ + orderBy: undefined, + limit: undefined, + offset: undefined, + forUpdate: true, + }); + }); }); }); diff --git a/packages/entity-database-adapter-knex/src/__tests__/fixtures/StubPostgresDatabaseAdapter.ts b/packages/entity-database-adapter-knex/src/__tests__/fixtures/StubPostgresDatabaseAdapter.ts index af77a4d5e7..943d8d4b96 100644 --- a/packages/entity-database-adapter-knex/src/__tests__/fixtures/StubPostgresDatabaseAdapter.ts +++ b/packages/entity-database-adapter-knex/src/__tests__/fixtures/StubPostgresDatabaseAdapter.ts @@ -213,7 +213,7 @@ export class StubPostgresDatabaseAdapter< tableName, tableFieldSingleValueEqualityOperands, tableFieldMultiValueEqualityOperands, - { orderBy: undefined, offset: undefined, limit: undefined }, + { orderBy: undefined, offset: undefined, limit: undefined, forUpdate: undefined }, ); return results.length; } diff --git a/packages/entity-database-adapter-knex/src/internal/EntityKnexDataManager.ts b/packages/entity-database-adapter-knex/src/internal/EntityKnexDataManager.ts index ee5eec03d0..19f4a72595 100644 --- a/packages/entity-database-adapter-knex/src/internal/EntityKnexDataManager.ts +++ b/packages/entity-database-adapter-knex/src/internal/EntityKnexDataManager.ts @@ -177,6 +177,7 @@ export class EntityKnexDataManager< querySelectionModifiers: PostgresQuerySelectionModifiers, ): Promise[]> { EntityKnexDataManager.validateOrderByClauses(querySelectionModifiers.orderBy); + EntityKnexDataManager.validateForUpdate(queryContext, querySelectionModifiers.forUpdate); return await timeAndLogLoadEventAsync( this.metricsAdapter, @@ -215,6 +216,7 @@ export class EntityKnexDataManager< querySelectionModifiers: PostgresQuerySelectionModifiers, ): Promise[]> { EntityKnexDataManager.validateOrderByClauses(querySelectionModifiers.orderBy); + EntityKnexDataManager.validateForUpdate(queryContext, querySelectionModifiers.forUpdate); return await timeAndLogLoadEventAsync( this.metricsAdapter, @@ -531,6 +533,24 @@ export class EntityKnexDataManager< } } + /** + * `SELECT ... FOR UPDATE` row locks are released at the end of the transaction. Outside of a + * transaction the lock is released as soon as the statement completes, which makes it useless, + * so require a transactional query context. + */ + private static validateForUpdate( + queryContext: EntityQueryContext, + forUpdate: boolean | undefined, + ): void { + if (!forUpdate) { + return; + } + assert( + queryContext.isInTransaction(), + 'forUpdate requires a transactional query context since row locks are released at the end of the transaction.', + ); + } + /** * Cursor-based pagination uses Postgres tuple comparison (e.g., (a, b) \> (x, y)) which * applies a single comparison direction to all columns. Mixed ordering directions would diff --git a/packages/entity-database-adapter-knex/src/internal/__tests__/EntityKnexDataManager-test.ts b/packages/entity-database-adapter-knex/src/internal/__tests__/EntityKnexDataManager-test.ts index 56643e109d..466111e0cc 100644 --- a/packages/entity-database-adapter-knex/src/internal/__tests__/EntityKnexDataManager-test.ts +++ b/packages/entity-database-adapter-knex/src/internal/__tests__/EntityKnexDataManager-test.ts @@ -7,6 +7,7 @@ import { anyNumber, anything, deepEqual, instance, mock, verify, when } from 'ts import { OrderByOrdering } from '../../BasePostgresEntityDatabaseAdapter.ts'; import { PaginationStrategy } from '../../PaginationStrategy.ts'; import { PostgresEntityDatabaseAdapter } from '../../PostgresEntityDatabaseAdapter.ts'; +import { sql } from '../../SQLOperator.ts'; import type { TestFields } from '../../__tests__/fixtures/TestEntity.ts'; import { TestEntity, testEntityConfiguration } from '../../__tests__/fixtures/TestEntity.ts'; import { EntityKnexDataManager } from '../EntityKnexDataManager.ts'; @@ -193,6 +194,146 @@ describe(EntityKnexDataManager, () => { }); }); + describe('forUpdate', () => { + const fieldObject = { + customIdField: '1', + testIndexedField: 'unique1', + stringField: 'hello', + intField: 1, + dateField: new Date(), + nullableField: null, + }; + + it('throws when loading by field equality conjunction outside of a transaction', async () => { + const queryContext = new StubQueryContextProvider().getQueryContext(); + const databaseAdapterMock = mock>( + PostgresEntityDatabaseAdapter, + ); + const entityDataManager = new EntityKnexDataManager( + testEntityConfiguration, + instance(databaseAdapterMock), + new NoOpEntityMetricsAdapter(), + TestEntity.name, + ); + + await expect( + entityDataManager.loadManyByFieldEqualityConjunctionAsync(queryContext, [], { + forUpdate: true, + }), + ).rejects.toThrow('forUpdate requires a transactional query context'); + verify( + databaseAdapterMock.fetchManyByFieldEqualityConjunctionAsync( + anything(), + anything(), + anything(), + ), + ).never(); + }); + + it('throws when loading by SQL fragment outside of a transaction', async () => { + const queryContext = new StubQueryContextProvider().getQueryContext(); + const databaseAdapterMock = mock>( + PostgresEntityDatabaseAdapter, + ); + const entityDataManager = new EntityKnexDataManager( + testEntityConfiguration, + instance(databaseAdapterMock), + new NoOpEntityMetricsAdapter(), + TestEntity.name, + ); + + await expect( + entityDataManager.loadManyBySQLFragmentAsync(queryContext, sql`TRUE`, { + forUpdate: true, + }), + ).rejects.toThrow('forUpdate requires a transactional query context'); + verify( + databaseAdapterMock.fetchManyBySQLFragmentAsync(anything(), anything(), anything()), + ).never(); + }); + + it('passes forUpdate through to the database adapter inside of a transaction', async () => { + const databaseAdapterMock = mock>( + PostgresEntityDatabaseAdapter, + ); + when( + databaseAdapterMock.fetchManyByFieldEqualityConjunctionAsync( + anything(), + anything(), + anything(), + ), + ).thenResolve([fieldObject]); + when( + databaseAdapterMock.fetchManyBySQLFragmentAsync(anything(), anything(), anything()), + ).thenResolve([fieldObject]); + + const entityDataManager = new EntityKnexDataManager( + testEntityConfiguration, + instance(databaseAdapterMock), + new NoOpEntityMetricsAdapter(), + TestEntity.name, + ); + + await new StubQueryContextProvider().runInTransactionAsync(async (queryContext) => { + const equalityResults = await entityDataManager.loadManyByFieldEqualityConjunctionAsync( + queryContext, + [{ fieldName: 'stringField', fieldValue: 'hello' }], + { forUpdate: true, limit: 1 }, + ); + expect(equalityResults).toHaveLength(1); + + const sqlResults = await entityDataManager.loadManyBySQLFragmentAsync( + queryContext, + sql`TRUE`, + { forUpdate: true }, + ); + expect(sqlResults).toHaveLength(1); + + verify( + databaseAdapterMock.fetchManyByFieldEqualityConjunctionAsync( + queryContext, + anything(), + deepEqual({ forUpdate: true, limit: 1 }), + ), + ).once(); + verify( + databaseAdapterMock.fetchManyBySQLFragmentAsync( + queryContext, + anything(), + deepEqual({ forUpdate: true }), + ), + ).once(); + }); + }); + + it('does not require a transaction when forUpdate is not set', async () => { + const queryContext = new StubQueryContextProvider().getQueryContext(); + const databaseAdapterMock = mock>( + PostgresEntityDatabaseAdapter, + ); + when( + databaseAdapterMock.fetchManyByFieldEqualityConjunctionAsync( + anything(), + anything(), + anything(), + ), + ).thenResolve([fieldObject]); + const entityDataManager = new EntityKnexDataManager( + testEntityConfiguration, + instance(databaseAdapterMock), + new NoOpEntityMetricsAdapter(), + TestEntity.name, + ); + + const results = await entityDataManager.loadManyByFieldEqualityConjunctionAsync( + queryContext, + [], + { forUpdate: false }, + ); + expect(results).toHaveLength(1); + }); + }); + describe('pagination', () => { describe('max page size validation', () => { it('should throw when first exceeds maxPageSize', async () => {