import { hashSignalCoordinatorConfiguration } from "@automate.ax/api-contract/runtime" import * as z from "zod/mini" import type { AutomationRunState, RuntimeValueContext } from "./runtime" import { AutomationCompatibilityError, formatHookLocation, isSameHookLocation, type HookLocation, } from "./runtime" import { buildSignalDerivation, isSignal, FailedSignalError, signalContextDependencyCollector, signalDependencyCollector, signalGlobal, signalIndex, signalOrigins, signalResolver, type KeyedSignal, type Signal, type SignalDependencyResolution, type SignalDependencySets, } from "./signal-protocol" import { collectBoundaryOrigins, findUniqueSignalBoundary, materializeSignal, selectRunDependencyContext, } from "./signal-resolution" const COLLECTION_COUNT_SCHEMA = z.int().check(z.positive()) const FUNNEL_DURATION_SCHEMA = z.union([ z.number().check(z.nonnegative()), z.string().check(z.minLength(1)), ]) interface PersistedSignalDependencies { actionDependencyIds: string[] eventDependencyIds: string[] signalDependencyIds: string[] } interface SourceCoordinationDefinition extends HookLocation { prerequisites?: Signal signal: Signal } interface CoordinatedSignalDefinition extends SourceCoordinationDefinition { primaryPosition: "first" | "last" } interface CollectionDefinition extends CoordinatedSignalDefinition { count: number | Signal occurrenceTtl?: number | string } interface FunnelDefinition extends CoordinatedSignalDefinition { maxBurstDuration?: number | string | Signal minGap?: number | string | Signal minQuietPeriod?: number | string | Signal selection: "all" | "first" | "last" triggerAt: "both" | "end" | "start" until?: Date | Signal } interface CorrelationDefinition extends HookLocation { inputPolicies?: readonly { consumption: "consume" | "retain" selection: "latest" | "oldest" }[] occurrenceTtl?: number | string ordered: boolean prerequisites?: Signal streams: KeyedSignal[] } interface GateDefinition extends HookLocation { condition: Signal prerequisites?: Signal value: Signal } interface RaceDefinition extends HookLocation { cohort?: Signal cohortOrigins: readonly string[] inputOrigins: readonly (readonly string[])[] inputs: ReadonlyArray prerequisiteOrigins: readonly string[] prerequisites?: Signal } interface MergeDefinition extends HookLocation { inputOrigins: readonly (readonly string[])[] inputs: ReadonlyArray prerequisites?: Signal } interface SerializationDefinition extends SourceCoordinationDefinition { releaseScopePath: readonly number[] releaseSlot: number } interface ConcurrencyDefinition extends SerializationDefinition { limit: number } interface TakeDefinition extends CoordinatedSignalDefinition { count: number ttl?: number | string | Signal } interface RateLimitDefinition extends CoordinatedSignalDefinition { interval: number | string | Signal limit: number overflow: "drop" | "wait" } /** * Collects one completed correlation, collection, or funnel boundary. * * @param location - Durable signal hook being consumed. * @param state - Current runtime traversal state. * @param dependencies - Consumer dependencies receiving the marker. * @throws {AutomationCompatibilityError} When the persisted outcome changed. */ export function collectCompletedSignalDependencies( location: HookLocation, state: AutomationRunState, dependencies: SignalDependencySets, ): SignalDependencyResolution { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, location), `signal hook ${formatHookLocation(location)}`, state.context, ) if (!completed) return { status: "pending" } dependencies.signalDependencyIds.add(completed.signalInvocationId) if (completed.status === "closed") { return { outcomeSeq: completed.outcomeSeq, status: "closed" } } if (completed.status === "failed") { return { failure: completed.failure, outcomeSeq: completed.outcomeSeq, status: "failed", } } if (completed.status === "succeeded") { return { outcomeSeq: completed.outcomeSeq, status: "succeeded" } } throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: terminal outcome changed.`, ) } /** * Collects one fan-out child marker as an available item dependency. * * @param location - Fan-out hook being consumed. * @param state - Current runtime traversal state. * @param dependencies - Consumer dependencies receiving the marker. * @throws {AutomationCompatibilityError} When the marker shape changed. */ export function collectFanoutDependencies( location: HookLocation, state: AutomationRunState, dependencies: SignalDependencySets, ): SignalDependencyResolution { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, location), `signal hook ${formatHookLocation(location)}`, state.context, ) if (!completed) return { status: "pending" } dependencies.signalDependencyIds.add(completed.signalInvocationId) if (completed.status === "closed") { return { outcomeSeq: completed.outcomeSeq, status: "closed" } } if (completed.status === "failed") { return { failure: completed.failure, outcomeSeq: completed.outcomeSeq, status: "failed", } } if ( completed.status === "succeeded" && completed.selectedIndex !== undefined ) { return { outcomeSeq: completed.outcomeSeq, status: "succeeded" } } throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: fan-out outcome changed.`, ) } /** * Registers one count-bounded occurrence offered by this context. * * @param definition - Collection definition and source signal. * @param state - Current runtime traversal state. */ export function registerCollection( definition: CollectionDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration( completed, state, definition, definition.signal[signalIndex]?.getPosition ? { ...(definition.occurrenceTtl !== undefined && { occurrenceTtl: definition.occurrenceTtl, }), ordered: true, policy: "collect", } : { ...(definition.occurrenceTtl !== undefined && { occurrenceTtl: definition.occurrenceTtl, }), policy: "collect", }, definition.primaryPosition, ) return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "collect" && isSameHookLocation(evaluation, definition), ) ) { return } const evaluation = planCoordinationOffer( definition, state, isSignal(definition.count) ? [definition.count] : [], ) if (!evaluation) return state.coordinationEvaluations.push({ ...evaluation, evaluate: (context) => { const { key, position } = evaluation.evaluate(context) return { count: COLLECTION_COUNT_SCHEMA.parse( isSignal(definition.count) ? materializeSignal(definition.count, context) : definition.count, ), key, ...(position !== undefined && { position }), } }, policy: "collect", ...(definition.occurrenceTtl !== undefined && { occurrenceTtl: definition.occurrenceTtl, ttl: definition.occurrenceTtl, }), }) } /** * Registers one bounded keyed admission occurrence. * * @param definition - Admission count, optional claim TTL, and source. * @param state - Current runtime traversal state. */ export function registerTake( definition: TakeDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { policy: "take", }) return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "take" && isSameHookLocation(evaluation, definition), ) ) { return } const evaluation = planCoordinationOffer( definition, state, isSignal(definition.ttl) ? [definition.ttl] : [], ) if (!evaluation) return state.coordinationEvaluations.push({ ...evaluation, count: definition.count, evaluate: (context) => { const { key } = evaluation.evaluate(context) return { count: definition.count, key, ...(definition.ttl !== undefined && { ttl: FUNNEL_DURATION_SCHEMA.parse( isSignal(definition.ttl) ? materializeSignal(definition.ttl, context) : definition.ttl, ), }), } }, policy: "take", }) } /** * Registers one rolling-window rate-limit occurrence. * * @param definition - Rolling capacity, interval, overflow, and source. * @param state - Current runtime traversal state. */ export function registerRateLimit( definition: RateLimitDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { overflow: definition.overflow, policy: "rateLimit", }) return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "rateLimit" && isSameHookLocation(evaluation, definition), ) ) { return } const evaluation = planCoordinationOffer( definition, state, isSignal(definition.interval) ? [definition.interval] : [], ) if (!evaluation) return state.coordinationEvaluations.push({ ...evaluation, evaluate: (context) => { const { key } = evaluation.evaluate(context) return { interval: FUNNEL_DURATION_SCHEMA.parse( isSignal(definition.interval) ? materializeSignal(definition.interval, context) : definition.interval, ), key, limit: definition.limit, } }, limit: definition.limit, overflow: definition.overflow, policy: "rateLimit", }) } /** * Registers one eager array fan-out declaration. * * @param definition - Fan-out source and deterministic hook slot. * @param state - Current runtime traversal state. * @throws {AutomationCompatibilityError} When replayed fan-out state changed. */ export function registerFanout( definition: SourceCoordinationDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { policy: "fanout", }) if ( completed.status !== "succeeded" || completed.selectedIndex === undefined ) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: fan-out outcome changed.`, ) } return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "fanout" && isSameHookLocation(evaluation, definition), ) ) { return } const evaluation = planCoordinationSource(definition, state) if (!evaluation) return const { resolve, ...offer } = evaluation state.coordinationEvaluations.push({ ...offer, evaluate: (context) => { return { count: resolve(context).length } }, policy: "fanout", }) } /** * Registers one temporally funneled occurrence offered by this context. * * @param definition - Funnel policy and source signal. * @param state - Current runtime traversal state. */ export function registerFunnel( definition: FunnelDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration( completed, state, definition, { policy: "funnel", selection: definition.selection, triggerAt: definition.triggerAt, }, definition.primaryPosition, ) return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "funnel" && isSameHookLocation(evaluation, definition), ) ) { return } const evaluation = planCoordinationOffer( definition, state, [ definition.maxBurstDuration, definition.minGap, definition.minQuietPeriod, definition.until, ].filter((value): value is Signal => isSignal(value)), ) if (!evaluation) return state.coordinationEvaluations.push({ ...evaluation, evaluate: (context) => { const { key } = evaluation.evaluate(context) const resolveDuration = ( duration: number | string | Signal | undefined, ) => duration === undefined ? undefined : FUNNEL_DURATION_SCHEMA.parse( isSignal(duration) ? materializeSignal(duration, context) : duration, ) return { key, ...(definition.maxBurstDuration !== undefined && { maxBurstDuration: resolveDuration(definition.maxBurstDuration), }), ...(definition.minGap !== undefined && { minGap: resolveDuration(definition.minGap), }), ...(definition.minQuietPeriod !== undefined && { minQuietPeriod: resolveDuration(definition.minQuietPeriod), }), ...(definition.until !== undefined && { until: isSignal(definition.until) ? materializeSignal(definition.until, context) : definition.until, }), } }, policy: "funnel", selection: definition.selection, triggerAt: definition.triggerAt, }) } /** * Plans value hydration for one generic coordination offer. * * @param definition - Coordinated source definition. * @param state - Current runtime traversal state. * @param parameters - Additional signals needed to configure the offer. */ function planCoordinationOffer( definition: CoordinatedSignalDefinition, state: AutomationRunState, parameters: Signal[] = [], ) { const source = planCoordinationSource(definition, state, parameters) if (!source) return const { resolve, ...offer } = source const index = definition.signal[signalIndex] return { ...offer, evaluate: (context: RuntimeValueContext) => { if (index) assertPublicCoordinatorSource(definition.signal, context) const item = resolve(context) return { key: index ? index.getKey(item, context) : null, ...(index?.getPosition && { position: index.getPosition(item, context), }), } }, primaryPosition: definition.primaryPosition, } } /** * Plans common prerequisites, source, and parameters for one coordinator. * * @param definition - Source signal and durable hook identity. * @param state - Current runtime traversal state. * @param parameters - Additional signals needed to configure the offer. */ function planCoordinationSource( definition: SourceCoordinationDefinition, state: AutomationRunState, parameters: Signal[] = [], ) { const itemDependencies = createSignalDependencySets() const prerequisiteOutcome = collectSignalPrerequisites( definition.prerequisites, state, itemDependencies, ) if (prerequisiteOutcome.status === "failed") { registerCoordinationFailure( definition, prerequisiteOutcome, itemDependencies, definition.prerequisites ? [definition.prerequisites] : [], state, ) return } if (prerequisiteOutcome.status !== "succeeded") return const itemOutcome = definition.signal[signalDependencyCollector]( state, itemDependencies, ) const parameterOutcomes = parameters.map((parameter) => parameter[signalDependencyCollector](state, itemDependencies), ) if ( itemOutcome.status === "pending" || parameterOutcomes.some((outcome) => outcome.status === "pending") ) { return } const settledParameterOutcomes = parameterOutcomes.filter( ( outcome, ): outcome is Exclude => outcome.status !== "pending", ) const outcomeSeq = Math.max( prerequisiteOutcome.outcomeSeq, itemOutcome.outcomeSeq, ...settledParameterOutcomes.map((outcome) => outcome.outcomeSeq), ) const failure = settledParameterOutcomes.find( ( outcome, ): outcome is Extract => outcome.status === "failed", ) ?? (itemOutcome.status === "failed" ? itemOutcome : undefined) if (failure) { registerCoordinationFailure( definition, { ...failure, outcomeSeq }, itemDependencies, [ definition.signal, ...parameters, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, ) return } if ( itemOutcome.status === "closed" || settledParameterOutcomes.some((outcome) => outcome.status === "closed") ) { return } return { actionDependencyIds: [...itemDependencies.actionDependencyIds], derivation: buildSignalDerivation( [ definition.signal, ...parameters, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, itemDependencies, ), eventDependencyIds: [...itemDependencies.eventDependencyIds], outcomeSeq, resolve: (context: RuntimeValueContext) => { const item = definition.signal[signalResolver](context) if (item.status === "failed") throw new FailedSignalError(item.failure) if (item.status === "closed") { throw new AutomationCompatibilityError( `Automation compatibility error at coordinated signal hook ${formatHookLocation(definition)}: source changed.`, ) } return item.value }, scopePath: [...definition.scopePath], signalDependencyIds: [...itemDependencies.signalDependencyIds], slot: definition.slot, } } /** * Persists a failed coordinator declaration without treating closure as final. * * @param location - Durable coordinator hook identity. * @param outcome - Failure produced by this context. * @param dependencies - Exact dependencies responsible for the failure. * @param signals - Signals whose evaluation produced the failure. * @param state - Current runtime traversal state. */ function registerCoordinationFailure( location: HookLocation, outcome: Extract, dependencies: SignalDependencySets, signals: readonly Signal[], state: AutomationRunState, ) { if ( state.signalInvocations.some((invocation) => isSameHookLocation(invocation, location), ) ) { return } state.signalInvocations.push({ actionDependencyIds: [...dependencies.actionDependencyIds], derivation: buildSignalDerivation(signals, state, dependencies), eventDependencyIds: [...dependencies.eventDependencyIds], failure: outcome.failure, scopePath: [...location.scopePath], signalDependencyIds: [...dependencies.signalDependencyIds], slot: location.slot, status: "failed", }) } /** * Registers one eager durable correlation declaration for this traversal. * * @param definition - Indexed streams and deterministic correlation slot. * @param state - Current runtime traversal state. */ export function registerCorrelation( definition: CorrelationDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { ...(definition.inputPolicies && { inputPolicies: definition.inputPolicies.map((policy) => ({ ...policy, })), }), ...(definition.occurrenceTtl !== undefined && { occurrenceTtl: definition.occurrenceTtl, }), ordered: definition.ordered, policy: "correlate", streams: definition.streams.map((stream) => stream[signalOrigins]), }) return } const prerequisiteDependencies = createSignalDependencySets() const prerequisiteOutcome = collectSignalPrerequisites( definition.prerequisites, state, prerequisiteDependencies, ) if (prerequisiteOutcome.status === "failed") { registerCoordinationFailure( definition, prerequisiteOutcome, prerequisiteDependencies, definition.prerequisites ? [definition.prerequisites] : [], state, ) return } if (prerequisiteOutcome.status !== "succeeded") return const streamOrigins = definition.streams.map((stream) => [ ...stream[signalOrigins], ]) const streams = definition.streams.map((stream, streamIndex) => { const streamDependencies = createSignalDependencySets() addDependencySets(streamDependencies, prerequisiteDependencies) return { dependencies: streamDependencies, outcome: stream[signalDependencyCollector](state, streamDependencies), stream, streamIndex, } }) for (const entry of streams) { if ( entry.outcome.status !== "succeeded" || state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "correlate" && isSameHookLocation(evaluation, definition) && evaluation.streamIndex === entry.streamIndex, ) ) { continue } state.coordinationEvaluations.push({ actionDependencyIds: [...entry.dependencies.actionDependencyIds], derivation: buildSignalDerivation( [ entry.stream, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, entry.dependencies, ), evaluate: (context) => { const resolution = entry.stream[signalResolver](context) if (resolution.status === "failed") { throw new FailedSignalError(resolution.failure) } if (resolution.status === "closed") { throw new AutomationCompatibilityError( `Automation compatibility error at correlation hook ${formatHookLocation(definition)}: stream ${entry.streamIndex} changed.`, ) } if (resolution.sensitivity) { throw new AutomationCompatibilityError( "Sensitive signals cannot be used as coordinator partition keys.", ) } return { key: entry.stream[signalIndex].getKey(resolution.value, context), } }, eventDependencyIds: [...entry.dependencies.eventDependencyIds], outcomeSeq: Math.max( prerequisiteOutcome.outcomeSeq, entry.outcome.outcomeSeq, ), policy: "correlate", ...(definition.inputPolicies && { inputPolicies: definition.inputPolicies.map((policy) => ({ ...policy, })), }), signalDependencyIds: [...entry.dependencies.signalDependencyIds], ordered: definition.ordered, scopePath: [...definition.scopePath], slot: definition.slot, streamIndex: entry.streamIndex, streamOrigins, ...(definition.occurrenceTtl !== undefined && { occurrenceTtl: definition.occurrenceTtl, ttl: definition.occurrenceTtl, }), }) } const failure = streams.find( ( entry, ): entry is typeof entry & { outcome: Extract } => entry.outcome.status === "failed", ) if ( failure && !state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "correlate" && isSameHookLocation(evaluation, definition), ) ) { registerCoordinationFailure( definition, failure.outcome, failure.dependencies, [ failure.stream, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, ) } } /** * Registers every successful input occurrence for independent stream union. * * @param definition - Merge inputs and deterministic hook slot. * @param state - Current run state. */ export function registerMerge( definition: MergeDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { inputs: definition.inputOrigins, policy: "merge", }) return } const prerequisiteDependencies = createSignalDependencySets() const prerequisiteOutcome = collectSignalPrerequisites( definition.prerequisites, state, prerequisiteDependencies, ) if (prerequisiteOutcome.status === "failed") { registerCoordinationFailure( definition, prerequisiteOutcome, prerequisiteDependencies, definition.prerequisites ? [definition.prerequisites] : [], state, ) return } if (prerequisiteOutcome.status !== "succeeded") return const inputs = definition.inputs.map((input, selectedIndex) => { const dependencies = createSignalDependencySets() addDependencySets(dependencies, prerequisiteDependencies) return { dependencies, input, outcome: input[signalDependencyCollector](state, dependencies), selectedIndex, } }) for (const entry of inputs) { if ( entry.outcome.status !== "succeeded" || state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "merge" && isSameHookLocation(evaluation, definition) && evaluation.selectedIndex === entry.selectedIndex, ) ) { continue } state.coordinationEvaluations.push({ actionDependencyIds: [...entry.dependencies.actionDependencyIds], derivation: buildSignalDerivation( [ entry.input, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, entry.dependencies, ), evaluate: (context) => { materializeSignal(entry.input, context) }, eventDependencyIds: [...entry.dependencies.eventDependencyIds], inputOrigins: definition.inputOrigins.map((origins) => [...origins]), outcomeSeq: Math.max( prerequisiteOutcome.outcomeSeq, entry.outcome.outcomeSeq, ), policy: "merge", scopePath: [...definition.scopePath], selectedIndex: entry.selectedIndex, signalDependencyIds: [...entry.dependencies.signalDependencyIds], slot: definition.slot, }) } const failure = inputs.find( ( entry, ): entry is typeof entry & { outcome: Extract } => entry.outcome.status === "failed", ) if ( failure && !state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "merge" && isSameHookLocation(evaluation, definition), ) ) { registerCoordinationFailure( definition, failure.outcome, failure.dependencies, [ failure.input, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, ) } } /** * Registers one partitioned occurrence waiting for serialized admission. * * @param definition - Partitioned source and durable release hook. * @param state - Current run state. * @throws {AutomationCompatibilityError} When the source has no partition. */ export function registerSerialization( definition: SerializationDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { policy: "serialize", release: { scopePath: definition.releaseScopePath, slot: definition.releaseSlot, }, source: definition.signal[signalOrigins], }) return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "serialize" && isSameHookLocation(evaluation, definition), ) ) { return } const source = planCoordinationSource(definition, state) if (!source) return const index = definition.signal[signalIndex] if (!index && !definition.signal[signalGlobal]) { throw new AutomationCompatibilityError( "serialize requires a keyed or global signal.", ) } const { resolve, ...offer } = source state.coordinationEvaluations.push({ ...offer, evaluate: (context) => { const item = resolve(context) if (index) assertPublicCoordinatorSource(definition.signal, context) return { key: index ? index.getKey(item, context) : null } }, policy: "serialize", releaseScopePath: [...definition.releaseScopePath], releaseSlot: definition.releaseSlot, sourceOrigins: [...definition.signal[signalOrigins]], }) } /** * Registers one occurrence waiting for bounded concurrent admission. * * @param definition - Capacity, source, and durable release hook. * @param state - Current runtime traversal state. * @throws {AutomationCompatibilityError} When the source is not partitioned. */ export function registerConcurrency( definition: ConcurrencyDefinition, state: AutomationRunState, ) { const completed = findUniqueSignalBoundary( state.context.signalOutputs, (candidate) => isSameHookLocation(candidate, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { validateCoordinatorConfiguration(completed, state, definition, { limit: definition.limit, policy: "concurrent", release: { scopePath: definition.releaseScopePath, slot: definition.releaseSlot, }, source: definition.signal[signalOrigins], }) return } if ( state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "concurrent" && isSameHookLocation(evaluation, definition), ) ) { return } const source = planCoordinationSource(definition, state) if (!source) return const index = definition.signal[signalIndex] if (!index && !definition.signal[signalGlobal]) { throw new AutomationCompatibilityError( "concurrent requires a keyed or global signal.", ) } const { resolve, ...offer } = source state.coordinationEvaluations.push({ ...offer, evaluate: (context) => { const item = resolve(context) if (index) assertPublicCoordinatorSource(definition.signal, context) return { key: index ? index.getKey(item, context) : null } }, limit: definition.limit, policy: "concurrent", releaseScopePath: [...definition.releaseScopePath], releaseSlot: definition.releaseSlot, sourceOrigins: [...definition.signal[signalOrigins]], }) } /** * Rejects values that would expose sensitive material as a durable key. * * @param signal - Coordinator source being evaluated. * @param context - Exact dependency values for the occurrence. * @throws {AutomationCompatibilityError} When the source is sensitive. */ function assertPublicCoordinatorSource( signal: Signal, context: RuntimeValueContext, ): void { const resolution = signal[signalResolver](context) if (resolution.status === "succeeded" && resolution.sensitivity) { throw new AutomationCompatibilityError( "Sensitive signals cannot be used as coordinator partition keys.", ) } } /** * Plans one causal selection offer and its selected-value evaluation. * * @param definition - Race inputs and deterministic hook slot. * @param state - Current run state. * @param dependencies - Consumer dependencies receiving the race boundary. * @throws {AutomationCompatibilityError} When replay changes the race. */ export function collectRaceDependencies( definition: RaceDefinition, state: AutomationRunState, dependencies: SignalDependencySets, ): SignalDependencyResolution { const inheritedDependencies = createSignalDependencySets() const inheritedOutcome = collectSignalPrerequisites( definition.prerequisites, state, inheritedDependencies, ) const cohortDependencies = createSignalDependencySets() const cohortOutcome = definition.cohort ? definition.cohort[signalContextDependencyCollector]( state, cohortDependencies, ) : collectBoundaryOrigins( definition.cohortOrigins, state, cohortDependencies, ) const completed = findUniqueSignalBoundary( state.context.signalOutputs, (output) => isSameHookLocation(output, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (completed) { dependencies.signalDependencyIds.add(completed.signalInvocationId) validateCoordinatorConfiguration(completed, state, definition, { cohort: definition.cohortOrigins, inputs: definition.inputOrigins, scope: definition.prerequisiteOrigins, policy: "race", }) if (inheritedOutcome.status === "pending") return inheritedOutcome if (cohortOutcome.status === "pending") return cohortOutcome if (completed.status === "closed") { return { outcomeSeq: completed.outcomeSeq, status: "closed" } } if (completed.status === "failed") { return { failure: completed.failure, outcomeSeq: completed.outcomeSeq, status: "failed", } } if ( inheritedOutcome.status !== "succeeded" || completed.status !== "succeeded" || completed.selectedIndex === undefined ) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: race outcome changed.`, ) } const selected = definition.inputs[completed.selectedIndex] if (!selected) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: race inputs changed.`, ) } const selectedDependencies = createSignalDependencySets() addDependencySets(selectedDependencies, inheritedDependencies) const selectedInputOutcome = selected[signalDependencyCollector]( state, selectedDependencies, ) if ( selectedInputOutcome.status === "pending" || selectedInputOutcome.status === "closed" ) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: race winner changed.`, ) } addDependencySets(dependencies, selectedDependencies) return selectedInputOutcome.status === "failed" ? { failure: selectedInputOutcome.failure, outcomeSeq: completed.outcomeSeq, status: "failed", } : { outcomeSeq: completed.outcomeSeq, status: "succeeded" } } if (inheritedOutcome.status === "pending") return inheritedOutcome if (cohortOutcome.status === "pending") return cohortOutcome if ( state.signalInvocations.some((invocation) => isSameHookLocation(invocation, definition), ) ) { return { status: "pending" } } if (inheritedOutcome.status !== "succeeded") { const invocation = { actionDependencyIds: [...inheritedDependencies.actionDependencyIds], derivation: buildSignalDerivation( definition.prerequisites ? [definition.prerequisites] : [], state, inheritedDependencies, ), eventDependencyIds: [...inheritedDependencies.eventDependencyIds], outcomeSeq: inheritedOutcome.outcomeSeq, scopePath: [...definition.scopePath], signalDependencyIds: [...inheritedDependencies.signalDependencyIds], slot: definition.slot, } state.signalInvocations.push( inheritedOutcome.status === "failed" ? { ...invocation, failure: inheritedOutcome.failure, status: "failed", } : { ...invocation, status: "closed" }, ) return { status: "pending" } } const inputs = definition.inputs.map((input, selectedIndex) => { const inputDependencies = createSignalDependencySets() addDependencySets(inputDependencies, inheritedDependencies) const inputOutcome = input[signalDependencyCollector]( state, inputDependencies, ) return { dependencies: inputDependencies, outcome: inputOutcome.status === "pending" ? inputOutcome : { ...inputOutcome, outcomeSeq: Math.max( inheritedOutcome.outcomeSeq, inputOutcome.outcomeSeq, ), }, selectedIndex, signal: input, } }) const winner = inputs .filter( ( input, ): input is typeof input & { outcome: Extract< SignalDependencyResolution, { status: "failed" | "succeeded" } > } => input.outcome.status === "succeeded" || input.outcome.status === "failed", ) .toSorted( (left, right) => left.outcome.outcomeSeq - right.outcome.outcomeSeq || left.selectedIndex - right.selectedIndex, )[0] if (winner) { if ( !state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "race" && isSameHookLocation(evaluation, definition) && evaluation.selectedIndex === winner.selectedIndex, ) ) { state.coordinationEvaluations.push({ actionDependencyIds: [...winner.dependencies.actionDependencyIds], cohortActionDependencyIds: [...cohortDependencies.actionDependencyIds], cohortEventDependencyIds: [...cohortDependencies.eventDependencyIds], cohortSignalDependencyIds: [...cohortDependencies.signalDependencyIds], derivation: buildSignalDerivation( [ winner.signal, ...(definition.prerequisites ? [definition.prerequisites] : []), ...(definition.cohort ? [definition.cohort] : []), ], state, winner.dependencies, cohortDependencies, ), eventDependencyIds: [...winner.dependencies.eventDependencyIds], ...(winner.outcome.status === "failed" ? { failure: winner.outcome.failure, status: "failed" as const, } : { evaluate: (context: RuntimeValueContext) => { materializeSignal(winner.signal, context) }, status: "succeeded" as const, }), inputOrigins: definition.inputOrigins.map((origins) => [...origins]), cohortOrigins: [...definition.cohortOrigins], outcomeSeq: winner.outcome.outcomeSeq, policy: "race", scopeOrigins: [...definition.prerequisiteOrigins], scopePath: [...definition.scopePath], selectedIndex: winner.selectedIndex, signalDependencyIds: [...winner.dependencies.signalDependencyIds], slot: definition.slot, }) } return { status: "pending" } } for (const input of inputs) { if ( input.outcome.status !== "closed" || state.coordinationEvaluations.some( (evaluation) => evaluation.policy === "race" && isSameHookLocation(evaluation, definition) && evaluation.selectedIndex === input.selectedIndex, ) ) { continue } state.coordinationEvaluations.push({ actionDependencyIds: [...input.dependencies.actionDependencyIds], cohortActionDependencyIds: [...cohortDependencies.actionDependencyIds], cohortEventDependencyIds: [...cohortDependencies.eventDependencyIds], cohortSignalDependencyIds: [...cohortDependencies.signalDependencyIds], derivation: buildSignalDerivation( [ input.signal, ...(definition.prerequisites ? [definition.prerequisites] : []), ...(definition.cohort ? [definition.cohort] : []), ], state, input.dependencies, cohortDependencies, ), eventDependencyIds: [...input.dependencies.eventDependencyIds], inputOrigins: definition.inputOrigins.map((origins) => [...origins]), cohortOrigins: [...definition.cohortOrigins], outcomeSeq: input.outcome.outcomeSeq, policy: "race", scopeOrigins: [...definition.prerequisiteOrigins], scopePath: [...definition.scopePath], selectedIndex: input.selectedIndex, signalDependencyIds: [...input.dependencies.signalDependencyIds], slot: definition.slot, status: "closed", }) } return { status: "pending" } } /** * Collects or plans the durable decision boundary for one gate. * * @param definition - Gate inputs and deterministic slot. * @param state - Current run state. * @param dependencies - Consumer dependencies receiving the gate boundary. * @throws When replayed gate structure is incompatible. */ export function collectGateDependencies( definition: GateDefinition, state: AutomationRunState, dependencies: SignalDependencySets, ): SignalDependencyResolution { const inheritedDependencies = createSignalDependencySets() const inheritedOutcome = collectSignalPrerequisites( definition.prerequisites, state, inheritedDependencies, ) const completed = findUniqueSignalBoundary( state.context.signalOutputs, (output) => isSameHookLocation(output, definition), `signal hook ${formatHookLocation(definition)}`, state.context, ) if (inheritedOutcome.status === "pending") return inheritedOutcome if (inheritedOutcome.status !== "succeeded") { if (completed) { if ( completed.status !== inheritedOutcome.status || completed.outcomeSeq !== inheritedOutcome.outcomeSeq || !equalDependencySets(completed, inheritedDependencies) ) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: gate dependencies changed.`, ) } dependencies.signalDependencyIds.add(completed.signalInvocationId) return inheritedOutcome } if ( !state.signalInvocations.some((invocation) => isSameHookLocation(invocation, definition), ) ) { const invocation = { actionDependencyIds: [...inheritedDependencies.actionDependencyIds], derivation: buildSignalDerivation( definition.prerequisites ? [definition.prerequisites] : [], state, inheritedDependencies, ), eventDependencyIds: [...inheritedDependencies.eventDependencyIds], outcomeSeq: inheritedOutcome.outcomeSeq, scopePath: [...definition.scopePath], signalDependencyIds: [...inheritedDependencies.signalDependencyIds], slot: definition.slot, } state.signalInvocations.push( inheritedOutcome.status === "failed" ? { ...invocation, failure: inheritedOutcome.failure, status: "failed", } : { ...invocation, status: "closed" }, ) } return { status: "pending" } } const conditionDependencies = createSignalDependencySets() addDependencySets(conditionDependencies, inheritedDependencies) const replayState = completed ? { ...state, context: selectRunDependencyContext(state.context, completed), } : state const conditionOutcome = definition.condition[signalDependencyCollector]( replayState, conditionDependencies, ) const condition = conditionOutcome.status === "pending" ? conditionOutcome : { ...conditionOutcome, outcomeSeq: Math.max( inheritedOutcome.outcomeSeq, conditionOutcome.outcomeSeq, ), } if (completed) { if (!equalDependencySets(completed, conditionDependencies)) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: gate inputs changed.`, ) } dependencies.signalDependencyIds.add(completed.signalInvocationId) if (completed.status === "closed") { return { outcomeSeq: completed.outcomeSeq, status: "closed" } } if (completed.status === "failed") { return { failure: completed.failure, outcomeSeq: completed.outcomeSeq, status: "failed", } } if (completed.status !== "open") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: gate outcome changed.`, ) } const selectedValue = definition.value[signalDependencyCollector]( replayState, dependencies, ) const value = selectedValue.status === "pending" ? definition.value[signalDependencyCollector](state, dependencies) : selectedValue return value.status === "pending" ? value : { ...value, outcomeSeq: Math.max(completed.outcomeSeq, value.outcomeSeq), } } if ( state.signalInvocations.some((invocation) => isSameHookLocation(invocation, definition), ) || state.signalEvaluations.some((evaluation) => isSameHookLocation(evaluation, definition), ) || condition.status === "pending" ) { return { status: "pending" } } const invocation = { actionDependencyIds: [...conditionDependencies.actionDependencyIds], derivation: buildSignalDerivation( [ definition.condition, ...(definition.prerequisites ? [definition.prerequisites] : []), ], state, conditionDependencies, ), eventDependencyIds: [...conditionDependencies.eventDependencyIds], scopePath: [...definition.scopePath], signalDependencyIds: [...conditionDependencies.signalDependencyIds], slot: definition.slot, } if (condition.status === "succeeded") { state.signalEvaluations.push({ ...invocation, evaluate: (context) => { const resolution = definition.condition[signalResolver](context) if (resolution.status === "closed") { return { status: "closed" } } if (resolution.status === "failed") return resolution return { status: resolution.value ? "open" : "closed" } }, }) } else { state.signalInvocations.push( condition.status === "failed" ? { ...invocation, failure: condition.failure, status: "failed" } : { ...invocation, status: "closed" }, ) } return { status: "pending" } } /** * Collects the composite dependency inherited from prerequisite sections. * * @param prerequisites - Composite inherited signal, when active. * @param state - Current runtime traversal state. * @param dependencies - Dependency sets receiving the prerequisite IDs. */ function collectSignalPrerequisites( prerequisites: Signal | undefined, state: AutomationRunState, dependencies: SignalDependencySets, ): SignalDependencyResolution { return prerequisites ? prerequisites[signalDependencyCollector](state, dependencies) : { outcomeSeq: 0, status: "succeeded" } } /** Creates empty ordered dependency sets for one signal traversal. */ function createSignalDependencySets(): SignalDependencySets { return { actionDependencyIds: new Set(), eventDependencyIds: new Set(), signalDependencyIds: new Set(), } } /** * Adds one ordered dependency collection to another. * * @param target - Dependency sets receiving IDs. * @param source - Dependency sets supplying IDs. */ function addDependencySets( target: SignalDependencySets, source: SignalDependencySets, ) { for (const id of source.actionDependencyIds) target.actionDependencyIds.add(id) for (const id of source.eventDependencyIds) target.eventDependencyIds.add(id) for (const id of source.signalDependencyIds) target.signalDependencyIds.add(id) } /** * Rejects replay when a coordinated decision belongs to another definition. * * @param completed - Persisted signal boundary being replayed. * @param completed.coordinatorConfigurationHash - Stored coordinator identity. * @param completed.status - Persisted lifecycle used to recognize a direct * terminal marker. * @param state - Runtime traversal collecting asynchronous checks. * @param location - Current durable hook location. * @param configuration - Current structural coordinator definition. * @param primaryPosition - Optional inherited-parent precedence. * @throws {AutomationCompatibilityError} When coordinator identity is absent or * changed. */ function validateCoordinatorConfiguration( completed: { coordinatorConfigurationHash?: string status: "closed" | "failed" | "open" | "succeeded" }, state: AutomationRunState, location: HookLocation, configuration: unknown, primaryPosition?: "first" | "last", ) { if (!completed.coordinatorConfigurationHash) { if (completed.status === "closed" || completed.status === "failed") return throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: coordinator identity is missing.`, ) } state.compatibilityChecks.push( hashSignalCoordinatorConfiguration(configuration, primaryPosition).then( (configurationHash) => { if (configurationHash !== completed.coordinatorConfigurationHash) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: coordinator definition changed.`, ) } }, ), ) } /** * Compares persisted dependency arrays with one current traversal. * * @param actual - Persisted ordered dependency IDs. * @param expected - IDs from the current signal traversal. */ function equalDependencySets( actual: PersistedSignalDependencies, expected: SignalDependencySets, ) { return ( new Set(actual.actionDependencyIds).symmetricDifference( expected.actionDependencyIds, ).size === 0 && new Set(actual.eventDependencyIds).symmetricDifference( expected.eventDependencyIds, ).size === 0 && new Set(actual.signalDependencyIds).symmetricDifference( expected.signalDependencyIds, ).size === 0 ) }