// --------------------------------------------------------------------------- // Memory Graph — Extraction job handler // // Wraps runGraphExtraction for the jobs worker. Handles both: // - Mid-conversation batch extraction (incremental, from checkpoint) // - End-of-conversation extraction (full transcript) // --------------------------------------------------------------------------- import type { AssistantConfig } from "../../../../../config/types.js"; import { getMemoryCheckpoint, setMemoryCheckpoint, } from "../../../../../persistence/checkpoints.js"; import { asString } from "../../../../../persistence/job-utils.js"; import type { MemoryJob } from "../../../../../persistence/jobs-store.js"; import { getLogger } from "../../logging.js"; import { runGraphExtraction } from "./extraction.js"; const log = getLogger("graph-extraction-job"); /** * Job handler for `graph_extract`. Runs incremental or full extraction * depending on whether a checkpoint exists for this conversation. * * Checkpoint key: `graph_extract::last_ts` * Value: epoch ms of the most recent message processed. * * Trigger sources: * - Indexer after batchSize messages (default 10) * - Indexer idle debounce (default 300s) * - Conversation dispose (end of conversation) */ export async function graphExtractJob( job: MemoryJob, config: AssistantConfig, ): Promise { const conversationId = asString(job.payload.conversationId); if (!conversationId) { return; } // Read checkpoint for incremental extraction const checkpointKey = `graph_extract:${conversationId}:last_ts`; const lastTs = getMemoryCheckpoint(checkpointKey); const afterTimestamp = lastTs ? parseInt(lastTs, 10) : undefined; const activeContextNodeIds = Array.isArray(job.payload.activeContextNodeIds) ? (job.payload.activeContextNodeIds as string[]) : undefined; try { const result = await runGraphExtraction(conversationId, config, { afterTimestamp, activeContextNodeIds, }); // Update checkpoint to the newest message actually processed — using // Date.now() could skip messages that arrived during extraction. if (result.lastProcessedTimestamp) { setMemoryCheckpoint(checkpointKey, String(result.lastProcessedTimestamp)); } log.info( { conversationId, incremental: !!afterTimestamp, ...result, }, "Graph extraction job complete", ); } catch (err) { log.error( { conversationId, err: err instanceof Error ? err.message : String(err) }, "Graph extraction job failed", ); throw err; } }