diff --git a/.changeset/fn-8932-memory-agent.md b/.changeset/fn-8932-memory-agent.md new file mode 100644 index 0000000000..50275e22a5 --- /dev/null +++ b/.changeset/fn-8932-memory-agent.md @@ -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. diff --git a/AGENTS.md b/AGENTS.md index 4854c17501..6fb6581e0a 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -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. diff --git a/docs/agents.md b/docs/agents.md index 9c301fba1a..99d3e741c2 100644 --- a/docs/agents.md +++ b/docs/agents.md @@ -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`. diff --git a/docs/settings-reference.md b/docs/settings-reference.md index 2a3a2ed0b1..2967f22e51 100644 --- a/docs/settings-reference.md +++ b/docs/settings-reference.md @@ -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. diff --git a/docs/storage.md b/docs/storage.md index 7fc8766c30..57d60f5bbf 100644 --- a/docs/storage.md +++ b/docs/storage.md @@ -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: `/.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. diff --git a/packages/core/src/__tests__/memoryConsolidationEnabled-default.test.ts b/packages/core/src/__tests__/memoryConsolidationEnabled-default.test.ts new file mode 100644 index 0000000000..6104cd2ba8 --- /dev/null +++ b/packages/core/src/__tests__/memoryConsolidationEnabled-default.test.ts @@ -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 = {}): 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); + }); +}); diff --git a/packages/core/src/__tests__/postgres/memory-recall-graph-cross-reference.pg.test.ts b/packages/core/src/__tests__/postgres/memory-recall-graph-cross-reference.pg.test.ts new file mode 100644 index 0000000000..15cb04ce71 --- /dev/null +++ b/packages/core/src/__tests__/postgres/memory-recall-graph-cross-reference.pg.test.ts @@ -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 { + 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 { + 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; resolve: () => void } { + let resolve!: () => void; + const promise = new Promise((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" }); + }); +}); diff --git a/packages/core/src/agents/__tests__/memory-agent-provisioning.test.ts b/packages/core/src/agents/__tests__/memory-agent-provisioning.test.ts new file mode 100644 index 0000000000..a7d5ff33cf --- /dev/null +++ b/packages/core/src/agents/__tests__/memory-agent-provisioning.test.ts @@ -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 = {}): 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; + 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; writeAgent: ReturnType }; +} + +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(); + }); +}); diff --git a/packages/core/src/agents/agent-store.ts b/packages/core/src/agents/agent-store.ts index ff88c3483d..11682ab6b2 100644 --- a/packages/core/src/agents/agent-store.ts +++ b/packages/core/src/agents/agent-store.ts @@ -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 { + const provision = async (executor?: QueryHandle): Promise => { + 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 = { + 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. diff --git a/packages/core/src/agents/memory-agent-defaults.ts b/packages/core/src/agents/memory-agent-defaults.ts new file mode 100644 index 0000000000..0d4e85eaeb --- /dev/null +++ b/packages/core/src/agents/memory-agent-defaults.ts @@ -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; +} + +/* +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, +}); diff --git a/packages/core/src/index.gate.ts b/packages/core/src/index.gate.ts index 1b59c819e4..1e00d27817 100644 --- a/packages/core/src/index.gate.ts +++ b/packages/core/src/index.gate.ts @@ -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, diff --git a/packages/core/src/index.ts b/packages/core/src/index.ts index 543e2f6469..d6d729ecfe 100644 --- a/packages/core/src/index.ts +++ b/packages/core/src/index.ts @@ -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, diff --git a/packages/core/src/memory/recall/index.ts b/packages/core/src/memory/recall/index.ts index 601c7857d4..c932621e2f 100644 --- a/packages/core/src/memory/recall/index.ts +++ b/packages/core/src/memory/recall/index.ts @@ -6,5 +6,6 @@ export { deleteRecallRecord, getRecallRecord, listRecall, + mergeRecallGraphNodeIds, searchRecall, } from "./recall-store.js"; diff --git a/packages/core/src/memory/recall/recall-dedup.ts b/packages/core/src/memory/recall/recall-dedup.ts index 80594d2b0a..6486c6f6e2 100644 --- a/packages/core/src/memory/recall/recall-dedup.ts +++ b/packages/core/src/memory/recall/recall-dedup.ts @@ -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"); } diff --git a/packages/core/src/memory/recall/recall-store.ts b/packages/core/src/memory/recall/recall-store.ts index 3db09c8ca0..9810ab4255 100644 --- a/packages/core/src/memory/recall/recall-store.ts +++ b/packages/core/src/memory/recall/recall-store.ts @@ -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; }; +type RecallGraphNodeIdsTestHooks = { afterRowRead?: () => void | Promise }; // 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 { 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 { 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 { 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]); } diff --git a/packages/core/src/workflows/builtin-workflow-settings.ts b/packages/core/src/workflows/builtin-workflow-settings.ts index 5319cfda74..bcbcbf1e21 100644 --- a/packages/core/src/workflows/builtin-workflow-settings.ts +++ b/packages/core/src/workflows/builtin-workflow-settings.ts @@ -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", diff --git a/packages/core/src/workflows/workflow-settings-resolver.ts b/packages/core/src/workflows/workflow-settings-resolver.ts index f7fea23e08..1dd5a10f91 100644 --- a/packages/core/src/workflows/workflow-settings-resolver.ts +++ b/packages/core/src/workflows/workflow-settings-resolver.ts @@ -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 | boolean | null | undefined, +): boolean { + const value = typeof workflowEffective === "object" && workflowEffective !== null + ? workflowEffective[MEMORY_CONSOLIDATION_ENABLED_SETTING_ID] + : workflowEffective; + return value !== false; +} diff --git a/packages/dashboard/app/components/__tests__/WorkflowSettingsPanel.test.tsx b/packages/dashboard/app/components/__tests__/WorkflowSettingsPanel.test.tsx index a026ec45e2..50c10977db 100644 --- a/packages/dashboard/app/components/__tests__/WorkflowSettingsPanel.test.tsx +++ b/packages/dashboard/app/components/__tests__/WorkflowSettingsPanel.test.tsx @@ -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(); + 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 }, diff --git a/packages/dashboard/app/components/__tests__/workflow-setting-display.test.ts b/packages/dashboard/app/components/__tests__/workflow-setting-display.test.ts index ac461f21e4..e78cb25c9f 100644 --- a/packages/dashboard/app/components/__tests__/workflow-setting-display.test.ts +++ b/packages/dashboard/app/components/__tests__/workflow-setting-display.test.ts @@ -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" }, diff --git a/packages/dashboard/app/components/workflow-setting-display.ts b/packages/dashboard/app/components/workflow-setting-display.ts index ae37892ecd..c777e0c143 100644 --- a/packages/dashboard/app/components/workflow-setting-display.ts +++ b/packages/dashboard/app/components/workflow-setting-display.ts @@ -181,6 +181,11 @@ const DISPLAY: Record = { * 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", diff --git a/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts b/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts new file mode 100644 index 0000000000..e79a5b61c6 --- /dev/null +++ b/packages/engine/src/__tests__/memory-consolidation-heartbeat-hook.test.ts @@ -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("../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 = {}) { + let sequence = 0; const audits: Array> = []; + 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(); + 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) => 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)), getAsyncLayer: vi.fn(() => ({ projectId: "project" })) } as unknown as TaskStore; + return { agent, store, taskStore, audits }; +} + +function monitor(f: ReturnType) { + 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"); + }); +}); diff --git a/packages/engine/src/__tests__/memory-consolidation-ports.test.ts b/packages/engine/src/__tests__/memory-consolidation-ports.test.ts new file mode 100644 index 0000000000..d6c9c2a800 --- /dev/null +++ b/packages/engine/src/__tests__/memory-consolidation-ports.test.ts @@ -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("@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 = {}) => ({ 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" }); + }); +}); diff --git a/packages/engine/src/__tests__/memory-consolidation-run-audit-metadata.test.ts b/packages/engine/src/__tests__/memory-consolidation-run-audit-metadata.test.ts new file mode 100644 index 0000000000..d00d1244b7 --- /dev/null +++ b/packages/engine/src/__tests__/memory-consolidation-run-audit-metadata.test.ts @@ -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/); + }); +}); diff --git a/packages/engine/src/__tests__/memory-consolidation-tick.test.ts b/packages/engine/src/__tests__/memory-consolidation-tick.test.ts new file mode 100644 index 0000000000..1937a368a2 --- /dev/null +++ b/packages/engine/src/__tests__/memory-consolidation-tick.test.ts @@ -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 { + 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((_, 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); + 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); + }); +}); diff --git a/packages/engine/src/agent-heartbeat.ts b/packages/engine/src/agent-heartbeat.ts index f796bba121..b70be8dd94 100644 --- a/packages/engine/src/agent-heartbeat.ts +++ b/packages/engine/src/agent-heartbeat.ts @@ -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 { + 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; getSelfImprovePrompt(agentId: string): Promise; @@ -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) => { + 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); diff --git a/packages/engine/src/index.ts b/packages/engine/src/index.ts index 625de62c85..1a74b5f929 100644 --- a/packages/engine/src/index.ts +++ b/packages/engine/src/index.ts @@ -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 { diff --git a/packages/engine/src/memory/index.ts b/packages/engine/src/memory/index.ts new file mode 100644 index 0000000000..5403f0b6b5 --- /dev/null +++ b/packages/engine/src/memory/index.ts @@ -0,0 +1,3 @@ +export * from "./memory-consolidation.js"; +export * from "./memory-consolidation-material.js"; +export * from "./memory-consolidation-adapters.js"; diff --git a/packages/engine/src/memory/memory-consolidation-adapters.ts b/packages/engine/src/memory/memory-consolidation-adapters.ts new file mode 100644 index 0000000000..525cc377d4 --- /dev/null +++ b/packages/engine/src/memory/memory-consolidation-adapters.ts @@ -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 }; 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 { + 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, + } }; +} diff --git a/packages/engine/src/memory/memory-consolidation-material.ts b/packages/engine/src/memory/memory-consolidation-material.ts new file mode 100644 index 0000000000..4e9c3ea528 --- /dev/null +++ b/packages/engine/src/memory/memory-consolidation-material.ts @@ -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); +} diff --git a/packages/engine/src/memory/memory-consolidation.ts b/packages/engine/src/memory/memory-consolidation.ts new file mode 100644 index 0000000000..1ae8386364 --- /dev/null +++ b/packages/engine/src/memory/memory-consolidation.ts @@ -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; + 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(); +const empty = (): Omit => ({ graphChanged:false, graphRecoveryReason:null, parsedFiles:0, reusedFiles:0, prunedFiles:0, nodeCount:0, edgeCount:0, recallCandidates:0, recallCreated:0, recallDuplicate:0, crossRefUpdated:0, crossRefUnchanged:0, crossRefMissing:0, durationMs:0, changed:false }); + +/* 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 { + 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>(); 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(); 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); } + } +} diff --git a/packages/engine/src/util/run-audit.ts b/packages/engine/src/util/run-audit.ts index dafba5950a..9816965403 100644 --- a/packages/engine/src/util/run-audit.ts +++ b/packages/engine/src/util/run-audit.ts @@ -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"