FN-8933: capture semantic memory from completed work
Add inferred graph relationships and detached recall capture across completed-work memory surfaces. - Validate and persist deterministic inferred semantic edges with audit outcomes. - Capture completed tasks, research findings, and insights as bounded recall entries. - Wire memory semantics through engine, research APIs, and coverage tests. Files changed: .changeset/fn-8933-memory-semantics-capture.md | 7 + AGENTS.md | 1 + docs/knowledge-graph.md | 10 +- .../src/__tests__/memory/recall-capture.test.ts | 111 +++++++++++++ .../__tests__/postgres/insight-store.pg.test.ts | 25 +++ .../postgres/research-execution.pg.test.ts | 38 +++++ .../__tests__/research-feature-promotion.test.ts | 39 +++++ .../core/src/async-stores/async-insight-store.ts | 23 ++- packages/core/src/index.ts | 1 + .../__tests__/graph-builder-incremental.test.ts | 28 ++++ .../__tests__/inferred-edge-writer.test.ts | 74 +++++++++ packages/core/src/knowledge-graph/graph-builder.ts | 27 ++- .../src/knowledge-graph/graph-serialization.ts | 2 +- packages/core/src/knowledge-graph/graph-store.ts | 12 ++ packages/core/src/knowledge-graph/graph-types.ts | 2 +- packages/core/src/knowledge-graph/index.ts | 1 + .../src/knowledge-graph/inferred-edge-writer.ts | 96 +++++++++++ packages/core/src/memory/index.ts | 1 + packages/core/src/memory/recall-capture.ts | 184 +++++++++++++++++++++ .../src/research/research-feature-promotion.ts | 23 ++- packages/core/src/task-store/task-store-helpers.ts | 13 +- .../src/__tests__/research-routes.test.ts | 80 ++++++++- packages/dashboard/src/research-routes.ts | 27 ++- .../src/__tests__/agent-mission-tools.test.ts | 60 ++++++- .../src/__tests__/in-process-runtime.pg.test.ts | 28 +++- .../memory-consolidation-heartbeat-hook.test.ts | 8 +- .../__tests__/memory-consolidation-ports.test.ts | 17 +- .../src/__tests__/memory-semantics-pass.test.ts | 115 +++++++++++++ .../engine/src/__tests__/project-engine.test.ts | 65 +++++++- packages/engine/src/agent-heartbeat.ts | 14 +- packages/engine/src/agent-tools.ts | 12 +- packages/engine/src/agents/agent-reflection.ts | 9 +- packages/engine/src/memory/index.ts | 1 + .../src/memory/memory-consolidation-adapters.ts | 18 +- packages/engine/src/memory/memory-consolidation.ts | 13 +- packages/engine/src/memory/memory-semantics.ts | 67 ++++++++ packages/engine/src/project-engine.ts | 3 + .../engine/src/research/research-orchestrator.ts | 20 ++- packages/engine/src/runtimes/in-process-runtime.ts | 12 ++ packages/engine/src/util/run-audit.ts | 11 ++ 40 files changed, 1250 insertions(+), 48 deletions(-) Fusion-Task-Id: FN-8933 Fusion-Task-Lineage: b75a6b23-906f-4255-92f1-7684beb742b8 Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
7
.changeset/fn-8933-memory-semantics-capture.md
Normal file
7
.changeset/fn-8933-memory-semantics-capture.md
Normal file
@@ -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.
|
||||
@@ -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.
|
||||
|
||||
@@ -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
|
||||
|
||||
|
||||
111
packages/core/src/__tests__/memory/recall-capture.test.ts
Normal file
111
packages/core/src/__tests__/memory/recall-capture.test.ts
Normal file
@@ -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<never>(() => {}));
|
||||
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<ReturnType<typeof created>>((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<string, RecallAppendInput>();
|
||||
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);
|
||||
});
|
||||
});
|
||||
@@ -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" });
|
||||
|
||||
@@ -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({
|
||||
|
||||
@@ -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<void>((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();
|
||||
});
|
||||
});
|
||||
@@ -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<InsightStoreEvents> {
|
||||
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<InsightStoreEvents> {
|
||||
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;
|
||||
}
|
||||
|
||||
|
||||
@@ -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
|
||||
|
||||
@@ -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 });
|
||||
|
||||
@@ -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<string> {
|
||||
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 });
|
||||
});
|
||||
});
|
||||
@@ -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
|
||||
|
||||
@@ -47,7 +47,7 @@ export const serializeManifest = (manifest: GraphManifest) => canonicalJson({
|
||||
|
||||
const owners = new Set<GraphOwner>(["file", "derived"]);
|
||||
const nodeKinds = new Set<GraphNodeKind>(["file", "module", "symbol", "doc-concept", "rationale"]);
|
||||
const edgeKinds = new Set<EdgeKind>(["contains", "imports", "re-exports"]);
|
||||
const edgeKinds = new Set<EdgeKind>(["contains", "imports", "re-exports", "relates-to", "rationale-supports"]);
|
||||
const hasStringMap = (value: unknown): value is Record<string, string> => !!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;
|
||||
|
||||
@@ -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));
|
||||
|
||||
@@ -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<string,string> }
|
||||
export interface GraphEdge { id: string; kind: EdgeKind; from: string; to: string; provenance: EdgeProvenance; owner: GraphOwner; ownerPath: string; source: SourceLocation; attributes: Record<string,string> }
|
||||
|
||||
@@ -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";
|
||||
|
||||
96
packages/core/src/knowledge-graph/inferred-edge-writer.ts
Normal file
96
packages/core/src/knowledge-graph/inferred-edge-writer.ts
Normal file
@@ -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<string, string>;
|
||||
}
|
||||
|
||||
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<InferredEdgeProposal>;
|
||||
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<AddInferredEdgesResult> {
|
||||
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 };
|
||||
}
|
||||
@@ -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";
|
||||
|
||||
184
packages/core/src/memory/recall-capture.ts
Normal file
184
packages/core/src/memory/recall-capture.ts
Normal file
@@ -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<void>;
|
||||
}
|
||||
|
||||
export type RecallCaptureWriterWithTestDrain = RecallCaptureWriter & RecallCaptureWriterTestDrain;
|
||||
|
||||
/** The recall kind is part of each origin's durable contract. */
|
||||
export const RECALL_CAPTURE_KIND_BY_ORIGIN: Readonly<Record<RecallCaptureOrigin, RecallKind>> = {
|
||||
"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<Record<RecallCaptureOrigin, RecallOrigin>> = {
|
||||
"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<Logger, "warn">;
|
||||
/** Test seam only; production omits this and calls FN-8922 appendRecall directly. */
|
||||
append?: (input: RecallAppendInput) => Promise<RecallAppendResult>;
|
||||
/** Optional audit adapter; production defaults to the core run-audit persistence seam. */
|
||||
audit?: (input: { type: "memory:capture-recorded" | "memory:capture-failed"; metadata: Record<string, string> }) => Promise<void>;
|
||||
}
|
||||
|
||||
/*
|
||||
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<Promise<void>>();
|
||||
const append = deps.append ?? ((input: RecallAppendInput) => appendRecall(deps.layer, input));
|
||||
const recordAudit = async (type: "memory:capture-recorded" | "memory:capture-failed", input: RecallCaptureInput, metadata: Record<string, string>) => {
|
||||
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]);
|
||||
},
|
||||
};
|
||||
}
|
||||
@@ -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<AsyncResearchStore, "getRun">,
|
||||
missionStore: Pick<AsyncMissionStore, "addResearchFeature">,
|
||||
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 };
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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);
|
||||
|
||||
@@ -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<unknown>; 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");
|
||||
|
||||
@@ -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<string, unknown> = {}) {
|
||||
@@ -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 () => {
|
||||
|
||||
@@ -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<typeof import("@fusion/core")>("@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<string, unknown> = {}) => ({ 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<unknown> }) => 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" });
|
||||
|
||||
115
packages/engine/src/__tests__/memory-semantics-pass.test.ts
Normal file
115
packages/engine/src/__tests__/memory-semantics-pass.test.ts
Normal file
@@ -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" });
|
||||
});
|
||||
});
|
||||
@@ -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<string, unknown>;
|
||||
});
|
||||
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<void>;
|
||||
};
|
||||
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();
|
||||
}
|
||||
});
|
||||
});
|
||||
|
||||
@@ -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<string, unknown>) => {
|
||||
const emit = async (type: "memory:consolidation-completed" | "memory:consolidation-skipped" | "memory:consolidation-failed" | "memory:semantics-inferred" | "memory:semantics-skipped", metadata: Record<string, unknown>) => {
|
||||
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) {
|
||||
|
||||
@@ -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<ReturnType<typeof fusionCore.promoteResearchFinding>>;
|
||||
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<typeof provider> => 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") }) } : {}),
|
||||
});
|
||||
}
|
||||
|
||||
|
||||
@@ -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 } : {}),
|
||||
|
||||
@@ -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";
|
||||
|
||||
@@ -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<Memor
|
||||
const fusionDir = resolve(deps.rootDir, ".fusion"); const rel = relative(fusionDir, graphDir);
|
||||
/* FNXC:MemoryAgent 2026-08-11-10:17: Graph artifacts are committable, so the exact .fusion directory is forbidden alongside all of its children. */
|
||||
if (rel === "" || (!rel.startsWith(`..${sep}`) && rel !== ".." && !rel.includes(`${sep}..${sep}`))) return { status: "unavailable", reason: "knowledge-graph-dir-unresolved" };
|
||||
let latestGraph: KnowledgeGraph | undefined;
|
||||
return { status: "ready", projectId: layer.projectId, ports: {
|
||||
refreshGraph: async () => { 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,
|
||||
|
||||
@@ -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<RecallAppendResult>;
|
||||
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<string>();
|
||||
const empty = (): Omit<MemoryConsolidationOutcome, "skipped"> => ({ 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<MemoryConsolidationOutcome, "skipped"> => ({ 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<string>(); 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); }
|
||||
}
|
||||
}
|
||||
|
||||
67
packages/engine/src/memory/memory-semantics.ts
Normal file
67
packages/engine/src/memory/memory-semantics.ts
Normal file
@@ -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<string, unknown>;
|
||||
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<MemorySemanticsResult & { written?: number; deduped?: number; droppedUnresolved?: number }> {
|
||||
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 };
|
||||
}
|
||||
@@ -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,
|
||||
|
||||
@@ -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<string, ActiveRunState>();
|
||||
private readonly cancellation = new Map<string, ResearchCancellationState>();
|
||||
|
||||
@@ -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<string> {
|
||||
@@ -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<void> {
|
||||
|
||||
@@ -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) {
|
||||
|
||||
@@ -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"
|
||||
|
||||
Reference in New Issue
Block a user