diff --git a/src/arts/orca/consts.ts b/src/arts/orca/consts.ts new file mode 100644 index 0000000..519bc43 --- /dev/null +++ b/src/arts/orca/consts.ts @@ -0,0 +1,94 @@ +/** + * Engine-internal constants for `arts/orca` v0.0. + * + * Public values that consumers can reach (stages, error policies, result + * statuses, run statuses) are exported here. Module-local strings (log + * messages) live alongside. + */ + +export const ORCA_MODULE = 'orca' as const; + +// ── Stages ────────────────────────────────────────────────────────────── + +export const ORCA_STAGE_GUARD = 'guard' as const; +export const ORCA_STAGE_PRE = 'pre' as const; +export const ORCA_STAGE_MAIN = 'main' as const; +export const ORCA_STAGE_POST = 'post' as const; +export const ORCA_STAGE_CLEANUP = 'cleanup' as const; +export const ORCA_STAGE_FINALLY = 'finally' as const; + +export const ORCA_STAGES_CANONICAL_ORDER = [ + ORCA_STAGE_GUARD, + ORCA_STAGE_PRE, + ORCA_STAGE_MAIN, + ORCA_STAGE_POST, + ORCA_STAGE_CLEANUP, + ORCA_STAGE_FINALLY +] as const; + +// ── Result statuses ───────────────────────────────────────────────────── + +export const ORCA_RESULT_SUCCESS = 'success' as const; +export const ORCA_RESULT_SKIPPED = 'skipped' as const; +export const ORCA_RESULT_ERROR = 'error' as const; + +/** @v0.1+ Accepted in the Result type, never produced by the v0.0 engine. */ +export const ORCA_RESULT_TIMEOUT = 'timeout' as const; +/** @v0.1+ Accepted in the Result type, never produced by the v0.0 engine. */ +export const ORCA_RESULT_FATAL = 'fatal' as const; + +// ── Action statuses (run trace) ───────────────────────────────────────── + +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_TIMEOUT = 'timeout' as const; +export const ORCA_ACTION_STATUS_FATAL = 'fatal' as const; + +// ── Run statuses ──────────────────────────────────────────────────────── + +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_FATAL = 'fatal' as const; +export const ORCA_RUN_TIMEOUT = 'timeout' as const; + +// ── Error policies ────────────────────────────────────────────────────── +// v0.0 implements CONTINUE and ABORT_RUN. The remaining policies are +// accepted in the type and treated as CONTINUE. + +export const ORCA_ON_ERROR_CONTINUE = 'continue' as const; +/** @v0.1+ Accepted, treated as CONTINUE in v0.0. */ +export const ORCA_ON_ERROR_ABORT_ACTION = 'abort-action' as const; +/** @v0.1+ Accepted, treated as CONTINUE in v0.0. */ +export const ORCA_ON_ERROR_ABORT_STAGE = 'abort-stage' as const; +export const ORCA_ON_ERROR_ABORT_RUN = 'abort-run' as const; + +// ── Diagnostic event names ────────────────────────────────────────────── + +export const ORCA_DIAGNOSTIC_EVENTS = { + RUN_STARTED: 'orca.run.started', + RUN_COMPLETED: 'orca.run.completed', + RUN_ABORTED: 'orca.run.aborted', + ACTION_STARTED: 'orca.action.started', + ACTION_COMPLETED: 'orca.action.completed', + ACTION_FAILED: 'orca.action.failed', + ACTION_SKIPPED: 'orca.action.skipped', + CONFIGURATION_INVALID: 'orca.configuration.invalid' +} as const; + +// ── Log messages ──────────────────────────────────────────────────────── + +export const ORCA_LOG_MSG_RUN_STARTED = 'orca run started'; +export const ORCA_LOG_MSG_RUN_COMPLETED = 'orca run completed'; +export const ORCA_LOG_MSG_RUN_ABORTED = 'orca run aborted'; +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_CONFIGURATION_INVALID = 'orca configuration invalid'; + +// ── Defaults ──────────────────────────────────────────────────────────── + +export const ORCA_DEFAULT_MAX_RUNS = 256; diff --git a/src/arts/orca/diagnostics.ts b/src/arts/orca/diagnostics.ts new file mode 100644 index 0000000..cfaa962 --- /dev/null +++ b/src/arts/orca/diagnostics.ts @@ -0,0 +1,121 @@ +import { + LogLevel, + createCatalogDiagnostics, + type DiagnosticCatalog, + type DiagnosticEvent, + type Diagnostics, + type Logger +} from '$libs/logger'; +import { + ORCA_DIAGNOSTIC_EVENTS, + ORCA_LOG_MSG_ACTION_COMPLETED, + ORCA_LOG_MSG_ACTION_FAILED, + ORCA_LOG_MSG_ACTION_SKIPPED, + ORCA_LOG_MSG_ACTION_STARTED, + ORCA_LOG_MSG_CONFIGURATION_INVALID, + ORCA_LOG_MSG_RUN_ABORTED, + ORCA_LOG_MSG_RUN_COMPLETED, + ORCA_LOG_MSG_RUN_STARTED, + ORCA_MODULE +} from './consts.ts'; +import type { OrcaActionId, OrcaRunId, OrcaStage } from './types.ts'; + +export type OrcaDiagnosticType = + (typeof ORCA_DIAGNOSTIC_EVENTS)[keyof typeof ORCA_DIAGNOSTIC_EVENTS]; + +export type OrcaDiagnosticMeta = + | { + readonly runId: OrcaRunId; + readonly event: string; + } + | { + readonly runId: OrcaRunId; + readonly event: string; + readonly status: string; + readonly durationMs: number; + readonly actionCount: number; + } + | { + readonly runId: OrcaRunId; + readonly event: string; + readonly cause: unknown; + } + | { + readonly runId: OrcaRunId; + readonly actionId: OrcaActionId; + readonly stage: OrcaStage; + } + | { + readonly runId: OrcaRunId; + readonly actionId: OrcaActionId; + readonly durationMs: number; + } + | { + readonly runId: OrcaRunId; + readonly actionId: OrcaActionId; + readonly error: unknown; + readonly thrown?: boolean; + } + | { + readonly runId: OrcaRunId; + readonly actionId: OrcaActionId; + readonly reason?: string; + } + | { + readonly reason: string; + }; + +export type OrcaDiagnosticEvent = DiagnosticEvent; +export type OrcaDiagnostics = Diagnostics; + +const ORCA_DIAGNOSTIC_LOGS: DiagnosticCatalog = { + [ORCA_DIAGNOSTIC_EVENTS.RUN_STARTED]: { level: LogLevel.DEBUG, message: ORCA_LOG_MSG_RUN_STARTED }, + [ORCA_DIAGNOSTIC_EVENTS.RUN_COMPLETED]: { + level: LogLevel.DEBUG, + message: ORCA_LOG_MSG_RUN_COMPLETED + }, + [ORCA_DIAGNOSTIC_EVENTS.RUN_ABORTED]: { + level: LogLevel.WARN, + message: ORCA_LOG_MSG_RUN_ABORTED + }, + [ORCA_DIAGNOSTIC_EVENTS.ACTION_STARTED]: { + level: LogLevel.TRACE, + message: ORCA_LOG_MSG_ACTION_STARTED + }, + [ORCA_DIAGNOSTIC_EVENTS.ACTION_COMPLETED]: { + level: LogLevel.TRACE, + message: ORCA_LOG_MSG_ACTION_COMPLETED + }, + [ORCA_DIAGNOSTIC_EVENTS.ACTION_FAILED]: { + level: LogLevel.ERROR, + message: ORCA_LOG_MSG_ACTION_FAILED + }, + [ORCA_DIAGNOSTIC_EVENTS.ACTION_SKIPPED]: { + level: LogLevel.DEBUG, + message: ORCA_LOG_MSG_ACTION_SKIPPED + }, + [ORCA_DIAGNOSTIC_EVENTS.CONFIGURATION_INVALID]: { + level: LogLevel.ERROR, + message: ORCA_LOG_MSG_CONFIGURATION_INVALID + } +}; + +export function createOrcaDiagnostics(logger?: Logger): OrcaDiagnostics { + return createCatalogDiagnostics({ + logger, + defaultCategory: ORCA_MODULE, + catalog: ORCA_DIAGNOSTIC_LOGS + }); +} + +export function emitOrcaDiagnostic( + diagnostics: OrcaDiagnostics, + type: OrcaDiagnosticType, + meta: OrcaDiagnosticMeta +): void { + diagnostics.emit({ + artifact: ORCA_MODULE, + type, + meta + }); +} diff --git a/src/arts/orca/engine-orca.ts b/src/arts/orca/engine-orca.ts new file mode 100644 index 0000000..8b9a5ff --- /dev/null +++ b/src/arts/orca/engine-orca.ts @@ -0,0 +1,431 @@ +/** + * `EngineOrca` v0.0 — orchestration engine. + * + * 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. + * + * 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. + */ + +import { + ORCA_ACTION_STATUS_BLOCKED, + ORCA_ACTION_STATUS_ERROR, + ORCA_ACTION_STATUS_SKIPPED, + ORCA_ACTION_STATUS_SUCCESS, + ORCA_DEFAULT_MAX_RUNS, + ORCA_DIAGNOSTIC_EVENTS, + ORCA_ON_ERROR_ABORT_RUN, + ORCA_ON_ERROR_CONTINUE, + ORCA_RESULT_ERROR, + ORCA_RESULT_FATAL, + ORCA_RESULT_SKIPPED, + ORCA_RESULT_SUCCESS, + ORCA_RUN_ABORTED, + ORCA_RUN_PARTIAL, + ORCA_RUN_SUCCESS, + ORCA_STAGE_FINALLY, + ORCA_STAGES_CANONICAL_ORDER +} from './consts.ts'; +import { + OrcaDisposedError, + OrcaDuplicateActionIdError, + OrcaInvalidActionError, + OrcaInvalidStageError +} from './errors.ts'; +import { createOrcaDiagnostics, emitOrcaDiagnostic } from './diagnostics.ts'; +import type { + EngineOrca, + EngineOrcaOptions, + OrcaAction, + OrcaActionContext, + OrcaActionRun, + OrcaResult, + OrcaRunId, + OrcaRunResult, + OrcaStage, + OrcaToken +} from './types.ts'; +import type { Logger } from '$libs/logger'; +import type { BusSubscription } from '$libs/bus'; + +interface RegisteredAction { + readonly action: OrcaAction; + readonly registeredAt: number; +} + +export function createEngineOrca(options: EngineOrcaOptions): EngineOrca { + const { bus, timers } = options; + const maxRuns = options.maxRuns ?? ORCA_DEFAULT_MAX_RUNS; + const diagnostics = createOrcaDiagnostics(options.logger); + const actionLogger: Logger = options.logger ?? createNoopLogger(); + + const actionsByEvent = new Map(); + const busSubscriptions = new Map(); + const recentRuns: OrcaRunResult[] = []; + const runQueue: Array<() => Promise> = []; + + let disposed = false; + let running = false; + let activeRunController: AbortController | null = null; + let registrationCounter = 0; + + function ensureNotDisposed(): void { + if (disposed) throw new OrcaDisposedError(); + } + + function validateAction(action: OrcaAction): void { + if (action === null || typeof action !== 'object') { + throw new OrcaInvalidActionError('action must be an object'); + } + if (typeof action.id !== 'string' || action.id.length === 0) { + throw new OrcaInvalidActionError('action.id must be a non-empty string'); + } + if (typeof action.action !== 'function') { + throw new OrcaInvalidActionError('action.action must be a function'); + } + if (!ORCA_STAGES_CANONICAL_ORDER.includes(action.stage)) { + throw new OrcaInvalidStageError(action.stage as string); + } + } + + function onEvent( + event: string, + action: OrcaAction + ): () => void { + ensureNotDisposed(); + validateAction(action as OrcaAction); + + const list = actionsByEvent.get(event) ?? []; + if (list.some((entry) => entry.action.id === action.id)) { + throw new OrcaDuplicateActionIdError(event, action.id); + } + + registrationCounter += 1; + list.push({ action: action as OrcaAction, registeredAt: registrationCounter }); + actionsByEvent.set(event, list); + + // 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); + }); + busSubscriptions.set(event, subscription); + } + + return () => detach(event, action.id); + } + + function detach(event: string, actionId: string): void { + const current = actionsByEvent.get(event); + if (!current) return; + const filtered = current.filter((entry) => entry.action.id !== actionId); + if (filtered.length === 0) { + actionsByEvent.delete(event); + busSubscriptions.get(event)?.unsubscribe(); + busSubscriptions.delete(event); + } else { + actionsByEvent.set(event, filtered); + } + } + + function enqueueRun(event: string, payload: unknown): void { + if (disposed) return; + const registered = actionsByEvent.get(event); + if (!registered || registered.length === 0) return; + + // Snapshot taken at dispatch. Actions registered during the run + // do not participate in it. + const snapshot = registered.map((entry) => entry.action); + + runQueue.push(() => executeRun(event, payload, snapshot)); + void drainQueue(); + } + + async function drainQueue(): Promise { + if (running || disposed) return; + const next = runQueue.shift(); + if (!next) return; + + running = true; + try { + await next(); + } finally { + running = false; + if (!disposed && runQueue.length > 0) { + queueMicrotask(() => { + void drainQueue(); + }); + } + } + } + + async function executeRun( + event: string, + payload: unknown, + actions: OrcaAction[] + ): Promise { + const runId = generateRunId(); + const startedAt = timers.clock.now(); + const tokens = new Set(); + const actionRuns: OrcaActionRun[] = []; + const controller = new AbortController(); + activeRunController = controller; + + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_STARTED, { + runId, + event + }); + + let aborted = false; + const byStage = groupByStage(actions); + + stageLoop: for (const stage of ORCA_STAGES_CANONICAL_ORDER) { + const stageActions = byStage.get(stage); + if (!stageActions || stageActions.length === 0) continue; + + // FINALLY always runs, even after abort. + if (aborted && stage !== ORCA_STAGE_FINALLY) continue; + + for (const action of stageActions) { + if (controller.signal.aborted && stage !== ORCA_STAGE_FINALLY) break; + if (disposed) break stageLoop; + + const actionRun = await runAction(action, payload, { + runId, + event, + stage, + tokens, + signal: controller.signal, + logger: actionLogger + }); + + actionRuns.push(actionRun); + + if (actionRun.status === ORCA_ACTION_STATUS_ERROR) { + const policy = action.onError ?? ORCA_ON_ERROR_CONTINUE; + if (policy === ORCA_ON_ERROR_ABORT_RUN) { + aborted = true; + controller.abort(); + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_ABORTED, { + runId, + event, + cause: actionRun.error + }); + break; + } + } + } + } + + const endedAt = timers.clock.now(); + const status = computeRunStatus(aborted, actionRuns); + + const runResult: OrcaRunResult = { + id: runId, + event, + status, + startedAt, + endedAt, + durationMs: endedAt - startedAt, + tokens: Array.from(tokens), + actions: actionRuns + }; + + recordRun(runResult); + activeRunController = null; + + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.RUN_COMPLETED, { + runId, + event, + status, + durationMs: runResult.durationMs, + actionCount: actionRuns.length + }); + } + + async function runAction( + action: OrcaAction, + payload: unknown, + context: OrcaActionContext + ): Promise { + const startedAt = timers.clock.now(); + + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_STARTED, { + runId: context.runId, + actionId: action.id, + stage: action.stage + }); + + // If the run is already aborted (e.g. between stages, only finally + // remains), and we are not in finally, mark blocked. + if (context.signal.aborted && context.stage !== ORCA_STAGE_FINALLY) { + return { + id: action.id, + stage: action.stage, + status: ORCA_ACTION_STATUS_BLOCKED, + startedAt, + endedAt: startedAt, + durationMs: 0, + emitted: [] + }; + } + + try { + const result = await action.action(payload, context); + const endedAt = timers.clock.now(); + const emitted = result.emits ?? []; + for (const token of emitted) (context.tokens as Set).add(token); + + const status = mapResultToActionStatus(result); + + if (status === ORCA_ACTION_STATUS_ERROR) { + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_FAILED, { + runId: context.runId, + actionId: action.id, + 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, + reason: (result as { reason?: string }).reason + }); + } else { + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_COMPLETED, { + runId: context.runId, + actionId: action.id, + durationMs: endedAt - startedAt + }); + } + + return { + id: action.id, + stage: action.stage, + status, + startedAt, + endedAt, + durationMs: endedAt - startedAt, + emitted: Array.from(emitted), + error: + status === ORCA_ACTION_STATUS_ERROR + ? (result as { error?: unknown }).error + : undefined + }; + } catch (thrown) { + // Uncaught exceptions become OrcaError, regardless of intent. + const endedAt = timers.clock.now(); + emitOrcaDiagnostic(diagnostics, ORCA_DIAGNOSTIC_EVENTS.ACTION_FAILED, { + runId: context.runId, + actionId: action.id, + error: thrown, + thrown: true + }); + return { + id: action.id, + stage: action.stage, + status: ORCA_ACTION_STATUS_ERROR, + startedAt, + endedAt, + durationMs: endedAt - startedAt, + emitted: [], + error: thrown + }; + } + } + + function recordRun(run: OrcaRunResult): void { + recentRuns.push(run); + while (recentRuns.length > maxRuns) recentRuns.shift(); + } + + function dispose(): void { + if (disposed) return; + disposed = true; + activeRunController?.abort(); + for (const subscription of busSubscriptions.values()) subscription.unsubscribe(); + busSubscriptions.clear(); + actionsByEvent.clear(); + runQueue.length = 0; + } + + return { + onEvent, + actionCount(event: string): number { + return actionsByEvent.get(event)?.length ?? 0; + }, + recentRuns(): readonly OrcaRunResult[] { + return recentRuns.slice(); + }, + get running(): boolean { + return running; + }, + get disposed(): boolean { + return disposed; + }, + dispose + }; +} + +// ── Helpers ─────────────────────────────────────────────────────────── + +function groupByStage(actions: OrcaAction[]): Map { + const result = new Map(); + for (const action of actions) { + const list = result.get(action.stage) ?? []; + list.push(action); + result.set(action.stage, list); + } + return result; +} + +function mapResultToActionStatus(result: OrcaResult) { + switch (result.status) { + case ORCA_RESULT_SUCCESS: + return ORCA_ACTION_STATUS_SUCCESS; + case ORCA_RESULT_SKIPPED: + return ORCA_ACTION_STATUS_SKIPPED; + case ORCA_RESULT_ERROR: + return ORCA_ACTION_STATUS_ERROR; + case ORCA_RESULT_FATAL: + // v0.0 treats fatal as error. + return ORCA_ACTION_STATUS_ERROR; + default: + // timeout is not produced in v0.0; defensive fallback. + return ORCA_ACTION_STATUS_ERROR; + } +} + +function computeRunStatus(aborted: boolean, actions: OrcaActionRun[]) { + if (aborted) return ORCA_RUN_ABORTED; + const anyError = actions.some((a) => a.status === ORCA_ACTION_STATUS_ERROR); + if (anyError) return ORCA_RUN_PARTIAL; + return ORCA_RUN_SUCCESS; +} + +function generateRunId(): OrcaRunId { + return `run_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`; +} + +function createNoopLogger(): Logger { + const noop = () => {}; + return { + trace: noop, + debug: noop, + info: noop, + warn: noop, + error: noop, + fatal: noop + }; +} + diff --git a/src/arts/orca/errors.ts b/src/arts/orca/errors.ts new file mode 100644 index 0000000..365b542 --- /dev/null +++ b/src/arts/orca/errors.ts @@ -0,0 +1,93 @@ +import { + CodeError, + errCode, + moduleSeed, + type ErrCode, + type ErrorMessages, + type ModuleSeed +} from '$libs/errs'; +import { ORCA_MODULE, ORCA_STAGES_CANONICAL_ORDER } from './consts.ts'; + +// ── Error codes ──────────────────────────────────────────────────────── + +export const ORCA_ERR: ModuleSeed = moduleSeed(ORCA_MODULE); +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'); + +// ── Error message strings ────────────────────────────────────────────── + +const validStages = ORCA_STAGES_CANONICAL_ORDER.join(', '); + +export const ORCA_ERROR_MSG_DISPOSED = `[${ORCA_MODULE}] engine has been disposed`; + +export const orcaDuplicateActionIdMessage = (event: string, actionId: string): string => + `[${ORCA_MODULE}] action "${actionId}" already registered for event "${event}"`; + +export const orcaInvalidStageMessage = (stage: string): string => + `[${ORCA_MODULE}] unknown stage "${stage}". Valid stages: ${validStages}`; + +export const orcaInvalidActionMessage = (reason: string): string => + `[${ORCA_MODULE}] invalid action: ${reason}`; + +// ── 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 +}; + +// ── Error classes ────────────────────────────────────────────────────── + +export class OrcaDisposedError extends CodeError { + constructor() { + super(ORCA_ERR_DISPOSED, { message: ORCA_ERROR_MSG_DISPOSED }); + } +} + +export class OrcaDuplicateActionIdError extends CodeError { + readonly event: string; + readonly actionId: string; + constructor(event: string, actionId: string) { + super(ORCA_ERR_DUPLICATE_ACTION_ID, { + message: orcaDuplicateActionIdMessage(event, actionId) + }); + this.event = event; + this.actionId = actionId; + } +} + +export class OrcaInvalidStageError extends CodeError { + readonly stage: string; + constructor(stage: string) { + super(ORCA_ERR_INVALID_STAGE, { message: orcaInvalidStageMessage(stage) }); + this.stage = stage; + } +} + +export class OrcaInvalidActionError extends CodeError { + constructor(reason: string) { + super(ORCA_ERR_INVALID_ACTION, { message: orcaInvalidActionMessage(reason) }); + } +} + +// ── Type guards ──────────────────────────────────────────────────────── + +export function isOrcaDisposedError(error: unknown): error is OrcaDisposedError { + return error instanceof OrcaDisposedError; +} + +export function isOrcaDuplicateActionIdError(error: unknown): error is OrcaDuplicateActionIdError { + return error instanceof OrcaDuplicateActionIdError; +} + +export function isOrcaInvalidStageError(error: unknown): error is OrcaInvalidStageError { + return error instanceof OrcaInvalidStageError; +} + +export function isOrcaInvalidActionError(error: unknown): error is OrcaInvalidActionError { + return error instanceof OrcaInvalidActionError; +} diff --git a/src/arts/orca/index.ts b/src/arts/orca/index.ts new file mode 100644 index 0000000..d0653a7 --- /dev/null +++ b/src/arts/orca/index.ts @@ -0,0 +1,90 @@ +// Public surface of the orca artifact (v0.0). + +export { createEngineOrca } from './engine-orca.ts'; +export { + orcaSuccess, + orcaSkipped, + orcaError, + orcaTimeout, + orcaFatal +} from './result.ts'; +export { createOrcaDiagnostics, emitOrcaDiagnostic } from './diagnostics.ts'; + +export { + ORCA_MODULE, + ORCA_STAGE_GUARD, + ORCA_STAGE_PRE, + ORCA_STAGE_MAIN, + ORCA_STAGE_POST, + ORCA_STAGE_CLEANUP, + ORCA_STAGE_FINALLY, + ORCA_STAGES_CANONICAL_ORDER, + ORCA_RESULT_SUCCESS, + ORCA_RESULT_SKIPPED, + ORCA_RESULT_ERROR, + 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_TIMEOUT, + ORCA_ACTION_STATUS_FATAL, + ORCA_RUN_SUCCESS, + ORCA_RUN_PARTIAL, + ORCA_RUN_ABORTED, + 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_DIAGNOSTIC_EVENTS, + ORCA_DEFAULT_MAX_RUNS +} from './consts.ts'; + +export { + ORCA_ERR, + ORCA_ERR_DISPOSED, + ORCA_ERR_DUPLICATE_ACTION_ID, + ORCA_ERR_INVALID_STAGE, + ORCA_ERR_INVALID_ACTION, + ORCA_ERROR_MESSAGES, + OrcaDisposedError, + OrcaDuplicateActionIdError, + OrcaInvalidStageError, + OrcaInvalidActionError, + isOrcaDisposedError, + isOrcaDuplicateActionIdError, + isOrcaInvalidStageError, + isOrcaInvalidActionError +} from './errors.ts'; + +export type { + EngineOrca, + EngineOrcaOptions, + OrcaAction, + OrcaActionContext, + OrcaActionFn, + OrcaActionId, + OrcaActionRun, + OrcaBus, + OrcaError, + OrcaErrorPolicy, + OrcaFatal, + OrcaResult, + OrcaRunId, + OrcaRunResult, + OrcaSkipped, + OrcaStage, + OrcaSuccess, + OrcaTimeout, + OrcaToken +} from './types.ts'; + +export type { + OrcaDiagnosticEvent, + OrcaDiagnosticMeta, + OrcaDiagnosticType, + OrcaDiagnostics +} from './diagnostics.ts'; diff --git a/src/arts/orca/result.ts b/src/arts/orca/result.ts new file mode 100644 index 0000000..905cb10 --- /dev/null +++ b/src/arts/orca/result.ts @@ -0,0 +1,84 @@ +/** + * Helpers to construct `OrcaResult` values without writing the literal + * shape at every action site. Action authors return `orcaSuccess()`, + * `orcaError(error)`, etc. The engine inspects `.status` to route the + * result. + */ + +import { + ORCA_RESULT_SUCCESS, + ORCA_RESULT_SKIPPED, + ORCA_RESULT_ERROR, + ORCA_RESULT_TIMEOUT, + ORCA_RESULT_FATAL +} from './consts.ts'; +import type { + OrcaSuccess, + OrcaSkipped, + OrcaError, + OrcaTimeout, + OrcaFatal, + OrcaToken +} from './types.ts'; + +export function orcaSuccess( + options: { value?: TValue; emits?: readonly OrcaToken[] } = {} +): OrcaSuccess { + return { + ok: true, + status: ORCA_RESULT_SUCCESS, + value: options.value, + emits: options.emits + }; +} + +export function orcaSkipped( + reason?: string, + options: { emits?: readonly OrcaToken[] } = {} +): OrcaSkipped { + return { + ok: true, + status: ORCA_RESULT_SKIPPED, + reason, + emits: options.emits + }; +} + +export function orcaError( + error: unknown, + options: { emits?: readonly OrcaToken[]; recoverable?: boolean } = {} +): OrcaError { + return { + ok: false, + status: ORCA_RESULT_ERROR, + error, + recoverable: options.recoverable, + emits: options.emits + }; +} + +/** @v0.1+ */ +export function orcaTimeout( + timeoutMs: number, + options: { emits?: readonly OrcaToken[] } = {} +): OrcaTimeout { + return { + ok: false, + status: ORCA_RESULT_TIMEOUT, + timeoutMs, + emits: options.emits + }; +} + +/** @v0.1+ */ +export function orcaFatal( + error: unknown, + options: { emits?: readonly OrcaToken[] } = {} +): OrcaFatal { + return { + ok: false, + status: ORCA_RESULT_FATAL, + error, + emits: options.emits + }; +} diff --git a/src/arts/orca/test/engine-orca.test.ts b/src/arts/orca/test/engine-orca.test.ts new file mode 100644 index 0000000..335058a --- /dev/null +++ b/src/arts/orca/test/engine-orca.test.ts @@ -0,0 +1,861 @@ +import { afterEach, beforeEach, describe, expect, it, vi } from 'vitest'; +import { createEngineOrca } from '../engine-orca.ts'; +import { + OrcaDisposedError, + OrcaDuplicateActionIdError, + OrcaInvalidActionError, + OrcaInvalidStageError +} from '../errors.ts'; +import { orcaError, orcaSkipped, orcaSuccess } from '../result.ts'; +import { + ORCA_ACTION_STATUS_BLOCKED, + ORCA_ACTION_STATUS_ERROR, + ORCA_ACTION_STATUS_SKIPPED, + ORCA_ACTION_STATUS_SUCCESS, + ORCA_ON_ERROR_ABORT_ACTION, + ORCA_ON_ERROR_ABORT_RUN, + ORCA_ON_ERROR_ABORT_STAGE, + ORCA_ON_ERROR_CONTINUE, + ORCA_RUN_ABORTED, + ORCA_RUN_PARTIAL, + ORCA_RUN_SUCCESS, + ORCA_STAGE_CLEANUP, + ORCA_STAGE_FINALLY, + ORCA_STAGE_GUARD, + ORCA_STAGE_MAIN, + ORCA_STAGE_POST, + ORCA_STAGE_PRE +} from '../consts.ts'; +import type { OrcaAction, OrcaBus, OrcaToken } from '../types.ts'; +import type { TimerScheduler } from '$libs/timer'; +import type { + BusEnvelope, + BusListenOptions, + BusListener, + BusSubscription +} from '$libs/bus'; + +// ── Fakes ────────────────────────────────────────────────────────────── + +interface FakeBus extends OrcaBus { + publish(event: string, payload: unknown): void; + listenerCount(event: string): number; +} + +function createFakeBus(): FakeBus { + const listeners = new Map>>(); + let id = 0; + + function on( + event: string, + listener: BusListener, + _options?: BusListenOptions + ): BusSubscription { + const subId = `sub_${++id}`; + const map = listeners.get(event) ?? new Map(); + map.set(subId, listener); + listeners.set(event, map); + return { + id: subId, + type: event, + get active() { + return listeners.get(event)?.has(subId) ?? false; + }, + unsubscribe() { + const m = listeners.get(event); + if (m) { + m.delete(subId); + if (m.size === 0) listeners.delete(event); + } + } + }; + } + + return { + on: on as unknown as OrcaBus['on'], + once: on as unknown as OrcaBus['once'], + publish(event: string, payload: unknown) { + const map = listeners.get(event); + if (!map) return; + const envelope: BusEnvelope = { + id: `env_${event}_${Date.now()}`, + type: event, + payload, + at: Date.now(), + source: 'test' + }; + const ctx = { signal: new AbortController().signal }; + for (const fn of map.values()) { + void fn(envelope, ctx); + } + }, + listenerCount(event: string) { + return listeners.get(event)?.size ?? 0; + } + }; +} + +function createFakeTimers(initialNow = 0): TimerScheduler & { advance(ms: number): void } { + let now = initialNow; + return { + clock: { + now() { + return now; + } + }, + size: 0, + schedule() { + throw new Error('not implemented in fake'); + }, + scheduleAt() { + throw new Error('not implemented in fake'); + }, + interval() { + throw new Error('not implemented in fake'); + }, + cancel() { + return false; + }, + cancelAll() { + return 0; + }, + has() { + return false; + }, + dispose() {}, + advance(ms: number) { + now += ms; + } + } as unknown as TimerScheduler & { advance(ms: number): void }; +} + +// ── Helpers ──────────────────────────────────────────────────────────── + +const noopAction: OrcaAction = { + id: 'noop', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess() +}; + +async function flush(rounds = 4) { + // Multiple rounds: bus listener → enqueueRun → drainQueue → executeRun and + // any await chain inside actions. Each round yields to both microtasks + // and the macrotask queue so async actions awaiting setTimeout settle. + for (let i = 0; i < rounds; i++) { + await new Promise((r) => queueMicrotask(() => r(undefined))); + await new Promise((r) => setTimeout(r, 0)); + } +} + +async function sleep(ms: number) { + await new Promise((r) => setTimeout(r, ms)); +} + +// ── Tests ────────────────────────────────────────────────────────────── + +describe('EngineOrca v0.0 — registration', () => { + let bus: FakeBus; + let timers: TimerScheduler; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers(); + }); + + it('registers and detaches actions per event', () => { + const orca = createEngineOrca({ bus, timers }); + const detach = orca.onEvent('e1', { ...noopAction, id: 'a1' }); + expect(orca.actionCount('e1')).toBe(1); + detach(); + expect(orca.actionCount('e1')).toBe(0); + }); + + it('throws OrcaDuplicateActionIdError on duplicate id within same event', () => { + const orca = createEngineOrca({ bus, timers }); + orca.onEvent('e1', { ...noopAction, id: 'a1' }); + expect(() => orca.onEvent('e1', { ...noopAction, id: 'a1' })).toThrow( + OrcaDuplicateActionIdError + ); + }); + + it('allows the same id under different events', () => { + const orca = createEngineOrca({ bus, timers }); + orca.onEvent('e1', { ...noopAction, id: 'a1' }); + orca.onEvent('e2', { ...noopAction, id: 'a1' }); + expect(orca.actionCount('e1')).toBe(1); + expect(orca.actionCount('e2')).toBe(1); + }); + + it('throws OrcaInvalidStageError on unknown stage', () => { + const orca = createEngineOrca({ bus, timers }); + expect(() => + orca.onEvent('e1', { + id: 'a1', + stage: 'invalid' as never, + action: () => orcaSuccess() + }) + ).toThrow(OrcaInvalidStageError); + }); + + it('throws OrcaInvalidActionError on missing id or action', () => { + const orca = createEngineOrca({ bus, timers }); + expect(() => + orca.onEvent('e1', { id: '', stage: ORCA_STAGE_MAIN, action: () => orcaSuccess() }) + ).toThrow(OrcaInvalidActionError); + expect(() => + orca.onEvent('e1', { + id: 'a1', + stage: ORCA_STAGE_MAIN, + action: undefined as unknown as () => never + }) + ).toThrow(OrcaInvalidActionError); + }); + + it('subscribes to the bus on first registration per event', () => { + const orca = createEngineOrca({ bus, timers }); + expect(bus.listenerCount('e1')).toBe(0); + orca.onEvent('e1', { ...noopAction, id: 'a1' }); + expect(bus.listenerCount('e1')).toBe(1); + // Second action reuses the same subscription + orca.onEvent('e1', { ...noopAction, id: 'a2' }); + expect(bus.listenerCount('e1')).toBe(1); + }); + + it('unsubscribes from the bus when the last action for the event is removed', () => { + const orca = createEngineOrca({ bus, timers }); + const detachA = orca.onEvent('e1', { ...noopAction, id: 'a1' }); + const detachB = orca.onEvent('e1', { ...noopAction, id: 'a2' }); + expect(bus.listenerCount('e1')).toBe(1); + detachA(); + expect(bus.listenerCount('e1')).toBe(1); + detachB(); + expect(bus.listenerCount('e1')).toBe(0); + }); +}); + +describe('EngineOrca v0.0 — execution', () => { + let bus: FakeBus; + let timers: TimerScheduler & { advance(ms: number): void }; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers() as TimerScheduler & { advance(ms: number): void }; + }); + + it('executes actions in canonical stage order', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'cleanup-action', + stage: ORCA_STAGE_CLEANUP, + action: () => { + calls.push('cleanup'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'main-action', + stage: ORCA_STAGE_MAIN, + action: () => { + calls.push('main'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'guard-action', + stage: ORCA_STAGE_GUARD, + action: () => { + calls.push('guard'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'finally-action', + stage: ORCA_STAGE_FINALLY, + action: () => { + calls.push('finally'); + return orcaSuccess(); + } + }); + + bus.publish('e1', { foo: 'bar' }); + await flush(); + + expect(calls).toEqual(['guard', 'main', 'cleanup', 'finally']); + }); + + it('executes same-stage actions in registration order', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'first', + stage: ORCA_STAGE_MAIN, + action: () => { + calls.push('first'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'second', + stage: ORCA_STAGE_MAIN, + action: () => { + calls.push('second'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'third', + stage: ORCA_STAGE_MAIN, + action: () => { + calls.push('third'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['first', 'second', 'third']); + }); + + it('captures thrown exceptions as OrcaError', async () => { + const orca = createEngineOrca({ bus, timers }); + orca.onEvent('e1', { + id: 'thrower', + stage: ORCA_STAGE_MAIN, + action: () => { + throw new Error('boom'); + } + }); + + bus.publish('e1', null); + await flush(); + + const runs = orca.recentRuns(); + expect(runs).toHaveLength(1); + expect(runs[0].actions[0].status).toBe(ORCA_ACTION_STATUS_ERROR); + expect((runs[0].actions[0].error as Error).message).toBe('boom'); + }); + + it('aggregates emitted tokens into the run context', async () => { + const orca = createEngineOrca({ bus, timers }); + const seen: ReadonlySet[] = []; + + orca.onEvent('e1', { + id: 'first', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess({ emits: ['token-a'] }) + }); + orca.onEvent('e1', { + id: 'second', + stage: ORCA_STAGE_POST, + action: (_payload, ctx) => { + seen.push(new Set(ctx.tokens)); + return orcaSuccess({ emits: ['token-b'] }); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(seen[0].has('token-a')).toBe(true); + + const runs = orca.recentRuns(); + expect([...runs[0].tokens].sort()).toEqual(['token-a', 'token-b']); + }); + + it('runs FINALLY even after run is aborted', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'main', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('oops')) + }); + orca.onEvent('e1', { + id: 'post', + stage: ORCA_STAGE_POST, + action: () => { + calls.push('post'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'finally', + stage: ORCA_STAGE_FINALLY, + action: () => { + calls.push('finally'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['finally']); + const runs = orca.recentRuns(); + expect(runs[0].status).toBe(ORCA_RUN_ABORTED); + }); + + it('does not include actions registered during a run in that same run', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'first', + stage: ORCA_STAGE_MAIN, + action: () => { + calls.push('first'); + orca.onEvent('e1', { + id: 'late', + stage: ORCA_STAGE_POST, + action: () => { + calls.push('late'); + return orcaSuccess(); + } + }); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['first']); + + // On the next event, 'late' is registered and runs. + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['first', 'first', 'late']); + }); +}); + +describe('EngineOrca v0.0 — error policies', () => { + let bus: FakeBus; + let timers: TimerScheduler; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers(); + }); + + it('CONTINUE keeps running subsequent actions after an error', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'erroring', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_CONTINUE, + action: () => orcaError(new Error('x')) + }); + orca.onEvent('e1', { + id: 'after', + stage: ORCA_STAGE_MAIN, + action: () => { + calls.push('after'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['after']); + expect(orca.recentRuns()[0].status).toBe(ORCA_RUN_PARTIAL); + }); + + it('ABORT_RUN stops subsequent actions and marks run as aborted', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'erroring', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('x')) + }); + orca.onEvent('e1', { + id: 'after', + stage: ORCA_STAGE_POST, + action: () => { + calls.push('after'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual([]); + expect(orca.recentRuns()[0].status).toBe(ORCA_RUN_ABORTED); + }); + + it('treats ABORT_ACTION and ABORT_STAGE as CONTINUE in v0.0', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'a1', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_ACTION, + action: () => orcaError(new Error('x')) + }); + orca.onEvent('e1', { + id: 'a2', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_STAGE, + action: () => { + calls.push('a2'); + return orcaSuccess(); + } + }); + orca.onEvent('e1', { + id: 'a3', + stage: ORCA_STAGE_POST, + action: () => { + calls.push('a3'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['a2', 'a3']); + }); +}); + +describe('EngineOrca v0.0 — concurrency', () => { + let bus: FakeBus; + let timers: TimerScheduler; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers(); + }); + + it('queues events of the same type that arrive while a run is in flight', async () => { + const orca = createEngineOrca({ bus, timers }); + let resolveFirst!: () => void; + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'long', + stage: ORCA_STAGE_MAIN, + action: async (payload) => { + calls.push(`start:${(payload as { n: number }).n}`); + await new Promise((r) => { + resolveFirst = r; + }); + calls.push(`end:${(payload as { n: number }).n}`); + return orcaSuccess(); + } + }); + + bus.publish('e1', { n: 1 }); + bus.publish('e1', { n: 2 }); + + // Let the first run start + await new Promise((r) => queueMicrotask(() => r(undefined))); + await new Promise((r) => queueMicrotask(() => r(undefined))); + + expect(calls).toEqual(['start:1']); + + // Resolve the first; second should start after. + resolveFirst!(); + // Allow the second run to start + await sleep(20); + // Resolve the second one too (it uses the same resolveFirst slot + // because the action body re-binds it). + resolveFirst!(); + await sleep(20); + + expect(calls).toEqual(['start:1', 'end:1', 'start:2', 'end:2']); + }); +}); + +describe('EngineOrca v0.0 — run trace', () => { + let bus: FakeBus; + let timers: TimerScheduler & { advance(ms: number): void }; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers() as TimerScheduler & { advance(ms: number): void }; + }); + + it('produces OrcaRunResult with startedAt / endedAt / durationMs', async () => { + const orca = createEngineOrca({ bus, timers }); + orca.onEvent('e1', { + id: 'a1', + stage: ORCA_STAGE_MAIN, + action: () => { + timers.advance(50); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + const runs = orca.recentRuns(); + expect(runs[0].startedAt).toBeDefined(); + expect(runs[0].endedAt).toBeGreaterThanOrEqual(runs[0].startedAt); + expect(runs[0].durationMs).toBe(runs[0].endedAt - runs[0].startedAt); + }); + + it('lists every executed action with status', async () => { + const orca = createEngineOrca({ bus, timers }); + orca.onEvent('e1', { + id: 'success', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess() + }); + orca.onEvent('e1', { + id: 'skipped', + stage: ORCA_STAGE_MAIN, + action: () => orcaSkipped('reason') + }); + orca.onEvent('e1', { + id: 'errored', + stage: ORCA_STAGE_MAIN, + action: () => orcaError(new Error('x')) + }); + + bus.publish('e1', null); + await flush(); + + const runs = orca.recentRuns(); + expect(runs[0].actions.map((a) => [a.id, a.status])).toEqual([ + ['success', ORCA_ACTION_STATUS_SUCCESS], + ['skipped', ORCA_ACTION_STATUS_SKIPPED], + ['errored', ORCA_ACTION_STATUS_ERROR] + ]); + }); + + it('respects maxRuns in recentRuns()', async () => { + const orca = createEngineOrca({ bus, timers, maxRuns: 3 }); + orca.onEvent('e1', { ...noopAction, id: 'a1' }); + + for (let i = 0; i < 5; i++) { + bus.publish('e1', { i }); + await flush(); + } + + expect(orca.recentRuns()).toHaveLength(3); + }); + + it('marks an action BLOCKED when run is already aborted at the start of finally', async () => { + const orca = createEngineOrca({ bus, timers }); + + orca.onEvent('e1', { + id: 'aborter', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('x')) + }); + // Action in POST should not run; not even reach BLOCKED. + orca.onEvent('e1', { + id: 'post', + stage: ORCA_STAGE_POST, + action: () => orcaSuccess() + }); + + bus.publish('e1', null); + await flush(); + + const runs = orca.recentRuns(); + const ids = runs[0].actions.map((a) => a.id); + expect(ids).not.toContain('post'); + }); +}); + +describe('EngineOrca v0.0 — disposal', () => { + let bus: FakeBus; + let timers: TimerScheduler; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers(); + }); + + it('dispose() is idempotent', () => { + const orca = createEngineOrca({ bus, timers }); + expect(orca.disposed).toBe(false); + orca.dispose(); + expect(orca.disposed).toBe(true); + expect(() => orca.dispose()).not.toThrow(); + expect(orca.disposed).toBe(true); + }); + + it('aborts in-flight run on dispose', async () => { + const orca = createEngineOrca({ bus, timers }); + let signaled = false; + + orca.onEvent('e1', { + id: 'long', + stage: ORCA_STAGE_MAIN, + action: async (_p, ctx) => { + ctx.signal.addEventListener('abort', () => { + signaled = true; + }); + await new Promise((r) => setTimeout(r, 50)); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await new Promise((r) => queueMicrotask(() => r(undefined))); + await new Promise((r) => queueMicrotask(() => r(undefined))); + orca.dispose(); + + expect(signaled).toBe(true); + }); + + it('throws OrcaDisposedError when registering after dispose', () => { + const orca = createEngineOrca({ bus, timers }); + orca.dispose(); + expect(() => orca.onEvent('e1', { ...noopAction, id: 'a1' })).toThrow(OrcaDisposedError); + }); + + it('events received after dispose do not execute actions', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls = vi.fn(); + orca.onEvent('e1', { + id: 'a1', + stage: ORCA_STAGE_MAIN, + action: () => { + calls(); + return orcaSuccess(); + } + }); + orca.dispose(); + bus.publish('e1', null); + await flush(); + expect(calls).not.toHaveBeenCalled(); + }); +}); + +describe('EngineOrca v0.0 — accepted-but-ignored fields (forward-compat)', () => { + let bus: FakeBus; + let timers: TimerScheduler; + + beforeEach(() => { + bus = createFakeBus(); + timers = createFakeTimers(); + }); + + it('accepts after without waiting for tokens', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'declares-after', + stage: ORCA_STAGE_MAIN, + after: ['nonexistent-token'], + action: () => { + calls.push('ran'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['ran']); // ran despite missing token + }); + + it('accepts unless without skipping', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'producer', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess({ emits: ['present'] }) + }); + orca.onEvent('e1', { + id: 'declares-unless', + stage: ORCA_STAGE_POST, + unless: ['present'], + action: () => { + calls.push('ran'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['ran']); // ran despite token being present + }); + + it('accepts abortOn without blocking', async () => { + const orca = createEngineOrca({ bus, timers }); + const calls: string[] = []; + + orca.onEvent('e1', { + id: 'producer', + stage: ORCA_STAGE_MAIN, + action: () => orcaSuccess({ emits: ['danger'] }) + }); + orca.onEvent('e1', { + id: 'declares-abortOn', + stage: ORCA_STAGE_POST, + abortOn: ['danger'], + action: () => { + calls.push('ran'); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await flush(); + + expect(calls).toEqual(['ran']); + }); + + it('accepts actionTimeoutMs without enforcing timeout', async () => { + const orca = createEngineOrca({ bus, timers }); + + orca.onEvent('e1', { + id: 'slow', + stage: ORCA_STAGE_MAIN, + actionTimeoutMs: 1, // ignored + action: async () => { + await sleep(5); + return orcaSuccess(); + } + }); + + bus.publish('e1', null); + await sleep(20); + await flush(); + + const runs = orca.recentRuns(); + expect(runs[0].actions[0].status).toBe(ORCA_ACTION_STATUS_SUCCESS); + }); + + it('accepts compensate without invoking it', async () => { + const orca = createEngineOrca({ bus, timers }); + const compensate = vi.fn(() => orcaSuccess()); + + orca.onEvent('e1', { + id: 'with-compensate', + stage: ORCA_STAGE_MAIN, + onError: ORCA_ON_ERROR_ABORT_RUN, + action: () => orcaError(new Error('x')), + compensate + }); + + bus.publish('e1', null); + await flush(); + + expect(compensate).not.toHaveBeenCalled(); + }); +}); diff --git a/src/arts/orca/test/result.test.ts b/src/arts/orca/test/result.test.ts new file mode 100644 index 0000000..2e2179e --- /dev/null +++ b/src/arts/orca/test/result.test.ts @@ -0,0 +1,56 @@ +import { describe, expect, it } from 'vitest'; +import { orcaError, orcaFatal, orcaSkipped, orcaSuccess, orcaTimeout } from '../result.ts'; +import { + ORCA_RESULT_ERROR, + ORCA_RESULT_FATAL, + ORCA_RESULT_SKIPPED, + ORCA_RESULT_SUCCESS, + ORCA_RESULT_TIMEOUT +} from '../consts.ts'; + +describe('orca result helpers', () => { + it('orcaSuccess() has ok:true and SUCCESS status', () => { + const r = orcaSuccess(); + expect(r.ok).toBe(true); + expect(r.status).toBe(ORCA_RESULT_SUCCESS); + }); + + it('orcaSuccess({ value, emits }) carries value and tokens', () => { + const r = orcaSuccess({ value: 42, emits: ['t1', 't2'] }); + expect(r.value).toBe(42); + expect(r.emits).toEqual(['t1', 't2']); + }); + + it('orcaSkipped(reason) carries reason and SKIPPED status', () => { + const r = orcaSkipped('not applicable'); + expect(r.ok).toBe(true); + expect(r.status).toBe(ORCA_RESULT_SKIPPED); + expect(r.reason).toBe('not applicable'); + }); + + it('orcaError(error) carries error and ERROR status', () => { + const err = new Error('boom'); + const r = orcaError(err); + expect(r.ok).toBe(false); + expect(r.status).toBe(ORCA_RESULT_ERROR); + expect(r.error).toBe(err); + }); + + it('orcaError supports recoverable flag', () => { + const r = orcaError(new Error('x'), { recoverable: true }); + expect(r.recoverable).toBe(true); + }); + + it('orcaTimeout returns TIMEOUT status (v0.1+ shape)', () => { + const r = orcaTimeout(1000); + expect(r.ok).toBe(false); + expect(r.status).toBe(ORCA_RESULT_TIMEOUT); + expect(r.timeoutMs).toBe(1000); + }); + + it('orcaFatal returns FATAL status (v0.1+ shape)', () => { + const r = orcaFatal(new Error('boom')); + expect(r.ok).toBe(false); + expect(r.status).toBe(ORCA_RESULT_FATAL); + }); +}); diff --git a/src/arts/orca/types.ts b/src/arts/orca/types.ts new file mode 100644 index 0000000..c607a53 --- /dev/null +++ b/src/arts/orca/types.ts @@ -0,0 +1,308 @@ +import type { EventSubscriber } from '$libs/bus'; +import type { Logger } from '$libs/logger'; +import type { TimerScheduler } from '$libs/timer'; +import type { + ORCA_STAGE_GUARD, + ORCA_STAGE_PRE, + ORCA_STAGE_MAIN, + ORCA_STAGE_POST, + ORCA_STAGE_CLEANUP, + ORCA_STAGE_FINALLY, + ORCA_RESULT_SUCCESS, + ORCA_RESULT_SKIPPED, + ORCA_RESULT_ERROR, + 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_TIMEOUT, + ORCA_ACTION_STATUS_FATAL, + ORCA_RUN_SUCCESS, + ORCA_RUN_PARTIAL, + ORCA_RUN_ABORTED, + 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 +} from './consts.ts'; + +// ── Identifiers ─────────────────────────────────────────────────────── + +export type OrcaStage = + | typeof ORCA_STAGE_GUARD + | typeof ORCA_STAGE_PRE + | typeof ORCA_STAGE_MAIN + | typeof ORCA_STAGE_POST + | typeof ORCA_STAGE_CLEANUP + | typeof ORCA_STAGE_FINALLY; + +export type OrcaToken = string; +export type OrcaActionId = string; +export type OrcaRunId = string; + +export type OrcaErrorPolicy = + | typeof ORCA_ON_ERROR_CONTINUE + | typeof ORCA_ON_ERROR_ABORT_ACTION + | typeof ORCA_ON_ERROR_ABORT_STAGE + | typeof ORCA_ON_ERROR_ABORT_RUN; + +// ── Result types ────────────────────────────────────────────────────── + +export interface OrcaSuccess { + readonly ok: true; + readonly status: typeof ORCA_RESULT_SUCCESS; + readonly value?: TValue; + readonly emits?: readonly OrcaToken[]; +} + +export interface OrcaSkipped { + readonly ok: true; + readonly status: typeof ORCA_RESULT_SKIPPED; + readonly reason?: string; + readonly emits?: readonly OrcaToken[]; +} + +export interface OrcaError { + readonly ok: false; + readonly status: typeof ORCA_RESULT_ERROR; + readonly error: unknown; + readonly recoverable?: boolean; + readonly emits?: readonly OrcaToken[]; +} + +/** + * @v0.1+ Returned by the engine when an action exceeds `actionTimeoutMs`. + * @v0.0 Shape exists so action authors can type results, but the engine + * never produces one (timeouts are not enforced). + */ +export interface OrcaTimeout { + readonly ok: false; + readonly status: typeof ORCA_RESULT_TIMEOUT; + readonly timeoutMs: number; + readonly emits?: readonly OrcaToken[]; +} + +/** + * @v0.1+ Distinguished from `OrcaError` for run-aborting failures. + * @v0.0 Engine treats fatal as error. + */ +export interface OrcaFatal { + readonly ok: false; + readonly status: typeof ORCA_RESULT_FATAL; + readonly error: unknown; + readonly emits?: readonly OrcaToken[]; +} + +export type OrcaResult = + | OrcaSuccess + | OrcaSkipped + | OrcaError + | OrcaTimeout + | OrcaFatal; + +// ── Action context ──────────────────────────────────────────────────── + +export interface OrcaActionContext { + readonly runId: OrcaRunId; + readonly event: string; + readonly stage: OrcaStage; + /** + * Tokens already emitted in the current run by previous actions. + * @v0.0 Populated correctly, but not consumed by the engine + * (after/unless/abortOn are ignored). Action authors MAY read it. + */ + readonly tokens: ReadonlySet; + /** + * 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. + */ + readonly signal: AbortSignal; + /** + * Logger for ad-hoc diagnostics inside the action. Catalogued events + * are emitted by the engine itself. + */ + readonly logger: Logger; +} + +// ── Action definition ───────────────────────────────────────────────── + +export type OrcaActionFn = ( + payload: TPayload, + context: OrcaActionContext +) => OrcaResult | Promise>; + +export interface OrcaAction { + readonly id: OrcaActionId; + readonly stage: OrcaStage; + + /** + * @v0.1+ Tokens this action waits for before running. + * @v0.0 Accepted in type, ignored by engine. Authors should still + * declare it accurately for forward compatibility. + */ + readonly after?: readonly OrcaToken[]; + + /** + * @v0.1+ If any of these tokens is present when this action would run, + * the action is skipped. + * @v0.0 Accepted, ignored. + */ + readonly unless?: readonly OrcaToken[]; + + /** + * @v0.1+ If any of these tokens is emitted before this action runs, + * the action is blocked (not skipped — appears in trace). + * @v0.0 Accepted, ignored. + */ + readonly abortOn?: readonly OrcaToken[]; + + /** + * Tokens this action may emit. Documented for consumers; honored by + * the engine (the engine adds emitted tokens to the run context + * regardless of this declaration — `provides` is metadata for + * static validation, future v0.1+). + */ + readonly provides?: readonly OrcaToken[]; + + /** + * @v0.1+ Maximum wall-clock duration of this action before it is + * cancelled and reported as ORCA_RESULT_TIMEOUT. + * @v0.0 Accepted, ignored. Long actions hang. + */ + readonly actionTimeoutMs?: number; + + /** + * @v0.0 Honored. Only CONTINUE and ABORT_RUN are distinct; + * ABORT_ACTION and ABORT_STAGE behave as CONTINUE. + */ + readonly onError?: OrcaErrorPolicy; + + /** + * @v1+ Compensating action invoked when the run aborts after this + * action completed successfully. + * @v0.0 Accepted, never invoked. + */ + readonly compensate?: OrcaActionFn; + + readonly action: OrcaActionFn; +} + +// ── Run trace ───────────────────────────────────────────────────────── + +export interface OrcaActionRun { + readonly id: OrcaActionId; + readonly stage: OrcaStage; + readonly status: + | typeof ORCA_ACTION_STATUS_SUCCESS + | typeof ORCA_ACTION_STATUS_SKIPPED + | typeof ORCA_ACTION_STATUS_BLOCKED + | typeof ORCA_ACTION_STATUS_ERROR + | typeof ORCA_ACTION_STATUS_TIMEOUT + | typeof ORCA_ACTION_STATUS_FATAL; + readonly startedAt: number; + readonly endedAt: number; + readonly durationMs: number; + readonly emitted: readonly OrcaToken[]; + readonly error?: unknown; +} + +export interface OrcaRunResult { + readonly id: OrcaRunId; + readonly event: string; + readonly status: + | typeof ORCA_RUN_SUCCESS + | typeof ORCA_RUN_PARTIAL + | typeof ORCA_RUN_ABORTED + | typeof ORCA_RUN_FATAL + | typeof ORCA_RUN_TIMEOUT; + readonly startedAt: number; + readonly endedAt: number; + readonly durationMs: number; + readonly tokens: readonly OrcaToken[]; + readonly actions: readonly OrcaActionRun[]; +} + +// ── Engine ──────────────────────────────────────────────────────────── + +/** + * orca subscribes to arbitrary event names with `unknown` payloads. The + * default `BusEventMap = object` makes `keyof TEvents & string = never` + * which blocks `bus.on(event, …)` for arbitrary strings, so we widen + * the contract to a record of string keys. + */ +export type OrcaBus = EventSubscriber>; + +export interface EngineOrcaOptions { + /** + * Subscriber for bus events. orca only listens; it never publishes. + */ + readonly bus: OrcaBus; + /** + * Source of monotonic time. orca uses `clock.now()` for run/action + * timestamps. Tests inject a fake scheduler with a controllable clock. + */ + readonly timers: TimerScheduler; + /** + * Optional logger for catalogued diagnostics. If omitted, the engine + * runs silently. Action context still receives a logger (a noop one + * when this option is missing). + */ + readonly logger?: Logger; + /** + * Maximum number of runs retained in `recentRuns()`. + * @default 256 + */ + readonly maxRuns?: number; +} + +export interface EngineOrca { + /** + * Register an action against an event. The engine subscribes to the + * bus on first registration per event, and unsubscribes when the + * last action for that event is removed. + * + * Returns a detach function. Calling detach removes the action from + * the engine. + * + * Throws `OrcaDisposedError` if the engine has been disposed. + * Throws `OrcaDuplicateActionIdError` if `action.id` is already + * registered for the same event. + * Throws `OrcaInvalidStageError` if `action.stage` is not canonical. + * Throws `OrcaInvalidActionError` if `action.id` or `action.action` + * is missing or invalid. + */ + onEvent( + event: string, + action: OrcaAction + ): () => void; + + /** + * Number of actions registered for a given event. Useful for tests + * and devtools. + */ + actionCount(event: string): number; + + /** + * Snapshot of recent runs. Bounded by `maxRuns`. + */ + recentRuns(): readonly OrcaRunResult[]; + + /** + * True while the engine is currently executing a run. + */ + readonly running: boolean; + + readonly disposed: boolean; + + /** + * Tear down. Idempotent. Cancels any in-flight run by aborting its + * signal. Releases the bus subscriptions. Subsequent calls to + * `onEvent()` throw `OrcaDisposedError`. + */ + dispose(): void; +} diff --git a/svelte.config.js b/svelte.config.js index 070fb02..f79b0aa 100644 --- a/svelte.config.js +++ b/svelte.config.js @@ -29,6 +29,7 @@ const config = { $http: 'src/arts/http', $lang: 'src/arts/lang', $logger: 'src/arts/logger', + $orca: 'src/arts/orca', $perm: 'src/arts/perm', $session: 'src/arts/session', $sium: 'src/arts/sium',