Tokens with payload — `emits` accepts `{ token, payload }`

Token emissions can now carry arbitrary payload data. The `emits`
array accepts either a bare token name (existing form) or
`{ token, payload }`; both forms mix freely. Downstream actions
read payloads via `ctx.tokenPayloads.get(name)`, with
`ctx.tokens.has()` still answering name-presence. The two views
can diverge: a string-form emission has presence but no payload.

Gates (`after` / `unless` / `abortOn` / `provides`) keep comparing
names only — payload semantics are opt-in metadata. Wave snapshots
extend to payloads, so parallel siblings never read each other's
payloads mid-flight. Last-write-wins on duplicate names.

Surface: `OrcaTokenWithPayload`, `OrcaTokenEmit`,
`OrcaActionContext.tokenPayloads`, `OrcaActionRun.emittedPayloads?`,
`OrcaRunResult.tokenPayloads`. Compensation contexts also expose
`tokenPayloads` so rollback paths can read the run-level state.

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

@ -1347,7 +1347,6 @@ Marcado como `@v1+` en el código fuente:
| Pieza | Qué falta para v1 |
|---|---|
| `commit()` / `replace` cola | Políticas de cola por evento (drop-prev / replace / parallel). |
| 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. |
| `createActiveOrca()` | Wrapper reactivo (`$state` para `running`, `recentRuns`, etc.) — el `EngineOrca` ya está. |
@ -1417,6 +1416,21 @@ Marcado como `@v1+` en el código fuente:
navegación de trace. `validate()` emite warning
`transaction-without-compensate` cuando todos los miembros de un
tx no declaran `compensate` (rollback no-op).
- ✅ **Tokens con payload** — el array `emits` acepta entradas
`{ token, payload }` además de strings sueltos; ambos formatos
se mezclan en el mismo array. `ctx.tokenPayloads` (mapa
`ReadonlyMap<OrcaToken, unknown>`) expone los payloads por
nombre — los tokens emitidos como string suelto **no** aparecen
en el mapa, así `ctx.tokens.has('x')` y
`ctx.tokenPayloads.has('x')` pueden divergir (presencia vs.
payload). Los gates (`after`/`unless`/`abortOn`/`provides`)
siguen comparando solo nombres. La instantánea de `tokenPayloads`
por wave es independiente: paralelas hermanas nunca ven los
payloads emitidos por sus pares mid-flight. Last-write-wins en
colisiones de nombre. `OrcaActionRun.emittedPayloads?` carga el
mapa por acción (ausente si la acción solo emitió strings).
`OrcaRunResult.tokenPayloads` es la unión final del run (Map
vacío cuando ningún token llevó payload).
- ✅ **`compensate`** — cada acción puede declarar una función
compensatoria. Cuando el run aborta (FATAL / ABORT_RUN onError /
trace-aborted), el motor recorre la pila LIFO de acciones que ya

@ -407,6 +407,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
const runId = generateRunId();
const startedAt = timers.clock.now();
const tokens = new Set<OrcaToken>();
const tokenPayloads = new Map<OrcaToken, unknown>();
const actionRuns: OrcaActionRun[] = [];
// Stack of actions that completed success AND declared a compensate.
// Pushed in completion order so reverse iteration is LIFO. The
@ -457,7 +458,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
entry.payload,
envelope,
runId,
tokens
tokens,
tokenPayloads
);
compensations.push(compResult);
entry.compensated = true;
@ -484,9 +486,10 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
// wave, every action sees the same view; their emissions
// are merged into the run-level set after the wave settles.
const waveTokens: ReadonlySet<OrcaToken> = new Set(tokens);
const waveTokenPayloads: ReadonlyMap<OrcaToken, unknown> = new Map(tokenPayloads);
const waveItems = await Promise.all(
wave.map((action) => runActionInWave(action, stage, waveTokens, runId, envelope, controller))
wave.map((action) => runActionInWave(action, stage, waveTokens, waveTokenPayloads, runId, envelope, controller))
);
for (const item of waveItems) {
@ -504,6 +507,9 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
});
}
for (const token of item.actionRun.emitted) tokens.add(token);
if (item.actionRun.emittedPayloads) {
for (const [k, v] of item.actionRun.emittedPayloads) tokenPayloads.set(k, v);
}
}
// Detect transaction failure first: any tx member that
@ -535,7 +541,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
entry.payload,
envelope,
runId,
tokens
tokens,
tokenPayloads
);
compensations.push(compResult);
entry.compensated = true;
@ -607,6 +614,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
endedAt,
durationMs: endedAt - startedAt,
tokens: Array.from(tokens),
tokenPayloads,
actions: actionRuns,
compensations
};
@ -636,6 +644,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
action: OrcaAction,
stage: OrcaStage,
waveTokens: ReadonlySet<OrcaToken>,
waveTokenPayloads: ReadonlyMap<OrcaToken, unknown>,
runId: OrcaRunId,
envelope: OrcaEnvelope,
runController: AbortController
@ -689,6 +698,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
envelope,
action.stage,
waveTokens,
waveTokenPayloads,
actionController.signal,
action.id
);
@ -703,6 +713,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
envelope: OrcaEnvelope,
stage: OrcaStage,
tokens: ReadonlySet<OrcaToken>,
tokenPayloads: ReadonlyMap<OrcaToken, unknown>,
signal: AbortSignal,
actionId: OrcaActionId
): OrcaActionContext {
@ -715,6 +726,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
parentEventId: envelope.meta.parentEventId,
depth: envelope.meta.depth,
tokens,
tokenPayloads,
signal,
logger: actionLogger,
emit<TPayload>(
@ -772,7 +784,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
payload: unknown,
envelope: OrcaEnvelope,
runId: OrcaRunId,
tokens: ReadonlySet<OrcaToken>
tokens: ReadonlySet<OrcaToken>,
tokenPayloads: ReadonlyMap<OrcaToken, unknown>
): Promise<OrcaActionRun> {
const startedAt = timers.clock.now();
emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.COMPENSATION_STARTED, {
@ -793,6 +806,7 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
parentEventId: envelope.meta.parentEventId,
depth: envelope.meta.depth,
tokens,
tokenPayloads,
signal: compController.signal,
logger: actionLogger,
emit() {
@ -932,11 +946,21 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
try {
const result = await runActionWithTimeout(action, payload, context, actionController);
const endedAt = timers.clock.now();
const emitted = result.emits ?? [];
// Tokens are NOT merged into the run-level set here. The wave
// coordinator merges `actionRun.emitted` after the surrounding
// wave settles, so parallel siblings never see each other's
// tokens mid-flight.
// Normalise the heterogeneous `emits` array into name-only and
// payload-keyed views. The wave coordinator merges these into
// the run-level state once the surrounding wave settles, so
// parallel siblings never see each other's emissions mid-flight.
const rawEmits = result.emits ?? [];
const emittedNames: OrcaToken[] = [];
const emittedPayloads = new Map<OrcaToken, unknown>();
for (const e of rawEmits) {
if (typeof e === 'string') {
emittedNames.push(e);
} else {
emittedNames.push(e.token);
emittedPayloads.set(e.token, e.payload);
}
}
const status = mapResultToActionStatus(result);
@ -1004,7 +1028,8 @@ export function createEngineOrca(options: EngineOrcaOptions): EngineOrca {
startedAt,
endedAt,
durationMs: endedAt - startedAt,
emitted: Array.from(emitted),
emitted: emittedNames,
emittedPayloads: emittedPayloads.size > 0 ? emittedPayloads : undefined,
error:
status === ORCA_ACTION_STATUS_ERROR ||
status === ORCA_ACTION_STATUS_FATAL

@ -122,6 +122,8 @@ export type {
OrcaSuccess,
OrcaTimeout,
OrcaToken,
OrcaTokenEmit,
OrcaTokenWithPayload,
OrcaTraceId,
OrcaValidateIssue,
OrcaValidateIssueKind,

@ -21,11 +21,11 @@ import type {
OrcaReentryReason,
OrcaTimeout,
OrcaFatal,
OrcaToken
OrcaTokenEmit
} from './types.ts';
export function orcaSuccess<TValue = void>(
options: { value?: TValue; emits?: readonly OrcaToken[] } = {}
options: { value?: TValue; emits?: readonly OrcaTokenEmit[] } = {}
): OrcaSuccess<TValue> {
return {
ok: true,
@ -37,7 +37,7 @@ export function orcaSuccess<TValue = void>(
export function orcaSkipped(
reason?: string,
options: { emits?: readonly OrcaToken[] } = {}
options: { emits?: readonly OrcaTokenEmit[] } = {}
): OrcaSkipped {
return {
ok: true,
@ -49,7 +49,7 @@ export function orcaSkipped(
export function orcaError(
error: unknown,
options: { emits?: readonly OrcaToken[]; recoverable?: boolean } = {}
options: { emits?: readonly OrcaTokenEmit[]; recoverable?: boolean } = {}
): OrcaError {
return {
ok: false,
@ -67,7 +67,7 @@ export function orcaError(
*/
export function orcaInterrupted(
reason: OrcaReentryReason | string,
options: { emits?: readonly OrcaToken[] } = {}
options: { emits?: readonly OrcaTokenEmit[] } = {}
): OrcaInterrupted {
return {
ok: false,
@ -80,7 +80,7 @@ export function orcaInterrupted(
/** @v0.1+ */
export function orcaTimeout(
timeoutMs: number,
options: { emits?: readonly OrcaToken[] } = {}
options: { emits?: readonly OrcaTokenEmit[] } = {}
): OrcaTimeout {
return {
ok: false,
@ -93,7 +93,7 @@ export function orcaTimeout(
/** @v0.1+ */
export function orcaFatal(
error: unknown,
options: { emits?: readonly OrcaToken[] } = {}
options: { emits?: readonly OrcaTokenEmit[] } = {}
): OrcaFatal {
return {
ok: false,

@ -3303,3 +3303,255 @@ describe('EngineOrca v1 — transaction', () => {
});
});
describe('EngineOrca v1 — tokens with payload', () => {
let bus: FakeBus;
let timers: TimerScheduler;
beforeEach(() => {
bus = createFakeBus();
timers = createFakeTimers();
});
it('emits a token with payload and a downstream action reads it via ctx.tokenPayloads', async () => {
const orca = createEngineOrca({ bus, timers });
let seen: unknown;
orca.onEvent('e', {
id: 'producer',
stage: ORCA_STAGE_PRE,
provides: ['user'],
action: () =>
orcaSuccess({ emits: [{ token: 'user', payload: { id: 'u-1', name: 'ana' } }] })
});
orca.onEvent('e', {
id: 'consumer',
stage: ORCA_STAGE_MAIN,
after: ['user'],
action: (_p, ctx) => {
seen = ctx.tokenPayloads.get('user');
return orcaSuccess();
}
});
bus.publish('e', null);
await flush();
expect(seen).toEqual({ id: 'u-1', name: 'ana' });
});
it('a bare-name emission has no payload entry; .has() reports correctly', async () => {
const orca = createEngineOrca({ bus, timers });
let hasInTokens = false;
let hasInPayloads = false;
orca.onEvent('e', {
id: 'producer',
stage: ORCA_STAGE_PRE,
provides: ['plain'],
action: () => orcaSuccess({ emits: ['plain'] })
});
orca.onEvent('e', {
id: 'consumer',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
hasInTokens = ctx.tokens.has('plain');
hasInPayloads = ctx.tokenPayloads.has('plain');
return orcaSuccess();
}
});
bus.publish('e', null);
await flush();
expect(hasInTokens).toBe(true);
expect(hasInPayloads).toBe(false);
});
it('OrcaActionRun.emittedPayloads carries the payload map for that action', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
action: () =>
orcaSuccess({
emits: ['plain', { token: 'data', payload: { n: 42 } }]
})
});
bus.publish('e', null);
await flush();
const action = orca.recentRuns()[0].actions.find((r) => r.id === 'a');
expect(action?.emitted.sort()).toEqual(['data', 'plain']);
expect(action?.emittedPayloads).toBeDefined();
expect(action?.emittedPayloads?.get('data')).toEqual({ n: 42 });
expect(action?.emittedPayloads?.has('plain')).toBe(false);
});
it('OrcaActionRun.emittedPayloads is undefined when no token carries a payload', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess({ emits: ['plain'] })
});
bus.publish('e', null);
await flush();
const action = orca.recentRuns()[0].actions.find((r) => r.id === 'a');
expect(action?.emittedPayloads).toBeUndefined();
});
it('OrcaRunResult.tokenPayloads carries the union of all emitted payloads', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_PRE,
action: () => orcaSuccess({ emits: [{ token: 't1', payload: 1 }] })
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess({ emits: [{ token: 't2', payload: 'two' }] })
});
bus.publish('e', null);
await flush();
const run = orca.recentRuns()[0];
expect(run.tokenPayloads.get('t1')).toBe(1);
expect(run.tokenPayloads.get('t2')).toBe('two');
expect(run.tokenPayloads.size).toBe(2);
});
it('OrcaRunResult.tokenPayloads is empty when no token-with-payload was emitted', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess({ emits: ['plain'] })
});
bus.publish('e', null);
await flush();
expect(orca.recentRuns()[0].tokenPayloads.size).toBe(0);
});
it('last write wins when the same token is emitted twice with different payloads', async () => {
const orca = createEngineOrca({ bus, timers });
orca.onEvent('e', {
id: 'first',
stage: ORCA_STAGE_PRE,
action: () => orcaSuccess({ emits: [{ token: 'k', payload: 'first' }] })
});
orca.onEvent('e', {
id: 'second',
stage: ORCA_STAGE_MAIN,
action: () => orcaSuccess({ emits: [{ token: 'k', payload: 'second' }] })
});
bus.publish('e', null);
await flush();
expect(orca.recentRuns()[0].tokenPayloads.get('k')).toBe('second');
});
it('parallel siblings see the wave-start snapshot — never each other payloads', async () => {
const orca = createEngineOrca({ bus, timers });
let seenByB: unknown = 'untouched';
orca.onEvent('e', {
id: 'a',
stage: ORCA_STAGE_MAIN,
parallel: true,
action: async () => {
await sleep(2);
return orcaSuccess({ emits: [{ token: 't', payload: 'from-a' }] });
}
});
orca.onEvent('e', {
id: 'b',
stage: ORCA_STAGE_MAIN,
parallel: true,
action: async (_p, ctx) => {
await sleep(5);
seenByB = ctx.tokenPayloads.get('t');
return orcaSuccess();
}
});
bus.publish('e', null);
await sleep(20);
await flush();
expect(seenByB).toBeUndefined();
});
it('a sequential action after a parallel wave sees merged payloads from the wave', async () => {
const orca = createEngineOrca({ bus, timers });
let merged: Record<string, unknown> = {};
orca.onEvent('e', {
id: 'p1',
stage: ORCA_STAGE_MAIN,
parallel: true,
action: () => orcaSuccess({ emits: [{ token: 'x', payload: 1 }] })
});
orca.onEvent('e', {
id: 'p2',
stage: ORCA_STAGE_MAIN,
parallel: true,
action: () => orcaSuccess({ emits: [{ token: 'y', payload: 2 }] })
});
orca.onEvent('e', {
id: 'after',
stage: ORCA_STAGE_MAIN,
action: (_p, ctx) => {
merged = {
x: ctx.tokenPayloads.get('x'),
y: ctx.tokenPayloads.get('y')
};
return orcaSuccess();
}
});
bus.publish('e', null);
await flush();
expect(merged).toEqual({ x: 1, y: 2 });
});
it('compensation context exposes the run-level tokenPayloads on rollback', async () => {
const orca = createEngineOrca({ bus, timers });
let payloadAtCompensate: unknown;
orca.onEvent('e', {
id: 'producer',
stage: ORCA_STAGE_PRE,
action: () => orcaSuccess({ emits: [{ token: 'tx', payload: { id: 7 } }] }),
compensate: (_p, ctx) => {
payloadAtCompensate = ctx.tokenPayloads.get('tx');
return orcaSuccess();
}
});
orca.onEvent('e', {
id: 'aborter',
stage: ORCA_STAGE_MAIN,
onError: ORCA_ON_ERROR_ABORT_RUN,
action: () => orcaError(new Error('x'))
});
bus.publish('e', null);
await flush();
expect(payloadAtCompensate).toEqual({ id: 7 });
});
});

@ -66,6 +66,29 @@ export type OrcaRunId = string;
export type OrcaEventId = string;
export type OrcaTraceId = string;
/**
* Emit form that pairs a token name with a payload. Downstream actions
* read the payload via `ctx.tokenPayloads.get(token)`. Tokens emitted
* as a bare string carry no payload; the two forms can be mixed in the
* same `emits` array.
*
* If two emissions share the same token name within a run, the **last
* one in registration order** wins. Within a single parallel wave,
* sibling actions resolve ties by registration order, but we recommend
* not having parallel siblings emit the same token-with-payload to
* avoid surprise.
*/
export interface OrcaTokenWithPayload<TPayload = unknown> {
readonly token: OrcaToken;
readonly payload: TPayload;
}
/**
* One entry in `OrcaResult.emits`. Either a bare token name or a
* token paired with a payload.
*/
export type OrcaTokenEmit<TPayload = unknown> = OrcaToken | OrcaTokenWithPayload<TPayload>;
// ── Event envelope ─────────────────────────────────────────────────────
//
// Every event the engine processes is wrapped in an envelope before it
@ -105,14 +128,14 @@ export interface OrcaSuccess<TValue = unknown> {
readonly ok: true;
readonly status: typeof ORCA_RESULT_SUCCESS;
readonly value?: TValue;
readonly emits?: readonly OrcaToken[];
readonly emits?: readonly OrcaTokenEmit[];
}
export interface OrcaSkipped {
readonly ok: true;
readonly status: typeof ORCA_RESULT_SKIPPED;
readonly reason?: string;
readonly emits?: readonly OrcaToken[];
readonly emits?: readonly OrcaTokenEmit[];
}
export interface OrcaError {
@ -120,7 +143,7 @@ export interface OrcaError {
readonly status: typeof ORCA_RESULT_ERROR;
readonly error: unknown;
readonly recoverable?: boolean;
readonly emits?: readonly OrcaToken[];
readonly emits?: readonly OrcaTokenEmit[];
}
/**
@ -132,7 +155,7 @@ export interface OrcaTimeout {
readonly ok: false;
readonly status: typeof ORCA_RESULT_TIMEOUT;
readonly timeoutMs: number;
readonly emits?: readonly OrcaToken[];
readonly emits?: readonly OrcaTokenEmit[];
}
/**
@ -146,7 +169,7 @@ export interface OrcaInterrupted {
readonly ok: false;
readonly status: typeof ORCA_RESULT_INTERRUPTED;
readonly reason: OrcaReentryReason | string;
readonly emits?: readonly OrcaToken[];
readonly emits?: readonly OrcaTokenEmit[];
}
/**
@ -157,7 +180,7 @@ export interface OrcaFatal {
readonly ok: false;
readonly status: typeof ORCA_RESULT_FATAL;
readonly error: unknown;
readonly emits?: readonly OrcaToken[];
readonly emits?: readonly OrcaTokenEmit[];
}
export type OrcaResult<TValue = unknown> =
@ -258,6 +281,15 @@ export interface OrcaActionContext {
* (after/unless/abortOn are ignored). Action authors MAY read it.
*/
readonly tokens: ReadonlySet<OrcaToken>;
/**
* Payloads keyed by token name for tokens that were emitted in the
* `{ token, payload }` form. Tokens emitted as bare names are not
* present in this map. Last-write-wins on collisions.
*
* Like `tokens`, this is a snapshot taken at wave start, so parallel
* siblings within the same wave never see each other's payloads.
*/
readonly tokenPayloads: ReadonlyMap<OrcaToken, unknown>;
/**
* Abort signal for the current action. Aborts when the run is aborted
* (via ORCA_ON_ERROR_ABORT_RUN from another action), when the trace
@ -418,6 +450,12 @@ export interface OrcaActionRun {
readonly endedAt: number;
readonly durationMs: number;
readonly emitted: readonly OrcaToken[];
/**
* Payloads keyed by token name for tokens this action emitted in
* `{ token, payload }` form. Absent when the action emitted no
* tokens-with-payload. Use `emitted` for name-only iteration.
*/
readonly emittedPayloads?: ReadonlyMap<OrcaToken, unknown>;
readonly error?: unknown;
/**
* Short tag describing why the action ended in this status.
@ -458,6 +496,13 @@ export interface OrcaRunResult {
readonly endedAt: number;
readonly durationMs: number;
readonly tokens: readonly OrcaToken[];
/**
* Payloads accumulated through the run, keyed by token name. Only
* tokens emitted in the `{ token, payload }` form appear here;
* bare-name emissions are absent. Empty Map when no token carried
* a payload.
*/
readonly tokenPayloads: ReadonlyMap<OrcaToken, unknown>;
readonly actions: readonly OrcaActionRun[];
/**
* Compensation runs invoked when the run aborted (FATAL, ABORT_RUN

Loading…
Cancel
Save

Powered by TurnKey Linux.