Promote orca v0 from callback runner to orchestration kernel

Implements the v0-kernel design captured in docs/orca_minds.txt: the
engine now carries every event through an internal envelope, exposes
trace identity to actions, and bounds derived-event chains through
configurable reentry guards. The previous engine was effectively a
callback runner with stages; this commit turns it into a runtime that
can answer "where does this event come from, how deep is it, and when
should I stop?" without polluting user-defined payloads with runtime
metadata.

New public surface:
  - OrcaEnvelope<TPayload> + OrcaEventMeta (eventId, traceId,
    parentEventId, parentRunId, emittedByAction, depth, stack,
    publishedAt, dedupeKey)
  - OrcaActionContext gains eventId / traceId / parentEventId / depth
    plus emit(event, payload, options?) -> OrcaEventId | null
  - OrcaResult adds OrcaInterrupted (with orcaInterrupted() helper)
  - OrcaActionRun.status and OrcaRunResult.status add 'interrupted'
  - OrcaRunResult exposes eventId / traceId / parentEventId / depth
  - OrcaReentryOptions on EngineOrcaOptions: maxDepth (16),
    maxEventsPerTrace (128), repeatedEventLimit (2),
    repeatedEventPolicy (skip / abort-trace / error)
  - Constants for reentry policies and reasons; OrcaReentryError class

Engine semantics:
  - Bus publishes are roots: fresh traceId, depth=0, no parent. They
    never enter the reentry counters.
  - ctx.emit() builds a child envelope inheriting the parent's traceId
    and incrementing depth. The child is enqueued, never executed
    inline.
  - Reentry guards apply only to derived envelopes. Crossing maxDepth,
    maxEventsPerTrace, repeatedEventLimit, or matching a previous
    dedupeKey triggers the configured policy. Skip blocks just that
    envelope; abort-trace marks the trace and skips every queued event
    that belongs to it; error throws OrcaReentryError synchronously.
  - dispose() marks every live trace aborted with reason 'disposed' so
    late ctx.emit() calls (e.g. from compensating cleanup) get a clean
    rejection instead of an exception.

Diagnostics gain four new event types
  (action.interrupted, event.emitted, reentry.blocked, trace.aborted)
and every existing one carries the envelope identifiers
(runId, eventId, traceId, parentEventId, depth) where applicable, so a
log sink can correlate runs without parsing variant tags.

Tests: 10 new tests covering envelope identity (root depth=0, child
depth+1), ctx.emit() trace inheritance, run-trace correlation, orphan
emit, all four reentry guards (maxDepth with disjoint event names so
the same-name limit doesn't interfere, repeatedEventLimit, error
policy throwing OrcaReentryError, abort-trace, dedupeKey), and the
invariant that bus publishes start fresh traces.

Active-app presets keep working unchanged (they don't call ctx.emit
yet); the API extension is additive on the OrcaActionContext side
(presets still type-check against the wider context shape).

Total: 1357/1357 vitest tests pass (47 in orca, +10 from this commit).

Roadmap v1 documented at the foot of the orca README and parked under
@v1+ in the source: actionTimeoutMs runtime, compensate invocation,
after/unless/abortOn gating, OrcaFatal distinction, queue policies
(commit/replace/parallel), transaction/atomic groups, payload-bearing
tokens, fan-in, validate()/commit() static graph validation,
createActiveOrca() reactive wrapper, and the bus.publish interception
(Option B) that would let modules' direct publishes attach
emittedByAction perfectly.

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

@ -0,0 +1,691 @@
Clippy-pet verdict: `orca` es una ballena preciosa, pero si intentamos montarla entera en v0 nos hunde la barca.
Lo correcto:
```txt
v0 = kernel
v1 = enterprise orchestration
```
`orca` está bien pensado porque tiene tres aciertos enormes:
- Desacopla módulos: `orca` no sabe qué es `cach`, `perm`, `sess` o `connection`.
- Convierte flujos críticos en artefactos trazables: `OrcaRunResult`.
- Usa tokens como hechos semánticos per-run, que es la parte diferencial.
Pero ahora mismo el README mezcla kernel y futuro. Para que sea implementable, v0 debería recortarse así:
```txt
entra en v0:
setupOrca
createEngineOrca
App.Orchestration inerte
onEvent
stages
queue/drop
tokens per-run flag
after/provides
OrcaResult success/skipped/error/fatal/interrupted
OrcaRunResult
validate()
diagnostics
sale de v0:
commit()
replace
timeouts
transaction
compensate
parallel
abortOn/unless avanzados si complican
tokens con payload
fan-in
```
La pregunta decisiva no es “¿está bien diseñado?”. Sí, lo está.
La pregunta correcta es: **¿cuál es el núcleo mínimo que demuestra que `orca` merece existir?**
Respuesta: `event -> staged actions -> tokens -> result trace`.
Si eso funciona, lo demás se puede crecer con calma.
como evitamos la reentrada en orca , si no establecemos en el payload del evento los eventos que le preceden en el stack ? al disparar el componente verifica si esta gestionando ya un evento de ese stack y lo omitiria , como lo suelen hacer el evitar loops infinitos ?
La forma sana no es meter “eventos precedentes” en el payload de negocio. Eso contaminaría todos los eventos con metadata de orquestación y acabaría mezclando dominio con runtime.
Lo correcto es que `orca` envuelva cada evento en un **OrcaEnvelope** interno:
```ts
interface OrcaEnvelope<TPayload = unknown> {
readonly event: string;
readonly payload: TPayload;
readonly meta: {
readonly eventId: string;
readonly parentEventId?: string;
readonly runId?: string;
readonly traceId: string;
readonly depth: number;
readonly stack: readonly string[];
readonly emittedBy?: OrcaActionId;
};
}
```
El payload sigue limpio:
```ts
{ previousActorId, nextActorId }
```
Y `orca` conserva el contexto:
```ts
{
event: APP_EVENT_CACHE_CLEARED,
meta: {
parentEventId: 'evt-1',
traceId: 'trace-7',
depth: 2,
stack: [
APP_EVENT_USER_IDENTITY_CHANGED,
APP_EVENT_CACHE_CLEARED
],
emittedBy: ORCA_ACTION_CLEAR_PRIVATE_CACHE
}
}
```
Entonces la prevención de loops vive en `orca`, no en los módulos.
Reglas típicas para evitar loops infinitos:
```txt
1. No ejecutar inline eventos emitidos durante un run.
2. Encolarlos con traceId y parentEventId.
3. Mantener stack interno por trace.
4. Rechazar o saltar eventos si superan maxDepth.
5. Rechazar o saltar si el mismo evento aparece demasiado en el mismo trace.
6. Aplicar dedupeKey opcional.
7. Tener maxEventsPerTrace.
8. Tener diagnostics cuando se corta el loop.
```
Ejemplo:
```ts
const Orca = createEngineOrca({
bus,
timers,
logger,
reentry: {
maxDepth: 16,
maxEventsPerTrace: 128,
repeatedEventPolicy: ORCA_REENTRY_SKIP,
repeatedEventLimit: 2
}
});
```
Si ocurre esto:
```txt
A -> action publica B
B -> action publica A
A -> action publica B
...
```
`orca` ve el trace:
```txt
A, B, A, B, A
```
Y puede cortar:
```ts
ORCA_RESULT_INTERRUPTED
ORCA_REENTRY_REASON_REPEATED_EVENT
```
o marcar el run como:
```ts
ORCA_RUN_ABORTED
```
Los patrones habituales en otros sistemas:
- Redux middleware suele evitar loops por convención y con guards manuales.
- Redux-Saga usa `takeLatest`, cancelación y control de tasks.
- EventEmitter no te protege: si emites recursivo, te comes el loop.
- State machines evitan loops por transición/estado: si ya estás en cierto estado, el evento no transiciona.
- Sistemas workflow usan `trace/correlationId`, depth limits, idempotency keys y dedupe.
Para `orca`, yo fijaría esto:
```txt
payload = dominio, limpio
envelope/meta = runtime, trace, stack, depth
tokens = coordinación intra-run
bus events = comunicación inter-run
```
Y añadiría al contrato v0:
```ts
interface OrcaReentryOptions {
readonly maxDepth?: number; // default 16
readonly maxEventsPerTrace?: number; // default 128
readonly repeatedEventLimit?: number; // default 1 o 2
readonly repeatedEventPolicy?:
| typeof ORCA_REENTRY_SKIP
| typeof ORCA_REENTRY_ABORT_TRACE
| typeof ORCA_REENTRY_ERROR;
}
```
Así evitas reentrada sin ensuciar los payloads y sin obligar a cada módulo/componente a recordar “quién vino antes”.
cual estrategia es mas arquitectonicamente mejor ?
Arquitectónicamente, la mejor estrategia es:
```txt
payload limpio + envelope runtime interno + cola de orca + límites de reentrada
```
No metería el stack en el payload del evento.
Por qué:
- El payload pertenece al dominio: `userId`, `tenantId`, `orderId`, etc.
- El stack pertenece al runtime de orquestación.
- Si mezclas ambos, todos los módulos empiezan a conocer detalles de `orca`.
- Si mañana cambias cómo trazas reentrada, rompes contratos de eventos.
- Los eventos pueden venir de `connection`, `auth`, `sess`, `http`, etc.; no todos deberían cargar metadata de runtime.
La arquitectura correcta sería:
```txt
buss event:
type + payload limpio
orca internal envelope:
eventId
traceId
parentEventId
runId
depth
stack
emittedByAction
dedupeKey
```
Flujo:
```txt
1. buss publica EVENT_A(payload)
2. orca crea envelope A con traceId
3. action de A publica EVENT_B
4. orca no ejecuta B inline
5. orca crea envelope B con parent=A, mismo traceId, depth+1
6. orca comprueba stack/repeated/depth
7. si es seguro, encola B
8. si detecta loop, corta con diagnostic
```
Esto es más limpio que “cada componente mira si ya está gestionando un evento de ese stack”, porque esa opción reparte lógica de orquestación por toda la app. Es justo lo que `orca` intenta evitar.
Estrategia recomendada para v0:
```ts
createEngineOrca({
bus,
logger,
reentry: {
maxDepth: 16,
maxEventsPerTrace: 128,
repeatedEventLimit: 2,
repeatedEventPolicy: ORCA_REENTRY_SKIP
}
});
```
Reglas:
- Eventos publicados durante un run se encolan, no se ejecutan inline.
- El stack es interno de `orca`.
- Si el mismo evento aparece demasiadas veces en el mismo trace, se salta o aborta.
- Si depth supera límite, aborta el trace.
- Si eventsPerTrace supera límite, aborta el trace.
- Se emite diagnostic estructurado.
- El `OrcaRunResult` registra `traceId`, `parentEventId`, `depth`, `interruptedReason`.
Mi voto: **envelope interno + queue + trace guards**. Es la opción más limpia, menos acoplada y más mantenible.
me refiero a los otros sistemas de referencia
En los sistemas de referencia, la estrategia arquitectónicamente mejor **no es meter el stack en el payload**. Casi todos los sistemas buenos separan:
```txt
payload = datos de negocio
runtime context = control, cancelación, stack, tracing, concurrencia
```
**Comparativa**
| Sistema | Cómo evita loops/reentrada | Lectura para `orca` |
|---|---|---|
| RTK Listener Middleware | Usa runtime API: `signal`, `cancelActiveListeners()`, `cancel()`, `take`, `condition`, `delay`, `pause`, `fork`. No mete stack en la action. | Muy buen referente. `orca` debería pasar `signal`, `traceId`, `runId` y control de cancelación en `OrcaActionContext`. Fuente: [RTK listener middleware](https://redux-toolkit.js.org/api/createListenerMiddleware). |
| Redux-Saga | Usa tareas y cancelación: `takeLatest` cancela la tarea anterior; `takeEvery` permite concurrencia; `race`, `cancel`, `fork`, `join`. | Bueno para políticas `queue/drop/replace/parallel`, pero no copiaría generators. Fuente: [Redux-Saga concurrency](https://redux-saga.js.org/docs/advanced/Concurrency/). |
| XState | Evita loops mediante estado, guards y transiciones. Si ya estás en un estado, un evento puede no transicionar. El control vive en la máquina, no en el payload. | Útil para flows cerrados; para `orca`, copiaría `setup()` tipado y la idea de runtime/state como dueño del control. Fuente: [XState setup](https://stately.ai/docs/setup). |
| Effect Workflow | Distingue error/cause/interruption y tiene compensación estructurada. El runtime sabe qué falló y cómo compensar; no lo mete en payload. | Muy buen referente para `OrcaInterrupted`, compensaciones futuras y `Cause`. Fuente: [Effect Workflow compensation](https://effect-ts.github.io/effect/workflow/Workflow.ts.html). |
| EventEmitter simple | No evita loops. Si un listener emite el mismo evento recursivamente, problema del usuario. | Esto es justo lo que `orca` no debe ser. |
**Veredicto**
La mejor estrategia para `orca` es una mezcla de RTK Listener + Saga + XState:
```txt
payload limpio
OrcaEnvelope interno
cola propia de orca
traceId / parentEventId / runId
AbortSignal en OrcaActionContext
política de concurrencia por evento
límites maxDepth / maxEventsPerTrace
diagnostics cuando corta loops
```
No pondría el stack en el payload. Eso es arquitectura más débil porque obliga a todos los módulos a conocer el runtime de `orca`.
La estrategia superior es:
```ts
interface OrcaActionContext {
readonly runId: OrcaRunId;
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly depth: number;
readonly signal: AbortSignal;
}
```
Y el envelope interno:
```ts
interface OrcaEnvelope<TPayload> {
readonly event: string;
readonly payload: TPayload;
readonly meta: OrcaEventMeta;
}
```
Así `orca` se parece a los sistemas buenos: el **runtime** controla reentrada, cancelación y concurrencia; el **payload** sigue siendo dominio limpio.
ahora mismo no se parece a eso en nada ?
Correcto: **el diseño actual apunta en esa dirección, pero todavía no se parece lo suficiente en contrato concreto**.
Ahora mismo el README ya dice cosas buenas:
```txt
payload limpio
no ejecutar inline
queue/drop/replace
tokens per-run
RunResult
App.Orchestration inerte
```
Pero faltan las piezas que lo harían parecerse de verdad a RTK/Saga/XState/Effect en lo importante:
```txt
OrcaEnvelope interno
OrcaActionContext con traceId/runId/eventId/depth
AbortSignal en acciones
cola interna explícita de orca
políticas de reentrada: maxDepth, maxEventsPerTrace
reason/result para reentry-skip / interrupted
relación formal buss -> orca listener -> queue
```
Así que mi diagnóstico es:
```txt
Conceptualmente: sí está en el camino correcto.
Contractualmente: todavía no está cerrado.
Implementablemente: aún tendría ambigüedad peligrosa.
```
La sección que habría que añadir al README antes de implementar es algo como:
```ts
interface OrcaEnvelope<TPayload = unknown> {
readonly event: string;
readonly payload: TPayload;
readonly meta: OrcaEventMeta;
}
interface OrcaEventMeta {
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly parentRunId?: OrcaRunId;
readonly emittedByAction?: OrcaActionId;
readonly depth: number;
readonly stack: readonly string[];
}
interface OrcaActionContext {
readonly runId: OrcaRunId;
readonly event: string;
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly depth: number;
readonly signal: AbortSignal;
readonly emit: <TPayload>(event: string, payload: TPayload) => void;
}
```
Y reglas:
```txt
1. Las acciones no llaman bus.publish directamente si quieren preservar trace.
2. Usan ctx.emit() para eventos derivados.
3. ctx.emit() crea envelope hijo.
4. Eventos derivados se encolan, no inline.
5. Orca corta loops con maxDepth/maxEventsPerTrace/repeatedEventLimit.
6. Si una acción llama App.Bus directamente, orca lo trata como evento raíz nuevo.
```
Esto último es importante: **si queremos trazabilidad real, las acciones deben emitir eventos derivados vía `ctx.emit()`**, no vía `App.Bus.publish()` directamente. Ahí está la diferencia entre “tenemos bus” y “tenemos runtime de orquestación”.
Así que sí: ahora mismo aún no está del todo. Tiene el esqueleto filosófico, pero le falta el contrato de runtime/envelope/context para ser una pieza robusta.
ctx deberiamos definirlo, seria otra capa mas ?
Sí, `ctx` hay que definirlo. Pero no lo vería como “otra capa” al nivel de `buss`/`orca`/`timr`; lo vería como **el runtime context de una ejecución de `orca`**.
No es un artefacto nuevo. Es parte del contrato de `orca`.
```txt
orca = engine
run = ejecución de un evento
ctx = contexto que orca entrega a cada action dentro de ese run
```
El `ctx` es donde ponemos todo lo que no debe ir en el payload:
```ts
interface OrcaActionContext {
readonly runId: OrcaRunId;
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly depth: number;
readonly signal: AbortSignal;
emit<TPayload>(
event: string,
payload: TPayload,
options?: OrcaEmitOptions
): void;
token(token: OrcaToken): void;
hasToken(token: OrcaToken): boolean;
}
```
La separación queda:
```txt
payload
Datos del evento.
Pertenece al dominio.
ctx
Datos/control de ejecución.
Pertenece a orca.
App modules
Servicios reales.
La action los cierra por closure.
```
Ejemplo:
```ts
App.Orchestration.onEvent(APP_EVENT_USER_IDENTITY_CHANGED, {
id: ORCA_ACTION_CLEAR_PRIVATE_CACHE,
stage: ORCA_STAGE_MAIN,
provides: [ORCA_TOKEN_CACHE_OK, ORCA_TOKEN_CACHE_ERROR],
action: async (payload, ctx) => {
if (ctx.signal.aborted) {
return orcaInterrupted('run aborted');
}
const result = await App.Cache.clearActorScope(payload.previousActorId);
if (!result.ok) {
return orcaError(result.error, {
emits: [ORCA_TOKEN_CACHE_ERROR]
});
}
ctx.emit(APP_EVENT_PRIVATE_CACHE_CLEARED, {
actorId: payload.previousActorId
});
return orcaSuccess({
emits: [ORCA_TOKEN_CACHE_OK]
});
}
});
```
`ctx.emit()` es importante porque preserva traza:
```txt
run A
action clear-cache
ctx.emit(B)
orca crea envelope B:
traceId = A.traceId
parentEventId = A.eventId
emittedByAction = clear-cache
depth = A.depth + 1
```
Si la action llama directamente:
```ts
App.Bus.publish(B, payload)
```
entonces `orca` no sabe que B viene de A salvo que intercepte todo, lo cual es más frágil. Por eso yo pondría esta regla:
```txt
Dentro de acciones orca, usa ctx.emit() para eventos derivados.
Usa App.Bus.publish() solo para eventos raíz externos al run.
```
Entonces no es una capa adicional. Es el equivalente a:
- `listenerApi` en RTK listener middleware.
- `saga effects` en redux-saga.
- `context/event` en XState.
- `workflow context/cause` en Effect/Temporal.
Mi recomendación: sí, definir `OrcaActionContext` como contrato central de v0. Sin `ctx`, `orca` será solo un callback runner. Con `ctx`, empieza a ser runtime de orquestación.
y como sabe una accion de un modulo a la hora de emitir un evento donde lo debe hacer, por ejemplo
auth.onChangeUser (newUser) {
bus.publish(AUTH_CHANGE_EVENT,newUser)
o
ctx.publish(AUTH_CHANGE_EVENT,newUser)
}
La regla tiene que ser muy clara:
```txt
Un módulo nunca usa ctx.
Una acción de orca sí puede usar ctx.
```
Tu ejemplo:
```ts
auth.onChangeUser((newUser) => {
bus.publish(AUTH_CHANGE_EVENT, newUser);
});
```
Eso está bien si ocurre **dentro del módulo auth**. `auth` no debe saber que existe `orca`, por tanto no puede recibir ni usar `ctx`.
En cambio esto:
```ts
auth.onChangeUser((newUser) => {
ctx.publish(AUTH_CHANGE_EVENT, newUser);
});
```
Eso estaría mal si `auth.onChangeUser` es código del módulo. Meterías `orca` dentro de `auth`.
La separación correcta:
```txt
Módulo -> Bus
Orca action -> Ctx
```
Ejemplo completo:
```ts
// Dentro de auth
function setUser(nextUser: User) {
currentUser = nextUser;
bus.publish(AUTH_EVENT_USER_CHANGED, {
userId: nextUser.id
});
}
```
`auth` publica en `bus`. Punto. No conoce `orca`.
Ahora una acción de `orca`:
```ts
App.Orchestration.onEvent(AUTH_EVENT_USER_CHANGED, {
id: ORCA_ACTION_CLEAR_PRIVATE_STATE,
stage: ORCA_STAGE_MAIN,
action: async (payload, ctx) => {
await App.Cache.clearActorScope(payload.userId);
ctx.emit(APP_EVENT_PRIVATE_STATE_CLEARED, {
userId: payload.userId
});
return orcaSuccess({
emits: [ORCA_TOKEN_CACHE_OK]
});
}
});
```
Aquí sí usamos `ctx.emit()` porque estamos dentro de una action y queremos preservar trace.
Entonces hay dos tipos de emisión:
```txt
bus.publish()
Para módulos, componentes, servidor, código externo a orca.
Crea evento raíz desde la perspectiva de orca.
ctx.emit()
Solo dentro de una OrcaAction.
Crea evento hijo dentro del trace actual.
```
Si un módulo emite durante una acción, por ejemplo:
```ts
action: async () => {
await App.Auth.changeUser(nextUser); // Auth internamente hace bus.publish
}
```
`orca` lo verá como evento publicado durante un run, pero no sabrá con tanta precisión qué acción lo emitió salvo que tenga un contexto activo global interno. Hay dos opciones:
**Opción A, simple y limpia v0:**
```txt
Módulo publica por bus.
Orca detecta que hay run activo y lo encola como child event genérico.
parentRunId sí, emittedByAction quizá no.
```
**Opción B, más precisa:**
`orca` ejecuta cada action dentro de un “current action scope” interno. Si durante esa ejecución algún módulo hace `bus.publish`, `orca` intercepta el publish y lo marca como:
```ts
emittedByAction: ORCA_ACTION_CLEAR_PRIVATE_STATE
```
Eso requiere integración fuerte entre `buss` y `orca`.
Mi recomendación:
```txt
v0:
módulos siempre bus.publish
actions pueden ctx.emit
si un módulo publica durante un run, orca lo trata como evento derivado del run si puede detectarlo, pero no exigimos emittedByAction perfecto
v1:
integración bus/orca con current action scope para trazabilidad completa
```
Regla para desarrolladores:
```txt
Si estás escribiendo código de módulo: usa bus.publish.
Si estás escribiendo una OrcaAction: usa ctx.emit para eventos derivados del flujo.
```
Así evitamos acoplar módulos a `orca`, pero damos a las acciones una vía correcta para preservar trazabilidad.

@ -24,17 +24,40 @@ que el acoplamiento inter-modulo quede escondido en `buss`, `connection` o
## Estado Del Documento
Este README es un boceto preliminar de diseno. Todavia no describe una API
implementada. Sirve para cerrar el contrato antes de escribir codigo.
Las decisiones aqui son intencionadas:
- `buss` transporta eventos locales.
Este README es la referencia de diseño de orca. La **v0-kernel** está
implementada y testeada (114 archivos / 1357 tests pasan, 47 de ellos
sobre orca). El kernel expone:
- `createEngineOrca({ bus, timers, logger?, maxRuns?, reentry? })`
- `onEvent(event, action) → detach`
- Stages canónicos (`guard / pre / main / post / cleanup / finally`)
- `OrcaEnvelope` interno + `OrcaActionContext` con `eventId`, `traceId`,
`parentEventId`, `depth`, `signal`, `emit()`, `tokens`
- `OrcaResult`: `success / skipped / error / interrupted` (`fatal`,
`timeout` aceptados en el tipo, no producidos por v0)
- `OrcaRunResult` con `eventId`, `traceId`, `parentEventId`, `depth`
- Reentry guards: `maxDepth`, `maxEventsPerTrace`,
`repeatedEventLimit` con políticas `skip / abort-trace / error`
- `dedupeKey` por trace
- Diagnostics estructurados (`orca.run.*`, `orca.action.*`,
`orca.event.emitted`, `orca.reentry.blocked`, `orca.trace.aborted`)
- `applyStandardOrca(App)` agregador de presets en
`arts/active-app/presets/`
Lo que sigue siendo **v1** (aceptado en el tipo, no honrado por el
motor) está marcado con `@v1+` en el código fuente. Listado en la
sección "Roadmap v1" al final.
Las decisiones arquitectónicas que sostienen el diseño:
- `bus` transporta eventos locales.
- `orca` ejecuta acciones asociadas a eventos del bus.
- `timr` coordina timers, timeouts y clocks.
- `logr` recibe diagnostics/logs.
- La aplicacion decide que modulos toca cada accion.
- `timer` coordina timers, timeouts y clocks.
- `logger` recibe diagnostics/logs.
- La aplicación decide qué módulos toca cada acción.
- Nada destructivo se ejecuta por defecto sin estar registrado.
- Payload limpio, envelope runtime separado: el módulo nunca conoce a
`orca`; sólo las `OrcaAction` reciben `ctx`.
## Naming
@ -504,29 +527,134 @@ valida y congela la configuracion para ejecucion. En desarrollo, `orca` puede
ejecutar `validate()` automaticamente antes del primer evento si el usuario no
lo hizo.
## Reentrada Y Eventos Durante Un Run
## OrcaEnvelope y OrcaActionContext
Cada evento que el motor procesa se envuelve en un `OrcaEnvelope`
interno antes de llegar a la cola de runs:
```ts
interface OrcaEnvelope<TPayload> {
readonly event: string;
readonly payload: TPayload;
readonly meta: OrcaEventMeta;
}
interface OrcaEventMeta {
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly parentRunId?: OrcaRunId;
readonly emittedByAction?: OrcaActionId;
readonly depth: number;
readonly stack: readonly string[];
readonly publishedAt: number;
readonly dedupeKey?: string;
}
```
El `payload` permanece limpio (`{ userId, tenantId, ... }`). El `meta`
es propiedad del runtime de orca.
Cada `OrcaAction` recibe un `OrcaActionContext` con la metadata útil
para esa ejecución:
```ts
interface OrcaActionContext {
readonly runId: OrcaRunId;
readonly event: string;
readonly stage: OrcaStage;
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly depth: number;
readonly tokens: ReadonlySet<OrcaToken>;
readonly signal: AbortSignal;
readonly logger: Logger;
emit<TPayload>(
event: string,
payload: TPayload,
options?: OrcaEmitOptions
): OrcaEventId | null;
}
```
Una accion puede llamar a un modulo que publique nuevos eventos en `buss`.
Ejemplo: `App.Cache.clear()` puede publicar `APP_EVENT_CACHE_CLEARED`.
## Regla módulo → bus, action → ctx.emit
Regla v0:
La **frontera de responsabilidad** queda nítida:
```txt
eventos publicados durante un run no se ejecutan inline
Módulo → bus.publish() (nunca conoce orca)
Action → ctx.emit() (preserva traceId, parentEventId, depth)
```
No deben ejecutarse dentro del mismo stack ni mezclarse con el stage actual. Se
encolan y se procesan despues de que el run actual alcance un punto seguro,
respetando FIFO del bus y la politica de concurrencia del evento destino.
Un módulo nunca recibe `ctx`. Si una `OrcaAction` quiere emitir un
evento derivado durante un run, usa `ctx.emit()` para que orca cree un
envelope hijo (mismo `traceId`, `parentEventId = ctx.eventId`,
`depth = ctx.depth + 1`, `emittedByAction` = id de la acción).
Esto evita reentrada confusa:
Si en su lugar la acción llama a un módulo que internamente hace
`bus.publish(...)`, orca lo verá como evento durante el run pero como
**raíz** desde la perspectiva de orca: nuevo `traceId`, `depth = 0`. No
se mezcla con la traza del run actual. Esa es la **opción A** de
v0-kernel; la opción B (interceptar `bus.publish` durante un run para
atribuir `emittedByAction` perfecto) queda pospuesta a v1.
## Reentrada Y Eventos Durante Un Run
Cualquier evento que un módulo publique en el bus durante un run se
encola, no se ejecuta inline. Lo mismo aplica a los eventos derivados
vía `ctx.emit()`. El motor mantiene una cola FIFO por evento y nunca
ejecuta acciones recursivamente dentro del mismo stack.
```txt
main(action A)
-> publica event B
-> ctx.emit(B) -> child envelope B encolado
-> NO ejecuta pipeline B dentro de action A
```
### Reentry guards
`createEngineOrca` acepta `reentry: OrcaReentryOptions`:
```ts
interface OrcaReentryOptions {
readonly maxDepth?: number; // default 16
readonly maxEventsPerTrace?: number; // default 128
readonly repeatedEventLimit?: number; // default 2
readonly repeatedEventPolicy?:
| typeof ORCA_REENTRY_SKIP // default
| typeof ORCA_REENTRY_ABORT_TRACE
| typeof ORCA_REENTRY_ERROR;
}
```
Cuando un envelope **derivado** (vía `ctx.emit()`) supera uno de los
límites, el motor aplica la política configurada:
- `skip` — el envelope no se encola; la traza continúa para otros
eventos. Diagnostic `orca.reentry.blocked` con la `reason` adecuada.
- `abort-trace` — la traza queda marcada como abortada; cualquier
evento ya encolado para esa traza se descarta cuando le toca turno;
los runs en vuelo de esa traza terminan, pero los siguientes
registran `ORCA_RUN_INTERRUPTED`. Diagnostic `orca.trace.aborted`.
- `error` — se lanza `OrcaReentryError` síncronamente desde
`ctx.emit()`. La acción puede capturarlo o dejarlo escalar.
Razones (`OrcaReentryReason`):
- `max-depth` — `depth > maxDepth`
- `max-events-per-trace` — `eventCount + 1 > maxEventsPerTrace`
- `repeated-event` — un mismo nombre supera `repeatedEventLimit`
- `deduped` — el `dedupeKey` ya existe en la traza
- `trace-aborted` — la traza fue marcada como abortada antes
- `run-aborted` — el run fue abortado por `ORCA_ON_ERROR_ABORT_RUN`
- `disposed` — el motor está disposed
**Importante**: los publishes directos en el bus desde código de módulo
NO entran en estos contadores. Cada `bus.publish(event, payload)`
genera un envelope raíz con `traceId` fresco y `depth = 0`. Los
contadores nacen y mueren con cada traza.
`buss` conserva la entrega local. `orca` decide cuando consume y arranca runs.
Si un evento nuevo requiere ejecucion inmediata, debe modelarse como token del
run actual, no como evento reentrante.
@ -1169,18 +1297,39 @@ Tests compuestos con ecosistema:
## Invariantes
- `orca` no importa artefactos concretos salvo contratos comunes.
- `orca` no conoce modulos de negocio.
- los artefactos no consumen `orca`; solo publican eventos en `buss`.
- la aplicacion registra acciones en `orca`.
- `App.Orchestration` existe en `aapp`, pero no ejecuta nada sin acciones.
- `App.Bus` debe existir si existe `App.Orchestration`.
- Todas las strings publicas viven en constantes.
- `orca` no conoce módulos de negocio.
- Los artefactos no consumen `orca`; solo publican eventos en `bus`.
- La aplicación registra acciones en `orca`.
- `App.Orca` existe siempre, pero no ejecuta nada sin acciones.
- `App.Bus` debe existir si existe `App.Orca`.
- Todas las strings públicas viven en constantes.
- Los eventos son constantes, no strings inline.
- Los tokens son constantes, no strings inline.
- `setupOrca()` no conoce modulos, solo contratos de eventos/tokens/acciones.
- Los logs usan `Logger` comun.
- Los logs usan `Logger` común.
- Los diagnostics son catalogados.
- Los timers los ejecuta `timr`.
- El resultado de una accion siempre queda representado.
- Los timers los ejecuta `timer`.
- El resultado de una acción siempre queda representado.
- Los payloads de eventos no contienen credenciales.
- Las acciones destructivas son explicitas.
- Las acciones destructivas son explícitas.
- Los módulos nunca reciben `ctx`. Sólo las `OrcaAction` lo reciben, y
`ctx.emit()` es el camino que preserva trazabilidad.
## Roadmap v1
Aceptado en el contrato público, no honrado por el motor todavía.
Marcado como `@v1+` en el código fuente:
| Pieza | Qué falta para v1 |
|---|---|
| `actionTimeoutMs` | Ejecutar la acción con timeout real (`signal.abort()` por timer); producir `ORCA_RESULT_TIMEOUT`. |
| `compensate` | Invocar la función compensatoria cuando el run aborta tras un éxito previo. |
| `after` / `unless` / `abortOn` | Honrar los gates de tokens entre acciones. |
| `OrcaFatal` | Distinguir fatal de error en políticas downstream. |
| `commit()` / `replace` | Políticas de cola por evento (drop-prev / replace / parallel). |
| `transaction` (atómico) | Grupos atómicos cuyo fallo lanza compensaciones en orden inverso. |
| `parallel` | Ejecución concurrente dentro de un mismo stage cuando no hay `after`. |
| 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. |
| `validate()` / `commit()` configuración | Validación estática del grafo (ciclos, tokens imposibles, IDs duplicados). |
| `createActiveOrca()` | Wrapper reactivo (`$state` para `running`, `recentRuns`, etc.) — el `EngineOrca` ya está. |

@ -31,6 +31,12 @@ export const ORCA_STAGES_CANONICAL_ORDER = [
export const ORCA_RESULT_SUCCESS = 'success' as const;
export const ORCA_RESULT_SKIPPED = 'skipped' as const;
export const ORCA_RESULT_ERROR = 'error' as const;
/**
* The action was prevented from completing by a runtime guard (reentry
* limit, trace abort) or aborted via its `signal`. Produced by the v0.0
* engine when the run/trace is cut short.
*/
export const ORCA_RESULT_INTERRUPTED = 'interrupted' as const;
/** @v0.1+ Accepted in the Result type, never produced by the v0.0 engine. */
export const ORCA_RESULT_TIMEOUT = 'timeout' as const;
@ -43,6 +49,7 @@ export const ORCA_ACTION_STATUS_SUCCESS = 'success' as const;
export const ORCA_ACTION_STATUS_SKIPPED = 'skipped' as const;
export const ORCA_ACTION_STATUS_BLOCKED = 'blocked' as const;
export const ORCA_ACTION_STATUS_ERROR = 'error' as const;
export const ORCA_ACTION_STATUS_INTERRUPTED = 'interrupted' as const;
export const ORCA_ACTION_STATUS_TIMEOUT = 'timeout' as const;
export const ORCA_ACTION_STATUS_FATAL = 'fatal' as const;
@ -51,9 +58,37 @@ export const ORCA_ACTION_STATUS_FATAL = 'fatal' as const;
export const ORCA_RUN_SUCCESS = 'success' as const;
export const ORCA_RUN_PARTIAL = 'partial' as const;
export const ORCA_RUN_ABORTED = 'aborted' as const;
export const ORCA_RUN_INTERRUPTED = 'interrupted' as const;
export const ORCA_RUN_FATAL = 'fatal' as const;
export const ORCA_RUN_TIMEOUT = 'timeout' as const;
// ── Reentry guards ──────────────────────────────────────────────────────
//
// Per-trace bookkeeping prevents `ctx.emit()` chains from cycling
// indefinitely. Modules that publish on the bus directly do not enter
// these counters because they create root events (depth=0, fresh
// traceId) — only events emitted through `ctx.emit()` count against
// their parent's trace.
/**
* Skip the offending event (it is not enqueued; an interrupted action
* appears in the trace with a reason). The trace continues for other
* events.
*/
export const ORCA_REENTRY_SKIP = 'skip' as const;
/** Abort every pending event for the trace and stop processing it. */
export const ORCA_REENTRY_ABORT_TRACE = 'abort-trace' as const;
/** Throw `OrcaReentryError` synchronously when the guard fires. */
export const ORCA_REENTRY_ERROR = 'error' as const;
export const ORCA_REENTRY_REASON_MAX_DEPTH = 'max-depth' as const;
export const ORCA_REENTRY_REASON_MAX_EVENTS_PER_TRACE = 'max-events-per-trace' as const;
export const ORCA_REENTRY_REASON_REPEATED_EVENT = 'repeated-event' as const;
export const ORCA_REENTRY_REASON_DEDUPED = 'deduped' as const;
export const ORCA_REENTRY_REASON_TRACE_ABORTED = 'trace-aborted' as const;
export const ORCA_REENTRY_REASON_RUN_ABORTED = 'run-aborted' as const;
export const ORCA_REENTRY_REASON_DISPOSED = 'disposed' as const;
// ── Error policies ──────────────────────────────────────────────────────
// v0.0 implements CONTINUE and ABORT_RUN. The remaining policies are
// accepted in the type and treated as CONTINUE.
@ -75,6 +110,10 @@ export const ORCA_DIAGNOSTIC_EVENTS = {
ACTION_COMPLETED: 'orca.action.completed',
ACTION_FAILED: 'orca.action.failed',
ACTION_SKIPPED: 'orca.action.skipped',
ACTION_INTERRUPTED: 'orca.action.interrupted',
EVENT_EMITTED: 'orca.event.emitted',
REENTRY_BLOCKED: 'orca.reentry.blocked',
TRACE_ABORTED: 'orca.trace.aborted',
CONFIGURATION_INVALID: 'orca.configuration.invalid'
} as const;
@ -87,8 +126,15 @@ export const ORCA_LOG_MSG_ACTION_STARTED = 'orca action started';
export const ORCA_LOG_MSG_ACTION_COMPLETED = 'orca action completed';
export const ORCA_LOG_MSG_ACTION_FAILED = 'orca action failed';
export const ORCA_LOG_MSG_ACTION_SKIPPED = 'orca action skipped';
export const ORCA_LOG_MSG_ACTION_INTERRUPTED = 'orca action interrupted';
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';
export const ORCA_LOG_MSG_CONFIGURATION_INVALID = 'orca configuration invalid';
// ── Defaults ────────────────────────────────────────────────────────────
export const ORCA_DEFAULT_MAX_RUNS = 256;
export const ORCA_DEFAULT_MAX_DEPTH = 16;
export const ORCA_DEFAULT_MAX_EVENTS_PER_TRACE = 128;
export const ORCA_DEFAULT_REPEATED_EVENT_LIMIT = 2;

@ -10,56 +10,103 @@ import {
ORCA_DIAGNOSTIC_EVENTS,
ORCA_LOG_MSG_ACTION_COMPLETED,
ORCA_LOG_MSG_ACTION_FAILED,
ORCA_LOG_MSG_ACTION_INTERRUPTED,
ORCA_LOG_MSG_ACTION_SKIPPED,
ORCA_LOG_MSG_ACTION_STARTED,
ORCA_LOG_MSG_CONFIGURATION_INVALID,
ORCA_LOG_MSG_EVENT_EMITTED,
ORCA_LOG_MSG_REENTRY_BLOCKED,
ORCA_LOG_MSG_RUN_ABORTED,
ORCA_LOG_MSG_RUN_COMPLETED,
ORCA_LOG_MSG_RUN_STARTED,
ORCA_LOG_MSG_TRACE_ABORTED,
ORCA_MODULE
} from './consts.ts';
import type { OrcaActionId, OrcaRunId, OrcaStage } from './types.ts';
import type {
OrcaActionId,
OrcaEventId,
OrcaReentryReason,
OrcaRunId,
OrcaStage,
OrcaTraceId
} from './types.ts';
export type OrcaDiagnosticType =
(typeof ORCA_DIAGNOSTIC_EVENTS)[keyof typeof ORCA_DIAGNOSTIC_EVENTS];
/**
* Common envelope identity carried on every run-scoped diagnostic.
* Engines that consume diagnostics can correlate run/event/trace
* without parsing the variant tags.
*/
interface OrcaDiagnosticEnvelopeRef {
readonly runId?: OrcaRunId;
readonly eventId?: OrcaEventId;
readonly traceId?: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly depth?: number;
}
export type OrcaDiagnosticMeta =
| {
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly event: string;
}
| {
})
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly event: string;
readonly status: string;
readonly durationMs: number;
readonly actionCount: number;
}
| {
})
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly event: string;
readonly cause: unknown;
}
| {
})
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly actionId: OrcaActionId;
readonly stage: OrcaStage;
}
| {
})
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly actionId: OrcaActionId;
readonly durationMs: number;
}
| {
})
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly actionId: OrcaActionId;
readonly error: unknown;
readonly thrown?: boolean;
}
| {
})
| (OrcaDiagnosticEnvelopeRef & {
readonly runId: OrcaRunId;
readonly actionId: OrcaActionId;
readonly reason?: string;
})
| {
/** ctx.emit() of a derived event — the run hasn't started yet. */
readonly parentRunId: OrcaRunId;
readonly parentEventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly event: string;
readonly eventId: OrcaEventId;
readonly depth: number;
readonly emittedByAction: OrcaActionId;
}
| {
/** Reentry guard fired for an event being enqueued. */
readonly traceId: OrcaTraceId;
readonly event: string;
readonly parentEventId?: OrcaEventId;
readonly depth: number;
readonly reason: OrcaReentryReason;
}
| {
/** Trace marked aborted; queued events for that trace will be cut. */
readonly traceId: OrcaTraceId;
readonly reason: OrcaReentryReason;
}
| {
readonly reason: string;
@ -94,6 +141,22 @@ const ORCA_DIAGNOSTIC_LOGS: DiagnosticCatalog<OrcaDiagnosticEvent> = {
level: LogLevel.DEBUG,
message: ORCA_LOG_MSG_ACTION_SKIPPED
},
[ORCA_DIAGNOSTIC_EVENTS.ACTION_INTERRUPTED]: {
level: LogLevel.WARN,
message: ORCA_LOG_MSG_ACTION_INTERRUPTED
},
[ORCA_DIAGNOSTIC_EVENTS.EVENT_EMITTED]: {
level: LogLevel.TRACE,
message: ORCA_LOG_MSG_EVENT_EMITTED
},
[ORCA_DIAGNOSTIC_EVENTS.REENTRY_BLOCKED]: {
level: LogLevel.WARN,
message: ORCA_LOG_MSG_REENTRY_BLOCKED
},
[ORCA_DIAGNOSTIC_EVENTS.TRACE_ABORTED]: {
level: LogLevel.WARN,
message: ORCA_LOG_MSG_TRACE_ABORTED
},
[ORCA_DIAGNOSTIC_EVENTS.CONFIGURATION_INVALID]: {
level: LogLevel.ERROR,
message: ORCA_LOG_MSG_CONFIGURATION_INVALID

@ -1,36 +1,62 @@
/**
* `EngineOrca` v0.0 — orchestration engine.
* `EngineOrca` v0 — orchestration kernel.
*
* Surface complete, engine minimal:
* - All public types and fields are accepted at registration.
* - Honored at runtime: stages (canonical order), `onError` (CONTINUE
* vs ABORT_RUN), exception capture, run trace, run-queue per event,
* dispose lifecycle, lazy bus subscription per event.
* - Accepted but ignored at runtime: `after`, `unless`, `abortOn`,
* `actionTimeoutMs`, `compensate`. The engine populates the token
* set in the action context so authors can read it, but does not
* act on it.
* Surface complete, motor focused on the kernel 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.
* - 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.
*
* The "queue per event" concurrency is implicit: events of the same
* type that arrive while a run is in flight queue FIFO and are
* processed when the active run completes. There is no concurrent
* execution.
* 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.
*/
import {
ORCA_ACTION_STATUS_BLOCKED,
ORCA_ACTION_STATUS_ERROR,
ORCA_ACTION_STATUS_INTERRUPTED,
ORCA_ACTION_STATUS_SKIPPED,
ORCA_ACTION_STATUS_SUCCESS,
ORCA_DEFAULT_MAX_DEPTH,
ORCA_DEFAULT_MAX_EVENTS_PER_TRACE,
ORCA_DEFAULT_MAX_RUNS,
ORCA_DEFAULT_REPEATED_EVENT_LIMIT,
ORCA_DIAGNOSTIC_EVENTS,
ORCA_ON_ERROR_ABORT_RUN,
ORCA_ON_ERROR_CONTINUE,
ORCA_REENTRY_ABORT_TRACE,
ORCA_REENTRY_ERROR,
ORCA_REENTRY_REASON_DEDUPED,
ORCA_REENTRY_REASON_DISPOSED,
ORCA_REENTRY_REASON_MAX_DEPTH,
ORCA_REENTRY_REASON_MAX_EVENTS_PER_TRACE,
ORCA_REENTRY_REASON_REPEATED_EVENT,
ORCA_REENTRY_REASON_RUN_ABORTED,
ORCA_REENTRY_REASON_TRACE_ABORTED,
ORCA_REENTRY_SKIP,
ORCA_RESULT_ERROR,
ORCA_RESULT_FATAL,
ORCA_RESULT_INTERRUPTED,
ORCA_RESULT_SKIPPED,
ORCA_RESULT_SUCCESS,
ORCA_RUN_ABORTED,
ORCA_RUN_INTERRUPTED,
ORCA_RUN_PARTIAL,
ORCA_RUN_SUCCESS,
ORCA_STAGE_FINALLY,
@ -40,7 +66,8 @@ import {
OrcaDisposedError,
OrcaDuplicateActionIdError,
OrcaInvalidActionError,
OrcaInvalidStageError
OrcaInvalidStageError,
OrcaReentryError
} from './errors.ts';
import { createOrcaDiagnostics, emitOrcaDiagnostic } from './diagnostics.ts';
import type {
@ -48,12 +75,20 @@ import type {
EngineOrcaOptions,
OrcaAction,
OrcaActionContext,
OrcaActionId,
OrcaActionRun,
OrcaEmitOptions,
OrcaEnvelope,
OrcaEventId,
OrcaReentryOptions,
OrcaReentryPolicy,
OrcaReentryReason,
OrcaResult,
OrcaRunId,
OrcaRunResult,
OrcaStage,
OrcaToken
OrcaToken,
OrcaTraceId
} from './types.ts';
import type { Logger } from '$libs/logger';
import type { BusSubscription } from '$libs/bus';
@ -63,16 +98,38 @@ interface RegisteredAction {
readonly registeredAt: number;
}
interface QueuedRun {
readonly envelope: OrcaEnvelope;
readonly actions: readonly OrcaAction[];
}
interface TraceState {
eventCount: number;
readonly eventCounts: Map<string, number>;
readonly dedupeKeys: Set<string>;
aborted: boolean;
abortedReason?: OrcaReentryReason;
}
interface ResolvedReentry {
readonly maxDepth: number;
readonly maxEventsPerTrace: number;
readonly repeatedEventLimit: number;
readonly repeatedEventPolicy: OrcaReentryPolicy;
}
export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
const { bus, timers } = options;
const maxRuns = options.maxRuns ?? ORCA_DEFAULT_MAX_RUNS;
const reentry = resolveReentry(options.reentry);
const diagnostics = createOrcaDiagnostics(options.logger);
const actionLogger: Logger = options.logger ?? createNoopLogger();
const actionsByEvent = new Map<string, RegisteredAction[]>();
const busSubscriptions = new Map<string, BusSubscription>();
const recentRuns: OrcaRunResult[] = [];
const runQueue: Array<() => Promise<void>> = [];
const runQueue: QueuedRun[] = [];
const traceStates = new Map<OrcaTraceId, TraceState>();
let disposed = false;
let running = false;
@ -117,8 +174,12 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
// Lazy bus subscription: only when the first action for this event
// registers. A single subscription per event services all actions.
if (!busSubscriptions.has(event)) {
const subscription = bus.on(event, (envelope) => {
enqueueRun(event, envelope.payload);
const subscription = bus.on(event, (busEnvelope) => {
// Bus events are roots from orca's perspective: fresh
// traceId, depth = 0, no parent. Whatever counters the
// trace accumulates start here.
const envelope = buildRootEnvelope(event, busEnvelope.payload);
enqueueEnvelope(envelope);
});
busSubscriptions.set(event, subscription);
}
@ -139,27 +200,174 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
}
}
function enqueueRun(event: string, payload: unknown): void {
if (disposed) return;
const registered = actionsByEvent.get(event);
if (!registered || registered.length === 0) return;
function buildRootEnvelope(event: string, payload: unknown): OrcaEnvelope {
return {
event,
payload,
meta: {
eventId: generateEventId(),
traceId: generateTraceId(),
depth: 0,
stack: [event],
publishedAt: timers.clock.now()
}
};
}
// Snapshot taken at dispatch. Actions registered during the run
// do not participate in it.
const snapshot = registered.map((entry) => entry.action);
function buildChildEnvelope(
event: string,
payload: unknown,
parent: OrcaEnvelope,
emittedByAction: OrcaActionId,
parentRunId: OrcaRunId,
dedupeKey?: string
): OrcaEnvelope {
return {
event,
payload,
meta: {
eventId: generateEventId(),
traceId: parent.meta.traceId,
parentEventId: parent.meta.eventId,
parentRunId,
emittedByAction,
depth: parent.meta.depth + 1,
stack: [...parent.meta.stack, event],
publishedAt: timers.clock.now(),
dedupeKey
}
};
}
/**
* Reentry-guarded enqueue. Applies maxDepth/maxEventsPerTrace/
* repeatedEventLimit/dedupe checks before adding the envelope to the
* run queue. Returns `null` when the envelope is rejected by SKIP
* or ABORT_TRACE; throws when policy is ERROR.
*/
function enqueueEnvelope(envelope: OrcaEnvelope): OrcaEventId | null {
if (disposed) {
emitReentryBlocked(envelope, ORCA_REENTRY_REASON_DISPOSED);
return null;
}
const registered = actionsByEvent.get(envelope.event);
if (!registered || registered.length === 0) {
// No actions registered for this event — silently drop. Note
// that we still consider the envelope as observed for trace
// counting purposes only when it has a parent (depth > 0);
// orphan derived events would otherwise allow infinite
// emission of unsubscribed names.
if (envelope.meta.depth > 0) accountTraceEntry(envelope);
return envelope.meta.eventId;
}
runQueue.push(() => executeRun(event, payload, snapshot));
// Trace state lookup. Root envelopes always create a fresh entry.
const trace = ensureTraceState(envelope.meta.traceId);
if (trace.aborted) {
emitReentryBlocked(envelope, trace.abortedReason ?? ORCA_REENTRY_REASON_TRACE_ABORTED);
return null;
}
if (envelope.meta.depth > 0) {
const denial = applyReentryGuards(envelope, trace);
if (denial !== null) {
emitReentryBlocked(envelope, denial);
if (reentry.repeatedEventPolicy === ORCA_REENTRY_ABORT_TRACE) {
trace.aborted = true;
trace.abortedReason = denial;
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.TRACE_ABORTED, {
traceId: envelope.meta.traceId,
reason: denial
});
}
if (reentry.repeatedEventPolicy === ORCA_REENTRY_ERROR) {
throw new OrcaReentryError(denial, envelope.meta.traceId, envelope.event);
}
return null;
}
}
// Accepted: account counters and queue. The snapshot of actions
// is taken at dispatch time so registrations during the run do
// not affect this run.
accountTraceEntry(envelope);
runQueue.push({
envelope,
actions: registered.map((entry) => entry.action)
});
void drainQueue();
return envelope.meta.eventId;
}
function ensureTraceState(traceId: OrcaTraceId): TraceState {
let trace = traceStates.get(traceId);
if (!trace) {
trace = {
eventCount: 0,
eventCounts: new Map(),
dedupeKeys: new Set(),
aborted: false
};
traceStates.set(traceId, trace);
}
return trace;
}
function accountTraceEntry(envelope: OrcaEnvelope): void {
const trace = ensureTraceState(envelope.meta.traceId);
trace.eventCount += 1;
trace.eventCounts.set(
envelope.event,
(trace.eventCounts.get(envelope.event) ?? 0) + 1
);
if (envelope.meta.dedupeKey) trace.dedupeKeys.add(envelope.meta.dedupeKey);
}
function applyReentryGuards(
envelope: OrcaEnvelope,
trace: TraceState
): OrcaReentryReason | null {
if (envelope.meta.depth > reentry.maxDepth) return ORCA_REENTRY_REASON_MAX_DEPTH;
if (trace.eventCount + 1 > reentry.maxEventsPerTrace) {
return ORCA_REENTRY_REASON_MAX_EVENTS_PER_TRACE;
}
const repeats = (trace.eventCounts.get(envelope.event) ?? 0) + 1;
if (repeats > reentry.repeatedEventLimit) return ORCA_REENTRY_REASON_REPEATED_EVENT;
if (envelope.meta.dedupeKey && trace.dedupeKeys.has(envelope.meta.dedupeKey)) {
return ORCA_REENTRY_REASON_DEDUPED;
}
return null;
}
function emitReentryBlocked(envelope: OrcaEnvelope, reason: OrcaReentryReason): void {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.REENTRY_BLOCKED, {
traceId: envelope.meta.traceId,
event: envelope.event,
parentEventId: envelope.meta.parentEventId,
depth: envelope.meta.depth,
reason
});
}
async function drainQueue(): Promise<void> {
if (running || disposed) return;
// Skip head-of-queue entries whose trace has been aborted while
// queued. They do not produce a run trace; the abort was already
// recorded on the trace itself.
while (runQueue.length > 0) {
const next = runQueue[0];
const trace = traceStates.get(next.envelope.meta.traceId);
if (!trace || !trace.aborted) break;
runQueue.shift();
}
const next = runQueue.shift();
if (!next) return;
running = true;
try {
await next();
await executeRun(next.envelope, next.actions);
} finally {
running = false;
if (!disposed && runQueue.length > 0) {
@ -171,9 +379,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
}
async function executeRun(
event: string,
payload: unknown,
actions: OrcaAction[]
envelope: OrcaEnvelope,
actions: readonly OrcaAction[]
): Promise<void> {
const runId = generateRunId();
const startedAt = timers.clock.now();
@ -184,7 +391,11 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_STARTED, {
runId,
event
event: envelope.event,
eventId: envelope.meta.eventId,
traceId: envelope.meta.traceId,
parentEventId: envelope.meta.parentEventId,
depth: envelope.meta.depth
});
let aborted = false;
@ -201,15 +412,28 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
if (controller.signal.aborted && stage !== ORCA_STAGE_FINALLY) break;
if (disposed) break stageLoop;
const actionRun = await runAction(action, payload, {
const trace = traceStates.get(envelope.meta.traceId);
if (trace?.aborted && stage !== ORCA_STAGE_FINALLY) {
actionRuns.push(
interruptedActionRun(
action,
timers.clock.now(),
trace.abortedReason ?? ORCA_REENTRY_REASON_TRACE_ABORTED
)
);
continue;
}
const context = buildActionContext(
runId,
event,
stage,
envelope,
action.stage,
tokens,
signal: controller.signal,
logger: actionLogger
});
controller.signal,
action.id
);
const actionRun = await runAction(action, envelope.payload, context);
actionRuns.push(actionRun);
if (actionRun.status === ORCA_ACTION_STATUS_ERROR) {
@ -219,7 +443,10 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
controller.abort();
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_ABORTED, {
runId,
event,
event: envelope.event,
eventId: envelope.meta.eventId,
traceId: envelope.meta.traceId,
depth: envelope.meta.depth,
cause: actionRun.error
});
break;
@ -229,11 +456,15 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
}
const endedAt = timers.clock.now();
const status = computeRunStatus(aborted, actionRuns);
const status = computeRunStatus(aborted, actionRuns, envelope, traceStates);
const runResult: OrcaRunResult = {
id: runId,
event,
event: envelope.event,
eventId: envelope.meta.eventId,
traceId: envelope.meta.traceId,
parentEventId: envelope.meta.parentEventId,
depth: envelope.meta.depth,
status,
startedAt,
endedAt,
@ -247,11 +478,77 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_COMPLETED, {
runId,
event,
event: envelope.event,
eventId: envelope.meta.eventId,
traceId: envelope.meta.traceId,
depth: envelope.meta.depth,
status,
durationMs: runResult.durationMs,
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);
}
function buildActionContext(
runId: OrcaRunId,
envelope: OrcaEnvelope,
stage: OrcaStage,
tokens: Set<OrcaToken>,
signal: AbortSignal,
actionId: OrcaActionId
): OrcaActionContext {
return {
runId,
event: envelope.event,
stage,
eventId: envelope.meta.eventId,
traceId: envelope.meta.traceId,
parentEventId: envelope.meta.parentEventId,
depth: envelope.meta.depth,
tokens,
signal,
logger: actionLogger,
emit<TPayload>(
event: string,
payload: TPayload,
options: OrcaEmitOptions = {}
): OrcaEventId | null {
if (disposed) return null;
const child = buildChildEnvelope(
event,
payload,
envelope,
actionId,
runId,
options.dedupeKey
);
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.EVENT_EMITTED, {
parentRunId: runId,
parentEventId: envelope.meta.eventId,
traceId: child.meta.traceId,
event: child.event,
eventId: child.meta.eventId,
depth: child.meta.depth,
emittedByAction: actionId
});
return enqueueEnvelope(child);
}
};
}
function maybeReleaseTrace(traceId: OrcaTraceId): void {
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);
}
async function runAction(
@ -264,7 +561,10 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_STARTED, {
runId: context.runId,
actionId: action.id,
stage: action.stage
stage: action.stage,
eventId: context.eventId,
traceId: context.traceId,
depth: context.depth
});
// If the run is already aborted (e.g. between stages, only finally
@ -293,18 +593,36 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_FAILED, {
runId: context.runId,
actionId: action.id,
eventId: context.eventId,
traceId: context.traceId,
depth: context.depth,
error: (result as { error?: unknown }).error
});
} else if (status === ORCA_ACTION_STATUS_SKIPPED) {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_SKIPPED, {
runId: context.runId,
actionId: action.id,
eventId: context.eventId,
traceId: context.traceId,
depth: context.depth,
reason: (result as { reason?: string }).reason
});
} else if (status === ORCA_ACTION_STATUS_INTERRUPTED) {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_INTERRUPTED, {
runId: context.runId,
actionId: action.id,
eventId: context.eventId,
traceId: context.traceId,
depth: context.depth,
reason: (result as { reason?: string }).reason
});
} else {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_COMPLETED, {
runId: context.runId,
actionId: action.id,
eventId: context.eventId,
traceId: context.traceId,
depth: context.depth,
durationMs: endedAt - startedAt
});
}
@ -320,6 +638,10 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
error:
status === ORCA_ACTION_STATUS_ERROR
? (result as { error?: unknown }).error
: undefined,
interruptedReason:
status === ORCA_ACTION_STATUS_INTERRUPTED
? (result as { reason?: string }).reason
: undefined
};
} catch (thrown) {
@ -328,6 +650,9 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_FAILED, {
runId: context.runId,
actionId: action.id,
eventId: context.eventId,
traceId: context.traceId,
depth: context.depth,
error: thrown,
thrown: true
});
@ -357,6 +682,12 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
busSubscriptions.clear();
actionsByEvent.clear();
runQueue.length = 0;
// Mark every live trace aborted so any late `ctx.emit()` calls
// (from compensating cleanup) get ORCA_REENTRY_REASON_DISPOSED.
for (const trace of traceStates.values()) {
trace.aborted = true;
trace.abortedReason = ORCA_REENTRY_REASON_DISPOSED;
}
}
return {
@ -379,7 +710,16 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
// ── Helpers ───────────────────────────────────────────────────────────
function groupByStage(actions: OrcaAction[]): Map<OrcaStage, OrcaAction[]> {
function resolveReentry(options: OrcaReentryOptions | undefined): ResolvedReentry {
return {
maxDepth: options?.maxDepth ?? ORCA_DEFAULT_MAX_DEPTH,
maxEventsPerTrace: options?.maxEventsPerTrace ?? ORCA_DEFAULT_MAX_EVENTS_PER_TRACE,
repeatedEventLimit: options?.repeatedEventLimit ?? ORCA_DEFAULT_REPEATED_EVENT_LIMIT,
repeatedEventPolicy: options?.repeatedEventPolicy ?? ORCA_REENTRY_SKIP
};
}
function groupByStage(actions: readonly OrcaAction[]): Map<OrcaStage, OrcaAction[]> {
const result = new Map<OrcaStage, OrcaAction[]>();
for (const action of actions) {
const list = result.get(action.stage) ?? [];
@ -397,26 +737,64 @@ function mapResultToActionStatus(result: OrcaResult) {
return ORCA_ACTION_STATUS_SKIPPED;
case ORCA_RESULT_ERROR:
return ORCA_ACTION_STATUS_ERROR;
case ORCA_RESULT_INTERRUPTED:
return ORCA_ACTION_STATUS_INTERRUPTED;
case ORCA_RESULT_FATAL:
// v0.0 treats fatal as error.
// v0 treats fatal as error.
return ORCA_ACTION_STATUS_ERROR;
default:
// timeout is not produced in v0.0; defensive fallback.
// timeout is not produced in v0; defensive fallback.
return ORCA_ACTION_STATUS_ERROR;
}
}
function computeRunStatus(aborted: boolean, actions: OrcaActionRun[]) {
function computeRunStatus(
aborted: boolean,
actions: OrcaActionRun[],
envelope: OrcaEnvelope,
traceStates: Map<OrcaTraceId, TraceState>
) {
if (aborted) return ORCA_RUN_ABORTED;
const trace = traceStates.get(envelope.meta.traceId);
if (trace?.aborted && trace.abortedReason !== ORCA_REENTRY_REASON_RUN_ABORTED) {
return ORCA_RUN_INTERRUPTED;
}
const anyError = actions.some((a) => a.status === ORCA_ACTION_STATUS_ERROR);
if (anyError) return ORCA_RUN_PARTIAL;
const anyInterrupted = actions.some((a) => a.status === ORCA_ACTION_STATUS_INTERRUPTED);
if (anyInterrupted) return ORCA_RUN_INTERRUPTED;
return ORCA_RUN_SUCCESS;
}
function interruptedActionRun(
action: OrcaAction,
now: number,
reason: OrcaReentryReason
): OrcaActionRun {
return {
id: action.id,
stage: action.stage,
status: ORCA_ACTION_STATUS_INTERRUPTED,
startedAt: now,
endedAt: now,
durationMs: 0,
emitted: [],
interruptedReason: reason
};
}
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)}`;
}
function createNoopLogger(): Logger {
const noop = () => {};
return {
@ -428,4 +806,3 @@ function createNoopLogger(): Logger {
fatal: noop
};
}

@ -15,6 +15,7 @@ export const ORCA_ERR_DISPOSED: ErrCode = errCode(ORCA_ERR, 'disposed');
export const ORCA_ERR_DUPLICATE_ACTION_ID: ErrCode = errCode(ORCA_ERR, 'duplicate_action_id');
export const ORCA_ERR_INVALID_STAGE: ErrCode = errCode(ORCA_ERR, 'invalid_stage');
export const ORCA_ERR_INVALID_ACTION: ErrCode = errCode(ORCA_ERR, 'invalid_action');
export const ORCA_ERR_REENTRY: ErrCode = errCode(ORCA_ERR, 'reentry');
// ── Error message strings ──────────────────────────────────────────────
@ -31,13 +32,17 @@ export const orcaInvalidStageMessage = (stage: string): string =>
export const orcaInvalidActionMessage = (reason: string): string =>
`[${ORCA_MODULE}] invalid action: ${reason}`;
export const orcaReentryMessage = (reason: string, traceId: string, event: string): string =>
`[${ORCA_MODULE}] reentry guard fired (reason=${reason}, trace=${traceId}, event=${event})`;
// ── Error messages ─────────────────────────────────────────────────────
export const ORCA_ERROR_MESSAGES: ErrorMessages = {
[ORCA_ERR_DISPOSED]: ORCA_ERROR_MSG_DISPOSED,
[ORCA_ERR_DUPLICATE_ACTION_ID]: orcaDuplicateActionIdMessage,
[ORCA_ERR_INVALID_STAGE]: orcaInvalidStageMessage,
[ORCA_ERR_INVALID_ACTION]: orcaInvalidActionMessage
[ORCA_ERR_INVALID_ACTION]: orcaInvalidActionMessage,
[ORCA_ERR_REENTRY]: orcaReentryMessage
};
// ── Error classes ──────────────────────────────────────────────────────
@ -74,6 +79,20 @@ export class OrcaInvalidActionError extends CodeError {
}
}
export class OrcaReentryError extends CodeError {
readonly reason: string;
readonly traceId: string;
readonly event: string;
constructor(reason: string, traceId: string, event: string) {
super(ORCA_ERR_REENTRY, {
message: orcaReentryMessage(reason, traceId, event)
});
this.reason = reason;
this.traceId = traceId;
this.event = event;
}
}
// ── Type guards ────────────────────────────────────────────────────────
export function isOrcaDisposedError(error: unknown): error is OrcaDisposedError {
@ -91,3 +110,7 @@ export function isOrcaInvalidStageError(error: unknown): error is OrcaInvalidSta
export function isOrcaInvalidActionError(error: unknown): error is OrcaInvalidActionError {
return error instanceof OrcaInvalidActionError;
}
export function isOrcaReentryError(error: unknown): error is OrcaReentryError {
return error instanceof OrcaReentryError;
}

@ -5,6 +5,7 @@ export {
orcaSuccess,
orcaSkipped,
orcaError,
orcaInterrupted,
orcaTimeout,
orcaFatal
} from './result.ts';
@ -22,25 +23,41 @@ export {
ORCA_RESULT_SUCCESS,
ORCA_RESULT_SKIPPED,
ORCA_RESULT_ERROR,
ORCA_RESULT_INTERRUPTED,
ORCA_RESULT_TIMEOUT,
ORCA_RESULT_FATAL,
ORCA_ACTION_STATUS_SUCCESS,
ORCA_ACTION_STATUS_SKIPPED,
ORCA_ACTION_STATUS_BLOCKED,
ORCA_ACTION_STATUS_ERROR,
ORCA_ACTION_STATUS_INTERRUPTED,
ORCA_ACTION_STATUS_TIMEOUT,
ORCA_ACTION_STATUS_FATAL,
ORCA_RUN_SUCCESS,
ORCA_RUN_PARTIAL,
ORCA_RUN_ABORTED,
ORCA_RUN_INTERRUPTED,
ORCA_RUN_FATAL,
ORCA_RUN_TIMEOUT,
ORCA_ON_ERROR_CONTINUE,
ORCA_ON_ERROR_ABORT_ACTION,
ORCA_ON_ERROR_ABORT_STAGE,
ORCA_ON_ERROR_ABORT_RUN,
ORCA_REENTRY_SKIP,
ORCA_REENTRY_ABORT_TRACE,
ORCA_REENTRY_ERROR,
ORCA_REENTRY_REASON_MAX_DEPTH,
ORCA_REENTRY_REASON_MAX_EVENTS_PER_TRACE,
ORCA_REENTRY_REASON_REPEATED_EVENT,
ORCA_REENTRY_REASON_DEDUPED,
ORCA_REENTRY_REASON_TRACE_ABORTED,
ORCA_REENTRY_REASON_RUN_ABORTED,
ORCA_REENTRY_REASON_DISPOSED,
ORCA_DIAGNOSTIC_EVENTS,
ORCA_DEFAULT_MAX_RUNS
ORCA_DEFAULT_MAX_RUNS,
ORCA_DEFAULT_MAX_DEPTH,
ORCA_DEFAULT_MAX_EVENTS_PER_TRACE,
ORCA_DEFAULT_REPEATED_EVENT_LIMIT
} from './consts.ts';
export {
@ -49,15 +66,18 @@ export {
ORCA_ERR_DUPLICATE_ACTION_ID,
ORCA_ERR_INVALID_STAGE,
ORCA_ERR_INVALID_ACTION,
ORCA_ERR_REENTRY,
ORCA_ERROR_MESSAGES,
OrcaDisposedError,
OrcaDuplicateActionIdError,
OrcaInvalidStageError,
OrcaInvalidActionError,
OrcaReentryError,
isOrcaDisposedError,
isOrcaDuplicateActionIdError,
isOrcaInvalidStageError,
isOrcaInvalidActionError
isOrcaInvalidActionError,
isOrcaReentryError
} from './errors.ts';
export type {
@ -69,9 +89,17 @@ export type {
OrcaActionId,
OrcaActionRun,
OrcaBus,
OrcaEmitOptions,
OrcaEnvelope,
OrcaError,
OrcaErrorPolicy,
OrcaEventId,
OrcaEventMeta,
OrcaFatal,
OrcaInterrupted,
OrcaReentryOptions,
OrcaReentryPolicy,
OrcaReentryReason,
OrcaResult,
OrcaRunId,
OrcaRunResult,
@ -79,7 +107,8 @@ export type {
OrcaStage,
OrcaSuccess,
OrcaTimeout,
OrcaToken
OrcaToken,
OrcaTraceId
} from './types.ts';
export type {

@ -9,6 +9,7 @@ import {
ORCA_RESULT_SUCCESS,
ORCA_RESULT_SKIPPED,
ORCA_RESULT_ERROR,
ORCA_RESULT_INTERRUPTED,
ORCA_RESULT_TIMEOUT,
ORCA_RESULT_FATAL
} from './consts.ts';
@ -16,6 +17,8 @@ import type {
OrcaSuccess,
OrcaSkipped,
OrcaError,
OrcaInterrupted,
OrcaReentryReason,
OrcaTimeout,
OrcaFatal,
OrcaToken
@ -57,6 +60,23 @@ export function orcaError(
};
}
/**
* The action surrendered because a guard told it to stop. Typical use:
* an action checks `ctx.signal.aborted` early and returns
* `orcaInterrupted('run aborted')` instead of doing further work.
*/
export function orcaInterrupted(
reason: OrcaReentryReason | string,
options: { emits?: readonly OrcaToken[] } = {}
): OrcaInterrupted {
return {
ok: false,
status: ORCA_RESULT_INTERRUPTED,
reason,
emits: options.emits
};
}
/** @v0.1+ */
export function orcaTimeout(
timeoutMs: number,

@ -859,3 +859,328 @@ describe('EngineOrca v0.0 — accepted-but-ignored fields (forward-compat)', ()
expect(compensate).not.toHaveBeenCalled();
});
});
// ── v0-kernel: envelope, ctx.emit, reentry guards ─────────────────────
describe('EngineOrca v0 — envelope and OrcaActionContext', () => {
let bus: FakeBus;
let timers: ReturnType<typeof createFakeTimers>;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('exposes eventId, traceId and depth=0 on root events', async () => {
const orca = createEngineOrca({ bus, timers });
const seen: Array<{ eventId: string; traceId: string; depth: number }> = [];
orca.onEvent('root', {
id: 'observer',
stage: ORCA_STAGE_MAIN,
action: (_payload, ctx) => {
seen.push({ eventId: ctx.eventId, traceId: ctx.traceId, depth: ctx.depth });
return orcaSuccess();
}
});
bus.publish('root', null);
await flush();
expect(seen).toHaveLength(1);
expect(seen[0].depth).toBe(0);
expect(seen[0].eventId).toMatch(/^evt_/);
expect(seen[0].traceId).toMatch(/^trc_/);
});
it('ctx.emit() creates a child envelope inheriting traceId and incrementing depth', async () => {
const orca = createEngineOrca({ bus, timers });
const observed: Array<{
event: string;
traceId: string;
parentEventId?: string;
depth: number;
}> = [];
orca.onEvent('parent', {
id: 'parent-action',
stage: ORCA_STAGE_MAIN,
action: (_payload, ctx) => {
observed.push({
event: ctx.event,
traceId: ctx.traceId,
parentEventId: ctx.parentEventId,
depth: ctx.depth
});
ctx.emit('child', { ok: true });
return orcaSuccess();
}
});
orca.onEvent('child', {
id: 'child-action',
stage: ORCA_STAGE_MAIN,
action: (_payload, ctx) => {
observed.push({
event: ctx.event,
traceId: ctx.traceId,
parentEventId: ctx.parentEventId,
depth: ctx.depth
});
return orcaSuccess();
}
});
bus.publish('parent', null);
await flush();
expect(observed).toHaveLength(2);
expect(observed[0].event).toBe('parent');
expect(observed[0].depth).toBe(0);
expect(observed[0].parentEventId).toBeUndefined();
expect(observed[1].event).toBe('child');
expect(observed[1].depth).toBe(1);
expect(observed[1].traceId).toBe(observed[0].traceId);
// `recentRuns` correlates parent/child by id.
const runs = orca.recentRuns();
expect(observed[1].parentEventId).toBe(runs[0].eventId);
});
it('records eventId, traceId and depth on each run trace', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('a', {
id: 'a-action',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
ctx.emit('b', null);
return orcaSuccess();
}
});
orca.onEvent('b', {
id: 'b-action',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess()
});
bus.publish('a', null);
await flush();
const runs = orca.recentRuns();
expect(runs).toHaveLength(2);
expect(runs[0].depth).toBe(0);
expect(runs[1].depth).toBe(1);
expect(runs[1].traceId).toBe(runs[0].traceId);
expect(runs[1].parentEventId).toBe(runs[0].eventId);
});
it('ctx.emit() with no listeners returns the eventId without enqueueing a run', async () => {
const orca = createEngineOrca({ bus, timers });
let returnedId: string | null = null;
orca.onEvent('only', {
id: 'only-action',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
returnedId = ctx.emit('orphan', null);
return orcaSuccess();
}
});
bus.publish('only', null);
await flush();
expect(returnedId).toMatch(/^evt_/);
expect(orca.recentRuns()).toHaveLength(1); // only the parent run, no orphan run
});
});
describe('EngineOrca v0 — reentry guards', () => {
let bus: FakeBus;
let timers: ReturnType<typeof createFakeTimers>;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('blocks events that exceed maxDepth (default policy: skip)', async () => {
// Use distinct event names per depth so the same-name limit does
// not interfere with the maxDepth check we want to exercise.
const orca = createEngineOrca({ bus, timers, reentry: { maxDepth: 2 } });
const log: number[] = [];
const action = (label: string) => ({
id: `${label}-action`,
stage: ORCA_STAGE_MAIN,
action: (_p: unknown, ctx: { depth: number; emit: (e: string, p: unknown) => unknown }) => {
log.push(ctx.depth);
ctx.emit(`level-${ctx.depth + 1}`, null);
return orcaSuccess();
}
});
orca.onEvent('level-0', action('level-0'));
orca.onEvent('level-1', action('level-1'));
orca.onEvent('level-2', action('level-2'));
orca.onEvent('level-3', action('level-3'));
bus.publish('level-0', null);
await flush();
// Depths 0, 1, 2 accepted (root through maxDepth); depth 3 blocked.
expect(log).toEqual([0, 1, 2]);
});
it('blocks events whose name repeats more than repeatedEventLimit times', async () => {
const orca = createEngineOrca({
bus,
timers,
reentry: { repeatedEventLimit: 1, maxDepth: 32 }
});
const log: number[] = [];
orca.onEvent('chain', {
id: 'recurser',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
log.push(ctx.depth);
ctx.emit('chain', null);
return orcaSuccess();
}
});
bus.publish('chain', null);
await flush();
// Root counts as occurrence 1; second 'chain' event hits the limit.
expect(log).toEqual([0]);
});
it('throws OrcaReentryError when repeatedEventPolicy is error', async () => {
const orca = createEngineOrca({
bus,
timers,
reentry: {
repeatedEventLimit: 1,
repeatedEventPolicy: 'error',
maxDepth: 32
}
});
let captured: unknown;
orca.onEvent('chain', {
id: 'recurser',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
try {
ctx.emit('chain', null);
} catch (err) {
captured = err;
}
return orcaSuccess();
}
});
bus.publish('chain', null);
await flush();
expect(captured).toBeDefined();
expect((captured as Error).message).toMatch(/reentry guard fired/);
});
it('aborts the trace and records ORCA_RUN_INTERRUPTED with abort-trace policy', async () => {
const orca = createEngineOrca({
bus,
timers,
reentry: {
maxDepth: 1,
repeatedEventPolicy: 'abort-trace'
}
});
orca.onEvent('a', {
id: 'a-action',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
ctx.emit('b', null);
return orcaSuccess();
}
});
orca.onEvent('b', {
id: 'b-action',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
ctx.emit('c', null); // exceeds maxDepth=1; trace aborts
return orcaSuccess();
}
});
orca.onEvent('c', {
id: 'c-action',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess()
});
bus.publish('a', null);
await flush();
const runs = orca.recentRuns();
// Two runs: parent + child. Grandchild was rejected; trace then aborted.
expect(runs).toHaveLength(2);
});
it('dedupeKey skips a second emit with the same key in the same trace', async () => {
const orca = createEngineOrca({ bus, timers });
const log: string[] = [];
orca.onEvent('seed', {
id: 'seeder',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
ctx.emit('child', { id: 1 }, { dedupeKey: 'k1' });
ctx.emit('child', { id: 2 }, { dedupeKey: 'k1' }); // skipped
ctx.emit('child', { id: 3 }, { dedupeKey: 'k2' }); // distinct key, accepted
return orcaSuccess();
}
});
orca.onEvent('child', {
id: 'child-action',
stage: ORCA_STAGE_MAIN,
action: (payload) => {
log.push(JSON.stringify(payload));
return orcaSuccess();
}
});
bus.publish('seed', null);
await flush();
expect(log).toEqual(['{"id":1}', '{"id":3}']);
});
it('module-level bus publishes always start fresh traces (not subject to reentry)', async () => {
const orca = createEngineOrca({
bus,
timers,
reentry: { repeatedEventLimit: 1 }
});
const traces: string[] = [];
orca.onEvent('e', {
id: 'observer',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
traces.push(ctx.traceId);
return orcaSuccess();
}
});
bus.publish('e', null);
bus.publish('e', null);
bus.publish('e', null);
await flush();
expect(traces).toHaveLength(3);
expect(new Set(traces).size).toBe(3); // each is a fresh trace
});
});

@ -11,23 +11,36 @@ import type {
ORCA_RESULT_SUCCESS,
ORCA_RESULT_SKIPPED,
ORCA_RESULT_ERROR,
ORCA_RESULT_INTERRUPTED,
ORCA_RESULT_TIMEOUT,
ORCA_RESULT_FATAL,
ORCA_ACTION_STATUS_SUCCESS,
ORCA_ACTION_STATUS_SKIPPED,
ORCA_ACTION_STATUS_BLOCKED,
ORCA_ACTION_STATUS_ERROR,
ORCA_ACTION_STATUS_INTERRUPTED,
ORCA_ACTION_STATUS_TIMEOUT,
ORCA_ACTION_STATUS_FATAL,
ORCA_RUN_SUCCESS,
ORCA_RUN_PARTIAL,
ORCA_RUN_ABORTED,
ORCA_RUN_INTERRUPTED,
ORCA_RUN_FATAL,
ORCA_RUN_TIMEOUT,
ORCA_ON_ERROR_CONTINUE,
ORCA_ON_ERROR_ABORT_ACTION,
ORCA_ON_ERROR_ABORT_STAGE,
ORCA_ON_ERROR_ABORT_RUN
ORCA_ON_ERROR_ABORT_RUN,
ORCA_REENTRY_SKIP,
ORCA_REENTRY_ABORT_TRACE,
ORCA_REENTRY_ERROR,
ORCA_REENTRY_REASON_MAX_DEPTH,
ORCA_REENTRY_REASON_MAX_EVENTS_PER_TRACE,
ORCA_REENTRY_REASON_REPEATED_EVENT,
ORCA_REENTRY_REASON_DEDUPED,
ORCA_REENTRY_REASON_TRACE_ABORTED,
ORCA_REENTRY_REASON_RUN_ABORTED,
ORCA_REENTRY_REASON_DISPOSED
} from './consts.ts';
// ── Identifiers ───────────────────────────────────────────────────────
@ -43,6 +56,35 @@ export type OrcaStage =
export type OrcaToken = string;
export type OrcaActionId = string;
export type OrcaRunId = string;
export type OrcaEventId = string;
export type OrcaTraceId = string;
// ── Event envelope ─────────────────────────────────────────────────────
//
// Every event the engine processes is wrapped in an envelope before it
// reaches the run queue. The envelope carries runtime metadata
// (identifiers, depth, parent linkage, originating action) so the
// engine can reason about reentry and trace propagation without
// polluting the user-defined payload.
export interface OrcaEventMeta {
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly parentRunId?: OrcaRunId;
readonly emittedByAction?: OrcaActionId;
readonly depth: number;
/** Sequence of event names from the trace root to this event. */
readonly stack: readonly string[];
readonly publishedAt: number;
readonly dedupeKey?: string;
}
export interface OrcaEnvelope<TPayload = unknown> {
readonly event: string;
readonly payload: TPayload;
readonly meta: OrcaEventMeta;
}
export type OrcaErrorPolicy =
| typeof ORCA_ON_ERROR_CONTINUE
@ -86,6 +128,20 @@ export interface OrcaTimeout {
readonly emits?: readonly OrcaToken[];
}
/**
* Action did not complete because a runtime guard cut it short. The
* action either returned `orcaInterrupted(...)` explicitly (on reading
* `ctx.signal.aborted`) or the engine wrote it into the run trace
* because a reentry guard, trace abort or dispose blocked the action
* before it could run.
*/
export interface OrcaInterrupted {
readonly ok: false;
readonly status: typeof ORCA_RESULT_INTERRUPTED;
readonly reason: OrcaReentryReason | string;
readonly emits?: readonly OrcaToken[];
}
/**
* @v0.1+ Distinguished from `OrcaError` for run-aborting failures.
* @v0.0 Engine treats fatal as error.
@ -101,15 +157,94 @@ export type OrcaResult<TValue = unknown> =
| OrcaSuccess<TValue>
| OrcaSkipped
| OrcaError
| OrcaInterrupted
| OrcaTimeout
| OrcaFatal;
// ── Reentry policy ────────────────────────────────────────────────────
export type OrcaReentryPolicy =
| typeof ORCA_REENTRY_SKIP
| typeof ORCA_REENTRY_ABORT_TRACE
| typeof ORCA_REENTRY_ERROR;
export type OrcaReentryReason =
| typeof ORCA_REENTRY_REASON_MAX_DEPTH
| typeof ORCA_REENTRY_REASON_MAX_EVENTS_PER_TRACE
| typeof ORCA_REENTRY_REASON_REPEATED_EVENT
| typeof ORCA_REENTRY_REASON_DEDUPED
| typeof ORCA_REENTRY_REASON_TRACE_ABORTED
| typeof ORCA_REENTRY_REASON_RUN_ABORTED
| typeof ORCA_REENTRY_REASON_DISPOSED;
/**
* Tunables that bound the size and shape of a single trace
* (a trace = one root event + every event derived from it via
* `ctx.emit()`).
*/
export interface OrcaReentryOptions {
/**
* Maximum nesting depth a single trace can reach. Depth 0 is the
* root event published on the bus; depth 1 is an event emitted from
* an action handling the root; etc.
* @default 16
*/
readonly maxDepth?: number;
/**
* Maximum total number of events allowed to belong to one trace.
* @default 128
*/
readonly maxEventsPerTrace?: number;
/**
* Maximum number of times the same event name may appear in a single
* trace before the policy fires.
* @default 2
*/
readonly repeatedEventLimit?: number;
/**
* What the engine does when one of the limits above is hit.
* @default 'skip'
*/
readonly repeatedEventPolicy?: OrcaReentryPolicy;
}
// ── Action context ────────────────────────────────────────────────────
export interface OrcaEmitOptions {
/**
* Optional dedupe key. Two events emitted into the same trace with
* the same dedupeKey trigger the configured reentry policy on the
* second occurrence. Independent of `repeatedEventLimit` (which
* counts event names, not arbitrary keys).
*/
readonly dedupeKey?: string;
}
export interface OrcaActionContext {
readonly runId: OrcaRunId;
readonly event: string;
readonly stage: OrcaStage;
/**
* Identity of the envelope being processed. Stable across all
* actions of the same run.
*/
readonly eventId: OrcaEventId;
/**
* Trace this run belongs to. Root events get a fresh traceId;
* events emitted via `ctx.emit()` inherit their parent's traceId.
*/
readonly traceId: OrcaTraceId;
/**
* Identity of the envelope that produced this one via `ctx.emit()`,
* or `undefined` when this run is processing a root event published
* directly on the bus.
*/
readonly parentEventId?: OrcaEventId;
/**
* Distance from the trace root. Root events have `depth = 0`; the
* first child via `ctx.emit()` has `depth = 1`; etc.
*/
readonly depth: number;
/**
* Tokens already emitted in the current run by previous actions.
* @v0.0 Populated correctly, but not consumed by the engine
@ -118,8 +253,8 @@ export interface OrcaActionContext {
readonly tokens: ReadonlySet<OrcaToken>;
/**
* Abort signal for the current action. Aborts when the run is aborted
* (via ORCA_ON_ERROR_ABORT_RUN from another action) or when the engine
* is disposed.
* (via ORCA_ON_ERROR_ABORT_RUN from another action), when the trace
* is aborted by a reentry guard, or when the engine is disposed.
*/
readonly signal: AbortSignal;
/**
@ -127,6 +262,28 @@ export interface OrcaActionContext {
* are emitted by the engine itself.
*/
readonly logger: Logger;
/**
* Publish a derived event back into the engine. The new envelope
* inherits this run's `traceId`, declares this envelope as parent,
* sets `depth = ctx.depth + 1`, and stamps `emittedByAction` with
* the active action's id.
*
* Reentry guards apply: if the resulting envelope would exceed
* `maxDepth`, `maxEventsPerTrace`, `repeatedEventLimit`, or matches
* a previous `dedupeKey` in the trace, the configured policy fires
* and the event is not enqueued.
*
* `ctx.emit()` returns the assigned `eventId` (whether or not the
* event was ultimately enqueued — the id participates in
* diagnostics either way) or `null` when the engine refused to
* enqueue and the policy is `skip` / `abort-trace`. Throws
* `OrcaReentryError` when the policy is `error`.
*/
emit<TPayload>(
event: string,
payload: TPayload,
options?: OrcaEmitOptions
): OrcaEventId | null;
}
// ── Action definition ─────────────────────────────────────────────────
@ -202,6 +359,7 @@ export interface OrcaActionRun {
| typeof ORCA_ACTION_STATUS_SKIPPED
| typeof ORCA_ACTION_STATUS_BLOCKED
| typeof ORCA_ACTION_STATUS_ERROR
| typeof ORCA_ACTION_STATUS_INTERRUPTED
| typeof ORCA_ACTION_STATUS_TIMEOUT
| typeof ORCA_ACTION_STATUS_FATAL;
readonly startedAt: number;
@ -209,15 +367,26 @@ export interface OrcaActionRun {
readonly durationMs: number;
readonly emitted: readonly OrcaToken[];
readonly error?: unknown;
/** Populated when status === 'interrupted'. */
readonly interruptedReason?: OrcaReentryReason | string;
}
export interface OrcaRunResult {
readonly id: OrcaRunId;
readonly event: string;
/**
* Identity of the envelope this run processed. Stable across the
* actions in `actions[]`.
*/
readonly eventId: OrcaEventId;
readonly traceId: OrcaTraceId;
readonly parentEventId?: OrcaEventId;
readonly depth: number;
readonly status:
| typeof ORCA_RUN_SUCCESS
| typeof ORCA_RUN_PARTIAL
| typeof ORCA_RUN_ABORTED
| typeof ORCA_RUN_INTERRUPTED
| typeof ORCA_RUN_FATAL
| typeof ORCA_RUN_TIMEOUT;
readonly startedAt: number;
@ -258,6 +427,12 @@ export interface EngineOrcaOptions {
* @default 256
*/
readonly maxRuns?: number;
/**
* Reentry guards applied to events emitted via `ctx.emit()`. Module
* publishes via the bus directly are root events and do not enter
* these counters. See `OrcaReentryOptions` for the individual knobs.
*/
readonly reentry?: OrcaReentryOptions;
}
export interface EngineOrca {

Loading…
Cancel
Save

Powered by TurnKey Linux.