import { CredentialResolverFactory } from "@codemation/core/bootstrap"; import { CoreTokens, inject, injectable, PollingTriggerDedupWindow } from "@codemation/core"; import type { EngineExecutionLimitsPolicy } from "@codemation/core/bootstrap"; import type { ActivationIdFactory, ExecutionContextFactory, Items, NodeResolver, PersistedTriggerSetupState, PollingTriggerHandle, RunDataFactory, RunIdFactory, TriggerInstanceId, TriggerInvoker, TriggerNode, TriggerNodeConfig, TriggerSetupContext, TriggerSetupStateRepository, } from "@codemation/core"; @injectable() export class InProcessTriggerInvoker implements TriggerInvoker { constructor( @inject(CoreTokens.NodeResolver) private readonly nodeResolver: NodeResolver, @inject(CredentialResolverFactory) private readonly credentialResolverFactory: CredentialResolverFactory, @inject(CoreTokens.TriggerSetupStateRepository) private readonly triggerSetupStateRepository: TriggerSetupStateRepository, @inject(CoreTokens.RunIdFactory) private readonly runIdFactory: RunIdFactory, @inject(CoreTokens.ActivationIdFactory) private readonly activationIdFactory: ActivationIdFactory, @inject(CoreTokens.RunDataFactory) private readonly runDataFactory: RunDataFactory, @inject(CoreTokens.ExecutionContextFactory) private readonly executionContextFactory: ExecutionContextFactory, @inject(CoreTokens.EngineExecutionLimitsPolicy) private readonly executionLimitsPolicy: EngineExecutionLimitsPolicy, @inject(PollingTriggerDedupWindow) private readonly dedupWindow: PollingTriggerDedupWindow, ) {} async invoke(trigger: TriggerInstanceId, config: TriggerNodeConfig, _lastRanAt: string | undefined): Promise { const node = this.nodeResolver.resolve(config.type) as TriggerNode; const runId = this.runIdFactory.makeRunId(); const data = this.runDataFactory.create(); const rootLimits = this.executionLimitsPolicy.createRootExecutionOptions(); const baseCtx = this.executionContextFactory.create({ runId, workflowId: trigger.workflowId, parent: undefined, subworkflowDepth: rootLimits.subworkflowDepth ?? 0, engineMaxNodeActivations: rootLimits.maxNodeActivations!, engineMaxSubworkflowDepth: rootLimits.maxSubworkflowDepth!, data, getCredential: this.credentialResolverFactory.create(trigger.workflowId, trigger.nodeId, config), }); const ctx = { ...baseCtx, binary: baseCtx.binary.forNode({ nodeId: trigger.nodeId, activationId: this.activationIdFactory.makeActivationId(), }), }; const stateRecord: PersistedTriggerSetupState | undefined = await this.triggerSetupStateRepository.load(trigger); const collectedItems: Items[] = []; const emit = async (items: Items): Promise => { if (items.length > 0) { collectedItems.push(items); } }; const stateRepository = this.triggerSetupStateRepository; const dedupWindow = this.dedupWindow; const pollingHandle: PollingTriggerHandle = { dedup: dedupWindow, start: async (args) => { const previousState = stateRecord !== undefined ? stateRecord.state : args.seedState; const controller = new AbortController(); const { items, nextState } = await args.runCycle({ previousState: previousState as never, signal: controller.signal, }); await stateRepository.save({ trigger, updatedAt: new Date().toISOString(), state: nextState as never, }); if ((items as Items).length > 0) { await emit(items as Items); } return nextState; }, }; await node.setup({ ...ctx, trigger, config, previousState: stateRecord?.state as never, registerCleanup: () => {}, emit, polling: pollingHandle, } satisfies TriggerSetupContext); return collectedItems.flat() as Items; } }