/** * Debounced embedding scheduler for web UI. * Accumulates docIds from sync operations and runs embedding after debounce. * * @module src/serve/embed-scheduler */ import type { Database } from "bun:sqlite"; import type { EmbeddingPort } from "../llm/types"; import type { VectorIndexPort } from "../store/vector"; import { embedBacklog } from "../embed"; import { withBackgroundInference, withOwnedInferenceScope, } from "../llm/inference-scope"; import { createVectorStatsPort } from "../store/vector"; // ───────────────────────────────────────────────────────────────────────────── // Constants // ───────────────────────────────────────────────────────────────────────────── const DEBOUNCE_MS = 30_000; // 30 seconds const MAX_WAIT_MS = 300_000; // 5 minutes const BATCH_SIZE = 32; // ───────────────────────────────────────────────────────────────────────────── // Types // ───────────────────────────────────────────────────────────────────────────── export interface EmbedSchedulerState { pendingDocCount: number; running: boolean; nextRunAt?: number; lastRunAt?: number; lastResult?: EmbedResult; } export interface EmbedResult { embedded: number; errors: number; } export interface EmbedSchedulerDeps { db: Database; /** Getter for current embed port (survives context reloads) */ getEmbedPort: () => EmbeddingPort | null; /** Getter for current vector index (survives context reloads) */ getVectorIndex: () => VectorIndexPort | null; /** Getter for current model URI (survives preset changes) */ getModelUri: () => string; onEmbedded?: (result: EmbedResult) => void; embedBacklogFn?: typeof embedBacklog; } // ───────────────────────────────────────────────────────────────────────────── // Scheduler // ───────────────────────────────────────────────────────────────────────────── export interface EmbedScheduler { /** Called after sync with list of changed doc IDs (docid not id) */ notifySyncComplete(docIds: string[]): void; /** Force immediate embed (for Cmd+S). Returns null if no embedPort. */ triggerNow(): Promise; /** Get current state (for debugging/status) */ getState(): EmbedSchedulerState; /** Cleanup on server shutdown */ dispose(): Promise; /** Stop new turns while allowing the current pass to settle before cancellation. */ stop?(): Promise; } /** * Create an embed scheduler for debounced background embedding. * Uses getters to resolve dependencies at execution time (survives context reloads). */ export function createEmbedScheduler(deps: EmbedSchedulerDeps): EmbedScheduler { const { db, getEmbedPort, getVectorIndex, getModelUri, onEmbedded, embedBacklogFn = embedBacklog, } = deps; // State let pendingCount = 0; // Track pending triggers (not actual docIds - we embed full backlog) let timer: ReturnType | null = null; let running = false; let needsRerun = false; let firstPendingAt: number | null = null; let nextRunAt: number | null = null; // Accurate timer due time let disposed = false; let currentRun: Promise | null = null; let lastRunAt: number | null = null; let lastResult: EmbedResult | null = null; const controller = new AbortController(); const stats = createVectorStatsPort(db); /** * Run embedding for pending docs using shared helper. * Uses global backlog - we don't filter by docIds since: * 1. Backlog query is already efficient (only unembedded chunks) * 2. Filtering by docId would require joining through documents table * 3. Simpler to just embed all backlog when triggered */ async function runEmbed(): Promise { // Resolve dependencies at execution time (survives context reloads) const embedPort = getEmbedPort(); const vectorIndex = getVectorIndex(); const modelUri = getModelUri(); if (!embedPort || !vectorIndex) { return { embedded: 0, errors: 0 }; } let result: Awaited>; try { result = await withBackgroundInference(() => withOwnedInferenceScope({ signal: controller.signal }, () => embedBacklogFn({ statsPort: stats, embedPort, vectorIndex, modelUri, batchSize: BATCH_SIZE, identityStillCurrent: () => getEmbedPort() === embedPort && getVectorIndex() === vectorIndex && getModelUri() === modelUri, }) ) ); } catch (cause) { if (controller.signal.aborted) return { embedded: 0, errors: 0 }; throw cause; } if (!result.ok) { needsRerun = true; console.error("[embed-scheduler] Embed failed:", result.error.message); return { embedded: 0, errors: 0 }; } if ( getEmbedPort() !== embedPort || getVectorIndex() !== vectorIndex || getModelUri() !== modelUri ) needsRerun = true; if ((result.value.contentionErrors ?? 0) > 0 || result.value.errors > 0) { // Provider failures and contended checkpoints remain durably pending. console.error( `[embed-scheduler] ${result.value.errors} embedding errors, ${result.value.contentionErrors ?? 0} contended writes; rescheduling` ); needsRerun = true; } if (result.value.embedded > 0) onEmbedded?.(result.value); return result.value; } /** * Schedule or reschedule the debounced embed run. */ function scheduleRun(): void { if (disposed) { return; } // If currently running, mark for rerun instead of scheduling if (running) { needsRerun = true; return; } // Calculate delay const now = Date.now(); let delay = DEBOUNCE_MS; // Check max-wait if (firstPendingAt !== null) { const elapsed = now - firstPendingAt; if (elapsed >= MAX_WAIT_MS) { // Max wait reached, run immediately delay = 0; } else { // Don't exceed max wait delay = Math.min(delay, MAX_WAIT_MS - elapsed); } } // Clear existing timer if (timer) { clearTimeout(timer); } // Track accurate due time nextRunAt = now + delay; timer = setTimeout(() => { nextRunAt = null; void executeRun(); }, delay); } /** * Execute the embed run with concurrency guard. */ function executeRun(): Promise { if (disposed || running) { needsRerun = true; return Promise.resolve(null); } const operation = (async (): Promise => { running = true; timer = null; nextRunAt = null; // Clear pending state at START so new notifications accumulate pendingCount = 0; firstPendingAt = null; let result: EmbedResult; try { result = await runEmbed(); lastRunAt = Date.now(); lastResult = result; } finally { running = false; } // Check if we need to rerun (notifications arrived while running) // Must be AFTER running=false so scheduleRun() actually schedules if ((needsRerun || pendingCount > 0) && !disposed) { needsRerun = false; // Set firstPendingAt if we have pending work if (pendingCount > 0 && firstPendingAt === null) { firstPendingAt = Date.now(); } scheduleRun(); } return result; })(); currentRun = operation; void operation.then( () => { if (currentRun === operation) currentRun = null; }, () => { if (currentRun === operation) currentRun = null; } ); return operation; } const stop = async (): Promise => { disposed = true; if (timer) clearTimeout(timer); timer = null; nextRunAt = null; await currentRun; }; return { stop, notifySyncComplete(docIds: string[]): void { // Resolve embedPort at call time to check availability if (disposed || !getEmbedPort()) { return; } // Track first pending time for max-wait if (pendingCount === 0 && firstPendingAt === null) { firstPendingAt = Date.now(); } // Count pending triggers (we don't track individual docIds) pendingCount += docIds.length; // Schedule/reschedule debounced run (or mark needsRerun if running) scheduleRun(); }, async triggerNow(): Promise { if (disposed || !getEmbedPort()) { return null; } // Cancel pending timer if (timer) { clearTimeout(timer); timer = null; nextRunAt = null; } // If already running, mark for rerun if (running) { needsRerun = true; return { embedded: 0, errors: 0 }; } return executeRun(); }, getState(): EmbedSchedulerState { const state: EmbedSchedulerState = { pendingDocCount: pendingCount, running, }; // Use accurate nextRunAt from timer scheduling if (nextRunAt !== null) { state.nextRunAt = nextRunAt; } if (lastRunAt !== null) { state.lastRunAt = lastRunAt; } if (lastResult) { state.lastResult = lastResult; } return state; }, async dispose(): Promise { controller.abort(); await stop(); }, }; }