From b959101bc218155f5044004807917f7cc97e0973 Mon Sep 17 00:00:00 2001 From: reddeye1337 Date: Mon, 7 Sep 2026 21:41:36 +0300 Subject: [PATCH 1/4] fix(graph): bound sourceObservationIds in graph-extract and project it out of graph-query Node creation capped provenance but mergeNode and mergeEdge re-unioned with no cap, so re-observing an entity grew the array for the life of the graph. graph-query then returned node objects verbatim, making every consumer pay for the accumulation on every call. Measured on a 71.5K-node / 191K-edge corpus: 500 nodes weighed 9.4 MB, of which sourceObservationIds was 9.55 MB - 98.9 percent. One package.json node held 4,707 ids in 133 KB. A full survey found 2,628,533 ids past a cap of 10, across 26,686 nodes and 28,884 edges. Applies both fixes suggested on the issue, since they solve different halves: - Bound the write path to GRAPH_MAX_SOURCE_IDS (default 10, matching the existing create-time cap), keeping the newest ids. Creation turned out to be unbounded too when a batch carried more observations than the cap, as did the heuristic extraction path, so the cap is applied on all four write paths rather than only on merge. - Project the array out of graph-query: callers get a three-id sample plus sourceObservationCount, and opt into the full array with includeSources. Storage growth and payload size stay independent knobs. REST /graph/query at the default limit drops from 10.7 MB to 940 KB. Complements #1294, which bounds the same field on the temporal-graph path in temporal-graph.ts; this covers mem::graph-extract in graph.ts and the read side, which neither that PR nor #1295 touches. Refs #1171 Refs #1168 Signed-off-by: reddeye1337 --- src/config.ts | 20 ++++ src/functions/graph.ts | 82 +++++++++++--- src/types.ts | 8 ++ test/graph-provenance-cap.test.ts | 174 ++++++++++++++++++++++++++++++ 4 files changed, 267 insertions(+), 17 deletions(-) create mode 100644 test/graph-provenance-cap.test.ts diff --git a/src/config.ts b/src/config.ts index a6c62ce55e..eaaa293c25 100644 --- a/src/config.ts +++ b/src/config.ts @@ -385,6 +385,26 @@ export function getGraphBatchSize(): number { return safeParseInt(getMergedEnv()["GRAPH_EXTRACTION_BATCH_SIZE"], 10); } +// Upstream #1168 / #1171: `sourceObservationIds` on graph nodes (and the +// same field on edges) was capped at creation but re-unioned without a +// cap on every merge. Because extraction re-observes the same entities +// continuously, the array grew monotonically for the life of the graph +// and ended up as ~97-99% of all bytes in the collection — measured here +// at 9.55 MB of 9.4 MB across 500 nodes, with one `package.json` node +// holding 4,707 ids in 133 KB. Provenance past the most recent handful +// carries almost no signal, so bound it and keep the newest. +const GRAPH_MAX_SOURCE_IDS_DEFAULT = 10; + +export function getGraphMaxSourceIds(): number { + return Math.max( + 1, + safeParseInt( + getMergedEnv()["GRAPH_MAX_SOURCE_IDS"], + GRAPH_MAX_SOURCE_IDS_DEFAULT, + ), + ); +} + // window for the smart-search followup-rate diagnostic. A second // search arriving within this many seconds (with disjoint results) // counts as a "follow-up" — a directional signal that the first result diff --git a/src/functions/graph.ts b/src/functions/graph.ts index 76340d60fc..19d5df9c4b 100644 --- a/src/functions/graph.ts +++ b/src/functions/graph.ts @@ -13,7 +13,10 @@ import { GRAPH_EXTRACTION_SYSTEM, buildGraphExtractionPrompt, } from "../prompts/graph-extraction.js"; -import { isGraphExtractionEnabled } from "../config.js"; +import { + isGraphExtractionEnabled, + getGraphMaxSourceIds, +} from "../config.js"; import { recordAudit } from "./audit.js"; import { logger } from "../logger.js"; @@ -22,6 +25,37 @@ import { logger } from "../logger.js"; // reported 11k-node / 28k-edge corpus, and 5,000 is the upper bound a // caller can request explicitly. Tuned conservatively because edges // fan out faster than nodes. +// How many provenance ids a projected graph-query response keeps as a +// sample. The full array is available with `includeSources: true`. +const GRAPH_QUERY_SOURCE_SAMPLE = 3; + +/** Bound a provenance array, keeping the most recent ids. */ +function capSourceIds(ids: string[]): string[] { + const max = getGraphMaxSourceIds(); + return ids.length <= max ? ids : ids.slice(-max); +} + +/** + * Strip the provenance array down to a count plus a small sample. + * graph-query used to return node objects verbatim, so every consumer + * paid for accumulated provenance on every call — for an agent that + * cost is context window (upstream #1171). + */ +function projectSources( + row: T, + includeSources: boolean, +): T & { sourceObservationCount?: number } { + const ids = row.sourceObservationIds ?? []; + if (includeSources || ids.length <= GRAPH_QUERY_SOURCE_SAMPLE) { + return { ...row, sourceObservationCount: ids.length }; + } + return { + ...row, + sourceObservationIds: ids.slice(-GRAPH_QUERY_SOURCE_SAMPLE), + sourceObservationCount: ids.length, + }; +} + const DEFAULT_GRAPH_QUERY_LIMIT = 500; const MAX_GRAPH_QUERY_LIMIT = 5000; @@ -146,6 +180,7 @@ function paginateFromSnapshot( filterType: string | undefined, limit: number, offset: number, + includeSources = false, ): GraphQueryResult { const filteredNodes = filterType ? snap.topNodes.filter((n) => n.type === filterType) @@ -159,8 +194,8 @@ function paginateFromSnapshot( (e) => pageIds.has(e.sourceNodeId) && pageIds.has(e.targetNodeId), ); return { - nodes: pageNodes, - edges: pageEdges, + nodes: pageNodes.map((n) => projectSources(n, includeSources)), + edges: pageEdges.map((e) => projectSources(e, includeSources)), depth: 0, totalNodes: total, totalEdges: snap.stats.totalEdges, @@ -280,13 +315,15 @@ function mergeNode( ): GraphNode { return { ...existing, - sourceObservationIds: [ + // Newest ids sort last through the Set, so the cap keeps the most + // recent provenance and drops the oldest (upstream #1171). + sourceObservationIds: capSourceIds([ ...new Set([ ...existing.sourceObservationIds, ...incoming.sourceObservationIds, ...obsIds, ]), - ], + ]), properties: { ...existing.properties, ...incoming.properties }, updatedAt: capturedAt, }; @@ -298,9 +335,9 @@ function mergeEdge( ): GraphEdge { return { ...existing, - sourceObservationIds: [ + sourceObservationIds: capSourceIds([ ...new Set([...existing.sourceObservationIds, ...obsIds]), - ], + ]), }; } @@ -327,6 +364,7 @@ function paginate( depth: number, limit: number, offset: number, + includeSources = false, ): GraphQueryResult { const totalNodes = nodes.length; const pageNodes = nodes.slice(offset, offset + limit); @@ -349,8 +387,8 @@ function paginate( 0, ); return { - nodes: pageNodes, - edges: pageEdges, + nodes: pageNodes.map((n) => projectSources(n, includeSources)), + edges: pageEdges.map((e) => projectSources(e, includeSources)), depth, totalNodes, totalEdges, @@ -411,7 +449,7 @@ function parseGraphXml( type, name, properties, - sourceObservationIds: observationIds, + sourceObservationIds: capSourceIds(observationIds), createdAt: now, }); }; @@ -443,7 +481,7 @@ function parseGraphXml( sourceNodeId: sourceNode.id, targetNodeId: targetNode.id, weight: Math.max(0, Math.min(1, weight)), - sourceObservationIds: observationIds, + sourceObservationIds: capSourceIds(observationIds), createdAt: now, }); } @@ -484,7 +522,10 @@ export function extractGraphHeuristics( nodeByKey.set(key, node); nodes.push(node); } else if (!node.sourceObservationIds.includes(obsId)) { - node.sourceObservationIds.push(obsId); + node.sourceObservationIds = capSourceIds([ + ...node.sourceObservationIds, + obsId, + ]); } return node; }; @@ -497,7 +538,10 @@ export function extractGraphHeuristics( const existing = edgeByPair.get(pair); if (existing) { if (!existing.sourceObservationIds.includes(obs.id)) { - existing.sourceObservationIds.push(obs.id); + existing.sourceObservationIds = capSourceIds([ + ...existing.sourceObservationIds, + obs.id, + ]); } return; } @@ -784,9 +828,13 @@ export function registerGraphFunction( query?: string; limit?: number; offset?: number; + includeSources?: boolean; }): Promise => { const maxDepth = Math.min(data.maxDepth || 3, 5); const { limit, offset } = resolvePagination(data.limit, data.offset); + // Off by default: the full provenance array is ~99% of the bytes + // and almost never what the caller wanted (upstream #1171). + const includeSources = data.includeSources === true; // #814 v2: the empty-body / nodeType-only path NEVER enumerates. // It reads the snapshot exclusively. The snapshot is updated @@ -799,7 +847,7 @@ export function registerGraphFunction( if (noWalk) { const snap = await readSnapshot(kv); if (snap && snap.stats.totalNodes > 0) { - return paginateFromSnapshot(snap, data.nodeType, limit, offset); + return paginateFromSnapshot(snap, data.nodeType, limit, offset, includeSources); } return { nodes: [], @@ -875,7 +923,7 @@ export function registerGraphFunction( (v) => typeof v === "string" && v.toLowerCase().includes(lower), ), ); - return paginate(matchingNodes, allEdges, 0, limit, offset); + return paginate(matchingNodes, allEdges, 0, limit, offset, includeSources); } if (data.startNodeId) { @@ -917,11 +965,11 @@ export function registerGraphFunction( } } - return paginate(resultNodes, resultEdges, maxDepth, limit, offset); + return paginate(resultNodes, resultEdges, maxDepth, limit, offset, includeSources); } // Unreachable — noWalk branch handles the rest. - return paginate([], [], 0, limit, offset); + return paginate([], [], 0, limit, offset, includeSources); }, ); diff --git a/src/types.ts b/src/types.ts index d2c63efa61..b5b786a8fb 100644 --- a/src/types.ts +++ b/src/types.ts @@ -409,6 +409,10 @@ export interface GraphNode { updatedAt?: string; aliases?: string[]; stale?: boolean; + /** Present on graph-query responses: true length of sourceObservationIds + * before projection. The array itself is a recent sample unless the + * caller passed includeSources (upstream #1171). */ + sourceObservationCount?: number; } export type GraphEdgeType = @@ -445,6 +449,10 @@ export interface GraphEdge { supersededBy?: string; isLatest?: boolean; stale?: boolean; + /** Present on graph-query responses: true length of sourceObservationIds + * before projection. The array itself is a recent sample unless the + * caller passed includeSources (upstream #1171). */ + sourceObservationCount?: number; } export interface EdgeContext { diff --git a/test/graph-provenance-cap.test.ts b/test/graph-provenance-cap.test.ts new file mode 100644 index 0000000000..73dded2c71 --- /dev/null +++ b/test/graph-provenance-cap.test.ts @@ -0,0 +1,174 @@ +import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; + +vi.mock("../src/logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn() }, +})); + +import { registerGraphFunction } from "../src/functions/graph.js"; +import type { + CompressedObservation, + GraphNode, + GraphQueryResult, +} from "../src/types.js"; + +function mockKV() { + const store = new Map>(); + return { + get: async (scope: string, key: string): Promise => { + return (store.get(scope)?.get(key) as T) ?? null; + }, + set: async (scope: string, key: string, data: T): Promise => { + if (!store.has(scope)) store.set(scope, new Map()); + store.get(scope)!.set(key, data); + return data; + }, + delete: async (scope: string, key: string): Promise => { + store.get(scope)?.delete(key); + }, + list: async (scope: string): Promise => { + const entries = store.get(scope); + return entries ? (Array.from(entries.values()) as T[]) : []; + }, + }; +} + +function mockSdk() { + const functions = new Map(); + return { + registerFunction: (idOrOpts: string | { id: string }, handler: Function) => { + const id = typeof idOrOpts === "string" ? idOrOpts : idOrOpts.id; + functions.set(id, handler); + }, + registerTrigger: () => {}, + trigger: async ( + idOrInput: string | { function_id: string; payload: unknown }, + data?: unknown, + ) => { + const id = + typeof idOrInput === "string" ? idOrInput : idOrInput.function_id; + const payload = typeof idOrInput === "string" ? data : idOrInput.payload; + const fn = functions.get(id); + if (!fn) throw new Error(`No function: ${id}`); + return fn(payload); + }, + }; +} + +// One entity, no relationships: every extract re-observes the same node, +// which is the shape that made the provenance array grow without bound. +const mockProvider = { + name: "test", + compress: vi.fn().mockResolvedValue(` +src/hot-path.ts + + +`), + summarize: vi.fn(), +}; + +function obs(n: number): CompressedObservation { + return { + id: `obs_${n}`, + sessionId: "ses_1", + timestamp: `2026-02-01T10:00:${String(n % 60).padStart(2, "0")}Z`, + type: "file_edit", + title: `Edit ${n}`, + facts: [`change ${n}`], + narrative: `Edited the hot path, revision ${n}`, + concepts: [], + files: [], + importance: 5, + }; +} + +describe("graph provenance bounds (#1168 / #1171)", () => { + let sdk: ReturnType; + let kv: ReturnType; + const ORIG_FLAG = process.env["GRAPH_EXTRACTION_ENABLED"]; + const ORIG_MAX = process.env["GRAPH_MAX_SOURCE_IDS"]; + + beforeEach(() => { + sdk = mockSdk(); + kv = mockKV(); + vi.clearAllMocks(); + process.env["GRAPH_EXTRACTION_ENABLED"] = "true"; + process.env["GRAPH_MAX_SOURCE_IDS"] = "5"; + registerGraphFunction(sdk as never, kv as never, mockProvider as never); + }); + + afterEach(() => { + if (ORIG_FLAG === undefined) delete process.env["GRAPH_EXTRACTION_ENABLED"]; + else process.env["GRAPH_EXTRACTION_ENABLED"] = ORIG_FLAG; + if (ORIG_MAX === undefined) delete process.env["GRAPH_MAX_SOURCE_IDS"]; + else process.env["GRAPH_MAX_SOURCE_IDS"] = ORIG_MAX; + }); + + // The regression: creation capped at ten ids, mergeNode re-unioned with + // no cap, so re-observing an entity grew the array for the life of the + // graph. + it("keeps sourceObservationIds bounded across repeated merges", async () => { + for (let i = 1; i <= 20; i++) { + await sdk.trigger("mem::graph-extract", { observations: [obs(i)] }); + } + + const nodes = await kv.list("mem:graph:nodes"); + const hot = nodes.find((n) => n.name === "src/hot-path.ts")!; + expect(hot).toBeDefined(); + expect(hot.sourceObservationIds.length).toBeLessThanOrEqual(5); + // Newest kept, oldest dropped. + expect(hot.sourceObservationIds).toContain("obs_20"); + expect(hot.sourceObservationIds).not.toContain("obs_1"); + }); + + it("caps provenance at creation too", async () => { + await sdk.trigger("mem::graph-extract", { + observations: [obs(1), obs(2), obs(3), obs(4), obs(5), obs(6), obs(7)], + }); + + const nodes = await kv.list("mem:graph:nodes"); + for (const node of nodes) { + expect(node.sourceObservationIds.length).toBeLessThanOrEqual(5); + } + }); + + it("graph-query projects provenance to a count plus a sample", async () => { + for (let i = 1; i <= 20; i++) { + await sdk.trigger("mem::graph-extract", { observations: [obs(i)] }); + } + + const result = (await sdk.trigger("mem::graph-query", { + query: "hot-path", + })) as GraphQueryResult; + + const node = result.nodes.find((n) => n.name === "src/hot-path.ts")!; + expect(node).toBeDefined(); + expect(node.sourceObservationCount).toBeGreaterThan(0); + expect(node.sourceObservationIds.length).toBeLessThanOrEqual(3); + }); + + it("graph-query returns the full array when includeSources is set", async () => { + for (let i = 1; i <= 20; i++) { + await sdk.trigger("mem::graph-extract", { observations: [obs(i)] }); + } + + const projected = (await sdk.trigger("mem::graph-query", { + query: "hot-path", + })) as GraphQueryResult; + const full = (await sdk.trigger("mem::graph-query", { + query: "hot-path", + includeSources: true, + })) as GraphQueryResult; + + const projectedNode = projected.nodes.find( + (n) => n.name === "src/hot-path.ts", + )!; + const fullNode = full.nodes.find((n) => n.name === "src/hot-path.ts")!; + + expect(fullNode.sourceObservationIds.length).toBe( + fullNode.sourceObservationCount, + ); + expect(fullNode.sourceObservationIds.length).toBeGreaterThanOrEqual( + projectedNode.sourceObservationIds.length, + ); + }); +}); From 160be41423363e77b680702f211b26862de82665 Mon Sep 17 00:00:00 2001 From: reddeye1337 Date: Mon, 7 Sep 2026 21:46:11 +0300 Subject: [PATCH 2/4] feat(graph): add mem::graph-prune with in-place provenance compaction The knowledge graph has no collector. Every extraction appends nodes and edges and nothing removes them, so a corpus passes the 25K-node ceiling that mem::graph-snapshot-rebuild refuses to run above and there is no path back short of graph-reset, which discards the whole graph. #1171 names the second half of this: extract, build, reset and snapshot-rebuild are the only writers, so an operator whose nodes have already accumulated millions of provenance ids cannot compact them in place. Bounding the write path only helps rows written afterwards. mem::graph-prune, exposed at POST /agentmemory/graph/prune, collects three classes that are provably dead - rows already tombstoned as stale, edges whose endpoint node no longer exists, and temporal edges marked isLatest:false past a retention cutoff - and clears the matching name-index, edge-key and degree entries so graph-extract dedup cannot resolve into deleted rows afterwards. With compactSourceIds it also trims oversized provenance arrays to the newest ids, using the same cap the write path applies, so an existing corpus can be brought in line without losing the graph. Defaults to a dry run, which reports what each class would collect before anything is written. Duplicate-name merging rewrites live rows, so it stays behind its own flag. On a 71.5K-node / 191K-edge corpus: 26,686 nodes and 28,884 edges compacted in 22s, 2,628,533 provenance ids dropped, 0 errors, daemon RSS 1636 MB to 922 MB. Stacked on the graph provenance cap, whose GRAPH_MAX_SOURCE_IDS the compaction pass reuses so the backfill and the write path cannot disagree. Refs #1171 Refs #1168 Signed-off-by: reddeye1337 --- AGENTS.md | 2 +- README.md | 2 +- src/functions/graph-prune.ts | 369 +++++++++++++++++++++++++++++++++++ src/functions/graph.ts | 6 +- src/index.ts | 4 +- src/triggers/api.ts | 31 +++ test/graph-prune.test.ts | 170 ++++++++++++++++ 7 files changed, 578 insertions(+), 6 deletions(-) create mode 100644 src/functions/graph-prune.ts create mode 100644 test/graph-prune.test.ts diff --git a/AGENTS.md b/AGENTS.md index dc08f431ca..29f2e7ed78 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -117,7 +117,7 @@ Hook scripts in `src/hooks/` are standalone Node.js scripts (no iii-sdk import). ## Current Stats (v0.9.29) - 54 MCP tools (8 visible by default, `AGENTMEMORY_TOOLS=all` for all) -- 130 REST endpoints +- 131 REST endpoints - 6 MCP resources, 3 MCP prompts - 12 hooks, 17 skills - 260+ iii functions diff --git a/README.md b/README.md index 959d68c7e4..eb07f2064d 100644 --- a/README.md +++ b/README.md @@ -1609,7 +1609,7 @@ Create `~/.agentmemory/.env`:

API

-130 endpoints on port `3111`. The REST API binds to `127.0.0.1` by default. Protected endpoints require `Authorization: Bearer ` when `AGENTMEMORY_SECRET` is set, and mesh sync endpoints require `AGENTMEMORY_SECRET` on both peers. +131 endpoints on port `3111`. The REST API binds to `127.0.0.1` by default. Protected endpoints require `Authorization: Bearer ` when `AGENTMEMORY_SECRET` is set, and mesh sync endpoints require `AGENTMEMORY_SECRET` on both peers.
Key endpoints diff --git a/src/functions/graph-prune.ts b/src/functions/graph-prune.ts new file mode 100644 index 0000000000..c95e2aea1e --- /dev/null +++ b/src/functions/graph-prune.ts @@ -0,0 +1,369 @@ +import type { ISdk } from "iii-sdk"; +import type { GraphEdge, GraphNode, GraphSnapshot } from "../types.js"; +import type { StateKV } from "../state/kv.js"; +import { KV } from "../state/schema.js"; +import { logger } from "../logger.js"; +import { edgeIndexKey, nameIndexKey, SNAPSHOT_KEY } from "./graph.js"; +import { getGraphMaxSourceIds } from "../config.js"; + +// The knowledge graph had no GC. Every extraction appended nodes and +// edges, nothing ever removed them, and at 69K nodes / 185K edges the +// retrieval path (and mem::graph-snapshot-rebuild, whose safe ceiling +// is 25K nodes) had already gone past what a single kv.list can carry. +// +// This is the missing collector. It runs off its own enumeration +// rather than the search-path graph view, because the view +// deliberately excludes `stale` rows and those are the first thing +// worth collecting. Deletions route through StateKV, so the shared +// view is patched as they happen. +// +// Every class it removes is provably dead: +// stale — already tombstoned by cascade/mesh, kept only because +// nothing swept them. +// dangling — an edge whose endpoint node no longer exists; it can +// never be traversed. +// superseded— a temporal edge explicitly marked isLatest:false and +// older than the retention cutoff, i.e. history that has +// already been replaced by a newer revision. +// +// Duplicate-name merging is real cleanup but it rewrites live rows, so +// it stays opt-in (`mergeDuplicateNames: true`) and reports its count +// in dry-run first. + +const DEFAULT_SUPERSEDED_RETENTION_DAYS = 90; + +export interface GraphPruneReport { + success: true; + dryRun: boolean; + before: { nodes: number; edges: number }; + staleNodes: number; + staleEdges: number; + danglingEdges: number; + supersededEdges: number; + duplicateNodes: number; + oversizedNodes: number; + oversizedEdges: number; + droppableSourceIds: number; + compactedNodes: number; + compactedEdges: number; + mergedNodes: number; + repointedEdges: number; + deletedNodes: number; + deletedEdges: number; + errors: number; + ms: number; +} + +function edgeTimestamp(edge: GraphEdge): number { + const raw = edge.tcommit || edge.createdAt; + const parsed = raw ? new Date(raw).getTime() : NaN; + return Number.isFinite(parsed) ? parsed : 0; +} + +export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { + sdk.registerFunction( + "mem::graph-prune", + async (data?: { + dryRun?: boolean; + supersededOlderThanDays?: number; + mergeDuplicateNames?: boolean; + compactSourceIds?: boolean; + }): Promise => { + const started = Date.now(); + // Default to a dry run: this deletes graph rows, and the caller + // should have to say so explicitly. The scheduled sweep passes + // dryRun:false. + const dryRun = data?.dryRun !== false; + const retentionDays = + typeof data?.supersededOlderThanDays === "number" && + data.supersededOlderThanDays >= 0 + ? data.supersededOlderThanDays + : DEFAULT_SUPERSEDED_RETENTION_DAYS; + const mergeDuplicates = data?.mergeDuplicateNames === true; + const compactSources = data?.compactSourceIds === true; + const maxSourceIds = getGraphMaxSourceIds(); + const cutoff = Date.now() - retentionDays * 24 * 60 * 60 * 1000; + + // Sequential, matching the graph view build: two multi-megabyte + // frames in flight at once doubles the parse cost on the worker + // event loop. + const allNodes = await kv.list(KV.graphNodes); + const allEdges = await kv.list(KV.graphEdges); + + const report: GraphPruneReport = { + success: true, + dryRun, + before: { nodes: allNodes.length, edges: allEdges.length }, + staleNodes: 0, + staleEdges: 0, + danglingEdges: 0, + supersededEdges: 0, + duplicateNodes: 0, + oversizedNodes: 0, + oversizedEdges: 0, + droppableSourceIds: 0, + compactedNodes: 0, + compactedEdges: 0, + mergedNodes: 0, + repointedEdges: 0, + deletedNodes: 0, + deletedEdges: 0, + errors: 0, + ms: 0, + }; + + const liveNodes = new Map(); + const nodesToDelete = new Map(); + for (const node of allNodes) { + if (!node?.id) continue; + if (node.stale) { + report.staleNodes++; + nodesToDelete.set(node.id, node); + } else { + liveNodes.set(node.id, node); + } + } + + const edgesToDelete = new Map(); + const markEdge = (edge: GraphEdge): void => { + if (!edgesToDelete.has(edge.id)) edgesToDelete.set(edge.id, edge); + }; + for (const edge of allEdges) { + if (!edge?.id) continue; + if (edge.stale) { + report.staleEdges++; + markEdge(edge); + continue; + } + if ( + !liveNodes.has(edge.sourceNodeId) || + !liveNodes.has(edge.targetNodeId) + ) { + report.danglingEdges++; + markEdge(edge); + continue; + } + if (edge.isLatest === false && edgeTimestamp(edge) < cutoff) { + report.supersededEdges++; + markEdge(edge); + } + } + + // Duplicate detection always runs so the dry-run report shows + // what merging would recover; the rewrite itself is opt-in. + const byName = new Map(); + for (const node of liveNodes.values()) { + const key = nameIndexKey(node.type, node.name); + const group = byName.get(key); + if (group) group.push(node); + else byName.set(key, [node]); + } + const duplicateGroups: GraphNode[][] = []; + for (const group of byName.values()) { + if (group.length < 2) continue; + // Oldest row wins: it is the one the name index and existing + // edges are most likely already pointing at. + group.sort((a, b) => (a.createdAt || "").localeCompare(b.createdAt || "")); + duplicateGroups.push(group); + report.duplicateNodes += group.length - 1; + } + + // Provenance survey. Counted for every live row so the dry run + // reports what compaction would recover before anything is written. + const oversizedNodeRows: GraphNode[] = []; + for (const node of liveNodes.values()) { + const extra = (node.sourceObservationIds?.length ?? 0) - maxSourceIds; + if (extra > 0) { + report.oversizedNodes++; + report.droppableSourceIds += extra; + oversizedNodeRows.push(node); + } + } + const oversizedEdgeRows: GraphEdge[] = []; + for (const edge of allEdges) { + if (!edge?.id || edgesToDelete.has(edge.id)) continue; + const extra = (edge.sourceObservationIds?.length ?? 0) - maxSourceIds; + if (extra > 0) { + report.oversizedEdges++; + report.droppableSourceIds += extra; + oversizedEdgeRows.push(edge); + } + } + + if (dryRun) { + report.ms = Date.now() - started; + logger.info("Graph prune (dry run)", { ...report }); + return report; + } + + // ---- compact provenance ---------------------------------------------- + // Keeps the newest ids, matching the cap the merge path now applies + // at write time. Nothing else on the row is touched. + if (compactSources) { + for (const node of oversizedNodeRows) { + if (nodesToDelete.has(node.id)) continue; + try { + const trimmed = node.sourceObservationIds.slice(-maxSourceIds); + await kv.set(KV.graphNodes, node.id, { + ...node, + sourceObservationIds: trimmed, + }); + liveNodes.set(node.id, { ...node, sourceObservationIds: trimmed }); + report.compactedNodes++; + } catch (err) { + report.errors++; + logger.warn("Graph prune node compaction failed", { + nodeId: node.id, + error: err instanceof Error ? err.message : String(err), + }); + } + } + for (const edge of oversizedEdgeRows) { + if (edgesToDelete.has(edge.id)) continue; + try { + await kv.set(KV.graphEdges, edge.id, { + ...edge, + sourceObservationIds: edge.sourceObservationIds.slice(-maxSourceIds), + }); + report.compactedEdges++; + } catch (err) { + report.errors++; + logger.warn("Graph prune edge compaction failed", { + edgeId: edge.id, + error: err instanceof Error ? err.message : String(err), + }); + } + } + } + + // ---- merge duplicates ------------------------------------------------- + const idRemap = new Map(); + if (mergeDuplicates) { + for (const group of duplicateGroups) { + const [keeper, ...losers] = group; + const obsIds = new Set(keeper.sourceObservationIds ?? []); + for (const loser of losers) { + for (const obsId of loser.sourceObservationIds ?? []) { + obsIds.add(obsId); + } + idRemap.set(loser.id, keeper.id); + nodesToDelete.set(loser.id, loser); + } + const merged: GraphNode = { + ...keeper, + sourceObservationIds: Array.from(obsIds), + updatedAt: new Date().toISOString(), + }; + try { + await kv.set(KV.graphNodes, keeper.id, merged); + await kv.set(KV.graphNameIndex, nameIndexKey(keeper.type, keeper.name), keeper.id); + liveNodes.set(keeper.id, merged); + report.mergedNodes += losers.length; + } catch (err) { + report.errors++; + logger.warn("Graph prune merge failed", { + nodeId: keeper.id, + error: err instanceof Error ? err.message : String(err), + }); + } + } + + // Repoint surviving edges off the merged-away nodes. An edge + // that collapses to a self-loop is dropped rather than kept. + for (const edge of allEdges) { + if (!edge?.id || edgesToDelete.has(edge.id)) continue; + const src = idRemap.get(edge.sourceNodeId) ?? edge.sourceNodeId; + const tgt = idRemap.get(edge.targetNodeId) ?? edge.targetNodeId; + if (src === edge.sourceNodeId && tgt === edge.targetNodeId) continue; + if (src === tgt) { + markEdge(edge); + continue; + } + try { + await kv.delete(KV.graphEdgeKey, edgeIndexKey(edge.sourceNodeId, edge.targetNodeId, edge.type)); + const repointed: GraphEdge = { ...edge, sourceNodeId: src, targetNodeId: tgt }; + await kv.set(KV.graphEdges, edge.id, repointed); + await kv.set(KV.graphEdgeKey, edgeIndexKey(src, tgt, edge.type), edge.id); + report.repointedEdges++; + } catch (err) { + report.errors++; + logger.warn("Graph prune repoint failed", { + edgeId: edge.id, + error: err instanceof Error ? err.message : String(err), + }); + } + } + } + + // ---- delete ---------------------------------------------------------- + // Edges first: deleting a node before its edges would turn those + // edges into dangling rows if the run is interrupted. + for (const edge of edgesToDelete.values()) { + try { + await kv.delete(KV.graphEdges, edge.id); + const key = edgeIndexKey(edge.sourceNodeId, edge.targetNodeId, edge.type); + const indexed = await kv.get(KV.graphEdgeKey, key).catch(() => null); + // Only clear the index entry if it still points at this edge; + // a newer edge may have claimed the same endpoint triple. + if (indexed === edge.id) await kv.delete(KV.graphEdgeKey, key); + report.deletedEdges++; + } catch (err) { + report.errors++; + logger.warn("Graph prune edge delete failed", { + edgeId: edge.id, + error: err instanceof Error ? err.message : String(err), + }); + } + } + + for (const node of nodesToDelete.values()) { + try { + await kv.delete(KV.graphNodes, node.id); + await kv.delete(KV.graphNodeDegree, node.id).catch(() => {}); + const key = nameIndexKey(node.type, node.name); + const indexed = await kv.get(KV.graphNameIndex, key).catch(() => null); + if (indexed === node.id) await kv.delete(KV.graphNameIndex, key); + report.deletedNodes++; + } catch (err) { + report.errors++; + logger.warn("Graph prune node delete failed", { + nodeId: node.id, + error: err instanceof Error ? err.message : String(err), + }); + } + } + + // Keep the precomputed snapshot honest. Its per-type breakdown + // will drift, so flag it dirty; the totals are corrected here + // because /graph/stats reports them directly. + if (report.deletedNodes > 0 || report.deletedEdges > 0) { + try { + const snap = await kv.get(KV.graphSnapshot, SNAPSHOT_KEY); + if (snap) { + await kv.set(KV.graphSnapshot, SNAPSHOT_KEY, { + ...snap, + topNodes: (snap.topNodes ?? []).filter((n) => !nodesToDelete.has(n.id)), + topEdges: (snap.topEdges ?? []).filter((e) => !edgesToDelete.has(e.id)), + stats: { + ...snap.stats, + totalNodes: Math.max(0, snap.stats.totalNodes - report.deletedNodes), + totalEdges: Math.max(0, snap.stats.totalEdges - report.deletedEdges), + }, + dirty: true, + updatedAt: new Date().toISOString(), + }); + } + } catch (err) { + report.errors++; + logger.warn("Graph prune snapshot update failed", { + error: err instanceof Error ? err.message : String(err), + }); + } + } + + report.ms = Date.now() - started; + logger.info("Graph prune complete", { ...report }); + return report; + }, + ); +} diff --git a/src/functions/graph.ts b/src/functions/graph.ts index 19d5df9c4b..35eee3ec21 100644 --- a/src/functions/graph.ts +++ b/src/functions/graph.ts @@ -66,7 +66,7 @@ const MAX_GRAPH_QUERY_LIMIT = 5000; // enumeration. Aggregate stats (nodesByType / edgesByType) are computed // fresh during rebuild and stored alongside. const SNAPSHOT_TOP_NODES = DEFAULT_GRAPH_QUERY_LIMIT; -const SNAPSHOT_KEY = "current"; +export const SNAPSHOT_KEY = "current"; // `state::list` over a 75K-node scope can exceed the iii invocation // timeout. The query handler races the enumeration against this budget @@ -217,11 +217,11 @@ function paginateFromSnapshot( // future extracts rebuild incrementally. const REBUILD_SAFE_NODE_CEILING = 25000; -function nameIndexKey(type: string, name: string): string { +export function nameIndexKey(type: string, name: string): string { return `${type}|${name}`; } -function edgeIndexKey( +export function edgeIndexKey( sourceNodeId: string, targetNodeId: string, type: string, diff --git a/src/index.ts b/src/index.ts index 717419f7df..b67e290d6f 100644 --- a/src/index.ts +++ b/src/index.ts @@ -59,6 +59,7 @@ import { registerExportImportFunction } from "./functions/export-import.js"; import { registerEnrichFunction } from "./functions/enrich.js"; import { registerClaudeBridgeFunction } from "./functions/claude-bridge.js"; import { registerGraphFunction } from "./functions/graph.js"; +import { registerGraphPruneFunction } from "./functions/graph-prune.js"; import { registerGraphImportFunction } from "./functions/graph-import.js"; import { registerConsolidationPipelineFunction } from "./functions/consolidation-pipeline.js"; import { registerTeamFunction } from "./functions/team.js"; @@ -268,6 +269,7 @@ async function main() { } registerGraphFunction(sdk, kv, provider); + registerGraphPruneFunction(sdk, kv); registerGraphImportFunction(sdk, kv); bootLog( `Knowledge graph: structural extraction on (LLM relations ${isGraphExtractionEnabled() ? "enabled" : "off"})`, @@ -536,7 +538,7 @@ async function main() { `Ready. ${embeddingProvider ? "Triple-stream (BM25+Vector+Graph)" : "BM25+Graph"} search active.`, ); bootLog( - `REST API: 130 endpoints at http://localhost:${config.restPort}/agentmemory/*`, + `REST API: 131 endpoints at http://localhost:${config.restPort}/agentmemory/*`, ); bootLog( `MCP surface (opt-in via \`npx @agentmemory/mcp\`): ${getAllTools().length} tools · 6 resources · 3 prompts`, diff --git a/src/triggers/api.ts b/src/triggers/api.ts index 56fad4f0de..c3d06ce169 100644 --- a/src/triggers/api.ts +++ b/src/triggers/api.ts @@ -1577,6 +1577,37 @@ export function registerApiTriggers( config: { api_path: "/agentmemory/graph/reset", http_method: "POST" }, }); + // Graph GC. Defaults to a dry run so an accidental POST reports what + // it would collect instead of collecting it; the scheduled sweep + // passes dryRun:false explicitly. + sdk.registerFunction("api::graph-prune", + async ( + req: ApiRequest<{ + dryRun?: boolean; + supersededOlderThanDays?: number; + mergeDuplicateNames?: boolean; + compactSourceIds?: boolean; + }>, + ): Promise => { + const authErr = checkAuth(req, secret); + if (authErr) return authErr; + try { + const result = await sdk.trigger({ + function_id: "mem::graph-prune", + payload: req.body ?? {}, + }); + return { status_code: 200, body: result }; + } catch { + return graphDisabledResponse(); + } + }, + ); + sdk.registerTrigger({ + type: "http", + function_id: "api::graph-prune", + config: { api_path: "/agentmemory/graph/prune", http_method: "POST" }, + }); + sdk.registerFunction("api::graph-extract", async (req: ApiRequest<{ observations: unknown[] }>): Promise => { const authErr = checkAuth(req, secret); diff --git a/test/graph-prune.test.ts b/test/graph-prune.test.ts new file mode 100644 index 0000000000..57919b27fa --- /dev/null +++ b/test/graph-prune.test.ts @@ -0,0 +1,170 @@ +import { describe, it, expect, beforeEach, afterEach, vi } from "vitest"; + +vi.mock("../src/logger.js", () => ({ + logger: { info: vi.fn(), warn: vi.fn(), error: vi.fn() }, +})); + +import { registerGraphPruneFunction } from "../src/functions/graph-prune.js"; +import type { GraphNode, GraphEdge } from "../src/types.js"; + +function mockKV() { + const store = new Map>(); + return { + raw: store, + get: async (scope: string, key: string): Promise => + (store.get(scope)?.get(key) as T) ?? null, + set: async (scope: string, key: string, data: T): Promise => { + if (!store.has(scope)) store.set(scope, new Map()); + store.get(scope)!.set(key, data); + return data; + }, + delete: async (scope: string, key: string): Promise => { + store.get(scope)?.delete(key); + }, + list: async (scope: string): Promise => { + const entries = store.get(scope); + return entries ? (Array.from(entries.values()) as T[]) : []; + }, + }; +} + +function mockSdk() { + const functions = new Map(); + return { + registerFunction: (id: string, handler: Function) => { + functions.set(id, handler); + }, + registerTrigger: () => {}, + trigger: async (input: { function_id: string; payload: unknown }) => { + const fn = functions.get(input.function_id); + if (!fn) throw new Error(`No function: ${input.function_id}`); + return fn(input.payload); + }, + }; +} + +function node(id: string, over: Partial = {}): GraphNode { + return { + id, + type: "concept", + name: id, + properties: {}, + sourceObservationIds: ["obs_1"], + createdAt: "2026-02-01T10:00:00Z", + ...over, + }; +} + +function edge(id: string, s: string, t: string, over: Partial = {}): GraphEdge { + return { + id, + type: "related_to", + sourceNodeId: s, + targetNodeId: t, + weight: 0.5, + sourceObservationIds: ["obs_1"], + createdAt: "2026-02-01T10:00:00Z", + ...over, + }; +} + +type Report = { + staleNodes: number; + staleEdges: number; + danglingEdges: number; + supersededEdges: number; + oversizedNodes: number; + droppableSourceIds: number; + deletedNodes: number; + deletedEdges: number; + compactedNodes: number; + compactedEdges: number; + errors: number; +}; + +describe("mem::graph-prune", () => { + let sdk: ReturnType; + let kv: ReturnType; + const ORIG_MAX = process.env["GRAPH_MAX_SOURCE_IDS"]; + + const run = (payload: unknown) => + sdk.trigger({ function_id: "mem::graph-prune", payload }) as Promise; + + beforeEach(async () => { + sdk = mockSdk(); + kv = mockKV(); + process.env["GRAPH_MAX_SOURCE_IDS"] = "3"; + registerGraphPruneFunction(sdk as never, kv as never); + + await kv.set("mem:graph:nodes", "live", node("live")); + await kv.set("mem:graph:nodes", "gone", node("gone", { stale: true })); + await kv.set( + "mem:graph:nodes", + "fat", + node("fat", { + sourceObservationIds: ["o1", "o2", "o3", "o4", "o5", "o6"], + }), + ); + await kv.set("mem:graph:edges", "keep", edge("keep", "live", "fat")); + // Endpoint was tombstoned, so this edge can never be traversed. + await kv.set("mem:graph:edges", "dangling", edge("dangling", "live", "gone")); + }); + + afterEach(() => { + if (ORIG_MAX === undefined) delete process.env["GRAPH_MAX_SOURCE_IDS"]; + else process.env["GRAPH_MAX_SOURCE_IDS"] = ORIG_MAX; + }); + + it("dry run reports what it would collect and writes nothing", async () => { + const report = await run({ dryRun: true }); + + expect(report.staleNodes).toBe(1); + expect(report.danglingEdges).toBe(1); + expect(report.oversizedNodes).toBe(1); + expect(report.droppableSourceIds).toBe(3); + expect(report.deletedNodes).toBe(0); + expect(report.compactedNodes).toBe(0); + + expect((await kv.list("mem:graph:nodes")).length).toBe(3); + expect((await kv.list("mem:graph:edges")).length).toBe(2); + }); + + it("deletes stale nodes and edges that can never be traversed", async () => { + const report = await run({ dryRun: false }); + + expect(report.errors).toBe(0); + expect(report.deletedNodes).toBe(1); + expect(report.deletedEdges).toBe(1); + + const nodes = await kv.list("mem:graph:nodes"); + const edges = await kv.list("mem:graph:edges"); + expect(nodes.map((n) => n.id).sort()).toEqual(["fat", "live"]); + expect(edges.map((e) => e.id)).toEqual(["keep"]); + }); + + it("leaves provenance alone unless compaction is asked for", async () => { + await run({ dryRun: false }); + const fat = await kv.get("mem:graph:nodes", "fat"); + expect(fat!.sourceObservationIds.length).toBe(6); + }); + + it("compacts oversized provenance to the newest ids", async () => { + const report = await run({ dryRun: false, compactSourceIds: true }); + + expect(report.compactedNodes).toBe(1); + expect(report.errors).toBe(0); + + const fat = await kv.get("mem:graph:nodes", "fat"); + expect(fat!.sourceObservationIds).toEqual(["o4", "o5", "o6"]); + }); + + it("is idempotent: a second pass finds nothing left", async () => { + await run({ dryRun: false, compactSourceIds: true }); + const second = await run({ dryRun: true }); + + expect(second.staleNodes).toBe(0); + expect(second.danglingEdges).toBe(0); + expect(second.oversizedNodes).toBe(0); + expect(second.droppableSourceIds).toBe(0); + }); +}); From a4e367cacd5314f6f688bc030c3342497615e6ae Mon Sep 17 00:00:00 2001 From: reddeye1337 Date: Mon, 7 Sep 2026 22:09:15 +0300 Subject: [PATCH 3/4] feat(graph): drop duplicate merging, reconcile the snapshot after compaction Review on this PR raised five issues. Four were on the duplicate-merge path: an uncapped union on the keeper, losers marked for deletion before the keeper write succeeded, edge-key collisions when two duplicates both had a same-type edge to the same node, and degree reconciliation. That path was opt-in and never exercised against a real corpus, and designing merge semantics properly does not belong inside a PR about collection and compaction. It is removed; the survey still reports how many duplicate names exist, which is the useful half. The remaining two applied to what this PR does ship, and are fixed: The snapshot serves /graph/query on the no-argument path, so compacted rows have to be reflected there. Previously a compaction run reported success while queries kept returning the provenance it had just trimmed. The snapshot was reconciled against the delete candidate map rather than the rows actually deleted, so one failed delete in a run could drop a still-persisted row out of snapshot-backed queries. Successful deletions are now tracked separately and drive the filter. Signed-off-by: reddeye1337 --- src/functions/graph-prune.ts | 124 ++++++++++++++--------------------- src/triggers/api.ts | 1 - 2 files changed, 49 insertions(+), 76 deletions(-) diff --git a/src/functions/graph-prune.ts b/src/functions/graph-prune.ts index c95e2aea1e..93ffdea643 100644 --- a/src/functions/graph-prune.ts +++ b/src/functions/graph-prune.ts @@ -40,14 +40,13 @@ export interface GraphPruneReport { staleEdges: number; danglingEdges: number; supersededEdges: number; + /** Reported by the survey only; merging duplicates is not part of this function. */ duplicateNodes: number; oversizedNodes: number; oversizedEdges: number; droppableSourceIds: number; compactedNodes: number; compactedEdges: number; - mergedNodes: number; - repointedEdges: number; deletedNodes: number; deletedEdges: number; errors: number; @@ -66,7 +65,6 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { async (data?: { dryRun?: boolean; supersededOlderThanDays?: number; - mergeDuplicateNames?: boolean; compactSourceIds?: boolean; }): Promise => { const started = Date.now(); @@ -79,7 +77,6 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { data.supersededOlderThanDays >= 0 ? data.supersededOlderThanDays : DEFAULT_SUPERSEDED_RETENTION_DAYS; - const mergeDuplicates = data?.mergeDuplicateNames === true; const compactSources = data?.compactSourceIds === true; const maxSourceIds = getGraphMaxSourceIds(); const cutoff = Date.now() - retentionDays * 24 * 60 * 60 * 1000; @@ -104,8 +101,6 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { droppableSourceIds: 0, compactedNodes: 0, compactedEdges: 0, - mergedNodes: 0, - repointedEdges: 0, deletedNodes: 0, deletedEdges: 0, errors: 0, @@ -199,16 +194,17 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { // ---- compact provenance ---------------------------------------------- // Keeps the newest ids, matching the cap the merge path now applies // at write time. Nothing else on the row is touched. + const compactedNodeRows: GraphNode[] = []; + const compactedEdgeRows: GraphEdge[] = []; if (compactSources) { for (const node of oversizedNodeRows) { if (nodesToDelete.has(node.id)) continue; try { const trimmed = node.sourceObservationIds.slice(-maxSourceIds); - await kv.set(KV.graphNodes, node.id, { - ...node, - sourceObservationIds: trimmed, - }); - liveNodes.set(node.id, { ...node, sourceObservationIds: trimmed }); + const compacted = { ...node, sourceObservationIds: trimmed }; + await kv.set(KV.graphNodes, node.id, compacted); + liveNodes.set(node.id, compacted); + compactedNodeRows.push(compacted); report.compactedNodes++; } catch (err) { report.errors++; @@ -221,10 +217,12 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { for (const edge of oversizedEdgeRows) { if (edgesToDelete.has(edge.id)) continue; try { - await kv.set(KV.graphEdges, edge.id, { + const compacted = { ...edge, sourceObservationIds: edge.sourceObservationIds.slice(-maxSourceIds), - }); + }; + await kv.set(KV.graphEdges, edge.id, compacted); + compactedEdgeRows.push(compacted); report.compactedEdges++; } catch (err) { report.errors++; @@ -236,68 +234,15 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { } } - // ---- merge duplicates ------------------------------------------------- - const idRemap = new Map(); - if (mergeDuplicates) { - for (const group of duplicateGroups) { - const [keeper, ...losers] = group; - const obsIds = new Set(keeper.sourceObservationIds ?? []); - for (const loser of losers) { - for (const obsId of loser.sourceObservationIds ?? []) { - obsIds.add(obsId); - } - idRemap.set(loser.id, keeper.id); - nodesToDelete.set(loser.id, loser); - } - const merged: GraphNode = { - ...keeper, - sourceObservationIds: Array.from(obsIds), - updatedAt: new Date().toISOString(), - }; - try { - await kv.set(KV.graphNodes, keeper.id, merged); - await kv.set(KV.graphNameIndex, nameIndexKey(keeper.type, keeper.name), keeper.id); - liveNodes.set(keeper.id, merged); - report.mergedNodes += losers.length; - } catch (err) { - report.errors++; - logger.warn("Graph prune merge failed", { - nodeId: keeper.id, - error: err instanceof Error ? err.message : String(err), - }); - } - } - - // Repoint surviving edges off the merged-away nodes. An edge - // that collapses to a self-loop is dropped rather than kept. - for (const edge of allEdges) { - if (!edge?.id || edgesToDelete.has(edge.id)) continue; - const src = idRemap.get(edge.sourceNodeId) ?? edge.sourceNodeId; - const tgt = idRemap.get(edge.targetNodeId) ?? edge.targetNodeId; - if (src === edge.sourceNodeId && tgt === edge.targetNodeId) continue; - if (src === tgt) { - markEdge(edge); - continue; - } - try { - await kv.delete(KV.graphEdgeKey, edgeIndexKey(edge.sourceNodeId, edge.targetNodeId, edge.type)); - const repointed: GraphEdge = { ...edge, sourceNodeId: src, targetNodeId: tgt }; - await kv.set(KV.graphEdges, edge.id, repointed); - await kv.set(KV.graphEdgeKey, edgeIndexKey(src, tgt, edge.type), edge.id); - report.repointedEdges++; - } catch (err) { - report.errors++; - logger.warn("Graph prune repoint failed", { - edgeId: edge.id, - error: err instanceof Error ? err.message : String(err), - }); - } - } - } - // ---- delete ---------------------------------------------------------- // Edges first: deleting a node before its edges would turn those // edges into dangling rows if the run is interrupted. + // Reconcile the snapshot against rows that were really removed. A + // failed delete left in the candidate map would otherwise drop a row + // that is still persisted out of snapshot-backed queries. + const deletedNodeIds = new Set(); + const deletedEdgeIds = new Set(); + for (const edge of edgesToDelete.values()) { try { await kv.delete(KV.graphEdges, edge.id); @@ -306,6 +251,7 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { // Only clear the index entry if it still points at this edge; // a newer edge may have claimed the same endpoint triple. if (indexed === edge.id) await kv.delete(KV.graphEdgeKey, key); + deletedEdgeIds.add(edge.id); report.deletedEdges++; } catch (err) { report.errors++; @@ -323,6 +269,7 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { const key = nameIndexKey(node.type, node.name); const indexed = await kv.get(KV.graphNameIndex, key).catch(() => null); if (indexed === node.id) await kv.delete(KV.graphNameIndex, key); + deletedNodeIds.add(node.id); report.deletedNodes++; } catch (err) { report.errors++; @@ -336,14 +283,41 @@ export function registerGraphPruneFunction(sdk: ISdk, kv: StateKV): void { // Keep the precomputed snapshot honest. Its per-type breakdown // will drift, so flag it dirty; the totals are corrected here // because /graph/stats reports them directly. - if (report.deletedNodes > 0 || report.deletedEdges > 0) { + if ( + report.deletedNodes > 0 || + report.deletedEdges > 0 || + report.compactedNodes > 0 || + report.compactedEdges > 0 + ) { try { const snap = await kv.get(KV.graphSnapshot, SNAPSHOT_KEY); if (snap) { + // The snapshot serves /graph/query on the no-argument path, so a + // compacted row has to be reflected here too. Otherwise the run + // reports success while queries keep returning the provenance it + // just trimmed. + const compactedNodeIds = new Map( + compactedNodeRows.map((n) => [n.id, n.sourceObservationIds]), + ); + const compactedEdgeIds = new Map( + compactedEdgeRows.map((e) => [e.id, e.sourceObservationIds]), + ); await kv.set(KV.graphSnapshot, SNAPSHOT_KEY, { ...snap, - topNodes: (snap.topNodes ?? []).filter((n) => !nodesToDelete.has(n.id)), - topEdges: (snap.topEdges ?? []).filter((e) => !edgesToDelete.has(e.id)), + topNodes: (snap.topNodes ?? []) + .filter((n) => !deletedNodeIds.has(n.id)) + .map((n) => + compactedNodeIds.has(n.id) + ? { ...n, sourceObservationIds: compactedNodeIds.get(n.id)! } + : n, + ), + topEdges: (snap.topEdges ?? []) + .filter((e) => !deletedEdgeIds.has(e.id)) + .map((e) => + compactedEdgeIds.has(e.id) + ? { ...e, sourceObservationIds: compactedEdgeIds.get(e.id)! } + : e, + ), stats: { ...snap.stats, totalNodes: Math.max(0, snap.stats.totalNodes - report.deletedNodes), diff --git a/src/triggers/api.ts b/src/triggers/api.ts index c3d06ce169..f2ba6885d8 100644 --- a/src/triggers/api.ts +++ b/src/triggers/api.ts @@ -1585,7 +1585,6 @@ export function registerApiTriggers( req: ApiRequest<{ dryRun?: boolean; supersededOlderThanDays?: number; - mergeDuplicateNames?: boolean; compactSourceIds?: boolean; }>, ): Promise => { From 30a9be62d8e5a9ad562f3050bba68f093fe1af25 Mon Sep 17 00:00:00 2001 From: reddeye1337 Date: Mon, 7 Sep 2026 22:37:22 +0300 Subject: [PATCH 4/4] docs(graph-prune): drop the stale reference to duplicate merging The header still described mergeDuplicateNames as an opt-in flag after the merge path was removed from this PR. Signed-off-by: reddeye1337 --- src/functions/graph-prune.ts | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/src/functions/graph-prune.ts b/src/functions/graph-prune.ts index 93ffdea643..2146bb00f5 100644 --- a/src/functions/graph-prune.ts +++ b/src/functions/graph-prune.ts @@ -26,9 +26,9 @@ import { getGraphMaxSourceIds } from "../config.js"; // older than the retention cutoff, i.e. history that has // already been replaced by a newer revision. // -// Duplicate-name merging is real cleanup but it rewrites live rows, so -// it stays opt-in (`mergeDuplicateNames: true`) and reports its count -// in dry-run first. +// Duplicate names are counted by the survey but not merged here: merging +// rewrites live rows rather than removing dead ones, and needs its own +// design for edge-key collisions and two-phase keeper writes. const DEFAULT_SUPERSEDED_RETENTION_DAYS = 90;