A production incident: a native-provider op ground 38-40s inside a transaction, blew the apply-phase budget, rolled back, and a downstream pipeline hot-retried the identical operation into a 6-minute CPU storm. Brainy itself never auto-retried the timeout; the gap was that TransactionTimeoutError only said "retryable" in prose, with nothing machine-readable for a caller to branch on. - TransactionTimeoutError gains two typed, always-true fields: retryable (a later attempt may succeed once the slowness resolves or the budget is raised) and hotRetryUnsafe (an immediate identical retry re-pays the full cost that just timed out and can cascade into a CPU storm -- callers must latch and back off, never loop). context's existing telemetry fields (timeoutMs, operationIndex, elapsedMs, totalOperations, operationName) are now documented as the caller's backoff inputs. - Updated the "retryable" doc-prose sites (transact()'s timeoutMs option, transactionBudgetFloorMs, Transaction.execute()'s contract) to point at the new fields instead of bare prose. - Regression pin (tests/unit/transaction/timeout-never-internally-retried.test.ts): an execution counter proves the engine never re-drives a timed-out operation, through both the single-op engine TransactionManager/Transaction drives for every single-record write, and add()'s upsert-race retry loop (which must exit on the first TransactionTimeoutError, never treat it like the lost-insert-race signal it retries on). - Removed TransactionManager.executeTransactionWithResult -- zero callers anywhere in the codebase.
321 lines
11 KiB
TypeScript
321 lines
11 KiB
TypeScript
/**
|
||
* Transaction - Atomic Unit of Work
|
||
*
|
||
* Executes operations atomically: all succeed or all rollback.
|
||
* Prevents partial failures that leave system in inconsistent state.
|
||
*
|
||
* Usage:
|
||
* ```typescript
|
||
* const tx = new Transaction()
|
||
* tx.addOperation(operation1)
|
||
* tx.addOperation(operation2)
|
||
* await tx.execute() // Both succeed or both rollback
|
||
* ```
|
||
*/
|
||
|
||
import {
|
||
Operation,
|
||
RollbackAction,
|
||
TransactionState,
|
||
TransactionContext,
|
||
TransactionOptions
|
||
} from './types.js'
|
||
import {
|
||
InvalidTransactionStateError,
|
||
TransactionExecutionError,
|
||
TransactionRollbackError,
|
||
TransactionTimeoutError
|
||
} from './errors.js'
|
||
import { prodLog } from '../utils/logger.js'
|
||
|
||
/**
|
||
* Default transaction options
|
||
*/
|
||
const DEFAULT_OPTIONS: Required<TransactionOptions> = {
|
||
timeout: 30000, // 30 seconds
|
||
logging: false,
|
||
maxRollbackRetries: 3
|
||
}
|
||
|
||
/**
|
||
* The floor term of {@link transactTimeoutBudget}'s scaling formula when the
|
||
* caller supplies neither an explicit override nor `config.transactionBudgetFloorMs`.
|
||
*/
|
||
const DEFAULT_BUDGET_FLOOR_MS = 30_000
|
||
|
||
/**
|
||
* The apply-phase budget for a batch of `opCount` operations.
|
||
*
|
||
* An explicit `override` wins untouched (a caller-specified deadline for this
|
||
* one batch). Otherwise the budget SCALES with the batch:
|
||
* `max(floorMs, opCount × 2 000)`, where `floorMs` defaults to `30 000` and
|
||
* is overridable via `BrainyConfig.transactionBudgetFloorMs` for the whole
|
||
* brain. The per-op term is calibrated from field data — bulk imports on
|
||
* network-attached disks measure ~2 s per operation (each op pays canonical
|
||
* writes + fsync + index maintenance) — so a flat 30 s budget silently capped
|
||
* honest work at ~15 operations while looking generous for small batches. A
|
||
* 16-op batch, for example, scales to `max(30 000, 16 × 2 000) = 32 000`.
|
||
* Scaling keeps small transacts fast-failing and gives bulk ones a budget
|
||
* proportional to the work they actually asked for.
|
||
*
|
||
* This is a **mid-batch** budget only: it governs whether the transaction's
|
||
* NEXT operation may start (see {@link Transaction.execute}), never whether
|
||
* already-completed work is rolled back after the fact. A trip mid-batch
|
||
* still rolls back every operation applied so far, atomically, and throws a
|
||
* fully-labeled `TransactionTimeoutError` — retryable-with-latch, never
|
||
* hot-retry (see its `retryable` and `hotRetryUnsafe` fields) — that
|
||
* zero-loss guarantee doesn't change; only the point at which the clock
|
||
* stops mattering does (at the last operation, not one check later).
|
||
*
|
||
* @param opCount - Number of operations in the batch.
|
||
* @param override - A full override for this call; wins over everything else.
|
||
* @param floorMs - The scaling floor for this call (e.g. `config.transactionBudgetFloorMs`).
|
||
* Ignored when `override` is set. Defaults to `30 000` when omitted.
|
||
*/
|
||
export function transactTimeoutBudget(
|
||
opCount: number,
|
||
override?: number,
|
||
floorMs?: number
|
||
): number {
|
||
if (override !== undefined) return override
|
||
return Math.max(floorMs ?? DEFAULT_BUDGET_FLOOR_MS, opCount * 2_000)
|
||
}
|
||
|
||
/**
|
||
* Transaction class
|
||
*/
|
||
export class Transaction implements TransactionContext {
|
||
private operations: Operation[] = []
|
||
private rollbackActions: RollbackAction[] = []
|
||
private state: TransactionState = 'pending'
|
||
private readonly options: Required<TransactionOptions>
|
||
private startTime?: number
|
||
private endTime?: number
|
||
|
||
constructor(options: TransactionOptions = {}) {
|
||
this.options = { ...DEFAULT_OPTIONS, ...options }
|
||
}
|
||
|
||
/**
|
||
* Add an operation to the transaction
|
||
*/
|
||
addOperation(operation: Operation): void {
|
||
if (this.state !== 'pending') {
|
||
throw new InvalidTransactionStateError(
|
||
this.state,
|
||
'add operation'
|
||
)
|
||
}
|
||
this.operations.push(operation)
|
||
}
|
||
|
||
/**
|
||
* Get current transaction state
|
||
*/
|
||
getState(): TransactionState {
|
||
return this.state
|
||
}
|
||
|
||
/**
|
||
* Get number of operations in transaction
|
||
*/
|
||
getOperationCount(): number {
|
||
return this.operations.length
|
||
}
|
||
|
||
/**
|
||
* Execute all operations atomically.
|
||
*
|
||
* Budget semantics: the configured budget gates starting the next
|
||
* operation; completed work is never rolled back for elapsed time. Elapsed
|
||
* time is checked ONLY before starting operation `i` for `i ≥ 1` — never
|
||
* before operation 0 (an operation always gets to start) and never again
|
||
* after the final operation completes. Concretely: a single-op transaction
|
||
* can never time out post-hoc — it either runs (and its result stands, no
|
||
* matter how long it took) or a `TransactionTimeoutError` is thrown before
|
||
* it starts, which cannot happen since there is no operation before it. A
|
||
* multi-op transaction whose op `i` overruns the budget stops op `i+1` from
|
||
* starting, rolls back operations `0..i` (reverse order), and throws
|
||
* `TransactionTimeoutError` — that zero-loss rollback guarantee for
|
||
* mid-batch trips is unchanged; only the post-completion check is gone.
|
||
*/
|
||
async execute(): Promise<void> {
|
||
if (this.state !== 'pending') {
|
||
throw new InvalidTransactionStateError(
|
||
this.state,
|
||
'execute'
|
||
)
|
||
}
|
||
|
||
this.state = 'executing'
|
||
this.startTime = Date.now()
|
||
|
||
if (this.options.logging) {
|
||
prodLog.info(`[Transaction] Executing ${this.operations.length} operations`)
|
||
}
|
||
|
||
try {
|
||
// Execute each operation in order. This loop is the sole rollback-guarded
|
||
// region: ANY error that escapes it — an operation failure OR a
|
||
// mid-flight timeout — is caught below and rolls back every operation
|
||
// applied so far, in reverse order. That single guarantee is the
|
||
// transaction's atomicity contract, and the generation-store commit path
|
||
// depends on it: its abort cleanup assumes a throw from here already
|
||
// restored the applied operations byte-identically (it only discards the
|
||
// uncommitted staging directory). A rollback per exit path — the previous
|
||
// design — let the timeout throw slip past rollback and strand the
|
||
// already-applied writes as torn, generation-less state in canonical
|
||
// storage.
|
||
for (let i = 0; i < this.operations.length; i++) {
|
||
// Budget gates STARTING the next operation; completed work is never
|
||
// rolled back for elapsed time. Skipped for i === 0 (operation 0
|
||
// always gets to start — nothing has run yet to have overrun) and
|
||
// never re-checked after the final operation completes (there is no
|
||
// "next operation" left to gate). A trip here throws into the catch
|
||
// below and rolls back like any other failure — it must never bypass
|
||
// rollback.
|
||
if (i > 0 && Date.now() - this.startTime > this.options.timeout) {
|
||
throw new TransactionTimeoutError(this.options.timeout, i, {
|
||
elapsedMs: Date.now() - this.startTime,
|
||
totalOperations: this.operations.length,
|
||
operationName: this.operations[i]?.name
|
||
})
|
||
}
|
||
|
||
const operation = this.operations[i]
|
||
|
||
if (this.options.logging) {
|
||
prodLog.info(`[Transaction] Executing operation ${i}: ${operation.name || 'unnamed'}`)
|
||
}
|
||
|
||
let rollbackAction: RollbackAction | undefined
|
||
try {
|
||
rollbackAction = await operation.execute()
|
||
} catch (error) {
|
||
// Normalize an operation failure to a TransactionExecutionError; the
|
||
// outer catch performs the (single) rollback and surfaces it.
|
||
throw new TransactionExecutionError(
|
||
`Operation ${i} failed: ${(error as Error).message}`,
|
||
i,
|
||
operation.name,
|
||
error as Error
|
||
)
|
||
}
|
||
|
||
// Record rollback action (if the operation provided one)
|
||
if (rollbackAction) {
|
||
this.rollbackActions.push(rollbackAction)
|
||
}
|
||
}
|
||
} catch (error) {
|
||
// Single rollback point: undo everything applied so far, in reverse
|
||
// order, then surface the original error. A rollback failure supersedes
|
||
// it (rollback() throws TransactionRollbackError wrapping this error).
|
||
await this.rollback(error as Error)
|
||
throw error
|
||
}
|
||
|
||
// All operations succeeded — commit.
|
||
this.commit()
|
||
}
|
||
|
||
/**
|
||
* Commit the transaction
|
||
*/
|
||
private commit(): void {
|
||
this.state = 'committed'
|
||
this.endTime = Date.now()
|
||
|
||
if (this.options.logging) {
|
||
const duration = this.endTime - (this.startTime || this.endTime)
|
||
prodLog.info(`[Transaction] Committed successfully in ${duration}ms`)
|
||
}
|
||
|
||
// Clear rollback actions - no longer needed
|
||
this.rollbackActions = []
|
||
}
|
||
|
||
/**
|
||
* Rollback all executed operations in reverse order
|
||
*/
|
||
private async rollback(originalError: Error): Promise<void> {
|
||
this.state = 'rolling_back'
|
||
|
||
if (this.options.logging) {
|
||
prodLog.info(`[Transaction] Rolling back ${this.rollbackActions.length} operations`)
|
||
}
|
||
|
||
const rollbackErrors: Error[] = []
|
||
|
||
// Execute rollback actions in REVERSE order
|
||
for (let i = this.rollbackActions.length - 1; i >= 0; i--) {
|
||
const action = this.rollbackActions[i]
|
||
|
||
// Retry rollback with exponential backoff
|
||
let attempts = 0
|
||
let success = false
|
||
|
||
while (attempts < this.options.maxRollbackRetries && !success) {
|
||
try {
|
||
await action()
|
||
success = true
|
||
|
||
if (this.options.logging) {
|
||
prodLog.info(`[Transaction] Rolled back operation ${i}`)
|
||
}
|
||
|
||
} catch (error) {
|
||
attempts++
|
||
|
||
if (attempts >= this.options.maxRollbackRetries) {
|
||
// Max retries exceeded - log error and continue
|
||
const rollbackError = error as Error
|
||
rollbackErrors.push(rollbackError)
|
||
|
||
prodLog.error(
|
||
`[Transaction] Rollback failed for operation ${i} after ${attempts} attempts: ${rollbackError.message}`
|
||
)
|
||
} else {
|
||
// Retry with exponential backoff
|
||
const delayMs = Math.pow(2, attempts) * 100
|
||
await new Promise(resolve => setTimeout(resolve, delayMs))
|
||
}
|
||
}
|
||
}
|
||
}
|
||
|
||
// Honest terminal state: only claim 'rolled_back' when EVERY undo applied.
|
||
// If any undo failed, the store is not cleanly rolled back — mark
|
||
// 'inconsistent' so the commit orchestration reconciles rather than trusts
|
||
// a false 'rolled_back'. (Fixes the state half of the post-commit response
|
||
// lie: the record whose undo failed is still durable.)
|
||
this.state = rollbackErrors.length > 0 ? 'inconsistent' : 'rolled_back'
|
||
this.endTime = Date.now()
|
||
|
||
if (this.options.logging) {
|
||
const duration = this.endTime - (this.startTime || this.endTime)
|
||
prodLog.info(`[Transaction] ${this.state === 'inconsistent' ? 'Rollback INCOMPLETE' : 'Rolled back'} in ${duration}ms`)
|
||
}
|
||
|
||
// If rollback encountered errors, wrap them with original error. The
|
||
// orchestration keys on `instanceof TransactionRollbackError` to know the
|
||
// undo did not fully apply and reconciliation is required.
|
||
if (rollbackErrors.length > 0) {
|
||
throw new TransactionRollbackError(
|
||
`Transaction rollback encountered ${rollbackErrors.length} errors during cleanup`,
|
||
originalError,
|
||
rollbackErrors
|
||
)
|
||
}
|
||
}
|
||
|
||
/**
|
||
* Get transaction execution time in milliseconds
|
||
*/
|
||
getExecutionTimeMs(): number | undefined {
|
||
if (this.startTime && this.endTime) {
|
||
return this.endTime - this.startTime
|
||
}
|
||
return undefined
|
||
}
|
||
}
|