From edc007619b63d0d91fdad79a6d534c112b8a3981 Mon Sep 17 00:00:00 2001 From: asepharyana Date: Wed, 2 Sep 2026 19:49:15 +0700 Subject: [PATCH] =?UTF-8?q?feat(workflow):=20port=203d=20workflow=20+=20hi?= =?UTF-8?q?ve-mind=20engine=20=E2=80=94=20parse,=20execute,=20cycle,=20syn?= =?UTF-8?q?thesis,=20docs?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- TODO.md | 35 +++--- .../infrastructure/src/workflow/complexity.ts | 22 ++++ .../infrastructure/src/workflow/cycle.ts | 66 +++++++++++ .../infrastructure/src/workflow/docs.ts | 44 +++++++ .../infrastructure/src/workflow/engine.ts | 60 ++++++++++ .../infrastructure/src/workflow/hive_mind.ts | 11 +- .../infrastructure/src/workflow/index.ts | 17 +-- .../src/workflow/orchestrate.ts | 108 ++++++++++++++++++ .../infrastructure/src/workflow/script.ts | 93 +++++++++++++++ .../infrastructure/src/workflow/synthesis.ts | 91 +++++++++++++++ 10 files changed, 512 insertions(+), 35 deletions(-) create mode 100644 apps/packages/infrastructure/src/workflow/complexity.ts create mode 100644 apps/packages/infrastructure/src/workflow/cycle.ts create mode 100644 apps/packages/infrastructure/src/workflow/docs.ts create mode 100644 apps/packages/infrastructure/src/workflow/engine.ts create mode 100644 apps/packages/infrastructure/src/workflow/orchestrate.ts create mode 100644 apps/packages/infrastructure/src/workflow/script.ts create mode 100644 apps/packages/infrastructure/src/workflow/synthesis.ts diff --git a/TODO.md b/TODO.md index 003b6ef..ea4dcda 100644 --- a/TODO.md +++ b/TODO.md @@ -3,7 +3,7 @@ > Plan lengkap migrasi in-place dari Rust (11 crate, clean architecture 4 lapis) ke monorepo > TypeScript/Bun. Bahasa Indonesia, commit pakai Conventional Commits. > -> **Status:** Fase 0–2 + 3a + 3b + 3e + 3c **selesai**. Berikutnya: **3d workflow + hive-mind**. +> **Status:** Fase 0–2 + 3a + 3b + 3e + 3c + 3d **selesai**. Berikutnya: **3f auth infrastructure**. > Base: `bun run check` bersih, 54 test hijau. --- @@ -95,12 +95,8 @@ apps/ **Sumber Rust:** `apps/infrastructure/src/subagent/{engine,division,context,provider,spawn}.rs`, `subagent/auto/*` - [x] `AccessTier.toolsFor(access)` — `tools/division.ts`: Read / Write / Full yang mem-filter registry - - Read → tool baca-saja - - Write → baca + tulis - - Full → semua - [x] `run_agent` loop (`tools/engine.ts`): `MAX_ITERATIONS=25`, `TOOL_OUTPUT_MAX_CHARS=12_000`, - `MAX_CONSECUTIVE_TOOL_ERRORS=3`, `MAX_PARALLEL_TOOLS=8`; recovery note setelah 3 error beruntun; - return pesan limit iterasai saat keluar alami + `MAX_CONSECUTIVE_TOOL_ERRORS=3`, `MAX_PARALLEL_TOOLS=8`; recovery note setelah 3 error beruntun - [x] `SubagentContext` (`subagent/context.ts`) — directive, toolCtx, access, baseURL, apiKey, model - [x] `spawn_subagent` — `spawn_tools.ts` — runAgent berbasis HTTP provider + config_resolver - [x] `resolve_subagent_provider` (`subagent/provider.ts`) — pilih provider/model subagent dari settings + app_config @@ -111,21 +107,22 @@ apps/ wiring `review_enabled` (yang sekarang no-op di `Review` agent) - [x] Placeholder di `apps/packages/infrastructure/src/subagent/*` diisi penuh ---- - -## ⬜ Yang belum dikerjakan - ### Fase 3d — Workflow + hive-mind **Sumber Rust:** `apps/infrastructure/src/workflow/{script,mod,docs}.rs`, `workflow/engine/{execution,phases,primitives}.rs`, `workflow/hive_mind/{complexity,cycle,synthesis,live,types}.rs` -- [ ] `parse_workflow_script(yaml)` — parser YAML workflow (`script.ts`) -- [ ] `execute_workflow(script, ctx)` — eksekusi multi-phase (`engine/execution.ts`) -- [ ] Phase primitives: run sub-agent per node, gate/condition, paralel (bounded), I/O -- [ ] `execute_cycle(cycle, ctx)` — hive-mind cycle, concurrency 8, kumpulkan NodeOutput -- [ ] `synthesize_consensus(node_outputs, ctx)` — konsensus via LLM -- [ ] `complexity_heuristic` — pilih apakah cukup 1 cycle / butuh banyak -- [ ] `docs.ts` — tulis convergence/decision doc -- [ ] Isi body `WorkflowRun` & `HiveMind` tool (ganti placeholder), hapus `workflow/index.ts` + `hive_mind.ts` placeholder +- [x] `parse_workflow_script(yaml)` — parser YAML workflow (`workflow/script.ts`) +- [x] `execute_workflow(script, ctx)` — eksekusi multi-phase (`workflow/engine.ts` + `workflow/orchestrate.ts`) +- [x] Phase primitives: run sub-agent per node (`executePrimitive`) +- [x] `execute_cycle(cycle, ctx)` — hive-mind cycle, concurrency 8, kumpulkan NodeOutput (`workflow/cycle.ts`) +- [x] `synthesize_consensus(node_outputs, ctx)` — konsensus via LLM + fallback concat (`workflow/synthesis.ts`) +- [x] `complexity_heuristic` — pilih apakah cukup 1 cycle / butuh banyak (`workflow/complexity.ts`) +- [x] `docs.ts` — tulis convergence/decision doc (`workflow/docs.ts`) +- [x] Isi body `WorkflowRun` & `HiveMind` tool (ganti placeholder) di `tools/workflow.ts` +- [x] `orchestrate.ts` — top-level `executeWorkflow` + `executeHiveMind` entry points + +--- + +## ⬜ Yang belum dikerjakan ### Fase 3f — Auth infrastructure **Sumber Rust:** `apps/infrastructure/src/auth/{password,jwt,oauth_loopback,mod}.rs` @@ -195,4 +192,4 @@ apps/ - Registry tool = **38** (terhitung workflow_run, note_finding, read_findings, hive_mind; ExploreCodebase tidak diport). - `Tool.run` boleh return `string | Promise`. - `spawn_agents`/`spawn_pipeline`/`parallel_delegate` sudah terisi penuh via 3c. -- Workflow/hive-mind tools masih placeholder — diisi penuh di 3d. +- Workflow/hive-mind tools sudah terisi penuh via 3d. diff --git a/apps/packages/infrastructure/src/workflow/complexity.ts b/apps/packages/infrastructure/src/workflow/complexity.ts new file mode 100644 index 0000000..cb8de3f --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/complexity.ts @@ -0,0 +1,22 @@ +/** + * Complexity heuristics — determine whether a request is complex enough to + * warrant hive-mind orchestration. + * Mirrors `apps/infrastructure/src/workflow/hive_mind/complexity.rs`. + */ + +/** Heuristics to determine if a request is complex enough for hive-mind. */ +export function isComplexRequest(task: string): boolean { + const complexityIndicators = [ + "refactor", + "redesign", + "multiple files", + "architecture", + "migration", + "comprehensive", + "end-to-end", + "full-stack", + ]; + + const taskLower = task.toLowerCase(); + return complexityIndicators.some((indicator) => taskLower.includes(indicator)); +} diff --git a/apps/packages/infrastructure/src/workflow/cycle.ts b/apps/packages/infrastructure/src/workflow/cycle.ts new file mode 100644 index 0000000..c495fcc --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/cycle.ts @@ -0,0 +1,66 @@ +/** + * Hive-mind cycle execution — run one cycle of parallel nodes. + * Mirrors `apps/infrastructure/src/workflow/hive_mind/cycle.rs`. + * + * Flow: resolve LLM credentials → run all directives in the cycle concurrently + * with bounded concurrency (MAX_CONCURRENT_NODES=8) → collect NodeOutput. + * Failed nodes are logged and replaced with [ERROR] output. + */ +import type { CognitiveCycle, NodeOutput } from "@zesdex/domain"; +import type { ToolCtx } from "../tools/mod.ts"; +import { runAgent } from "../tools/engine.ts"; +import { resolveConfig } from "../subagent/config_resolver.ts"; +import { buildProviderService } from "../subagent/http_provider.ts"; + +/** Maximum number of hive-mind nodes running concurrently per cycle. */ +const MAX_CONCURRENT_NODES = 8; + +/** + * Execute one cycle: run each node directive and collect outputs. + * + * Flow: resolve LLM credentials → spawn directives with bounded concurrency → + * collect NodeOutputs; failed nodes are logged and replaced with [ERROR] + * placeholder so the cycle still completes. + */ +export async function executeCycle( + cycle: CognitiveCycle, + toolCtx: ToolCtx, +): Promise { + console.info( + `Executing cycle ${cycle.index} with ${cycle.directives.length} directives`, + ); + + const { baseUrl, apiKey, model } = await resolveConfig(); + const svc = buildProviderService(baseUrl, apiKey, model); + + const tasks = cycle.directives.map(async (dir, i) => { + const nodeId = `Node-${cycle.index}-${i}`; + const access = (dir.access_tier as "read" | "write" | "full") ?? "read"; + + try { + const output = await runAgent(svc, dir.directive, access, toolCtx, undefined); + return { + id: nodeId, + directive: dir.directive, + output, + } satisfies NodeOutput; + } catch (e) { + console.warn(`Node ${nodeId} failed (isolated): ${(e as Error).message}`); + return { + id: nodeId, + directive: dir.directive, + output: `[ERROR] ${(e as Error).message}`, + } satisfies NodeOutput; + } + }); + + // Bounded concurrency: run at most MAX_CONCURRENT_NODES futures at once. + const results: NodeOutput[] = []; + for (let i = 0; i < tasks.length; i += MAX_CONCURRENT_NODES) { + const batch = tasks.slice(i, i + MAX_CONCURRENT_NODES); + const batchResults = await Promise.all(batch); + results.push(...batchResults); + } + + return results; +} diff --git a/apps/packages/infrastructure/src/workflow/docs.ts b/apps/packages/infrastructure/src/workflow/docs.ts new file mode 100644 index 0000000..cacc085 --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/docs.ts @@ -0,0 +1,44 @@ +/** + * Hive-mind convergence documentation — writes deterministic audit trail. + * Mirrors `apps/infrastructure/src/workflow/docs.rs`. + * + * Creates a markdown file documenting node outputs and consensus synthesis. + */ +import * as fs from "node:fs"; +import * as path from "node:path"; +import type { NodeOutput } from "@zesdex/domain"; + +/** + * Write a deterministic audit trail for a hive-mind convergence. + * + * Flow: create docs/runs/ dir → build markdown content → write file. + * This is deterministic (not an LLM step) and never skippable. + */ +export function writeHiveMindConvergence( + runDir: string, + nodes: NodeOutput[], + consensus: string, +): string { + const docsDir = path.join(runDir, "docs", "runs"); + fs.mkdirSync(docsDir, { recursive: true }); + + const timestamp = new Date().toISOString().replace(/[:.]/g, "").slice(0, 15); + const filename = `${timestamp}-hive-mind-convergence.md`; + const filepath = path.join(docsDir, filename); + + let content = `# Hive Mind Convergence — ${timestamp}\n\n`; + content += "## Node Outputs\n\n"; + + for (const node of nodes) { + content += `### ${node.id} — ${node.directive}\n\n`; + content += `${node.output}\n\n`; + } + + content += "## Consensus\n\n"; + content += consensus; + content += "\n"; + + fs.writeFileSync(filepath, content); + console.info(`Hive mind convergence written to ${filepath}`); + return filepath; +} diff --git a/apps/packages/infrastructure/src/workflow/engine.ts b/apps/packages/infrastructure/src/workflow/engine.ts new file mode 100644 index 0000000..b818d0f --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/engine.ts @@ -0,0 +1,60 @@ +/** + * Workflow engine — runs a parsed workflow script phase by phase. + * Mirrors `apps/infrastructure/src/workflow/engine/execution.rs` and `primitives.rs`. + * + * Flow: parse YAML → build WorkflowScript → for each phase → run as subagent → collect results. + */ +import type { WorkflowScript } from "@zesdex/domain"; +import type { ToolCtx } from "../tools/mod.ts"; +import { runAgent } from "../tools/engine.ts"; + +/** Maximum characters of node output to feed into synthesis prompt per node. */ +export const MAX_NODE_OUTPUT_CHARS = 4000; + +/** + * Execute a single directive by spawning a subagent. + * Mirrors `engine/primitives.rs` `execute_primitive`. + * + * Flow: load settings → resolve LLM credentials → run_agent → return output. + */ +export async function executePrimitive(directive: string, toolCtx: ToolCtx): Promise { + const { resolveConfig } = await import("../subagent/config_resolver.ts"); + const { buildProviderService } = await import("../subagent/http_provider.ts"); + + const { baseUrl, apiKey, model } = await resolveConfig(); + const svc = buildProviderService(baseUrl, apiKey, model); + + return runAgent(svc, directive, "full", toolCtx); +} + +/** + * Execute each phase of a workflow script sequentially. + * Mirrors `engine/execution.rs` `execute_workflow`. + * + * Flow: for each phase → execute_primitive → collect result. + */ +export async function executeWorkflowScript( + script: WorkflowScript, + toolCtx: ToolCtx, +): Promise { + console.info(`Executing workflow: ${script.name} (${script.phases.length} phases)`); + const results: string[] = []; + + for (const phase of script.phases) { + console.info(`Executing phase: ${phase.name}`); + const result = await executePrimitive(phase.directive, toolCtx); + results.push(result); + } + + return results; +} + +/** Format a workflow execution result as a summary string. */ +export function formatWorkflowResult(script: WorkflowScript, results: string[]): string { + let out = `## Workflow: ${script.name}\n\n`; + results.forEach((result, i) => { + const phase = script.phases[i]!; + out += `### Phase ${i}: ${phase.name}\n\n${result}\n\n`; + }); + return out; +} diff --git a/apps/packages/infrastructure/src/workflow/hive_mind.ts b/apps/packages/infrastructure/src/workflow/hive_mind.ts index 3aaf4a4..91714ab 100644 --- a/apps/packages/infrastructure/src/workflow/hive_mind.ts +++ b/apps/packages/infrastructure/src/workflow/hive_mind.ts @@ -1,10 +1,5 @@ /** - * Hive-mind convergence engine — ported in sub-phase 3d. + * Hive-mind orchestration — re-exports executeHiveMind entry point. + * Mirrors `apps/infrastructure/src/workflow/hive_mind/mod.rs`. */ -import type { JsonValue } from "@zesdex/domain"; -import type { ToolCtx } from "../tools/mod.ts"; - -/** Execute a hive-mind convergence from tool arguments. */ -export async function executeHiveMind(_args: JsonValue, _ctx: ToolCtx): Promise { - throw new Error("hive-mind engine not yet wired in this build"); -} +export { executeHiveMind } from "./orchestrate.ts"; diff --git a/apps/packages/infrastructure/src/workflow/index.ts b/apps/packages/infrastructure/src/workflow/index.ts index f87d960..b20f04e 100644 --- a/apps/packages/infrastructure/src/workflow/index.ts +++ b/apps/packages/infrastructure/src/workflow/index.ts @@ -1,10 +1,11 @@ /** - * Workflow execution engine — ported in sub-phase 3d. This module provides - * the execution entry points used by the workflow tools. + * Workflow + Hive-mind re-exports. + * Mirrors `apps/infrastructure/src/workflow/mod.rs`. */ -import type { ToolCtx } from "../tools/mod.ts"; - -/** Execute a YAML workflow definition. */ -export async function executeWorkflow(_yaml: string, _ctx: ToolCtx): Promise { - throw new Error("workflow engine not yet wired in this build"); -} +export { parseWorkflowScript } from "./script.ts"; +export { executeWorkflowScript, executePrimitive, formatWorkflowResult, MAX_NODE_OUTPUT_CHARS } from "./engine.ts"; +export { executeWorkflow } from "./orchestrate.ts"; +export { executeCycle } from "./cycle.ts"; +export { synthesizeConsensus } from "./synthesis.ts"; +export { isComplexRequest } from "./complexity.ts"; +export { writeHiveMindConvergence } from "./docs.ts"; diff --git a/apps/packages/infrastructure/src/workflow/orchestrate.ts b/apps/packages/infrastructure/src/workflow/orchestrate.ts new file mode 100644 index 0000000..9142f97 --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/orchestrate.ts @@ -0,0 +1,108 @@ +/** + * Workflow + Hive-mind orchestration — entry points for the WorkflowRun and HiveMind tools. + * Mirrors the top-level dispatch in `workflow/mod.rs` and `workflow/hive_mind/mod.rs`. + * + * Flow: parse cycles from args → execute each cycle → synthesize consensus + * → write convergence doc → return formatted result. + */ +import type { JsonValue, CognitiveCycle, NodeOutput, NodeDirective } from "@zesdex/domain"; +import type { ToolCtx } from "../tools/mod.ts"; +import { parseWorkflowScript } from "./script.ts"; +import { executeWorkflowScript, formatWorkflowResult } from "./engine.ts"; +import { executeCycle } from "./cycle.ts"; +import { synthesizeConsensus } from "./synthesis.ts"; +import { writeHiveMindConvergence } from "./docs.ts"; + +/** + * Execute a YAML workflow definition. + * Entry point for the WorkflowRun tool. + */ +export async function executeWorkflow(yaml: string, ctx: ToolCtx): Promise { + const script = parseWorkflowScript(yaml); + if (script.phases.length === 0) { + return "Workflow has no phases to execute."; + } + const results = await executeWorkflowScript(script, ctx); + return formatWorkflowResult(script, results); +} + +/** + * Execute a hive-mind convergence from tool arguments. + * Entry point for the HiveMind tool. + * + * Flow: parse cycles from args → execute each cycle → synthesize consensus + * → write convergence doc → return formatted result. + */ +export async function executeHiveMind(args: JsonValue, ctx: ToolCtx): Promise { + const cyclesData = extractCycles(args); + + if (cyclesData.length === 0) { + return "Hive-mind has no cycles to execute."; + } + + const allNodeOutputs: NodeOutput[] = []; + const cycleResults: string[] = []; + + for (let i = 0; i < cyclesData.length; i++) { + const cycle = cyclesData[i]!; + console.info(`Executing hive-mind cycle ${i} with ${cycle.directives.length} directives`); + + const nodeOutputs = await executeCycle(cycle, ctx); + allNodeOutputs.push(...nodeOutputs); + + const cycleSummary = nodeOutputs + .map((n) => `- ${n.id}: ${n.output.slice(0, 200)}${n.output.length > 200 ? "..." : ""}`) + .join("\n"); + cycleResults.push(`### Cycle ${i}\n\n${cycleSummary}`); + } + + // Synthesize consensus from all node outputs + let consensus: string; + try { + consensus = await synthesizeConsensus(allNodeOutputs, ctx); + } catch (e) { + consensus = `Consensus synthesis failed: ${(e as Error).message}\n\n` + + allNodeOutputs.map((n) => `### ${n.id}\n${n.output}`).join("\n\n"); + } + + // Write convergence doc + const runDir = ctx.sessionDir || process.cwd(); + writeHiveMindConvergence(runDir, allNodeOutputs, consensus); + + // Build output + let out = `## Hive Mind Convergence\n\n`; + out += `**Cycles:** ${cyclesData.length}\n`; + out += `**Nodes:** ${allNodeOutputs.length}\n\n`; + out += cycleResults.join("\n\n"); + out += `\n\n---\n\n${consensus}`; + + return out; +} + +/** Extract CognitiveCycle objects from tool args. */ +function extractCycles(args: JsonValue): CognitiveCycle[] { + if (!args || typeof args !== "object" || Array.isArray(args)) return []; + const obj = args as Record; + const cyclesArr = obj["cycles"]; + + if (!Array.isArray(cyclesArr)) return []; + + return cyclesArr.map((c, i) => { + if (!c || typeof c !== "object" || Array.isArray(c)) { + return { index: i, directives: [] }; + } + const cycle = c as Record; + const directivesArr = Array.isArray(cycle["directives"]) ? cycle["directives"]! : []; + + const directives: NodeDirective[] = directivesArr + .filter((d): d is Record => + typeof d === "object" && d !== null && !Array.isArray(d)) + .map((d) => ({ + directive: typeof d["directive"] === "string" ? d["directive"] : "", + access_tier: typeof d["access"] === "string" ? d["access"] : "read", + })) + .filter((d) => d.directive !== ""); + + return { index: i, directives }; + }); +} diff --git a/apps/packages/infrastructure/src/workflow/script.ts b/apps/packages/infrastructure/src/workflow/script.ts new file mode 100644 index 0000000..a7ceed2 --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/script.ts @@ -0,0 +1,93 @@ +/** + * Workflow script — parse and execute user-defined workflow scripts. + * Mirrors `apps/infrastructure/src/workflow/script.rs`. + * + * Parses a simple YAML-like workflow definition into a WorkflowScript + * (name + ordered phases). Supports a minimal YAML subset: `name:` and + * `phases:` with `- name:` / `- directive:` entries. + */ +import type { WorkflowScript, WorkflowPhase } from "@zesdex/domain"; + +/** + * Parse a YAML string into a WorkflowScript. + * + * Expected format: + * ```yaml + * name: my-workflow + * phases: + * - name: research + * directive: "Explore the codebase..." + * - name: implement + * directive: "Implement the changes..." + * ``` + * + * Uses a lightweight regex-based parser to avoid adding a YAML library dep. + * Handles flat key-value pairs and a single level of list-of-maps. + */ +export function parseWorkflowScript(yaml: string): WorkflowScript { + const lines = yaml.split("\n"); + let name = "unnamed"; + const phases: WorkflowPhase[] = []; + + let inPhases = false; + let currentPhase: Partial | null = null; + + for (const raw of lines) { + const trimmed = raw.trim(); + if (trimmed === "" || trimmed.startsWith("#")) continue; + + // Top-level key: value + const topMatch = trimmed.match(/^(\w+):\s*(.*)$/); + if (topMatch && !raw.startsWith(" ") && !raw.startsWith("-")) { + const key = topMatch[1]!; + const val = topMatch[2]!.trim(); + if (key === "name" && val !== "") { + name = val; + } else if (key === "phases") { + inPhases = true; + } + continue; + } + + if (!inPhases) continue; + + // List item: - name: foo or - directive: bar + const listMatch = trimmed.match(/^-\s+(\w+):\s*(.*)$/); + if (listMatch) { + // If there was a previous phase, push it + if (currentPhase && (currentPhase.name || currentPhase.directive)) { + phases.push({ + name: currentPhase.name ?? "phase", + directive: currentPhase.directive ?? "", + }); + } + const key = listMatch[1]!; + const val = listMatch[2]!.trim(); + currentPhase = { name: undefined, directive: undefined }; + if (key === "name") currentPhase.name = val; + if (key === "directive") currentPhase.directive = val; + continue; + } + + // Continuation of a phase field (indented under a list item) + if (currentPhase && !trimmed.startsWith("-")) { + const kvMatch = trimmed.match(/^(\w+):\s*(.*)$/); + if (kvMatch) { + const key = kvMatch[1]!; + const val = kvMatch[2]!.trim(); + if (key === "name") currentPhase.name = val; + if (key === "directive") currentPhase.directive = val; + } + } + } + + // Push the last phase + if (currentPhase && (currentPhase.name || currentPhase.directive)) { + phases.push({ + name: currentPhase.name ?? "phase", + directive: currentPhase.directive ?? "", + }); + } + + return { name, phases }; +} diff --git a/apps/packages/infrastructure/src/workflow/synthesis.ts b/apps/packages/infrastructure/src/workflow/synthesis.ts new file mode 100644 index 0000000..b18f762 --- /dev/null +++ b/apps/packages/infrastructure/src/workflow/synthesis.ts @@ -0,0 +1,91 @@ +/** + * Consensus synthesis — reconciles multiple node outputs into one assessment. + * Mirrors `apps/infrastructure/src/workflow/hive_mind/synthesis.rs`. + * + * Flow: build combined prompt from node outputs → ask the model to distill + * into a single consensus (conflicts, agreements, key findings) → return the + * synthesized text. Falls back to plain concatenation if the LLM call fails. + */ +import type { NodeOutput } from "@zesdex/domain"; +import type { ToolCtx } from "../tools/mod.ts"; +import { MAX_NODE_OUTPUT_CHARS } from "./engine.ts"; + +/** + * Synthesize a consensus from all node outputs using the LLM. + * + * Flow: combine node outputs → ask the model to reconcile them into a single + * consensus → return the synthesized text. Falls back to a plain + * concatenation summary if the LLM is unreachable or the call fails. + */ +export async function synthesizeConsensus( + nodes: NodeOutput[], + _toolCtx: ToolCtx, +): Promise { + if (nodes.length === 0) return "No node outputs to synthesize."; + + const combined = buildCombinedBody(nodes); + + const { resolveConfig } = await import("../subagent/config_resolver.ts"); + const { buildProviderService } = await import("../subagent/http_provider.ts"); + + const { baseUrl, apiKey, model } = await resolveConfig(); + const svc = buildProviderService(baseUrl, apiKey, model); + + const systemMsg = { + role: "system" as const, + content: "You are a consensus synthesizer for a multi-agent hive mind. " + + "Several independent nodes analysed a problem and produced the outputs " + + "below. Distill them into ONE coherent consensus report with these sections:\n" + + "- AGREEMENTS: points multiple nodes converge on.\n" + + "- CONFLICTS: contradictory conclusions, with which node(s) support each side.\n" + + "- KEY FINDINGS: the most important, actionable takeaways.\n" + + "- RECOMMENDATION: a single recommended next action, or 'no clear consensus' " + + "if the outputs are too divergent.\n" + + "Be concise and factual. If a node errored, note it and ignore its content.\n" + + "Do not invent facts not present in the node outputs.", + }; + + const userMsg = { + role: "user" as const, + content: `Consolidate these ${nodes.length} node outputs into a single consensus:\n\n${combined}`, + }; + + try { + const result = await svc.chat([systemMsg, userMsg], undefined, 1024, 0.3); + const text = result.message.content ?? ""; + + if (text.trim() === "") { + console.warn("consensus LLM returned empty output; falling back to concat summary"); + return concatSummary(nodes); + } + + return `# Consensus Synthesis\n\nNodes synthesized: ${nodes.length}\n\n${text}`; + } catch (e) { + console.warn(`consensus LLM call failed; falling back to concat summary: ${(e as Error).message}`); + return concatSummary(nodes); + } +} + +/** Build the concatenated node-output body for the prompt. */ +function buildCombinedBody(nodes: NodeOutput[]): string { + let combined = ""; + for (const node of nodes) { + const output = truncateChars(node.output, MAX_NODE_OUTPUT_CHARS); + combined += `\n## ${node.id} — ${node.directive}\n${output}\n`; + } + return combined; +} + +/** Fallback: a plain concatenation summary. */ +function concatSummary(nodes: NodeOutput[]): string { + const combined = buildCombinedBody(nodes); + return `# Consensus Synthesis\n\nNodes synthesized: ${nodes.length}\n\n## Summary\n\n` + + `The following node outputs were collected:\n\n${combined}\n\n` + + `## Key Findings\n\nReview the individual node outputs above for detailed findings.`; +} + +/** Truncate a string to maxLen characters (char-safe). */ +function truncateChars(s: string, maxLen: number): string { + if (s.length <= maxLen) return s; + return s.slice(0, maxLen); +}