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
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}

/**
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -134,6 +134,11 @@ export interface PostgresQuerySelectionModifiers<TFields extends Record<string,
* Limit the number of entities returned.
*/
limit?: number;

/**
* Lock the selected rows for update using `SELECT ... FOR UPDATE`.
*/
forUpdate?: boolean;
}

export type TableOrderByClause<TFields extends Record<string, any>> =
Expand All @@ -152,6 +157,7 @@ export interface TableQuerySelectionModifiers<TFields extends Record<string, any
orderBy: TableOrderByClause<TFields>[] | undefined;
offset: number | undefined;
limit: number | undefined;
forUpdate: boolean | undefined;
}

export abstract class BasePostgresEntityDatabaseAdapter<
Expand Down Expand Up @@ -343,6 +349,7 @@ export abstract class BasePostgresEntityDatabaseAdapter<
: undefined,
offset: querySelectionModifiers.offset,
limit: querySelectionModifiers.limit,
forUpdate: querySelectionModifiers.forUpdate,
};
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,7 @@ export abstract class BaseSQLQueryBuilder<
limit?: number;
offset?: number;
orderBy?: readonly EntityLoaderOrderByClause<TFields, TSelectedFields>[];
forUpdate?: boolean;
},
) {}

Expand All @@ -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.
*/
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand All @@ -125,7 +125,7 @@ export class PostgresEntityDatabaseAdapter<
query: Knex.QueryBuilder,
querySelectionModifiers: TableQuerySelectionModifiers<TFields>,
): Knex.QueryBuilder {
const { orderBy, offset, limit } = querySelectionModifiers;
const { orderBy, offset, limit, forUpdate } = querySelectionModifiers;

let ret = query;

Expand Down Expand Up @@ -161,6 +161,10 @@ export class PostgresEntityDatabaseAdapter<
ret = ret.limit(limit);
}

if (forUpdate) {
ret = ret.forUpdate();
}

return ret;
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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<string | null> => {
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<void> => {
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));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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';

Expand All @@ -22,6 +24,12 @@ class TestEntityDatabaseAdapter extends BasePostgresEntityDatabaseAdapter<
private readonly fetchEqualityConditionResults: object[];
private readonly fetchSQLFragmentResults: object[];
private readonly deleteCount: number;
public lastEqualityConditionQuerySelectionModifiers:
| TableQuerySelectionModifiers<TestFields>
| undefined;
public lastSQLFragmentQuerySelectionModifiers:
| TableQuerySelectionModifiers<TestFields>
| undefined;

constructor({
fetchResults = [],
Expand Down Expand Up @@ -76,7 +84,9 @@ class TestEntityDatabaseAdapter extends BasePostgresEntityDatabaseAdapter<
_queryInterface: any,
_tableName: string,
_sqlFragment: any,
querySelectionModifiers: TableQuerySelectionModifiers<TestFields>,
): Promise<object[]> {
this.lastSQLFragmentQuerySelectionModifiers = querySelectionModifiers;
return this.fetchSQLFragmentResults;
}

Expand All @@ -85,7 +95,9 @@ class TestEntityDatabaseAdapter extends BasePostgresEntityDatabaseAdapter<
_tableName: string,
_tableFieldSingleValueEqualityOperands: TableFieldSingleValueEqualityCondition[],
_tableFieldMultiValueEqualityOperands: TableFieldMultiValueEqualityCondition[],
querySelectionModifiers: TableQuerySelectionModifiers<TestFields>,
): Promise<object[]> {
this.lastEqualityConditionQuerySelectionModifiers = querySelectionModifiers;
return this.fetchEqualityConditionResults;
}

Expand Down Expand Up @@ -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,
});
});
});
});
Original file line number Diff line number Diff line change
Expand Up @@ -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;
}
Expand Down
Loading
Loading