import { inject, injectable, type NodeId, type RunEventBus, type RunEventSubscription, type TestSuiteRunId, type TriggerNodeConfig, type TypeToken, type WorkflowDefinition, type WorkflowId, } from "@codemation/core"; import { TestSuiteOrchestrator } from "@codemation/core/bootstrap"; import type { TestSuiteChildRunSummary, TestSuiteRunRecord, TestSuiteRunRepository, } from "../../domain/runs/TestSuiteRunRepository"; import { TestSuiteRunRepositoryToken, TestSuiteRunTrackerFactory } from "./TestSuiteRunTrackerFactory"; export interface TestRunnerWorkflowLookup { resolveWorkflow(workflowId: WorkflowId): WorkflowDefinition | undefined; } export const TestRunnerWorkflowLookupToken = Symbol.for( "codemation.application.testing.TestRunnerWorkflowLookup", ) as unknown as TypeToken; export const TestRunnerEventBusToken = Symbol.for( "codemation.application.testing.TestRunnerEventBus", ) as unknown as TypeToken; export interface StartTestSuiteRunResult { readonly testSuiteRunId: TestSuiteRunId; readonly status: "running"; } @injectable() export class TestRunnerService { constructor( @inject(TestSuiteOrchestrator) private readonly orchestrator: TestSuiteOrchestrator, @inject(TestRunnerWorkflowLookupToken) private readonly workflowLookup: TestRunnerWorkflowLookup, @inject(TestRunnerEventBusToken) private readonly eventBus: RunEventBus, @inject(TestSuiteRunRepositoryToken) private readonly suiteRepo: TestSuiteRunRepository, @inject(TestSuiteRunTrackerFactory) private readonly trackerFactory: TestSuiteRunTrackerFactory, ) {} async startTestSuiteRun(args: { workflowId: WorkflowId; triggerNodeId: NodeId; concurrency?: number; signal?: AbortSignal; }): Promise { const workflow = this.workflowLookup.resolveWorkflow(args.workflowId); if (!workflow) { throw new Error(`Unknown workflowId: ${args.workflowId}`); } const triggerDef = workflow.nodes.find((n) => n.id === args.triggerNodeId); if (!triggerDef || triggerDef.kind !== "trigger") { throw new Error(`Node ${args.triggerNodeId} is not a trigger`); } const triggerConfig = triggerDef.config as TriggerNodeConfig; if (triggerConfig.triggerKind !== "test") { throw new Error( `Node ${args.triggerNodeId} is not a test trigger (triggerKind="${triggerConfig.triggerKind ?? "live"}")`, ); } const startedAt = new Date().toISOString(); const placeholderId = `tsr_pending_${Date.now().toString(36)}_${Math.random().toString(36).slice(2, 8)}`; const tracker = this.trackerFactory.create({ workflow }); const subscription: RunEventSubscription = await this.eventBus.subscribeToWorkflow(args.workflowId, (event) => { void tracker.onEvent(event); }); await this.suiteRepo.create({ id: placeholderId, workflowId: workflow.id, triggerNodeId: triggerDef.id, triggerNodeName: triggerDef.name ?? triggerConfig.name, concurrency: args.concurrency ?? 4, startedAt, }); tracker.adopt(placeholderId); const finalize = async (): Promise => { try { const orchestratorResult = await this.orchestrator.runSuite({ workflow, triggerNodeId: triggerDef.id, testSuiteRunId: placeholderId, ...(args.concurrency !== undefined ? { concurrency: args.concurrency } : {}), ...(args.signal ? { signal: args.signal } : {}), }); await tracker.finalize(orchestratorResult); } catch (err) { const message = err instanceof Error ? err.message : String(err); await this.suiteRepo.update(placeholderId, { status: "errored", finishedAt: new Date().toISOString(), errorMessage: message, }); } finally { await subscription.close(); } }; void finalize(); return { testSuiteRunId: placeholderId, status: "running" }; } async getTestSuiteRun(id: TestSuiteRunId): Promise { return await this.suiteRepo.findById(id); } async listTestSuiteRuns(workflowId: WorkflowId, limit?: number): Promise> { return await this.suiteRepo.listByWorkflow({ workflowId, limit }); } async listChildRuns(testSuiteRunId: TestSuiteRunId): Promise> { return await this.suiteRepo.listChildRuns(testSuiteRunId); } }