diff --git a/.coderabbit.yaml b/.coderabbit.yaml new file mode 100644 index 00000000..0bf7448c --- /dev/null +++ b/.coderabbit.yaml @@ -0,0 +1,6 @@ +# yaml-language-server: $schema=https://coderabbit.ai/integrations/schema.v2.json +reviews: + auto_review: + # Review stacked pull requests too, not only those that target main. + base_branches: + - ".*" diff --git a/AGENTS.md b/AGENTS.md index fae060de..433918bc 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -29,6 +29,7 @@ OMMS (npm `om-memory-system`) is a memory plugin for AI coding agents. One share Keep these boundaries: - `src/core/` and `src/services/` must not import `@opencode-ai/*`, `@earendil-works/*`, or `src/adapters/*`. +- `src/importer/` must not import `src/adapters/*`. Each host's history reader lives in `src/importer/`, and adapters import from it. - `tests/host-neutral-capture-boundary.test.ts` and `tests/pi-adapter-boundary.test.ts` enforce that rule. - An adapter must not import the other host's adapter modules. - Load host SDKs and heavy modules with dynamic `import()`. `tests/plugin-bundle-boundary.test.ts` checks the plugin bundle. diff --git a/docs/tdr/012-tag-migration-touches-only-untagged-memories.md b/docs/tdr/012-tag-migration-touches-only-untagged-memories.md new file mode 100644 index 00000000..f3ff04bd --- /dev/null +++ b/docs/tdr/012-tag-migration-touches-only-untagged-memories.md @@ -0,0 +1,55 @@ +# TDR-012: Tag migration touches only untagged memories + +**Date:** 2026-09-28 +**Status:** Proposed +**Deciders:** OMMS maintainers +**Tags:** web-ui, embeddings, migration + +## Context + +The web UI's **Memory Tagging Migration** dialog reports how many memories have no tags ("Found 1 memories needing technical tags") and offers **Start Migration**. On a store with 2,236 memories and one untagged memory, the run showed `/2236` and worked through every memory. + +### Root Cause Analysis + +`handleDetectTagMigration` counted only rows with empty `tags`, but `handleRunTagMigrationBatch` loaded `SELECT * FROM memories` from every project shard and re-embedded each row's content and tags, calling the model only for the untagged ones. Tagged memories got identical vectors back, so no data was lost, but the run took far longer than the dialog implied and ran the local embedding model thousands of times. A second defect: a memory that threw an error did not advance `processed`, so every later batch restarted at the same memory and the run could never finish. + +## Decision + +The run builds its work list once, when it starts, from `SELECT id FROM memories WHERE tags IS NULL OR tags = ''` in each project shard, so the total matches the dialog's count. Each memory is re-read by id; one that gained tags or was deleted since is passed over. Only a memory that receives tags is re-embedded. Its tags and both vectors are saved in one `UPDATE`, so a failed embedding leaves it untagged and a later run tries it again. A failure is recorded in `errors` and the run always moves on. When a run completes, the next run starts from a fresh list of whatever is still untagged. + +## Consequences + +### Positive + +- The dialog's count and the run's total agree, and tagged memories are never rewritten. +- A failing memory cannot stall the run; it stays untagged for a later run. + +### Negative + +- A memory whose model call fails stays untagged, so the dialog reappears until a run succeeds. + +### Neutral + +- Migration state is still held in the web server's memory; restarting the server restarts the run. + +## Alternatives Considered + +| Option | Rejected Because | +| ----------------------------------- | ---------------------------------------------------------------- | +| Keep re-embedding everything | Wastes time and CPU, and contradicts the dialog's count | +| Select untagged rows on every batch | Rows that fail stay untagged and would be selected again forever | + +## How to Recognise / Handle This Again + +1. A migration or backfill total is much larger than the count shown before it starts. +2. Compare the detect query with the query the run iterates over. +3. Build the run's list from the same filter as the detect step, and always advance past failures. + +## Revisit Triggers + +- Tags or tag vectors change format and existing tagged memories need re-embedding on purpose. + +## References + +- `src/services/api-handlers.ts` (`handleDetectTagMigration`, `handleRunTagMigrationBatch`) +- `tests/tag-migration.test.ts` diff --git a/docs/tdr/README.md b/docs/tdr/README.md index a0c0f427..06fd430c 100644 --- a/docs/tdr/README.md +++ b/docs/tdr/README.md @@ -33,6 +33,7 @@ TDRs capture **implementation-level technical decisions** such as platform-speci | [009](./009-protect-capture-traces-on-windows.md) | Protect capture traces with Windows access-control lists | Accepted | 2026-09-27 | | [010](./010-match-windows-native-import-source-paths.md) | Match Windows import-source tests to native canonical paths | Proposed | 2026-09-27 | | [011](./011-use-execfilesync-in-windows-git-wrapper-test.md) | Use execFileSync in the Windows Git wrapper test | Proposed | 2026-09-27 | +| [012](./012-tag-migration-touches-only-untagged-memories.md) | Tag migration touches only untagged memories | Proposed | 2026-09-28 | ## Status values diff --git a/src/adapters/opencode/user-prompt.ts b/src/adapters/opencode/user-prompt.ts index 91c03217..2ba304a8 100644 --- a/src/adapters/opencode/user-prompt.ts +++ b/src/adapters/opencode/user-prompt.ts @@ -1,15 +1,8 @@ +import { isStructuredSummaryPromptMessage } from "../../core/internal-prompt.js"; import { isInternalStructuredSession } from "../../services/ai/opencode-provider.js"; import { userPromptManager } from "../../services/user-prompt/user-prompt-manager.js"; -export function isStructuredSummaryPromptMessage(userMessage: string): boolean { - // This is the plugin's own structured-summary or profile-analysis request. - // OpenCode echoes it through chat.message like a normal user message, but - // capturing it would create self-referential memories / an infinite learning loop. - if (userMessage.includes("# User Profile Analysis")) { - return true; - } - return userMessage.includes("Analyze this conversation.") && userMessage.includes('type="skip"'); -} +export { isStructuredSummaryPromptMessage }; /** True when a prompt is omms's own internal traffic and must not be recorded or searched. */ export function isInternalPrompt(sessionID: string, userMessage: string): boolean { diff --git a/src/adapters/pi/capture.ts b/src/adapters/pi/capture.ts index 4c2b13c7..18333f90 100644 --- a/src/adapters/pi/capture.ts +++ b/src/adapters/pi/capture.ts @@ -3,7 +3,7 @@ import { captureConversation } from "../../core/capture.js"; import type { CaptureSummaryProvider } from "../../core/host.js"; import { log } from "../../services/logger.js"; import { memoryClient } from "../../services/client.js"; -import { extractPiConversation, type PiSessionEntry } from "./conversation.js"; +import { extractPiConversation, type PiSessionEntry } from "../../importer/pi-conversation.js"; export interface PiCaptureState { /** User entry IDs terminally handled by a settled capture for this session. */ diff --git a/src/adapters/pi/extension.ts b/src/adapters/pi/extension.ts index 46af1ea4..e9f827d1 100644 --- a/src/adapters/pi/extension.ts +++ b/src/adapters/pi/extension.ts @@ -7,7 +7,7 @@ import { getLanguageName } from "../../services/language-detector.js"; import { log } from "../../services/logger.js"; import { memoryClient } from "../../services/client.js"; import { capturePiSettledWorkUnit, createPiCaptureState } from "./capture.js"; -import type { PiSessionEntry } from "./conversation.js"; +import type { PiSessionEntry } from "../../importer/pi-conversation.js"; import { createPiLiveModels } from "./live-model.js"; import { registerPiHistoryImportCommand } from "./import-command.js"; import { performPiProfileLearning } from "./profile.js"; diff --git a/src/core/internal-prompt.ts b/src/core/internal-prompt.ts new file mode 100644 index 00000000..c94cda60 --- /dev/null +++ b/src/core/internal-prompt.ts @@ -0,0 +1,11 @@ +/** + * True when a prompt is omms's own structured-summary or profile-analysis + * request. Both hosts echo it back like a user message; recording or learning + * from it would create self-referential memories and a learning loop. + */ +export function isStructuredSummaryPromptMessage(userMessage: string): boolean { + if (userMessage.includes("# User Profile Analysis")) { + return true; + } + return userMessage.includes("Analyze this conversation.") && userMessage.includes('type="skip"'); +} diff --git a/src/importer/importer.ts b/src/importer/importer.ts index aa882874..248aa596 100644 --- a/src/importer/importer.ts +++ b/src/importer/importer.ts @@ -13,7 +13,7 @@ import { extractPiConversationWindows, type PiConversationWindow, type PiSessionEntry, -} from "../adapters/pi/conversation.js"; +} from "./pi-conversation.js"; import { discoverPiSessions } from "./discovery.js"; import { resolveImportProject } from "./import-project.js"; import { PiImportLedger, importLedgerDbPath, type ImportLedgerRow } from "./ledger.js"; diff --git a/src/adapters/pi/conversation.ts b/src/importer/pi-conversation.ts similarity index 98% rename from src/adapters/pi/conversation.ts rename to src/importer/pi-conversation.ts index 13d630ab..89bd1832 100644 --- a/src/adapters/pi/conversation.ts +++ b/src/importer/pi-conversation.ts @@ -1,4 +1,4 @@ -import type { CaptureConversation, CaptureToolCall } from "../../core/host.js"; +import type { CaptureConversation, CaptureToolCall } from "../core/host.js"; /** * Minimal structural view of Pi session entries (see Pi session-format docs). diff --git a/src/importer/profile-import.ts b/src/importer/profile-import.ts index abe9ba16..f15f57d0 100644 --- a/src/importer/profile-import.ts +++ b/src/importer/profile-import.ts @@ -1,4 +1,4 @@ -import { isInternalPrompt } from "../adapters/opencode/user-prompt.js"; +import { isStructuredSummaryPromptMessage } from "../core/internal-prompt.js"; import { analyzeProfile, type ModelPort } from "../core/profile-analysis.js"; import { isFullyPrivate, stripPrivateContent } from "../services/privacy.js"; import { getTags } from "../services/tags.js"; @@ -77,7 +77,8 @@ export async function importProfileFromHistory( if ( !prompt || isFullyPrivate(unit.userPrompt) || - isInternalPrompt(session.sessionId, prompt) + // OpenCode history already leaves out omms's own capture sessions by title. + isStructuredSummaryPromptMessage(prompt) ) { continue; } diff --git a/src/importer/session-loader.ts b/src/importer/session-loader.ts index 63323379..8feed3d5 100644 --- a/src/importer/session-loader.ts +++ b/src/importer/session-loader.ts @@ -1,5 +1,5 @@ import { SessionManager } from "@earendil-works/pi-coding-agent"; -import type { PiSessionEntry } from "../adapters/pi/conversation.js"; +import type { PiSessionEntry } from "./pi-conversation.js"; /** * Loads a Pi session file for import through Pi's own exported session model diff --git a/src/services/api-handlers.ts b/src/services/api-handlers.ts index 730d8353..592d1d5e 100644 --- a/src/services/api-handlers.ts +++ b/src/services/api-handlers.ts @@ -151,7 +151,8 @@ export async function handleListMemories( tag?: string, page: number = 1, pageSize: number = 20, - includePrompts: boolean = true + includePrompts: boolean = true, + keyword?: string ): Promise>> { try { await ensureTursoReady(); @@ -211,11 +212,22 @@ export async function handleListMemories( }; }); - let timeline: any[] = memoriesWithType; + // A keyword filter keeps memories carrying that keyword (case-insensitive) + // and only the prompts linked to them. + const wanted = keyword?.trim().toLowerCase(); + const keptMemories = wanted + ? memoriesWithType.filter((m) => m.tags.some((t: string) => t.toLowerCase() === wanted)) + : memoriesWithType; + const keptIds = new Set(keptMemories.map((m) => m.id)); + + let timeline: any[] = keptMemories; if (includePrompts) { const projectPath = tag ? await getProjectPathFromTag(tag) : undefined; const prompts = await userPromptManager.getCapturedPrompts(projectPath); - const promptsWithType = prompts.map((p) => ({ + const visiblePrompts = wanted + ? prompts.filter((p) => p.linkedMemoryId && keptIds.has(p.linkedMemoryId)) + : prompts; + const promptsWithType = visiblePrompts.map((p) => ({ type: "prompt", id: p.id, sessionId: p.sessionId, @@ -224,7 +236,7 @@ export async function handleListMemories( projectPath: p.projectPath, linkedMemoryId: p.linkedMemoryId, })); - timeline = [...memoriesWithType, ...promptsWithType]; + timeline = [...keptMemories, ...promptsWithType]; } const linkedPairs = new Map(); @@ -1434,19 +1446,28 @@ interface MigrationProgress { errors: string[]; } -const migrationProgress: MigrationProgress = { +const idleMigration = (): MigrationProgress => ({ processed: 0, total: 0, currentBatch: 0, totalBatches: 0, isComplete: true, errors: [], -}; +}); + +let migrationProgress: MigrationProgress = idleMigration(); +/** The untagged memories this migration run covers, fixed when the run starts. */ +let migrationQueue: Array<{ id: string; dbPath: string }> = []; export async function handleGetTagMigrationProgress(): Promise> { return { success: true, data: migrationProgress }; } +/** + * Tag the memories that have no tags, a few per request. Only untagged + * memories are touched: tagged ones keep their vectors. A memory that fails + * is recorded and passed over, so one bad row cannot stall the run. + */ export async function handleRunTagMigrationBatch( batchSize: number = 5 ): Promise> { @@ -1459,91 +1480,76 @@ export async function handleRunTagMigrationBatch( iterationTimeout: 30000, }); const provider = AIProviderFactory.createProvider(CONFIG.memoryProvider, providerConfig); - const projectShards = await tursoShardManager.getAllShards("project", ""); - - const allMemories: { memory: any; shard: any }[] = []; - for (const shard of projectShards) { - const db = await tursoConnectionManager.getConnection(shard.dbPath); - const memories = await db.all("SELECT * FROM memories"); - for (const m of memories) { - allMemories.push({ memory: m, shard }); + if (migrationProgress.isComplete) { + migrationQueue = []; + for (const shard of await tursoShardManager.getAllShards("project", "")) { + const db = await tursoConnectionManager.getConnection(shard.dbPath); + const rows = await db.all( + "SELECT id FROM memories WHERE tags IS NULL OR tags = '' ORDER BY id" + ); + for (const row of rows) migrationQueue.push({ id: String(row.id), dbPath: shard.dbPath }); } + migrationProgress = { + ...idleMigration(), + total: migrationQueue.length, + totalBatches: Math.ceil(migrationQueue.length / batchSize), + isComplete: migrationQueue.length === 0, + }; } - if (migrationProgress.total === 0) { - migrationProgress.total = allMemories.length; - migrationProgress.totalBatches = Math.ceil(allMemories.length / batchSize); - migrationProgress.isComplete = false; - } - - const startIdx = migrationProgress.processed; - const endIdx = Math.min(startIdx + batchSize, allMemories.length); - - for (let i = startIdx; i < endIdx; i++) { - const item = allMemories[i]; - if (!item) continue; - const { memory: m, shard } = item; - const db = await tursoConnectionManager.getConnection(shard.dbPath); - + const batch = migrationQueue.slice( + migrationProgress.processed, + migrationProgress.processed + batchSize + ); + for (const item of batch) { try { - let currentTags = m.tags - ? m.tags - .split(",") - .map((t: string) => t.trim().toLowerCase()) - .filter((t: string) => t) - : []; - - if (currentTags.length === 0) { - const prompt = `Generate 2-4 short technical tags for this memory content:\n\n${m.content}\n\nReturn ONLY a comma-separated list of tags.`; - const result = await provider.executeToolCall( - "You are a technical tagger.", - prompt, - { - type: "function", - function: { - name: "save_tags", - description: "Save generated tags", - parameters: { - type: "object", - properties: { tags: { type: "array", items: { type: "string" } } }, - required: ["tags"], - }, + const db = await tursoConnectionManager.getConnection(item.dbPath); + const m: any = await db.get("SELECT * FROM memories WHERE id = ?", [item.id]); + // Tagged or deleted since the run started: nothing to do. + if (!m || (m.tags && String(m.tags).trim())) continue; + const prompt = `Generate 2-4 short technical tags for this memory content:\n\n${m.content}\n\nReturn ONLY a comma-separated list of tags.`; + const result = await provider.executeToolCall( + "You are a technical tagger.", + prompt, + { + type: "function", + function: { + name: "save_tags", + description: "Save generated tags", + parameters: { + type: "object", + properties: { tags: { type: "array", items: { type: "string" } } }, + required: ["tags"], }, }, - `migration_${m.id}` - ); - if (result.success && result.data?.tags) { - currentTags = result.data.tags; - await db.run("UPDATE memories SET tags = ? WHERE id = ?", [ - currentTags.join(","), - m.id, - ]); - } + }, + `migration_${m.id}` + ); + const tags: string[] = Array.isArray(result.data?.tags) + ? result.data.tags.map((t: unknown) => String(t).trim().toLowerCase()).filter(Boolean) + : []; + if (!result.success || tags.length === 0) { + throw new Error(result.error ?? "The model returned no tags"); } - const vector = await embeddingService.embedWithTimeout(m.content, { task: "document" }); - const tagsVector = currentTags.length - ? await embeddingService.embedWithTimeout(formatTagsForEmbedding(currentTags), { - task: "document", - }) - : undefined; - await tursoVectorSearch.updateVector(db, m.id, vector, tagsVector); - - migrationProgress.processed++; + const tagsVector = await embeddingService.embedWithTimeout(formatTagsForEmbedding(tags), { + task: "document", + }); + // One statement, so a memory never ends up tagged with stale vectors. + await tursoVectorSearch.updateVector(db, m.id, vector, tagsVector, tags.join(",")); } catch (e) { const errorMsg = String(e); migrationProgress.errors.push(errorMsg); - log("Migration error for memory", { id: m.id, error: errorMsg }); + log("Migration error for memory", { id: item.id, error: errorMsg }); + } finally { + migrationProgress.processed++; } } migrationProgress.currentBatch++; const hasMore = migrationProgress.processed < migrationProgress.total; - - if (!hasMore) { - migrationProgress.isComplete = true; - } + if (!hasMore) migrationProgress.isComplete = true; return { success: true, diff --git a/src/services/turso/vector-search.ts b/src/services/turso/vector-search.ts index 32b34d66..03af1ae2 100644 --- a/src/services/turso/vector-search.ts +++ b/src/services/turso/vector-search.ts @@ -355,10 +355,16 @@ export class TursoVectorSearch { db: TursoDb, memoryId: string, vector: Float32Array, - tagsVector?: Float32Array + tagsVector?: Float32Array, + tags?: string ): Promise { const contentVector = vectorToJson(vector); - if (tagsVector) { + if (tagsVector && tags !== undefined) { + await db.execute( + `UPDATE memories SET tags = ?, vector = vector32(?), tags_vector = vector32(?) WHERE id = ?`, + [tags, contentVector, vectorToJson(tagsVector), memoryId] + ); + } else if (tagsVector) { await db.execute( `UPDATE memories SET vector = vector32(?), tags_vector = vector32(?) WHERE id = ?`, [contentVector, vectorToJson(tagsVector), memoryId] diff --git a/tests/pi-adapter-boundary.test.ts b/tests/pi-adapter-boundary.test.ts index f76e6217..5d9277ef 100644 --- a/tests/pi-adapter-boundary.test.ts +++ b/tests/pi-adapter-boundary.test.ts @@ -22,8 +22,8 @@ describe("Pi adapter boundary", () => { } }); - it("does not leak the Pi adapter into the shared core or services", () => { - for (const dir of ["src/core", "src/services", "src/types"]) { + it("keeps host adapters out of the shared core, services, and importer", () => { + for (const dir of ["src/core", "src/services", "src/types", "src/importer"]) { const entries = readdirSync(join(import.meta.dir, "..", dir), { recursive: true, }) as string[]; diff --git a/tests/pi-conversation.test.ts b/tests/pi-conversation.test.ts index 6a972887..a1e51298 100644 --- a/tests/pi-conversation.test.ts +++ b/tests/pi-conversation.test.ts @@ -1,5 +1,5 @@ import { describe, expect, it } from "bun:test"; -import { extractPiConversation, type PiSessionEntry } from "../src/adapters/pi/conversation.js"; +import { extractPiConversation, type PiSessionEntry } from "../src/importer/pi-conversation.js"; function userEntry(id: string, text: string): PiSessionEntry { return { diff --git a/tests/tag-migration.test.ts b/tests/tag-migration.test.ts new file mode 100644 index 00000000..2e16ae57 --- /dev/null +++ b/tests/tag-migration.test.ts @@ -0,0 +1,94 @@ +import { afterEach, expect, it } from "bun:test"; +import { mkdtempSync } from "node:fs"; +import { join } from "node:path"; +import { tmpdir } from "node:os"; +import { cleanupTursoTestDirectory } from "./turso-test-utils.js"; + +let baseDir: string; +afterEach(async () => { + await cleanupTursoTestDirectory(baseDir); +}); + +it("tags only untagged memories, keeps tagged vectors, and moves past a failure", async () => { + baseDir = mkdtempSync(join(tmpdir(), "omms-tag-migration-")); + const { CONFIG } = await import("../src/config.js"); + Object.assign(CONFIG, { + storagePath: baseDir, + embeddingDimensions: 768, + memoryProvider: "openai-chat", + memoryModel: "m", + memoryApiUrl: "https://x.invalid/v1", + memoryApiKey: "k", + }); + const { closeTursoAndInvalidateCaches } = await import("../src/services/turso/lifecycle.js"); + await closeTursoAndInvalidateCaches(); + const { tursoShardManager } = await import("../src/services/turso/shard-manager.js"); + const { tursoConnectionManager } = await import("../src/services/turso/connection-manager.js"); + const { tursoVectorSearch } = await import("../src/services/turso/vector-search.js"); + const { embeddingService } = await import("../src/services/embedding.js"); + const { AIProviderFactory } = await import("../src/services/ai/ai-provider-factory.js"); + const { handleDetectTagMigration, handleRunTagMigrationBatch } = + await import("../src/services/api-handlers.js"); + + const vector = new Float32Array(768); + vector[0] = 1; + const shard = await tursoShardManager.createShard("project", "a1b2c3d4e5f60718", 0); + const db = await tursoConnectionManager.getConnection(shard.dbPath); + for (const [id, tags] of [ + ["tagged", ["kept"]], + ["untagged-ok", undefined], + ["untagged-fails", undefined], + ["untagged-embed-fails", undefined], + ] as const) { + await tursoVectorSearch.insertVector(db, { + id, + content: `content of ${id}`, + vector, + containerTag: "omms_project_a1b2c3d4e5f60718", + createdAt: Date.now(), + updatedAt: Date.now(), + ...(tags ? { tags: tags.join(",") } : {}), + }); + } + + const embedded: string[] = []; + const originalEmbed = embeddingService.embedWithTimeout; + const originalCreate = AIProviderFactory.createProvider; + embeddingService.embedWithTimeout = (async (text: string) => { + if (text.includes("untagged-embed-fails")) throw new Error("embedding timed out"); + embedded.push(text); + return vector; + }) as typeof originalEmbed; + AIProviderFactory.createProvider = (() => ({ + executeToolCall: async (_s: string, prompt: string) => + prompt.includes("untagged-fails") + ? { success: false, error: "model refused" } + : { success: true, data: { tags: ["Bun", "sqlite"] } }, + })) as unknown as typeof originalCreate; + try { + expect((await handleDetectTagMigration()).data).toEqual({ needsMigration: true, count: 3 }); + const first = await handleRunTagMigrationBatch(1); + const second = await handleRunTagMigrationBatch(2); + expect(first.data).toEqual({ processed: 1, total: 3, hasMore: true }); + expect(second.data).toEqual({ processed: 3, total: 3, hasMore: false }); + + const rows = await db.all("SELECT id, tags FROM memories ORDER BY id"); + expect(rows.map((row) => [row.id, row.tags])).toEqual([ + ["tagged", "kept"], + // Tags are saved only together with fresh vectors, so a failed embedding stays untagged. + ["untagged-embed-fails", null], + ["untagged-fails", null], + ["untagged-ok", "bun,sqlite"], + ]); + // Only the memory that gained tags was re-embedded. + expect(embedded.some((text) => text.includes("content of tagged"))).toBe(false); + expect(embedded.some((text) => text.includes("untagged-ok"))).toBe(true); + + // A later run starts over with only what is still untagged. + const retry = await handleRunTagMigrationBatch(5); + expect(retry.data).toEqual({ processed: 2, total: 2, hasMore: false }); + } finally { + embeddingService.embedWithTimeout = originalEmbed; + AIProviderFactory.createProvider = originalCreate; + } +});