diff --git a/cloud/apps/relay/src/assignment-store.ts b/cloud/apps/relay/src/assignment-store.ts index 14c3e50a0dc..fac9c514916 100644 --- a/cloud/apps/relay/src/assignment-store.ts +++ b/cloud/apps/relay/src/assignment-store.ts @@ -7426,6 +7426,18 @@ export class RelayAssignmentStore { [sourceCellId, targetCellId], 'pool-default' ) + // The counter is written as an absolute, so its units are read after the + // cell lock: a lease committed since the read above already bumped it. + const cellUnits = new Map( + ( + await transaction.query( + `SELECT cell_id, SUM(request_units) AS units + FROM relay_assignment_activity_leases + WHERE cell_id IN (?, ?) GROUP BY cell_id`, + [sourceCellId, targetCellId] + ) + ).map((row) => [text(row, 'cell_id'), integer(row, 'units')]) + ) const assignmentKeys = new Set( assignments.map((row) => assignmentKey(text(row, 'user_id'), text(row, 'relay_host_id')) @@ -7435,7 +7447,6 @@ export class RelayAssignmentStore { string, { counts: Record; leaseExpiresAt: number } >() - const cellUnits = new Map() for (const lease of leases) { const key = assignmentKey(text(lease, 'user_id'), text(lease, 'relay_host_id')) @@ -7451,7 +7462,6 @@ export class RelayAssignmentStore { current.counts[kind]++ current.leaseExpiresAt = Math.max(current.leaseExpiresAt, integer(lease, 'expires_at')) assignmentCounts.set(key, current) - cellUnits.set(cellId, (cellUnits.get(cellId) ?? 0) + integer(lease, 'request_units')) } for (const row of assignments) { diff --git a/cloud/apps/relay/src/postgres-transaction-recovery.test.ts b/cloud/apps/relay/src/postgres-transaction-recovery.test.ts index ae9d52a7c86..a6a5ecbfc79 100644 --- a/cloud/apps/relay/src/postgres-transaction-recovery.test.ts +++ b/cloud/apps/relay/src/postgres-transaction-recovery.test.ts @@ -794,6 +794,57 @@ describePostgres('PostgreSQL transaction recovery', () => { ]) }, 15_000) + it('keeps a unit placed between reconciliation reading leases and locking cells', async () => { + const now = 1_350_000_000_000 + const cells = [ + { id: 'reconcile-cell-a', url: 'https://reconcile-a.example.com', capacityRequests: 100 }, + { id: 'reconcile-cell-b', url: 'https://reconcile-b.example.com', capacityRequests: 100 } + ] + const seedStore = new RelayAssignmentStore(database, () => now) + await seedStore.reconcileCells(cells) + + const leasesRead = signal() + const continueReconciliation = signal() + let gateReconciliation = true + const reconcileDatabase = new TransactionProbeDatabase(database, async (phase, sql) => { + if ( + gateReconciliation && + phase === 'after' && + sql.includes('SELECT lease.* FROM relay_assignment_activity_leases lease') + ) { + gateReconciliation = false + leasesRead.resolve() + await continueReconciliation.promise + } + }) + const reconciliation = new RelayAssignmentStore( + reconcileDatabase, + () => now + ).cellEvacuationStatus('reconcile-cell-a', 'reconcile-cell-b', true) + await leasesRead.promise + // A brand-new host holds no row reconciliation locked, so it commits in the gap. + const grant = await seedStore.assign({ + userId: 'reconcile-user-gap', + relayHostId: 'reconcilehostgap' + }) + continueReconciliation.resolve() + await reconciliation + + expect(['reconcile-cell-a', 'reconcile-cell-b']).toContain(grant.cellId) + const reservations = await database.query( + `SELECT cell.cell_id, cell.reserved_requests, + COALESCE(SUM(lease.request_units), 0) AS lease_units + FROM relay_cells cell + LEFT JOIN relay_assignment_activity_leases lease ON lease.cell_id = cell.cell_id + WHERE cell.cell_id = ? + GROUP BY cell.cell_id, cell.reserved_requests`, + [grant.cellId] + ) + expect(reservations).toEqual([ + { cell_id: grant.cellId, reserved_requests: '1', lease_units: '1' } + ]) + }, 15_000) + it('fences a source activity queued behind evacuation completion', async () => { const now = 1_400_000_000_000 const identity = {