import { unique } from "remeda" import type { RequireAtLeastOne, Simplify } from "type-fest" import * as z from "zod/mini" import { planSignalContextBoundary } from "./context-boundary-planning" import type { AutomationRunState, RuntimeValueContext } from "./runtime" import { AutomationCompatibilityError, consumeHookLocation, formatHookLocation, getAutomationPlanningContext, getCurrentHookScopePath, getAutomationRuntimeState, isSameHookLocation, type HookLocation, withHookScope, } from "./runtime" import { collectCompletedSignalDependencies, collectFanoutDependencies, collectGateDependencies, collectRaceDependencies, registerCollection, registerConcurrency, registerCorrelation, registerFanout, registerFunnel, registerMerge, registerRateLimit, registerSerialization, registerTake, } from "./signal-planning" import { combineSignalDependencyOutcomes, createSignal, FailedSignalError, inheritSignalCoordination, isSignal, signalAncestorOrigins, signalContextDependencyCollector, signalContextOrigins, signalDependencyBinder, signalDependencyCollector, signalGlobal, signalIndex, signalOrigins, signalResolver, signalValueType, transform, type DurationString, type GlobalSignal, type InheritSignalCoordination, type KeyedSignal, type InferSignal, type Signal, type SignalDependencySets, type AutomationError, type SignalIndexDefinition, type SignalOutcome, type Synchronous, } from "./signal-protocol" import { findUniqueSignalBoundary, selectDependencyContext, } from "./signal-resolution" const PREREQUISITE_SCOPES: (Signal | undefined)[] = [] const CONTEXT_PREREQUISITE_SCOPES: (Signal | undefined)[] = [] const POSITIVE_INTEGER_SCHEMA = z.int().check(z.positive()) export interface CorrelationOptions { /** How long an unmatched occurrence remains eligible. */ occurrenceTtl?: number | DurationString /** Whether streams must arrive in declaration order. */ ordered?: boolean } export interface CorrelationInputOptions { /** Whether a selected occurrence remains available for later matches. */ consumption?: "consume" | "retain" /** Which eligible occurrence wins when several share one key. */ selection?: "latest" | "oldest" } export type CorrelationInput = | TSignal | ({ signal: TSignal } & CorrelationInputOptions) /** * Matches one value from every keyed stream. * * Declaration eagerly registers durable correlation work. Matching is * one-to-one and oldest-first, and a match creates a child context that can * consume the original streams through their merged parent histories. * * @param inputs - Keyed signals and per-input policies in this correlation. * @param options - Optional temporal ordering. * @param options.ordered - Whether streams must arrive in array order. * @throws {TypeError} When streams share a durable origin. */ export function correlate( inputs: readonly [CorrelationInput, CorrelationInput, ...CorrelationInput[]], options: CorrelationOptions = {}, ): Signal { const normalized = inputs.map((input) => isSignal(input) ? { consumption: "consume" as const, selection: "oldest" as const, signal: input, } : { consumption: input.consumption ?? ("consume" as const), selection: input.selection ?? ("oldest" as const), signal: input.signal, }, ) const streams = normalized.map(({ signal }) => signal) const origins = streams.flatMap((signal) => signal[signalOrigins]) if (new Set(origins).size !== origins.length) { throw new TypeError( "correlate cannot use the same signal more than once; collect distinct occurrences before correlating them.", ) } const inputPolicies = normalized.map(({ consumption, selection }) => ({ consumption, selection, })) const definition = { ...(inputPolicies.some( ({ consumption, selection }) => consumption !== "consume" || selection !== "oldest", ) && { inputPolicies }), occurrenceTtl: options.occurrenceTtl, ordered: options.ordered ?? false, prerequisites: getCurrentSignalPrerequisites(), ...consumeHookLocation(), streams, } const planningContext = getAutomationPlanningContext() const boundaryDependencies = [ ...streams, ...(definition.prerequisites ? [definition.prerequisites] : []), ] if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, true, ) } const state = getAutomationRuntimeState() if (state) registerCorrelation(definition, state) return createSignal({ ancestorOrigins: unique([ ...boundaryDependencies.flatMap( (signal) => signal[signalAncestorOrigins], ), `signal:${formatHookLocation(definition)}`, ]), collect: (state, dependencies) => collectCompletedSignalDependencies(definition, state, dependencies), origins: [`signal:${formatHookLocation(definition)}`], resolve: (context) => { const outcome = findSignalBoundary(definition, context) if (outcome.status === "closed") return { status: "closed" } if (outcome.status === "failed") { return { failure: outcome.failure, status: "failed" } } if (outcome.status !== "succeeded") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: correlation outcome changed.`, ) } return { status: "succeeded", value: null } }, }) } /** Controls which parent supplies a value when merged histories conflict. */ export interface ValueInheritanceOptions { /** Parent whose value is inherited at otherwise ambiguous hook locations. */ inheritConflictingValuesFrom?: "first" | "last" } /** Fixed-size collection admission and parent-value behavior. */ export interface CollectionOptions extends ValueInheritanceOptions { /** How long an unmatched occurrence remains eligible. */ occurrenceTtl?: number | DurationString } export interface AdmissionOptions { /** How long the admission count remains claimed after its first emission. */ ttl?: Duration } export interface RateLimitOptions { /** Rolling interval in which at most `limit` occurrences may emit. */ interval: Duration /** Maximum emissions within the rolling interval. */ limit: number /** Whether excess occurrences wait FIFO or close without emitting. */ overflow?: "drop" | "wait" } export interface ConcurrencyOptions { /** Maximum active structured regions in each partition. */ limit: number } /** Compact or millisecond duration, optionally derived from another signal. */ export type Duration = | number | DurationString | (Signal & { readonly [signalValueType]: () => number | DurationString }) export type { DurationString } from "./signal-protocol" /** Absolute instant at which a durable timing policy ends. */ export type Deadline = Date | Signal /** * Marks the containing causal run as significant when the signal emits. * * The returned signal is the original value-preserving reference. * * @param signal - Signal whose successful emission makes its run significant. */ export function markSignificant( signal: TSignal, ): TSignal { const planningContext = getAutomationPlanningContext() if (planningContext) planningContext.usesMarkSignificant = true const state = getAutomationRuntimeState() if (!state) return signal const dependencies: SignalDependencySets = { actionDependencyIds: new Set(), eventDependencyIds: new Set(), signalDependencyIds: new Set(), } if ( signal[signalDependencyCollector](state, dependencies).status === "succeeded" ) { state.significanceEvaluations.push({ actionDependencyIds: [...dependencies.actionDependencyIds], evaluate: (context) => signal[signalResolver](context).status === "succeeded", eventDependencyIds: [...dependencies.eventDependencyIds], signalDependencyIds: [...dependencies.signalDependencyIds], }) } return signal } type AggregationSignal = GlobalSignal | KeyedSignal type KeyedAggregationSignal = TSignal & { readonly [signalIndex]: SignalIndexDefinition> } type GlobalAggregationSignal = TSignal & { readonly [signalGlobal]: true } /** Array containing at least one value. */ export type NonEmptyArray = [T, ...T[]] export type FunnelValue = [ Extract, ] extends [never] ? T : [Exclude] extends [never] ? NonEmptyArray : T | NonEmptyArray type FunnelResult< T, TMode extends "global" | "keyed", TOptions extends FunnelOptions, > = TMode extends "keyed" ? KeyedSignal> : GlobalSignal> type FunnelMaximumOptions = | { maxBurstDuration?: Duration; until?: never } | { maxBurstDuration?: never; until?: Deadline } type FunnelTimingOptions = FunnelMaximumOptions & RequireAtLeastOne<{ maxBurstDuration: Duration minGap: Duration minQuietPeriod: Duration until: Deadline }> /** Timing controls that can produce a trailing funnel emission. */ type FunnelEndTimingOptions = FunnelMaximumOptions & RequireAtLeastOne<{ maxBurstDuration: Duration minQuietPeriod: Duration until: Deadline }> /** Mutually exclusive single-value and buffered output policies. */ type FunnelOutputOptions = | (ValueInheritanceOptions & { buffer: true; select?: never }) | { buffer?: false inheritConflictingValuesFrom?: never select?: "first" | "latest" } /** Declarative timing and selection policy for a durable signal funnel. */ export type FunnelOptions = | Simplify< FunnelEndTimingOptions & { minGap?: Duration } & FunnelOutputOptions & { triggerAt?: "end" } > | Simplify | Simplify< FunnelTimingOptions & { buffer?: never inheritConflictingValuesFrom?: never select?: never triggerAt: "start" } > /** Options for an absolute-deadline buffered window. */ export interface WindowUntilOptions extends ValueInheritanceOptions { /** Absolute instant at which the window emits. */ until: Deadline } interface FunnelImplementationOptions extends ValueInheritanceOptions { buffer?: boolean maxBurstDuration?: number | string | Signal minGap?: number | string | Signal minQuietPeriod?: number | string | Signal select?: "first" | "latest" triggerAt?: "both" | "end" | "start" until?: Deadline } /** * Collects a fixed or signal-derived number of occurrences into one array. * * Occurrences are consumed oldest-first. A keyed source creates an independent * collection for every key; an explicitly global source shares one collection. * The nonempty result preserves that partition and inherits every selected * parent. * * @param signal - One-shot signal whose occurrences are collected. * @param count - Positive literal or signal-derived collection size. * @param options - Parent value precedence for the resulting child context. * @throws When a literal count is not a positive integer. */ export function collect( signal: KeyedAggregationSignal, count: number | Signal, options?: CollectionOptions, ): KeyedSignal>> /** @inheritdoc */ export function collect( signal: GlobalAggregationSignal, count: number | Signal, options?: CollectionOptions, ): GlobalSignal>> export function collect( signal: AggregationSignal, count: number | Signal, options: CollectionOptions = {}, ): AggregationSignal> { if (typeof count === "number") POSITIVE_INTEGER_SCHEMA.parse(count) const definition = { count, occurrenceTtl: options.occurrenceTtl, primaryPosition: options.inheritConflictingValuesFrom ?? ("last" as const), prerequisites: getCurrentSignalPrerequisites(), signal, ...consumeHookLocation(), } const boundaryDependencies = [ signal, ...(isSignal(count) ? [count] : []), ...(definition.prerequisites ? [definition.prerequisites] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) registerCollection(definition, state) return createCoordinatedSignal(signal, definition, true, boundaryDependencies) } /** * Emits only the first `count` occurrences in each source partition. * * @param signal - Keyed source whose partitions own independent counts. * @param count - Positive admission limit. * @param options - Optional lifetime for the complete count claim. */ export function take( signal: KeyedAggregationSignal, count: number, options?: AdmissionOptions, ): KeyedSignal> /** @inheritdoc */ export function take( signal: GlobalAggregationSignal, count: number, options?: AdmissionOptions, ): GlobalSignal> export function take( signal: AggregationSignal, count: number, options: AdmissionOptions = {}, ): AggregationSignal { return createTake(signal, count, options) } /** * Builds the shared `once` and `take` boundary. * * @param signal - Partitioned source occurrence. * @param count - Maximum admissions in one claim lifetime. * @param options - Optional claim lifetime. */ function createTake( signal: AggregationSignal, count: number, options: AdmissionOptions, ): AggregationSignal { POSITIVE_INTEGER_SCHEMA.parse(count) const definition = { count, primaryPosition: "last" as const, prerequisites: getCurrentSignalPrerequisites(), signal, ttl: options.ttl, ...consumeHookLocation(), } const boundaryDependencies = [ signal, ...(isSignal(definition.ttl) ? [definition.ttl] : []), ...(definition.prerequisites ? [definition.prerequisites] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) registerTake(definition, state) return createCoordinatedSignal( signal, definition, false, boundaryDependencies, ) } /** * Emits the first occurrence in each source partition. * * @param signal - Keyed source admitted once per partition. * @param options - Optional lifetime for the claim. */ export function once( signal: KeyedAggregationSignal, options?: AdmissionOptions, ): KeyedSignal> /** @inheritdoc */ export function once( signal: GlobalAggregationSignal, options?: AdmissionOptions, ): GlobalSignal> export function once( signal: AggregationSignal, options?: AdmissionOptions, ): AggregationSignal { return createTake(signal, 1, options ?? {}) } /** * Emits at most `limit` occurrences per rolling interval and partition. * * @param signal - Keyed source whose partitions own independent quotas. * @param options - Rolling capacity, interval, and overflow behavior. */ export function rateLimit( signal: KeyedAggregationSignal, options: RateLimitOptions, ): KeyedSignal> /** @inheritdoc */ export function rateLimit( signal: GlobalAggregationSignal, options: RateLimitOptions, ): GlobalSignal> export function rateLimit( signal: AggregationSignal, options: RateLimitOptions, ): AggregationSignal { POSITIVE_INTEGER_SCHEMA.parse(options.limit) const definition = { interval: options.interval, limit: options.limit, overflow: options.overflow ?? ("wait" as const), primaryPosition: "last" as const, prerequisites: getCurrentSignalPrerequisites(), signal, ...consumeHookLocation(), } const boundaryDependencies = [ signal, ...(isSignal(definition.interval) ? [definition.interval] : []), ...(definition.prerequisites ? [definition.prerequisites] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) registerRateLimit(definition, state) return createCoordinatedSignal( signal, definition, false, boundaryDependencies, ) } /** * Runs one durable section for every value and gathers its signal results. * * @param signal - Array emitted by one source context. * @param section - Symbolic item section whose returned signal is gathered. */ export function each( signal: Signal, section: (item: Signal) => Signal, ): Signal /** * Runs one durable side-effect section for every value without gathering. * * @param signal - Array emitted by one source context. * @param section - Symbolic item section traversed in every child context. */ export function each( signal: Signal, section: (item: Signal) => Signal | void, ): void export function each( signal: Signal, section: (item: Signal) => Signal | void, ): Signal | void { return group({ name: "Each" }, () => { const definition = { prerequisites: getCurrentSignalPrerequisites(), signal, ...consumeHookLocation(), } const boundaryDependencies = [ signal, ...(definition.prerequisites ? [definition.prerequisites] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) registerFanout(definition, state) const fanoutIndex = { getKey: (_value, context) => resolveFanoutItemMetadata(definition, context).batchContextId, getPosition: (_value, context) => resolveFanoutItemMetadata(definition, context).position, } satisfies SignalIndexDefinition const item = createSignal( { ancestorOrigins: unique([ ...boundaryDependencies.flatMap( (dependency) => dependency[signalAncestorOrigins], ), `signal:${formatHookLocation(definition)}`, ]), collect: (state, dependencies) => collectFanoutDependencies(definition, state, dependencies), origins: [`signal:${formatHookLocation(definition)}`], resolve: (context) => { const { position, source: sourceSelection } = resolveFanoutItemMetadata(definition, context) const source = signal[signalResolver]( selectDependencyContext(context, sourceSelection), ) if (source.status === "failed") return source if (source.status === "closed" || position >= source.value.length) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: fan-out source changed.`, ) } /* oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The persisted position is range-checked against this source array. */ const value = source.value[position] as T return { status: "succeeded", value, } }, }, { [signalIndex]: fanoutIndex }, ) const result = scope(() => withPrerequisites(item, () => section(item))) if (result === undefined) return const terminal = outcome(result) const count = signal.transform((items) => items.length) const gathered = collect( createSignal< SignalOutcome, { readonly [signalIndex]: typeof fanoutIndex } >( { ancestorOrigins: terminal[signalAncestorOrigins], collect: terminal[signalDependencyCollector], collectContext: terminal[signalContextDependencyCollector], contextOrigins: terminal[signalContextOrigins], origins: terminal[signalOrigins], resolve: terminal[signalResolver], }, { [signalIndex]: fanoutIndex }, ), count, ) return selectFirstNonclosing([ gate( signal.transform((): TResult[] => []), count.transform((value) => value === 0), ), gate( gathered.transform((results) => results.flatMap((result) => result.status === "succeeded" ? [result.value] : [], ), ), gathered.transform((results) => { const failedResult = results.find( (result) => result.status === "failed", ) if (failedResult?.status === "failed") { throw new FailedSignalError(failedResult.failure) } return results.every((result) => result.status === "succeeded") }), ), ]) }) } /** * Reshapes cross-context occurrences through a durable temporal policy. * * A selecting funnel emits one occurrence from every burst. Set `buffer` to * emit every occurrence in arrival order instead. Keyed sources maintain an * independent funnel for every key, and every result preserves that partition. * A start-only funnel emits the first occurrence immediately and ignores the * remainder of its burst. * * @param signal - One-shot signal whose occurrences enter the funnel. * @param options - Timing, edge, and selection policy. */ export function funnel< TSignal extends Signal, const TOptions extends FunnelOptions, >( signal: KeyedAggregationSignal, options: TOptions, ): FunnelResult>, "keyed", TOptions> /** @inheritdoc */ export function funnel< TSignal extends Signal, const TOptions extends FunnelOptions, >( signal: GlobalAggregationSignal, options: TOptions, ): FunnelResult>, "global", TOptions> export function funnel( signal: AggregationSignal, options: FunnelOptions, ): unknown { return createFunnel(signal, { buffer: options.buffer, inheritConflictingValuesFrom: options.inheritConflictingValuesFrom, maxBurstDuration: options.maxBurstDuration, minGap: options.minGap, minQuietPeriod: options.minQuietPeriod, select: options.select, triggerAt: options.triggerAt, until: options.until, }) } /** * Emits the latest occurrence after the source remains quiet. * * @param signal - One-shot signal whose occurrences are debounced. * @param duration - Required quiet period before the latest value emits. */ export function debounce( signal: KeyedAggregationSignal, duration: Duration, ): KeyedSignal> /** @inheritdoc */ export function debounce( signal: GlobalAggregationSignal, duration: Duration, ): GlobalSignal> export function debounce( signal: AggregationSignal, duration: Duration, ): unknown { return createFunnel(signal, { minQuietPeriod: duration }) } /** * Collects occurrences that arrive within a window opened by the first item. * * A keyed source creates an independent window for every key. The window emits * once its duration elapses and always contains at least its opening * occurrence. The nonempty result preserves the source partition. * * @param signal - One-shot signal whose occurrences are collected. * @param duration - Nonnegative millisecond or compact duration literal, or a * signal carrying either representation. * @param options - Parent value precedence for the resulting child context. * @throws When a literal duration is invalid. */ export function window( signal: KeyedAggregationSignal, duration: Duration, options?: ValueInheritanceOptions, ): KeyedSignal>> /** @inheritdoc */ export function window( signal: GlobalAggregationSignal, duration: Duration, options?: ValueInheritanceOptions, ): GlobalSignal>> /** @inheritdoc */ export function window( signal: KeyedAggregationSignal, options: WindowUntilOptions, ): KeyedSignal>> /** @inheritdoc */ export function window( signal: GlobalAggregationSignal, options: WindowUntilOptions, ): GlobalSignal>> export function window( signal: AggregationSignal, durationOrOptions: Duration | WindowUntilOptions, options: ValueInheritanceOptions = {}, ): AggregationSignal> { if (typeof durationOrOptions === "object" && !isSignal(durationOrOptions)) { return createFunnel(signal, { ...durationOrOptions, buffer: true, }) } return createFunnel(signal, { ...options, buffer: true, maxBurstDuration: durationOrOptions, }) } /** * Adds a boolean dependency that conditionally materializes another signal. * * @param value - Signal whose value is preserved when the gate opens. * @param condition - Boolean signal controlling whether the value emits. */ export function gate( value: TSignal, condition: Signal, ): InheritSignalCoordination> export function gate( value: Signal, condition: Signal, ): Signal { const definition = { condition, prerequisites: getCurrentSignalPrerequisites(), ...consumeHookLocation(), value, } const state = getAutomationRuntimeState() if (state) { collectGateDependencies(definition, state, { actionDependencyIds: new Set(), eventDependencyIds: new Set(), signalDependencyIds: new Set(), }) } return inheritSignalCoordination( value, createSignal({ ancestorOrigins: unique([ ...value[signalAncestorOrigins], ...condition[signalAncestorOrigins], ...(definition.prerequisites?.[signalAncestorOrigins] ?? []), ]), collect: (state, dependencies) => collectGateDependencies(definition, state, dependencies), collectContext: (state, dependencies) => combineSignalDependencyOutcomes( [ value, condition, ...(definition.prerequisites ? [definition.prerequisites] : []), ].map((signal) => signal[signalContextDependencyCollector](state, dependencies), ), ), contextOrigins: unique([ ...value[signalContextOrigins], ...condition[signalContextOrigins], ...(definition.prerequisites?.[signalContextOrigins] ?? []), ]), origins: [`signal:${formatHookLocation(definition)}`], resolve: (context) => { const outcome = findSignalBoundary(definition, context) if (outcome.status === "closed") return { status: "closed" } if (outcome.status === "failed") { return { failure: outcome.failure, status: "failed" } } if (outcome.status !== "open") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: gate outcome changed.`, ) } return value[signalResolver](context) }, }), ) } /** * Selects the first signal by declaration order that does not close. * * A pending earlier signal blocks later signals. Failure is selected just like * an emission, while closure advances to the next input. * * @param inputs - Ordered signals to inspect. */ function selectFirstNonclosing< const TInputs extends readonly [Signal, Signal, ...Signal[]], >(inputs: TInputs): Signal> { return createSignal>({ ancestorOrigins: unique( inputs.flatMap((input) => input[signalAncestorOrigins]), ), collect: (state, dependencies) => { let outcomeSeq = 0 for (const input of inputs) { const outcome = input[signalDependencyCollector](state, dependencies) if (outcome.status !== "closed") return outcome outcomeSeq = Math.max(outcomeSeq, outcome.outcomeSeq) } return { outcomeSeq, status: "closed" } }, collectContext: (state, dependencies) => { let outcomeSeq = 0 for (const input of inputs) { const outcome = input[signalContextDependencyCollector]( state, dependencies, ) if (outcome.status !== "closed") return outcome outcomeSeq = Math.max(outcomeSeq, outcome.outcomeSeq) } return { outcomeSeq, status: "closed" } }, contextOrigins: unique( inputs.flatMap((input) => input[signalContextOrigins]), ), origins: unique(inputs.flatMap((input) => input[signalOrigins])), resolve: (context) => { for (const input of inputs) { const outcome = input[signalResolver](context) if (outcome.status !== "closed") { // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The selected input is a member of TInputs. return outcome as SignalOutcome> } } return { status: "closed" } }, }) } /** * Emits every successful occurrence from independent input streams. * * Each occurrence creates its own inherited child context. Inputs never wait * for or consume one another. * * @param inputs - Independent one-shot signals contributing occurrences. */ export function merge< const TInputs extends readonly [Signal, Signal, ...Signal[]], >(inputs: TInputs): Signal> { const prerequisites = getCurrentSignalPrerequisites() const definition = { inputOrigins: inputs.map((input) => input[signalOrigins]), inputs, prerequisites, ...consumeHookLocation(), } const boundaryDependencies = [ ...inputs, ...(prerequisites ? [prerequisites] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) registerMerge(definition, state) return createSignal>({ ancestorOrigins: unique([ ...boundaryDependencies.flatMap( (signal) => signal[signalAncestorOrigins], ), `signal:${formatHookLocation(definition)}`, ]), collect: (state, dependencies) => collectCompletedSignalDependencies(definition, state, dependencies), origins: [`signal:${formatHookLocation(definition)}`], resolve: (context) => { const outcome = findSignalBoundary(definition, context) if (outcome.status === "closed") return { status: "closed" } if (outcome.status === "failed") { return { failure: outcome.failure, status: "failed" } } if ( outcome.status !== "succeeded" || outcome.selectedIndex === undefined ) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: merge outcome changed.`, ) } const selected = inputs[outcome.selectedIndex] if (!selected) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: merge inputs changed.`, ) } // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The selected input is a member of TInputs. return selected[signalResolver](context) as SignalOutcome< InferSignal > }, }) } /** * Runs one structured work region at a time within each source partition. * * Independent declarations inside the section remain concurrent. The next * occurrence is admitted only after every context created by the current * section becomes terminal. * * @param signal - Keyed or global occurrence stream defining the queue. * @param section - Synchronous declaration region protected by admission. */ export function serialize< TSignal extends Signal, TSection extends () => unknown, >( signal: KeyedAggregationSignal, section: Synchronous extends never ? never : TSection, ): ReturnType /** @inheritdoc */ export function serialize< TSignal extends Signal, TSection extends () => unknown, >( signal: GlobalAggregationSignal, section: Synchronous extends never ? never : TSection, ): ReturnType export function serialize( signal: AggregationSignal, section: () => unknown, ): unknown { return createExecutionRegion(signal, section, { policy: "serialize" }) } /** * Runs a bounded number of structured work regions per partition. * * @param signal - Keyed source whose occurrences enter the queue. * @param options - Maximum active regions per partition. * @param section - Synchronous declaration traversed inside each region. */ export function concurrent< TSignal extends Signal, TSection extends () => unknown, >( signal: KeyedAggregationSignal, options: ConcurrencyOptions, section: Synchronous extends never ? never : TSection, ): ReturnType /** @inheritdoc */ export function concurrent< TSignal extends Signal, TSection extends () => unknown, >( signal: GlobalAggregationSignal, options: ConcurrencyOptions, section: Synchronous extends never ? never : TSection, ): ReturnType export function concurrent( signal: AggregationSignal, options: ConcurrencyOptions, section: () => unknown, ): unknown { POSITIVE_INTEGER_SCHEMA.parse(options.limit) return createExecutionRegion(signal, section, { limit: options.limit, policy: "concurrent", }) } /** * Builds the shared serialized or bounded-concurrency region. * * @param signal - Partitioned source occurrence. * @param section - Synchronous declarations inside the region. * @param policy - Single or bounded active-region capacity. */ function createExecutionRegion( signal: AggregationSignal, section: () => unknown, policy: { policy: "serialize" } | { limit: number; policy: "concurrent" }, ): unknown { const prerequisites = getCurrentSignalPrerequisites() const admissionLocation = consumeHookLocation() const releaseLocation = consumeHookLocation() const definition = { ...admissionLocation, ...(policy.policy === "concurrent" && { limit: policy.limit }), prerequisites, releaseScopePath: releaseLocation.scopePath, releaseSlot: releaseLocation.slot, signal, } const boundaryDependencies = [ signal, ...(prerequisites ? [prerequisites] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) { if (policy.policy === "concurrent") { registerConcurrency({ ...definition, limit: policy.limit }, state) } else { registerSerialization(definition, state) } } const admission = createSignal({ ancestorOrigins: unique([ ...boundaryDependencies.flatMap( (dependency) => dependency[signalAncestorOrigins], ), `signal:${formatHookLocation(admissionLocation)}`, ]), collect: (state, dependencies) => collectCompletedSignalDependencies( admissionLocation, state, dependencies, ), origins: [`signal:${formatHookLocation(admissionLocation)}`], resolve: (context) => { const outcome = findSignalBoundary(admissionLocation, context) if (outcome.status === "closed") return { status: "closed" } if (outcome.status === "failed") { return { failure: outcome.failure, status: "failed" } } if (outcome.status !== "succeeded") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(admissionLocation)}: execution-region admission changed.`, ) } return signal[signalResolver](context) }, }) const result = scope(() => withPrerequisites(admission, section)) const releaseDependencies = [admission, ...(isSignal(result) ? [result] : [])] if (planningContext) { planSignalContextBoundary( planningContext, releaseLocation, releaseDependencies, false, ) } // Keep the durable release marker distinct from its admission guard. const release = createSignal({ ancestorOrigins: unique([ ...releaseDependencies.flatMap( (dependency) => dependency[signalAncestorOrigins], ), `signal:${formatHookLocation(releaseLocation)}`, ]), collect: (state, dependencies) => collectCompletedSignalDependencies(releaseLocation, state, dependencies), origins: [`signal:${formatHookLocation(releaseLocation)}`], resolve: (context) => { const outcome = findSignalBoundary(releaseLocation, context) if (outcome.status === "failed") { return { failure: outcome.failure, status: "failed" } } if (outcome.status === "closed") return { status: "closed" } if (outcome.status !== "succeeded") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(releaseLocation)}: execution-region release changed.`, ) } return { status: "succeeded", value: null } }, }) const completion = dependentOn(release, admission) return isSignal(result) ? dependentOn(result, completion) : result } /** * Selects the first causally related input to emit or fail. * * Inputs sharing the same activation boundary compete in one durable race. * Closed inputs leave the race, and it closes only when every input closes. * * @param inputs - One-shot signals competing to supply the result. * @throws {TypeError} When the inputs don't share an activation boundary. */ export function race< const TInputs extends readonly [Signal, Signal, ...Signal[]], >(inputs: TInputs): Signal> { const prerequisites = getCurrentSignalPrerequisites() const contextPrerequisite = getCurrentContextPrerequisite() const explicitCohort = contextPrerequisite ?? prerequisites const sharedAncestorOrigins = inputs[0][signalAncestorOrigins].filter( (origin) => inputs .slice(1) .every((input) => input[signalAncestorOrigins].includes(origin)), ) const definition = { cohort: explicitCohort, cohortOrigins: explicitCohort ? explicitCohort[signalContextOrigins] : sharedAncestorOrigins.slice(-1), inputOrigins: inputs.map((input) => input[signalOrigins]), inputs, prerequisiteOrigins: prerequisites?.[signalOrigins] ?? [], prerequisites, ...consumeHookLocation(), } const boundaryDependencies = [ ...inputs, ...(prerequisites ? [prerequisites] : []), ...(contextPrerequisite ? [contextPrerequisite] : []), ] const planningContext = getAutomationPlanningContext() if (planningContext) { if (definition.cohortOrigins.length === 0) { throw new TypeError( "race inputs must share one activation boundary; use merge for independent trigger streams.", ) } planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) { collectRaceDependencies(definition, state, { actionDependencyIds: new Set(), eventDependencyIds: new Set(), signalDependencyIds: new Set(), }) } return createSignal>({ ancestorOrigins: unique([ ...boundaryDependencies.flatMap( (signal) => signal[signalAncestorOrigins], ), `signal:${formatHookLocation(definition)}`, ]), collect: (state, dependencies) => collectRaceDependencies(definition, state, dependencies), origins: [`signal:${formatHookLocation(definition)}`], resolve: (context) => { const outcome = findSignalBoundary(definition, context) if (outcome.status === "closed") return { status: "closed" } if (outcome.status === "failed") { return { failure: outcome.failure, status: "failed" } } if ( outcome.status !== "succeeded" || outcome.selectedIndex === undefined ) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: race outcome changed.`, ) } const selected = inputs[outcome.selectedIndex] if (!selected) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(definition)}: race inputs changed.`, ) } // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The selected input is a member of TInputs. return selected[signalResolver](context) as SignalOutcome< InferSignal > }, }) } /** * Materializes a signal only when its value satisfies a predicate. * * @param signal - Signal whose value is tested and preserved. * @param predicate - Pure predicate evaluated after the value materializes. */ export function filter< TSignal extends Signal, TNarrowed extends InferSignal, >( signal: TSignal, predicate: (value: InferSignal) => value is TNarrowed, ): InheritSignalCoordination /** * Materializes a signal only when its value satisfies a predicate. * * @param signal - Signal whose value is tested and preserved. * @param predicate - Pure predicate evaluated after the value materializes. */ export function filter( signal: TSignal, predicate: (value: InferSignal) => boolean, ): InheritSignalCoordination> export function filter( signal: Signal, predicate: (value: T) => boolean, ): unknown { return gate(signal, signal.transform(predicate)) } /** * Splits a signal into narrowed matching and nonmatching materialization paths. * * @param signal - Signal whose value is routed to exactly one output. * @param predicate - Type guard selecting the first output when true. */ export function partition< TSignal extends Signal, TNarrowed extends InferSignal, >( signal: TSignal, predicate: (value: InferSignal) => value is TNarrowed, ): readonly [ matched: InheritSignalCoordination, unmatched: InheritSignalCoordination< TSignal, Exclude, TNarrowed> >, ] /** * Splits a signal into matching and nonmatching materialization paths. * * @param signal - Signal whose value is routed to exactly one output. * @param predicate - Pure predicate selecting the first output when true. */ export function partition( signal: TSignal, predicate: (value: InferSignal) => boolean, ): readonly [ matched: InheritSignalCoordination>, unmatched: InheritSignalCoordination>, ] export function partition( signal: Signal, predicate: (value: T) => boolean, ): unknown { const condition = signal.transform(predicate) return [ gate(signal, condition), gate( signal, condition.transform((matches) => !matches), ), ] } /** Presentation options for an author-facing group. */ export interface GroupOptions { /** Author-facing group name. */ name: string /** Initial UI treatment for declarations nested in this group. */ presentation?: "collapsed" | "expanded" | "hidden" } type SignalPrerequisites = Signal | readonly [Signal, ...Signal[]] /** * Traverses a synchronous section in an isolated durable hook namespace. * * @param fn - Synchronous automation section to traverse. */ export function scope unknown>( fn: Synchronous extends never ? never : TSection, ): ReturnType export function scope(fn: () => unknown): unknown { return withHookScope(fn) } /** * Traverses an author-facing group in an isolated durable hook namespace. * * @param options - Group name and initial presentation. * @param fn - Synchronous automation section to traverse. */ export function group unknown>( options: GroupOptions, fn: Synchronous extends never ? never : TSection, ): ReturnType export function group(options: GroupOptions, fn: () => unknown): unknown { return scope(() => { getAutomationPlanningContext()?.scopes.push({ name: options.name, path: [...getCurrentHookScopePath()], presentation: options.presentation ?? "collapsed", }) return fn() }) } /** * Preserves a signal's value while attaching additional dependencies. * * @param signal - Value-producing signal to preserve. * @param prerequisites - Signals that must settle successfully first. */ export function dependentOn( signal: TSignal, prerequisites: SignalPrerequisites, ): InheritSignalCoordination> export function dependentOn( signal: Signal, prerequisites: SignalPrerequisites, ): Signal { return signal[signalDependencyBinder]( isSignal(prerequisites) ? prerequisites : transform(prerequisites, () => null), ) } /** * Makes every durable declaration in a synchronous section wait for signals. * * @param prerequisites - Signals required by declarations in the section. * @param fn - Synchronous automation section to traverse. */ export function withPrerequisites unknown>( prerequisites: SignalPrerequisites, fn: Synchronous extends never ? never : TSection, ): ReturnType export function withPrerequisites( prerequisites: SignalPrerequisites, fn: () => unknown, ): unknown { const local = isSignal(prerequisites) ? prerequisites : transform(prerequisites, () => null) const parent = getCurrentSignalPrerequisites() const dependency = parent ? dependentOn(local, parent) : local PREREQUISITE_SCOPES.push(dependency) try { const result = fn() return isSignal(result) ? dependentOn(result, dependency) : result } finally { PREREQUISITE_SCOPES.pop() } } /** * Runs a section with context placement but without waiting for a value. * * @param prerequisite - Signal whose activation context places the section. * @param fn - Synchronous section to traverse. */ export function withContextPrerequisite( prerequisite: Signal, fn: () => TResult, ): TResult { CONTEXT_PREREQUISITE_SCOPES.push(prerequisite) try { return fn() } finally { CONTEXT_PREREQUISITE_SCOPES.pop() } } /** * Runs internal coordination outside the caller's ambient prerequisite. * * @param fn - Internal composition section to traverse. */ export function withoutSignalPrerequisites( fn: () => TResult, ): TResult { PREREQUISITE_SCOPES.push(undefined) CONTEXT_PREREQUISITE_SCOPES.push(undefined) try { return fn() } finally { CONTEXT_PREREQUISITE_SCOPES.pop() PREREQUISITE_SCOPES.pop() } } /** * Traverses complementary automation sections selected by a boolean signal. * * @param condition - Boolean signal selecting one branch. * @param whenTrue - Automation section traversed under the truthy partition. * @param tail - Alternating else-if conditions and sections, followed by an * optional final else section. */ export function branch< TTrueSection extends BranchSection, const TTail extends BranchTail, >( condition: Signal, whenTrue: TTrueSection, ...tail: TTail & ValidBranchTail, TTail> ): BranchResult, BranchTailReturn> export function branch( condition: Signal, whenTrue: BranchSection, ...tail: readonly (Signal | BranchSection)[] ): Signal | void { return branchNested(condition, whenTrue, tail) } /** * Converts every terminal state into an ordinary emitted value. * * @param signal - Signal whose terminal outcome is observed. */ export function outcome(signal: Signal): Signal> { return createSignal>({ ancestorOrigins: signal[signalAncestorOrigins], collect: (state, dependencies) => { const source = signal[signalDependencyCollector](state, dependencies) return source.status === "pending" ? source : { outcomeSeq: source.outcomeSeq, status: "succeeded" } }, collectContext: signal[signalContextDependencyCollector], contextOrigins: signal[signalContextOrigins], origins: signal[signalOrigins], resolve: (context) => { const source = signal[signalResolver](context) return { ...(source.status === "succeeded" && source.sensitivity && { sensitivity: true as const }), status: "succeeded", value: source, } }, }) } /** * Emits whether a source eventually fails. * * @param signal - Signal whose terminal outcome is tested. */ export function failed(signal: Signal): Signal { return outcome(signal).transform((result) => result.status === "failed") } /** * Emits whether a source eventually succeeds with a value. * * @param signal - Signal whose terminal outcome is tested. */ export function succeeded(signal: Signal): Signal { return outcome(signal).transform((result) => result.status === "succeeded") } /** * Emits whether a source eventually closes without a value. * * @param signal - Signal whose terminal outcome is tested. */ export function closed(signal: Signal): Signal { return outcome(signal).transform((result) => result.status === "closed") } /** * Emits the failure only when a source fails, and closes otherwise. * * @param signal - Signal whose failure is selected. */ export function onFailure(signal: Signal): Signal /** @inheritdoc */ export function onFailure< TSignal extends Signal, TSection extends ( failure: Signal, source: TSignal, ) => unknown, >( signal: TSignal, section: Synchronous extends never ? never : TSection, ): ReturnType export function onFailure( signal: Signal, section?: (failure: Signal, source: Signal) => unknown, ): unknown { const failure = filter( outcome(signal), (result) => result.status === "failed", ).transform((result) => result.failure) return section ? withPrerequisites(failure, () => section(failure, signal)) : failure } /** * Emits the source value only when it succeeds, and closes otherwise. * * @param signal - Signal whose successful value is selected. * @param section - Optional synchronous section selected by success. */ export function onSuccess< TSignal extends Signal, TSection extends ( value: InheritSignalCoordination>, source: TSignal, ) => unknown, >( signal: TSignal, section: Synchronous extends never ? never : TSection, ): ReturnType /** @inheritdoc */ export function onSuccess(signal: KeyedSignal): KeyedSignal /** @inheritdoc */ export function onSuccess(signal: GlobalSignal): GlobalSignal /** @inheritdoc */ export function onSuccess(signal: Signal): Signal export function onSuccess( signal: Signal, section?: (value: Signal, source: Signal) => unknown, ): unknown { const value = inheritSignalCoordination( signal, filter( outcome(signal), (result) => result.status === "succeeded", ).transform((result) => result.value), ) return section ? withPrerequisites(value, () => section(value, signal)) : value } /** * Emits `null` only when a source closes without a value. * * @param signal - Signal whose empty closure is selected. */ export function onClose(signal: Signal): Signal /** @inheritdoc */ export function onClose< TSignal extends Signal, TSection extends (closure: Signal, source: TSignal) => unknown, >( signal: TSignal, section: Synchronous extends never ? never : TSection, ): ReturnType export function onClose( signal: Signal, section?: (closure: Signal, source: Signal) => unknown, ): unknown { const closure = filter( outcome(signal), (result) => result.status === "closed", ).transform(() => null) return section ? withPrerequisites(closure, () => section(closure, signal)) : closure } /** Returns the prerequisites inherited from active synchronous sections. */ export function getCurrentSignalPrerequisites(): Signal | undefined { return PREREQUISITE_SCOPES.at(-1) } /** Returns the context-only placement prerequisite active for actions. */ export function getCurrentContextPrerequisite(): Signal | undefined { return CONTEXT_PREREQUISITE_SCOPES.at(-1) } /** * Declares a temporal coordinator and returns its selected or buffered signal. * * @param signal - Explicitly partitioned occurrence signal. * @param options - Timing and output policy. */ function createFunnel( signal: AggregationSignal, options: FunnelImplementationOptions & { buffer: true }, ): AggregationSignal> /** @inheritdoc */ function createFunnel( signal: AggregationSignal, options: FunnelImplementationOptions, ): Signal | AggregationSignal> function createFunnel( signal: AggregationSignal, options: FunnelImplementationOptions, ): Signal | AggregationSignal> { const selection = options.buffer ? ("all" as const) : options.triggerAt === "start" || options.select === "first" ? ("first" as const) : ("last" as const) const definition = { maxBurstDuration: options.maxBurstDuration, minGap: options.minGap, minQuietPeriod: options.minQuietPeriod, primaryPosition: selection === "all" ? (options.inheritConflictingValuesFrom ?? ("last" as const)) : selection, prerequisites: getCurrentSignalPrerequisites(), selection, signal, triggerAt: options.triggerAt ?? ("end" as const), until: options.until, ...consumeHookLocation(), } const boundaryDependencies = [ signal, ...[ definition.maxBurstDuration, definition.minGap, definition.minQuietPeriod, definition.until, definition.prerequisites, ].filter((value): value is Signal => isSignal(value)), ] const planningContext = getAutomationPlanningContext() if (planningContext) { planSignalContextBoundary( planningContext, definition, boundaryDependencies, false, ) } const state = getAutomationRuntimeState() if (state) registerFunnel(definition, state) return options.buffer ? createCoordinatedSignal(signal, definition, true, boundaryDependencies) : createCoordinatedSignal(signal, definition, false, boundaryDependencies) } /** * Creates the shared returned signal for a collection or funnel. * * @param source - Occurrence signal reconstructed in each selected parent. * @param location - Durable coordinated hook identity. * @param buffer - Whether every selected value is returned. * @param ancestorSignals - Signals inherited by the coordinated child. */ function createCoordinatedSignal( source: AggregationSignal, location: HookLocation, buffer: true, ancestorSignals: readonly Signal[], ): AggregationSignal> /** @inheritdoc */ function createCoordinatedSignal( source: AggregationSignal, location: HookLocation, buffer: false, ancestorSignals: readonly Signal[], ): AggregationSignal function createCoordinatedSignal( source: AggregationSignal, location: HookLocation, buffer: boolean, ancestorSignals: readonly Signal[], ): Signal | AggregationSignal> { const collect = ( state: AutomationRunState, dependencies: SignalDependencySets, ) => collectCompletedSignalDependencies(location, state, dependencies) const origins = [`signal:${formatHookLocation(location)}`] const ancestorOrigins = unique([ ...ancestorSignals.flatMap((signal) => signal[signalAncestorOrigins]), ...origins, ]) const resolveSelection = (context: RuntimeValueContext) => { const outcome = findSignalBoundary(location, context) if (outcome.status === "closed") return { status: "closed" } as const if (outcome.status === "failed") return outcome if (outcome.status !== "succeeded") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: coordinated outcome changed.`, ) } const selection = context.coordinationSelections.find( (candidate) => candidate.signalInvocationId === outcome.signalInvocationId, ) if (!selection) { throw new Error( `Missing coordinated values for hook ${formatHookLocation(location)}.`, ) } return { status: "succeeded", value: selection } as const } if (buffer) { return inheritAggregateSignalCoordination( source, createSignal>({ ancestorOrigins, collect, origins, resolve: (context) => { const selection = resolveSelection(context) if (selection.status !== "succeeded") return selection const [first, ...remaining] = selection.value.items if (!first) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: buffered selection changed.`, ) } const outcomes = [first, ...remaining].map((item) => source[signalResolver](selectDependencyContext(context, item)), ) const failed = outcomes.find( ( outcome, ): outcome is Extract, { status: "failed" }> => outcome.status === "failed", ) if (failed) return failed if (outcomes.some((outcome) => outcome.status === "closed")) { return { status: "closed" } } const [firstOutcome, ...remainingOutcomes] = outcomes if (firstOutcome?.status !== "succeeded") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: buffered source changed.`, ) } return { ...(outcomes.some( (outcome) => outcome.status === "succeeded" && outcome.sensitivity, ) && { sensitivity: true as const }), status: "succeeded", value: [ firstOutcome.value, ...remainingOutcomes.map((outcome) => { if (outcome.status !== "succeeded") { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: buffered source changed.`, ) } return outcome.value }), ], } }, }), ) } return inheritSignalCoordination( source, createSignal({ ancestorOrigins, collect, origins, resolve: (context) => { const selection = resolveSelection(context) if (selection.status !== "succeeded") return selection const item = selection.value.items[0] if (!item || selection.value.items.length !== 1) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: funnel selection changed.`, ) } return source[signalResolver](selectDependencyContext(context, item)) }, }), ) } /** * Preserves a source partition on an aggregate of values from that partition. * * @param source - Partitioned signal whose values formed the aggregate. * @param result - Nonempty aggregate reconstructed from selected occurrences. */ function inheritAggregateSignalCoordination( source: AggregationSignal, result: Signal>, ): AggregationSignal> { const index = source[signalIndex] const behavior = { ancestorOrigins: result[signalAncestorOrigins], collect: result[signalDependencyCollector], collectContext: result[signalContextDependencyCollector], contextOrigins: result[signalContextOrigins], derivation: { inputs: [result], type: "transparent" as const }, origins: result[signalOrigins], resolve: result[signalResolver], } if (!index) return createSignal(behavior, { [signalGlobal]: true }) return createSignal(behavior, { [signalIndex]: { getKey: (values: NonEmptyArray, context: RuntimeValueContext) => index.getKey(values[0], context), } satisfies SignalIndexDefinition>, }) } /** * Finds the persisted output for one durable signal hook. * * @param location - Durable signal hook identity. * @param context - Current hydrated context. * @throws When the boundary is missing or ambiguous. */ function findSignalBoundary( location: HookLocation, context: RuntimeValueContext, ) { const formatted = formatHookLocation(location) const outcome = findUniqueSignalBoundary( context.signalOutputs, (candidate) => isSameHookLocation(candidate, location), `signal hook ${formatted}`, context, ) if (!outcome) { throw new Error(`Missing signal dependency for hook ${formatted}.`) } return outcome } /** * Resolves the hidden source batch and item position for one fan-out child. * * @param location - Fan-out hook identifying the child marker. * @param context - Materialized child context and coordinator selections. * @throws {AutomationCompatibilityError} When child provenance changed. */ function resolveFanoutItemMetadata( location: HookLocation, context: RuntimeValueContext, ) { const outcome = findSignalBoundary(location, context) if (outcome.status !== "succeeded" || outcome.selectedIndex === undefined) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: fan-out outcome changed.`, ) } const selection = context.coordinationSelections.find( (candidate) => candidate.signalInvocationId === outcome.signalInvocationId, ) const source = selection?.items[0] if (!source || selection.items.length !== 1) { throw new AutomationCompatibilityError( `Automation compatibility error at signal hook ${formatHookLocation(location)}: fan-out source changed.`, ) } return { batchContextId: source.contextId, position: outcome.selectedIndex, source, } } /** * Expands a flat else-if tail through the binary branch primitive. * * @param condition - Current branch condition. * @param whenTrue - Section selected by a true condition. * @param tail - Remaining else-if pairs and optional else section. */ function branchNested( condition: Signal, whenTrue: BranchSection, tail: readonly (Signal | BranchSection)[], ): Signal | void { const [thenDependency, elseDependency] = partition( condition, (value) => value, ) const trueResult = scope(() => withPrerequisites(thenDependency, whenTrue)) const [next, nextSection, ...remaining] = tail const falseResult = typeof next === "function" ? scope(() => withPrerequisites(elseDependency, next)) : next && typeof nextSection === "function" ? scope(() => withPrerequisites(elseDependency, () => branchNested(next, nextSection, remaining), ), ) : undefined if (isSignal(trueResult)) { return isSignal(falseResult) ? selectFirstNonclosing([trueResult, falseResult]) : trueResult } } /** Synchronous automation section accepted by a branch arm. */ export type BranchSection = () => Signal | void /** Alternating else-if conditions and sections with an optional final else. */ export type BranchTail = readonly (Signal | BranchSection)[] /** Returns the union produced by every section in a branch tail. */ export type BranchTailReturn = TTail extends readonly [] ? undefined : TTail extends readonly [infer TElse extends BranchSection] ? ReturnType : TTail extends readonly [ Signal, infer TSection extends BranchSection, ...infer TRest extends BranchTail, ] ? ReturnType | BranchTailReturn : never /** Rejects branch tails that mix signal-returning and side-effect sections. */ export type ValidBranchTail< TFirst, TTail extends BranchTail, > = TTail extends readonly [] ? unknown : TTail extends readonly [infer TElse extends BranchSection] ? ValidBranchPair> : TTail extends readonly [ Signal, infer TSection extends BranchSection, ...infer TRest extends BranchTail, ] ? ValidBranchPair> & ValidBranchTail : never /** Rejects two branch sections that disagree about returning a signal. */ type ValidBranchPair = TOther extends undefined ? unknown : TFirst extends Signal ? TOther extends Signal ? unknown : never : TOther extends Signal ? never : unknown /** Computes the signal union returned by compatible branch sections. */ export type BranchResult = TTrue extends Signal ? Signal> : void /** Extracts the emitted value union from signal-returning branch sections. */ type BranchSignalValue = TResult extends Signal ? TValue : never