import { and, eq } from "drizzle-orm"; import { db } from "../db/client.js"; import { qscoreSnapshots, workflowArtifacts, workflowEvents, workflowRunModules, workflowRuns } from "../db/schema.js"; import { runServiceAgentProbe } from "../services/service-agents.js"; import { getWorkflowDefinition } from "./registry.js"; import { prepareOpenCodeWorkflowModule } from "./executors/opencode-executor.js"; export type ModuleRunResult = { status: "done" | "blocked" | "manual_required" | "coming_soon" | "opencode_required"; summary: string; output?: Record; }; export async function executeWorkflowModule(input: { userId: string; runId: string; moduleId: string }): Promise { const [run] = await db.select().from(workflowRuns).where(and(eq(workflowRuns.id, input.runId), eq(workflowRuns.userId, input.userId))).limit(1); if (!run) throw new Error(`workflow run not found: ${input.runId}`); const workflow = getWorkflowDefinition(run.workflowId); if (!workflow) throw new Error(`workflow definition not found: ${run.workflowId}`); const mod = workflow.modules.find((m) => m.id === input.moduleId); if (!mod) throw new Error(`workflow module not found: ${input.moduleId}`); await db.update(workflowRunModules).set({ status: "running", startedAt: new Date() }).where(and(eq(workflowRunModules.runId, input.runId), eq(workflowRunModules.moduleId, input.moduleId))); if (mod.execution === "service" && mod.service) { const result = await runServiceAgentProbe({ id: mod.id, name: mod.title, role: mod.role, kind: "microservice", description: mod.description, service: mod.service }, { userId: input.userId, goal: run.goal ?? "" }); const status = result.status === "ok" ? "done" : "blocked"; const output = result.detail as Record | undefined; await db.update(workflowRunModules).set({ status, outputSummary: result.summary, output, completedAt: new Date() }).where(and(eq(workflowRunModules.runId, input.runId), eq(workflowRunModules.moduleId, input.moduleId))); if (mod.service === "qscore-service" && output) { const rqScore = extractRqScore(output); await db.insert(qscoreSnapshots).values({ userId: input.userId, runId: input.runId, snapshotType: "module", rqScore, payload: output }); await db.update(workflowRuns).set({ qscoreAfter: output, updatedAt: new Date() }).where(eq(workflowRuns.id, input.runId)); } await db.insert(workflowEvents).values({ runId: input.runId, userId: input.userId, type: status === "done" ? "module.completed" : "module.blocked", payload: { moduleId: input.moduleId, summary: result.summary, output } }); await updateRunProgress(input.runId); return { status, summary: result.summary, output }; } if (mod.execution === "opencode") { try { const prepared = await prepareOpenCodeWorkflowModule({ userId: input.userId, runId: input.runId, workflow, module: mod, goal: run.goal ?? undefined }); const status = prepared.status === "ok" ? "done" : prepared.status === "blocked_service_unavailable" ? "blocked" : "opencode_required"; const output = { artifacts: prepared.artifacts }; if (prepared.artifacts.length) { await db.insert(workflowArtifacts).values(prepared.artifacts.map((a) => ({ runId: input.runId, moduleId: input.moduleId, type: a.type, title: a.title, repoPath: a.repoPath, metadata: a.metadata }))).onConflictDoNothing(); } await db.update(workflowRunModules).set({ status, outputSummary: prepared.summary, output, completedAt: new Date() }).where(and(eq(workflowRunModules.runId, input.runId), eq(workflowRunModules.moduleId, input.moduleId))); await db.insert(workflowEvents).values({ runId: input.runId, userId: input.userId, type: status === "done" ? "artifact.generated" : "artifact.contract_created", payload: { moduleId: input.moduleId, status, artifacts: prepared.artifacts } }); await updateRunProgress(input.runId); return { status, summary: prepared.summary, output }; } catch (err) { const summary = `OpenCode execution unavailable: ${err instanceof Error ? err.message : String(err)}`; await db.update(workflowRunModules).set({ status: "blocked", outputSummary: summary, error: summary, completedAt: new Date() }).where(and(eq(workflowRunModules.runId, input.runId), eq(workflowRunModules.moduleId, input.moduleId))); await db.insert(workflowEvents).values({ runId: input.runId, userId: input.userId, type: "module.blocked", payload: { moduleId: input.moduleId, status: "blocked", summary } }); await updateRunProgress(input.runId); return { status: "blocked", summary }; } } const status = mod.execution === "coming_soon" ? "coming_soon" : "manual_required"; const summary = mod.execution === "coming_soon" ? `${mod.title} is coming soon.` : `${mod.title} requires manual input or approval.`; await db.update(workflowRunModules).set({ status, outputSummary: summary, completedAt: new Date() }).where(and(eq(workflowRunModules.runId, input.runId), eq(workflowRunModules.moduleId, input.moduleId))); await db.insert(workflowEvents).values({ runId: input.runId, userId: input.userId, type: "module.blocked", payload: { moduleId: input.moduleId, status, summary } }); await updateRunProgress(input.runId); return { status, summary }; } export async function updateRunProgress(runId: string) { const modules = await db.select().from(workflowRunModules).where(eq(workflowRunModules.runId, runId)); if (!modules.length) return; const terminal = new Set(["done", "blocked", "manual_required", "coming_soon", "opencode_required"]); const finished = modules.filter((m) => terminal.has(m.status)).length; const done = modules.filter((m) => m.status === "done").length; const progressPercent = Math.round((finished / modules.length) * 100); const status = finished === modules.length ? (done === modules.length ? "completed" : "failed") : "running"; await db.update(workflowRuns).set({ progressPercent, status, completedAt: status === "completed" ? new Date() : null, updatedAt: new Date() }).where(eq(workflowRuns.id, runId)); if (status === "completed") { const [run] = await db.select().from(workflowRuns).where(eq(workflowRuns.id, runId)).limit(1); if (run) await db.insert(workflowEvents).values({ runId, userId: run.userId, type: "workflow.completed", payload: { progressPercent, modules: modules.map((m) => ({ moduleId: m.moduleId, status: m.status, summary: m.outputSummary })) } }); } } function extractRqScore(output: Record): number | undefined { const direct = output.rq_score; if (typeof direct === "number") return Math.round(direct); const compute = output.compute as Record | undefined; if (typeof compute?.rq_score === "number") return Math.round(compute.rq_score); return undefined; }