diff --git a/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt b/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt index 4f1dcd66..159ac1da 100644 --- a/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt +++ b/apps/server/src/main/kotlin/com/correx/apps/server/Main.kt @@ -195,12 +195,26 @@ fun main() { ?.let { Path.of(it) } ?: toolsConfig.workspaceRoot.takeIf { it.isNotEmpty() }?.let { Path.of(it) } ?: workingDir + // One shared HTTP client backs both the default and per-workspace registries' research tools + // (web_search/web_fetch). Built only when research is enabled, so the static path stays offline. + val researchToolConfig = com.correx.infrastructure.tools.ResearchToolConfig( + enabled = toolsConfig.research.enabled, + searxngUrl = toolsConfig.research.searxngUrl, + maxResults = toolsConfig.research.maxResults, + maxFetchBytes = toolsConfig.research.maxFetchBytes, + httpClient = if (toolsConfig.research.enabled) { + io.ktor.client.HttpClient(io.ktor.client.engine.cio.CIO) + } else { + null + }, + ) val toolRegistry = InfrastructureModule.createToolRegistry( buildToolConfig( workspaceRoot, workingDir, shellAllowedExecutables, toolsConfig, + researchToolConfig, ), ) val toolExecutor = InfrastructureModule.createToolExecutor( @@ -228,7 +242,7 @@ fun main() { val wsToolRegistryProvider = WorkspaceToolRegistryProvider { workspace -> val wsRegistry = InfrastructureModule.createToolRegistry( - buildToolConfigForWorkspace(workspace, shellAllowedExecutables, toolsConfig), + buildToolConfigForWorkspace(workspace, shellAllowedExecutables, toolsConfig, researchToolConfig), ) val wsExecutor = DispatchingToolExecutor(wsRegistry) WorkspaceTools(registry = wsRegistry, executor = wsExecutor) @@ -626,6 +640,7 @@ private fun buildToolConfig( workingDir: Path, shellAllowedExecutables: Set, toolsConfig: com.correx.core.config.ToolsConfig, + research: com.correx.infrastructure.tools.ResearchToolConfig, ): ToolConfig { val allowed = setOf(workspaceRoot, workingDir) return ToolConfig( @@ -648,6 +663,7 @@ private fun buildToolConfig( allowedPaths = allowed, workingDir = workingDir, ), + research = research, ) } @@ -655,7 +671,9 @@ private fun buildToolConfigForWorkspace( workspace: WorkspaceContext, shellAllowedExecutables: Set, toolsConfig: com.correx.core.config.ToolsConfig, + research: com.correx.infrastructure.tools.ResearchToolConfig, ): ToolConfig = ToolConfig( + research = research, shell = ShellConfig( enabled = toolsConfig.shellEnabled, allowedExecutables = shellAllowedExecutables, diff --git a/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt b/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt index 57b787d4..36f2b725 100644 --- a/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt +++ b/core/config/src/main/kotlin/com/correx/core/config/CorrexConfig.kt @@ -125,6 +125,7 @@ data class ToolsConfig( val networkAllowedHosts: List = emptyList(), val networkDeniedHosts: List = emptyList(), val allowedWorkspaceRoots: List = emptyList(), + val research: ResearchConfig = ResearchConfig(), ) { companion object { val DEFAULT_PRIVILEGED_LOCATIONS: List = listOf( @@ -138,6 +139,19 @@ data class ToolsConfig( } } +/** + * Deep-research workflow (research-workflow-spec). Off by default: enabling it registers the + * network tools web_search (T1, hits only [searxngUrl]) and web_fetch (T2, operator-approved). + * [searxngUrl] is the self-hosted SearXNG endpoint; [maxFetchBytes] caps a single fetched page. + */ +@Serializable +data class ResearchConfig( + val enabled: Boolean = false, + val searxngUrl: String = "http://localhost:8888", + val maxResults: Int = 8, + val maxFetchBytes: Long = 10_000_000, +) + @Serializable data class ProviderConfig( val id: String, diff --git a/docs/schemas/research_plan.json b/docs/schemas/research_plan.json new file mode 100644 index 00000000..0a609544 --- /dev/null +++ b/docs/schemas/research_plan.json @@ -0,0 +1,23 @@ +{ + "type": "object", + "properties": { + "question": { + "type": "string", + "description": "the research question, restated in your own words" + }, + "sub_questions": { + "type": "array", + "description": "the question decomposed into concrete, independently-answerable sub-questions", + "minItems": 1, + "items": { "type": "string" } + }, + "search_queries": { + "type": "array", + "description": "concrete web search queries to run, one or more per sub-question", + "minItems": 1, + "items": { "type": "string" } + } + }, + "required": ["question", "sub_questions", "search_queries"], + "additionalProperties": true +} diff --git a/docs/schemas/research_report.json b/docs/schemas/research_report.json new file mode 100644 index 00000000..c25ba17d --- /dev/null +++ b/docs/schemas/research_report.json @@ -0,0 +1,23 @@ +{ + "type": "object", + "properties": { + "summary": { + "type": "string", + "description": "the executive answer to the research question, synthesized across all sources" + }, + "findings": { + "type": "array", + "description": "the key findings, one per item; each should be supported by at least one cited source", + "minItems": 1, + "items": { "type": "string" } + }, + "sources": { + "type": "array", + "description": "the source URLs cited in the report (citations into the source dossier)", + "minItems": 1, + "items": { "type": "string" } + } + }, + "required": ["summary", "findings", "sources"], + "additionalProperties": true +} diff --git a/docs/schemas/source_dossier.json b/docs/schemas/source_dossier.json new file mode 100644 index 00000000..c5e8ee7f --- /dev/null +++ b/docs/schemas/source_dossier.json @@ -0,0 +1,29 @@ +{ + "type": "object", + "properties": { + "sources": { + "type": "array", + "description": "one entry per fetched source — each summarized on its own (per-source synthesis). This is the compression step: raw page text must never be copied here, only what the source contributes.", + "minItems": 1, + "items": { + "type": "object", + "properties": { + "url": { "type": "string", "description": "the fetched source URL (a citation handle)" }, + "title": { "type": "string" }, + "summary": { + "type": "string", + "description": "what this source contributes to the question, in your own words — not a copy of the page" + }, + "relevance": { + "type": "string", + "description": "which sub-question(s) this source bears on, and how strongly" + } + }, + "required": ["url", "summary"], + "additionalProperties": true + } + } + }, + "required": ["sources"], + "additionalProperties": true +} diff --git a/examples/workflows/prompts/research_decompose.md b/examples/workflows/prompts/research_decompose.md new file mode 100644 index 00000000..7f31bc2e --- /dev/null +++ b/examples/workflows/prompts/research_decompose.md @@ -0,0 +1,20 @@ +You are the **Research Planner** — the first stage of a deep-research workflow. + +Your job is to turn one research question into a concrete plan to answer it. You do not search +or browse yet; you decompose. + +Steps: +1. Restate the research question in your own words so intent is unambiguous. +2. Break it into a small set of sub-questions that can each be answered independently. Cover the + whole question; avoid overlap. Prefer 3–6 sub-questions over a long list. +3. For each sub-question, write one or more concrete web search queries — the actual strings you + would type into a search engine, specific enough to surface authoritative sources. + +The decision history above (steering, approvals) is ground truth — honour it. + +Emit your result as the `research_plan` artifact (JSON, schema provided): +- `question`: the research question, restated. +- `sub_questions`: the decomposed sub-questions, one per item. +- `search_queries`: concrete search queries to run, one per item. + +Do not answer the question here. Plan only. diff --git a/examples/workflows/prompts/research_gather.md b/examples/workflows/prompts/research_gather.md new file mode 100644 index 00000000..8b2b8611 --- /dev/null +++ b/examples/workflows/prompts/research_gather.md @@ -0,0 +1,21 @@ +You are the **Research Gatherer** — you run the plan and build a dossier of summarized sources. + +You have two tools: +- `web_search(query)` — searches the local SearXNG instance; returns result titles, URLs, snippets. +- `web_fetch(url)` — fetches a URL and returns its main content as clean markdown. Fetches are + approved by the operator before they run, so choose sources deliberately — quality over quantity. + +Steps: +1. For each search query in the `research_plan` above, call `web_search`. +2. From the results, pick the most promising, authoritative sources. Skip duplicates, SEO spam, + and pages unlikely to contain primary information. +3. `web_fetch` each chosen source. If a fetch comes back empty or clearly low-quality (a paywall + or JS-only page), drop it and move on — do not retry it. +4. **Summarize each fetched source on its own**, in your own words: what it contributes to the + question and which sub-question(s) it bears on. This is the most important step — the next + stage sees only your summaries, never the raw pages, so a faithful summary is the whole product. + +Never copy raw page text into a summary. Cite each source by its exact URL. + +Emit your result as the `source_dossier` artifact (JSON, schema provided): +- `sources`: one entry per fetched source, each with `url`, `title`, `summary`, and `relevance`. diff --git a/examples/workflows/prompts/research_report.md b/examples/workflows/prompts/research_report.md new file mode 100644 index 00000000..bc26e852 --- /dev/null +++ b/examples/workflows/prompts/research_report.md @@ -0,0 +1,19 @@ +You are the **Research Synthesizer** — you write the final report from the source dossier. + +You see only the `source_dossier` above: per-source summaries with their URLs. You do not have the +raw pages, and you do not search or fetch. Work from the summaries. + +Steps: +1. Read across all source summaries and find where they agree, disagree, or leave gaps. +2. Answer the research question directly and concisely in the `summary`. +3. State the key findings, each supported by at least one source. Where sources conflict, say so + rather than picking silently. +4. Cite the sources you actually relied on by their exact URLs. + +Do not invent facts that no source supports. If the dossier is insufficient to answer the +question, say what is missing in the `summary`. + +Emit your result as the `research_report` artifact (JSON, schema provided): +- `summary`: the executive answer to the question. +- `findings`: the key findings, one per item, each grounded in a cited source. +- `sources`: the source URLs you cited. diff --git a/examples/workflows/research.toml b/examples/workflows/research.toml new file mode 100644 index 00000000..d6d9a080 --- /dev/null +++ b/examples/workflows/research.toml @@ -0,0 +1,79 @@ +# Deep research workflow: decompose -> gather (search + fetch) -> report +# +# A native Correx workflow (research-workflow-spec): one graph, two network tools, no new +# orchestration machinery. The fan-out the spec describes (search per sub-question, fetch per +# source) happens *inside* the gather stage as repeated tool calls — Correx runs stages linearly +# (no parallel agents), so the model drives the fan-out itself. +# +# Per-source synthesis is the compression mechanism, not a nicety: the gather stage summarizes +# each fetched page into the source_dossier and the raw page text is discarded. The report stage +# consumes only those summaries, so research survives the token budget. +# +# Requires (in ~/.config/correx/config.toml): +# [tools.research] +# enabled = true +# searxng_url = "http://localhost:8888" # your self-hosted SearXNG (format=json enabled) +# [[artifacts]] +# id = "research_plan"; schema_path = "schemas/research_plan.json"; llm_emitted = true +# [[artifacts]] +# id = "source_dossier"; schema_path = "schemas/source_dossier.json"; llm_emitted = true +# [[artifacts]] +# id = "research_report"; schema_path = "schemas/research_report.json"; llm_emitted = true +# +# web_fetch is T2: the operator approves fetches before they leave the machine. web_search is T1 +# (it only ever touches the configured SearXNG instance). + +id = "research" +start = "decompose" +description = "Deep research: decompose a question, search + fetch sources, synthesize a cited report." + +# 1. Question -> sub-questions + search queries. Reuses the planning discipline; read-only. +[[stages]] +id = "decompose" +prompt = "prompts/research_decompose.md" +produces = [{ name = "research_plan", kind = "research_plan" }] +token_budget = 8192 +max_retries = 2 + +# 2. Run the searches, fetch promising sources, and summarize each one (per-source synthesis). +# web_search (T1) hits only SearXNG; web_fetch (T2) is operator-approved per source. +[[stages]] +id = "gather" +prompt = "prompts/research_gather.md" +needs = ["research_plan"] +produces = [{ name = "source_dossier", kind = "source_dossier" }] +allowed_tools = ["web_search", "web_fetch"] +token_budget = 32768 +max_retries = 2 + +# 3. Cross-source synthesis into a cited report. Consumes the dossier summaries only — never raw pages. +[[stages]] +id = "report" +prompt = "prompts/research_report.md" +needs = ["source_dossier"] +produces = [{ name = "research_report", kind = "research_report" }] +token_budget = 16384 +max_retries = 2 + +# --- forward edges --- + +[[transitions]] +id = "decompose-to-gather" +from = "decompose" +to = "gather" +condition_type = "artifact_validated" +condition_artifact_id = "research_plan" + +[[transitions]] +id = "gather-to-report" +from = "gather" +to = "report" +condition_type = "artifact_validated" +condition_artifact_id = "source_dossier" + +[[transitions]] +id = "report-done" +from = "report" +to = "done" +condition_type = "artifact_validated" +condition_artifact_id = "research_report" diff --git a/infrastructure/tools/src/main/kotlin/com/correx/infrastructure/tools/ToolConfig.kt b/infrastructure/tools/src/main/kotlin/com/correx/infrastructure/tools/ToolConfig.kt index 7ac06aa2..cab7c3f4 100644 --- a/infrastructure/tools/src/main/kotlin/com/correx/infrastructure/tools/ToolConfig.kt +++ b/infrastructure/tools/src/main/kotlin/com/correx/infrastructure/tools/ToolConfig.kt @@ -5,6 +5,9 @@ import com.correx.infrastructure.tools.filesystem.FileEditTool import com.correx.infrastructure.tools.filesystem.FileReadTool import com.correx.infrastructure.tools.filesystem.FileWriteTool import com.correx.infrastructure.tools.shell.ShellTool +import com.correx.infrastructure.tools.web.WebFetchTool +import com.correx.infrastructure.tools.web.WebSearchTool +import io.ktor.client.HttpClient import java.nio.file.Path data class ToolConfig( @@ -12,6 +15,20 @@ data class ToolConfig( val fileRead: FileReadConfig = FileReadConfig(), val fileWrite: FileWriteConfig = FileWriteConfig(), val fileEdit: FileEditConfig = FileEditConfig(), + val research: ResearchToolConfig = ResearchToolConfig(), +) + +/** + * Research tool registration. [httpClient] is the single shared client threaded in from the + * composition root (one per process, not per registry); when null the tools are not built even + * if [enabled], so tests and the static path never spin up a network client. + */ +class ResearchToolConfig( + val enabled: Boolean = false, + val searxngUrl: String = "http://localhost:8888", + val maxResults: Int = 8, + val maxFetchBytes: Long = 10_000_000, + val httpClient: HttpClient? = null, ) data class FileReadConfig( @@ -73,4 +90,8 @@ fun ToolConfig.buildTools(): List = buildList { ), ) } + research.httpClient?.takeIf { research.enabled }?.let { client -> + add(WebSearchTool(client, searxngBaseUrl = research.searxngUrl, maxResults = research.maxResults)) + add(WebFetchTool(client, maxBytes = research.maxFetchBytes)) + } } diff --git a/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/ResearchWorkflowTest.kt b/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/ResearchWorkflowTest.kt new file mode 100644 index 00000000..ee652b72 --- /dev/null +++ b/infrastructure/workflow/src/test/kotlin/com/correx/infrastructure/workflow/ResearchWorkflowTest.kt @@ -0,0 +1,53 @@ +package com.correx.infrastructure.workflow + +import com.correx.core.artifacts.kind.ConfigArtifactKind +import com.correx.core.artifacts.kind.DefaultArtifactKindRegistry +import com.correx.core.artifacts.kind.JsonSchema +import org.junit.jupiter.api.Assertions.assertEquals +import org.junit.jupiter.api.Assertions.assertTrue +import org.junit.jupiter.api.Test +import java.nio.file.Path +import kotlin.io.path.exists + +/** + * Validates that the shipped research workflow (examples/workflows/research.toml) parses into a + * structurally valid WorkflowGraph: reachable stages, every `needs` produced upstream, the gather + * stage wired to the two network tools. Synthesis quality is inference-dependent and out of scope; + * this guards the graph contract. + */ +class ResearchWorkflowTest { + + private val registry = DefaultArtifactKindRegistry().apply { + listOf("research_plan", "source_dossier", "research_report").forEach { id -> + register(ConfigArtifactKind(id = id, schema = JsonSchema(type = "object"), llmEmitted = true)) + } + } + + private fun repoFile(relative: String): Path { + var dir: Path? = Path.of("").toAbsolutePath() + while (dir != null) { + val candidate = dir.resolve(relative) + if (candidate.exists()) return candidate + dir = dir.parent + } + error("could not locate $relative from ${Path.of("").toAbsolutePath()}") + } + + @Test + fun `research workflow builds a valid graph`() { + val graph = TomlWorkflowLoader(registry).load(repoFile("examples/workflows/research.toml")) + + assertEquals("research", graph.id) + assertEquals("decompose", graph.start.value) + assertEquals(setOf("decompose", "gather", "report"), graph.stages.keys.map { it.value }.toSet()) + + val gather = graph.stages.values.first { it.produces.any { a -> a.name.value == "source_dossier" } } + assertEquals(setOf("research_plan"), gather.needs.map { it.value }.toSet()) + assertEquals(setOf("web_search", "web_fetch"), gather.allowedTools) + + val report = graph.stages.values.first { it.produces.any { a -> a.name.value == "research_report" } } + assertEquals(setOf("source_dossier"), report.needs.map { it.value }.toSet()) + + assertTrue(graph.transitions.any { it.to.value == "done" }, "must terminate at done") + } +}