import type { ConnectionInvocationRecord, EngineRunCounters, PendingResumeEntry, PersistedSuspensionEntry, PreparedNodeActivationDispatch, NodeActivationRequest, NodeActivationScheduler, NodeExecutionSnapshot, NodeId, ParentExecutionRef, PendingNodeExecution, PersistedRunControlState, RunDataFactory, RunExecutionOptions, RunId, RunQueueEntry, RunResult, WorkflowExecutionRepository, WorkflowId, } from "../types"; import { RunQueuePlanner } from "../planning/RunQueuePlanner"; import { NodeEventPublisher } from "../events/NodeEventPublisher"; import type { NodeActivationRequestInputPreparer } from "./NodeActivationRequestInputPreparer"; import { NodeExecutionSnapshotFactory } from "./NodeExecutionSnapshotFactory"; import { NodeInputsByPortFactory } from "./NodeInputsByPortFactory"; type PersistedRunStateRecord = NonNullable>>; type ActivationSchedulerPort = Pick; export type ActivationEnqueueRequest = { runId: RunId; workflowId: WorkflowId; startedAt: string; parent?: ParentExecutionRef; executionOptions?: RunExecutionOptions; control: PersistedRunControlState | undefined; workflowSnapshot: PersistedRunStateRecord["workflowSnapshot"]; mutableState: PersistedRunStateRecord["mutableState"]; policySnapshot: PersistedRunStateRecord["policySnapshot"]; pendingQueue: RunQueueEntry[]; request: NodeActivationRequest; previousNodeSnapshotsByNodeId: Record; planner: RunQueuePlanner; engineCounters?: EngineRunCounters; connectionInvocations?: ReadonlyArray; suspension?: ReadonlyArray; pendingResume?: PendingResumeEntry; }; export class ActivationEnqueueService { constructor( private readonly activationScheduler: ActivationSchedulerPort, private readonly workflowExecutionRepository: WorkflowExecutionRepository, private readonly nodeEventPublisher: NodeEventPublisher, private readonly nodeActivationRequestInputPreparer: NodeActivationRequestInputPreparer, ) {} async enqueueActivation(args: ActivationEnqueueRequest): Promise { const { result, queuedSnapshot } = await this.enqueueActivationWithSnapshot(args); await this.nodeEventPublisher.publish("nodeQueued", queuedSnapshot); return result; } async enqueueActivationWithSnapshot( args: ActivationEnqueueRequest, ): Promise<{ result: RunResult; queuedSnapshot: NodeExecutionSnapshot }> { const preparedRequest = await this.nodeActivationRequestInputPreparer.prepare(args.request); const preparedDispatch = await this.activationScheduler.prepareDispatch(preparedRequest); const inputsByPort = NodeInputsByPortFactory.fromRequest(preparedRequest); const itemsIn = preparedRequest.kind === "multi" ? args.planner.sumItemsByPort(preparedRequest.inputsByPort) : preparedRequest.input.length; const enqueuedAt = new Date().toISOString(); const pending: PendingNodeExecution = { runId: args.runId, activationId: args.request.activationId, workflowId: args.workflowId, nodeId: args.request.nodeId, itemsIn, inputsByPort, receiptId: preparedDispatch.receipt.receiptId, queue: preparedDispatch.receipt.queue, batchId: args.request.batchId, enqueuedAt, }; const queuedSnapshot = NodeExecutionSnapshotFactory.queued({ runId: args.runId, workflowId: args.workflowId, nodeId: args.request.nodeId, activationId: args.request.activationId, parent: args.parent, queuedAt: enqueuedAt, inputsByPort, }); await this.workflowExecutionRepository.save({ runId: args.runId, workflowId: args.workflowId, startedAt: args.startedAt, parent: args.parent, executionOptions: args.executionOptions, control: args.control, workflowSnapshot: args.workflowSnapshot, mutableState: args.mutableState, policySnapshot: args.policySnapshot, engineCounters: args.engineCounters, connectionInvocations: args.connectionInvocations ? [...args.connectionInvocations] : [], status: "pending", pending, queue: args.pendingQueue.map((entry) => ({ ...entry })), outputsByNode: (args.request.ctx.data as ReturnType).dump(), nodeSnapshotsByNodeId: { ...args.previousNodeSnapshotsByNodeId, [args.request.nodeId]: queuedSnapshot, }, ...(args.suspension !== undefined ? { suspension: args.suspension } : {}), ...(args.pendingResume !== undefined ? { pendingResume: args.pendingResume } : {}), }); await this.dispatchPreparedActivation(preparedDispatch); return { result: { runId: args.runId, workflowId: args.workflowId, startedAt: args.startedAt, status: "pending", pending }, queuedSnapshot, }; } private async dispatchPreparedActivation(preparedDispatch: PreparedNodeActivationDispatch): Promise { await preparedDispatch.dispatch(); } }