import { projectSensitivityMask, type Encodable, type SensitivityMask, } from "@automate.ax/codec" import { AUTOMATION_ERROR_CAUSE_DEPTH, type AutomationErrorValue, } from "@automate.ax/api-contract/errors" import type { SignalDerivation } from "@automate.ax/api-contract/runtime" import type ms from "ms" import { unique } from "remeda" import * as z from "zod/mini" import type { AutomationRunState, RuntimeValueContext } from "./runtime" import type { SignalCoreMethods, SignalMethods, SignalRegisteredMethods, SignalReservedMethods, } from "./signal-methods" export const signalDependencyCollector: unique symbol = Symbol( "automation.signalDependencyCollector", ) export const signalDependencyBinder: unique symbol = Symbol( "automation.signalDependencyBinder", ) const SIGNAL_DERIVATION_DEFINITION: unique symbol = Symbol( "automation.signalDerivationDefinition", ) export const signalContextDependencyCollector: unique symbol = Symbol( "automation.signalContextDependencyCollector", ) export const signalContextOrigins: unique symbol = Symbol( "automation.signalContextOrigins", ) export const signalAncestorOrigins: unique symbol = Symbol( "automation.signalAncestorOrigins", ) export const signalGlobal: unique symbol = Symbol("automation.signalGlobal") export const signalIndex: unique symbol = Symbol("automation.signalIndex") export const signalOrigins: unique symbol = Symbol("automation.signalOrigins") export const signalResolver: unique symbol = Symbol("automation.signalResolver") export const signalValueType: unique symbol = Symbol( "automation.signalValueType", ) const JSON_VALUE_SCHEMA = z.json() /** Author-safe normalized error exposed by signal lifecycle APIs. */ export type AutomationError = AutomationErrorValue /** A signal's mutually exclusive terminal outcome. */ export type SignalOutcome = | { status: "closed" } | { failure: AutomationError; status: "failed" } | { sensitivity?: SensitivityMask; status: "succeeded"; value: T } export type SignalDependencyResolution = | { outcomeSeq: number; status: "closed" } | { failure: AutomationError; outcomeSeq: number; status: "failed" } | { status: "pending" } | { outcomeSeq: number; status: "succeeded" } export interface SignalDependencySets { actionDependencyIds: Set eventDependencyIds: Set signalDependencyIds: Set } /** Compact duration literal accepted by the platform duration parser. */ export type DurationString = ms.StringValue /** Rejects callbacks whose inferred result is asynchronous. */ export type Synchronous unknown> = [ReturnType] extends [never] ? unknown : ReturnType extends PromiseLike ? never : unknown interface SignalBehavior { ancestorOrigins?: readonly string[] collect( this: void, state: AutomationRunState, dependencies: SignalDependencySets, ): SignalDependencyResolution collectContext?: SignalBehavior["collect"] contextOrigins?: readonly string[] derivation?: SignalDerivationDefinition origins: readonly string[] resolve(this: void, context: RuntimeValueContext): SignalOutcome } type SignalDerivationDefinition = | { inputs: readonly Signal[]; property: string; type: "property" } | { inputs: readonly Signal[]; type: "transform" | "transparent" } const SIGNAL_METHOD_IMPLEMENTATIONS: Record = {} export interface SignalIndexDefinition { getKey(value: T, context: RuntimeValueContext): Encodable getPosition?(value: T, context: RuntimeValueContext): number } /** Maps a payload's properties to child signals. */ type SignalProperties = { readonly [TKey in keyof NonNullable as TKey extends | keyof SignalProtocol | keyof SignalCoreMethods | keyof SignalReservedMethods ? never : TKey extends string | number ? TKey : never]-?: NonNullable[TKey] extends ( ...arguments_: never[] ) => unknown ? never : Signal[TKey]> } type SignalProtocol = { readonly [signalAncestorOrigins]: readonly string[] readonly [signalContextDependencyCollector]: SignalBehavior["collect"] readonly [signalContextOrigins]: readonly string[] readonly [signalDependencyBinder]: (dependency: Signal) => Signal readonly [signalDependencyCollector]: SignalBehavior["collect"] readonly [SIGNAL_DERIVATION_DEFINITION]?: SignalDerivationDefinition readonly [signalGlobal]?: true readonly [signalIndex]?: SignalIndexDefinition readonly [signalOrigins]: readonly string[] readonly [signalResolver]: SignalBehavior["resolve"] readonly [signalValueType]: () => T } /** * A durable symbolic reference to a value that will be materialized during * automation execution. * * Signals compose synchronously while their values are unavailable. Accessing a * property creates a derived signal, so `request.path` is equivalent to * `request.transform((value) => value.path)`. Property access can be chained * across objects and primitives; for example, `request.path.length` is a * `Signal`. Resolving a property derived from a `null` or `undefined` * value throws a `TypeError`. * * Function-valued properties are reserved for signal methods rather than * projected as callable values. Use `transform` to invoke a method on the * materialized value. * * @template T - Value represented by the signal. */ export type Signal = SignalProtocol & SignalCoreMethods & SignalMethods & SignalProperties & object /** A signal carrying a durable cross-context key definition. */ export type KeyedSignal = Signal & { readonly [signalIndex]: SignalIndexDefinition } /** A signal explicitly using one shared partition per cross-context operator. */ export type GlobalSignal = Signal & { readonly [signalGlobal]: true } /** Preserves a source's explicit keyed or global coordination mode. */ export type InheritSignalCoordination< TSignal extends Signal, TValue, > = TSignal extends KeyedSignal ? KeyedSignal : TSignal extends GlobalSignal ? GlobalSignal : Signal export type InferSignal = TSignal extends TSignal ? ReturnType : never /** Internal control-flow marker for a signal that closed without emitting. */ export class ClosedSignalError extends Error { override name = "ClosedSignalError" } /** Internal control-flow marker carrying a failed signal's details. */ export class FailedSignalError extends Error { override name = "FailedSignalError" /** @param failure - Failure propagated from the resolved signal. */ constructor(readonly failure: AutomationError) { super(failure.message) } } /** Raised when merged parent histories expose multiple values at one slot. */ export class AmbiguousSignalError extends Error { override name = "AmbiguousSignalError" } /** * Determines whether a value implements the internal signal protocol. * * @param value - Value to inspect. */ export function isSignal(value: unknown): value is Signal { return ( typeof value === "object" && value !== null && signalValueType in value && typeof value[signalValueType] === "function" ) } /** * Installs the public fluent adapters used by every signal proxy. * * @param implementations - Source-first functions keyed by fluent method name. */ export function registerSignalMethodImplementations( implementations: Readonly>, ): void { Object.assign(SIGNAL_METHOD_IMPLEMENTATIONS, implementations) } /** * Combines dependencies that must all succeed. * * @param outcomes - Dependency states to combine. */ export function combineSignalDependencyOutcomes( outcomes: SignalDependencyResolution[], ): SignalDependencyResolution { if (outcomes.some((outcome) => outcome.status === "pending")) { return { status: "pending" } } const outcomeSeq = Math.max( 0, ...outcomes.map((outcome) => outcome.status === "pending" ? 0 : outcome.outcomeSeq, ), ) const failure = outcomes.find( ( outcome, ): outcome is Extract => outcome.status === "failed", ) if (failure) return { ...failure, outcomeSeq } return outcomes.some((outcome) => outcome.status === "closed") ? { outcomeSeq, status: "closed" } : { outcomeSeq, status: "succeeded" } } /** * Materializes the explanatory source, property, and transform graph consumed * by one durable declaration. * * @param signals - Signal roots consumed by the declaration. * @param state - Current runtime occurrence used to resolve durable sources. * @param selectedDependencies - Exact durable sources selected by dependency * traversal. */ export function buildSignalDerivation( signals: readonly Signal[], state: AutomationRunState, ...selectedDependencies: [SignalDependencySets, ...SignalDependencySets[]] ): SignalDerivation | null { const nodes: SignalDerivation["nodes"] = [] const rootsBySignal = new WeakMap() const sourceIndexes = new Map() const addSource = (source: "action" | "event" | "signal", id: string) => { if ( !selectedDependencies.some((dependencies) => source === "action" ? dependencies.actionDependencyIds.has(id) : source === "event" ? dependencies.eventDependencyIds.has(id) : dependencies.signalDependencyIds.has(id), ) ) { return undefined } const key = `${source}:${id}` const existing = sourceIndexes.get(key) if (existing !== undefined) return existing const index = nodes.length nodes.push({ id, source, type: "source" }) sourceIndexes.set(key, index) return index } const visit = (signal: Signal): number[] => { const existing = rootsBySignal.get(signal) if (existing) return existing const definition = signal[SIGNAL_DERIVATION_DEFINITION] if (definition) { const inputs = unique(definition.inputs.flatMap(visit)) if (definition.type === "transparent" || inputs.length === 0) { rootsBySignal.set(signal, inputs) return inputs } nodes.push( definition.type === "property" ? { inputs, property: definition.property, type: definition.type, } : { inputs, type: definition.type }, ) const roots = [nodes.length - 1] rootsBySignal.set(signal, roots) return roots } const dependencies: SignalDependencySets = { actionDependencyIds: new Set(), eventDependencyIds: new Set(), signalDependencyIds: new Set(), } signal[signalDependencyCollector](state, dependencies) const roots = [ ...[...dependencies.actionDependencyIds].toSorted().flatMap((id) => { const source = addSource("action", id) return source === undefined ? [] : [source] }), ...[...dependencies.eventDependencyIds].toSorted().flatMap((id) => { const source = addSource("event", id) return source === undefined ? [] : [source] }), ...[...dependencies.signalDependencyIds].toSorted().flatMap((id) => { const source = addSource("signal", id) return source === undefined ? [] : [source] }), ] rootsBySignal.set(signal, roots) return roots } const roots = unique(signals.flatMap(visit)) return roots.length === 0 ? null : { nodes, roots } } /** * Creates a derived signal from explicit signal inputs and a pure transformer. * * Dependency collection traverses the inputs without invoking the transformer. * The transformer runs only after the runtime materializes every input. * * @param inputs - Signals consumed by the transformer. * @param transformer - Pure function that computes the derived value. */ export function transform< const TValues extends readonly [unknown, ...unknown[]], TTransformer extends (...inputs: TValues) => unknown, >( inputs: { readonly [TIndex in keyof TValues]: Signal }, transformer: Synchronous extends never ? never : TTransformer, ): Signal> /** @inheritdoc */ export function transform< TValue, TTransformer extends (input: TValue) => unknown, >( input: Signal, transformer: Synchronous extends never ? never : TTransformer, ): Signal> export function transform( inputOrInputs: Signal | ReadonlyArray, transformer: (...inputs: unknown[]) => unknown, ): Signal { return createDerivedSignal( isSignal(inputOrInputs) ? [inputOrInputs] : inputOrInputs, { type: "transform" }, transformer, ) } /** * Creates one derived signal while retaining its explanatory operation. * * @param inputs - Signals consumed by the operation. * @param operation - Explanatory property or transform step. * @param transformer - Pure function that resolves the derived value. */ function createDerivedSignal( inputs: ReadonlyArray, operation: { property: string; type: "property" } | { type: "transform" }, transformer: (...inputs: unknown[]) => unknown, ): Signal { return createSignal({ ancestorOrigins: unique( inputs.flatMap((input) => input[signalAncestorOrigins]), ), collect: (state, dependencies) => combineSignalDependencyOutcomes( inputs.map((input) => input[signalDependencyCollector](state, dependencies), ), ), collectContext: (state, dependencies) => combineSignalDependencyOutcomes( inputs.map((input) => input[signalContextDependencyCollector](state, dependencies), ), ), contextOrigins: unique( inputs.flatMap((input) => input[signalContextOrigins]), ), derivation: { ...operation, inputs }, origins: unique(inputs.flatMap((input) => input[signalOrigins])), resolve: (context) => { const outcomes = inputs.map((input) => input[signalResolver](context)) const failure = outcomes.find( (outcome): outcome is Extract => outcome.status === "failed", ) if (failure) return failure if (outcomes.some((outcome) => outcome.status === "closed")) { return { status: "closed" } } try { const sensitivity = operation.type === "property" ? projectSensitivityMask( outcomes[0]?.status === "succeeded" ? outcomes[0].sensitivity : undefined, operation.property, ) : outcomes.some( (outcome) => outcome.status === "succeeded" && outcome.sensitivity, ) ? true : undefined return { ...(sensitivity && { sensitivity }), status: "succeeded", value: transformer( ...outcomes.map((outcome) => outcome.status === "succeeded" ? outcome.value : undefined, ), ), } } catch (error) { return { failure: normalizeAutomationError(error, "expression_failed"), status: "failed", } } }, }) } /** * Normalizes an unknown thrown value for public automation lifecycle handling. * * @param error - Thrown value to normalize. * @param fallbackCode - Classification used when the error has no stable code. */ export function normalizeAutomationError( error: unknown, fallbackCode = "automation_failed", ): AutomationError { if (error instanceof FailedSignalError) return error.failure return normalizeAutomationErrorCause( error, fallbackCode, AUTOMATION_ERROR_CAUSE_DEPTH, new Set(), ) } /** * Copies an explicit keyed or global coordination mode to a value-preserving * signal result. * * @param source - Signal whose coordination mode is retained. * @param result - Value-preserving result signal. */ export function inheritSignalCoordination( source: TSignal, result: Signal, ): InheritSignalCoordination { const index = source[signalIndex] const properties = index ? { [signalIndex]: index } : source[signalGlobal] ? { [signalGlobal]: true as const } : undefined if (!properties) { /* oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The conditional type reduces to a plain signal for an unmarked source. */ return result as InheritSignalCoordination } /* oxlint-disable-next-line typescript/no-unsafe-type-assertion -- The copied marker establishes the conditional keyed or global result. */ return createSignal( { ancestorOrigins: result[signalAncestorOrigins], collect: result[signalDependencyCollector], collectContext: result[signalContextDependencyCollector], contextOrigins: result[signalContextOrigins], derivation: { inputs: [result], type: "transparent" }, origins: result[signalOrigins], resolve: result[signalResolver], }, properties, ) as InheritSignalCoordination } /** * Exposes static properties on a transparent wrapper around a signal. * * @param signal - Signal whose value and dependency behavior are retained. * @param staticProperties - Properties exposed directly on the returned proxy. */ export function withSignalStaticProperties< TValue, const TStatic extends object, >(signal: Signal, staticProperties: TStatic) { return createSignal( { ancestorOrigins: signal[signalAncestorOrigins], collect: signal[signalDependencyCollector], collectContext: signal[signalContextDependencyCollector], contextOrigins: signal[signalContextOrigins], derivation: { inputs: [signal], type: "transparent" }, origins: signal[signalOrigins], resolve: signal[signalResolver], }, staticProperties, ) } /** * Creates a proxy implementing behavior shared by every signal. * * @param behavior - Dependency, origin, and resolution behavior. * @param staticProperties - Values exposed directly instead of as projections. */ export function createSignal( behavior: SignalBehavior, staticProperties?: TStatic, ): Signal & Readonly { const staticPropertyKeys = new Set( staticProperties ? Reflect.ownKeys(staticProperties) : [], ) /** * Prepends this signal to optional additional transform inputs. * * @param inputOrInputs - Transformer or additional signal inputs. * @param transformer - Transformer for a multi-signal call. * @throws {TypeError} When a multi-signal call omits its transformer. */ function transformSignal( inputOrInputs: | Signal | readonly [Signal, ...Signal[]] | ((value: T) => unknown), transformer?: (value: T, ...others: unknown[]) => unknown, ): Signal { if (typeof inputOrInputs === "function") { return transform(signal, inputOrInputs) } if (!transformer) { throw new TypeError("A multi-signal transform requires a transformer.") } return transform( [signal, ...(isSignal(inputOrInputs) ? [inputOrInputs] : inputOrInputs)], transformer, ) } /** Protocol properties stored directly instead of projected from the value. */ const target = { [signalDependencyBinder](dependency: Signal) { return createSignal( { ancestorOrigins: unique([ ...dependency[signalAncestorOrigins], ...signal[signalAncestorOrigins], ]), collect: (state, dependencies) => { const dependencyOutcome = dependency[signalDependencyCollector]( state, dependencies, ) if (dependencyOutcome.status !== "succeeded") { return dependencyOutcome } const signalOutcome = signal[signalDependencyCollector]( state, dependencies, ) return signalOutcome.status === "pending" ? signalOutcome : { ...signalOutcome, outcomeSeq: Math.max( dependencyOutcome.outcomeSeq, signalOutcome.outcomeSeq, ), } }, collectContext: (state, dependencies) => combineSignalDependencyOutcomes([ dependency[signalContextDependencyCollector](state, dependencies), signal[signalContextDependencyCollector](state, dependencies), ]), contextOrigins: unique([ ...dependency[signalContextOrigins], ...signal[signalContextOrigins], ]), derivation: { inputs: [signal, dependency], type: "transparent", }, origins: unique([ ...dependency[signalOrigins], ...signal[signalOrigins], ]), resolve: (context) => { const dependencyOutcome = dependency[signalResolver](context) if (dependencyOutcome.status === "closed") { return { status: "closed" } } if (dependencyOutcome.status === "failed") { return dependencyOutcome } return signal[signalResolver](context) }, }, staticProperties, ) }, [signalAncestorOrigins]: behavior.ancestorOrigins ?? behavior.contextOrigins ?? behavior.origins, [signalContextDependencyCollector]: behavior.collectContext ?? behavior.collect, [signalContextOrigins]: behavior.contextOrigins ?? behavior.origins, [signalDependencyCollector]: behavior.collect, [SIGNAL_DERIVATION_DEFINITION]: behavior.derivation, [signalOrigins]: behavior.origins, [signalResolver]: behavior.resolve, [signalValueType]: () => { throw new Error("Signal type markers cannot be called.") }, globally(): GlobalSignal { return createSignal( { ancestorOrigins: signal[signalAncestorOrigins], collect: signal[signalDependencyCollector], collectContext: signal[signalContextDependencyCollector], contextOrigins: signal[signalContextOrigins], derivation: { inputs: [signal], type: "transparent" }, origins: signal[signalOrigins], resolve: signal[signalResolver], }, { [signalGlobal]: true }, ) }, keyBy(getKey: (value: T) => Encodable): KeyedSignal { return createSignal( { ancestorOrigins: signal[signalAncestorOrigins], collect: signal[signalDependencyCollector], collectContext: signal[signalContextDependencyCollector], contextOrigins: signal[signalContextOrigins], derivation: { inputs: [signal], type: "transparent" }, origins: signal[signalOrigins], resolve: signal[signalResolver], }, { [signalIndex]: { getKey }, }, ) }, toString(): never { throw new TypeError( "Signals cannot be converted directly to strings. Import `t` from `automate.ax` and use it as a tagged template to interpolate signals.", ) }, transform: transformSignal, } // oxlint-disable-next-line typescript/no-unsafe-type-assertion -- Proxy traps provide the mapped and static properties. const signal = new Proxy(target, { get(target, property, receiver) { if ( property === "globally" || property === "transform" || property === "keyBy" || property === "toString" || (typeof property === "symbol" && Object.hasOwn(target, property)) ) { return Reflect.get(target, property, receiver) } if (typeof property === "string") { const implementation = SIGNAL_METHOD_IMPLEMENTATIONS[property] if (typeof implementation === "function") { return (...arguments_: unknown[]) => Reflect.apply(implementation, undefined, [signal, ...arguments_]) } } if (staticProperties && staticPropertyKeys.has(property)) { return Reflect.get(staticProperties, property) } if (typeof property !== "string") return undefined return createDerivedSignal( [signal], { property, type: "property" }, (value) => { if (value == null) { throw new TypeError( `Cannot read property ${property} from ${String(value)}.`, ) } return Reflect.get(Object(value), property) }, ) }, }) as Signal & Readonly return signal } /** * Normalizes one error and its bounded causal chain. * * @param error - Thrown value at this level of the causal chain. * @param fallbackCode - Classification used without an explicit error code. * @param remainingDepth - Number of nested causes still exposed publicly. * @param seen - Values already visited while traversing causes. */ function normalizeAutomationErrorCause( error: unknown, fallbackCode: string, remainingDepth: number, seen: Set, ): AutomationError { const record = isRecord(error) ? error : undefined const name = error instanceof Error ? error.name : undefined const details = normalizeAutomationErrorDetails(record) const cause = record?.cause seen.add(error) return { code: name === "ZodError" ? "validation_failed" : (normalizeAutomationErrorCode(record?.code) ?? fallbackCode), message: error instanceof Error ? error.message : typeof error === "string" ? error : String(error), ...(name && { name }), ...(details !== undefined && { details }), ...(remainingDepth > 0 && cause !== undefined && !seen.has(cause) ? { cause: normalizeAutomationErrorCause( cause, fallbackCode, remainingDepth - 1, seen, ), } : {}), } } /** * Extracts explicitly structured or Zod validation diagnostics. * * @param error - Error-like record whose safe details are inspected. */ function normalizeAutomationErrorDetails( error: Record | undefined, ): AutomationError["details"] { if (!error) return if (error.name === "ZodError" && Array.isArray(error.issues)) { return { issues: error.issues.flatMap((issue) => { if (!isRecord(issue)) return [] const entry = issue return [ { code: typeof entry.code === "string" ? entry.code : "custom", message: typeof entry.message === "string" ? entry.message : "Validation failed.", path: Array.isArray(entry.path) ? entry.path.map((segment) => typeof segment === "string" || typeof segment === "number" ? segment : String(segment), ) : [], ...(typeof entry.format === "string" && { format: entry.format, }), }, ] }), } } const details = JSON_VALUE_SCHEMA.safeParse(error.details) return details.success ? details.data : undefined } /** * Converts provider-style error codes to the public lowercase format. * * @param code - Provider or JavaScript error code to normalize. */ function normalizeAutomationErrorCode(code: unknown): string | undefined { if (typeof code !== "string" || code.length === 0) return return ( code .replaceAll(/([a-z\d])([A-Z])/g, "$1_$2") .replaceAll(/[^a-zA-Z\d]+/g, "_") .replaceAll(/^_+|_+$/g, "") .toLowerCase() || undefined ) } /** * Narrows unknown values to safely inspectable object records. * * @param value - Candidate value to inspect. */ function isRecord(value: unknown): value is Record { return typeof value === "object" && value !== null }