import { createHash } from "node:crypto" import { AutomationCompatibilityError, consumeDeclarationSlot, formatHookLocation, isSameHookLocation, type AutomationRunState, type HookLocation, type RuntimeRunContext, type RuntimeValueContext, } from "./runtime" import { AmbiguousSignalError, ClosedSignalError, combineSignalDependencyOutcomes, createSignal, FailedSignalError, signalAncestorOrigins, signalContextDependencyCollector, signalContextOrigins, signalDependencyCollector, signalIndex, signalOrigins, signalResolver, type KeyedSignal, type Signal, type SignalDependencySets, type SignalIndexDefinition, } from "./signal-protocol" const CORRELATION_ID_INDEX = { getKey: (value: string) => value, } satisfies SignalIndexDefinition /** * Materializes a signal or throws its terminal control-flow marker. * * @param signal - Signal to materialize. * @param context - Exact dependency values loaded for this execution. * @throws {ClosedSignalError} When the signal closed without a value. * @throws {FailedSignalError} When the signal failed. */ export function materializeSignal( signal: Signal, context: RuntimeValueContext, ): T { const resolution = signal[signalResolver](context) if (resolution.status === "closed") throw new ClosedSignalError() if (resolution.status === "failed") { throw new FailedSignalError(resolution.failure) } return resolution.value } /** * Creates a signal for an action output at a deterministic hook slot. * * @param hook - Hook location assigned to the action invocation. * @param contextDependencies - Signals that place the action in a context. */ export function createActionSignal( hook: HookLocation, contextDependencies: readonly Signal[] = [], ): Signal { const formatted = formatHookLocation(hook) return createSignal({ ancestorOrigins: [ ...new Set([ ...contextDependencies.flatMap( (signal) => signal[signalAncestorOrigins], ), `action:${formatted}`, ]), ], collect: (state, dependencies) => { const match = findUniqueBoundary( state.context.actions, (output) => isSameHookLocation(output, hook), `action hook ${formatted}`, state.context, ) if (!match || match.status === "pending") return { status: "pending" } dependencies.actionDependencyIds.add(match.actionInvocationId) if (match.status === "failed") { return { failure: match.failure, outcomeSeq: match.outcomeSeq, status: "failed", } } if (match.status === "skipped") { return match.reason === "closed-dependency" ? { outcomeSeq: match.outcomeSeq, status: "closed" } : { failure: match.failure, outcomeSeq: match.outcomeSeq, status: "failed", } } return { outcomeSeq: match.outcomeSeq, status: "succeeded" } }, collectContext: (state, dependencies) => { if (contextDependencies.length === 0) { const invocation = findUniqueBoundary( state.context.actions, (output) => isSameHookLocation(output, hook), `action hook ${formatted}`, state.context, ) ?? state.actionInvocations.find((candidate) => isSameHookLocation(candidate, hook), ) if (!invocation) return { status: "pending" } for (const id of invocation.actionDependencyIds) { dependencies.actionDependencyIds.add(id) } for (const id of invocation.eventDependencyIds) { dependencies.eventDependencyIds.add(id) } for (const id of invocation.signalDependencyIds) { dependencies.signalDependencyIds.add(id) } if ( invocation.actionDependencyIds.length + invocation.eventDependencyIds.length + invocation.signalDependencyIds.length === 0 ) { return { status: "pending" } } return { outcomeSeq: Math.max( 0, ...state.context.actions .filter((candidate) => invocation.actionDependencyIds.includes( candidate.actionInvocationId, ), ) .flatMap((candidate) => candidate.status === "pending" ? [] : [candidate.outcomeSeq], ), ...state.context.events .filter((candidate) => invocation.eventDependencyIds.includes( candidate.automationEventId, ), ) .map((candidate) => candidate.outcomeSeq), ...state.context.signalOutputs .filter((candidate) => invocation.signalDependencyIds.includes( candidate.signalInvocationId, ), ) .map((candidate) => candidate.outcomeSeq), ), status: "succeeded", } } return combineSignalDependencyOutcomes( contextDependencies.map((dependency) => dependency[signalContextDependencyCollector](state, dependencies), ), ) }, contextOrigins: [ ...new Set( contextDependencies.flatMap((signal) => signal[signalContextOrigins]), ), ], origins: [`action:${formatted}`], resolve: (context) => { const match = findUniqueBoundary( context.actionOutputs, (output) => isSameHookLocation(output, hook), `action hook ${formatted}`, context, ) if (!match) { throw new Error(`Missing action dependency for hook ${formatted}.`) } if (match.status === "failed") { return { failure: match.failure, status: "failed" } } if (match.status === "skipped") { return match.reason === "closed-dependency" ? { status: "closed" } : { failure: match.failure, status: "failed" } } return { ...(match.outputSensitivity && { sensitivity: match.outputSensitivity, }), status: "succeeded", // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The action's output schema establishes T. value: match.output as T, } }, }) } /** * Creates a signal for a subscription event at a deterministic hook slot. * * @param hook - Hook location assigned to the subscription. * @param staticProperties - Values exposed directly instead of as projections. */ export function createSubscriptionSignal< T, const TStatic extends object = object, >( hook: HookLocation, staticProperties?: TStatic, ): Signal & Readonly { const formatted = formatHookLocation(hook) return createSignal( { collect: (state, dependencies) => { const match = findUniqueBoundary( state.context.events, (event) => isSameHookLocation(event, hook), `subscription hook ${formatted}`, state.context, ) if (!match) return { status: "pending" } dependencies.eventDependencyIds.add(match.automationEventId) return { outcomeSeq: match.outcomeSeq, status: "succeeded" } }, origins: [`event:${formatted}`], resolve: (context) => { const match = findUniqueBoundary( context.events, (event) => isSameHookLocation(event, hook), `subscription hook ${formatted}`, context, ) if (!match) { throw new Error( `Missing subscription dependency for hook ${formatted}.`, ) } return { status: "succeeded", // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The subscription's event definition establishes T. value: match.payload as T, } }, }, staticProperties, ) } /** * Resolves the durable delivery ID for one subscription occurrence. * * @param anchor - Subscription signal whose delivery owns the ID. */ export function subscriptionEventId(anchor: Signal): Signal { const origins = anchor[signalOrigins] let selectedDependencies: SignalDependencySets | undefined return createSignal({ ancestorOrigins: anchor[signalAncestorOrigins], collect: (state, dependencies) => { const collected = collectAnchorDependencies(anchor, state, dependencies) selectedDependencies = collected.dependencies return collected.outcome }, collectContext: anchor[signalContextDependencyCollector], contextOrigins: anchor[signalContextOrigins], derivation: { inputs: [anchor], type: "transform" }, origins, resolve: (context) => { const outcome = anchor[signalResolver](context) if (outcome.status !== "succeeded") return outcome const boundaries = resolveAnchorBoundaries( context, origins, selectedDependencies, ) if (boundaries.length !== 1 || boundaries[0]?.kind !== "event") { throw new Error( "Cannot resolve a subscription delivery ID without exactly one event boundary.", ) } return { status: "succeeded", value: boundaries[0].id } }, }) } /** * Creates a stable opaque identity anchored to a signal occurrence. * * By default each declaration receives a distinct identity. Pass `false` to * derive identity only from the anchor's durable dependency boundary, so * equivalent transforms of the same occurrence resolve to the same value. The * returned signal is keyed by the identity itself and can be passed directly to * cross-context coordination operators. * * @param anchor - Signal occurrence that owns the correlation identity. * @param perDeclaration - Whether the declaration contributes to identity. */ export function correlationId( anchor: Signal, perDeclaration = true, ): KeyedSignal { const declaration = String(consumeDeclarationSlot()) const origins = anchor[signalOrigins] let selectedDependencies: SignalDependencySets | undefined return createSignal( { ancestorOrigins: anchor[signalAncestorOrigins], collect: (state, dependencies) => { const collected = collectAnchorDependencies(anchor, state, dependencies) selectedDependencies = collected.dependencies const { outcome } = collected return outcome }, collectContext: anchor[signalContextDependencyCollector], contextOrigins: anchor[signalContextOrigins], derivation: { inputs: [anchor], type: "transform" }, origins, resolve: (context) => { const outcome = anchor[signalResolver](context) return outcome.status === "succeeded" ? { status: "succeeded", value: hashCorrelationId( context, origins, selectedDependencies, perDeclaration ? declaration : undefined, ), } : outcome }, }, { [signalIndex]: CORRELATION_ID_INDEX }, ) } /** * Resolves the durable instant at which one signal occurrence emitted. * * Pure derivations use their latest contributing boundary. Actions use their * completion time, triggers use platform ingestion time, and coordinated * signals use their persisted decision boundary. * * @param anchor - Signal occurrence whose emission time is requested. */ export function timestamp(anchor: Signal): Signal { const origins = anchor[signalOrigins] let selectedDependencies: SignalDependencySets | undefined return createSignal({ ancestorOrigins: anchor[signalAncestorOrigins], collect: (state, dependencies) => { const collected = collectAnchorDependencies(anchor, state, dependencies) selectedDependencies = collected.dependencies return collected.outcome }, collectContext: anchor[signalContextDependencyCollector], contextOrigins: anchor[signalContextOrigins], derivation: { inputs: [anchor], type: "transform" }, origins, resolve: (context) => { const outcome = anchor[signalResolver](context) if (outcome.status !== "succeeded") return outcome const latest = resolveAnchorBoundaries( context, origins, selectedDependencies, ).toSorted( (left, right) => right.timestamp.getTime() - left.timestamp.getTime(), )[0] if (!latest) { throw new Error( "Cannot resolve a timestamp without a durable boundary.", ) } return { status: "succeeded", value: latest.timestamp } }, }) } /** * Collects exact durable dependencies for structural boundary origins. * * @param origins - Structural origins to resolve in the current frontier. * @param state - Current symbolic runtime traversal. * @param dependencies - Exact durable dependencies receiving the result. * @throws When an origin is malformed or resolves ambiguously. */ export function collectBoundaryOrigins( origins: readonly string[], state: AutomationRunState, dependencies: SignalDependencySets, ) { let outcomeSeq = 0 for (const origin of origins) { const [kind, formatted] = origin.split(":") const segments = formatted?.split(".").map(Number) ?? [] const slot = segments.pop() if (slot === undefined || segments.some(Number.isNaN)) { throw new AutomationCompatibilityError( `Automation compatibility error: invalid signal origin ${origin}.`, ) } const hook = { scopePath: segments, slot } if (kind === "action") { const boundary = findUniqueBoundary( state.context.actions, (output) => isSameHookLocation(output, hook), `action hook ${formatted}`, state.context, ) if (!boundary || boundary.status === "pending") { return { status: "pending" } as const } dependencies.actionDependencyIds.add(boundary.actionInvocationId) outcomeSeq = Math.max(outcomeSeq, boundary.outcomeSeq) continue } if (kind === "event") { const boundary = findUniqueBoundary( state.context.events, (event) => isSameHookLocation(event, hook), `subscription hook ${formatted}`, state.context, ) if (!boundary) return { status: "pending" } as const dependencies.eventDependencyIds.add(boundary.automationEventId) outcomeSeq = Math.max(outcomeSeq, boundary.outcomeSeq) continue } if (kind === "signal") { const boundary = findUniqueSignalBoundary( state.context.signalOutputs, (output) => isSameHookLocation(output, hook), `signal hook ${formatted}`, state.context, ) if (!boundary) return { status: "pending" } as const dependencies.signalDependencyIds.add(boundary.signalInvocationId) outcomeSeq = Math.max(outcomeSeq, boundary.outcomeSeq) continue } throw new AutomationCompatibilityError( `Automation compatibility error: invalid signal origin ${origin}.`, ) } return { outcomeSeq, status: "succeeded" } as const } /** * Restricts a hydrated context to one selected source occurrence. * * @param context - Fully hydrated coordinated child context. * @param item - Exact dependency membership of one selected occurrence. */ export function selectDependencyContext( context: RuntimeValueContext, item: RuntimeValueContext["coordinationSelections"][number]["items"][number], ): RuntimeValueContext { const actionDependencyIds = new Set(item.actionDependencyIds) const eventDependencyIds = new Set(item.eventDependencyIds) const signalDependencyIds = new Set(item.signalDependencyIds) const pendingSignalDependencyIds = [...signalDependencyIds] for (const signalDependencyId of pendingSignalDependencyIds) { const selection = context.coordinationSelections.find( (candidate) => candidate.signalInvocationId === signalDependencyId, ) if (!selection) continue for (const nestedItem of selection.items) { for (const id of nestedItem.actionDependencyIds) { actionDependencyIds.add(id) } for (const id of nestedItem.eventDependencyIds) { eventDependencyIds.add(id) } for (const id of nestedItem.signalDependencyIds) { if (signalDependencyIds.has(id)) continue signalDependencyIds.add(id) pendingSignalDependencyIds.push(id) } } } return { actionOutputs: context.actionOutputs.filter((output) => actionDependencyIds.has(output.actionInvocationId), ), coordinationSelections: context.coordinationSelections, contextId: item.contextId, contextParents: context.contextParents, events: context.events.filter((event) => eventDependencyIds.has(event.automationEventId), ), signalOutputs: context.signalOutputs.filter((output) => signalDependencyIds.has(output.signalInvocationId), ), } } /** * Restricts symbolic traversal to one persisted signal's exact dependencies. * * @param context - Merged durable frontier being traversed. * @param boundary - Persisted signal whose dependencies select one occurrence. */ export function selectRunDependencyContext( context: RuntimeRunContext, boundary: RuntimeRunContext["signalOutputs"][number], ): RuntimeRunContext { const actionDependencyIds = new Set(boundary.actionDependencyIds) const eventDependencyIds = new Set(boundary.eventDependencyIds) const signalDependencyIds = new Set(boundary.signalDependencyIds) return { ...context, actions: context.actions.filter((action) => actionDependencyIds.has(action.actionInvocationId), ), events: context.events.filter((event) => eventDependencyIds.has(event.automationEventId), ), signalOutputs: context.signalOutputs.filter((output) => signalDependencyIds.has(output.signalInvocationId), ), } } /** * Selects one durable signal occurrence while treating closure as absence. * * Merged contexts may contain closed sibling occurrences at the same hook. A * single non-closed occurrence still supplies one unambiguous signal value. * * @param boundaries - Merged durable signal occurrences to inspect. * @param predicate - Boundary identity filter. * @param label - Signal description included in ambiguity errors. * @param context - Current context and inherited parent relations. */ export function findUniqueSignalBoundary< T extends { contextId: string; status: string }, >( boundaries: T[], predicate: (boundary: T) => boolean, label: string, context: Pick, ) { const matches = boundaries.filter(predicate) const nonClosed = matches.filter((boundary) => boundary.status !== "closed") return findUniqueBoundary( nonClosed.length > 0 ? nonClosed : matches, () => true, label, context, ) } /** * Selects a boundary only when merged parent history contains one occurrence. * * @param boundaries - Merged durable occurrences to inspect. * @param predicate - Boundary identity filter. * @param label - Signal description included in ambiguity errors. * @param context - Current context and inherited parent relations. * @throws {AmbiguousSignalError} When several occurrences match. */ export function findUniqueBoundary( boundaries: T[], predicate: (boundary: T) => boolean, label: string, context: Pick, ) { const matches = boundaries.filter(predicate) if (matches.length < 2) return matches[0] const resolveContext = (contextId: string): T | undefined => { const direct = matches.filter((match) => match.contextId === contextId) if (direct.length > 1) { throw new AmbiguousSignalError( `Cannot resolve ${label}: merged parent contexts contain ${direct.length} distinct occurrences.`, ) } if (direct[0]) return direct[0] const parents = context.contextParents .filter((relation) => relation.contextId === contextId) .toSorted((left, right) => left.position - right.position) const primary = parents.find((relation) => relation.isPrimary) if (primary) { const primaryMatch = resolveContext(primary.parentContextId) if (primaryMatch) return primaryMatch } const inherited = [ ...new Set( parents.flatMap((relation) => { const match = resolveContext(relation.parentContextId) return match ? [match] : [] }), ), ] if (inherited.length < 2) return inherited[0] throw new AmbiguousSignalError( `Cannot resolve ${label}: merged parent contexts contain ${inherited.length} distinct occurrences.`, ) } return resolveContext(context.contextId) } /** * Hashes the exact durable boundaries hydrated for one anchor occurrence. * * @param context - Hydrated values containing the anchor's boundaries. * @param origins - Durable hooks contributing to the anchor. * @param selectedDependencies - Exact dependencies selected during collection. * @param declaration - Optional declaration hook distinguishing sibling IDs. */ function hashCorrelationId( context: RuntimeValueContext, origins: readonly string[], selectedDependencies: SignalDependencySets | undefined, declaration?: string, ) { const boundaries = resolveAnchorBoundaries( context, origins, selectedDependencies, ) .map((boundary) => `${boundary.kind}:${boundary.id}`) .toSorted() return createHash("sha256") .update( JSON.stringify({ boundaries: boundaries.length > 0 ? boundaries : [`context:${context.contextId}`], declaration: declaration ?? null, }), ) .digest() .subarray(0, 16) .toString("base64url") } /** * Collects an anchor into isolated sets before forwarding its dependencies. * * @param anchor - Signal whose exact dependencies are collected. * @param state - Active automation traversal state. * @param dependencies - Outer dependency sets receiving the collected IDs. */ function collectAnchorDependencies( anchor: Signal, state: AutomationRunState, dependencies: SignalDependencySets, ) { const selected = { actionDependencyIds: new Set(), eventDependencyIds: new Set(), signalDependencyIds: new Set(), } // Collection must run before the populated sets can be forwarded outward. const outcome = anchor[signalDependencyCollector](state, selected) for (const id of selected.actionDependencyIds) { dependencies.actionDependencyIds.add(id) } for (const id of selected.eventDependencyIds) { dependencies.eventDependencyIds.add(id) } for (const id of selected.signalDependencyIds) { dependencies.signalDependencyIds.add(id) } return { dependencies: selected, outcome } } type TimestampedBoundary = | { id: string kind: "action" timestamp: Date } | { id: string kind: "event" timestamp: Date } | { id: string kind: "signal" timestamp: Date } /** * Resolves exact durable boundaries selected by an anchor occurrence. * * @param context - Hydrated durable values visible to the occurrence. * @param origins - Fallback hook origins used before dependency collection. * @param selectedDependencies - Exact dependencies selected during collection. */ function resolveAnchorBoundaries( context: RuntimeValueContext, origins: readonly string[], selectedDependencies?: SignalDependencySets, ): TimestampedBoundary[] { if (selectedDependencies) { return [ ...[...selectedDependencies.actionDependencyIds].map((id) => { const boundary = context.actionOutputs.find( (output) => output.actionInvocationId === id, ) if (!boundary) throw new Error(`Missing action dependency ${id}.`) return { id, kind: "action" as const, timestamp: boundary.timestamp } }), ...[...selectedDependencies.eventDependencyIds].map((id) => { const boundary = context.events.find( (event) => event.automationEventId === id, ) if (!boundary) throw new Error(`Missing event dependency ${id}.`) return { id, kind: "event" as const, timestamp: boundary.timestamp } }), ...[...selectedDependencies.signalDependencyIds].map((id) => { const boundary = context.signalOutputs.find( (output) => output.signalInvocationId === id, ) if (!boundary) throw new Error(`Missing signal dependency ${id}.`) return { id, kind: "signal" as const, timestamp: boundary.timestamp } }), ] } return origins.reduce((boundaries, origin) => { const [kind, formatted] = origin.split(":") const segments = formatted?.split(".").map(Number) ?? [] const slot = segments.pop() if (slot === undefined) return boundaries const hook = { scopePath: segments, slot } if (kind === "action") { const boundary = findUniqueBoundary( context.actionOutputs, (output) => isSameHookLocation(output, hook), `action hook ${formatted}`, context, ) if (boundary) { boundaries.push({ id: boundary.actionInvocationId, kind, timestamp: boundary.timestamp, }) } return boundaries } if (kind === "event") { const boundary = findUniqueBoundary( context.events, (event) => isSameHookLocation(event, hook), `subscription hook ${formatted}`, context, ) if (boundary) { boundaries.push({ id: boundary.automationEventId, kind, timestamp: boundary.timestamp, }) } return boundaries } if (kind === "signal") { const boundary = findUniqueSignalBoundary( context.signalOutputs, (output) => isSameHookLocation(output, hook), `signal hook ${formatted}`, context, ) if (boundary) { boundaries.push({ id: boundary.signalInvocationId, kind, timestamp: boundary.timestamp, }) } return boundaries } return boundaries }, []) }