diff --git a/.changeset/fn-8933-memory-semantics-capture.md b/.changeset/fn-8933-memory-semantics-capture.md new file mode 100644 index 0000000000..a7ec62c144 --- /dev/null +++ b/.changeset/fn-8933-memory-semantics-capture.md @@ -0,0 +1,7 @@ +--- +"@runfusion/fusion": minor +--- + +summary: Add provenance-tagged memory semantics and automatic recall capture. +category: feature +dev: Adds the inferred-edge writer, detached task/research/insight capture roots, and memory semantics audit events. diff --git a/AGENTS.md b/AGENTS.md index 8a5d9731e8..11ac444081 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -282,6 +282,7 @@ Scoped exception (FN-5819/FN-8823): while project auto-merge is On, shared-branc - FN-7158: agent performance reflections emit `reflection:generated`, `reflection:skipped`, and `reflection:failed` with ids/counts/outcomes-only metadata; never persist reflection prose or prompt text in run-audit. - FN-7528: a deterministic, non-LLM post-task performance capture (`AgentReflectionService.captureTaskPerformance`) runs once per completed task and emits `reflection:captured` with ids/counts/outcomes-only metadata (`retryReworkCount?`, `filesTouchedCount?`, `packagesTouchedCount?`, `verificationFileScoped?`, `durationMs?`); never persists `verificationScopeReason` free-text or summary prose in run-audit. - FN-8932: Memory Keeper emits `memory:consolidation-completed`, `memory:consolidation-skipped`, and `memory:consolidation-failed`. Metadata is ids/counts/outcomes-only: `reason`/`unavailableReason`, `stage`, and `graphRecoveryReason` are fixed enums; it never includes error class/message, memory content, paths, or graph identifiers. No-op ticks emit no audit row; only disabled, unavailable, or genuinely overlapping ticks emit `-skipped`. +- FN-8933: Memory semantics emits `memory:semantics-inferred`/`memory:semantics-skipped`, while automatic recall capture reserves `memory:capture-recorded`/`memory:capture-failed`; metadata is ids/counts/fixed outcomes only and never includes labels, recalled prose, prompts, model output, or reasoning. Inferred graph writes stamp provenance at the sole core seam; detached task, research, and insight capture writers never block their source lifecycle. - FN-7787: `createResolvedAgentSession` enriches `session:runtime-resolved` with `noModelResolved: true` and `runtimeBuiltInFallbackModel` when a non-mock/non-test session reaches runtime creation without a complete provider/model pair; this is a visibility signal for runtime built-in fallback usage, not a fabricated model-resolution verdict. - FN-8661: `session:runtime-resolved` records `credentialInstanceId` for a resolved explicit credential, plus `credentialInstanceMissing`, requested, and resolved instance ids when a dangling selection falls back to the provider default. Metadata is ids/outcomes-only and never includes credential material. - FN-7835/FN-7844/FN-7859/FN-7878: durable-agent error-state recovery emits `agent:auto-recover-error-state` when either the heartbeat timer or the self-healing sweep clears a recoverable, non-operator-actionable `error` and retries; metadata stays ids/counts/outcomes-only (`agentId`, attempt, limit, source), where `source` is `timer`/`automation`/`self-healing`. Generic/unknown heartbeat failures are recoverable by default because manual Retry often proves they were transient; both entry paths share the `heartbeatErrorRecovery` budget (self-healing keeps `durableErrorRecovery` only for cooldown/stale-path bookkeeping) and emit `agent:error-retry-exhausted` when the shared budget is exhausted and the agent is parked `paused` with `pauseReason:"error-retry-exhausted"`. Only operator-actionable durable heartbeat errors (credentials/OAuth scope, model access, billing/quota, excluding transient auth rotation), plus stale worktree/module-resolution errors handled by their dedicated suppression path, skip the retry budget and emit `agent:error-parked-unrecoverable` with ids/counts/outcomes-only metadata (`agentId`, `source`, optional `attempts`, `limit`) before parking `paused` with `pauseReason:"error-unrecoverable"` for human repair. diff --git a/docs/knowledge-graph.md b/docs/knowledge-graph.md index 8bf2b49e8d..b5022a1eb2 100644 --- a/docs/knowledge-graph.md +++ b/docs/knowledge-graph.md @@ -2,7 +2,7 @@ # Knowledge graph -`fn knowledge-graph build` creates a deterministic, committable structure graph for the FN-8920 memory epic. It is the embedding-free first layer: no LLM, vector recall, MCP API, inferred relationships, or capability bundle is included. +`fn knowledge-graph build` creates the deterministic structure graph for the FN-8920 memory epic. Memory Keeper may subsequently add clearly provenance-tagged inferred semantic relationships; deterministic extraction itself remains LLM-free. ## Artifact and configuration @@ -14,7 +14,7 @@ An operator who explicitly wants a snapshot in history can still force one with ## Model -Nodes are `file`, `module`, `symbol`, `doc-concept`, or `rationale`. Edges are `contains`, `imports`, and `re-exports`, and always include source, owner (`file` or `derived`), and provenance (`extracted` or reserved `inferred`). IDs are path-derived (`file:path`, `module:dir`, `symbol:path#name`, `doc:path#slug~index`, and `rationale:path#area@stamp~index`) with reserved separators percent-escaped. +Nodes are `file`, `module`, `symbol`, `doc-concept`, or `rationale`. Structural edges are `contains`, `imports`, and `re-exports`; semantic edges are `relates-to` and `rationale-supports`. Every edge includes source, owner (`file` or `derived`), and provenance (`extracted` or `inferred`). IDs are path-derived (`file:path`, `module:dir`, `symbol:path#name`, `doc:path#slug~index`, and `rationale:path#area@stamp~index`) with reserved separators percent-escaped. TypeScript/TSX parsing is parser-only. Exported declarations become symbols; duplicate exports collapse to one earliest-position node with `declarationCount`, including invalid source. Syntax errors remain best-effort. `export *` records a re-export relationship but cannot expand names without a checker. Relative imports are resolved lexically using `.ts`, `.tsx`, and index candidates; package and tsconfig aliases are out of scope. @@ -34,7 +34,11 @@ Every real source position is recorded as a repository-relative path, line, and FNXC rationale is first-class data. TypeScript-family comment ranges come from the parsed tree, not a raw scanner, which prevents strings, regexes, template text, and JSX text from becoming rationale. Markdown recognizes HTML comments outside fenced or narrowly defined indented code; code state gates the comment opener only, so an already-open multi-header comment is not truncated by indentation. Each stamped header starts a separate rationale node and runs to the next header in its comment. -`queryNodes(filter)`, `neighbors(id, options)`, and `shortestPath(from, to)` are deterministic in-process APIs. Neighbor and path results retain complete edge objects, including source, ownership, and `extracted` provenance. `inferred` is reserved in the schema for the later memory-agent layer and is not emitted by this layer. +`queryNodes(filter)`, `neighbors(id, options)`, and `shortestPath(from, to)` are deterministic in-process APIs. Neighbor and path results retain complete edge objects, including source, ownership, and provenance. The sole inferred-edge writer stamps `inferred` unconditionally and accepts an input type with no provenance field, so model output can never masquerade as deterministic extraction. + +## Automatic recall capture + +Memory capture is optional and detached. `RecallCaptureWriter.capture()` returns `void`, so task completion, research finalization/promotion, and insight upsert cannot await recall persistence. Completed tasks and research findings capture `solution` records with task-completion/deep-research provenance; insights capture `decision` records with the available `other` provenance. Live writers are composed at the in-process reflection runtime, research orchestrator construction, research-promotion callers, and the lazy async insight-store factory; absent memory falls back to a shared no-op writer. ## Non-goals diff --git a/packages/core/src/__tests__/memory/recall-capture.test.ts b/packages/core/src/__tests__/memory/recall-capture.test.ts new file mode 100644 index 0000000000..14dd70cd9e --- /dev/null +++ b/packages/core/src/__tests__/memory/recall-capture.test.ts @@ -0,0 +1,111 @@ +import { describe, expect, it, vi } from "vitest"; +import { + buildRecallCaptureContent, + createRecallCaptureWriter, + RECALL_CAPTURE_CONTENT_MAX_BYTES, + RECALL_CAPTURE_KIND_BY_ORIGIN, + RECALL_CAPTURE_SOURCE_ORIGIN_BY_ORIGIN, +} from "../../memory/recall-capture.js"; +import type { RecallAppendInput } from "../../memory/recall/recall-types.js"; + +const layer = {} as never; +const logger = { warn: vi.fn() }; + +function created(input: RecallAppendInput) { + return { + status: "created" as const, + record: { + id: "recall-1", projectId: "project", kind: input.kind, content: input.content, + contentHash: "hash", source: input.source, tags: input.tags ?? [], graphNodeIds: input.graphNodeIds ?? [], + createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", + }, + }; +} + +describe("recall capture writer", () => { + it("maps every automatic origin to FN-8922 kinds and source provenance", async () => { + const append = vi.fn(async (input: RecallAppendInput) => created(input)); + const writer = createRecallCaptureWriter({ layer, logger, append }); + writer.capture({ origin: "task-completion", summary: "Implemented the safe retry.", taskId: "FN-1", agentId: "agent-1" }); + writer.capture({ origin: "research-finding", summary: "The upstream API rejects empty cursors." }); + writer.capture({ origin: "insight", summary: "Prefer bounded paging." }); + await writer.flushPendingCaptures(); + + expect(RECALL_CAPTURE_KIND_BY_ORIGIN).toEqual({ "task-completion": "solution", "research-finding": "solution", insight: "decision" }); + expect(RECALL_CAPTURE_SOURCE_ORIGIN_BY_ORIGIN).toEqual({ "task-completion": "task-completion", "research-finding": "deep-research", insight: "other" }); + expect(append.mock.calls.map(([input]) => ({ kind: input.kind, origin: input.source.origin }))).toEqual([ + { kind: "solution", origin: "task-completion" }, + { kind: "solution", origin: "deep-research" }, + { kind: "decision", origin: "other" }, + ]); + expect(append.mock.calls[0]?.[0].source).toMatchObject({ taskId: "FN-1", agentId: "agent-1" }); + }); + + it("builds structured UTF-8-safe content within its durable byte budget", () => { + const content = buildRecallCaptureContent({ origin: "insight", title: "✨".repeat(3_000), summary: "é".repeat(3_000) }); + expect(content).toContain("[insight]\nTitle:"); + expect(Buffer.byteLength(content, "utf8")).toBeLessThanOrEqual(RECALL_CAPTURE_CONTENT_MAX_BYTES); + expect(Buffer.from(content, "utf8").toString("utf8")).toBe(content); + }); + + it("returns undefined synchronously even when persistence never settles", () => { + const append = vi.fn(() => new Promise(() => {})); + const writer = createRecallCaptureWriter({ layer, logger, append }); + expect(writer.capture({ origin: "insight", summary: "Never block the insight seam." })).toBeUndefined(); + expect(append).toHaveBeenCalledTimes(1); + }); + + it("swallows a failed append and logs exactly once", async () => { + const warn = vi.fn(); + const writer = createRecallCaptureWriter({ layer, logger: { warn }, append: async () => { throw new Error("recall unavailable"); } }); + writer.capture({ origin: "research-finding", summary: "Failure remains best effort." }); + await writer.flushPendingCaptures(); + expect(warn).toHaveBeenCalledTimes(1); + expect(warn).toHaveBeenCalledWith("Automatic recall capture failed for research-finding"); + }); + + it("records ids-only capture audit metadata after persistence", async () => { + const audit = vi.fn(async () => {}); + const append = vi.fn(async (input: RecallAppendInput) => created(input)); + const writer = createRecallCaptureWriter({ layer, logger, append, audit }); + writer.capture({ origin: "insight", summary: "distinctive model prose must not reach audit", insightId: "INS-1", sessionId: "INSR-1" }); + await writer.flushPendingCaptures(); + expect(audit).toHaveBeenCalledWith(expect.objectContaining({ + type: "memory:capture-recorded", + metadata: { recallRecordId: "recall-1", outcome: "created" }, + })); + expect(JSON.stringify(audit.mock.calls)).not.toContain("distinctive model prose"); + expect(append.mock.calls[0]?.[0].tags).toContain("insight:INS-1"); + }); + + it("drains in-flight writes deterministically", async () => { + let release!: () => void; + const append = vi.fn(() => new Promise>((resolve) => { release = () => resolve(created({ kind: "decision", content: "x", source: { origin: "other" } })); })); + const writer = createRecallCaptureWriter({ layer, logger, append }); + writer.capture({ origin: "insight", summary: "Await only in tests." }); + let drained = false; + const drain = writer.flushPendingCaptures().then(() => { drained = true; }); + await Promise.resolve(); + expect(drained).toBe(false); + release(); + await drain; + expect(drained).toBe(true); + }); + + it("relies on FN-8922 append deduplication for repeated identical captures", async () => { + const records = new Map(); + const append = vi.fn(async (input: RecallAppendInput) => { + const prior = records.get(input.content); + if (prior) return { status: "duplicate" as const, duplicateOf: created(prior).record, similarity: 1 }; + records.set(input.content, input); + return created(input); + }); + const writer = createRecallCaptureWriter({ layer, logger, append }); + const input = { origin: "task-completion" as const, summary: "Use the existing retry policy.", taskId: "FN-2" }; + writer.capture(input); + writer.capture(input); + await writer.flushPendingCaptures(); + expect(records).toHaveLength(1); + expect(append).toHaveBeenCalledTimes(2); + }); +}); diff --git a/packages/core/src/__tests__/postgres/insight-store.pg.test.ts b/packages/core/src/__tests__/postgres/insight-store.pg.test.ts index 2915d0ed60..a16133d53b 100644 --- a/packages/core/src/__tests__/postgres/insight-store.pg.test.ts +++ b/packages/core/src/__tests__/postgres/insight-store.pg.test.ts @@ -20,12 +20,15 @@ import { } from "../../__test-utils__/pg-test-harness.js"; import { InsightLifecycleError } from "../../insights/insight-store.js"; import type { AsyncInsightStore } from "../../async-stores/async-insight-store.js"; +import { listRecall } from "../../memory/recall/recall-store.js"; +import type { RecallCaptureWriterWithTestDrain } from "../../memory/recall-capture.js"; const pgTest = pgDescribe; pgTest("InsightStore (PostgreSQL backend mode)", () => { const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_insight_store", + projectId: "P-RECALL", }); beforeAll(h.beforeAll); @@ -67,6 +70,28 @@ pgTest("InsightStore (PostgreSQL backend mode)", () => { expect(all).toHaveLength(1); }); + it("composes the live recall writer at the TaskStore insight factory", async () => { + const s = insights(); + const created = await s.upsertInsight("P-RECALL", { + title: "Capture this insight", + content: "durable outcome", + category: "architecture", + fingerprint: "recall-composition", + }); + + // This reaches getInsightStoreImpl's production constructor line, not an injected test hook. + const writer = (s as unknown as { recallCaptureWriter: RecallCaptureWriterWithTestDrain }).recallCaptureWriter; + await writer.flushPendingCaptures(); + const records = await listRecall(h.layer(), { limit: 10 }); + expect(records).toEqual(expect.arrayContaining([ + expect.objectContaining({ + kind: "decision", + source: expect.objectContaining({ origin: "other" }), + tags: expect.arrayContaining(["insight", `insight:${created.id}`]), + }), + ])); + }); + it("countInsights agrees with listInsights for a filtered set", async () => { const s = insights(); await s.upsertInsight("P-CNT", { title: "A", category: "quality", fingerprint: "A" }); diff --git a/packages/core/src/__tests__/postgres/research-execution.pg.test.ts b/packages/core/src/__tests__/postgres/research-execution.pg.test.ts index 108abc86be..bce9b827da 100644 --- a/packages/core/src/__tests__/postgres/research-execution.pg.test.ts +++ b/packages/core/src/__tests__/postgres/research-execution.pg.test.ts @@ -27,6 +27,8 @@ import { type SharedPgTaskStoreHarness, } from "../../__test-utils__/pg-test-harness.js"; import type { AsyncResearchStore } from "../../async-stores/async-research-store.js"; +import { createRecallCaptureWriter } from "../../memory/recall-capture.js"; +import { listRecall } from "../../memory/recall/recall-store.js"; import type { ResearchModelSettings, ResearchProviderConfig, @@ -77,6 +79,7 @@ function makeStubStepRunner(mode: "ok" | "no-sources"): ResearchStepRunnerApi { pgTest("Research run execution (PostgreSQL backend mode)", () => { const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ prefix: "fusion_research_exec", + projectId: "research-execution-capture", }); beforeAll(h.beforeAll); @@ -133,6 +136,41 @@ pgTest("Research run execution (PostgreSQL backend mode)", () => { expect(status.status).toBe("completed"); }); + /* + FNXC:MemoryRecallCapture 2026-08-11-12:06: + Research finalization is reachable through the production orchestrator and must receive a live + recall writer at its composition root. This uses the real async layer and recall store, rather + than an injected hook, so removing ProjectEngine's equivalent writer wiring cannot be mistaken + for a complete automatic-capture implementation. + */ + it("captures a finalized research outcome through the live recall store", async () => { + const store = research(); + const writer = createRecallCaptureWriter({ layer: h.layer(), logger: { warn: () => {} } }); + const orchestrator = new ResearchOrchestrator({ + store, + stepRunner: makeStubStepRunner("ok"), + maxConcurrentRuns: 1, + recallCaptureWriter: writer, + }); + + const runId = await orchestrator.createRun({ + providers: [{ type: "stub" }], + maxSources: 1, + maxSynthesisRounds: 1, + }); + await expect(orchestrator.startRun(runId, "Recall finalized research")).resolves.toMatchObject({ + status: "completed", + }); + await writer.flushPendingCaptures(); + + expect(await listRecall(h.layer(), { limit: 10 })).toEqual(expect.arrayContaining([ + expect.objectContaining({ + kind: "solution", + source: expect.objectContaining({ origin: "deep-research", sessionId: runId }), + }), + ])); + }); + it("persists a failed status when a step yields no sources (clean failure, no unhandled throw)", async () => { const store = research(); const orchestrator = new ResearchOrchestrator({ diff --git a/packages/core/src/__tests__/research-feature-promotion.test.ts b/packages/core/src/__tests__/research-feature-promotion.test.ts new file mode 100644 index 0000000000..1061d6fa71 --- /dev/null +++ b/packages/core/src/__tests__/research-feature-promotion.test.ts @@ -0,0 +1,39 @@ +import { describe, expect, it, vi } from "vitest"; +import { promoteResearchFinding } from "../research/research-feature-promotion.js"; + +describe("promoteResearchFinding recall capture", () => { + it("records the promoted finding without awaiting optional recall persistence", async () => { + let releaseCapture!: () => void; + const captureStarted = new Promise((resolve) => { releaseCapture = resolve; }); + const capture = vi.fn(() => { void captureStarted; }); + const researchStore = { + getRun: vi.fn(async () => ({ + id: "RR-1", + status: "completed", + tags: ["architecture"], + results: { findings: [{ id: "finding-1", heading: "Adopt a seam", content: "Use the durable seam.", sources: [] }] }, + })), + }; + const missionStore = { + addResearchFeature: vi.fn(async () => ({ feature: { id: "F-1" }, reused: false })), + }; + + const promoted = await promoteResearchFinding( + researchStore as never, + missionStore as never, + { runId: "RR-1", findingId: "finding-1", sliceId: "SL-1" }, + { capture }, + ); + + expect(promoted.feature.id).toBe("F-1"); + expect(capture).toHaveBeenCalledWith(expect.objectContaining({ + origin: "research-finding", + title: "Adopt a seam", + summary: "Research finding finding-1 from completed run RR-1 was promoted to the roadmap.", + researchRunId: "RR-1", + findingId: "finding-1", + tags: ["research", "promotion", "architecture"], + })); + releaseCapture(); + }); +}); diff --git a/packages/core/src/async-stores/async-insight-store.ts b/packages/core/src/async-stores/async-insight-store.ts index 1185694e41..74785aabfd 100644 --- a/packages/core/src/async-stores/async-insight-store.ts +++ b/packages/core/src/async-stores/async-insight-store.ts @@ -21,6 +21,7 @@ import { EventEmitter } from "node:events"; import { and, asc, desc, eq, inArray, lte, sql } from "drizzle-orm"; import * as schema from "../postgres/schema/index.js"; import type { AsyncDataLayer, DbTransaction } from "../postgres/data-layer.js"; +import { NOOP_RECALL_CAPTURE_WRITER, type RecallCaptureWriter } from "../memory/recall-capture.js"; import { InsightLifecycleError, TERMINAL_RUN_STATUSES, @@ -565,7 +566,10 @@ export async function listStalePendingRuns( * are not ported, so manual run execution/retry remain a sync-mode capability. */ export class AsyncInsightStore extends EventEmitter { - constructor(private readonly layer: AsyncDataLayer) { + constructor( + private readonly layer: AsyncDataLayer, + private readonly recallCaptureWriter: RecallCaptureWriter = NOOP_RECALL_CAPTURE_WRITER, + ) { super(); } @@ -600,6 +604,23 @@ export class AsyncInsightStore extends EventEmitter { provenance: input.provenance, }); this.emit(result.id === id ? "insight:created" : "insight:updated", result); + + /* + FNXC:InsightRecallCapture 2026-08-11-10:56: + Every async insight upsert shares this store seam, including dashboard and background origins. + The capture writer is deliberately void-only: durable insight persistence remains available when + optional recall is slow or unavailable. + */ + if (result.id === id) { + this.recallCaptureWriter.capture({ + origin: "insight", + title: result.title, + summary: `Insight outcome recorded in ${result.category} with ${result.status} status.`, + sessionId: result.lastRunId ?? undefined, + insightId: result.id, + tags: ["insight", result.category, result.status], + }); + } return result; } diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 0c8979d57b..0bff1fd539 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -2833,6 +2833,7 @@ export { AGENT_ACTIVITY_EVENT_TYPES, AGENT_ACTIVITY_ATTRIBUTIONS, AGENT_ACTIVITY export { appendAgentActivityEvent, queryAgentActivityEvents, getMaxAgentActivitySeq, pruneAgentActivityEvents } from "./task-store/async/async-agent-activity.js"; export { makeAgentActivityEventId, resolveAgentActivityAttribution, agentIdExistsInRoster, formatAgentActivitySummary, sanitizeAgentActivityMetadata } from "./task-store/agent-activity-outbox.js"; export * from "./memory/recall/index.js"; +export * from "./memory/recall-capture.js"; /* * FNXC:MemoryMcp 2026-08-11-00:19: * The dashboard aliases this barrel for browser code. Keep the pure descriptor here, while the diff --git a/packages/core/src/knowledge-graph/__tests__/graph-builder-incremental.test.ts b/packages/core/src/knowledge-graph/__tests__/graph-builder-incremental.test.ts index d90ea5c798..e281066c6f 100644 --- a/packages/core/src/knowledge-graph/__tests__/graph-builder-incremental.test.ts +++ b/packages/core/src/knowledge-graph/__tests__/graph-builder-incremental.test.ts @@ -3,6 +3,8 @@ import { join } from "node:path"; import { tmpdir } from "node:os"; import { afterEach, describe, expect, it, vi } from "vitest"; import { buildKnowledgeGraph } from "../graph-builder.js"; +import { addInferredEdges } from "../inferred-edge-writer.js"; +import { loadArtifacts } from "../graph-store.js"; import { extractFile as realExtractFile } from "../extract-file.js"; import { extractTypeScript as realTypeScript } from "../extract-typescript.js"; const roots:string[]=[]; @@ -35,6 +37,32 @@ describe("incremental graph builder", () => { await rm(join(root,"src","b.ts")); const extract=vi.fn(realExtractFile); const result=await buildKnowledgeGraph({projectRoot:root,graphDir:dir,extractFile:extract,discovery}); expect(extract).not.toHaveBeenCalled(); expect(result.graph.nodes.some(node=>node.id==="file:src/b.ts")).toBe(false); expect(result.graph.edges.some(edge=>edge.kind==="imports")).toBe(false); }); + it("re-anchors retained inferred edges when a source anchor moves", async () => { + const root = await fixture(), dir = join(root, ".fusion-knowledge/graph"); + await mkdir(join(root, "docs")); + await writeFile(join(root, "docs", "concepts.md"), "# First concept\n\n# Second concept\n"); + const documentationDiscovery = { sourceRoots: ["src"], markdownRoots: ["docs"] }; + await buildKnowledgeGraph({ projectRoot: root, graphDir: dir, discovery: documentationDiscovery }); + const initial = await loadArtifacts(dir); + expect(initial.ok).toBe(true); + if (!initial.ok) return; + const [from, to] = initial.graph.nodes.filter(node => node.kind === "doc-concept"); + expect(from).toBeTruthy(); + expect(to).toBeTruthy(); + await addInferredEdges(dir, [{ kind: "relates-to", from: from!.id, to: to!.id }]); + + await writeFile(join(root, "docs", "concepts.md"), "\n# First concept\n\n# Second concept\n"); + await buildKnowledgeGraph({ projectRoot: root, graphDir: dir, discovery: documentationDiscovery }); + + const rebuilt = await loadArtifacts(dir); + expect(rebuilt.ok).toBe(true); + if (!rebuilt.ok) return; + const currentAnchor = rebuilt.graph.nodes.find(node => node.id === from.id)!; + const inferred = rebuilt.graph.edges.find(edge => edge.provenance === "inferred")!; + expect(inferred.source).toEqual(currentAnchor.source); + expect(inferred.ownerPath).toBe(currentAnchor.ownerPath); + }); + it("rebuilds rather than reusing a cache with forged synthetic file provenance", async () => { const root = await fixture(), dir = join(root, ".fusion-knowledge/graph"); await buildKnowledgeGraph({ projectRoot: root, graphDir: dir, discovery }); diff --git a/packages/core/src/knowledge-graph/__tests__/inferred-edge-writer.test.ts b/packages/core/src/knowledge-graph/__tests__/inferred-edge-writer.test.ts new file mode 100644 index 0000000000..8e09572e57 --- /dev/null +++ b/packages/core/src/knowledge-graph/__tests__/inferred-edge-writer.test.ts @@ -0,0 +1,74 @@ +import { mkdtemp, rm } from "node:fs/promises"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { afterEach, describe, expect, it } from "vitest"; +import { addInferredEdges, type InferredEdgeProposal } from "../inferred-edge-writer.js"; +import { loadArtifacts, writeArtifacts } from "../graph-store.js"; +import type { GraphManifest, KnowledgeGraph } from "../graph-types.js"; + +const dirs: string[] = []; +const source = { path: "src/a.ts", line: 2, column: 1 }; +const docSource = { path: "docs/a.md", line: 2, column: 1 }; +const graph: KnowledgeGraph = { + schemaVersion: 2, + nodes: [ + { id: "file:src/a.ts", kind: "file", name: "a", owner: "file", ownerPath: "src/a.ts", source: { path: "src/a.ts", line: 1, column: 1 }, attributes: { ext: ".ts", syntheticSource: "true" } }, + { id: "file:docs/a.md", kind: "file", name: "a", owner: "file", ownerPath: "docs/a.md", source: { path: "docs/a.md", line: 1, column: 1 }, attributes: { ext: ".md", syntheticSource: "true" } }, + { id: "doc:docs/a.md#one~0", kind: "doc-concept", name: "one", owner: "file", ownerPath: "docs/a.md", source: docSource, attributes: {} }, + { id: "doc:docs/a.md#two~1", kind: "doc-concept", name: "two", owner: "file", ownerPath: "docs/a.md", source: docSource, attributes: {} }, + { id: "rationale:src/a.ts#Area@2026-08-11-10:56~0", kind: "rationale", name: "Area", owner: "file", ownerPath: "src/a.ts", source, attributes: {} }, + { id: "symbol:src/a.ts#a", kind: "symbol", name: "a", owner: "file", ownerPath: "src/a.ts", source, attributes: { symbolKind: "function", declarationCount: "1" } }, + ], + edges: [], +}; +const manifest: GraphManifest = { schemaVersion: 2, extractorVersion: 1, files: { "src/a.ts": { hash: "a".repeat(64) }, "docs/a.md": { hash: "b".repeat(64) } } }; + +async function fixture(): Promise { + const dir = await mkdtemp(join(tmpdir(), "kg-inferred-")); + dirs.push(dir); + await writeArtifacts(dir, graph, manifest); + return dir; +} + +afterEach(async () => { await Promise.all(dirs.splice(0).map(dir => rm(dir, { recursive: true, force: true }))); }); + +describe("addInferredEdges", () => { + it("anchors valid LLM proposals and always stamps them inferred", async () => { + const dir = await fixture(); + const proposal = { + kind: "relates-to", + from: "doc:docs/a.md#one~0", + to: "doc:docs/a.md#two~1", + attributes: { rationale: "shared-use" }, + provenance: "extracted", + source: { path: "forged.ts", line: 99, column: 99 }, + } as unknown as InferredEdgeProposal; + + expect(await addInferredEdges(dir, [proposal])).toEqual({ added: 1, deduped: 0, droppedUnresolved: 0 }); + const loaded = await loadArtifacts(dir); + expect(loaded).toMatchObject({ ok: true }); + if (!loaded.ok) return; + expect(loaded.graph.edges).toEqual([expect.objectContaining({ + id: "relates-to|doc:docs/a.md#one~0|doc:docs/a.md#two~1", + provenance: "inferred", + ownerPath: "docs/a.md", + source: docSource, + attributes: { rationale: "shared-use" }, + })]); + }); + + it("skips duplicate and hallucinated proposals without changing graph evidence", async () => { + const dir = await fixture(); + const proposal: InferredEdgeProposal = { kind: "relates-to", from: "doc:docs/a.md#one~0", to: "doc:docs/a.md#two~1" }; + expect(await addInferredEdges(dir, [proposal, proposal, { ...proposal, to: "symbol:missing#x" }])).toEqual({ added: 1, deduped: 1, droppedUnresolved: 1 }); + expect((await loadArtifacts(dir)).ok).toBe(true); + }); + + it("rejects valid-looking edge kinds when endpoint families are not semantic", async () => { + const dir = await fixture(); + expect(await addInferredEdges(dir, [ + { kind: "relates-to", from: "symbol:src/a.ts#a", to: "file:src/a.ts" }, + { kind: "rationale-supports", from: "doc:docs/a.md#one~0", to: "symbol:src/a.ts#a" }, + ])).toEqual({ added: 0, deduped: 0, droppedUnresolved: 2 }); + }); +}); diff --git a/packages/core/src/knowledge-graph/graph-builder.ts b/packages/core/src/knowledge-graph/graph-builder.ts index 6d3b0b91dc..315360f347 100644 --- a/packages/core/src/knowledge-graph/graph-builder.ts +++ b/packages/core/src/knowledge-graph/graph-builder.ts @@ -50,7 +50,9 @@ export async function buildKnowledgeGraph(options: BuildKnowledgeGraphOptions) { ? { schemaVersion: prior.graph.schemaVersion, nodes: prior.graph.nodes.filter(node => node.owner !== "derived"), - edges: prior.graph.edges.filter(edge => edge.owner !== "derived" && edge.kind !== "imports" && edge.kind !== "re-exports"), + // FNXC:KnowledgeGraphInferredEdges 2026-08-11-10:56: FN-8933 semantic edges are durable LLM + // conclusions, not parser-owned facts, so an incremental structural rebuild must retain them. + edges: prior.graph.edges.filter(edge => edge.provenance === "inferred" || (edge.owner !== "derived" && edge.kind !== "imports" && edge.kind !== "re-exports")), } : { schemaVersion: KNOWLEDGE_GRAPH_SCHEMA_VERSION, nodes: [], edges: [] }; const manifest: GraphManifest = prior.ok @@ -73,7 +75,7 @@ export async function buildKnowledgeGraph(options: BuildKnowledgeGraphOptions) { if (!previous) addedFiles++; parsedFiles++; graph.nodes = graph.nodes.filter(node => node.ownerPath !== relPath); - graph.edges = graph.edges.filter(edge => edge.ownerPath !== relPath); + graph.edges = graph.edges.filter(edge => edge.provenance === "inferred" || edge.ownerPath !== relPath); const result = (options.extractFile ?? extractFile)({ relPath, content }, options.deps); graph.nodes.push(...result.nodes); graph.edges.push(...result.edges); @@ -86,7 +88,7 @@ export async function buildKnowledgeGraph(options: BuildKnowledgeGraphOptions) { deletedFiles++; delete manifest.files[relPath]; graph.nodes = graph.nodes.filter(node => node.ownerPath !== relPath); - graph.edges = graph.edges.filter(edge => edge.ownerPath !== relPath); + graph.edges = graph.edges.filter(edge => edge.provenance === "inferred" || edge.ownerPath !== relPath); } const discovered = new Set(files); @@ -112,7 +114,24 @@ export async function buildKnowledgeGraph(options: BuildKnowledgeGraphOptions) { const derived = deriveModules(graph.nodes); graph.nodes.push(...derived.nodes); graph.edges.push(...derived.edges); - const nodeIds = new Set(graph.nodes.map(node => node.id)); + const nodesById = new Map(graph.nodes.map(node => [node.id, node])); + /* + FNXC:KnowledgeGraphInferredEdges 2026-08-11-11:53: + An incremental structural rebuild replaces parsed node anchors while retaining inferred semantic + edges by stable endpoint id. Re-anchor each retained inferred edge to its current `from` node so + the persisted artifact remains self-consistent after a heading, rationale, or symbol moves lines. + */ + graph.edges = graph.edges.map(edge => { + if (edge.provenance !== "inferred") return edge; + const anchor = nodesById.get(edge.from); + return anchor ? { + ...edge, + owner: anchor.owner, + ownerPath: anchor.ownerPath, + source: { ...anchor.source }, + } : edge; + }); + const nodeIds = new Set(nodesById.keys()); graph.edges = graph.edges.filter(edge => nodeIds.has(edge.from) && nodeIds.has(edge.to)); // Validate the deliberately narrow internal-invariant throw surface before any defensive diff --git a/packages/core/src/knowledge-graph/graph-serialization.ts b/packages/core/src/knowledge-graph/graph-serialization.ts index 71208722f3..87b06c9383 100644 --- a/packages/core/src/knowledge-graph/graph-serialization.ts +++ b/packages/core/src/knowledge-graph/graph-serialization.ts @@ -47,7 +47,7 @@ export const serializeManifest = (manifest: GraphManifest) => canonicalJson({ const owners = new Set(["file", "derived"]); const nodeKinds = new Set(["file", "module", "symbol", "doc-concept", "rationale"]); -const edgeKinds = new Set(["contains", "imports", "re-exports"]); +const edgeKinds = new Set(["contains", "imports", "re-exports", "relates-to", "rationale-supports"]); const hasStringMap = (value: unknown): value is Record => !!value && typeof value === "object" && !Array.isArray(value) && Object.values(value).every(entry => typeof entry === "string"); const validPath = (value: unknown) => { if (typeof value !== "string" || value.length === 0) return false; diff --git a/packages/core/src/knowledge-graph/graph-store.ts b/packages/core/src/knowledge-graph/graph-store.ts index c69f5556b4..03ee8988e0 100644 --- a/packages/core/src/knowledge-graph/graph-store.ts +++ b/packages/core/src/knowledge-graph/graph-store.ts @@ -62,8 +62,20 @@ function consistent(graph: KnowledgeGraph, manifest: GraphManifest): boolean { || !isSyntheticAnchor(node, moduleCanonicalPath(node.ownerPath, moduleFiles)); return node.owner !== "file" || node.source.path !== node.ownerPath || hasSyntheticMarker(node); })) return false; + const nodesById = new Map(graph.nodes.map(node => [node.id, node])); if (graph.edges.some(edge => { if (edge.id !== edgeId(edge.kind, edge.from, edge.to)) return true; + /* + FNXC:KnowledgeGraphInferredEdges 2026-08-11-10:56: + FN-8933 anchors an inferred LLM relation to its existing `from` node instead of accepting an + LLM-supplied location. It may cross ownership classes, so parser/derived edge restrictions do + not apply, but its persisted owner and source must remain exactly that canonical node anchor. + */ + if (edge.provenance === "inferred") { + const anchor = nodesById.get(edge.from); + return !anchor || !nodesById.has(edge.to) || edge.owner !== anchor.owner || edge.ownerPath !== anchor.ownerPath + || edge.source.path !== anchor.source.path || edge.source.line !== anchor.source.line || edge.source.column !== anchor.source.column; + } if (edge.owner === "derived") return edge.kind !== "contains" || edge.from !== moduleNodeId(edge.ownerPath) || !isSyntheticAnchor(edge, moduleCanonicalPath(edge.ownerPath, moduleFiles)); diff --git a/packages/core/src/knowledge-graph/graph-types.ts b/packages/core/src/knowledge-graph/graph-types.ts index e1747aed61..82fc245806 100644 --- a/packages/core/src/knowledge-graph/graph-types.ts +++ b/packages/core/src/knowledge-graph/graph-types.ts @@ -16,7 +16,7 @@ export type EdgeProvenance = "extracted" | "inferred"; export type GraphOwner = "file" | "derived"; export type GraphNodeKind = "file" | "module" | "symbol" | "doc-concept" | "rationale"; export type SymbolKind = "function" | "class" | "interface" | "type-alias" | "enum" | "variable" | "namespace" | "alias"; -export type EdgeKind = "contains" | "imports" | "re-exports"; +export type EdgeKind = "contains" | "imports" | "re-exports" | "relates-to" | "rationale-supports"; export interface SourceLocation { path: string; line: number; column: number } export interface GraphNode { id: string; kind: GraphNodeKind; name: string; owner: GraphOwner; ownerPath: string; source: SourceLocation; attributes: Record } export interface GraphEdge { id: string; kind: EdgeKind; from: string; to: string; provenance: EdgeProvenance; owner: GraphOwner; ownerPath: string; source: SourceLocation; attributes: Record } diff --git a/packages/core/src/knowledge-graph/index.ts b/packages/core/src/knowledge-graph/index.ts index de73f379fe..a1340f4920 100644 --- a/packages/core/src/knowledge-graph/index.ts +++ b/packages/core/src/knowledge-graph/index.ts @@ -12,3 +12,4 @@ export * from "./resolve-imports.js"; export * from "./derive-modules.js"; export * from "./graph-builder.js"; export * from "./graph-query.js"; +export * from "./inferred-edge-writer.js"; diff --git a/packages/core/src/knowledge-graph/inferred-edge-writer.ts b/packages/core/src/knowledge-graph/inferred-edge-writer.ts new file mode 100644 index 0000000000..c0164f9a15 --- /dev/null +++ b/packages/core/src/knowledge-graph/inferred-edge-writer.ts @@ -0,0 +1,96 @@ +import { loadArtifacts, writeArtifacts } from "./graph-store.js"; +import { edgeId, KnowledgeGraphError, type EdgeKind, type GraphEdge, type GraphNode } from "./graph-types.js"; + +/** An LLM may propose a relation, but cannot choose its persisted identity or provenance. */ +export interface InferredEdgeProposal { + kind: EdgeKind; + from: string; + to: string; + attributes?: Record; +} + +export interface AddInferredEdgesResult { + added: number; + /** Re-proposals whose deterministic edge identity already exists. */ + deduped: number; + /** Invalid proposals or proposals whose endpoint nodes/families cannot resolve. */ + droppedUnresolved: number; +} + +function validProposal(value: unknown): value is InferredEdgeProposal { + if (!value || typeof value !== "object" || Array.isArray(value)) return false; + const proposal = value as Partial; + return (proposal.kind === "relates-to" || proposal.kind === "rationale-supports") + && typeof proposal.from === "string" && typeof proposal.to === "string" + && (proposal.attributes === undefined || (!!proposal.attributes && typeof proposal.attributes === "object" + && !Array.isArray(proposal.attributes) && Object.values(proposal.attributes).every(value => typeof value === "string"))); +} + +function isAllowedSemanticRelationship(proposal: InferredEdgeProposal, from: GraphNode, to: GraphNode): boolean { + return (proposal.kind === "relates-to" && from.kind === "doc-concept" && to.kind === "doc-concept") + || (proposal.kind === "rationale-supports" && from.kind === "rationale" && to.kind === "symbol"); +} + +function inferredEdge(proposal: InferredEdgeProposal, anchor: GraphNode): GraphEdge { + return { + id: edgeId(proposal.kind, proposal.from, proposal.to), + kind: proposal.kind, + from: proposal.from, + to: proposal.to, + provenance: "inferred", + owner: anchor.owner, + ownerPath: anchor.ownerPath, + source: { ...anchor.source }, + attributes: { ...proposal.attributes }, + }; +} + +/** + * FNXC:KnowledgeGraphInferredEdges 2026-08-11-10:56: + * FN-8933 permits only concept-to-concept and rationale-to-symbol semantic relations between graph + * nodes that already exist. This seam owns edge identity, source anchoring, and the inferred stamp, + * so an LLM response cannot accidentally persist an extracted claim, fabricate a source location, + * or escape the two allowed semantic families. + * + * FNXC:KnowledgeGraphInferredEdges 2026-08-11-11:26: + * Audit consumers need retry-safe outcome truth: existing deterministic identities count as deduped, + * while malformed or unresolvable proposals count as droppedUnresolved. Never collapse these states. + */ +export async function addInferredEdges(graphDir: string, proposals: readonly InferredEdgeProposal[]): Promise { + const loaded = await loadArtifacts(graphDir); + if (!loaded.ok) throw new KnowledgeGraphError("Knowledge graph artifact is unavailable for inferred edges"); + + const nodes = new Map(loaded.graph.nodes.map(node => [node.id, node])); + const existing = new Map(loaded.graph.edges.map(edge => [edge.id, edge])); + let added = 0; + let deduped = 0; + let droppedUnresolved = 0; + + for (const value of proposals) { + if (!validProposal(value)) { + droppedUnresolved++; + continue; + } + const anchor = nodes.get(value.from); + const target = nodes.get(value.to); + if (!anchor || !target || !isAllowedSemanticRelationship(value, anchor, target)) { + droppedUnresolved++; + continue; + } + const edge = inferredEdge(value, anchor); + if (existing.has(edge.id)) { + deduped++; + continue; + } + existing.set(edge.id, edge); + added++; + } + + if (added > 0) { + await writeArtifacts(graphDir, { + ...loaded.graph, + edges: [...existing.values()], + }, loaded.manifest); + } + return { added, deduped, droppedUnresolved }; +} diff --git a/packages/core/src/memory/index.ts b/packages/core/src/memory/index.ts index eaa88cf698..7d9549c88a 100644 --- a/packages/core/src/memory/index.ts +++ b/packages/core/src/memory/index.ts @@ -9,5 +9,6 @@ export * from "./memory-dreams.js"; export * from "./memory-insights.js"; export * from "./project-memory.js"; export * from "./recall/index.js"; +export * from "./recall-capture.js"; export * from "./mcp/index.js"; export * from "./memory-pre-steering.js"; diff --git a/packages/core/src/memory/recall-capture.ts b/packages/core/src/memory/recall-capture.ts new file mode 100644 index 0000000000..105a14a15e --- /dev/null +++ b/packages/core/src/memory/recall-capture.ts @@ -0,0 +1,184 @@ +import { recordRunAuditEvent, type AsyncDataLayer } from "../postgres/data-layer.js"; +import type { Logger } from "../process/logger.js"; +import { appendRecall } from "./recall/recall-store.js"; +import type { RecallAppendInput, RecallAppendResult, RecallKind, RecallOrigin } from "./recall/recall-types.js"; + +/** The largest durable summary a capture origin may add to recall. */ +export const RECALL_CAPTURE_CONTENT_MAX_BYTES = 4_096; + +/** Origins that have an automatic-capture contract, rather than arbitrary recall origins. */ +export type RecallCaptureOrigin = "task-completion" | "research-finding" | "insight"; + +/** + * The durable, bounded material supplied by an automatic capture seam. + * + * Callers provide a summary rather than a raw prompt or model output: recall is a compact memory + * of an outcome, not a second transcript store. + */ +export interface RecallCaptureInput { + origin: RecallCaptureOrigin; + summary: string; + title?: string; + taskId?: string; + agentId?: string; + sessionId?: string; + /** Origin IDs are retained as compact tags because FN-8922 has one generic source shape. */ + researchRunId?: string; + findingId?: string; + insightId?: string; + tags?: readonly string[]; + /** A real graph node id already known by the caller; capture never invents one. */ + graphNodeId?: string; +} + +/** The capture surface is deliberately one-way: callers cannot await durable memory. */ +export interface RecallCaptureWriter { + capture(input: RecallCaptureInput): void; +} + +/** Test-only observability for work intentionally detached from its production seam. */ +export interface RecallCaptureWriterTestDrain { + flushPendingCaptures(): Promise; +} + +export type RecallCaptureWriterWithTestDrain = RecallCaptureWriter & RecallCaptureWriterTestDrain; + +/** The recall kind is part of each origin's durable contract. */ +export const RECALL_CAPTURE_KIND_BY_ORIGIN: Readonly> = { + "task-completion": "solution", + "research-finding": "solution", + insight: "decision", +}; + +/** Maps capture seams onto the real FN-8922 source-origin literals. */ +export const RECALL_CAPTURE_SOURCE_ORIGIN_BY_ORIGIN: Readonly> = { + "task-completion": "task-completion", + "research-finding": "deep-research", + insight: "other", +}; + +export interface RecallCaptureWriterDependencies { + /** The FN-8922 persistence layer used by the production appendRecall API. */ + layer: AsyncDataLayer; + logger: Pick; + /** Test seam only; production omits this and calls FN-8922 appendRecall directly. */ + append?: (input: RecallAppendInput) => Promise; + /** Optional audit adapter; production defaults to the core run-audit persistence seam. */ + audit?: (input: { type: "memory:capture-recorded" | "memory:capture-failed"; metadata: Record }) => Promise; +} + +/* +FNXC:MemoryRecallCapture 2026-08-11-10:57: +Automatic recall capture returns void so a completion, research, or insight call site cannot await +memory persistence and make optional memory load-bearing. The shared no-op is only the absent or +disabled-memory default; named composition roots must replace it with this factory's live writer. + +FNXC:MemoryRecallCapture 2026-08-11-10:57: +FN-8922 has no insight-specific RecallOrigin literal, so insight outcomes use its real "other" +origin while completed tasks and research use "task-completion" and "deep-research" respectively. +The test drain exists solely for deterministic detached-work tests and production must never call it. +*/ +export const NOOP_RECALL_CAPTURE_WRITER: RecallCaptureWriter = Object.freeze({ + capture: () => {}, +}); + +function clampUtf8(value: string, maxBytes: number): string { + if (Buffer.byteLength(value, "utf8") <= maxBytes) return value; + let output = ""; + for (const character of value) { + if (Buffer.byteLength(output + character, "utf8") > maxBytes) break; + output += character; + } + return output; +} + +/** Build a compact, origin-labelled recall summary without retaining raw model material. */ +export function buildRecallCaptureContent(input: RecallCaptureInput): string { + const lines = [`[${input.origin}]`]; + if (input.title?.trim()) lines.push(`Title: ${input.title.trim()}`); + lines.push(`Summary: ${input.summary.trim()}`); + return clampUtf8(lines.join("\n"), RECALL_CAPTURE_CONTENT_MAX_BYTES); +} + +function sourceIdentifierTags(input: RecallCaptureInput): string[] { + return [ + input.researchRunId?.trim() ? `research-run:${input.researchRunId.trim()}` : undefined, + input.findingId?.trim() ? `research-finding:${input.findingId.trim()}` : undefined, + input.insightId?.trim() ? `insight:${input.insightId.trim()}` : undefined, + ].filter((tag): tag is string => Boolean(tag)); +} + +function toAppendInput(input: RecallCaptureInput): RecallAppendInput { + const graphNodeIds = input.graphNodeId?.trim() ? [input.graphNodeId.trim()] : undefined; + return { + kind: RECALL_CAPTURE_KIND_BY_ORIGIN[input.origin], + content: buildRecallCaptureContent(input), + source: { + origin: RECALL_CAPTURE_SOURCE_ORIGIN_BY_ORIGIN[input.origin], + taskId: input.taskId, + agentId: input.agentId, + sessionId: input.sessionId, + }, + tags: [...new Set([input.origin, ...sourceIdentifierTags(input), ...(input.tags ?? [])])], + graphNodeIds, + }; +} + +/** + * Create the detached automatic-capture writer. Its drain is test-only: it exists so tests can + * observe background writes without sleeps or polling. + */ +export function createRecallCaptureWriter( + deps: RecallCaptureWriterDependencies, +): RecallCaptureWriterWithTestDrain { + const pending = new Set>(); + const append = deps.append ?? ((input: RecallAppendInput) => appendRecall(deps.layer, input)); + const recordAudit = async (type: "memory:capture-recorded" | "memory:capture-failed", input: RecallCaptureInput, metadata: Record) => { + if (deps.audit) { + await deps.audit({ type, metadata }); + } else if (!deps.append) { + await recordRunAuditEvent(deps.layer, { + agentId: input.agentId ?? "memory-capture", + runId: `memory-capture:${input.origin}:${Date.now()}`, + taskId: input.taskId, + domain: "database", + mutationType: type, + target: input.insightId ?? input.findingId ?? input.researchRunId ?? input.taskId ?? input.origin, + metadata: { origin: input.origin, ...metadata }, + }); + } + }; + + return { + capture(input) { + const operation = (async () => { + try { + const result = await append(toAppendInput(input)); + try { + await recordAudit("memory:capture-recorded", input, { + recallRecordId: result.status === "created" ? result.record.id : result.duplicateOf.id, + outcome: result.status, + }); + } catch { + deps.logger.warn(`Automatic recall capture audit failed for ${input.origin}`); + } + } catch (error) { + // Recall content can be sensitive, so diagnostics identify only the bounded origin. + deps.logger.warn(`Automatic recall capture failed for ${input.origin}`); + try { + await recordAudit("memory:capture-failed", input, { + errorClass: error instanceof Error ? error.name : "unknown", + }); + } catch { + deps.logger.warn(`Automatic recall capture audit failed for ${input.origin}`); + } + } + })(); + pending.add(operation); + void operation.finally(() => pending.delete(operation)); + }, + async flushPendingCaptures() { + await Promise.all([...pending]); + }, + }; +} diff --git a/packages/core/src/research/research-feature-promotion.ts b/packages/core/src/research/research-feature-promotion.ts index 9e1ea49529..93d76346d4 100644 --- a/packages/core/src/research/research-feature-promotion.ts +++ b/packages/core/src/research/research-feature-promotion.ts @@ -1,5 +1,6 @@ import type { AsyncResearchStore } from "../async-stores/async-research-store.js"; import type { AsyncMissionStore } from "../async-stores/async-mission-store.js"; +import { NOOP_RECALL_CAPTURE_WRITER, type RecallCaptureWriter } from "../memory/recall-capture.js"; import { resolveResearchFindingId } from "./research-types.js"; export type ResearchFeaturePromotionInput = { @@ -20,6 +21,7 @@ export async function promoteResearchFinding( researchStore: Pick, missionStore: Pick, input: ResearchFeaturePromotionInput, + recallCaptureWriter: RecallCaptureWriter = NOOP_RECALL_CAPTURE_WRITER, ) { const run = await researchStore.getRun(input.runId); if (!run) throw new Error(`Research run ${input.runId} not found`); @@ -28,11 +30,28 @@ export async function promoteResearchFinding( if (!finding) throw new Error(`Finding ${input.findingId} not found`); const findingId = resolveResearchFindingId(finding); const sourceUrls = [...new Set((finding.sources ?? []).map((url) => url.trim()).filter(Boolean))]; + const title = input.title?.trim() || finding.heading?.trim() || "Research finding"; + const description = input.description?.trim() || finding.content?.trim() || undefined; const promoted = await missionStore.addResearchFeature(input.sliceId, { - title: input.title?.trim() || finding.heading?.trim() || "Research finding", - description: input.description?.trim() || finding.content?.trim() || undefined, + title, + description, acceptanceCriteria: input.acceptanceCriteria?.trim() || undefined, researchProvenance: { researchRunId: run.id, findingId, sourceUrls }, }); + + /* + FNXC:ResearchRecallCapture 2026-08-11-10:56: + Promotion must return as soon as its canonical mission write commits. Recall is a best-effort + projection, so its void-only writer records the promoted finding without delaying or failing the + roadmap action. + */ + recallCaptureWriter.capture({ + origin: "research-finding", + title, + summary: `Research finding ${findingId} from completed run ${run.id} was promoted to the roadmap.`, + researchRunId: run.id, + findingId, + tags: ["research", "promotion", ...run.tags], + }); return { ...promoted, runId: run.id, findingId, citations: sourceUrls }; } diff --git a/packages/core/src/task-store/task-store-helpers.ts b/packages/core/src/task-store/task-store-helpers.ts index f575a8c182..12e6201ea9 100644 --- a/packages/core/src/task-store/task-store-helpers.ts +++ b/packages/core/src/task-store/task-store-helpers.ts @@ -24,6 +24,8 @@ import { TodoStore } from "../stores/todo-store.js"; import { AsyncTodoStore } from "../async-stores/async-todo-store.js"; import { AsyncInsightStore } from "../async-stores/async-insight-store.js"; import { AsyncResearchStore } from "../async-stores/async-research-store.js"; +import { createRecallCaptureWriter } from "../memory/recall-capture.js"; +import { createLogger } from "../process/logger.js"; import { assertColumnTraitsValid } from "../workflows/trait-registry.js"; import { BoardConfig, BranchGroup, MergeRequestRecord, Task, WorkflowStepTemplate, WorkflowWorkItem, WorkflowWorkItemKind } from "../types.js"; import { WorkflowFieldDefinition, WorkflowIr, WorkflowIrColumn } from "../workflows/workflow-ir-types.js"; @@ -494,7 +496,16 @@ export function getInsightStoreImpl(store: TaskStore): InsightStore | AsyncInsig if (!layer) { throw new Error("InsightStore is not available: AsyncDataLayer not initialized in backend mode"); } - store.insightStore = new AsyncInsightStore(layer); + /* + * FNXC:InsightRecallCapture 2026-08-11-10:56: + * This lazy factory is the sole backend AsyncInsightStore composition root. Supplying the + * live writer here keeps every upsert origin automatic while standalone constructors retain + * the writer's no-op default. + */ + store.insightStore = new AsyncInsightStore(layer, createRecallCaptureWriter({ + layer, + logger: createLogger("insight-recall-capture"), + })); } return store.insightStore; diff --git a/packages/dashboard/src/__tests__/research-routes.test.ts b/packages/dashboard/src/__tests__/research-routes.test.ts index f8d9f7e814..707a41c752 100644 --- a/packages/dashboard/src/__tests__/research-routes.test.ts +++ b/packages/dashboard/src/__tests__/research-routes.test.ts @@ -1,8 +1,14 @@ // @vitest-environment node -import { describe, it, expect, vi } from "vitest"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; import express from "express"; import { get as performGet, request as performRequest } from "../test-request.js"; -import { ResearchLifecycleError } from "@fusion/core"; +import * as fusionCore from "@fusion/core"; +import { ResearchLifecycleError, listRecall, resolveResearchFindingId, type RecallCaptureWriterWithTestDrain } from "@fusion/core"; +import { + createSharedPgTaskStoreTestHarness, + pgDescribe, + type SharedPgTaskStoreHarness, +} from "../../../core/src/__test-utils__/pg-test-harness.js"; import { createResearchRouter } from "../research-routes.js"; function createMockStore(options?: { @@ -87,6 +93,8 @@ function createMockStore(options?: { }), appendAgentLog: vi.fn(async () => undefined), log: vi.fn(async () => undefined), + getMissionStore: () => ({ addResearchFeature: vi.fn(async () => ({ feature: { id: "F-1" }, reused: false })) }), + getAsyncLayer: () => undefined, }; } @@ -220,6 +228,7 @@ describe("research-routes", () => { ); }); + it("enriches existing task from finding and returns revision", async () => { const store = createMockStore(); const app = express(); @@ -500,3 +509,70 @@ describe("research-routes", () => { expect(Array.isArray(search.body.runs)).toBe(true); }); }); + +/* +FNXC:MemoryRecallCapture 2026-08-11-12:31: +The dashboard route constructs its writer inline, so this production-shaped PG test wraps—not +replaces—the real factory and drains that real writer before reading recall persistence. A mocked +capture callback cannot prove the route's AsyncDataLayer reaches FN-8922's appendRecall store. +*/ +pgDescribe("research-routes recall capture composition", () => { + const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ + prefix: "fusion_dashboard_research_recall", + projectId: "dashboard-research-recall", + }); + + beforeAll(h.beforeAll); + beforeEach(h.beforeEach); + afterEach(h.afterEach); + afterAll(h.afterAll); + + it("persists a promoted dashboard finding through the live recall writer", async () => { + const store = h.store(); + const researchStore = store.getResearchStore(); + const run = await researchStore.createRun({ query: "Dashboard recall capture" }); + await researchStore.updateRun(run.id, { status: "running" }); + await researchStore.updateRun(run.id, { + status: "completed", + results: { + summary: "Dashboard promotion source", + findings: [{ id: "dashboard-recall-finding", heading: "Dashboard finding", content: "Persisted from the route.", sources: [] }], + }, + }); + const missionStore = store.getMissionStore(); + const mission = await missionStore.createMission({ title: "Dashboard recall mission" }); + const milestone = await missionStore.addMilestone(mission.id, { title: "Milestone" }); + const slice = await missionStore.addSlice(milestone.id, { title: "Slice" }); + const realFactory = fusionCore.createRecallCaptureWriter; + const writerFactory = vi.spyOn(fusionCore, "createRecallCaptureWriter"); + let writer: RecallCaptureWriterWithTestDrain | undefined; + writerFactory.mockImplementation((deps) => { + writer = realFactory(deps); + return writer; + }); + + try { + const app = express(); + app.use(express.json()); + app.use(createResearchRouter(store)); + const response = await performRequest( + app, + "POST", + `/runs/${run.id}/findings/dashboard-recall-finding/promote`, + JSON.stringify({ sliceId: slice.id }), + { "content-type": "application/json" }, + ); + expect(response.status).toBe(201); + await writer!.flushPendingCaptures(); + expect(await listRecall(h.layer(), { limit: 10 })).toEqual(expect.arrayContaining([ + expect.objectContaining({ + kind: "solution", + source: expect.objectContaining({ origin: "deep-research" }), + tags: expect.arrayContaining([`research-run:${run.id}`, "research-finding:dashboard-recall-finding"]), + }), + ])); + } finally { + writerFactory.mockRestore(); + } + }); +}); diff --git a/packages/dashboard/src/research-routes.ts b/packages/dashboard/src/research-routes.ts index d780322b81..8b35e708ba 100644 --- a/packages/dashboard/src/research-routes.ts +++ b/packages/dashboard/src/research-routes.ts @@ -12,6 +12,9 @@ import { buildResearchDocumentKey, resolveResearchFindingId, promoteResearchFinding, + createRecallCaptureWriter, + createLogger, + NOOP_RECALL_CAPTURE_WRITER, type ResearchRunListOptions, type ResearchRunStatus, } from "@fusion/core"; @@ -391,14 +394,22 @@ export function createResearchRouter(store: TaskStore, options?: ServerOptions): if (!sliceId) throw badRequest("sliceId is required"); const missionStore = scopedStore.getMissionStore(); if (!("addResearchFeature" in missionStore)) throw new ApiError(409, "Research promotion requires the PostgreSQL mission store"); - const promoted = await promoteResearchFinding(getStore() as never, missionStore, { - runId: req.params.runId, - findingId: req.params.findingId, - sliceId, - title: typeof req.body?.title === "string" ? req.body.title : undefined, - description: typeof req.body?.description === "string" ? req.body.description : undefined, - acceptanceCriteria: typeof req.body?.acceptanceCriteria === "string" ? req.body.acceptanceCriteria : undefined, - }); + const layer = scopedStore.getAsyncLayer(); + const promoted = await promoteResearchFinding( + getStore() as never, + missionStore, + { + runId: req.params.runId, + findingId: req.params.findingId, + sliceId, + title: typeof req.body?.title === "string" ? req.body.title : undefined, + description: typeof req.body?.description === "string" ? req.body.description : undefined, + acceptanceCriteria: typeof req.body?.acceptanceCriteria === "string" ? req.body.acceptanceCriteria : undefined, + }, + layer + ? createRecallCaptureWriter({ layer, logger: createLogger("research-recall-capture") }) + : NOOP_RECALL_CAPTURE_WRITER, + ); let feature = promoted.feature; if (typeof req.body?.taskId === "string" && req.body.taskId.trim()) feature = await missionStore.linkFeatureToTask(feature.id, req.body.taskId.trim()); if (req.body?.triage === true) feature = await missionStore.triageFeature(feature.id); diff --git a/packages/engine/src/__tests__/agent-mission-tools.test.ts b/packages/engine/src/__tests__/agent-mission-tools.test.ts index 3ea16b5587..6c7b2ecd46 100644 --- a/packages/engine/src/__tests__/agent-mission-tools.test.ts +++ b/packages/engine/src/__tests__/agent-mission-tools.test.ts @@ -1,5 +1,6 @@ import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; -import { RepairGroundTruthStaleError, type TaskStore } from "@fusion/core"; +import * as fusionCore from "@fusion/core"; +import { RepairGroundTruthStaleError, listRecall, type RecallCaptureWriterWithTestDrain, type TaskStore } from "@fusion/core"; import { createSharedPgTaskStoreTestHarness, pgDescribe, type SharedPgTaskStoreHarness } from "../../../core/src/__test-utils__/pg-test-harness.js"; import { createMissionTools } from "../agent-tools.js"; @@ -113,8 +114,9 @@ describe("createMissionTools", () => { it("promotes completed findings through the idempotent mission-store facade", async () => { const addResearchFeature = vi.fn().mockResolvedValue({ reused: false, feature: { id: "F-1", status: "defined" } }); const store = { - getResearchStore: () => ({ getRun: vi.fn().mockResolvedValue({ id: "R-1", status: "completed", results: { findings: [{ heading: "Finding", content: "Evidence", sources: ["https://source.example"] }] } }) }), + getResearchStore: () => ({ getRun: vi.fn().mockResolvedValue({ id: "R-1", status: "completed", tags: [], results: { findings: [{ id: "finding-b481c893", heading: "Finding", content: "Evidence", sources: ["https://source.example"] }] } }) }), getMissionStore: () => ({ addResearchFeature }), + getAsyncLayer: () => undefined, } as never; const tool = createMissionTools(store).find((candidate) => candidate.name === "fn_research_promote_finding")!; const result = await tool.execute("call", { runId: "R-1", findingId: "finding-b481c893", sliceId: "SL-1" }); @@ -259,6 +261,60 @@ pgDescribe("mission validation repair agent tool", () => { .toMatchObject({ status: "defined", loopState: "idle" }); }); + /* + FNXC:MemoryRecallCapture 2026-08-11-12:17: + The agent tool is a named production composition root for research-promotion capture. Retain the + real writer's test-only drain so this verifies both factory composition and the actual recall row, + rather than proving only that promoteResearchFinding accepts a hand-injected callback. + */ + it("captures a promoted finding through the agent-tool composition root", async () => { + const store = h.store(); + const researchStore = store.getResearchStore(); + const run = await researchStore.createRun({ query: "Capture through agent tools", tags: ["agent-tool"] }); + await researchStore.updateRun(run.id, { status: "running" }); + await researchStore.updateRun(run.id, { + status: "completed", + results: { + summary: "Agent-tool research result", + findings: [{ id: "finding-agent-capture", heading: "Capture finding", content: "Persist recall through the live writer.", sources: [] }], + }, + }); + const missionStore = store.getMissionStore(); + const mission = await missionStore.createMission({ title: "Promotion capture" }); + const milestone = await missionStore.addMilestone(mission.id, { title: "Milestone" }); + const slice = await missionStore.addSlice(milestone.id, { title: "Slice" }); + const realCreateRecallCaptureWriter = fusionCore.createRecallCaptureWriter; + const writerFactory = vi.spyOn(fusionCore, "createRecallCaptureWriter"); + let writer: RecallCaptureWriterWithTestDrain | undefined; + writerFactory.mockImplementation((deps) => { + writer = realCreateRecallCaptureWriter(deps); + return writer; + }); + + try { + const tool = createMissionTools(store, { agentId: "agent-capture" }) + .find((candidate) => candidate.name === "fn_research_promote_finding")!; + const result = await tool.execute("capture", { + runId: run.id, + findingId: "finding-agent-capture", + sliceId: slice.id, + }); + expect(result.isError).not.toBe(true); + expect(writerFactory).toHaveBeenCalledWith(expect.objectContaining({ layer: store.getAsyncLayer() })); + // The spy wraps the real factory, so its returned drain observes the actual detached insert. + await writer!.flushPendingCaptures(); + expect(await listRecall(store.getAsyncLayer()!, { limit: 10 })).toEqual(expect.arrayContaining([ + expect.objectContaining({ + kind: "solution", + source: expect.objectContaining({ origin: "deep-research" }), + tags: expect.arrayContaining([`research-run:${run.id}`]), + }), + ])); + } finally { + writerFactory.mockRestore(); + } + }); + it("re-resolves a real stale fence exactly once before clearing", async () => { const task = await h.store().createTask({ description: "Planner delivery", column: "todo" }); const feature = await blockedFeature(task.id); diff --git a/packages/engine/src/__tests__/in-process-runtime.pg.test.ts b/packages/engine/src/__tests__/in-process-runtime.pg.test.ts index 27e9e241ef..eb3523c110 100644 --- a/packages/engine/src/__tests__/in-process-runtime.pg.test.ts +++ b/packages/engine/src/__tests__/in-process-runtime.pg.test.ts @@ -46,7 +46,7 @@ vi.mock("@fusion/core", async (importOriginal) => { }; }); -import { CentralCore, type AsyncCentralClaimStore } from "@fusion/core"; +import { CentralCore, listRecall, type AsyncCentralClaimStore, type RecallCaptureWriterWithTestDrain } from "@fusion/core"; import { InProcessRuntime } from "../runtimes/in-process-runtime.js"; pgDescribe("InProcessRuntime PostgreSQL composition", () => { @@ -201,6 +201,32 @@ pgDescribe("InProcessRuntime PostgreSQL composition", () => { }); await assertConsumerMaterializedSecretsEnv(heartbeatTask.id); + /* + FNXC:MemoryRecallCapture 2026-08-11-11:53: + FN-8933's runtime root must supply the live recall writer to its reflection service; an + injected unit-test writer cannot detect a missing composition line. Drive the runtime-owned + service through a completed task and drain only its test seam before reading the real store. + */ + const reflectionAgent = await runtime.getAgentStore()!.createAgent({ + name: "Runtime recall composition agent", + role: "executor", + }); + const reflectedTask = await taskStore.createTask({ description: "runtime recall composition" }); + await taskStore.updateTask(reflectedTask.id, { + column: "done", + status: "completed", + assignedAgentId: reflectionAgent.id, + }); + const reflectionService = (runtime as unknown as { + executor?: { options?: { reflectionService?: { captureTaskPerformance(agentId: string, taskId: string): Promise; captureWriter: RecallCaptureWriterWithTestDrain } } }; + }).executor?.options?.reflectionService; + expect(reflectionService).toBeDefined(); + await reflectionService!.captureTaskPerformance(reflectionAgent.id, reflectedTask.id); + await reflectionService!.captureWriter.flushPendingCaptures(); + expect(await listRecall(layer!, { limit: 10 })).toEqual(expect.arrayContaining([ + expect.objectContaining({ source: expect.objectContaining({ taskId: reflectedTask.id, agentId: reflectionAgent.id }) }), + ])); + const missionStore = taskStore.getMissionStore(); const mission = await missionStore.createMission({ title: "Runtime composition" }); expect((await missionStore.getMission(mission.id))?.title).toBe("Runtime composition"); diff --git a/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts b/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts index e79a5b61c6..c8dd8c8409 100644 --- a/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts +++ b/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts @@ -22,7 +22,7 @@ import { HeartbeatMonitor } from "../agent-heartbeat.js"; import { MemoryConsolidationError } from "../memory/index.js"; function outcome(changed = false, skipped?: "in-progress") { - return { graphChanged: changed, graphRecoveryReason: changed ? "inconsistent-artifact" as const : null, parsedFiles: 0, reusedFiles: 1, prunedFiles: 0, nodeCount: 1, edgeCount: 0, recallCandidates: 1, recallCreated: changed ? 1 : 0, recallDuplicate: changed ? 0 : 1, crossRefUpdated: 0, crossRefUnchanged: 1, crossRefMissing: 0, durationMs: 1, changed, ...(skipped ? { skipped } : {}) }; + return { graphChanged: changed, graphRecoveryReason: changed ? "inconsistent-artifact" as const : null, parsedFiles: 0, reusedFiles: 1, prunedFiles: 0, nodeCount: 1, edgeCount: 0, recallCandidates: 1, recallCreated: changed ? 1 : 0, recallDuplicate: changed ? 0 : 1, crossRefUpdated: 0, crossRefUnchanged: 1, crossRefMissing: 0, semanticsWritten: changed ? 1 : 0, semanticsDeduped: changed ? 2 : 0, semanticsDroppedUnresolved: changed ? 3 : 0, durationMs: 1, changed, ...(skipped ? { skipped } : {}) }; } function fixture(enabled: unknown, metadata: Record = {}) { @@ -68,10 +68,12 @@ describe("Memory Keeper heartbeat hook", () => { const second = await monitor(f).executeHeartbeat({ agentId: "memory", source: "timer" }); expect(first?.status).toBe("completed"); expect(second?.status).toBe("completed"); expect(memory.run).toHaveBeenCalledWith({ agentId: "memory", projectId: "resolved-project" }); - expect(f.audits).toEqual([ + expect(f.audits).toEqual(expect.arrayContaining([ + expect.objectContaining({ mutationType: "memory:semantics-inferred", target: "memory", metadata: expect.objectContaining({ agentId: "memory", edgesWritten: 1, edgesDeduped: 2, edgesDroppedUnresolved: 3 }) }), expect.objectContaining({ mutationType: "memory:consolidation-completed", target: "memory", metadata: expect.objectContaining({ agentId: "memory", graphRecoveryReason: "inconsistent-artifact", recallCreated: 1 }) }), expect.objectContaining({ mutationType: "memory:consolidation-skipped", target: "memory", metadata: expect.objectContaining({ agentId: "memory", reason: "in-progress" }) }), - ]); + ])); + expect(JSON.stringify(f.audits)).not.toContain("distinctive model prose"); }); it("routes recoverable tick failures through completeRun and the shared exhaustion budget", async () => { diff --git a/packages/engine/src/__tests__/memory-consolidation-ports.test.ts b/packages/engine/src/__tests__/memory-consolidation-ports.test.ts index d6c9c2a800..f82af97518 100644 --- a/packages/engine/src/__tests__/memory-consolidation-ports.test.ts +++ b/packages/engine/src/__tests__/memory-consolidation-ports.test.ts @@ -1,16 +1,17 @@ import { beforeEach, describe, expect, it, vi } from "vitest"; -const mocks = vi.hoisted(() => ({ build: vi.fn(), append: vi.fn(), merge: vi.fn(), resolve: vi.fn() })); +const mocks = vi.hoisted(() => ({ build: vi.fn(), append: vi.fn(), merge: vi.fn(), resolve: vi.fn(), addInferred: vi.fn(), semantics: vi.fn() })); vi.mock("@fusion/core", async () => { const actual = await vi.importActual("@fusion/core"); - return { ...actual, buildKnowledgeGraph: mocks.build, appendRecall: mocks.append, mergeRecallGraphNodeIds: mocks.merge, resolveKnowledgeGraphDir: mocks.resolve }; + return { ...actual, buildKnowledgeGraph: mocks.build, appendRecall: mocks.append, mergeRecallGraphNodeIds: mocks.merge, resolveKnowledgeGraphDir: mocks.resolve, addInferredEdges: mocks.addInferred }; }); +vi.mock("../memory/memory-semantics.js", () => ({ runMemorySemanticsPass: mocks.semantics })); import { resolveMemoryConsolidationPorts } from "../memory/memory-consolidation-adapters.js"; const store = (over: Record = {}) => ({ getAsyncLayer: () => ({ projectId: "project" }), getSettings: vi.fn(async () => ({})), ...over }); describe("resolveMemoryConsolidationPorts", () => { - beforeEach(() => { mocks.resolve.mockReset().mockReturnValue("/repo/.fusion-knowledge/graph"); mocks.build.mockReset(); mocks.append.mockReset(); mocks.merge.mockReset(); }); + beforeEach(() => { mocks.resolve.mockReset().mockReturnValue("/repo/.fusion-knowledge/graph"); mocks.build.mockReset(); mocks.append.mockReset(); mocks.merge.mockReset(); mocks.addInferred.mockReset(); mocks.semantics.mockReset(); }); it.each([ ["no root", { taskStore: store(), rootDir: "", agentId: "memory" }, "no-root-dir"], ["no layer", { taskStore: {}, rootDir: "/repo", agentId: "memory" }, "no-data-layer"], @@ -27,6 +28,16 @@ describe("resolveMemoryConsolidationPorts", () => { mocks.build.mockResolvedValue({ graph: { nodes: [], edges: [] }, changed: false, stats: { parsedFiles: 0, reusedFiles: 1, prunedFiles: 0, recoveryReason: "inconsistent-artifact" } }); await expect(resolved.ports.refreshGraph()).resolves.toMatchObject({ recoveryReason: "inconsistent-artifact", nodeCount: 0 }); }); + it("forwards distinct inferred-edge dedupe and unresolved counts to consolidation", async () => { + const resolved = await resolveMemoryConsolidationPorts({ taskStore: store(), rootDir: "/repo", agentId: "memory" }); + expect(resolved.status).toBe("ready"); if (resolved.status !== "ready") return; + mocks.build.mockResolvedValue({ graph: { nodes: [], edges: [] }, changed: true, stats: { parsedFiles: 0, reusedFiles: 0, prunedFiles: 0, recoveryReason: null } }); + await resolved.ports.refreshGraph(); + mocks.addInferred.mockResolvedValue({ added: 1, deduped: 2, droppedUnresolved: 3 }); + mocks.semantics.mockImplementation(async (input: { write(proposals: unknown[]): Promise }) => input.write([])); + await expect(resolved.ports.runSemantics?.(true)).resolves.toEqual({ written: 1, deduped: 2, droppedUnresolved: 3 }); + }); + it("rejects graph paths under .fusion", async () => { mocks.resolve.mockReturnValue("/repo/.fusion/graph"); await expect(resolveMemoryConsolidationPorts({ taskStore: store(), rootDir: "/repo", agentId: "memory" })).resolves.toEqual({ status: "unavailable", reason: "knowledge-graph-dir-unresolved" }); diff --git a/packages/engine/src/__tests__/memory-semantics-pass.test.ts b/packages/engine/src/__tests__/memory-semantics-pass.test.ts new file mode 100644 index 0000000000..40bbe304e1 --- /dev/null +++ b/packages/engine/src/__tests__/memory-semantics-pass.test.ts @@ -0,0 +1,115 @@ +import { mkdtemp, rm, writeFile } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { join } from "node:path"; +import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; +import { loadArtifacts, neighbors, resolveKnowledgeGraphDir, type KnowledgeGraph } from "@fusion/core"; + +const mocks = vi.hoisted(() => ({ + create: vi.fn(), + prompt: vi.fn(), + mcp: vi.fn(async () => ({ servers: [] })), +})); +vi.mock("../pi.js", () => ({ createFnAgent: mocks.create, promptWithFallback: mocks.prompt })); +vi.mock("../mcp/mcp-resolution.js", () => ({ resolveMcpServersForStore: mocks.mcp })); + +import { parseMemorySemanticsResponse, runMemorySemanticsPass } from "../memory/memory-semantics.js"; +import { resolveMemoryConsolidationPorts } from "../memory/memory-consolidation-adapters.js"; + +const roots: string[] = []; +afterEach(async () => { await Promise.all(roots.splice(0).map(root => rm(root, { recursive: true, force: true }))); }); + +const graph: KnowledgeGraph = { + schemaVersion: 2, + nodes: [ + { id: "doc:docs/a.md#one~0", kind: "doc-concept", name: "one", owner: "file", ownerPath: "docs/a.md", source: { path: "docs/a.md", line: 1, column: 1 }, attributes: {} }, + { id: "doc:docs/a.md#two~1", kind: "doc-concept", name: "two", owner: "file", ownerPath: "docs/a.md", source: { path: "docs/a.md", line: 2, column: 1 }, attributes: {} }, + ], + edges: [], +}; + +const taskStore = { getSettings: vi.fn(async () => ({ defaultProvider: "mock", defaultModelId: "scripted" })) } as never; + +beforeEach(() => { + mocks.create.mockReset(); + mocks.prompt.mockReset(); + mocks.mcp.mockClear(); +}); + +describe("memory semantics pass", () => { + it("accepts only semantic relationship families and drops a model provenance claim", () => { + const result = parseMemorySemanticsResponse(JSON.stringify({ proposals: [ + { kind: "relates-to", from: "doc:a", to: "doc:b", provenance: "extracted" }, + { kind: "contains", from: "doc:a", to: "doc:b" }, + ] })); + expect(result.proposals).toEqual([{ kind: "relates-to", from: "doc:a", to: "doc:b" }]); + }); + + it("drives the readonly mock lane and forwards provenance-free proposals only to the inferred writer", async () => { + let onText: ((text: string) => void) | undefined; + const session = { state: {}, dispose: vi.fn() }; + mocks.create.mockImplementation(async (options) => { onText = options.onText; return { session }; }); + mocks.prompt.mockImplementation(async () => { + onText?.(JSON.stringify({ proposals: [{ kind: "relates-to", from: "doc:docs/a.md#one~0", to: "doc:docs/a.md#two~1", provenance: "extracted" }] })); + }); + const write = vi.fn(async (proposals) => ({ written: proposals.length, deduped: 0, droppedUnresolved: 0 })); + + await expect(runMemorySemanticsPass({ graph, graphChanged: true, taskStore, agentId: "memory", rootDir: "/repo", write })).resolves.toMatchObject({ written: 1, deduped: 0, droppedUnresolved: 0 }); + expect(mocks.create).toHaveBeenCalledWith(expect.objectContaining({ tools: "readonly", defaultProvider: "mock", defaultModelId: "scripted" })); + expect(write).toHaveBeenCalledWith([{ kind: "relates-to", from: "doc:docs/a.md#one~0", to: "doc:docs/a.md#two~1" }]); + expect(session.dispose).toHaveBeenCalledOnce(); + }); + + it("persists mocked model output through the production consolidation adapter as inferred graph evidence", async () => { + const root = await mkdtemp(join(tmpdir(), "memory-semantics-")); + roots.push(root); + await writeFile(join(root, "AGENTS.md"), "# One\n\n# Two\n"); + const layer = { projectId: "P-SEMANTICS" }; + const productionStore = { + getAsyncLayer: () => layer, + getSettings: async () => ({ defaultProvider: "mock", defaultModelId: "scripted" }), + } as never; + const resolved = await resolveMemoryConsolidationPorts({ taskStore: productionStore, rootDir: root, agentId: "memory" }); + expect(resolved.status).toBe("ready"); + if (resolved.status !== "ready") return; + + await resolved.ports.refreshGraph(); + const before = await loadArtifacts(resolveKnowledgeGraphDir(root)); + expect(before.ok).toBe(true); + if (!before.ok) return; + const [from, to] = before.graph.nodes.filter(node => node.kind === "doc-concept").map(node => node.id); + expect(from).toBeTruthy(); + expect(to).toBeTruthy(); + + let onText: ((text: string) => void) | undefined; + mocks.create.mockImplementation(async (options) => { onText = options.onText; return { session: { state: {}, dispose: vi.fn() } }; }); + mocks.prompt.mockImplementation(async () => { + onText?.(JSON.stringify({ proposals: [{ kind: "relates-to", from, to, provenance: "extracted" }] })); + }); + + await expect(resolved.ports.runSemantics(true)).resolves.toMatchObject({ written: 1 }); + const loaded = await loadArtifacts(resolveKnowledgeGraphDir(root)); + expect(loaded.ok).toBe(true); + if (!loaded.ok) return; + expect(neighbors(loaded.graph, from!)).toEqual(expect.arrayContaining([ + expect.objectContaining({ + node: expect.objectContaining({ id: to }), + edges: expect.arrayContaining([expect.objectContaining({ provenance: "inferred" })]), + }), + ])); + }); + + it("does not contact a model or write during an unchanged repeat tick", async () => { + const write = vi.fn(); + await expect(runMemorySemanticsPass({ graph, graphChanged: false, taskStore, agentId: "memory", rootDir: "/repo", write })).resolves.toEqual({ proposals: [], skipped: "unchanged" }); + expect(mocks.create).not.toHaveBeenCalled(); + expect(mocks.prompt).not.toHaveBeenCalled(); + expect(write).not.toHaveBeenCalled(); + }); + + it("treats malformed model output as a counted, non-throwing skip", async () => { + let onText: ((text: string) => void) | undefined; + mocks.create.mockImplementation(async (options) => { onText = options.onText; return { session: { state: {}, dispose: vi.fn() } }; }); + mocks.prompt.mockImplementation(async () => { onText?.("not json"); }); + await expect(runMemorySemanticsPass({ graph, graphChanged: true, taskStore, agentId: "memory", rootDir: "/repo", write: vi.fn() })).resolves.toMatchObject({ skipped: "malformed-response" }); + }); +}); diff --git a/packages/engine/src/__tests__/project-engine.test.ts b/packages/engine/src/__tests__/project-engine.test.ts index 2224f2a10e..2ca3fe9e24 100644 --- a/packages/engine/src/__tests__/project-engine.test.ts +++ b/packages/engine/src/__tests__/project-engine.test.ts @@ -1,5 +1,11 @@ -import { afterEach, beforeEach, describe, expect, it, vi } from "vitest"; -import type { Task } from "@fusion/core"; +import { afterAll, afterEach, beforeAll, beforeEach, describe, expect, it, vi } from "vitest"; +import * as fusionCore from "@fusion/core"; +import { listRecall, type RecallCaptureWriterWithTestDrain, type Task } from "@fusion/core"; +import { + createSharedPgTaskStoreTestHarness, + pgDescribe, + type SharedPgTaskStoreHarness, +} from "../../../core/src/__test-utils__/pg-test-harness.js"; import { ProjectEngine, __resetDeterministicMergerModeDeprecationWarned } from "../project-engine.js"; import { AgentSemaphore, projectAdmissionCoordinator} from "../concurrency/concurrency.js"; // Resolves to the vi.mock factory above (the mocked merger-ai exports the real-shaped @@ -410,6 +416,7 @@ beforeEach(() => { }); describe("ProjectEngine notification ownership wiring", () => { + beforeEach(() => { vi.clearAllMocks(); const mockStore = createMockStore(baseSettings); @@ -3991,3 +3998,57 @@ describe("U9 merge safeguards without prior coverage", () => { }); }); + +/* +FNXC:MemoryRecallCapture 2026-08-11-12:31: +ProjectEngine owns the long-lived research-orchestrator composition root. This fixture preserves +its real writer and AsyncDataLayer, then drains the real detached writer before checking recall +storage; a fake capture callback would not prove finalized production research is persisted. +*/ +pgDescribe("ProjectEngine research recall composition", () => { + const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({ + prefix: "fusion_project_engine_research_recall", + projectId: "project-engine-research-recall", + }); + + beforeAll(h.beforeAll); + beforeEach(async () => { + await h.beforeEach(); + mocks.currentStore = h.store() as unknown as Record; + }); + afterEach(h.afterEach); + afterAll(h.afterAll); + + it("persists finalized research through ProjectEngine's live recall composition", async () => { + const store = h.store(); + const run = await store.getResearchStore().createRun({ query: "ProjectEngine recall composition", tags: ["project-engine"] }); + const realFactory = fusionCore.createRecallCaptureWriter; + const writerFactory = vi.spyOn(fusionCore, "createRecallCaptureWriter"); + let writer: RecallCaptureWriterWithTestDrain | undefined; + writerFactory.mockImplementation((deps) => { + writer = realFactory(deps); + return writer; + }); + const engine = createEngine(); + + try { + await engine.start(); + const orchestrator = engine.getResearchOrchestrator() as unknown as { + runFinalizing(runId: string, output: string, citations: string[], confidence: number | undefined, signal: AbortSignal): Promise; + }; + expect(orchestrator).toBeDefined(); + await orchestrator.runFinalizing(run.id, "ProjectEngine final synthesis", ["https://example.test/project-engine"], 0.9, new AbortController().signal); + await writer!.flushPendingCaptures(); + expect(await listRecall(h.layer(), { limit: 10 })).toEqual(expect.arrayContaining([ + expect.objectContaining({ + kind: "solution", + source: expect.objectContaining({ origin: "deep-research", sessionId: run.id }), + tags: expect.arrayContaining([`research-run:${run.id}`]), + }), + ])); + } finally { + writerFactory.mockRestore(); + await engine.stop(); + } + }); +}); diff --git a/packages/engine/src/agent-heartbeat.ts b/packages/engine/src/agent-heartbeat.ts index f67fef7980..d03ba96018 100644 --- a/packages/engine/src/agent-heartbeat.ts +++ b/packages/engine/src/agent-heartbeat.ts @@ -2396,7 +2396,7 @@ export class HeartbeatMonitor { successful skips; runtime errors remain failed runs so completeRun owns shared recovery. */ if (agent.metadata?.builtInMemoryAgent === true) { - const emit = async (type: "memory:consolidation-completed" | "memory:consolidation-skipped" | "memory:consolidation-failed", metadata: Record) => { + const emit = async (type: "memory:consolidation-completed" | "memory:consolidation-skipped" | "memory:consolidation-failed" | "memory:semantics-inferred" | "memory:semantics-skipped", metadata: Record) => { try { await audit.database({ type, target: agentId, metadata }); } catch { /* audit is best effort */ } }; /* FNXC:MemoryAgent 2026-08-11-10:17: Procedure seeding preserves operator edits and is best-effort, so a filesystem failure cannot prevent deterministic upkeep. */ @@ -2419,10 +2419,14 @@ export class HeartbeatMonitor { try { const outcome = await new MemoryConsolidationService(resolution.ports).runConsolidationTick({ agentId, projectId: resolution.projectId }); if (outcome.skipped) await emit("memory:consolidation-skipped", { agentId, reason: outcome.skipped }); - else if (outcome.changed) { - // `changed` drives emission but is not part of the published audit metadata contract. - const { changed: _changed, skipped: _skipped, ...metadata } = outcome; - await emit("memory:consolidation-completed", { agentId, ...metadata }); + else { + if (outcome.semanticsWritten > 0) await emit("memory:semantics-inferred", { agentId, edgesWritten: outcome.semanticsWritten, edgesDeduped: outcome.semanticsDeduped, edgesDroppedUnresolved: outcome.semanticsDroppedUnresolved }); + if (outcome.semanticsSkipped) await emit("memory:semantics-skipped", { agentId, reason: outcome.semanticsSkipped }); + if (outcome.changed) { + // `changed` drives emission but is not part of the published audit metadata contract. + const { changed: _changed, skipped: _skipped, ...metadata } = outcome; + await emit("memory:consolidation-completed", { agentId, ...metadata }); + } } await this.completeRun(agentId, run.id, { status: "completed", resultJson: { reason: "memory_consolidation", ...outcome }, skipStateTransition: true }); } catch (error) { diff --git a/packages/engine/src/agent-tools.ts b/packages/engine/src/agent-tools.ts index 05e6412b9e..b302c27a4f 100644 --- a/packages/engine/src/agent-tools.ts +++ b/packages/engine/src/agent-tools.ts @@ -4592,7 +4592,15 @@ export function createMissionTools(store: TaskStore, context: MissionToolActorCo if (!("addResearchFeature" in missionStore)) return missionToolResult("Research promotion requires the PostgreSQL mission store", { code: "POSTGRES_REQUIRED" }, true); let promoted: Awaited>; try { - promoted = await fusionCore.promoteResearchFinding(store.getResearchStore() as never, missionStore, p); + const layer = store.getAsyncLayer(); + promoted = await fusionCore.promoteResearchFinding( + store.getResearchStore() as never, + missionStore, + p, + layer + ? fusionCore.createRecallCaptureWriter({ layer, logger: fusionCore.createLogger("research-recall-capture") }) + : fusionCore.NOOP_RECALL_CAPTURE_WRITER, + ); } catch (error) { const message = error instanceof Error ? error.message : String(error); return missionToolResult(message, { code: message.includes("not completed") ? "RUN_NOT_COMPLETED" : message.includes("not found") ? "FINDING_OR_RUN_NOT_FOUND" : "PROMOTION_FAILED", runId: p.runId, findingId: p.findingId }, true); @@ -5963,10 +5971,12 @@ export function createResearchTools(options: ResearchToolsOptions): ToolDefiniti .map((type) => registry.getProvider(type)) .filter((provider): provider is NonNullable => Boolean(provider)), }); + const layer = options.store.getAsyncLayer(); orchestratorState.orchestrator = new ResearchOrchestrator({ store: resolveResearchStore(), stepRunner, maxConcurrentRuns: resolved.limits.maxConcurrentRuns, + ...(layer ? { recallCaptureWriter: fusionCore.createRecallCaptureWriter({ layer, logger: fusionCore.createLogger("research-recall-capture") }) } : {}), }); } diff --git a/packages/engine/src/agents/agent-reflection.ts b/packages/engine/src/agents/agent-reflection.ts index 49430526a9..141af621c2 100644 --- a/packages/engine/src/agents/agent-reflection.ts +++ b/packages/engine/src/agents/agent-reflection.ts @@ -9,11 +9,12 @@ import type { ReflectionMetrics, ReflectionStore, ReflectionTrigger, + RecallCaptureWriter, Task, TaskStore, } from "@fusion/core"; import { createLogger } from "../logger.js"; -import { resolveProjectDefaultModel, columnsWithFlag, resolveWorkflowIrForTask} from "@fusion/core"; +import { NOOP_RECALL_CAPTURE_WRITER, resolveProjectDefaultModel, columnsWithFlag, resolveWorkflowIrForTask} from "@fusion/core"; import { createFnAgent, promptWithFallback } from "../pi.js"; import { resolveMcpServersForStore } from "../mcp/mcp-resolution.js"; import { createRunAuditor, generateSyntheticRunId, type EngineRunContext, type RunAuditor } from "../util/run-audit.js"; @@ -71,6 +72,7 @@ export interface AgentReflectionServiceOptions { rootDir: string; modelProvider?: string; modelId?: string; + captureWriter?: RecallCaptureWriter; } export class AgentReflectionService { @@ -81,6 +83,7 @@ export class AgentReflectionService { private readonly modelProvider?: string; private readonly modelId?: string; + private readonly captureWriter: RecallCaptureWriter; constructor(options: AgentReflectionServiceOptions) { this.agentStore = options.agentStore; @@ -89,6 +92,7 @@ export class AgentReflectionService { this.rootDir = options.rootDir; this.modelProvider = options.modelProvider; this.modelId = options.modelId; + this.captureWriter = options.captureWriter ?? NOOP_RECALL_CAPTURE_WRITER; } private async resolveReflectionModel(): Promise<{ provider?: string; modelId?: string }> { @@ -264,6 +268,9 @@ export class AgentReflectionService { summary: this.buildCapturedSummary(task, outcome, metrics), }); + // FNXC:AgentReflection 2026-08-11-10:55: completion invokes this method inside signal-task-complete's detached async wrapper; the void writer preserves that non-blocking boundary while retaining a deterministic outcome summary. + this.captureWriter.capture({ origin: "task-completion", taskId, agentId, title: task.title, summary: this.buildCapturedSummary(task, outcome, metrics) }); + await this.emitReflectionAudit(auditor, "reflection:captured", agentId, trigger, { taskId, ...options }, { reflectionId: reflection.id, ...(metrics.retryReworkCount !== undefined ? { retryReworkCount: metrics.retryReworkCount } : {}), diff --git a/packages/engine/src/memory/index.ts b/packages/engine/src/memory/index.ts index 5403f0b6b5..64f9272f89 100644 --- a/packages/engine/src/memory/index.ts +++ b/packages/engine/src/memory/index.ts @@ -1,3 +1,4 @@ export * from "./memory-consolidation.js"; export * from "./memory-consolidation-material.js"; export * from "./memory-consolidation-adapters.js"; +export * from "./memory-semantics.js"; diff --git a/packages/engine/src/memory/memory-consolidation-adapters.ts b/packages/engine/src/memory/memory-consolidation-adapters.ts index 525cc377d4..08e5542427 100644 --- a/packages/engine/src/memory/memory-consolidation-adapters.ts +++ b/packages/engine/src/memory/memory-consolidation-adapters.ts @@ -1,6 +1,7 @@ -import { buildKnowledgeGraph, mergeRecallGraphNodeIds, appendRecall, resolveKnowledgeGraphDir, type Settings } from "@fusion/core"; +import { addInferredEdges, buildKnowledgeGraph, mergeRecallGraphNodeIds, appendRecall, resolveKnowledgeGraphDir, type KnowledgeGraph, type Settings } from "@fusion/core"; import { relative, resolve, sep } from "node:path"; import type { MemoryConsolidationPorts } from "./memory-consolidation.js"; +import { runMemorySemanticsPass } from "./memory-semantics.js"; export type MemoryConsolidationUnavailableReason = "no-data-layer" | "no-project-id" | "no-root-dir" | "knowledge-graph-dir-unresolved"; export type MemoryConsolidationPortsResolution = { status: "ready"; projectId: string; ports: MemoryConsolidationPorts } | { status: "unavailable"; reason: MemoryConsolidationUnavailableReason }; @@ -20,8 +21,21 @@ export async function resolveMemoryConsolidationPorts(deps: Deps): Promise { const result = await buildKnowledgeGraph({ projectRoot: deps.rootDir, graphDir, force: false }); return { rationaleNodes: result.graph.nodes.filter((node) => node.kind === "rationale"), nodeCount: result.graph.nodes.length, edgeCount: result.graph.edges.length, changed: result.changed, recoveryReason: result.stats.recoveryReason, stats: { parsedFiles: result.stats.parsedFiles, reusedFiles: result.stats.reusedFiles, prunedFiles: result.stats.prunedFiles } }; }, + refreshGraph: async () => { const result = await buildKnowledgeGraph({ projectRoot: deps.rootDir, graphDir, force: false }); latestGraph = result.graph; return { rationaleNodes: result.graph.nodes.filter((node) => node.kind === "rationale"), nodeCount: result.graph.nodes.length, edgeCount: result.graph.edges.length, changed: result.changed, recoveryReason: result.stats.recoveryReason, stats: { parsedFiles: result.stats.parsedFiles, reusedFiles: result.stats.reusedFiles, prunedFiles: result.stats.prunedFiles } }; }, + runSemantics: async (graphChanged) => { + if (!latestGraph) return { skipped: "graph-unavailable" }; + const result = await runMemorySemanticsPass({ graph: latestGraph, graphChanged, taskStore: deps.taskStore as never, agentId: deps.agentId, rootDir: deps.rootDir, write: async (proposals) => { + const written = await addInferredEdges(graphDir, proposals); + return { + written: written.added, + deduped: written.deduped, + droppedUnresolved: written.droppedUnresolved, + }; + } }); + return result; + }, appendRecall: (input) => appendRecall(layer as never, input), mergeRecallGraphNodeIds: (id, ids) => mergeRecallGraphNodeIds(layer as never, id, ids), clock: Date.now, diff --git a/packages/engine/src/memory/memory-consolidation.ts b/packages/engine/src/memory/memory-consolidation.ts index 1ae8386364..a7b80e7f2c 100644 --- a/packages/engine/src/memory/memory-consolidation.ts +++ b/packages/engine/src/memory/memory-consolidation.ts @@ -2,17 +2,20 @@ import type { GraphNode, RecallAppendInput, RecallAppendResult, RecoveryReason } import { deriveRecallMaterial } from "./memory-consolidation-material.js"; export type MemoryConsolidationTickInput = { agentId: string; projectId: string }; +export type MemorySemanticsSkipReason = "unchanged" | "malformed-response" | "no-candidates" | "graph-unavailable"; export type MemoryConsolidationPorts = { refreshGraph: () => Promise<{ rationaleNodes: GraphNode[]; nodeCount: number; edgeCount: number; changed: boolean; recoveryReason: RecoveryReason | null; stats: { parsedFiles: number; reusedFiles: number; prunedFiles: number } }>; appendRecall: (input: RecallAppendInput) => Promise; mergeRecallGraphNodeIds: (recordId: string, ids: string[]) => Promise<{ status: "updated" | "unchanged" | "missing" }>; + /** Semantic enrichment is optional for test ports but production resolves it beside graph I/O. */ + runSemantics?: (graphChanged: boolean) => Promise<{ written?: number; deduped?: number; droppedUnresolved?: number; skipped?: MemorySemanticsSkipReason }>; clock?: () => number; logger?: { warn(message: string): void }; }; -export type MemoryConsolidationOutcome = { graphChanged: boolean; graphRecoveryReason: RecoveryReason | null; parsedFiles: number; reusedFiles: number; prunedFiles: number; nodeCount: number; edgeCount: number; recallCandidates: number; recallCreated: number; recallDuplicate: number; crossRefUpdated: number; crossRefUnchanged: number; crossRefMissing: number; durationMs: number; changed: boolean; skipped?: "in-progress" }; -export class MemoryConsolidationError extends Error { constructor(readonly stage: "graph" | "recall" | "cross-reference", message: string, options?: ErrorOptions) { super(message, options); this.name = "MemoryConsolidationError"; } } +export type MemoryConsolidationOutcome = { graphChanged: boolean; graphRecoveryReason: RecoveryReason | null; parsedFiles: number; reusedFiles: number; prunedFiles: number; nodeCount: number; edgeCount: number; recallCandidates: number; recallCreated: number; recallDuplicate: number; crossRefUpdated: number; crossRefUnchanged: number; crossRefMissing: number; semanticsWritten: number; semanticsDeduped: number; semanticsDroppedUnresolved: number; semanticsSkipped?: MemorySemanticsSkipReason; durationMs: number; changed: boolean; skipped?: "in-progress" }; +export class MemoryConsolidationError extends Error { constructor(readonly stage: "graph" | "recall" | "cross-reference" | "semantics", message: string, options?: ErrorOptions) { super(message, options); this.name = "MemoryConsolidationError"; } } const active = new Set(); -const empty = (): Omit => ({ graphChanged:false, graphRecoveryReason:null, parsedFiles:0, reusedFiles:0, prunedFiles:0, nodeCount:0, edgeCount:0, recallCandidates:0, recallCreated:0, recallDuplicate:0, crossRefUpdated:0, crossRefUnchanged:0, crossRefMissing:0, durationMs:0, changed:false }); +const empty = (): Omit => ({ graphChanged:false, graphRecoveryReason:null, parsedFiles:0, reusedFiles:0, prunedFiles:0, nodeCount:0, edgeCount:0, recallCandidates:0, recallCreated:0, recallDuplicate:0, crossRefUpdated:0, crossRefUnchanged:0, crossRefMissing:0, semanticsWritten:0, semanticsDeduped:0, semanticsDroppedUnresolved:0, durationMs:0, changed:false }); /* FNXC:MemoryAgent 2026-08-11-09:41: The synchronous claim fences manual same-process reentry before any await. It complements heartbeat's per-agent lock; cross-process recall writes serialize in their own advisory locks. No watermark exists: graph fingerprints, recall dedup, and no-write-on-equal reference merges make an unchanged tick silent. */ export class MemoryConsolidationService { @@ -34,8 +37,10 @@ export class MemoryConsolidationService { const ids = byRecord.get(record.id) ?? new Set(); material.graphNodeIds.forEach((id) => ids.add(id)); byRecord.set(record.id, ids); } let updated=0, unchanged=0, missing=0; + let semanticsWritten = 0, semanticsDeduped = 0, semanticsDroppedUnresolved = 0, semanticsSkipped: MemorySemanticsSkipReason | undefined; + if (this.ports.runSemantics) { try { const result = await this.ports.runSemantics(graph.changed); semanticsWritten = result.written ?? 0; semanticsDeduped = result.deduped ?? 0; semanticsDroppedUnresolved = result.droppedUnresolved ?? 0; semanticsSkipped = result.skipped; } catch (cause) { throw new MemoryConsolidationError("semantics", cause instanceof Error ? cause.message : String(cause), { cause }); } } for (const [id, ids] of byRecord) { try { const result = await this.ports.mergeRecallGraphNodeIds(id, [...ids].sort()); if(result.status === "updated") updated++; else if(result.status === "unchanged") unchanged++; else missing++; } catch (cause) { throw new MemoryConsolidationError("cross-reference", cause instanceof Error ? cause.message : String(cause), { cause }); } } - return { graphChanged:graph.changed, graphRecoveryReason:graph.recoveryReason, parsedFiles:graph.stats.parsedFiles, reusedFiles:graph.stats.reusedFiles, prunedFiles:graph.stats.prunedFiles, nodeCount:graph.nodeCount, edgeCount:graph.edgeCount, recallCandidates:created+duplicate, recallCreated:created, recallDuplicate:duplicate, crossRefUpdated:updated, crossRefUnchanged:unchanged, crossRefMissing:missing, durationMs:(this.ports.clock ?? Date.now)()-started, changed:graph.changed || created>0 || updated>0 }; + return { graphChanged:graph.changed, graphRecoveryReason:graph.recoveryReason, parsedFiles:graph.stats.parsedFiles, reusedFiles:graph.stats.reusedFiles, prunedFiles:graph.stats.prunedFiles, nodeCount:graph.nodeCount, edgeCount:graph.edgeCount, recallCandidates:created+duplicate, recallCreated:created, recallDuplicate:duplicate, crossRefUpdated:updated, crossRefUnchanged:unchanged, crossRefMissing:missing, semanticsWritten, semanticsDeduped, semanticsDroppedUnresolved, ...(semanticsSkipped ? { semanticsSkipped } : {}), durationMs:(this.ports.clock ?? Date.now)()-started, changed:graph.changed || created>0 || updated>0 || semanticsWritten>0 }; } finally { active.delete(key); } } } diff --git a/packages/engine/src/memory/memory-semantics.ts b/packages/engine/src/memory/memory-semantics.ts new file mode 100644 index 0000000000..1caebb79c1 --- /dev/null +++ b/packages/engine/src/memory/memory-semantics.ts @@ -0,0 +1,67 @@ +import { queryNodes, resolveProjectDefaultModel, type GraphNode, type KnowledgeGraph, type TaskStore } from "@fusion/core"; +import { createFnAgent, promptWithFallback } from "../pi.js"; +import { resolveMcpServersForStore } from "../mcp/mcp-resolution.js"; + +export type SemanticProposal = { kind: "relates-to" | "rationale-supports"; from: string; to: string }; +export type MemorySemanticsResult = { proposals: SemanticProposal[]; skipped?: "unchanged" | "malformed-response" | "no-candidates" }; + +const MAX_PROPOSALS = 24; +const MAX_CANDIDATES = 80; + +/* +FNXC:MemoryKnowledgeGraph 2026-08-11-10:55: +FN-8933 confines model work to semantic relationships. Structural extraction remains deterministic +and LLM-free; this pass sends bounded identifiers only, accepts strict JSON only, and delegates every +accepted proposal to the core inferred-edge seam rather than constructing graph edges itself. +*/ +const SEMANTICS_SYSTEM_PROMPT = `You infer only semantic graph links from candidate node identifiers. Return JSON only, no fences, as {"proposals":[{"kind":"relates-to","from":"id","to":"id"}]}. Allowed kinds: relates-to for concept↔concept and rationale-supports for rationale→symbol. Never emit contains, imports, re-exports, provenance, labels, or explanations. At most 24 proposals.`; + +function parseProposal(value: unknown): SemanticProposal | undefined { + if (!value || typeof value !== "object") return undefined; + const row = value as Record; + if ((row.kind !== "relates-to" && row.kind !== "rationale-supports") || typeof row.from !== "string" || typeof row.to !== "string") return undefined; + return { kind: row.kind, from: row.from, to: row.to }; +} + +export function parseMemorySemanticsResponse(text: string): MemorySemanticsResult { + try { + const decoded = JSON.parse(text) as { proposals?: unknown }; + if (!Array.isArray(decoded.proposals)) return { proposals: [], skipped: "malformed-response" }; + const proposals = decoded.proposals.map(parseProposal).filter((value): value is SemanticProposal => Boolean(value)).slice(0, MAX_PROPOSALS); + return { proposals }; + } catch { return { proposals: [], skipped: "malformed-response" }; } +} + +export function selectMemorySemanticCandidates(graph: KnowledgeGraph): GraphNode[] { + const concepts = queryNodes(graph, { kinds: ["doc-concept"], limit: MAX_CANDIDATES }); + const rationales = queryNodes(graph, { kinds: ["rationale"], limit: MAX_CANDIDATES }); + const symbols = queryNodes(graph, { kinds: ["symbol"], limit: MAX_CANDIDATES }); + return [...concepts, ...rationales, ...symbols].sort((left, right) => left.id.localeCompare(right.id)).slice(0, MAX_CANDIDATES); +} + +/** Runs the model only when graph material changed; unchanged ticks are provably model-free. */ +export async function runMemorySemanticsPass(input: { + graph: KnowledgeGraph; + graphChanged: boolean; + taskStore: TaskStore; + agentId: string; + rootDir: string; + write: (proposals: SemanticProposal[]) => Promise<{ written: number; deduped: number; droppedUnresolved: number }>; +}): Promise { + if (!input.graphChanged) return { proposals: [], skipped: "unchanged" }; + const candidates = selectMemorySemanticCandidates(input.graph); + if (!candidates.length) return { proposals: [], skipped: "no-candidates" }; + let response = ""; + let model: { provider?: string; modelId?: string } = {}; + try { const resolved = resolveProjectDefaultModel(await input.taskStore.getSettings()); if (resolved.provider && resolved.modelId) model = { provider: resolved.provider, modelId: resolved.modelId }; } catch { /* runtime default is the established fallback */ } + const { session } = await createFnAgent({ cwd: input.rootDir, systemPrompt: SEMANTICS_SYSTEM_PROMPT, tools: "readonly", ...(model.provider && model.modelId ? { defaultProvider: model.provider, defaultModelId: model.modelId } : {}), mcpServers: (await resolveMcpServersForStore(input.taskStore, { agentId: input.agentId })).servers, onText: (delta: string) => { response += delta; } }); + try { + await promptWithFallback(session, JSON.stringify({ candidates: candidates.map((node) => ({ id: node.id, kind: node.kind, area: node.attributes.fnxcArea })) })); + const state = session.state as { errorMessage?: string; error?: string }; + if (state.errorMessage ?? state.error) return { proposals: [], skipped: "malformed-response" }; + } finally { try { session.dispose(); } catch { /* best-effort */ } } + const parsed = parseMemorySemanticsResponse(response); + if (parsed.skipped) return parsed; + const written = await input.write(parsed.proposals); + return { ...parsed, ...written }; +} diff --git a/packages/engine/src/project-engine.ts b/packages/engine/src/project-engine.ts index a6dfd7119b..93f0d256fa 100644 --- a/packages/engine/src/project-engine.ts +++ b/packages/engine/src/project-engine.ts @@ -47,6 +47,7 @@ import { resolveReboundTargetForTask, REVIEW_ELIGIBLE_SENTINEL_COLUMN, clearMergeConfirmedTransientStatus, classifyGhError, + createRecallCaptureWriter, } from "@fusion/core"; import { assemblePlannerOverseerRuntimeSnapshot } from "./overseer/planner-overseer-runtime-snapshot.js"; import { resolveIntegrationBranch } from "./merge/integration-branch.js"; @@ -1031,10 +1032,12 @@ export class ProjectEngine { modelId: settings.researchGlobalDefaults?.synthesisModelId ?? settings.defaultModelId, }, signal) : undefined; + const layer = store.getAsyncLayer(); this.researchOrchestrator = new ResearchOrchestrator({ store: researchStore, stepRunner: new ResearchStepRunner({ providers, synthesisRunner }), maxConcurrentRuns: settings.researchMaxConcurrentRuns ?? 3, + ...(layer ? { recallCaptureWriter: createRecallCaptureWriter({ layer, logger: runtimeLog }) } : {}), }); this.researchDispatcher = new ResearchRunDispatcher({ store: researchStore, diff --git a/packages/engine/src/research/research-orchestrator.ts b/packages/engine/src/research/research-orchestrator.ts index 225c4f6cc9..946ee84ac4 100644 --- a/packages/engine/src/research/research-orchestrator.ts +++ b/packages/engine/src/research/research-orchestrator.ts @@ -1,4 +1,4 @@ -import type { AsyncResearchStore, ResearchStore } from "@fusion/core"; +import { NOOP_RECALL_CAPTURE_WRITER, type AsyncResearchStore, type RecallCaptureWriter, type ResearchStore } from "@fusion/core"; import type { ResearchCancellationState, ResearchOrchestrationConfig, @@ -45,6 +45,7 @@ export interface ResearchOrchestratorOptions { store: ResearchExecutorStore; stepRunner: ResearchStepRunnerApi; maxConcurrentRuns?: number; + recallCaptureWriter?: RecallCaptureWriter; } interface ActiveRunState { @@ -62,6 +63,7 @@ export class ResearchOrchestrator { private readonly store: ResearchExecutorStore; private readonly stepRunner: ResearchStepRunnerApi; private readonly semaphore: AgentSemaphore; + private readonly recallCaptureWriter: RecallCaptureWriter; private readonly activeRuns = new Map(); private readonly cancellation = new Map(); @@ -69,6 +71,7 @@ export class ResearchOrchestrator { this.store = options.store; this.stepRunner = options.stepRunner; this.semaphore = new AgentSemaphore(options.maxConcurrentRuns ?? 3); + this.recallCaptureWriter = options.recallCaptureWriter ?? NOOP_RECALL_CAPTURE_WRITER; } async createRun(config: ResearchOrchestrationConfig): Promise { @@ -417,6 +420,21 @@ export class ResearchOrchestrator { citations, synthesizedOutput: output, }); + + const run = await this.store.getRun(runId); + /* + FNXC:ResearchRecallCapture 2026-08-11-10:56: + Finalization owns the synthesized research outcome. Its detached capture must never extend the + research lifecycle or turn optional recall persistence into a completed-run failure. + */ + this.recallCaptureWriter.capture({ + origin: "research-finding", + title: run?.topic?.trim() || run?.query?.trim() || "Research synthesis", + summary: `Research synthesis completed with ${citations.length} cited sources${confidence === undefined ? "" : ` at confidence ${confidence}`}.`, + sessionId: runId, + researchRunId: runId, + tags: ["research", "synthesis", ...(run?.tags ?? [])], + }); } private async onCancelled(runId: string): Promise { diff --git a/packages/engine/src/runtimes/in-process-runtime.ts b/packages/engine/src/runtimes/in-process-runtime.ts index 00be56a4c9..7fc367c534 100644 --- a/packages/engine/src/runtimes/in-process-runtime.ts +++ b/packages/engine/src/runtimes/in-process-runtime.ts @@ -1318,6 +1318,17 @@ export class InProcessRuntime } // 5c. Initialize AgentReflectionService (requires agentStore and reflectionStore) + let recallCaptureWriter: import("@fusion/core").RecallCaptureWriter | undefined; + try { + const layer = this.taskStore.getAsyncLayer?.(); + if (layer) { + const { createRecallCaptureWriter } = await import("@fusion/core"); + recallCaptureWriter = createRecallCaptureWriter({ layer, logger: runtimeLog }); + } + } catch (captureInitError) { + // FNXC:MemoryRecallCapture 2026-08-11-10:55: optional recall initialization cannot block runtime startup. + runtimeLog.warn("Recall capture initialization failed; automatic capture remains disabled:", captureInitError instanceof Error ? captureInitError.message : captureInitError); + } let reflectionService: import("../agents/agent-reflection.js").AgentReflectionService | undefined; if (agentStoreForReflection && reflectionStoreForService) { try { @@ -1327,6 +1338,7 @@ export class InProcessRuntime taskStore: this.taskStore, reflectionStore: reflectionStoreForService, rootDir: this.config.workingDirectory, + ...(recallCaptureWriter ? { captureWriter: recallCaptureWriter } : {}), }); runtimeLog.log("AgentReflectionService initialized"); } catch (reflServiceErr) { diff --git a/packages/engine/src/util/run-audit.ts b/packages/engine/src/util/run-audit.ts index 9816965403..a675958605 100644 --- a/packages/engine/src/util/run-audit.ts +++ b/packages/engine/src/util/run-audit.ts @@ -809,6 +809,17 @@ export type DatabaseMutationType = | "memory:consolidation-completed" | "memory:consolidation-skipped" | "memory:consolidation-failed" + /* + * FNXC:MemoryAgent 2026-08-11-10:55: + * FN-8933 semantic and capture telemetry remains ids/counts/outcomes-only. `memory:semantics-inferred` + * carries edge/proposal counts; `memory:semantics-skipped` carries a closed reason; capture events + * carry a recall record id (when created), origin, and fixed outcome/error class only. Never record + * edge or node labels, recalled prose, prompt text, model output, or model reasoning in metadata. + */ + | "memory:semantics-inferred" + | "memory:semantics-skipped" + | "memory:capture-recorded" + | "memory:capture-failed" | "task:in-review-stall-deadlock-disposed" | "task:in-review-stall-terminal-provider-error" | "task:finalize-unproven-blocked"