Honor `transaction` — atomic groups with immediate LIFO rollback

Actions on the same event sharing a `transaction` tag form an atomic
group. When any member completes with `ERROR` or `FATAL`, the engine
immediately compensates that group's already-succeeded members in
LIFO order of completion, marks the run aborted, and emits
`RUN_ABORTED`. Transaction semantics override the failing member's
own `onError` — `abort-run` is implicit.

Compensators run at most once per action: tx-driven rollback marks
its entries as compensated, and the standard pre-`FINALLY` rollback
skips them. Non-tx compensable actions still compensate at the
global phase. `OrcaActionRun.transactionId` mirrors the tag for
trace navigation. `validate()` warns
`transaction-without-compensate` when no member of a transaction
declares a compensator (rollback would be a no-op).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
master
dev 5 months ago
parent a518dcac6b
commit 6f9b558238

@ -1347,7 +1347,6 @@ Marcado como `@v1+` en el código fuente:
| Pieza | Qué falta para v1 |
|---|---|
| `commit()` / `replace` cola | Políticas de cola por evento (drop-prev / replace / parallel). |
| `transaction` (atómico) | Grupos atómicos cuyo fallo lanza compensaciones en orden inverso. |
| Tokens con payload | Hoy son `string`; con payload abren coordinación más rica intra-run. |
| `fan-in` | Acción que dispara cuando _varios_ tokens están presentes. |
| Bus interception (Opción B) | Atribuir `emittedByAction` perfecto cuando un módulo hace `bus.publish` durante un run. |
@ -1405,6 +1404,19 @@ Marcado como `@v1+` en el código fuente:
`validate()` retiene `provides` hasta el final de la wave: dos
paralelas hermanas con `after`/`provides` cruzados se reportan
como `unsatisfiable-after`.
- ✅ **`transaction`** — grupos atómicos por evento (cualquier
stage). Si un miembro del tx termina en `ERROR` o `FATAL`, el
motor compensa **inmediatamente** los miembros del mismo tx que
ya tuvieron éxito en orden LIFO de finalización, marca al run
como abortado, y deja al flujo estándar pre-FINALLY que compense
el resto. Las semánticas tx **anulan** el `onError` del miembro
que falla (el tx siempre aborta el run; no hace falta declarar
`abort-run`). Los compensadores se invocan **a lo más una vez**:
el tx marca a sus miembros como `compensated` y la pasada global
los salta. Cada `OrcaActionRun` carga `transactionId?` para
navegación de trace. `validate()` emite warning
`transaction-without-compensate` cuando todos los miembros de un
tx no declaran `compensate` (rollback no-op).
- ✅ **`compensate`** — cada acción puede declarar una función
compensatoria. Cuando el run aborta (FATAL / ABORT_RUN onError /
trace-aborted), el motor recorre la pila LIFO de acciones que ya

@ -136,6 +136,13 @@ export const ORCA_VALIDATE_ORPHAN_ABORT_ON = 'orphan-abort-on' as const;
* event. Two actions block on each other; both would always skip.
*/
export const ORCA_VALIDATE_DEPENDENCY_CYCLE = 'dependency-cycle' as const;
/**
* A transaction id is shared by N actions but none of them declares
* `compensate`. Transaction failure would still abort the run, but
* the rollback would have nothing to undo — likely unintentional.
*/
export const ORCA_VALIDATE_TRANSACTION_WITHOUT_COMPENSATE =
'transaction-without-compensate' as const;
export const ORCA_VALIDATE_SEVERITY_ERROR = 'error' as const;
export const ORCA_VALIDATE_SEVERITY_WARN = 'warn' as const;

@ -75,6 +75,7 @@ import {
ORCA_VALIDATE_ORPHAN_UNLESS,
ORCA_VALIDATE_SEVERITY_ERROR,
ORCA_VALIDATE_SEVERITY_WARN,
ORCA_VALIDATE_TRANSACTION_WITHOUT_COMPENSATE,
ORCA_VALIDATE_UNSATISFIABLE_AFTER
} from './consts.ts';
import {
@ -408,8 +409,16 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
const tokens = new Set<OrcaToken>();
const actionRuns: OrcaActionRun[] = [];
// Stack of actions that completed success AND declared a compensate.
// Pushed in completion order so reverse iteration is LIFO.
const compensable: Array<{ action: OrcaAction; payload: unknown }> = [];
// Pushed in completion order so reverse iteration is LIFO. The
// `compensated` flag prevents an action from being compensated
// twice when its transaction failed mid-run and the global pre-
// FINALLY rollback later iterates the same list.
const compensable: Array<{
action: OrcaAction;
payload: unknown;
transactionId?: string;
compensated: boolean;
}> = [];
const compensations: OrcaActionRun[] = [];
let compensationsRun = false;
const controller = new AbortController();
@ -441,14 +450,17 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
) {
compensationsRun = true;
for (let i = compensable.length - 1; i >= 0; i--) {
const entry = compensable[i];
if (entry.compensated) continue;
const compResult = await runCompensation(
compensable[i].action,
compensable[i].payload,
entry.action,
entry.payload,
envelope,
runId,
tokens
);
compensations.push(compResult);
entry.compensated = true;
}
}
@ -484,13 +496,66 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
item.action.compensate !== undefined &&
stage !== ORCA_STAGE_FINALLY
) {
compensable.push({ action: item.action, payload: envelope.payload });
compensable.push({
action: item.action,
payload: envelope.payload,
transactionId: item.action.transaction,
compensated: false
});
}
for (const token of item.actionRun.emitted) tokens.add(token);
}
// Decide whether the run should abort. Fatal beats error;
// within a wave, registration order resolves ties.
// Detect transaction failure first: any tx member that
// errored or went fatal triggers an immediate compensation
// of its tx peers (LIFO over compensable entries belonging
// to that tx) and aborts the run. Tx semantics override
// the member's own `onError`.
const failedTxIds = new Set<string>();
let txCause: unknown | undefined;
for (const item of waveItems) {
const txId = item.action.transaction;
if (!txId) continue;
if (
item.actionRun.status === ORCA_ACTION_STATUS_ERROR ||
item.actionRun.status === ORCA_ACTION_STATUS_FATAL
) {
if (txCause === undefined) txCause = item.actionRun.error;
failedTxIds.add(txId);
}
}
if (failedTxIds.size > 0) {
for (const txId of failedTxIds) {
for (let i = compensable.length - 1; i >= 0; i--) {
const entry = compensable[i];
if (entry.compensated) continue;
if (entry.transactionId !== txId) continue;
const compResult = await runCompensation(
entry.action,
entry.payload,
envelope,
runId,
tokens
);
compensations.push(compResult);
entry.compensated = true;
}
}
aborted = true;
controller.abort();
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_ABORTED, {
runId,
event: envelope.event,
eventId: envelope.meta.eventId,
traceId: envelope.meta.traceId,
depth: envelope.meta.depth,
cause: txCause
});
break;
}
// Standard abort path. Fatal beats error; within a wave,
// registration order resolves ties.
let abortCause: unknown | undefined;
let abortTriggered = false;
for (const item of waveItems) {
@ -754,7 +819,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
startedAt,
endedAt,
durationMs: endedAt - startedAt,
emitted: []
emitted: [],
transactionId: action.transaction
};
} catch (thrown) {
const endedAt = timers.clock.now();
@ -775,7 +841,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
endedAt,
durationMs: endedAt - startedAt,
emitted: [],
error: thrown
error: thrown,
transactionId: action.transaction
};
}
}
@ -857,7 +924,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
endedAt: startedAt,
durationMs: 0,
emitted: [],
reason: ORCA_GATE_REASON_RUN_ABORTED
reason: ORCA_GATE_REASON_RUN_ABORTED,
transactionId: action.transaction
};
}
@ -946,7 +1014,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
status === ORCA_ACTION_STATUS_INTERRUPTED ||
status === ORCA_ACTION_STATUS_SKIPPED
? (result as { reason?: string }).reason
: undefined
: undefined,
transactionId: action.transaction
};
} catch (thrown) {
// Uncaught exceptions become OrcaError, regardless of intent.
@ -968,7 +1037,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
endedAt,
durationMs: endedAt - startedAt,
emitted: [],
error: thrown
error: thrown,
transactionId: action.transaction
};
}
}
@ -1134,7 +1204,8 @@ function interruptedActionRun(
endedAt: now,
durationMs: 0,
emitted: [],
reason
reason,
transactionId: action.transaction
};
}
@ -1165,7 +1236,8 @@ function evaluateGates(
endedAt: now,
durationMs: 0,
emitted: [],
reason: `${ORCA_GATE_REASON_UNLESS_TRIGGERED}:${matched}`
reason: `${ORCA_GATE_REASON_UNLESS_TRIGGERED}:${matched}`,
transactionId: action.transaction
};
}
}
@ -1180,7 +1252,8 @@ function evaluateGates(
endedAt: now,
durationMs: 0,
emitted: [],
reason: `${ORCA_GATE_REASON_ABORT_ON_TRIGGERED}:${matched}`
reason: `${ORCA_GATE_REASON_ABORT_ON_TRIGGERED}:${matched}`,
transactionId: action.transaction
};
}
}
@ -1195,7 +1268,8 @@ function evaluateGates(
endedAt: now,
durationMs: 0,
emitted: [],
reason: `${ORCA_GATE_REASON_AFTER_NOT_MET}:${missing.join(',')}`
reason: `${ORCA_GATE_REASON_AFTER_NOT_MET}:${missing.join(',')}`,
transactionId: action.transaction
};
}
}
@ -1251,6 +1325,7 @@ function validateConfiguration(
const sorted = sortActionsCanonical(registered);
validateGatesForEvent(event, sorted, issues);
validateCyclesForEvent(event, sorted, issues);
validateTransactionsForEvent(event, sorted, issues);
}
const ok = issues.every((i) => i.severity !== ORCA_VALIDATE_SEVERITY_ERROR);
return { ok, issues };
@ -1425,6 +1500,37 @@ function validateCyclesForEvent(
for (const id of adjacency.keys()) dfs(id);
}
/**
* Group actions by their `transaction` tag and report any group whose
* members all lack a `compensate` function — the rollback would have
* nothing to undo, so the transaction is effectively decorative.
* Reported once per offending tx id, not per member.
*/
function validateTransactionsForEvent(
event: string,
sorted: readonly RegisteredActionEntry[],
issues: OrcaValidateIssue[]
): void {
const byTx = new Map<string, OrcaAction[]>();
for (const { action } of sorted) {
if (!action.transaction) continue;
const list = byTx.get(action.transaction) ?? [];
list.push(action);
byTx.set(action.transaction, list);
}
for (const [txId, members] of byTx) {
const anyCompensable = members.some((m) => m.compensate !== undefined);
if (anyCompensable) continue;
issues.push({
kind: ORCA_VALIDATE_TRANSACTION_WITHOUT_COMPENSATE,
severity: ORCA_VALIDATE_SEVERITY_WARN,
event,
transactionId: txId,
message: `transaction "${txId}" on event "${event}" has no member with \`compensate\` — rollback would be a no-op`
});
}
}
/** Rotate the cycle so the smallest id comes first; used to dedupe. */
function canonicalCycleFingerprint(cycle: readonly OrcaActionId[]): string {
if (cycle.length === 0) return '';

@ -61,6 +61,7 @@ export {
ORCA_VALIDATE_ORPHAN_UNLESS,
ORCA_VALIDATE_ORPHAN_ABORT_ON,
ORCA_VALIDATE_DEPENDENCY_CYCLE,
ORCA_VALIDATE_TRANSACTION_WITHOUT_COMPENSATE,
ORCA_VALIDATE_SEVERITY_ERROR,
ORCA_VALIDATE_SEVERITY_WARN,
ORCA_DIAGNOSTIC_EVENTS,

@ -2946,3 +2946,360 @@ describe('EngineOrca v1 — parallel waves', () => {
});
});
describe('EngineOrca v1 — transaction', () => {
let bus: FakeBus;
let timers: TimerScheduler;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('compensates earlier tx members LIFO when a tx member errors', async () => {
const orca = createEngineOrca({ bus, timers });
const order: string[] = [];
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-1',
action: () => {
order.push('a');
return orcaSuccess();
},
compensate: () => {
order.push('comp:a');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => {
order.push('b');
return orcaSuccess();
},
compensate: () => {
order.push('comp:b');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'c',
stage: ORCA_STAGE_POST,
transaction: 'tx-1',
action: () => {
order.push('c');
return orcaError(new Error('fail'));
}
});
bus.publish('e', null);
await flush();
expect(order).toEqual(['a', 'b', 'c', 'comp:b', 'comp:a']);
const run = orca.recentRuns()[0];
expect(run.status).toBe(ORCA_RUN_ABORTED);
expect(run.compensations.map((c) => c.id)).toEqual(['b', 'a']);
});
it('tx semantics override the failing member onError (no need to set abort-run)', async () => {
const orca = createEngineOrca({ bus, timers });
const order: string[] = [];
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-1',
action: () => {
order.push('a');
return orcaSuccess();
},
compensate: () => {
order.push('comp:a');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'fail',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
// onError defaults to CONTINUE; transaction must override.
action: () => orcaError(new Error('fail'))
});
orca.onEvent('e', {
id: 'never',
stage: ORCA_STAGE_POST,
action: () => {
order.push('never');
return orcaSuccess();
}
});
bus.publish('e', null);
await flush();
expect(order).toEqual(['a', 'comp:a']);
expect(orca.recentRuns()[0].status).toBe(ORCA_RUN_ABORTED);
});
it('does not double-compensate when global rollback also fires for the same action', async () => {
const orca = createEngineOrca({ bus, timers });
let aCompensateCalls = 0;
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-1',
action: () => orcaSuccess(),
compensate: () => {
aCompensateCalls += 1;
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'fail',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => orcaError(new Error('fail'))
});
bus.publish('e', null);
await flush();
expect(aCompensateCalls).toBe(1);
const run = orca.recentRuns()[0];
expect(run.compensations.filter((c) => c.id === 'a')).toHaveLength(1);
});
it('aborts the run for ALL members of the failing tx, but other compensable actions still rollback at FINALLY', async () => {
const orca = createEngineOrca({ bus, timers });
const order: string[] = [];
orca.onEvent('e', {
id: 'outside',
stage: ORCA_STAGE_PRE,
action: () => {
order.push('outside');
return orcaSuccess();
},
compensate: () => {
order.push('comp:outside');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'tx-a',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => {
order.push('tx-a');
return orcaSuccess();
},
compensate: () => {
order.push('comp:tx-a');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'tx-fail',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => orcaError(new Error('fail'))
});
bus.publish('e', null);
await flush();
// tx-a compensates immediately; outside compensates at FINALLY rollback.
expect(order).toEqual(['outside', 'tx-a', 'comp:tx-a', 'comp:outside']);
});
it('tx member without compensate is silently skipped during rollback', async () => {
const orca = createEngineOrca({ bus, timers });
const order: string[] = [];
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-1',
action: () => {
order.push('a');
return orcaSuccess();
}
// no compensate
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => {
order.push('b');
return orcaSuccess();
},
compensate: () => {
order.push('comp:b');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'fail',
stage: ORCA_STAGE_POST,
transaction: 'tx-1',
action: () => orcaError(new Error('fail'))
});
bus.publish('e', null);
await flush();
expect(order).toEqual(['a', 'b', 'comp:b']);
expect(orca.recentRuns()[0].compensations.map((c) => c.id)).toEqual(['b']);
});
it('does NOT trigger when tx members all succeed', async () => {
const orca = createEngineOrca({ bus, timers });
let compCalls = 0;
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => orcaSuccess(),
compensate: () => {
compCalls += 1;
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_POST,
transaction: 'tx-1',
action: () => orcaSuccess()
});
bus.publish('e', null);
await flush();
expect(compCalls).toBe(0);
expect(orca.recentRuns()[0].status).toBe(ORCA_RUN_SUCCESS);
});
it('OrcaActionRun.transactionId carries the tag for trace navigation', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => orcaSuccess()
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_POST,
action: () => orcaSuccess()
});
bus.publish('e', null);
await flush();
const run = orca.recentRuns()[0];
expect(run.actions.find((r) => r.id === 'a')?.transactionId).toBe('tx-1');
expect(run.actions.find((r) => r.id === 'b')?.transactionId).toBeUndefined();
});
it('isolates failures: fail in tx-1 does NOT compensate tx-2 members', async () => {
const orca = createEngineOrca({ bus, timers });
const order: string[] = [];
orca.onEvent('e', {
id: 'tx2-a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-2',
action: () => {
order.push('tx2-a');
return orcaSuccess();
},
compensate: () => {
order.push('comp:tx2-a');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'tx1-a',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => {
order.push('tx1-a');
return orcaSuccess();
},
compensate: () => {
order.push('comp:tx1-a');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'tx1-fail',
stage: ORCA_STAGE_POST,
transaction: 'tx-1',
action: () => {
order.push('tx1-fail');
return orcaError(new Error('fail'));
}
});
bus.publish('e', null);
await flush();
// tx-1 rolls back immediately on failure (comp:tx1-a). tx-2 has
// no failure of its own but the run aborted, so the global
// rollback at FINALLY catches its tx-2-a too.
expect(order).toEqual(['tx2-a', 'tx1-a', 'tx1-fail', 'comp:tx1-a', 'comp:tx2-a']);
});
it('validate() warns when a transaction has no compensable member', () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-1',
action: () => orcaSuccess()
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => orcaSuccess()
});
const result = orca.validate();
expect(result.ok).toBe(true); // warning, not error
const issue = result.issues.find((i) => i.transactionId === 'tx-1');
expect(issue).toBeDefined();
expect(issue?.kind).toBe('transaction-without-compensate');
expect(issue?.severity).toBe('warn');
});
it('validate() does not warn when at least one tx member declares compensate', () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
transaction: 'tx-1',
action: () => orcaSuccess(),
compensate: () => orcaSuccess()
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_MAIN,
transaction: 'tx-1',
action: () => orcaSuccess()
});
const result = orca.validate();
const issue = result.issues.find(
(i) => i.kind === 'transaction-without-compensate'
);
expect(issue).toBeUndefined();
});
});

@ -46,6 +46,7 @@ import type {
ORCA_VALIDATE_ORPHAN_UNLESS,
ORCA_VALIDATE_SEVERITY_ERROR,
ORCA_VALIDATE_SEVERITY_WARN,
ORCA_VALIDATE_TRANSACTION_WITHOUT_COMPENSATE,
ORCA_VALIDATE_UNSATISFIABLE_AFTER
} from './consts.ts';
@ -339,6 +340,30 @@ export interface OrcaAction<TPayload = unknown, TValue = unknown> {
*/
readonly actionTimeoutMs?: number;
/**
* Atomic group tag. Actions on the same event sharing this id form
* an atomic transaction: if any member completes with `ERROR` or
* `FATAL`, the engine immediately compensates the members of that
* transaction that already succeeded — in LIFO order of completion —
* and aborts the run. The member's own `onError` is ignored when it
* is part of a transaction; transaction semantics always abort.
*
* Members may span any stages of the same event (transactions are
* not stage-bound). Compensators run at most once per action: a
* member compensated by transaction rollback is skipped during the
* standard pre-`FINALLY` rollback that other (non-transaction)
* compensable actions use. Transaction compensations are appended
* to `OrcaRunResult.compensations[]` in the same form as the
* standard ones.
*
* Tagging a member with `transaction` does not require it to
* declare `compensate` — a member without a compensator simply has
* nothing to undo on rollback. `validate()` reports a warning when
* no member of a transaction declares `compensate`, since the
* transaction would have no rollback effect at all.
*/
readonly transaction?: string;
/**
* Opt-in concurrency within a stage. Actions registered consecutively
* with `parallel: true` and the same stage form a wave that runs in
@ -404,6 +429,11 @@ export interface OrcaActionRun {
* `orcaInterrupted(reason)`.
*/
readonly reason?: OrcaReentryReason | string;
/**
* Mirrors `OrcaAction.transaction` for trace navigation. Present
* iff the action declared a transaction tag at registration.
*/
readonly transactionId?: string;
}
export interface OrcaRunResult {
@ -446,7 +476,8 @@ export type OrcaValidateIssueKind =
| typeof ORCA_VALIDATE_UNSATISFIABLE_AFTER
| typeof ORCA_VALIDATE_ORPHAN_UNLESS
| typeof ORCA_VALIDATE_ORPHAN_ABORT_ON
| typeof ORCA_VALIDATE_DEPENDENCY_CYCLE;
| typeof ORCA_VALIDATE_DEPENDENCY_CYCLE
| typeof ORCA_VALIDATE_TRANSACTION_WITHOUT_COMPENSATE;
export type OrcaValidateSeverity =
| typeof ORCA_VALIDATE_SEVERITY_ERROR
@ -466,6 +497,8 @@ export interface OrcaValidateIssue {
readonly token?: OrcaToken;
/** For dependency cycles: the action ids around the loop, in order. */
readonly cycle?: readonly OrcaActionId[];
/** For transaction issues: the offending transaction id. */
readonly transactionId?: string;
readonly message: string;
}

Loading…
Cancel
Save

Powered by TurnKey Linux.