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:
7
.changeset/fn-8932-memory-agent.md
Normal file
7
.changeset/fn-8932-memory-agent.md
Normal 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.
|
||||
@@ -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-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-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.
|
||||
|
||||
@@ -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.
|
||||
- **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
|
||||
|
||||
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`.
|
||||
|
||||
@@ -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". |
|
||||
| `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. |
|
||||
| `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.
|
||||
|
||||
|
||||
@@ -832,6 +832,6 @@ Revision listing defaults to 100 rows and clamps `limit` to 1–500. The API acc
|
||||
|
||||
### `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.
|
||||
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -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" });
|
||||
});
|
||||
});
|
||||
@@ -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();
|
||||
});
|
||||
});
|
||||
@@ -110,6 +110,12 @@ import {
|
||||
BUILTIN_WORKFLOW_ROLE_AGENT_DEFAULT_LIST,
|
||||
type BuiltinWorkflowRole,
|
||||
} 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");
|
||||
|
||||
@@ -558,6 +564,12 @@ export class AgentStore extends EventEmitter {
|
||||
still invoke them independently of the heartbeat scheduler.
|
||||
*/
|
||||
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;
|
||||
}
|
||||
|
||||
/*
|
||||
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.
|
||||
* Persists only the SHA-256 token hash; plaintext token is returned once.
|
||||
|
||||
30
packages/core/src/agents/memory-agent-defaults.ts
Normal file
30
packages/core/src/agents/memory-agent-defaults.ts
Normal 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,
|
||||
});
|
||||
@@ -297,6 +297,7 @@ export {
|
||||
DEFAULT_PLANNING_TIMEOUT_MS,
|
||||
DEFAULT_PLAN_REVIEW_REPLAN_CAP,
|
||||
PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID,
|
||||
MEMORY_CONSOLIDATION_ENABLED_SETTING_ID,
|
||||
renderTriagePolicyPlaceholders,
|
||||
} from "./workflows/builtin-workflow-settings.js";
|
||||
export {
|
||||
@@ -550,6 +551,7 @@ export {
|
||||
resolveOptionalReviewRevisionBudget,
|
||||
resolveEffectivePlannerOversightLevel,
|
||||
resolveEffectivePlannerHeartbeatPatrolEnabled,
|
||||
resolveEffectiveMemoryConsolidationEnabled,
|
||||
PLAN_REVIEW_MAX_REVISIONS_SETTING_ID,
|
||||
CODE_REVIEW_MAX_REVISIONS_SETTING_ID,
|
||||
PLAN_REVIEW_REPLAN_CAP_SETTING_ID,
|
||||
|
||||
@@ -46,6 +46,13 @@ export {
|
||||
BUILTIN_WORKFLOW_ROLE_AGENT_DEFAULT_LIST,
|
||||
} 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 {
|
||||
resolveEntryPointBranchAssignment,
|
||||
sanitizeBranchSegment,
|
||||
@@ -343,6 +350,7 @@ export {
|
||||
DEFAULT_PLAN_REVIEW_REPLAN_CAP,
|
||||
DEFAULT_PLANNER_OVERSEER_EXECUTOR_STUCK_AFTER_MS,
|
||||
PLANNER_HEARTBEAT_PATROL_ENABLED_SETTING_ID,
|
||||
MEMORY_CONSOLIDATION_ENABLED_SETTING_ID,
|
||||
renderTriagePolicyPlaceholders,
|
||||
} from "./workflows/builtin-workflow-settings.js";
|
||||
export {
|
||||
@@ -643,6 +651,7 @@ export {
|
||||
resolveOptionalReviewRevisionBudget,
|
||||
resolveEffectivePlannerOversightLevel,
|
||||
resolveEffectivePlannerHeartbeatPatrolEnabled,
|
||||
resolveEffectiveMemoryConsolidationEnabled,
|
||||
PLAN_REVIEW_MAX_REVISIONS_SETTING_ID,
|
||||
CODE_REVIEW_MAX_REVISIONS_SETTING_ID,
|
||||
PLAN_REVIEW_REPLAN_CAP_SETTING_ID,
|
||||
|
||||
@@ -6,5 +6,6 @@ export {
|
||||
deleteRecallRecord,
|
||||
getRecallRecord,
|
||||
listRecall,
|
||||
mergeRecallGraphNodeIds,
|
||||
searchRecall,
|
||||
} from "./recall-store.js";
|
||||
|
||||
@@ -12,6 +12,9 @@ outside that window, so the database exact-hash constraint remains a required ba
|
||||
export function normalizeRecallContent(content: string): string {
|
||||
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 {
|
||||
return createHash("sha256").update(`${kind}\0${normalizeRecallContent(content)}`).digest("hex");
|
||||
}
|
||||
|
||||
@@ -3,7 +3,7 @@ import { and, desc, eq, inArray, sql } from "drizzle-orm";
|
||||
import type { AsyncDataLayer } from "../../postgres/data-layer.js";
|
||||
import { project } from "../../postgres/schema/index.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 type { RecallAppendInput, RecallAppendResult, RecallRecord, RecallSearchResult } from "./recall-types.js";
|
||||
|
||||
@@ -12,6 +12,7 @@ type RecallRow = typeof project.memoryRecallRecords.$inferSelect;
|
||||
type RecallAppendTestHooks = {
|
||||
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.
|
||||
let appendRecallTestHooks: RecallAppendTestHooks | undefined;
|
||||
@@ -19,6 +20,12 @@ export function setRecallAppendTestHooksForTest(hooks: RecallAppendTestHooks | u
|
||||
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[] });
|
||||
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 };
|
||||
});
|
||||
}
|
||||
/*
|
||||
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 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]); }
|
||||
|
||||
@@ -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 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[] = [
|
||||
{
|
||||
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",
|
||||
name: "Planner oversight level",
|
||||
|
||||
@@ -33,7 +33,7 @@ import {
|
||||
type WorkflowIrResolverStore,
|
||||
} from "./workflow-ir-resolver.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 { PLANNER_OVERSIGHT_LEVELS, DEFAULT_PLANNER_OVERSIGHT_LEVEL, type PlannerOversightLevel } from "../types.js";
|
||||
|
||||
@@ -392,3 +392,13 @@ export function resolveEffectivePlannerHeartbeatPatrolEnabled(
|
||||
: workflowEffective;
|
||||
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;
|
||||
}
|
||||
|
||||
@@ -167,6 +167,14 @@ describe("WorkflowSettingsPanel — Definitions 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[] = [
|
||||
{ id: "timeout-ms", name: "Timeout", type: "number", default: 1000 },
|
||||
{ id: "new-sessions", name: "New sessions", type: "boolean", default: false },
|
||||
|
||||
@@ -4,6 +4,12 @@ import type { WorkflowSettingDefinition } from "../../api";
|
||||
import { getWorkflowSettingDisplay, groupWorkflowSettings } from "../workflow-setting-display";
|
||||
|
||||
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", () => {
|
||||
const titleSettings: WorkflowSettingDefinition[] = [
|
||||
{ id: "titleSummarizerProvider", name: "Title summarizer provider", type: "string" },
|
||||
|
||||
@@ -181,6 +181,11 @@ const DISPLAY: Record<string, WorkflowSettingDisplay> = {
|
||||
* 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).
|
||||
*/
|
||||
memoryConsolidationEnabled: {
|
||||
group: "oversight",
|
||||
label: "Memory consolidation enabled",
|
||||
description: "Enable the Memory Keeper's deterministic knowledge-graph and recall consolidation tick.",
|
||||
},
|
||||
plannerOversightLevel: {
|
||||
group: "oversight",
|
||||
label: "Planner oversight level",
|
||||
|
||||
@@ -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");
|
||||
});
|
||||
});
|
||||
@@ -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" });
|
||||
});
|
||||
});
|
||||
@@ -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/);
|
||||
});
|
||||
});
|
||||
@@ -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);
|
||||
});
|
||||
});
|
||||
@@ -35,6 +35,7 @@ import {
|
||||
formatAssignedTasksWakeDeltaSection,
|
||||
resolveEffectiveSettingsById,
|
||||
resolveEffectivePlannerHeartbeatPatrolEnabled,
|
||||
resolveEffectiveMemoryConsolidationEnabled,
|
||||
resolveReboundTarget,
|
||||
resolveWorkflowIrForTask,
|
||||
columnsWithFlag,
|
||||
@@ -51,6 +52,7 @@ import {
|
||||
resolveAgentInstructionsWithRatings,
|
||||
buildPluginPromptSection,
|
||||
resolveAgentHeartbeatProcedure,
|
||||
ensureDefaultHeartbeatProcedureFile,
|
||||
} from "./agents/agent-instructions.js";
|
||||
import { resolveHeartbeatPromptTemplate, resolveHeartbeatScopeDisciplineMode, selectHeartbeatProcedure } from "./agents/heartbeat-procedure-resolver.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 { countActiveAgentMembers, decideRoomCoordination, detectTaskFilingIntent, renderRoomCoordinationPromptBlock } from "./triage-domain/room-coordination.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):
|
||||
@@ -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 {
|
||||
shouldRunSelfImprove(agentId: string): Promise<boolean>;
|
||||
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)
|
||||
const agentHasIdentity = hasAgentIdentity(agent);
|
||||
const isAgentEphemeral = isEphemeralAgent(agent);
|
||||
|
||||
@@ -1194,6 +1194,7 @@ export {
|
||||
type CliAdapterDescriptor,
|
||||
} from "./cli-agent/adapters/index.js";
|
||||
export { installBaselineArchiveWorktreeDisposer } from "./healing/archive-worktree-disposer-install.js";
|
||||
export { MemoryConsolidationService, resolveMemoryConsolidationPorts } from "./memory/index.js";
|
||||
|
||||
// CLI Agent Executor — task ↔ session orchestration (U7).
|
||||
export {
|
||||
|
||||
3
packages/engine/src/memory/index.ts
Normal file
3
packages/engine/src/memory/index.ts
Normal file
@@ -0,0 +1,3 @@
|
||||
export * from "./memory-consolidation.js";
|
||||
export * from "./memory-consolidation-material.js";
|
||||
export * from "./memory-consolidation-adapters.js";
|
||||
29
packages/engine/src/memory/memory-consolidation-adapters.ts
Normal file
29
packages/engine/src/memory/memory-consolidation-adapters.ts
Normal 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,
|
||||
} };
|
||||
}
|
||||
22
packages/engine/src/memory/memory-consolidation-material.ts
Normal file
22
packages/engine/src/memory/memory-consolidation-material.ts
Normal 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);
|
||||
}
|
||||
41
packages/engine/src/memory/memory-consolidation.ts
Normal file
41
packages/engine/src/memory/memory-consolidation.ts
Normal 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); }
|
||||
}
|
||||
}
|
||||
@@ -799,6 +799,16 @@ export type DatabaseMutationType =
|
||||
| "reflection:skipped"
|
||||
| "reflection:failed"
|
||||
| "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-terminal-provider-error"
|
||||
| "task:finalize-unproven-blocked"
|
||||
|
||||
Reference in New Issue
Block a user