Address Codex audit P1.1 / P1.2 / P1.3 / P1.4 + leftover cleanup

P1.1 — trace cleanup safe under parallel runs.
`TraceState` gains `inFlightRuns`. `spawnRun` increments it
synchronously before the IIFE awaits `executeRun`, decrements it in
its `finally` and only then attempts `maybeReleaseTrace`. The
release helper waits for both the queue snapshot and `inFlightRuns`
to be empty before dropping the entry. Previously a parallel sibling
finishing first could erase counters another run still relied on
(reentry guards, dedupeKeys, abort flags).

P1.2 — `provides` becomes the actual contract. When an action
declares a non-empty `provides`, `runAction` checks each emitted
token against it and emits `orca.configuration.invalid` for every
undeclared token. Soft enforcement: the token is *not* dropped,
keeping runtime back-compat; the diagnostic flags drift between the
declaration and the runtime so authors notice. A v2 strict-drop
mode can opt in later.

P1.3 — `replace` queue policy renamed. The constant is now
`ORCA_QUEUE_REPLACE_QUEUED` (literal `'replace-queued'`). The old
name implied `takeLatest`-style "abort in-flight + queue new", which
the engine never did. The hard variant lives in Roadmap v2 as
`'replace-current'`. Tests updated; v1 has no external consumers
yet so no back-compat alias.

P1.4 — `idFactory` becomes injectable. New `OrcaIdFactory` type +
`EngineOrcaOptions.idFactory`. Default factory uses the injected
`timers.clock.now()` (no more direct `Date.now()` violating the
"all time via timr" rule); replay/snapshot tests pass a
deterministic counter. Eliminated `generateRunId` /
`generateEventId` / `generateTraceId` standalone helpers.

Cleanup leftovers from the audit:
- `engine-orca.ts` header rewritten — was still claiming
  `after`/`unless`/`abortOn`/`actionTimeoutMs`/`compensate` are
  "accepted, ignored". Now describes the real surface.
- `README.md` "Estado Del Documento" already updated; this commit
  also drops the legacy `## Roadmap` block, removes the `setupOrca`
  recommendation (moved to roadmap), rewrites "Tokens Flag" to
  cover the with-payload form, refreshes the Diagnostics list to
  match `consts.ts`, and replaces the "Tests Requeridos" wishlist
  with a snapshot of actual coverage + the pending ecosystem test.

Tests: +5 (149 in engine-orca.test.ts, 16 in active-orca, 0 in
result.test.ts → 165 in orca; 1475 / 1475 across the repo).

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
master
dev 5 months ago
parent 38462f3282
commit 7e2907a75b

@ -233,62 +233,31 @@ Orca.onEvent(APP_EVENT_USER_IDENTITY_CHANGED, {
La accion es la que decide llamar a `App.Cache`, `App.Permissions`,
`App.Connections` o cualquier otro servicio.
## Setup Tipado
## Setup Tipado *(roadmap v2)*
`orca` debe permitir dos niveles de ergonomia:
> No implementado. Ver "Roadmap v2" más abajo.
- API minima: registrar acciones directamente con `onEvent()`.
- API recomendada para apps grandes: `setupOrca()` con eventos, tokens y action
ids declarados antes del runtime.
La forma preliminar:
La API actual es directa: `createEngineOrca()` o `createActiveOrca()` con
`onEvent()` y `configureEvent()`. Los eventos, tokens y action ids viven como
constantes exportadas por la app o por los presets:
```ts
const orcaSetup = setupOrca({
events: {
userIdentityChanged: APP_EVENT_USER_IDENTITY_CHANGED
},
tokens: {
httpPrivateCancelled: ORCA_TOKEN_HTTP_PRIVATE_CANCELLED,
cacheOk: ORCA_TOKEN_CACHE_OK,
cacheError: ORCA_TOKEN_CACHE_ERROR
},
actions: {
cancelPrivateHttp: ORCA_ACTION_CANCEL_PRIVATE_HTTP,
clearPrivateCache: ORCA_ACTION_CLEAR_PRIVATE_CACHE
const Orca = createEngineOrca({ bus, timers, logger });
Orca.onEvent(APP_EVENT_USER_IDENTITY_CHANGED, {
id: ORCA_ACTION_CLEAR_PRIVATE_CACHE,
stage: ORCA_STAGE_MAIN,
provides: [ORCA_TOKEN_CACHE_CLEARED],
action: (payload, ctx) => {
// ...
}
});
const Orca = orcaSetup.createEngine({
bus: App.Bus,
timers: App.Timers,
logger: App.Logger
});
```
Ventajas:
- detectar tokens fantasma en `after`, `unless` o `abortOn`
- detectar action ids desconocidos
- inferir payload del evento
- limitar `emits` a tokens declarados
- generar documentacion/devtools del grafo
`setupOrca()` no debe conocer modulos. Solo declara lenguaje de orquestacion:
eventos, tokens, acciones, stages y policies.
Para prototipos y tests pequenos, la API directa sigue siendo valida:
```ts
const Orca = createEngineOrca({ bus, timers, logger });
Orca.onEvent(APP_EVENT_USER_IDENTITY_CHANGED, ACTION_CLEAR_CACHE);
```
La regla de producto:
```txt
v0 soporta API directa; proyectos serios deberian usar setupOrca().
```
Para v2 se valora un wrapper tipado tipo `setupOrca({ events, tokens,
actions })` que infiera payloads, valide tokens fantasma en gates,
restrinja `emits` a `provides` y emita un grafo navegable. Hoy esa
validación es estática y best-effort vía `Orca.validate()` (ver
sección "Validate").
## Eventos
@ -413,48 +382,38 @@ El motor debe decidir grupos de ejecucion usando:
## Concurrencia De Runs
Un run es una ejecucion concreta de un evento. El pipeline interno de un run no
resuelve por si solo que ocurre cuando el mismo evento entra varias veces
mientras una ejecucion anterior sigue viva. Esa decision debe ser parte del
contrato del evento.
Configuracion preliminar:
```ts
Orca.configureEvent(APP_EVENT_USER_IDENTITY_CHANGED, {
concurrency: ORCA_CONCURRENCY_QUEUE,
runTimeoutMs: 5_000
});
```
Un run es una ejecución concreta de un evento. El pipeline interno
de un run no resuelve por sí solo qué ocurre cuando el mismo evento
entra varias veces mientras una ejecución anterior sigue viva. Esa
decisión es parte del contrato del evento y se configura con
`Orca.configureEvent(event, { queuePolicy })`.
Modos previstos:
Constantes vivas (`consts.ts`):
```ts
export const ORCA_CONCURRENCY_QUEUE = 'queue' as const;
export const ORCA_CONCURRENCY_DROP = 'drop' as const;
export const ORCA_CONCURRENCY_REPLACE = 'replace' as const;
export const ORCA_CONCURRENCY_PARALLEL = 'parallel' as const;
ORCA_QUEUE_FIFO // 'fifo'
ORCA_QUEUE_REPLACE_QUEUED // 'replace-queued'
ORCA_QUEUE_DROP_LATEST // 'drop-latest'
ORCA_QUEUE_PARALLEL // 'parallel'
```
Semantica:
Semántica:
| Modo | Uso |
| ---- | --- |
| `queue` | Default obligatorio. No solapa runs del mismo evento; ejecuta el siguiente al terminar el actual. |
| `drop` | Ignora eventos repetidos mientras hay un run activo. Util para doble click o refresh redundante. |
| `replace` | Aborta el run activo y empieza uno nuevo. Util para ultima intencion gana. |
| `parallel` | Permite runs simultaneos del mismo evento. Footgun; queda fuera de v0. |
| `fifo` (default) | Cada evento encola un run; los runs no-paralelos se serializan globalmente. Garantía simple. |
| `replace-queued` | Si llega un evento mientras hay otro encolado, el queued se descarta y se sustituye por el nuevo. **No aborta in-flight.** Última intención pendiente gana. |
| `drop-latest` | Si hay un run del mismo evento in-flight o encolado, el incoming se descarta. Ignora retriggers durante trabajo. |
| `parallel` | Lanza runs concurrentes; eventos paralelos no toman el lock global. Footgun: solapa side-effects. |
`parallel` no debe ser el default. En orquestaciones destructivas como cambio de
identidad, logout, tenant switch, permisos o cache privada, solapar runs es una
fuente directa de condiciones de carrera.
`replace-queued` deja el run activo terminar (FINALLY incluido) — la
variante "fuerte" que aborta in-flight (`replace-current`, takeLatest)
queda en Roadmap v2.
`replace` debe ejecutar `cleanup` y `finally` del run abortado antes de cerrar
su `RunResult`. Si no lo hace, deja recursos/timers/overlays en un estado
ambiguo.
`parallel` queda documentado como posibilidad futura, pero no entra en v0. Si
se implementa, debe exigir opt-in explicito y diagnostics visibles.
`parallel` no es default por una razón: en orquestaciones destructivas
(cambio de identidad, logout, tenant switch, invalidación de permisos,
cache privada), solapar runs es una fuente directa de condiciones de
carrera. Opt-in explícito por evento y diagnostics visibles.
## Tokens
@ -514,20 +473,40 @@ Esta regla no debe ser configurable en v0. Si se necesita estado persistente,
la fuente de verdad es un artefacto (`sess`, `cach`, `perm`, `stor`) y no
`orca`.
### Tokens Flag En v0
### Tokens Flag Y Tokens Con Payload
En v0, los tokens son flags sin payload. Representan que un hecho ocurrio dentro
del run, no transportan datos.
Los tokens son nombres semánticos (`string`) que un run acumula a medida
que las acciones los emiten. Forma básica — flag puro, sin datos:
```ts
return orcaSuccess({ emits: [ORCA_TOKEN_CACHE_OK] });
```
Si una accion posterior necesita datos, debe recibirlos por el payload del
evento, por el `value` de su propia accion, o consultarlos en el modulo dueño
del estado. Los tokens con payload tipado son una extension posible, pero no
entran en v0 porque complican mucho el mapa `Token -> Payload` y la inferencia
del contexto.
Una acción aguas abajo decide si actuar consultando `ctx.tokens.has(...)`
directamente o vía gates declarativos (`after` / `unless` / `abortOn` /
`fanIn`). Si necesita transportar datos (más allá de "esto pasó"), usa
la forma con payload — `emits` acepta entradas
`{ token, payload }` mezcladas con strings sueltos:
```ts
return orcaSuccess({
emits: [
ORCA_TOKEN_CACHE_OK,
{ token: ORCA_TOKEN_USER_LOADED, payload: { id: 'u-7', tenantId: 't-3' } }
]
});
```
El payload se lee aguas abajo con `ctx.tokenPayloads.get(token)` (un
`ReadonlyMap<OrcaToken, unknown>`). Tokens emitidos como string suelto
**no** aparecen en el mapa, así que `ctx.tokens.has('x')` y
`ctx.tokenPayloads.has('x')` pueden divergir (presencia vs. payload).
Los gates siguen comparando solo nombres — el payload es un canal
ortogonal para coordinación intra-run.
Limitación tipada: `payload` está tipado como `unknown`; los authors
hacen cast. La inferencia tipo `setupOrca({ tokens: { Token:
SchemaPayload } })` queda para v2.
### Validacion Del Grafo
@ -1204,71 +1183,54 @@ Protecciones obligatorias:
## Diagnostics
`orca` debe usar `Logger` comun de `libs/logr` y una capa de diagnostics
catalogada, igual que el resto de artefactos.
Eventos esperados:
`orca` usa el `Logger` común de `libs/logger` con una capa de
diagnostics catalogada. Las claves vivas (`ORCA_DIAGNOSTIC_EVENTS` en
`consts.ts`):
```ts
ORCA_DIAGNOSTIC_RUN_STARTED
ORCA_DIAGNOSTIC_RUN_COMPLETED
ORCA_DIAGNOSTIC_RUN_ABORTED
ORCA_DIAGNOSTIC_ACTION_STARTED
ORCA_DIAGNOSTIC_ACTION_COMPLETED
ORCA_DIAGNOSTIC_ACTION_FAILED
ORCA_DIAGNOSTIC_ACTION_TIMEOUT
ORCA_DIAGNOSTIC_ACTION_BLOCKED
ORCA_DIAGNOSTIC_TOKEN_EMITTED
ORCA_DIAGNOSTIC_CONFIGURATION_INVALID
'orca.run.started'
'orca.run.completed'
'orca.run.aborted'
'orca.action.started'
'orca.action.completed'
'orca.action.failed'
'orca.action.fatal'
'orca.action.timeout'
'orca.action.skipped'
'orca.action.blocked'
'orca.action.interrupted'
'orca.compensation.started'
'orca.compensation.completed'
'orca.compensation.failed'
'orca.event.emitted'
'orca.reentry.blocked'
'orca.trace.aborted'
'orca.queue.dropped'
'orca.configuration.invalid'
```
Los mensajes/log categories deben estar en `consts.ts`, no hardcodeados en el
runtime.
## Tests Requeridos
Core minimo:
- registra y elimina acciones por evento
- ejecuta acciones en orden de stage
- respeta priority dentro de stage
- ejecuta acciones paralelas cuando no tienen dependencias
- no solapa runs del mismo evento con `queue`
- ignora repetidos con `drop`
- cancela/reemplaza con `replace`
- ejecuta `cleanup` y `finally` del run reemplazado
- no implementa `parallel` en v0
- no procesa eventos publicados durante un run inline
- espera tokens antes de ejecutar una accion
- mantiene tokens scoped por run
- no filtra tokens entre dos runs del mismo evento
- trata tokens como flags sin payload en v0
- valida grafo con tokens imposibles
- salta accion por `unless`
- bloquea/aborta accion por `abortOn`
- convierte throw en `OrcaError`
- respeta `onError: continue`
- respeta `onError: abort-run`
- respeta `fatal` como abort por defecto
- cancela timers al finalizar
- timeout de accion produce `ORCA_RESULT_TIMEOUT`
- timeout de espera por token bloquea o aborta segun policy
- detecta ciclos de dependencia
- rechaza `ORCA_TX_REQUIRED` sin transaction port
- respeta `maxRuns` en ActiveOrca
- `dispose()` es idempotente
- evento durante dispose no ejecuta acciones disposed
Tests compuestos con ecosistema:
- cambio de usuario cancela HTTP privado, limpia cache, invalida permisos y
reautentica conexiones en orden
- error de cache bloquea reauth de conexiones
- permisos cambiados desde servidor invalidan `perm` y cache permission-scoped
- session revoked cierra conexiones y limpia estado actor-scoped
- connection server event publica en `buss` y `orca` ejecuta pipeline
- `timr` controla todos los timeouts de forma determinista
- diagnostics no contienen credenciales
Cada evento carga meta con `runId` / `eventId` / `traceId` / `depth`
cuando aplica, más campos específicos (`durationMs`, `error`, `reason`,
etc). Los mensajes y niveles viven en `diagnostics.ts` — nunca
hardcodeados en el runtime.
## Tests
Cobertura unitaria del motor en
`src/arts/orca/test/engine-orca.test.ts` (registro, stages, error
policies, run trace, dispose, envelope/contexto, reentry guards,
diagnostics, `actionTimeoutMs`, `OrcaFatal`, gates `after`/`unless`/
`abortOn`/`fanIn`, `validate()`, `compensate`, `commit()`, `parallel`
waves, `transaction`, tokens con payload, queue policies,
bus interception). Cobertura del wrapper reactivo en
`src/arts/orca/test/active-orca.svelte.test.ts` (snapshots
reactivos vía `engine.onChange`).
**Pendiente — test compuesto del ecosistema** (P3.15 en
`docs/audit-codex.md`): cambio de usuario A → B que valida en orden
limpieza de cache, invalidación de perm y reauth/close de
connection, demostrando que ningún dato/credencial de A queda
observable. Ver "Roadmap v2 — pendientes desde la auditoría".
## Invariantes
@ -1392,10 +1354,12 @@ Aceptado en el contrato público y honrado por el motor — v1 cerrado al 100%.
- `'fifo'` (default): cada evento encola un run; los runs se
serializan globalmente con cualquier otro run no-paralelo
(preservando la invariante v0).
- `'replace'`: cuando llega un nuevo evento del mismo nombre y
ya hay uno encolado, el queued se descarta (se emite
- `'replace-queued'`: cuando llega un nuevo evento del mismo
nombre y ya hay uno encolado, el queued se descarta (se emite
`orca.queue.dropped` con `reason: 'replaced'`). El in-flight
NO se aborta; máximo "1 in-flight + 1 queued" por evento.
NO se aborta; máximo "1 in-flight + 1 queued" por evento. La
variante fuerte takeLatest (abortar in-flight + encolar nuevo)
queda reservada como `'replace-current'` para v2.
- `'drop-latest'`: si hay un run del mismo evento in-flight o
encolado, el incoming se descarta (`reason: 'drop-latest'`).
- `'parallel'`: los runs se lanzan concurrentemente vía
@ -1454,9 +1418,10 @@ forma libre antes de aterrizar.
### Diferidas explícitamente desde v1
- **`replace` con abort-in-flight** — variante "fuerte" de
`replace` (estilo `takeLatest`) que aborta el run en vuelo además
de descartar los queued. Hoy `replace` solo afecta a la cola.
- **`'replace-current'` policy** — variante "fuerte" de
`'replace-queued'` (estilo `takeLatest`) que aborta el run en
vuelo además de descartar los queued. Hoy `'replace-queued'` solo
afecta a la cola.
- **Transactions cross-event** — un `transaction` hoy vive en un
solo evento; spanning entre eventos necesita reconciliar tracing
y orden de compensación.

@ -72,11 +72,13 @@ export const ORCA_RUN_TIMEOUT = 'timeout' as const;
export const ORCA_QUEUE_FIFO = 'fifo' as const;
/**
* When a new event arrives for an event that already has a queued
* run, the queued run is **replaced** by the newer envelope. Does not
* abort an in-flight run — at most one in-flight + one queued at any
* point. Useful for "I only care about the latest" patterns.
* run, the queued run is **replaced** by the newer envelope. Does
* **not** abort an in-flight run — at most one in-flight + one
* queued at any point. Useful for "I only care about the latest
* pending intent" patterns. The takeLatest-style "abort in-flight
* + queue new" variant is reserved for `replace-current` (v2).
*/
export const ORCA_QUEUE_REPLACE = 'replace' as const;
export const ORCA_QUEUE_REPLACE_QUEUED = 'replace-queued' as const;
/**
* When an event arrives while another run for the same event is in
* flight or already queued, the **incoming** event is dropped. Useful

@ -1,30 +1,58 @@
/**
* `EngineOrca` v0 — orchestration kernel.
* `EngineOrca` — orchestration kernel for the orca artifact.
*
* Surface complete, motor focused on the kernel responsibilities:
* Responsibilities:
* - Stages run in canonical order; `finally` always runs.
* - Run-queue per trace: events derived via `ctx.emit()` are enqueued,
* never executed inline. Sibling events of the same trace run in
* emission order, after the producing run completes.
* - Envelope wrapping: every event the engine processes carries
* runtime metadata (eventId, traceId, parentEventId, depth, stack)
* separate from the user-defined payload.
* - Run queue with per-event policy (`fifo` / `replace` /
* `drop-latest` / `parallel`). Non-parallel events take a global
* serial lock; `parallel` events bypass it and can race with
* anyone.
* - Envelope wrapping: every event carries runtime metadata
* (eventId, traceId, parentEventId, depth, stack) separate from
* the user-defined payload.
* - Reentry guards: maxDepth, maxEventsPerTrace, repeatedEventLimit,
* dedupeKey. Guards apply only to derived events (`ctx.emit()`);
* module publishes via the bus directly create root envelopes with
* fresh traces and never enter the counters.
* - `onError`: CONTINUE vs ABORT_RUN. Other policies accepted, ignored.
* - Exception capture: thrown values become `OrcaError` results.
* - Run trace exposed via `recentRuns()`.
* - Lazy bus subscription: registered per event on first `onEvent`,
* released when the last action for that event is detached.
* - Dispose: cancels the current run via signal, drops the queue,
* unsubscribes the bus, marks all live traces aborted.
*
* Accepted but ignored at runtime (preserved in the public surface for
* forward compatibility): `after`, `unless`, `abortOn`,
* `actionTimeoutMs`, `compensate`. The token set in
* `OrcaActionContext` is populated correctly so authors can read it.
* dedupeKey, with `skip` / `abort-trace` / `error` policies.
* Module publishes on the bus directly create root envelopes with
* fresh traces and don't enter the counters; events emitted from
* inside an action body via `ctx.emit()` or attributed
* `bus.publish` (AsyncLocalStorage) do enter them.
* - Wave execution: consecutive `parallel: true` actions in the
* same stage form a wave run via `Promise.all`. Token sets are
* wave-scoped snapshots; emissions merge after settle.
* - Gates: `unless` → `abortOn` → `fanIn` → `after` (eval order).
* - Per-action `actionTimeoutMs` racing the action against a timer
* scheduled on the injected `TimerScheduler`; on expiry the
* action's signal is aborted and `ORCA_RESULT_TIMEOUT` is
* produced. Cooperative — the action body must honour
* `ctx.signal` to stop side effects.
* - Compensation: LIFO rollback before `FINALLY` for
* successfully-completed actions that declare `compensate`.
* Transactions (`transaction: 'tx-id'`) trigger an immediate
* scoped rollback when a member fails; tx members are marked
* compensated so the global rollback skips them.
* - `onError`: `CONTINUE` (default) and `ABORT_RUN` are honoured;
* `ABORT_ACTION` and `ABORT_STAGE` accepted in the type and
* treated as `CONTINUE`. Members of a `transaction` ignore
* `onError` — tx failure always aborts.
* - Validate / commit: `validate()` reports configuration issues
* (`unsatisfiable-after`, `unsatisfiable-fan-in`,
* `dependency-cycle`, `orphan-unless`, `orphan-abort-on`,
* `transaction-without-compensate`); `commit()` freezes the
* graph (subsequent `onEvent` / `configureEvent` throw
* `OrcaFrozenError`).
* - Bus interception: `bus.publish` from inside an action body is
* attributed to the active run via AsyncLocalStorage when
* available (Node/Bun, modern browsers with AsyncContext).
* Browsers without AsyncContext fall back to root-event
* semantics; authors should rely on `ctx.emit()` for
* attribution-critical paths.
* - Run trace exposed via `recentRuns()` (bounded by `maxRuns`)
* and reactive `onChange()` listeners drive
* `createActiveOrca()`.
* - Lazy bus subscription per event; released when the last action
* for that event detaches.
* - Dispose: aborts every in-flight controller, unsubscribes from
* the bus, drops the queue, marks live traces aborted.
*/
import {
@ -52,7 +80,7 @@ import {
ORCA_QUEUE_DROP_REASON_REPLACED,
ORCA_QUEUE_FIFO,
ORCA_QUEUE_PARALLEL,
ORCA_QUEUE_REPLACE,
ORCA_QUEUE_REPLACE_QUEUED,
ORCA_REENTRY_ABORT_TRACE,
ORCA_REENTRY_ERROR,
ORCA_REENTRY_REASON_DEDUPED,
@ -108,6 +136,7 @@ import type {
OrcaEnvelope,
OrcaEventId,
OrcaEventOptions,
OrcaIdFactory,
OrcaQueuePolicy,
OrcaReentryOptions,
OrcaReentryPolicy,
@ -152,6 +181,15 @@ interface TraceState {
readonly dedupeKeys: Set<string>;
aborted: boolean;
abortedReason?: OrcaReentryReason;
/**
* Runs of this trace currently executing. Incremented in `spawnRun`
* synchronously, decremented in its `finally`. `maybeReleaseTrace`
* waits for both this counter and the queue snapshot to be empty
* before dropping the entry — required because `parallel` queue
* policy lets sibling runs of the same trace coexist, and a run
* finishing first must not erase counters another run still needs.
*/
inFlightRuns: number;
}
interface ResolvedReentry {
@ -164,6 +202,7 @@ interface ResolvedReentry {
export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
const { bus, timers } = options;
const maxRuns = options.maxRuns ?? ORCA_DEFAULT_MAX_RUNS;
const idFactory = options.idFactory ?? createDefaultIdFactory(timers);
const reentry = resolveReentry(options.reentry);
const diagnostics = createOrcaDiagnostics(options.logger);
const actionLogger: Logger = options.logger ?? createNoopLogger();
@ -318,8 +357,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
event,
payload,
meta: {
eventId: generateEventId(),
traceId: generateTraceId(),
eventId: idFactory('evt'),
traceId: idFactory('trc'),
depth: 0,
stack: [event],
publishedAt: timers.clock.now()
@ -339,7 +378,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
event,
payload,
meta: {
eventId: generateEventId(),
eventId: idFactory('evt'),
traceId: parent.meta.traceId,
parentEventId: parent.meta.eventId,
parentRunId,
@ -417,7 +456,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
});
return envelope.meta.eventId;
}
} else if (policy === ORCA_QUEUE_REPLACE) {
} else if (policy === ORCA_QUEUE_REPLACE_QUEUED) {
// Walk the queue and remove every envelope of this event name.
// In-flight runs are not aborted: at-most one in-flight + one
// queued at any moment for replace events.
@ -450,7 +489,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
eventCount: 0,
eventCounts: new Map(),
dedupeKeys: new Set(),
aborted: false
aborted: false,
inFlightRuns: 0
};
traceStates.set(traceId, trace);
}
@ -542,6 +582,11 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
inFlightControllers.add(controller);
inFlightByEvent.set(event, (inFlightByEvent.get(event) ?? 0) + 1);
if (!isParallel) nonParallelInFlight += 1;
// Per-trace in-flight bookkeeping: ensures `maybeReleaseTrace`
// keeps the trace alive until *every* run of this trace has
// settled, even with parallel queue policy.
const traceForRun = ensureTraceState(envelope.meta.traceId);
traceForRun.inFlightRuns += 1;
notifyChange();
void (async () => {
try {
@ -552,6 +597,11 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
if (cur <= 1) inFlightByEvent.delete(event);
else inFlightByEvent.set(event, cur - 1);
if (!isParallel) nonParallelInFlight -= 1;
const traceAfter = traceStates.get(envelope.meta.traceId);
if (traceAfter) {
traceAfter.inFlightRuns -= 1;
maybeReleaseTrace(envelope.meta.traceId);
}
notifyChange();
if (!disposed) {
queueMicrotask(() => drainQueue());
@ -565,7 +615,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
actions: readonly OrcaAction[],
controller: AbortController
): Promise<void> {
const runId = generateRunId();
const runId = idFactory('run');
const startedAt = timers.clock.now();
const tokens = new Set<OrcaToken>();
const tokenPayloads = new Map<OrcaToken, unknown>();
@ -791,11 +841,11 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
actionCount: actionRuns.length
});
// Trace bookkeeping cleanup: when no more queued runs reference
// this trace and no run is in flight for it, drop the entry. We
// keep aborted traces around briefly so late `ctx.emit()` calls
// from compensating logic can still see the abort reason.
maybeReleaseTrace(envelope.meta.traceId);
// Trace bookkeeping cleanup is owned by `spawnRun`'s `finally`,
// which decrements `trace.inFlightRuns` and then calls
// `maybeReleaseTrace`. We deliberately don't release here so a
// run that completes while a parallel sibling of the same trace
// is still active doesn't strip counters the sibling needs.
}
async function runActionInWave(
@ -926,12 +976,11 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
}
function maybeReleaseTrace(traceId: OrcaTraceId): void {
const trace = traceStates.get(traceId);
if (!trace) return;
if (trace.inFlightRuns > 0) return;
const stillQueued = runQueue.some((entry) => entry.envelope.meta.traceId === traceId);
if (stillQueued) return;
// Run is still in flight if controller exists and refers to this
// trace; we cannot easily know without extra bookkeeping, but
// since drainQueue runs sequentially, by the time we reach this
// line the active run already finished.
traceStates.delete(traceId);
}
@ -1130,6 +1179,22 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
}
}
// Soft enforcement of `provides`: when the action declared a
// non-empty `provides`, any emitted token not listed there
// is reported as `orca.configuration.invalid`. The token is
// NOT dropped to keep behaviour back-compat — `validate()`
// is the authoritative declaration; the diagnostic flags
// drift between contract and runtime so authors notice.
if (action.provides && action.provides.length > 0 && emittedNames.length > 0) {
const declared = new Set(action.provides);
for (const token of emittedNames) {
if (declared.has(token)) continue;
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.CONFIGURATION_INVALID, {
reason: `action "${action.id}" emitted token "${token}" not declared in \`provides\``
});
}
}
const status = mapResultToActionStatus(result);
if (status === ORCA_ACTION_STATUS_FATAL) {
@ -1510,16 +1575,16 @@ function evaluateGates(
return null;
}
function generateRunId(): OrcaRunId {
return `run_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`;
}
function generateEventId(): OrcaEventId {
return `evt_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`;
}
function generateTraceId(): OrcaTraceId {
return `trc_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`;
/**
* Default id factory: timestamp from the injected clock + random
* suffix. Stays sortable and human-friendly without burning the
* `Date.now()` direct call at module scope (the engine's contract is
* "all time goes through `timr`"). Replay / snapshot tests pass a
* deterministic counter via `EngineOrcaOptions.idFactory`.
*/
function createDefaultIdFactory(timers: TimerScheduler): OrcaIdFactory {
return (kind) =>
`${kind}_${timers.clock.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`;
}
function createNoopLogger(): Logger {

@ -45,7 +45,7 @@ export {
ORCA_ON_ERROR_ABORT_STAGE,
ORCA_ON_ERROR_ABORT_RUN,
ORCA_QUEUE_FIFO,
ORCA_QUEUE_REPLACE,
ORCA_QUEUE_REPLACE_QUEUED,
ORCA_QUEUE_DROP_LATEST,
ORCA_QUEUE_PARALLEL,
ORCA_QUEUE_DROP_REASON_REPLACED,
@ -121,6 +121,7 @@ export type {
OrcaEventOptions,
OrcaFanInSpec,
OrcaFatal,
OrcaIdFactory,
OrcaInterrupted,
OrcaQueuePolicy,
OrcaReentryOptions,

@ -3852,7 +3852,7 @@ describe('EngineOrca v1 — queue policies', () => {
const orca = createEngineOrca({ bus, timers });
const order: number[] = [];
orca.configureEvent('e', { queuePolicy: 'replace' });
orca.configureEvent('e', { queuePolicy: 'replace-queued' });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
@ -4000,14 +4000,14 @@ describe('EngineOrca v1 — queue policies', () => {
it('configureEvent throws on conflicting policies for the same event', () => {
const orca = createEngineOrca({ bus, timers });
orca.configureEvent('e', { queuePolicy: 'replace' });
orca.configureEvent('e', { queuePolicy: 'replace-queued' });
expect(() => orca.configureEvent('e', { queuePolicy: 'parallel' })).toThrow();
});
it('configureEvent is idempotent when called with the same policy twice', () => {
const orca = createEngineOrca({ bus, timers });
orca.configureEvent('e', { queuePolicy: 'replace' });
expect(() => orca.configureEvent('e', { queuePolicy: 'replace' })).not.toThrow();
orca.configureEvent('e', { queuePolicy: 'replace-queued' });
expect(() => orca.configureEvent('e', { queuePolicy: 'replace-queued' })).not.toThrow();
});
it('configureEvent throws OrcaFrozenError after commit()', () => {
@ -4268,3 +4268,172 @@ describe('EngineOrca v1 — bus interception', () => {
});
});
describe('EngineOrca v1 — trace parallel safety', () => {
let bus: FakeBus;
let timers: TimerScheduler;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('keeps trace counters alive while parallel runs of the same trace are still in flight', async () => {
const orca = createEngineOrca({
bus,
timers,
reentry: { repeatedEventLimit: 2 }
});
const observed: { event: string; traceId: string }[] = [];
// Root action emits two children sequentially via ctx.emit; the
// children belong to the same trace. The slower child finishes
// after the fast one, so the fast one's `spawnRun.finally` would
// be the first candidate to release the trace — yet doing so
// would erase the counters the slower run still needs.
orca.onEvent('root', {
id: 'r',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
ctx.emit('child', null);
ctx.emit('child', null);
return orcaSuccess();
}
});
orca.configureEvent('child', { queuePolicy: 'parallel' });
orca.onEvent('child', {
id: 'fast',
stage: ORCA_STAGE_MAIN,
action: async (_p, ctx) => {
await sleep(2);
observed.push({ event: ctx.event, traceId: ctx.traceId });
return orcaSuccess();
}
});
orca.onEvent('child', {
id: 'slow',
stage: ORCA_STAGE_MAIN,
action: async (_p, ctx) => {
await sleep(20);
observed.push({ event: ctx.event, traceId: ctx.traceId });
return orcaSuccess();
}
});
bus.publish('root', null);
await sleep(40);
await flush();
// The repeatedEventLimit guard would normally cap `child` at 2
// per trace (the second emit hits the limit). The point of the
// test is the trace state is consistent — both runs see the
// same traceId, and the engine doesn't crash because the trace
// got dropped while the slow run was still active.
const traceIds = new Set(observed.map((o) => o.traceId));
expect(traceIds.size).toBe(1);
});
});
describe('EngineOrca v1 — provides soft enforcement', () => {
let bus: FakeBus;
let timers: TimerScheduler;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('emits orca.configuration.invalid when an action emits a token not in its provides', async () => {
const entries: DiagnosticEntry[] = [];
const orca = createEngineOrca({
bus,
timers,
logger: createCapturingLogger(entries) as Parameters<
typeof createEngineOrca
>[0]['logger']
});
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
provides: ['declared'],
action: () =>
orcaSuccess({
emits: ['declared', 'undeclared']
})
});
bus.publish('e', null);
await flush();
const violations = entries.filter(
(d) =>
d.type === 'orca.configuration.invalid' &&
typeof d.meta.reason === 'string' &&
(d.meta.reason as string).includes('undeclared')
);
expect(violations).toHaveLength(1);
// Token is NOT dropped — runtime stays back-compatible.
expect(orca.recentRuns()[0].tokens).toContain('undeclared');
});
it('does not flag an action that declares no provides', async () => {
const entries: DiagnosticEntry[] = [];
const orca = createEngineOrca({
bus,
timers,
logger: createCapturingLogger(entries) as Parameters<
typeof createEngineOrca
>[0]['logger']
});
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
// no `provides` field at all
action: () => orcaSuccess({ emits: ['anything'] })
});
bus.publish('e', null);
await flush();
const violations = entries.filter((d) => d.type === 'orca.configuration.invalid');
expect(violations).toEqual([]);
});
});
describe('EngineOrca v1 — idFactory', () => {
let bus: FakeBus;
let timers: TimerScheduler;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('honours an injected deterministic factory for run / event / trace ids', async () => {
const counters = { run: 0, evt: 0, trc: 0 };
const orca = createEngineOrca({
bus,
timers,
idFactory: (kind) => {
counters[kind] += 1;
return `${kind}-${counters[kind]}`;
}
});
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess()
});
bus.publish('e', null);
await flush();
const run = orca.recentRuns()[0];
expect(run.id).toBe('run-1');
expect(run.eventId).toBe('evt-1');
expect(run.traceId).toBe('trc-1');
});
});

@ -34,7 +34,7 @@ import type {
ORCA_QUEUE_DROP_LATEST,
ORCA_QUEUE_FIFO,
ORCA_QUEUE_PARALLEL,
ORCA_QUEUE_REPLACE,
ORCA_QUEUE_REPLACE_QUEUED,
ORCA_REENTRY_SKIP,
ORCA_REENTRY_ABORT_TRACE,
ORCA_REENTRY_ERROR,
@ -204,7 +204,7 @@ export type OrcaResult<TValue = unknown> =
export type OrcaQueuePolicy =
| typeof ORCA_QUEUE_FIFO
| typeof ORCA_QUEUE_REPLACE
| typeof ORCA_QUEUE_REPLACE_QUEUED
| typeof ORCA_QUEUE_DROP_LATEST
| typeof ORCA_QUEUE_PARALLEL;
@ -636,6 +636,16 @@ export interface OrcaValidateResult {
*/
export type OrcaBus = EventSubscriber<Record<string, unknown>>;
/**
* Factory used by the engine to mint a `run` / `event` / `trace` id.
* Default factories combine `timers.clock.now()` with `Math.random`;
* tests and replay tooling can inject deterministic counters here so
* snapshots remain stable. Each factory receives the artifact prefix
* (`'run' | 'evt' | 'trc'`) so a single counter can serve all three
* if desired.
*/
export type OrcaIdFactory = (kind: 'run' | 'evt' | 'trc') => string;
export interface EngineOrcaOptions {
/**
* Subscriber for bus events. orca only listens; it never publishes.
@ -663,6 +673,12 @@ export interface EngineOrcaOptions {
* these counters. See `OrcaReentryOptions` for the individual knobs.
*/
readonly reentry?: OrcaReentryOptions;
/**
* Override how run / event / trace ids are minted. Default uses
* `timers.clock.now()` plus a random suffix. Replay or
* snapshot-stable tests provide a deterministic factory here.
*/
readonly idFactory?: OrcaIdFactory;
}
export interface EngineOrca {

Loading…
Cancel
Save

Powered by TurnKey Linux.