FN-8932: add durable memory consolidation agent

Add a provisioned memory agent that consolidates durable recall material through idempotent heartbeat ticks.

- Provision and configure the durable memory agent with an enabled workflow setting.
- Add consolidation adapters, material collection, recall graph references, and run-audit metadata.
- Expose the setting in the dashboard and document the memory-agent behavior.
- Cover provisioning, heartbeat, consolidation, audit, and recall graph behavior with tests.

Files changed:
 .changeset/fn-8932-memory-agent.md                 |   7 ++
 AGENTS.md                                          |   1 +
 docs/agents.md                                     |  10 ++
 docs/settings-reference.md                         |   1 +
 docs/storage.md                                    |   2 +-
 .../memoryConsolidationEnabled-default.test.ts     |  22 ++++
 .../memory-recall-graph-cross-reference.pg.test.ts | 115 +++++++++++++++++++++
 .../__tests__/memory-agent-provisioning.test.ts    |  31 ++++++
 packages/core/src/agents/agent-store.ts            |  85 +++++++++++++++
 packages/core/src/agents/memory-agent-defaults.ts  |  30 ++++++
 packages/core/src/index.gate.ts                    |   2 +
 packages/core/src/index.ts                         |   9 ++
 packages/core/src/memory/recall/index.ts           |   1 +
 packages/core/src/memory/recall/recall-dedup.ts    |   3 +
 packages/core/src/memory/recall/recall-store.ts    |  35 ++++++-
 .../src/workflows/builtin-workflow-settings.ts     |  14 +++
 .../src/workflows/workflow-settings-resolver.ts    |  12 ++-
 .../__tests__/WorkflowSettingsPanel.test.tsx       |   8 ++
 .../__tests__/workflow-setting-display.test.ts     |   6 ++
 .../app/components/workflow-setting-display.ts     |   5 +
 .../memory-consolidation-heartbeat-hook.test.ts    |  92 +++++++++++++++++
 .../__tests__/memory-consolidation-ports.test.ts   |  34 ++++++
 ...memory-consolidation-run-audit-metadata.test.ts |  19 ++++
 .../__tests__/memory-consolidation-tick.test.ts    |  62 +++++++++++
 packages/engine/src/agent-heartbeat.ts             |  55 ++++++++++
 packages/engine/src/index.ts                       |   1 +
 packages/engine/src/memory/index.ts                |   3 +
 .../src/memory/memory-consolidation-adapters.ts    |  29 ++++++
 .../src/memory/memory-consolidation-material.ts    |  22 ++++
 packages/engine/src/memory/memory-consolidation.ts |  41 ++++++++
 packages/engine/src/util/run-audit.ts              |  10 ++
 31 files changed, 764 insertions(+), 3 deletions(-)

Fusion-Task-Id: FN-8932

Fusion-Task-Lineage: b4fdec50-1f42-4120-af17-0b6f3a94586e

Co-authored-by: Fusion (runfusion.ai) <noreply@runfusion.ai>
This commit is contained in:
gsxdsm
2026-08-11 03:53:11 -07:00
parent 0eaa3c7b9a
commit 637854ad36
31 changed files with 764 additions and 3 deletions

View File

@@ -0,0 +1,7 @@
---
"@runfusion/fusion": minor
---
summary: Add Memory Keeper for deterministic knowledge graph and recall consolidation.
category: feature
dev: Adds the Memory Keeper agent, memoryConsolidationEnabled setting, mergeRecallGraphNodeIds, and memory:consolidation audit events.

View File

@@ -279,6 +279,7 @@ Scoped exception (FN-5819/FN-8823): while project auto-merge is On, shared-branc
- FN-8948: `mission:reconcile-pass` records a bounded automatic reconciliation result. Metadata contains optional mission ID, source enum, and scan/write/skip/conflict/failure counters only; never roadmap prose, titles, reasons, or secrets. - FN-8948: `mission:reconcile-pass` records a bounded automatic reconciliation result. Metadata contains optional mission ID, source enum, and scan/write/skip/conflict/failure counters only; never roadmap prose, titles, reasons, or secrets.
- 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-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-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-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-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-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. - 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.

View File

@@ -1730,6 +1730,16 @@ Per-agent overrides via `runtimeConfig`:
- **Budgets**: per-agent token budget tracking; `HeartbeatMonitor.executeHeartbeat()` skips when `isOverBudget` or `isOverThreshold` (timer triggers). Hard caps pause the agent. - **Budgets**: per-agent token budget tracking; `HeartbeatMonitor.executeHeartbeat()` skips when `isOverBudget` or `isOverThreshold` (timer triggers). Hard caps pause the agent.
- **Performance ratings**: 1–5 scale with trend analysis, injected into system prompts. - **Performance ratings**: 1–5 scale with trend analysis, injected into system prompts.
## Memory Keeper (FN-8932)
Each project provisions a durable **Memory Keeper** custom agent for deterministic, hourly memory upkeep. It is heartbeat-enabled and has task auto-claim disabled, so it cannot claim board work or make product decisions. Provisioning identifies the owner by its provenance marker, not its display name: if an operator already owns `Memory Keeper`, Fusion creates `Memory Keeper (built-in)` instead; if both names are occupied, startup continues without a memory agent rather than renaming/adopting the operator agent or failing initialization.
When enabled, a heartbeat refreshes the knowledge graph incrementally, appends deterministic FNXC rationale decisions through recall deduplication, then merges rationale/file node identifiers into each resulting recall record. Cross-references only grow: the per-record PostgreSQL advisory lock reads, unions, and writes in one transaction, and equal unions perform no update. Pruning is intentionally out of scope. A fingerprint-stable graph, duplicate recall results, and equal cross-reference unions yield a no-write tick; an in-process `(agentId, projectId)` guard skips re-entry. The guard is defense-in-depth for manual callers and does not fence another process or CLI graph build.
Graph files are atomically replaced **per file**, not as an atomic three-file set. Concurrent builders can leave a transient mismatched set, but the manifest is written last and strict consistency validation reports `inconsistent-artifact` and triggers a full rebuild on the next load. This accepted residual risk can create temporary committable-tree artifact noise; `graphRecoveryReason` makes repeated recovery visible. The adapter alone imports graph filesystem/configuration APIs; the tick and material transformer may use only graph types and pure ID helpers so they remain fixture-testable. Missing data-layer/project/root-directory or invalid graph-directory environments are successful skips, while runtime failures use the existing shared heartbeat recovery budget.
The workflow-native `memoryConsolidationEnabled` setting defaults to `true` and is resolved from the project default workflow for no-task heartbeats. Disabling it stops the tick without deleting the agent or memory. The tick makes no LLM calls.
## Workflow role principals ## Workflow role principals
Permanent agents carry one or more normalized role tags: `triage`, `executor`, `reviewer`, `merger`, `scheduler`, `engineer`, and `custom`. Upgrades preserve legacy singular roles, and every project receives four distinct heartbeat-disabled built-ins for planning, execution, review, and merge. Heartbeat enablement and `maxConcurrentRuns` are independent from `runtimeConfig.maxWorkflowSessions`. Permanent agents carry one or more normalized role tags: `triage`, `executor`, `reviewer`, `merger`, `scheduler`, `engineer`, and `custom`. Upgrades preserve legacy singular roles, and every project receives four distinct heartbeat-disabled built-ins for planning, execution, review, and merge. Heartbeat enablement and `maxConcurrentRuns` are independent from `runtimeConfig.maxWorkflowSessions`.

View File

@@ -432,6 +432,7 @@ The built-in workflows also declare triage/spec policy settings that were **not*
| `plannerOverseerAdvisorModelId` | `""` | Session-advisor model id. Used only when `plannerOverseerAdvisorEnabled` is true. When enabled and both model fields are set, the advisor reviews executor agent-log deltas and may inject `[session-advisor]` steering comments at `steer`/`autonomous` (observe = log only). Discover project review priorities via `OVERSEER.md` / `WATCHDOG.md`. See `docs/architecture.md` → "Planner overseer session advisor". | | `plannerOverseerAdvisorModelId` | `""` | Session-advisor model id. Used only when `plannerOverseerAdvisorEnabled` is true. When enabled and both model fields are set, the advisor reviews executor agent-log deltas and may inject `[session-advisor]` steering comments at `steer`/`autonomous` (observe = log only). Discover project review priorities via `OVERSEER.md` / `WATCHDOG.md`. See `docs/architecture.md` → "Planner overseer session advisor". |
| `plannerOverseerExecutorStuckAfterMs` | `7200000` (2h) | Workflow-native executor-stage stall threshold (FN-7743). Milliseconds of executor-stage inactivity — no execution activity since the task's last column move/update (`columnMovedAt ?? updatedAt`) — before a non-paused `in-progress` task is reported `signal: "stuck"` instead of `"progressing"`, feeding the existing `decidePlannerRecovery` → bounded `inject_guidance` recovery path at the `autonomous` oversight level (no effect at `off`/`observe`/`steer`). Fixes the class of bug where a genuinely hung/idle executor (dead session, silent agent) was indistinguishable from a healthy one and was never nudged, retried, or escalated. A missing/malformed activity timestamp degrades to `"progressing"` (fail-safe — never fabricates a stall), and a user-paused/approval-blocked/`autoMerge:false` task is still fully withheld from any autonomous action regardless of this threshold. Resolves through the generic `resolveEffectiveSettings` default path alongside `plannerOversightLevel`. See `docs/architecture.md` → "Executor-stage stall detection (FN-7743)". | | `plannerOverseerExecutorStuckAfterMs` | `7200000` (2h) | Workflow-native executor-stage stall threshold (FN-7743). Milliseconds of executor-stage inactivity — no execution activity since the task's last column move/update (`columnMovedAt ?? updatedAt`) — before a non-paused `in-progress` task is reported `signal: "stuck"` instead of `"progressing"`, feeding the existing `decidePlannerRecovery` → bounded `inject_guidance` recovery path at the `autonomous` oversight level (no effect at `off`/`observe`/`steer`). Fixes the class of bug where a genuinely hung/idle executor (dead session, silent agent) was indistinguishable from a healthy one and was never nudged, retried, or escalated. A missing/malformed activity timestamp degrades to `"progressing"` (fail-safe — never fabricates a stall), and a user-paused/approval-blocked/`autoMerge:false` task is still fully withheld from any autonomous action regardless of this threshold. Resolves through the generic `resolveEffectiveSettings` default path alongside `plannerOversightLevel`. See `docs/architecture.md` → "Executor-stage stall detection (FN-7743)". |
| `plannerHeartbeatPatrolEnabled` | `true` | Workflow-native idle-heartbeat patrol switch (FN-7963). `true` preserves the existing no-task heartbeat/triage guidance that lets idle agents scan for gaps and create focused follow-up tasks. Set to `false` to remove proactive patrol task-creation guidance from idle/no-task heartbeat prompts; agents should then handle assigned work, direct messages, explicit operator requests, and safe read-only/logging coordination instead of opening new patrol tasks. This is separate from `plannerOversightLevel`: disabling heartbeat patrol does **not** disable stuck-task observation, steering, retry, or targeted-fix recovery for tasks already in flight. No-task heartbeats resolve this value from the project default workflow, falling back to `builtin:coding` when no default workflow is set. | | `plannerHeartbeatPatrolEnabled` | `true` | Workflow-native idle-heartbeat patrol switch (FN-7963). `true` preserves the existing no-task heartbeat/triage guidance that lets idle agents scan for gaps and create focused follow-up tasks. Set to `false` to remove proactive patrol task-creation guidance from idle/no-task heartbeat prompts; agents should then handle assigned work, direct messages, explicit operator requests, and safe read-only/logging coordination instead of opening new patrol tasks. This is separate from `plannerOversightLevel`: disabling heartbeat patrol does **not** disable stuck-task observation, steering, retry, or targeted-fix recovery for tasks already in flight. No-task heartbeats resolve this value from the project default workflow, falling back to `builtin:coding` when no default workflow is set. |
| `memoryConsolidationEnabled` | `true` | Workflow-native Memory Keeper switch (FN-8932), shown in the Workflow Settings Panel's **Oversight** group. No-task heartbeats resolve it from the project default workflow (or `builtin:coding` fallback); set it to `false` to stop deterministic knowledge-graph, recall, and recall-to-graph cross-reference consolidation without deleting the agent or stored memory. `knowledgeGraphDir` remains a separate project setting owned by the knowledge-graph feature and is only read by this tick. |
When `triageProactiveSubtaskSplittingEnabled` is `true` (the default), triage may proactively replace a large task with 2-5 child tasks when the size, step-count, package breadth, file-scope, or remediation-batch signals justify the coordination overhead. When it is `false`, those automatic oversized-task signals are advisory only for writing a realistic single-task spec; triage must not split solely because the task is large. The per-task `breakIntoSubtasks: true` flag is separate and remains mandatory: if a user explicitly asks for subtask breakdown, triage still evaluates and creates child tasks when the work is meaningfully decomposable. When `triageProactiveSubtaskSplittingEnabled` is `true` (the default), triage may proactively replace a large task with 2-5 child tasks when the size, step-count, package breadth, file-scope, or remediation-batch signals justify the coordination overhead. When it is `false`, those automatic oversized-task signals are advisory only for writing a realistic single-task spec; triage must not split solely because the task is large. The per-task `breakIntoSubtasks: true` flag is separate and remains mandatory: if a user explicitly asks for subtask breakdown, triage still evaluates and creates child tasks when the work is meaningfully decomposable.

View File

@@ -832,6 +832,6 @@ Revision listing defaults to 100 rows and clamps `limit` to 1–500. The API acc
### `project.memory_recall_records` ### `project.memory_recall_records`
Project-scoped structured recall records for durable decisions, preferences, and solutions. The table uses the composite `(project_id, id)` key, row-level security, created-at indexes, and a named `(project_id, kind, content_hash)` exact-hash backstop. Project-scoped structured recall records for durable decisions, preferences, and solutions. The table uses the composite `(project_id, id)` key, row-level security, created-at indexes, and a named `(project_id, kind, content_hash)` exact-hash backstop. `graph_node_ids` stores graph cross-references; Memory Keeper merges new IDs under a per-record advisory transaction lock, so identifiers only grow and an unchanged union does not update the row.
- Knowledge-graph artifact: `<rootDir>/.fusion-knowledge/graph/` (`nodes.json`, `edges.json`, and `manifest.json`). This is deliberately outside ignored `.fusion` and may be committed at the operator's discretion. - Knowledge-graph artifact: `<rootDir>/.fusion-knowledge/graph/` (`nodes.json`, `edges.json`, and `manifest.json`). This is deliberately outside ignored `.fusion` and may be committed at the operator's discretion.

View File

@@ -0,0 +1,22 @@
import { describe, expect, it, vi } from "vitest";
import { DEFAULT_PROJECT_SETTINGS } from "../config/settings-schema.js";
import { BUILTIN_OVERSIGHT_SETTINGS, BUILTIN_WORKFLOW_SETTINGS, MEMORY_CONSOLIDATION_ENABLED_SETTING_ID } from "../workflows/builtin-workflow-settings.js";
import { resolveEffectiveMemoryConsolidationEnabled, resolveEffectiveSettingsById, type WorkflowSettingsResolverStore } from "../workflows/workflow-settings-resolver.js";
function store(values: Record<string, unknown> = {}): WorkflowSettingsResolverStore {
return { getTaskWorkflowSelection: vi.fn(() => undefined), getWorkflowDefinition: vi.fn(async () => undefined), getWorkflowSettingValues: vi.fn(() => values), getWorkflowSettingsProjectId: vi.fn(() => "project") };
}
describe("memoryConsolidationEnabled workflow setting", () => {
it("is a default-on workflow-native oversight setting", () => {
expect(BUILTIN_OVERSIGHT_SETTINGS.find((setting) => setting.id === MEMORY_CONSOLIDATION_ENABLED_SETTING_ID)).toMatchObject({ type: "boolean", default: true, name: "Memory consolidation enabled" });
expect(BUILTIN_WORKFLOW_SETTINGS.some((setting) => setting.id === MEMORY_CONSOLIDATION_ENABLED_SETTING_ID)).toBe(true);
expect(DEFAULT_PROJECT_SETTINGS).not.toHaveProperty(MEMORY_CONSOLIDATION_ENABLED_SETTING_ID);
});
it("defaults on, honors explicit false, and treats garbage as enabled", async () => {
expect(resolveEffectiveMemoryConsolidationEnabled(await resolveEffectiveSettingsById(store(), "builtin:coding", "project"))).toBe(true);
expect(resolveEffectiveMemoryConsolidationEnabled(await resolveEffectiveSettingsById(store({ memoryConsolidationEnabled: false }), "builtin:coding", "project"))).toBe(false);
expect(resolveEffectiveMemoryConsolidationEnabled({ memoryConsolidationEnabled: "false" })).toBe(true);
});
});

View File

@@ -0,0 +1,115 @@
/*
FNXC:MemoryRecall 2026-08-11-10:15:
Cross-reference ids are merged under a per-record advisory transaction lock. The contention test
proves serialization only: PostgreSQL holds the lock until commit, so there is no stable A-only
state to observe after A settles and before B can write its union.
*/
import { afterAll, afterEach, beforeAll, beforeEach, expect, it } from "vitest";
import {
createSharedPgTaskStoreTestHarness,
pgDescribe,
type SharedPgTaskStoreHarness,
} from "../../__test-utils__/pg-test-harness.js";
import { createAsyncDataLayer, type AsyncDataLayer } from "../../postgres/data-layer.js";
import { createConnectionSetFromUrl } from "../../postgres/connection.js";
import type { ResolvedBackend } from "../../postgres/backend-resolver.js";
import {
appendRecall,
getRecallRecord,
mergeRecallGraphNodeIds,
setRecallGraphNodeIdsTestHooksForTest,
} from "../../memory/recall/recall-store.js";
pgDescribe("memory recall graph cross-references (PostgreSQL)", () => {
const h: SharedPgTaskStoreHarness = createSharedPgTaskStoreTestHarness({
prefix: "fusion_memory_recall_xref", projectId: "xref-project-a", poolMax: 2,
});
const layerA = () => Object.assign(h.layer(), { projectId: "xref-project-a" });
const layerB = () => ({ ...h.layer(), projectId: "xref-project-b" });
beforeAll(h.beforeAll);
beforeEach(h.beforeEach);
afterEach(async () => {
setRecallGraphNodeIdsTestHooksForTest(undefined);
await h.afterEach();
});
afterAll(h.afterAll);
async function record(): Promise<string> {
const appended = await appendRecall(layerA(), {
kind: "decision", content: "Cross-reference fixture", source: { origin: "manual" },
});
if (appended.status !== "created") throw new Error("Expected fixture record creation");
return appended.record.id;
}
async function independentLayer(projectId: string): Promise<AsyncDataLayer> {
const backend: ResolvedBackend = {
mode: "external", runtimeUrl: h.testUrl(), migrationUrl: h.testUrl(), migrationUrlOverridden: false,
};
const connections = await createConnectionSetFromUrl(backend, { poolMax: 1, connectTimeoutSeconds: 5 });
return createAsyncDataLayer(connections, { projectId });
}
function deferred(): { promise: Promise<void>; resolve: () => void } {
let resolve!: () => void;
const promise = new Promise<void>((resolvePromise) => { resolve = resolvePromise; });
return { promise, resolve };
}
it("merges sorted ids without updating an unchanged record", async () => {
const id = await record();
expect(await mergeRecallGraphNodeIds(layerA(), id, [" b ", "a", "a", ""])).toEqual({ status: "updated" });
const updated = await getRecallRecord(layerA(), id);
expect(updated?.graphNodeIds).toEqual(["a", "b"]);
const updatedAt = updated?.updatedAt;
expect(await mergeRecallGraphNodeIds(layerA(), id, ["b", "a", "a"])).toEqual({ status: "unchanged" });
const unchanged = await getRecallRecord(layerA(), id);
expect(unchanged?.graphNodeIds).toEqual(["a", "b"]);
expect(unchanged?.updatedAt).toBe(updatedAt);
expect(await mergeRecallGraphNodeIds(layerA(), id, ["c"])).toEqual({ status: "updated" });
expect((await getRecallRecord(layerA(), id))?.graphNodeIds).toEqual(["a", "b", "c"]);
});
it("serializes concurrent writers and retains their union", async () => {
const id = await record();
const firstLayer = await independentLayer("xref-project-a");
const secondLayer = await independentLayer("xref-project-a");
const enteredLock = deferred();
const releaseFirst = deferred();
let hooked = false;
setRecallGraphNodeIdsTestHooksForTest({ afterRowRead: async () => {
if (!hooked) {
hooked = true;
enteredLock.resolve();
await releaseFirst.promise;
}
} });
try {
const first = mergeRecallGraphNodeIds(firstLayer, id, ["a"]);
await enteredLock.promise;
let secondSettled = false;
const second = mergeRecallGraphNodeIds(secondLayer, id, ["b"]).finally(() => { secondSettled = true; });
// B cannot pass the lock while A is paused after its guarded row read.
await new Promise((resolve) => setTimeout(resolve, 30));
expect(secondSettled).toBe(false);
releaseFirst.resolve();
await expect(first).resolves.toEqual({ status: "updated" });
await expect(second).resolves.toEqual({ status: "updated" });
// Do not inspect after A alone: its commit immediately permits B to commit.
expect((await getRecallRecord(layerA(), id))?.graphNodeIds).toEqual(["a", "b"]);
} finally {
releaseFirst.resolve();
await Promise.all([firstLayer.close(), secondLayer.close()]);
}
});
it("returns missing for missing and cross-project records", async () => {
const id = await record();
await expect(mergeRecallGraphNodeIds(layerA(), "missing", ["a"])).resolves.toEqual({ status: "missing" });
await expect(mergeRecallGraphNodeIds(layerB(), id, ["a"])).resolves.toEqual({ status: "missing" });
});
});

View File

@@ -0,0 +1,31 @@
import { describe, expect, it, vi } from "vitest";
import { AgentStore } from "../agent-store.js";
import { BUILTIN_MEMORY_AGENT_FALLBACK_NAME, BUILTIN_MEMORY_AGENT_NAME, BUILTIN_MEMORY_AGENT_PROVENANCE_KEY } from "../memory-agent-defaults.js";
import type { Agent } from "../../types/agents/agents.js";
const agent = (name: string, metadata: Record<string, unknown> = {}): Agent => ({ id: name.toLowerCase().replaceAll(/[^a-z]/g, ""), name, role: "custom", roles: ["custom"], state: "idle", createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", metadata, runtimeConfig: { enabled: true } } as Agent);
function fakeStore(agents: Agent[]) {
const store = new AgentStore({ rootDir: process.cwd() }); const self = store as unknown as Record<string, unknown>;
self.listAgents = vi.fn(async () => agents);
self.findAgentByName = vi.fn(async (name: string) => agents.find((item) => item.name === name) ?? null);
self.createAgent = vi.fn(async (input: { name: string }) => { const created = agent(input.name, { [BUILTIN_MEMORY_AGENT_PROVENANCE_KEY]: true }); agents.push(created); return created; });
self.writeAgent = vi.fn(async () => undefined);
return store as AgentStore & { createAgent: ReturnType<typeof vi.fn>; writeAgent: ReturnType<typeof vi.fn> };
}
describe("Memory Keeper provisioning", () => {
it("creates exactly one custom, heartbeat-enabled owner", async () => {
const store = fakeStore([]); const first = await store.provisionBuiltinMemoryAgent(); const second = await store.provisionBuiltinMemoryAgent();
expect(first?.id).toBe(second?.id); expect(store.createAgent).toHaveBeenCalledTimes(1);
expect(store.createAgent).toHaveBeenCalledWith(expect.objectContaining({ name: BUILTIN_MEMORY_AGENT_NAME, roles: ["custom"], runtimeConfig: expect.objectContaining({ enabled: true, autoClaimRelevantTasks: false, heartbeatIntervalMs: 3_600_000 }) }), undefined);
});
it("does not adopt an operator agent with the canonical name", async () => {
const operator = agent(BUILTIN_MEMORY_AGENT_NAME); const store = fakeStore([operator]);
await store.provisionBuiltinMemoryAgent();
expect(operator.name).toBe(BUILTIN_MEMORY_AGENT_NAME); expect(store.createAgent).toHaveBeenCalledWith(expect.objectContaining({ name: BUILTIN_MEMORY_AGENT_FALLBACK_NAME }), undefined);
});
it("degrades safely when both reserved names are occupied", async () => {
const store = fakeStore([agent(BUILTIN_MEMORY_AGENT_NAME), agent(BUILTIN_MEMORY_AGENT_FALLBACK_NAME)]);
await expect(store.provisionBuiltinMemoryAgent()).resolves.toBeNull(); expect(store.createAgent).not.toHaveBeenCalled();
});
});

View File

@@ -110,6 +110,12 @@ import {
BUILTIN_WORKFLOW_ROLE_AGENT_DEFAULT_LIST, BUILTIN_WORKFLOW_ROLE_AGENT_DEFAULT_LIST,
type BuiltinWorkflowRole, type BuiltinWorkflowRole,
} from "./workflow-role-agent-defaults.js"; } from "./workflow-role-agent-defaults.js";
import {
BUILTIN_MEMORY_AGENT_DEFAULT,
BUILTIN_MEMORY_AGENT_FALLBACK_NAME,
BUILTIN_MEMORY_AGENT_NAME,
BUILTIN_MEMORY_AGENT_PROVENANCE_KEY,
} from "./memory-agent-defaults.js";
const agentStoreLog = createLogger("agent-store"); const agentStoreLog = createLogger("agent-store");
@@ -558,6 +564,12 @@ export class AgentStore extends EventEmitter {
still invoke them independently of the heartbeat scheduler. still invoke them independently of the heartbeat scheduler.
*/ */
await this.provisionBuiltinWorkflowRoleAgents(); await this.provisionBuiltinWorkflowRoleAgents();
// Memory upkeep is optional; its collision-safe provisioning must never prevent startup.
try {
await this.provisionBuiltinMemoryAgent();
} catch (error) {
agentStoreLog.warn(`Unable to provision built-in memory agent: ${error instanceof Error ? error.message : String(error)}`);
}
} }
/** /**
@@ -2139,6 +2151,79 @@ export class AgentStore extends EventEmitter {
return agents; return agents;
} }
/*
FNXC:MemoryAgent 2026-08-11-09:41:
Memory Keeper is custom because it is not a workflow-stage principal. Unlike the four routed
owners its heartbeat is enabled, while auto-claim remains off because it only maintains memory.
This runs during init, where createAgent name collisions would abort startup; preflight probes,
a fallback name, a null degraded result, and the init-side catch keep an operator's same-named
agent untouched and the project runnable.
*/
async provisionBuiltinMemoryAgent(): Promise<Agent | null> {
const provision = async (executor?: QueryHandle): Promise<Agent | null> => {
const existing = await this.listAgents({ includeEphemeral: true }, executor);
const owners = existing
.filter((agent) => agent.metadata?.[BUILTIN_MEMORY_AGENT_PROVENANCE_KEY] === true)
.sort(compareBuiltinWorkflowOwnerAge);
let owner = owners[0];
for (const loser of owners.slice(1)) {
const { [BUILTIN_MEMORY_AGENT_PROVENANCE_KEY]: _builtInMemoryAgent, ...metadata } = loser.metadata ?? {};
await this.writeAgent({ ...loser, metadata, updatedAt: new Date().toISOString() }, executor);
}
if (!owner) {
/*
FNXC:MemoryAgent 2026-08-11-10:17:
Probe through the provisioning transaction executor. Opening a second pooled query while the
startup advisory lock is held can exhaust a small pool behind concurrent startup callers.
*/
const canonicalTaken = (await this.findAgentByName(BUILTIN_MEMORY_AGENT_NAME, executor)) !== null;
const name = canonicalTaken ? BUILTIN_MEMORY_AGENT_FALLBACK_NAME : BUILTIN_MEMORY_AGENT_NAME;
if (canonicalTaken && (await this.findAgentByName(BUILTIN_MEMORY_AGENT_FALLBACK_NAME, executor)) !== null) {
agentStoreLog.warn(`Built-in memory agent not provisioned: both "${BUILTIN_MEMORY_AGENT_NAME}" and "${BUILTIN_MEMORY_AGENT_FALLBACK_NAME}" are in use`);
return null;
}
try {
owner = await this.createAgent({
name,
roles: [...BUILTIN_MEMORY_AGENT_DEFAULT.roles],
title: BUILTIN_MEMORY_AGENT_DEFAULT.title,
metadata: { [BUILTIN_MEMORY_AGENT_PROVENANCE_KEY]: true },
runtimeConfig: { enabled: true, autoClaimRelevantTasks: false, heartbeatIntervalMs: 3_600_000 },
instructionsText: BUILTIN_MEMORY_AGENT_DEFAULT.instructionsText,
soul: BUILTIN_MEMORY_AGENT_DEFAULT.soul,
bundleConfig: { ...BUILTIN_MEMORY_AGENT_DEFAULT.bundleConfig, files: [...BUILTIN_MEMORY_AGENT_DEFAULT.bundleConfig.files] },
}, executor);
} catch (error) {
if (error instanceof Error && error.message.includes("already exists")) {
agentStoreLog.warn(`Built-in memory agent not provisioned because its candidate name was claimed during creation`);
return null;
}
throw error;
}
return owner;
}
const metadata = { ...(owner.metadata ?? {}), [BUILTIN_MEMORY_AGENT_PROVENANCE_KEY]: true };
const runtimeConfig = { ...(owner.runtimeConfig ?? {}), enabled: true, autoClaimRelevantTasks: false, heartbeatIntervalMs: 3_600_000 };
const updates: Partial<Agent> = {
roles: [...BUILTIN_MEMORY_AGENT_DEFAULT.roles],
role: "custom",
title: owner.title ?? BUILTIN_MEMORY_AGENT_DEFAULT.title,
metadata,
runtimeConfig,
instructionsText: owner.instructionsText?.trim() ? owner.instructionsText : BUILTIN_MEMORY_AGENT_DEFAULT.instructionsText,
soul: owner.soul?.trim() ? owner.soul : BUILTIN_MEMORY_AGENT_DEFAULT.soul,
};
owner = { ...owner, ...updates, updatedAt: new Date().toISOString() };
await this.writeAgent(owner, executor);
return owner;
};
if (!this.asyncLayer) return provision();
return this.asyncLayer.transactionImmediate(async (tx) => {
await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtext(${this.backendProjectId}), hashtext('builtin-memory-agent-provisioning'))`);
return provision(tx);
});
}
/** /**
* Create an API key for an agent. * Create an API key for an agent.
* Persists only the SHA-256 token hash; plaintext token is returned once. * Persists only the SHA-256 token hash; plaintext token is returned once.

View File

@@ -0,0 +1,30 @@
import type { InstructionsBundleConfig } from "../types.js";
import { BUILTIN_WORKFLOW_AGENT_BUNDLE_CONFIG } from "./workflow-role-agent-defaults.js";
export const BUILTIN_MEMORY_AGENT_PROVENANCE_KEY = "builtInMemoryAgent";
export const BUILTIN_MEMORY_AGENT_NAME = "Memory Keeper";
export const BUILTIN_MEMORY_AGENT_FALLBACK_NAME = "Memory Keeper (built-in)";
export interface MemoryAgentDefault {
readonly name: string;
readonly title: string;
readonly roles: readonly ["custom"];
readonly instructionsText: string;
readonly soul: string;
readonly bundleConfig: Readonly<InstructionsBundleConfig>;
}
/*
FNXC:MemoryAgent 2026-08-11-09:41:
FN-8932 adds a durable owner for deterministic memory upkeep, not another workflow-stage
principal. AgentCapability is a closed union, so this owner uses the existing custom role and
performs no product decisions or LLM work.
*/
export const BUILTIN_MEMORY_AGENT_DEFAULT: MemoryAgentDefault = Object.freeze({
name: BUILTIN_MEMORY_AGENT_NAME,
title: "Built-in memory consolidation owner",
roles: ["custom"] as const,
instructionsText: "You are the Memory Keeper. Perform deterministic memory upkeep only: refresh the knowledge graph, consolidate already-extracted rationale into recall, and maintain recall graph references. Make no product decisions, do not claim board work, and do not invoke an LLM.",
soul: "Steady, conservative, and mechanical. You preserve durable project context through deterministic maintenance without inventing conclusions.",
bundleConfig: BUILTIN_WORKFLOW_AGENT_BUNDLE_CONFIG,
});

View File

@@ -297,6 +297,7 @@ export {
DEFAULT_PLANNING_TIMEOUT_MS, DEFAULT_PLANNING_TIMEOUT_MS,
DEFAULT_PLAN_REVIEW_REPLAN_CAP, DEFAULT_PLAN_REVIEW_REPLAN_CAP,
PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID, PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID,
MEMORY_CONSOLIDATION_ENABLED_SETTING_ID,
renderTriagePolicyPlaceholders, renderTriagePolicyPlaceholders,
} from "./workflows/builtin-workflow-settings.js"; } from "./workflows/builtin-workflow-settings.js";
export { export {
@@ -550,6 +551,7 @@ export {
resolveOptionalReviewRevisionBudget, resolveOptionalReviewRevisionBudget,
resolveEffectivePlannerOversightLevel, resolveEffectivePlannerOversightLevel,
resolveEffectivePlannerHeartbeatPatrolEnabled, resolveEffectivePlannerHeartbeatPatrolEnabled,
resolveEffectiveMemoryConsolidationEnabled,
PLAN_REVIEW_MAX_REVISIONS_SETTING_ID, PLAN_REVIEW_MAX_REVISIONS_SETTING_ID,
CODE_REVIEW_MAX_REVISIONS_SETTING_ID, CODE_REVIEW_MAX_REVISIONS_SETTING_ID,
PLAN_REVIEW_REPLAN_CAP_SETTING_ID, PLAN_REVIEW_REPLAN_CAP_SETTING_ID,

View File

@@ -46,6 +46,13 @@ export {
BUILTIN_WORKFLOW_ROLE_AGENT_DEFAULT_LIST, BUILTIN_WORKFLOW_ROLE_AGENT_DEFAULT_LIST,
} from "./agents/workflow-role-agent-defaults.js"; } from "./agents/workflow-role-agent-defaults.js";
export type { BuiltinWorkflowRole, WorkflowRoleAgentDefault } from "./agents/workflow-role-agent-defaults.js"; export type { BuiltinWorkflowRole, WorkflowRoleAgentDefault } from "./agents/workflow-role-agent-defaults.js";
export {
BUILTIN_MEMORY_AGENT_DEFAULT,
BUILTIN_MEMORY_AGENT_FALLBACK_NAME,
BUILTIN_MEMORY_AGENT_NAME,
BUILTIN_MEMORY_AGENT_PROVENANCE_KEY,
} from "./agents/memory-agent-defaults.js";
export type { MemoryAgentDefault } from "./agents/memory-agent-defaults.js";
export { export {
resolveEntryPointBranchAssignment, resolveEntryPointBranchAssignment,
sanitizeBranchSegment, sanitizeBranchSegment,
@@ -343,6 +350,7 @@ export {
DEFAULT_PLAN_REVIEW_REPLAN_CAP, DEFAULT_PLAN_REVIEW_REPLAN_CAP,
DEFAULT_PLANNER_OVERSEER_EXECUTOR_STUCK_AFTER_MS, DEFAULT_PLANNER_OVERSEER_EXECUTOR_STUCK_AFTER_MS,
PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID, PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID,
MEMORY_CONSOLIDATION_ENABLED_SETTING_ID,
renderTriagePolicyPlaceholders, renderTriagePolicyPlaceholders,
} from "./workflows/builtin-workflow-settings.js"; } from "./workflows/builtin-workflow-settings.js";
export { export {
@@ -643,6 +651,7 @@ export {
resolveOptionalReviewRevisionBudget, resolveOptionalReviewRevisionBudget,
resolveEffectivePlannerOversightLevel, resolveEffectivePlannerOversightLevel,
resolveEffectivePlannerHeartbeatPatrolEnabled, resolveEffectivePlannerHeartbeatPatrolEnabled,
resolveEffectiveMemoryConsolidationEnabled,
PLAN_REVIEW_MAX_REVISIONS_SETTING_ID, PLAN_REVIEW_MAX_REVISIONS_SETTING_ID,
CODE_REVIEW_MAX_REVISIONS_SETTING_ID, CODE_REVIEW_MAX_REVISIONS_SETTING_ID,
PLAN_REVIEW_REPLAN_CAP_SETTING_ID, PLAN_REVIEW_REPLAN_CAP_SETTING_ID,

View File

@@ -6,5 +6,6 @@ export {
deleteRecallRecord, deleteRecallRecord,
getRecallRecord, getRecallRecord,
listRecall, listRecall,
mergeRecallGraphNodeIds,
searchRecall, searchRecall,
} from "./recall-store.js"; } from "./recall-store.js";

View File

@@ -12,6 +12,9 @@ outside that window, so the database exact-hash constraint remains a required ba
export function normalizeRecallContent(content: string): string { export function normalizeRecallContent(content: string): string {
return content.trim().toLowerCase().replace(/\s+/g, " ").replace(/[.!?,;:]+$/g, ""); return content.trim().toLowerCase().replace(/\s+/g, " ").replace(/[.!?,;:]+$/g, "");
} }
/** Per-record cross-reference locks avoid table-wide serialization and never contend with kind dedup locks. */
export const recallGraphNodeIdsLockKey = (projectId: string, recordId: string): string => `fusion:memory-recall-xref:${projectId}:${recordId}`;
export function recallContentHash(kind: RecallKind, content: string): string { export function recallContentHash(kind: RecallKind, content: string): string {
return createHash("sha256").update(`${kind}\0${normalizeRecallContent(content)}`).digest("hex"); return createHash("sha256").update(`${kind}\0${normalizeRecallContent(content)}`).digest("hex");
} }

View File

@@ -3,7 +3,7 @@ import { and, desc, eq, inArray, sql } from "drizzle-orm";
import type { AsyncDataLayer } from "../../postgres/data-layer.js"; import type { AsyncDataLayer } from "../../postgres/data-layer.js";
import { project } from "../../postgres/schema/index.js"; import { project } from "../../postgres/schema/index.js";
import { MemoryBackendError } from "../memory-backend.js"; import { MemoryBackendError } from "../memory-backend.js";
import { RECALL_DEDUP_CANDIDATE_LIMIT, classifyRecallDuplicate, recallContentHash, recallDedupLockKey } from "./recall-dedup.js"; import { RECALL_DEDUP_CANDIDATE_LIMIT, classifyRecallDuplicate, recallContentHash, recallDedupLockKey, recallGraphNodeIdsLockKey } from "./recall-dedup.js";
import { applyRecallVectorRanking, clampRecallSearchLimit, resolveRecallCapabilities, searchRecallKeyword, type RecallSearchOptions } from "./recall-search.js"; import { applyRecallVectorRanking, clampRecallSearchLimit, resolveRecallCapabilities, searchRecallKeyword, type RecallSearchOptions } from "./recall-search.js";
import type { RecallAppendInput, RecallAppendResult, RecallRecord, RecallSearchResult } from "./recall-types.js"; import type { RecallAppendInput, RecallAppendResult, RecallRecord, RecallSearchResult } from "./recall-types.js";
@@ -12,6 +12,7 @@ type RecallRow = typeof project.memoryRecallRecords.$inferSelect;
type RecallAppendTestHooks = { type RecallAppendTestHooks = {
afterCandidateRead?: () => void | Promise<void>; afterCandidateRead?: () => void | Promise<void>;
}; };
type RecallGraphNodeIdsTestHooks = { afterRowRead?: () => void | Promise<void> };
// Intentionally omitted from the recall barrel: PG tests use this to force the lock interleaving. // Intentionally omitted from the recall barrel: PG tests use this to force the lock interleaving.
let appendRecallTestHooks: RecallAppendTestHooks | undefined; let appendRecallTestHooks: RecallAppendTestHooks | undefined;
@@ -19,6 +20,12 @@ export function setRecallAppendTestHooksForTest(hooks: RecallAppendTestHooks | u
appendRecallTestHooks = hooks; appendRecallTestHooks = hooks;
} }
// Intentionally omitted from the recall barrel: PG tests use this to prove lock serialization.
let recallGraphNodeIdsTestHooks: RecallGraphNodeIdsTestHooks | undefined;
export function setRecallGraphNodeIdsTestHooksForTest(hooks: RecallGraphNodeIdsTestHooks | undefined): void {
recallGraphNodeIdsTestHooks = hooks;
}
const map = (row: RecallRow): RecallRecord => ({ ...row, kind: row.kind as RecallRecord["kind"], source: row.source as RecallRecord["source"], tags: row.tags as string[], graphNodeIds: row.graphNodeIds as string[] }); const map = (row: RecallRow): RecallRecord => ({ ...row, kind: row.kind as RecallRecord["kind"], source: row.source as RecallRecord["source"], tags: row.tags as string[], graphNodeIds: row.graphNodeIds as string[] });
function projectId(layer: AsyncDataLayer): string { if (!layer.projectId) throw new MemoryBackendError("BACKEND_UNAVAILABLE", "Recall requires AsyncDataLayer.projectId", "recall"); return layer.projectId; } function projectId(layer: AsyncDataLayer): string { if (!layer.projectId) throw new MemoryBackendError("BACKEND_UNAVAILABLE", "Recall requires AsyncDataLayer.projectId", "recall"); return layer.projectId; }
@@ -44,6 +51,32 @@ export async function appendRecall(layer: AsyncDataLayer, input: RecallAppendInp
return { status: "duplicate", duplicateOf: map(exact[0]), similarity: 1 }; return { status: "duplicate", duplicateOf: map(exact[0]), similarity: 1 };
}); });
} }
/*
FNXC:MemoryRecall 2026-08-11-09:41:
FN-8922 reserved graph_node_ids for part 4. A per-record advisory lock keeps the read-union-write
atomic across processes: ids only grow, unrelated records and kind-dedup do not serialize, and an
equal union performs no UPDATE so unchanged consolidation ticks do not churn updated_at. Callers
must not move the read outside this transaction to fabricate an interleaving; lock holders serialize
until commit and the second writer unions the committed first value.
*/
export async function mergeRecallGraphNodeIds(layer: AsyncDataLayer, recordId: string, addGraphNodeIds: string[]): Promise<{ status: "updated" | "unchanged" | "missing" }> {
const pid = projectId(layer);
const table = project.memoryRecallRecords;
const normalize = (ids: readonly string[]) => [...new Set(ids.map((id) => id.trim()).filter(Boolean))].sort();
return layer.transactionImmediate(async (tx) => {
await tx.execute(sql`SELECT pg_advisory_xact_lock(hashtext(${recallGraphNodeIdsLockKey(pid, recordId)}))`);
const rows = await tx.select().from(table).where(and(eq(table.projectId, pid), eq(table.id, recordId))).limit(1);
await recallGraphNodeIdsTestHooks?.afterRowRead?.();
const row = rows[0];
if (!row) return { status: "missing" };
const stored = normalize(row.graphNodeIds as string[]);
const next = normalize([...stored, ...addGraphNodeIds]);
if (stored.length === next.length && stored.every((value, index) => value === next[index])) return { status: "unchanged" };
await tx.update(table).set({ graphNodeIds: next, updatedAt: new Date().toISOString() }).where(and(eq(table.projectId, pid), eq(table.id, recordId)));
return { status: "updated" };
});
}
export async function getRecallRecord(layer: AsyncDataLayer, id: string): Promise<RecallRecord | null> { const pid=projectId(layer); const row=await layer.db.select().from(project.memoryRecallRecords).where(and(eq(project.memoryRecallRecords.projectId,pid),eq(project.memoryRecallRecords.id,id))).limit(1); return row[0] ? map(row[0]) : null; } export async function getRecallRecord(layer: AsyncDataLayer, id: string): Promise<RecallRecord | null> { const pid=projectId(layer); const row=await layer.db.select().from(project.memoryRecallRecords).where(and(eq(project.memoryRecallRecords.projectId,pid),eq(project.memoryRecallRecords.id,id))).limit(1); return row[0] ? map(row[0]) : null; }
export async function listRecall(layer: AsyncDataLayer, options?: { kinds?: RecallRecord["kind"][]; limit?: number | null }): Promise<RecallRecord[]> { const pid=projectId(layer), t=project.memoryRecallRecords, limit=clampRecallSearchLimit(options?.limit); const rows=await layer.db.select().from(t).where(and(eq(t.projectId,pid), options?.kinds?.length ? inArray(t.kind, options.kinds) : undefined)).orderBy(desc(t.createdAt)).limit(limit); return rows.map(map); } export async function listRecall(layer: AsyncDataLayer, options?: { kinds?: RecallRecord["kind"][]; limit?: number | null }): Promise<RecallRecord[]> { const pid=projectId(layer), t=project.memoryRecallRecords, limit=clampRecallSearchLimit(options?.limit); const rows=await layer.db.select().from(t).where(and(eq(t.projectId,pid), options?.kinds?.length ? inArray(t.kind, options.kinds) : undefined)).orderBy(desc(t.createdAt)).limit(limit); return rows.map(map); }
export async function deleteRecallRecord(layer: AsyncDataLayer, id: string): Promise<boolean> { const pid=projectId(layer); const rows=await layer.db.delete(project.memoryRecallRecords).where(and(eq(project.memoryRecallRecords.projectId,pid),eq(project.memoryRecallRecords.id,id))).returning({ id: project.memoryRecallRecords.id }); return Boolean(rows[0]); } export async function deleteRecallRecord(layer: AsyncDataLayer, id: string): Promise<boolean> { const pid=projectId(layer); const rows=await layer.db.delete(project.memoryRecallRecords).where(and(eq(project.memoryRecallRecords.projectId,pid),eq(project.memoryRecallRecords.id,id))).returning({ id: project.memoryRecallRecords.id }); return Boolean(rows[0]); }

View File

@@ -682,8 +682,22 @@ export const DEFAULT_MAX_POST_REVIEW_FIXES = 10;
export const DEFAULT_PLANNER_OVERSEER_EXECUTOR_STUCK_AFTER_MS = 2 * 60 * 60 * 1000; export const DEFAULT_PLANNER_OVERSEER_EXECUTOR_STUCK_AFTER_MS = 2 * 60 * 60 * 1000;
export const PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID = "plannerHeartbeatPatrolEnabled"; export const PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID = "plannerHeartbeatPatrolEnabled";
/*
FNXC:MemoryAgent 2026-08-11-09:41:
Memory consolidation resolves through the default workflow on a no-task heartbeat, so this is
workflow-native like patrol rather than a ProjectSettings key. knowledgeGraphDir deliberately stays
project-scoped because FN-8921 owns the artifact location.
*/
export const MEMORY_CONSOLIDATION_ENABLED_SETTING_ID = "memoryConsolidationEnabled";
export const BUILTIN_OVERSIGHT_SETTINGS: WorkflowSettingDefinition[] = [ export const BUILTIN_OVERSIGHT_SETTINGS: WorkflowSettingDefinition[] = [
{
id: MEMORY_CONSOLIDATION_ENABLED_SETTING_ID,
name: "Memory consolidation enabled",
type: "boolean",
default: true,
description: "Enable the memory agent's deterministic knowledge-graph and recall consolidation tick. Disabling stops consolidation without deleting the agent or stored memory.",
},
{ {
id: "plannerOversightLevel", id: "plannerOversightLevel",
name: "Planner oversight level", name: "Planner oversight level",

View File

@@ -33,7 +33,7 @@ import {
type WorkflowIrResolverStore, type WorkflowIrResolverStore,
} from "./workflow-ir-resolver.js"; } from "./workflow-ir-resolver.js";
import { resolveEffectiveSettingValues, findOrphanedSettingValues } from "./workflow-settings.js"; import { resolveEffectiveSettingValues, findOrphanedSettingValues } from "./workflow-settings.js";
import { BUILTIN_WORKFLOW_SETTINGS, PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID } from "./builtin-workflow-settings.js"; import { BUILTIN_WORKFLOW_SETTINGS, MEMORY_CONSOLIDATION_ENABLED_SETTING_ID, PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID } from "./builtin-workflow-settings.js";
import type { WorkflowSettingDefinition, WorkflowIr, WorkflowOptionalGroupConfig } from "./workflow-ir-types.js"; import type { WorkflowSettingDefinition, WorkflowIr, WorkflowOptionalGroupConfig } from "./workflow-ir-types.js";
import { PLANNER_OVERSIGHT_LEVELS, DEFAULT_PLANNER_OVERSIGHT_LEVEL, type PlannerOversightLevel } from "../types.js"; import { PLANNER_OVERSIGHT_LEVELS, DEFAULT_PLANNER_OVERSIGHT_LEVEL, type PlannerOversightLevel } from "../types.js";
@@ -392,3 +392,13 @@ export function resolveEffectivePlannerHeartbeatPatrolEnabled(
: workflowEffective; : workflowEffective;
return value !== false; return value !== false;
} }
/** Normalize the workflow-native Memory Keeper switch; only explicit false disables upkeep. */
export function resolveEffectiveMemoryConsolidationEnabled(
workflowEffective: Record<string, unknown> | boolean | null | undefined,
): boolean {
const value = typeof workflowEffective === "object" && workflowEffective !== null
? workflowEffective[MEMORY_CONSOLIDATION_ENABLED_SETTING_ID]
: workflowEffective;
return value !== false;
}

View File

@@ -167,6 +167,14 @@ describe("WorkflowSettingsPanel — Definitions tab", () => {
}); });
describe("WorkflowSettingsPanel — Values tab", () => { describe("WorkflowSettingsPanel — Values tab", () => {
it("renders Memory consolidation enabled in Oversight with its default and override", async () => {
mockFetchValues.mockResolvedValue(payload({ effective: { memoryConsolidationEnabled: false } }));
render(<Host readOnly initial={[{ id: "memoryConsolidationEnabled", name: "Memory consolidation enabled", type: "boolean", default: true }]} />);
await waitFor(() => expect(mockFetchValues).toHaveBeenCalled());
expect(within(screen.getByTestId("wf-settings-group-oversight")).getByText("Planner Oversight")).toBeInTheDocument();
expect(screen.getByLabelText("Memory consolidation enabled")).not.toBeChecked();
});
const decls: WorkflowSettingDefinition[] = [ const decls: WorkflowSettingDefinition[] = [
{ id: "timeout-ms", name: "Timeout", type: "number", default: 1000 }, { id: "timeout-ms", name: "Timeout", type: "number", default: 1000 },
{ id: "new-sessions", name: "New sessions", type: "boolean", default: false }, { id: "new-sessions", name: "New sessions", type: "boolean", default: false },

View File

@@ -4,6 +4,12 @@ import type { WorkflowSettingDefinition } from "../../api";
import { getWorkflowSettingDisplay, groupWorkflowSettings } from "../workflow-setting-display"; import { getWorkflowSettingDisplay, groupWorkflowSettings } from "../workflow-setting-display";
describe("workflow setting display ownership", () => { describe("workflow setting display ownership", () => {
it("places the memory consolidation switch in Oversight", () => {
const setting: WorkflowSettingDefinition = { id: "memoryConsolidationEnabled", name: "ignored custom label", type: "boolean", default: true };
expect(getWorkflowSettingDisplay(setting)).toMatchObject({ group: "oversight", label: "Memory consolidation enabled" });
expect(groupWorkflowSettings([setting])).toEqual([{ group: "oversight", settings: [setting] }]);
});
it("does not classify title summarizer keys as workflow model settings", () => { it("does not classify title summarizer keys as workflow model settings", () => {
const titleSettings: WorkflowSettingDefinition[] = [ const titleSettings: WorkflowSettingDefinition[] = [
{ id: "titleSummarizerProvider", name: "Title summarizer provider", type: "string" }, { id: "titleSummarizerProvider", name: "Title summarizer provider", type: "string" },

View File

@@ -181,6 +181,11 @@ const DISPLAY: Record<string, WorkflowSettingDisplay> = {
* workflow's stored value here is the effective oversight level for every task * workflow's stored value here is the effective oversight level for every task
* under it that does not set a per-task override (FN-7515, TaskForm selector). * under it that does not set a per-task override (FN-7515, TaskForm selector).
*/ */
memoryConsolidationEnabled: {
group: "oversight",
label: "Memory consolidation enabled",
description: "Enable the Memory Keeper's deterministic knowledge-graph and recall consolidation tick.",
},
plannerOversightLevel: { plannerOversightLevel: {
group: "oversight", group: "oversight",
label: "Planner oversight level", label: "Planner oversight level",

View File

@@ -0,0 +1,92 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
import type { Agent, AgentHeartbeatRun, AgentStore, TaskStore } from "@fusion/core";
import {
HEARTBEAT_ERROR_RECOVERY_METADATA_KEY,
HEARTBEAT_ERROR_RETRY_EXHAUSTED_PAUSE_REASON,
} from "../agents/agent-heartbeat-error-recovery.js";
const memory = vi.hoisted(() => ({ resolve: vi.fn(), run: vi.fn(), ensure: vi.fn() }));
vi.mock("../memory/index.js", () => ({
resolveMemoryConsolidationPorts: memory.resolve,
MemoryConsolidationService: class { runConsolidationTick = memory.run; },
MemoryConsolidationError: class MemoryConsolidationError extends Error {
constructor(readonly stage: "graph" | "recall" | "cross-reference", message: string) { super(message); }
},
}));
vi.mock("../agents/agent-instructions.js", async () => {
const actual = await vi.importActual<typeof import("../agents/agent-instructions.js")>("../agents/agent-instructions.js");
return { ...actual, ensureDefaultHeartbeatProcedureFile: memory.ensure };
});
vi.mock("../logger.js", () => ({ heartbeatLog: { log: vi.fn(), warn: vi.fn(), error: vi.fn() }, createLogger: vi.fn(() => ({ log: vi.fn(), warn: vi.fn(), error: vi.fn() })) }));
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 } : {}) };
}
function fixture(enabled: unknown, metadata: Record<string, unknown> = {}) {
let sequence = 0; const audits: Array<Record<string, unknown>> = [];
const agent = { id: "memory", name: "Memory Keeper", role: "custom", roles: ["custom"], state: "active", createdAt: "2026-01-01T00:00:00.000Z", updatedAt: "2026-01-01T00:00:00.000Z", metadata: { builtInMemoryAgent: true, ...metadata }, runtimeConfig: { enabled: true } } as unknown as Agent;
const runs = new Map<string, AgentHeartbeatRun>();
const store = {
getAgent: vi.fn(async () => agent), getCachedAgent: vi.fn(() => agent), listAgents: vi.fn(async () => [agent]), on: vi.fn(), off: vi.fn(),
updateAgentState: vi.fn(async (_id: string, state: Agent["state"]) => { agent.state = state; }),
updateAgent: vi.fn(async (_id: string, patch: Partial<Agent>) => Object.assign(agent, patch, patch.metadata ? { metadata: patch.metadata } : {})),
startHeartbeatRun: vi.fn(async () => { const run = { id: `run-${++sequence}`, agentId: agent.id, source: "timer", startedAt: new Date().toISOString(), endedAt: null, status: "active" } as AgentHeartbeatRun; runs.set(run.id, run); return run; }),
saveRun: vi.fn(async (run: AgentHeartbeatRun) => runs.set(run.id, run)), getRunDetail: vi.fn(async (_id: string, runId: string) => runs.get(runId) ?? null), endHeartbeatRun: vi.fn(), appendRunLog: vi.fn(), getBudgetStatus: vi.fn(async () => ({ allowed: true })), getActiveHeartbeatRun: vi.fn(async () => null),
} as unknown as AgentStore;
const taskStore = { getSettings: vi.fn(async () => ({ defaultWorkflowId: "builtin:coding", heartbeatErrorRecoveryAttempts: 2 })), getWorkflowSettingsProjectId: () => "project", getTaskWorkflowSelection: vi.fn(() => undefined), getWorkflowDefinition: vi.fn(async () => undefined), getWorkflowSettingValues: vi.fn(() => ({ memoryConsolidationEnabled: enabled })), recordRunAuditEvent: vi.fn(async (event) => audits.push(event as Record<string, unknown>)), getAsyncLayer: vi.fn(() => ({ projectId: "project" })) } as unknown as TaskStore;
return { agent, store, taskStore, audits };
}
function monitor(f: ReturnType<typeof fixture>) {
return new HeartbeatMonitor({ store: f.store, taskStore: f.taskStore, rootDir: process.cwd() });
}
beforeEach(() => { memory.resolve.mockReset(); memory.run.mockReset(); memory.ensure.mockResolvedValue(undefined); });
describe("Memory Keeper heartbeat hook", () => {
it("suppresses ports and sessions when the workflow switch is disabled", async () => {
const f = fixture(false); const result = await monitor(f).executeHeartbeat({ agentId: "memory", source: "timer" });
expect(result?.status).toBe("completed"); expect(memory.resolve).not.toHaveBeenCalled(); expect(memory.run).not.toHaveBeenCalled();
expect(f.audits).toEqual([expect.objectContaining({ mutationType: "memory:consolidation-skipped", target: "memory", metadata: expect.objectContaining({ agentId: "memory", reason: "disabled" }) })]);
expect(vi.mocked(f.taskStore.getSettings)).toHaveBeenCalledTimes(1);
});
it("treats an unavailable adapter environment as a successful audited skip", async () => {
memory.resolve.mockResolvedValue({ status: "unavailable", reason: "no-data-layer" });
const f = fixture(true); const result = await monitor(f).executeHeartbeat({ agentId: "memory", source: "timer" });
expect(result?.status).toBe("completed"); expect(memory.run).not.toHaveBeenCalled();
expect(f.audits).toEqual([expect.objectContaining({ mutationType: "memory:consolidation-skipped", target: "memory", metadata: expect.objectContaining({ agentId: "memory", reason: "unavailable", unavailableReason: "no-data-layer" }) })]);
});
it("emits changed and in-progress production outcomes with their fixed audit shapes", async () => {
memory.resolve.mockResolvedValue({ status: "ready", projectId: "resolved-project", ports: {} });
memory.run.mockResolvedValueOnce(outcome(true)).mockResolvedValueOnce(outcome(false, "in-progress"));
const f = fixture(true); const first = await monitor(f).executeHeartbeat({ agentId: "memory", source: "timer" });
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.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" }) }),
]);
});
it("routes recoverable tick failures through completeRun and the shared exhaustion budget", async () => {
memory.resolve.mockResolvedValue({ status: "ready", projectId: "project", ports: {} });
memory.run.mockRejectedValue(new MemoryConsolidationError("recall", "temporary consolidation failure"));
const f = fixture(true, { [HEARTBEAT_ERROR_RECOVERY_METADATA_KEY]: { consecutiveAttempts: 2 } });
const result = await monitor(f).executeHeartbeat({ agentId: "memory", source: "timer" });
expect(result?.status).toBe("failed");
expect(f.agent.state).toBe("paused");
expect(f.agent.pauseReason).toBe(HEARTBEAT_ERROR_RETRY_EXHAUSTED_PAUSE_REASON);
expect(f.agent.metadata).toHaveProperty(HEARTBEAT_ERROR_RECOVERY_METADATA_KEY);
expect(f.audits).toEqual(expect.arrayContaining([
expect.objectContaining({ mutationType: "memory:consolidation-failed", target: "memory", metadata: expect.objectContaining({ agentId: "memory", stage: "recall", recoverable: true, priorRetryCount: 2, retryLimit: 2 }) }),
]));
const failure = f.audits.find((event) => event.mutationType === "memory:consolidation-failed")!;
expect(JSON.stringify(failure.metadata)).not.toContain("temporary consolidation failure");
});
});

View File

@@ -0,0 +1,34 @@
import { beforeEach, describe, expect, it, vi } from "vitest";
const mocks = vi.hoisted(() => ({ build: vi.fn(), append: vi.fn(), merge: vi.fn(), resolve: 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 };
});
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(); });
it.each([
["no root", { taskStore: store(), rootDir: "", agentId: "memory" }, "no-root-dir"],
["no layer", { taskStore: {}, rootDir: "/repo", agentId: "memory" }, "no-data-layer"],
["null layer", { taskStore: store({ getAsyncLayer: () => null }), rootDir: "/repo", agentId: "memory" }, "no-data-layer"],
["no project", { taskStore: store({ getAsyncLayer: () => ({}) }), rootDir: "/repo", agentId: "memory" }, "no-project-id"],
] as const)("reports %s as an unavailable environment", async (_name, deps, reason) => {
await expect(resolveMemoryConsolidationPorts(deps)).resolves.toEqual({ status: "unavailable", reason });
});
it("uses the supplied settings and forwards graph recovery reason", async () => {
const taskStore = store();
const resolved = await resolveMemoryConsolidationPorts({ taskStore, rootDir: "/repo", agentId: "memory", settings: { knowledgeGraphDir: ".graph" } as never });
expect(resolved.status).toBe("ready"); if (resolved.status !== "ready") return;
expect(taskStore.getSettings).not.toHaveBeenCalled();
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("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" });
});
});

View File

@@ -0,0 +1,19 @@
import { describe, expect, it } from "vitest";
import { readFileSync } from "node:fs";
import { resolve } from "node:path";
/** The audit declaration is the closed contract consumed by every heartbeat emission. */
describe("memory consolidation run-audit metadata contract", () => {
it("declares only fixed outcome fields and forbids prose-bearing fields", () => {
const source = readFileSync(resolve(import.meta.dirname, "../util/run-audit.ts"), "utf8");
const completed = source.indexOf("memory:consolidation-completed");
const block = source.slice(Math.max(0, completed - 1000), completed + 1000);
expect(block).toContain("memory:consolidation-completed");
expect(block).toContain("memory:consolidation-skipped");
expect(block).toContain("memory:consolidation-failed");
expect(block).toContain("closed graphRecoveryReason enum");
expect(block).toContain("skipped reasons and failure stages are closed enums");
expect(block).toMatch(/No memory content, paths,[\s\S]*node ids, error class\/message/);
expect(block).not.toMatch(/errorClass:\s*string/);
});
});

View File

@@ -0,0 +1,62 @@
import { describe, expect, it, vi } from "vitest";
import { fileNodeId, type GraphNode, type RecallRecord } from "@fusion/core";
import { MemoryConsolidationError, MemoryConsolidationService, type MemoryConsolidationPorts } from "../memory/memory-consolidation.js";
import { deriveRecallMaterial } from "../memory/memory-consolidation-material.js";
const rationale = (id: string, ownerPath = "src/a.ts"): GraphNode => ({
id, kind: "rationale", name: "Memory", owner: "file", ownerPath,
source: { path: ownerPath, line: 1, column: 1 },
attributes: { fnxcArea: "Memory", fnxcStamp: "2026-08-11-09:41", fnxcText: `decision ${id}` },
});
const record = (id: string): RecallRecord => ({ id, projectId: "project", kind: "decision", content: id, contentHash: id, source: { origin: "other" }, tags: [], graphNodeIds: [], createdAt: "", updatedAt: "" });
const graph = (changed: boolean, nodes = [rationale("rationale:src/a.ts#Memory@x~0")]) => ({ rationaleNodes: nodes, nodeCount: nodes.length, edgeCount: 0, changed, recoveryReason: null, stats: { parsedFiles: 1, reusedFiles: 0, prunedFiles: 0 } });
function ports(overrides: Partial<MemoryConsolidationPorts> = {}): MemoryConsolidationPorts {
return {
refreshGraph: vi.fn(async () => graph(true)),
appendRecall: vi.fn(async () => ({ status: "created" as const, record: record("rec-1") })),
mergeRecallGraphNodeIds: vi.fn(async () => ({ status: "unchanged" as const })),
clock: vi.fn().mockImplementationOnce(() => 1).mockReturnValue(2),
...overrides,
};
}
describe("MemoryConsolidationService", () => {
it("is silent on an unchanged second tick while still using dedup probes", async () => {
const p = ports(); const service = new MemoryConsolidationService(p);
expect((await service.runConsolidationTick({ agentId: "memory", projectId: "project" })).changed).toBe(true);
vi.mocked(p.refreshGraph).mockResolvedValue(graph(false));
vi.mocked(p.appendRecall).mockResolvedValue({ status: "duplicate", duplicateOf: record("rec-1"), similarity: 1 });
const unchanged = await service.runConsolidationTick({ agentId: "memory", projectId: "project" });
expect(unchanged).toMatchObject({ changed: false, recallCreated: 0, recallDuplicate: 1, crossRefUpdated: 0 });
expect(p.mergeRecallGraphNodeIds).toHaveBeenLastCalledWith("rec-1", expect.any(Array));
});
it("uses normalized graph ids and aggregates duplicate targets into one merge", async () => {
const nodes = [rationale("rationale:one", "src\\a.ts"), rationale("rationale:two", "src/b.ts")];
const p = ports({ refreshGraph: vi.fn(async () => graph(false, nodes)), appendRecall: vi.fn(async () => ({ status: "duplicate", duplicateOf: record("rec-1"), similarity: 1 })), mergeRecallGraphNodeIds: vi.fn(async () => ({ status: "updated" as const })) });
await new MemoryConsolidationService(p).runConsolidationTick({ agentId: "memory", projectId: "project" });
expect(p.mergeRecallGraphNodeIds).toHaveBeenCalledTimes(1);
expect(p.mergeRecallGraphNodeIds).toHaveBeenCalledWith("rec-1", ["rationale:one", "rationale:two", fileNodeId("src/a.ts"), fileNodeId("src/b.ts")].sort());
});
it("claims synchronous callers and releases after failures", async () => {
let reject!: (error: Error) => void;
const pending = new Promise<never>((_, r) => { reject = r; });
const p = ports({ refreshGraph: vi.fn(() => pending) }); const service = new MemoryConsolidationService(p);
const first = service.runConsolidationTick({ agentId: "memory", projectId: "project" });
const second = service.runConsolidationTick({ agentId: "memory", projectId: "project" });
expect(await second).toMatchObject({ skipped: "in-progress", changed: false });
reject(new Error("graph unavailable"));
await expect(first).rejects.toMatchObject({ stage: "graph" } satisfies Partial<MemoryConsolidationError>);
vi.mocked(p.refreshGraph).mockResolvedValue(graph(false));
vi.mocked(p.appendRecall).mockResolvedValue({ status: "duplicate", duplicateOf: record("rec-1"), similarity: 1 });
await expect(service.runConsolidationTick({ agentId: "memory", projectId: "project" })).resolves.toMatchObject({ changed: false });
});
it("derives deterministic material and skips blank rationale text", () => {
const node = rationale("rationale:one"); const blank = { ...rationale("rationale:blank"), attributes: {} };
expect(deriveRecallMaterial([node, blank], "memory")).toEqual(deriveRecallMaterial([node, blank], "memory"));
expect(deriveRecallMaterial([node, blank], "memory")).toHaveLength(1);
});
});

View File

@@ -35,6 +35,7 @@ import {
formatAssignedTasksWakeDeltaSection, formatAssignedTasksWakeDeltaSection,
resolveEffectiveSettingsById, resolveEffectiveSettingsById,
resolveEffectivePlannerHeartbeatPatrolEnabled, resolveEffectivePlannerHeartbeatPatrolEnabled,
resolveEffectiveMemoryConsolidationEnabled,
resolveReboundTarget, resolveReboundTarget,
resolveWorkflowIrForTask, resolveWorkflowIrForTask,
columnsWithFlag, columnsWithFlag,
@@ -51,6 +52,7 @@ import {
resolveAgentInstructionsWithRatings, resolveAgentInstructionsWithRatings,
buildPluginPromptSection, buildPluginPromptSection,
resolveAgentHeartbeatProcedure, resolveAgentHeartbeatProcedure,
ensureDefaultHeartbeatProcedureFile,
} from "./agents/agent-instructions.js"; } from "./agents/agent-instructions.js";
import { resolveHeartbeatPromptTemplate, resolveHeartbeatScopeDisciplineMode, selectHeartbeatProcedure } from "./agents/heartbeat-procedure-resolver.js"; import { resolveHeartbeatPromptTemplate, resolveHeartbeatScopeDisciplineMode, selectHeartbeatProcedure } from "./agents/heartbeat-procedure-resolver.js";
import { buildPromptLayers, collapsePromptLayers } from "./execution/prompt-layers.js"; import { buildPromptLayers, collapsePromptLayers } from "./execution/prompt-layers.js";
@@ -105,6 +107,7 @@ import { trimPromptMd, trimTaskDescription, trimTriggeringComments } from "./age
import { detectDeicticReference, extractAntecedentCandidates, renderAmbiguityPromptBlock, scoreReferentConfidence } from "./triage-domain/room-ambiguity.js"; import { detectDeicticReference, extractAntecedentCandidates, renderAmbiguityPromptBlock, scoreReferentConfidence } from "./triage-domain/room-ambiguity.js";
import { countActiveAgentMembers, decideRoomCoordination, detectTaskFilingIntent, renderRoomCoordinationPromptBlock } from "./triage-domain/room-coordination.js"; import { countActiveAgentMembers, decideRoomCoordination, detectTaskFilingIntent, renderRoomCoordinationPromptBlock } from "./triage-domain/room-coordination.js";
import { evaluateParkedAgentTaskLink, isParkedTaskColumn, type AgentTaskLinkExecutionProof } from "./agents/task-agent-sync.js"; import { evaluateParkedAgentTaskLink, isParkedTaskColumn, type AgentTaskLinkExecutionProof } from "./agents/task-agent-sync.js";
import { MemoryConsolidationError, MemoryConsolidationService, resolveMemoryConsolidationPorts } from "./memory/index.js";
/* /*
FNXC:WorkflowLifecycleColumns 2026-07-28-09:25 (U11 conversion): FNXC:WorkflowLifecycleColumns 2026-07-28-09:25 (U11 conversion):
@@ -163,6 +166,14 @@ async function resolveNoTaskHeartbeatPatrolEnabled(
} }
} }
async function resolveMemoryConsolidationEnabledForHeartbeat(taskStore: TaskStore, settings: Settings | undefined): Promise<boolean> {
try {
const projectId = typeof taskStore.getWorkflowSettingsProjectId === "function" ? taskStore.getWorkflowSettingsProjectId() : "default";
const effective = await resolveEffectiveSettingsById(taskStore, settings?.defaultWorkflowId || "builtin:coding", projectId);
return resolveEffectiveMemoryConsolidationEnabled(effective);
} catch { return resolveEffectiveMemoryConsolidationEnabled(undefined); }
}
interface SelfImproveServiceLike { interface SelfImproveServiceLike {
shouldRunSelfImprove(agentId: string): Promise<boolean>; shouldRunSelfImprove(agentId: string): Promise<boolean>;
getSelfImprovePrompt(agentId: string): Promise<string>; getSelfImprovePrompt(agentId: string): Promise<string>;
@@ -2372,6 +2383,50 @@ export class HeartbeatMonitor {
} }
} }
/*
FNXC:MemoryAgent 2026-08-11-09:41:
Memory Keeper runs before task/session assembly so 4a consumes no model quota. Provenance,
not its name, identifies fallback-named owners. Disabled or unavailable environments are
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>) => {
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. */
try {
await ensureDefaultHeartbeatProcedureFile(rootDir, agent.heartbeatProcedurePath ?? `.fusion/agents/${agentId}/HEARTBEAT.md`, "# Memory Keeper\n\nRun deterministic graph refresh, recall consolidation, and graph-reference merging. Merge graph references without dropping existing ids. Do not use an LLM. Unchanged inputs write nothing.");
} catch (error) {
heartbeatLog.warn(`Unable to seed Memory Keeper heartbeat procedure: ${error instanceof Error ? error.message : String(error)}`);
}
if (!await resolveMemoryConsolidationEnabledForHeartbeat(taskStore, heartbeatModelSettings)) {
await emit("memory:consolidation-skipped", { agentId, reason: "disabled" });
await this.completeRun(agentId, run.id, { status: "completed", resultJson: { reason: "memory_consolidation_disabled" }, skipStateTransition: true });
return (await this.store.getRunDetail(agentId, run.id))!;
}
const resolution = await resolveMemoryConsolidationPorts({ taskStore, rootDir, agentId, settings: heartbeatModelSettings });
if (resolution.status === "unavailable") {
await emit("memory:consolidation-skipped", { agentId, reason: "unavailable", unavailableReason: resolution.reason });
await this.completeRun(agentId, run.id, { status: "completed", resultJson: { reason: "memory_consolidation_unavailable", unavailableReason: resolution.reason }, skipStateTransition: true });
return (await this.store.getRunDetail(agentId, run.id))!;
}
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 });
}
await this.completeRun(agentId, run.id, { status: "completed", resultJson: { reason: "memory_consolidation", ...outcome }, skipStateTransition: true });
} catch (error) {
const message = error instanceof Error ? error.message : String(error);
await emit("memory:consolidation-failed", { agentId, stage: error instanceof MemoryConsolidationError ? error.stage : "unknown", recoverable: isHeartbeatErrorRecoverable({ lastError: message }), priorRetryCount: readHeartbeatErrorRetryCount(agent), retryLimit: resolveErrorRecoveryLimit(heartbeatModelSettings) });
await this.completeRun(agentId, run.id, { status: "failed", stderrExcerpt: message, errorMessage: message });
}
return (await this.store.getRunDetail(agentId, run.id))!;
}
// Check if agent has identity (used later for no-task run decisions) // Check if agent has identity (used later for no-task run decisions)
const agentHasIdentity = hasAgentIdentity(agent); const agentHasIdentity = hasAgentIdentity(agent);
const isAgentEphemeral = isEphemeralAgent(agent); const isAgentEphemeral = isEphemeralAgent(agent);

View File

@@ -1194,6 +1194,7 @@ export {
type CliAdapterDescriptor, type CliAdapterDescriptor,
} from "./cli-agent/adapters/index.js"; } from "./cli-agent/adapters/index.js";
export { installBaselineArchiveWorktreeDisposer } from "./healing/archive-worktree-disposer-install.js"; export { installBaselineArchiveWorktreeDisposer } from "./healing/archive-worktree-disposer-install.js";
export { MemoryConsolidationService, resolveMemoryConsolidationPorts } from "./memory/index.js";
// CLI Agent Executor — task ↔ session orchestration (U7). // CLI Agent Executor — task ↔ session orchestration (U7).
export { export {

View File

@@ -0,0 +1,3 @@
export * from "./memory-consolidation.js";
export * from "./memory-consolidation-material.js";
export * from "./memory-consolidation-adapters.js";

View File

@@ -0,0 +1,29 @@
import { buildKnowledgeGraph, mergeRecallGraphNodeIds, appendRecall, resolveKnowledgeGraphDir, type Settings } from "@fusion/core";
import { relative, resolve, sep } from "node:path";
import type { MemoryConsolidationPorts } from "./memory-consolidation.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 };
type Deps = { taskStore: { getAsyncLayer?: () => unknown; getSettings?: () => Promise<Settings> }; rootDir: string; agentId: string; settings?: Settings };
/* FNXC:MemoryAgent 2026-08-11-09:41: This is the only memory module binding graph I/O and project configuration. Missing collaborators are successful skips, not heartbeat failures; a blank graph directory uses FN-8921's default, while .fusion is rejected because graph artifacts are committable. */
export async function resolveMemoryConsolidationPorts(deps: Deps): Promise<MemoryConsolidationPortsResolution> {
if (!deps.rootDir) return { status: "unavailable", reason: "no-root-dir" };
if (typeof deps.taskStore.getAsyncLayer !== "function") return { status: "unavailable", reason: "no-data-layer" };
const layer = deps.taskStore.getAsyncLayer() as { projectId?: string } | null | undefined;
if (!layer) return { status: "unavailable", reason: "no-data-layer" };
if (!layer.projectId) return { status: "unavailable", reason: "no-project-id" };
let settings = deps.settings;
if (!settings && typeof deps.taskStore.getSettings === "function") { try { settings = await deps.taskStore.getSettings(); } catch { /* default directory remains valid */ } }
let graphDir: string;
try { graphDir = resolveKnowledgeGraphDir(deps.rootDir, settings?.knowledgeGraphDir); } catch { return { status: "unavailable", reason: "knowledge-graph-dir-unresolved" }; }
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" };
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 } }; },
appendRecall: (input) => appendRecall(layer as never, input),
mergeRecallGraphNodeIds: (id, ids) => mergeRecallGraphNodeIds(layer as never, id, ids),
clock: Date.now,
} };
}

View File

@@ -0,0 +1,22 @@
import { fileNodeId, recallContentHash, type GraphNode, type RecallAppendInput } from "@fusion/core";
export type MemoryConsolidationMaterial = { contentHash: string; append: RecallAppendInput; graphNodeIds: string[] };
/*
FNXC:MemoryAgent 2026-08-11-09:41:
This pure transformer may use graph types and fileNodeId but no graph I/O, settings, filesystem, or
database APIs. FNXC rationale is the deterministic 4a source; council/research capture and LLM
summaries belong to 4b, while fileNodeId preserves the graph's normalized ownership identity.
*/
export function deriveRecallMaterial(nodes: readonly GraphNode[], agentId: string): MemoryConsolidationMaterial[] {
return nodes.filter((node) => node.kind === "rationale").flatMap((node) => {
const area = node.attributes.fnxcArea ?? node.name;
const stamp = node.attributes.fnxcStamp ?? "";
const text = (node.attributes.fnxcText ?? "").trim();
if (!text) return [];
const content = `FNXC decision\nArea: ${area}\nStamp: ${stamp}\n${text}`;
const graphNodeIds = [...new Set([node.id, fileNodeId(node.ownerPath)])].sort();
const append: RecallAppendInput = { kind: "decision", content, source: { origin: "other", agentId }, tags: [area], graphNodeIds };
return [{ contentHash: recallContentHash("decision", content), append, graphNodeIds, nodeId: node.id }];
}).sort((a, b) => a.nodeId.localeCompare(b.nodeId)).map(({ nodeId: _nodeId, ...material }) => material);
}

View File

@@ -0,0 +1,41 @@
import type { GraphNode, RecallAppendInput, RecallAppendResult, RecoveryReason } from "@fusion/core";
import { deriveRecallMaterial } from "./memory-consolidation-material.js";
export type MemoryConsolidationTickInput = { agentId: string; projectId: string };
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" }>;
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"; } }
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 });
/* 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 {
constructor(private readonly ports: MemoryConsolidationPorts) {}
async runConsolidationTick(input: MemoryConsolidationTickInput): Promise<MemoryConsolidationOutcome> {
const key = `${input.agentId}:${input.projectId}`;
if (active.has(key)) return { ...empty(), skipped: "in-progress" };
active.add(key);
const started = (this.ports.clock ?? Date.now)();
try {
let graph;
try { graph = await this.ports.refreshGraph(); } catch (cause) { throw new MemoryConsolidationError("graph", cause instanceof Error ? cause.message : String(cause), { cause }); }
const byRecord = new Map<string, Set<string>>(); let created = 0; let duplicate = 0;
for (const material of deriveRecallMaterial(graph.rationaleNodes, input.agentId)) {
let result: RecallAppendResult;
try { result = await this.ports.appendRecall(material.append); } catch (cause) { throw new MemoryConsolidationError("recall", cause instanceof Error ? cause.message : String(cause), { cause }); }
const record = result.status === "created" ? result.record : result.duplicateOf;
if (result.status === "created") created++; else duplicate++;
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;
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 };
} finally { active.delete(key); }
}
}

View File

@@ -799,6 +799,16 @@ export type DatabaseMutationType =
| "reflection:skipped" | "reflection:skipped"
| "reflection:failed" | "reflection:failed"
| "reflection:captured" | "reflection:captured"
/*
* FNXC:MemoryAgent 2026-08-11-09:41:
* Memory consolidation telemetry is ids/counts/outcomes only. Completed metadata forwards the
* closed graphRecoveryReason enum so a manifest-last inconsistent-artifact rebuild loop is
* diagnosable; skipped reasons and failure stages are closed enums. No memory content, paths,
* node ids, error class/message, prompt text, or reasoning may be recorded. No-op ticks emit no row.
*/
| "memory:consolidation-completed"
| "memory:consolidation-skipped"
| "memory:consolidation-failed"
| "task:in-review-stall-deadlock-disposed" | "task:in-review-stall-deadlock-disposed"
| "task:in-review-stall-terminal-provider-error" | "task:in-review-stall-terminal-provider-error"
| "task:finalize-unproven-blocked" | "task:finalize-unproven-blocked"