diff --git a/src/arts/orca/README.md b/src/arts/orca/README.md index 033ee1f..f19dea7 100644 --- a/src/arts/orca/README.md +++ b/src/arts/orca/README.md @@ -25,7 +25,7 @@ que el acoplamiento inter-modulo quede escondido en `buss`, `connection` o ## Estado Del Documento Este README es la referencia de diseño de orca. La **v0-kernel** está -implementada y testeada (114 archivos / 1392 tests pasan, 82 de ellos +implementada y testeada (114 archivos / 1402 tests pasan, 92 de ellos sobre orca). El kernel expone: - `createEngineOrca({ bus, timers, logger?, maxRuns?, reentry? })` @@ -56,6 +56,13 @@ sobre orca). El kernel expone: upstream produce (error), `unless` y `abortOn` huérfanos (warnings), y ciclos en la cadena `after`/`provides` (error). No lanza; devuelve `{ ok, issues }`. +- **`compensate`:** cada acción puede declarar una función + compensatoria; cuando el run aborta (FATAL / ABORT_RUN / + trace-aborted), el motor invoca las compensaciones de las acciones + que ya completaron `success` en orden LIFO, antes del stage + `FINALLY`. Best-effort: si una compensación lanza, se registra y + la siguiente continúa. Las compensaciones aparecen en + `OrcaRunResult.compensations[]`. - Diagnostics estructurados (`orca.run.*`, `orca.action.*`, `orca.event.emitted`, `orca.reentry.blocked`, `orca.trace.aborted`, `orca.action.fatal`, `orca.action.timeout`, `orca.action.blocked`) @@ -1339,7 +1346,6 @@ Marcado como `@v1+` en el código fuente: | Pieza | Qué falta para v1 | |---|---| -| `compensate` | Invocar la función compensatoria cuando el run aborta tras un éxito previo. | | `commit()` configuración | Congelar el grafo (rechazar `onEvent` posteriores) tras `validate()`. | | `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. | @@ -1378,3 +1384,15 @@ Marcado como `@v1+` en el código fuente: decide si tratar las issues como bloqueantes — el motor no congela registros ni runs en función del resultado (el `commit()` configuracional queda para v1.x). +- ✅ **`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 + completaron `success` con compensador y las invoca en orden + inverso. Cada compensación corre con un `AbortController` propio + (el del run ya está abortado) y un `ctx` cuyo `emit()` devuelve + `null` — las compensaciones no fan-out eventos, son rollback puro. + Best-effort: si una compensación lanza, se registra y la cadena + continúa con la siguiente. Las compensaciones aparecen en + `OrcaRunResult.compensations[]` separadas de `actions[]`. Los + stages `FINALLY` no son compensables (finally es la limpieza + misma) y se ejecutan **después** de las compensaciones. diff --git a/src/arts/orca/consts.ts b/src/arts/orca/consts.ts index 47fad1e..cbced27 100644 --- a/src/arts/orca/consts.ts +++ b/src/arts/orca/consts.ts @@ -165,6 +165,9 @@ export const ORCA_DIAGNOSTIC_EVENTS = { ACTION_SKIPPED: 'orca.action.skipped', ACTION_BLOCKED: 'orca.action.blocked', ACTION_INTERRUPTED: 'orca.action.interrupted', + COMPENSATION_STARTED: 'orca.compensation.started', + COMPENSATION_COMPLETED: 'orca.compensation.completed', + COMPENSATION_FAILED: 'orca.compensation.failed', EVENT_EMITTED: 'orca.event.emitted', REENTRY_BLOCKED: 'orca.reentry.blocked', TRACE_ABORTED: 'orca.trace.aborted', @@ -184,6 +187,9 @@ export const ORCA_LOG_MSG_ACTION_TIMEOUT = 'orca action exceeded actionTimeoutMs export const ORCA_LOG_MSG_ACTION_SKIPPED = 'orca action skipped'; export const ORCA_LOG_MSG_ACTION_BLOCKED = 'orca action blocked by gate'; export const ORCA_LOG_MSG_ACTION_INTERRUPTED = 'orca action interrupted'; +export const ORCA_LOG_MSG_COMPENSATION_STARTED = 'orca compensation started'; +export const ORCA_LOG_MSG_COMPENSATION_COMPLETED = 'orca compensation completed'; +export const ORCA_LOG_MSG_COMPENSATION_FAILED = 'orca compensation failed'; export const ORCA_LOG_MSG_EVENT_EMITTED = 'orca derived event emitted'; export const ORCA_LOG_MSG_REENTRY_BLOCKED = 'orca reentry guard blocked an event'; export const ORCA_LOG_MSG_TRACE_ABORTED = 'orca trace aborted'; diff --git a/src/arts/orca/diagnostics.ts b/src/arts/orca/diagnostics.ts index 1f14f7f..d953c21 100644 --- a/src/arts/orca/diagnostics.ts +++ b/src/arts/orca/diagnostics.ts @@ -16,6 +16,9 @@ import { ORCA_LOG_MSG_ACTION_SKIPPED, ORCA_LOG_MSG_ACTION_STARTED, ORCA_LOG_MSG_ACTION_TIMEOUT, + ORCA_LOG_MSG_COMPENSATION_COMPLETED, + ORCA_LOG_MSG_COMPENSATION_FAILED, + ORCA_LOG_MSG_COMPENSATION_STARTED, ORCA_LOG_MSG_CONFIGURATION_INVALID, ORCA_LOG_MSG_EVENT_EMITTED, ORCA_LOG_MSG_REENTRY_BLOCKED, @@ -166,6 +169,18 @@ const ORCA_DIAGNOSTIC_LOGS: DiagnosticCatalog = { level: LogLevel.WARN, message: ORCA_LOG_MSG_ACTION_INTERRUPTED }, + [ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_STARTED]: { + level: LogLevel.DEBUG, + message: ORCA_LOG_MSG_COMPENSATION_STARTED + }, + [ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_COMPLETED]: { + level: LogLevel.DEBUG, + message: ORCA_LOG_MSG_COMPENSATION_COMPLETED + }, + [ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_FAILED]: { + level: LogLevel.ERROR, + message: ORCA_LOG_MSG_COMPENSATION_FAILED + }, [ORCA_DIAGNOSTIC_EVENTS.EVENT_EMITTED]: { level: LogLevel.TRACE, message: ORCA_LOG_MSG_EVENT_EMITTED diff --git a/src/arts/orca/engine-orca.ts b/src/arts/orca/engine-orca.ts index ac68a79..d845ecd 100644 --- a/src/arts/orca/engine-orca.ts +++ b/src/arts/orca/engine-orca.ts @@ -404,6 +404,11 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca { const startedAt = timers.clock.now(); const tokens = new Set(); 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 }> = []; + const compensations: OrcaActionRun[] = []; + let compensationsRun = false; const controller = new AbortController(); activeRunController = controller; @@ -420,6 +425,30 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca { const byStage = groupByStage(actions); stageLoop: for (const stage of ORCA_STAGES_CANONICAL_ORDER) { + // Compensation phase: when transitioning into FINALLY after an + // abort, replay the LIFO stack of compensable actions before + // running cleanup. Compensation runs even if there are no + // FINALLY actions registered (we still pass through this stage + // boundary because the stage loop iterates every stage). + if ( + stage === ORCA_STAGE_FINALLY && + aborted && + !compensationsRun && + compensable.length > 0 + ) { + compensationsRun = true; + for (let i = compensable.length - 1; i >= 0; i--) { + const compResult = await runCompensation( + compensable[i].action, + compensable[i].payload, + envelope, + runId, + tokens + ); + compensations.push(compResult); + } + } + const stageActions = byStage.get(stage); if (!stageActions || stageActions.length === 0) continue; @@ -492,6 +521,17 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca { controller.signal.removeEventListener('abort', onRunAbort); actionRuns.push(actionRun); + // Track for compensation if the action completed and has a + // compensator. Stage FINALLY actions are not compensable — + // finally is itself the cleanup pass. + if ( + actionRun.status === ORCA_ACTION_STATUS_SUCCESS && + action.compensate !== undefined && + stage !== ORCA_STAGE_FINALLY + ) { + compensable.push({ action, payload: envelope.payload }); + } + if (actionRun.status === ORCA_ACTION_STATUS_FATAL) { // Fatal always aborts the run, regardless of onError. aborted = true; @@ -540,7 +580,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca { endedAt, durationMs: endedAt - startedAt, tokens: Array.from(tokens), - actions: actionRuns + actions: actionRuns, + compensations }; recordRun(runResult); @@ -621,6 +662,96 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca { traceStates.delete(traceId); } + /** + * Invoke a compensator. Each compensation runs with a fresh + * `AbortController` (the run-level controller is already aborted + * by the time we reach this code path; compensations need their + * own signal). Emits inside `ctx.emit()` are silently dropped — + * compensation is a rollback, not a place to fan out new events. + * + * Errors thrown by the compensator are recorded but do not stop + * the next compensation in the LIFO chain. v1 chooses best-effort + * over fail-fast because rolling back N-1 entries when one of them + * failed is usually still more useful than rolling back zero. + */ + async function runCompensation( + action: OrcaAction, + payload: unknown, + envelope: OrcaEnvelope, + runId: OrcaRunId, + tokens: ReadonlySet + ): Promise { + const startedAt = timers.clock.now(); + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_STARTED, { + runId, + actionId: action.id, + eventId: envelope.meta.eventId, + traceId: envelope.meta.traceId, + depth: envelope.meta.depth, + stage: action.stage + }); + const compController = new AbortController(); + const compContext: OrcaActionContext = { + runId, + event: envelope.event, + stage: action.stage, + eventId: envelope.meta.eventId, + traceId: envelope.meta.traceId, + parentEventId: envelope.meta.parentEventId, + depth: envelope.meta.depth, + tokens, + signal: compController.signal, + logger: actionLogger, + emit() { + // Compensations cannot emit derived events. Returning null + // matches the disposed-engine and reentry-blocked paths. + return null; + } + }; + try { + await action.compensate!(payload, compContext); + const endedAt = timers.clock.now(); + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_COMPLETED, { + runId, + actionId: action.id, + eventId: envelope.meta.eventId, + traceId: envelope.meta.traceId, + depth: envelope.meta.depth, + durationMs: endedAt - startedAt + }); + return { + id: action.id, + stage: action.stage, + status: ORCA_ACTION_STATUS_SUCCESS, + startedAt, + endedAt, + durationMs: endedAt - startedAt, + emitted: [] + }; + } catch (thrown) { + const endedAt = timers.clock.now(); + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_FAILED, { + runId, + actionId: action.id, + eventId: envelope.meta.eventId, + traceId: envelope.meta.traceId, + depth: envelope.meta.depth, + error: thrown, + thrown: true + }); + return { + id: action.id, + stage: action.stage, + status: ORCA_ACTION_STATUS_ERROR, + startedAt, + endedAt, + durationMs: endedAt - startedAt, + emitted: [], + error: thrown + }; + } + } + /** * Race the action's promise against a timer scheduled on the App's * `TimerScheduler`. When the timer fires first, abort the action's diff --git a/src/arts/orca/test/engine-orca.test.ts b/src/arts/orca/test/engine-orca.test.ts index 3d2b356..a6b5284 100644 --- a/src/arts/orca/test/engine-orca.test.ts +++ b/src/arts/orca/test/engine-orca.test.ts @@ -743,34 +743,6 @@ describe('EngineOrca v0.0 — disposal', () => { }); }); -describe('EngineOrca v0.0 — accepted-but-ignored fields (forward-compat)', () => { - let bus: FakeBus; - let timers: TimerScheduler; - - beforeEach(() => { - bus = createFakeBus(); - timers = createFakeTimers(); - }); - - it('accepts compensate without invoking it', async () => { - const orca = createEngineOrca({ bus, timers }); - const compensate = vi.fn(() => orcaSuccess()); - - orca.onEvent('e1', { - id: 'with-compensate', - stage: ORCA_STAGE_MAIN, - onError: ORCA_ON_ERROR_ABORT_RUN, - action: () => orcaError(new Error('x')), - compensate - }); - - bus.publish('e1', null); - await flush(); - - expect(compensate).not.toHaveBeenCalled(); - }); -}); - // ── v0-kernel: envelope, ctx.emit, reentry guards ───────────────────── describe('EngineOrca v0 — envelope and OrcaActionContext', () => { @@ -2204,3 +2176,348 @@ describe('EngineOrca v1 — validate()', () => { }); }); +// ── v1: compensate ──────────────────────────────────────────────────── + +describe('EngineOrca v1 — compensate', () => { + let bus: FakeBus; + let timers: ReturnType; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers(); + }); + + it('does not invoke compensate when the run completes successfully', async () => { + const orca = createEngineOrca({ bus, timers }); + const compensate = vi.fn(); + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess(), + compensate + }); + + bus.publish('e', null); + await flush(); + + expect(compensate).not.toHaveBeenCalled(); + expect(orca.recentRuns()[0].compensations).toEqual([]); + }); + + it('invokes compensate of a previously successful action when the run aborts', async () => { + const orca = createEngineOrca({ bus, timers }); + const compensate = vi.fn(() => orcaSuccess()); + + orca.onEvent('e', { + id: 'first', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('boom')) + }); + + bus.publish('e', null); + await flush(); + + expect(compensate).toHaveBeenCalledTimes(1); + const runs = orca.recentRuns(); + expect(runs[0].compensations).toHaveLength(1); + expect(runs[0].compensations[0].id).toBe('first'); + expect(runs[0].compensations[0].status).toBe(ORCA_ACTION_STATUS_SUCCESS); + }); + + it('does not invoke compensate of an action that errored itself', async () => { + const orca = createEngineOrca({ bus, timers }); + const ownCompensate = vi.fn(); + + orca.onEvent('e', { + id: 'failing', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('boom')), + compensate: ownCompensate + }); + + bus.publish('e', null); + await flush(); + + expect(ownCompensate).not.toHaveBeenCalled(); + }); + + it('runs compensations in LIFO order', async () => { + const orca = createEngineOrca({ bus, timers }); + const order: string[] = []; + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: () => { + order.push('a'); + return orcaSuccess(); + } + }); + orca.onEvent('e', { + id: 'b', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: () => { + order.push('b'); + return orcaSuccess(); + } + }); + orca.onEvent('e', { + id: 'c', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: () => { + order.push('c'); + return orcaSuccess(); + } + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('boom')) + }); + + bus.publish('e', null); + await flush(); + + expect(order).toEqual(['c', 'b', 'a']); + }); + + it('runs compensations even when OrcaFatal aborts the run', async () => { + const orca = createEngineOrca({ bus, timers }); + const { orcaFatal } = await import('../result.ts'); + const compensate = vi.fn(() => orcaSuccess()); + + orca.onEvent('e', { + id: 'first', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate + }); + orca.onEvent('e', { + id: 'broken', + stage: ORCA_STAGE_MAIN, + action: () => orcaFatal(new Error('cannot continue')) + }); + + bus.publish('e', null); + await flush(); + + expect(compensate).toHaveBeenCalledTimes(1); + }); + + it('continues with remaining compensations even if one throws', async () => { + const orca = createEngineOrca({ bus, timers }); + const compA = vi.fn(() => orcaSuccess()); + const compB = vi.fn(() => { + throw new Error('rollback failed'); + }); + const compC = vi.fn(() => orcaSuccess()); + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: compA + }); + orca.onEvent('e', { + id: 'b', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: compB + }); + orca.onEvent('e', { + id: 'c', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: compC + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('boom')) + }); + + bus.publish('e', null); + await flush(); + + // c runs (LIFO), then b throws but the chain continues to a + expect(compC).toHaveBeenCalledTimes(1); + expect(compB).toHaveBeenCalledTimes(1); + expect(compA).toHaveBeenCalledTimes(1); + + const runs = orca.recentRuns(); + expect(runs[0].compensations).toHaveLength(3); + expect(runs[0].compensations[1].status).toBe(ORCA_ACTION_STATUS_ERROR); + expect(runs[0].compensations[1].id).toBe('b'); + }); + + it('runs FINALLY actions after compensations', async () => { + const orca = createEngineOrca({ bus, timers }); + const order: string[] = []; + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: () => { + order.push('compensate'); + return orcaSuccess(); + } + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('boom')) + }); + orca.onEvent('e', { + id: 'cleanup', + stage: ORCA_STAGE_FINALLY, + action: () => { + order.push('finally'); + return orcaSuccess(); + } + }); + + bus.publish('e', null); + await flush(); + + expect(order).toEqual(['compensate', 'finally']); + }); + + it('does not invoke compensate when an action has compensate but no abort happens', async () => { + const orca = createEngineOrca({ bus, timers }); + const compensate = vi.fn(); + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess(), + compensate + }); + orca.onEvent('e', { + id: 'b', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_CONTINUE, + action: () => orcaError(new Error('partial')) + }); + + bus.publish('e', null); + await flush(); + + // Run is PARTIAL, not aborted — compensate not invoked. + expect(compensate).not.toHaveBeenCalled(); + expect(orca.recentRuns()[0].status).toBe(ORCA_RUN_PARTIAL); + expect(orca.recentRuns()[0].compensations).toEqual([]); + }); + + it('emits orca.compensation.{started,completed,failed} diagnostics', async () => { + const entries: DiagnosticEntry[] = []; + const orca = createEngineOrca({ + bus, + timers, + logger: createCapturingLogger(entries) as Parameters< + typeof createEngineOrca + >[0]['logger'] + }); + + orca.onEvent('e', { + id: 'good', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: () => orcaSuccess() + }); + orca.onEvent('e', { + id: 'bad', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: () => { + throw new Error('rollback failed'); + } + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('boom')) + }); + + bus.publish('e', null); + await flush(); + + expect(entries.find((e) => e.type === 'orca.compensation.started')).toBeDefined(); + expect(entries.find((e) => e.type === 'orca.compensation.completed')).toBeDefined(); + expect(entries.find((e) => e.type === 'orca.compensation.failed')).toBeDefined(); + }); + + it('passes the original event payload to compensate', async () => { + const orca = createEngineOrca({ bus, timers }); + const seen: unknown[] = []; + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: (payload) => { + seen.push(payload); + return orcaSuccess(); + } + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('x')) + }); + + bus.publish('e', { userId: 'u-7', tenantId: 't-3' }); + await flush(); + + expect(seen).toEqual([{ userId: 'u-7', tenantId: 't-3' }]); + }); + + it('emits inside ctx during compensation are dropped (return null)', async () => { + const orca = createEngineOrca({ bus, timers }); + let emitResult: string | null | undefined; + + orca.onEvent('e', { + id: 'a', + stage: ORCA_STAGE_PRE, + action: () => orcaSuccess(), + compensate: (_p, ctx) => { + emitResult = ctx.emit('derived', null); + return orcaSuccess(); + } + }); + orca.onEvent('e', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('x')) + }); + orca.onEvent('derived', { + id: 'derived-action', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess() + }); + + bus.publish('e', null); + await flush(); + + expect(emitResult).toBeNull(); + }); +}); + diff --git a/src/arts/orca/types.ts b/src/arts/orca/types.ts index 712bc01..8cfe3aa 100644 --- a/src/arts/orca/types.ts +++ b/src/arts/orca/types.ts @@ -408,6 +408,15 @@ export interface OrcaRunResult { readonly durationMs: number; readonly tokens: readonly OrcaToken[]; readonly actions: readonly OrcaActionRun[]; + /** + * Compensation runs invoked when the run aborted (FATAL, ABORT_RUN + * onError, or trace-aborted) and at least one prior action that + * declared `compensate` had completed in success. Order is LIFO: + * the last successful action's compensation appears first. Empty + * for non-aborted runs and for aborted runs where no compensable + * action ran first. + */ + readonly compensations: readonly OrcaActionRun[]; } // ── Validation (Orca.validate) ────────────────────────────────────────