import { ConvexError, v } from "convex/values"; import { internal } from "./_generated/api"; import type { Doc } from "./_generated/dataModel"; import { internalAction, internalMutation, internalQuery, type MutationCtx, } from "./_generated/server"; /** * Append-only historical intelligence memory (#5694). * * The live seeder snapshots in Redis overwrite themselves every run, so * "what happened in this country three months ago" is unanswerable from * them. This module is the durable store behind those snapshots: seeders * append the events they published, reads come back either chronologically * (`timeline`) or by semantic similarity (`search`), and a daily cron ages * rows out (`prune`). * * Append-only in the steady state, with one deliberate exception: `retract` * (#5743) removes named events and tombstones their identity so the producing * seeder cannot re-add them. Nothing is ever updated in place. * * Everything here is `internal*` on purpose. The only ways in are the * secret-guarded HTTP routes in convex/http.ts — `/relay/intel-history` * (ingest, Railway seeders), `/relay/intel-history/retract` and * `/relay/intel-history/restore` (operator-driven retraction), * `/api/internal-intel-timeline` and `/api/internal-intel-search` (reads, * Vercel edge). No client-facing function touches this table. */ /** * Dimension of every stored and query vector. MUST equal `EMBED_DIMS` in * scripts/lib/brief-dedup-consts.mjs (openai/text-embedding-3-small@512) and * the `dimensions` of the `by_embedding` vector index in convex/schema.ts. * Changing it is a table migration plus a full re-embed, never an edit here — * see the schema comment on `intelHistory`. */ export const INTEL_HISTORY_EMBED_DIMS = 512; /** * Per-call ingest cap. Each record carries a 512-float vector (~4KB before * encoding), so this bounds both the mutation's write set and the request * body the relay route has to parse. Seeders chunk larger runs. */ export const INTEL_HISTORY_MAX_APPEND_RECORDS = 100; /** * Per-call cap on identifiers handed to `retract` / `restore` (#5743). * * Retraction is a hand-driven incident operation, not a bulk pipe: an operator * has identified specific poisoned or wrong rows and is removing them. The cap * keeps one call's write set bounded and makes a fat-fingered "retract * everything" fail at the boundary rather than half-succeed. */ export const INTEL_HISTORY_MAX_RETRACT_IDENTIFIERS = 100; /** Bound on one `listRetractions` page — a review read, not a bulk export. */ export const INTEL_HISTORY_MAX_RETRACTION_PAGE = 200; /** * Documented retention policy: history older than this is deleted by the * `intel-history-prune` daily cron. 180 days is the window the historical * comparisons are specified against ("versus six months ago"); it is also * what keeps the vector index — the expensive part of this table — bounded. */ export const INTEL_HISTORY_RETENTION_DAYS = 180; const RETENTION_MS = INTEL_HISTORY_RETENTION_DAYS * 24 * 60 * 60 * 1000; /** * Per-run delete cap, deliberately far below the apiPlanLimit prune's 500: * each row here carries a 512-float embedding, so 100 deletes stay well * inside Convex's per-mutation write limit. A larger backlog drains through * the self-reschedule below rather than over successive daily runs. */ const PRUNE_BATCH = 100; const TIMELINE_DEFAULT_LIMIT = 50; export const TIMELINE_MAX_LIMIT = 200; const SEARCH_DEFAULT_LIMIT = 20; export const SEARCH_MAX_LIMIT = 64; /** * When a read has to post-filter (a field that could not be pushed into the * index, or an occurredAt range on the vector path), fetch this multiple of * the caller's limit so the post-filter still has something to return. * Bounded by TIMELINE_MAX_SCAN / VECTOR_SEARCH_MAX_LIMIT — this widens the * read, it never unbounds it. */ const POST_FILTER_OVERFETCH = 4; /** Hard ceiling on documents a single `timeline` call may read. */ const TIMELINE_MAX_SCAN = 800; /** Convex's own ceiling on `vectorSearch` limit. */ const VECTOR_SEARCH_MAX_LIMIT = 256; /** The single pre-seeded OCC document used to serialize low-frequency appends. */ const APPEND_LOCK_KEY = "intel-history-append"; /** * Read the pre-seeded append lock or fail closed when deploy initialization * was skipped. An empty dedupe-key index range is not a document-backed OCC * dependency, so this document must be read and patched in every append * transaction before dedupe checks begin. */ async function readAppendLockOrThrow( ctx: MutationCtx, ): Promise> { const lock = await ctx.db .query("intelHistoryAppendLocks") .withIndex("by_lockKey", (q) => q.eq("lockKey", APPEND_LOCK_KEY)) .first(); if (!lock) { console.error( JSON.stringify({ breadcrumb: "intel_history_append_lock_not_seeded", lockKey: APPEND_LOCK_KEY, }), ); throw new ConvexError({ kind: "APPEND_LOCK_NOT_SEEDED" }); } return lock; } function clamp(value: number, min: number, max: number): number { if (!Number.isFinite(value)) return min; return Math.min(max, Math.max(min, Math.floor(value))); } /** * Projection returned by every read path. * * Two things are deliberately dropped: `embedding` (2KB+ of float noise no * caller can use, and the single biggest contributor to response size) and * the raw `_id`, replaced by a plain `id` string so the REST/MCP layers above * expose a stable opaque handle rather than a Convex document reference. * * Absent optionals are simply absent — Convex strips `undefined` object * fields during serialization rather than emitting nulls. */ function projectRecord(doc: Doc<"intelHistory">) { return { id: doc._id as string, domain: doc.domain, resource: doc.resource, country: doc.country, category: doc.category, title: doc.title, summary: doc.summary, sourceUrl: doc.sourceUrl, occurredAt: doc.occurredAt, ingestedAt: doc.ingestedAt, runId: doc.runId, dedupeKey: doc.dedupeKey, }; } export type IntelHistoryRecord = ReturnType; /** A `search` hit: the projection plus its cosine similarity to the query. */ export type IntelHistorySearchRecord = IntelHistoryRecord & { _score: number }; /** * Reject a vector that is the wrong length or carries a non-finite component. * The length check mirrors what Convex enforces on insert; the finiteness * check does not exist server-side and matters more: a single NaN silently * poisons the cosine similarity of every future search that ranks against * that row, and there is no error to trace it back from. */ function assertEmbedding(embedding: number[], context: string): void { if (embedding.length !== INTEL_HISTORY_EMBED_DIMS) { throw new Error( `${context}: embedding must have ${INTEL_HISTORY_EMBED_DIMS} dimensions, got ${embedding.length}`, ); } for (const component of embedding) { if (!Number.isFinite(component)) { throw new Error(`${context}: embedding components must be finite numbers`); } } } const appendRecordValidator = v.object({ dedupeKey: v.string(), country: v.optional(v.string()), category: v.optional(v.string()), title: v.string(), summary: v.optional(v.string()), sourceUrl: v.optional(v.string()), occurredAt: v.number(), embedding: v.array(v.float64()), }); /** * Append a run's events. Idempotent on `dedupeKey`: a seeder that republishes * the same event (every run does — the live snapshot is a rolling window) * adds nothing and reports it as skipped. * * The existence check and the insert share this mutation's transaction. Doing * the check via `ctx.runMutation` would put it in a SEPARATE transaction and * reopen exactly the race the dedupe is there to close. */ export const append = internalMutation({ args: { domain: v.string(), resource: v.string(), runId: v.string(), records: v.array(appendRecordValidator), }, handler: async (ctx, args) => { if (args.records.length > INTEL_HISTORY_MAX_APPEND_RECORDS) { throw new Error( `intelHistory.append: at most ${INTEL_HISTORY_MAX_APPEND_RECORDS} records per call, got ${args.records.length}`, ); } // Validate the whole batch up front. A throw rolls the transaction back // either way, but failing before the first write keeps the error about // the caller's payload rather than about a partially applied run. for (const rec of args.records) { assertEmbedding(rec.embedding, "intelHistory.append"); } // Must precede all dedupe reads. Patching an already-existing document // makes concurrent first-seen appends conflict and retry against the // winning transaction's inserted rows. const appendLock = await readAppendLockOrThrow(ctx); await ctx.db.patch(appendLock._id, { lastTouchedAt: Date.now() }); const ingestedAt = Date.now(); let inserted = 0; let skipped = 0; let retracted = 0; // Within-batch dedupe: two records sharing a key in one payload would // both miss the index lookup (neither is committed yet at read time). const seenInBatch = new Set(); const candidates = []; for (const rec of args.records) { if (seenInBatch.has(rec.dedupeKey)) { skipped += 1; continue; } seenInBatch.add(rec.dedupeKey); candidates.push(rec); } // The existence checks are independent of each other, so issue them // together: this is up to 100 indexed reads, and running them serially // holds the mutation's read set open far longer than needed while three // seeders write to this table on overlapping schedules. // // The retraction lookup rides in the same batch. It has to happen HERE, // inside the append transaction: an operator retracts a row because the // upstream feed served something poisoned or wrong, and that feed keeps // serving it — the seeders republish a rolling window, so without this // check the very next tick finds no row for the dedupeKey and re-inserts // it. See `retract` below. const [existing, tombstones] = await Promise.all([ Promise.all( candidates.map((rec) => ctx.db .query("intelHistory") .withIndex("by_dedupeKey", (q) => q.eq("dedupeKey", rec.dedupeKey)) .first(), ), ), Promise.all( candidates.map((rec) => ctx.db .query("intelHistoryRetractions") .withIndex("by_dedupeKey", (q) => q.eq("dedupeKey", rec.dedupeKey)) .first(), ), ), ]); for (const [index, rec] of candidates.entries()) { // Checked before the existence test so a retracted key is reported as // retracted rather than folded into the ordinary dedupe count — the two // are indistinguishable in the row state (both end with nothing written) // and only this counter tells an operator the tombstone is still doing // work. const tombstone = tombstones[index]; if (tombstone) { // Restart the tombstone's clock. Expiry is what eventually lets a // retracted identity back in, and the only question that should decide // it is "has the producer stopped offering this item?" — not "how long // ago did the operator act?". This record being in THIS run's payload // is direct evidence the feed is still serving it, so keying expiry on // the last suppressed attempt makes the tombstone outlive the item by // construction instead of by a bet on the retention constant. Only // fires for records an operator actually retracted, so the write cost // is proportional to tombstones, not to ingest volume. await ctx.db.patch(tombstone._id, { retractedAt: ingestedAt }); retracted += 1; continue; } if (existing[index] !== null) { skipped += 1; continue; } await ctx.db.insert("intelHistory", { domain: args.domain, resource: args.resource, runId: args.runId, country: rec.country, category: rec.category, title: rec.title, summary: rec.summary, sourceUrl: rec.sourceUrl, occurredAt: rec.occurredAt, ingestedAt, dedupeKey: rec.dedupeKey, embedding: rec.embedding, }); inserted += 1; } return { inserted, skipped, retracted }; }, }); /** * Idempotently initialize the document-backed append lock after deployment. * A missing lock is an operator/deploy failure, never a fallback to unsafe * index-range-only idempotency; `append` throws APPEND_LOCK_NOT_SEEDED until * this succeeds. */ export const _seedAppendLock = internalMutation({ args: {}, handler: async (ctx): Promise<{ seeded: number }> => { const existing = await ctx.db .query("intelHistoryAppendLocks") .withIndex("by_lockKey", (q) => q.eq("lockKey", APPEND_LOCK_KEY)) .first(); if (existing) return { seeded: 0 }; await ctx.db.insert("intelHistoryAppendLocks", { lockKey: APPEND_LOCK_KEY, lastTouchedAt: Date.now(), }); return { seeded: 1 }; }, }); /** * Best-effort append-lock touch for the retraction mutations. * * `retract` and `restore` write the same identity space `append` reads, so * they patch the same singleton to inherit its OCC serialization: a concurrent * append either commits before the retraction (and its row is then deleted) or * retries after it (and is then skipped by the tombstone). Neither can * interleave into a re-inserted row that no tombstone covers. * * Unlike `append`, a missing lock is NOT fatal here. `append` throws without * it, so an unseeded deploy has no concurrent writer to serialize against — * and refusing to retract in that state would block an incident cleanup for a * race that cannot happen. Log it and proceed. */ async function touchAppendLock(ctx: MutationCtx): Promise { const lock = await ctx.db .query("intelHistoryAppendLocks") .withIndex("by_lockKey", (q) => q.eq("lockKey", APPEND_LOCK_KEY)) .first(); if (!lock) { console.warn( JSON.stringify({ breadcrumb: "intel_history_retract_without_append_lock", lockKey: APPEND_LOCK_KEY, }), ); return; } await ctx.db.patch(lock._id, { lastTouchedAt: Date.now() }); } /** * Resolve the caller's identifiers to the set of `dedupeKey`s to act on. * * Document ids are what a reader has in hand — every retrieval path projects * `id` — but `dedupeKey` is what `append` matches on, so an id is only ever a * way to LOOK UP the key that actually has to be tombstoned. An id that no * longer resolves is reported back rather than ignored: it usually means the * row was already pruned or retracted, and an operator acting on a poisoned * record needs to know which of those it was. */ async function resolveRetractionKeys( ctx: MutationCtx, args: { ids?: string[]; dedupeKeys?: string[] }, ): Promise<{ keys: string[]; unresolvedIds: string[] }> { const keys = new Set(args.dedupeKeys ?? []); const unresolvedIds: string[] = []; const ids = args.ids ?? []; // normalizeId rather than a v.id() validator: a malformed id typed by an // operator should come back as "this id did not resolve", not as an opaque // argument-validation throw from the mutation boundary. const resolved = await Promise.all( ids.map(async (raw) => { const id = ctx.db.normalizeId("intelHistory", raw); return { raw, doc: id === null ? null : await ctx.db.get(id) }; }), ); for (const { raw, doc } of resolved) { if (doc === null) unresolvedIds.push(raw); else keys.add(doc.dedupeKey); } return { keys: [...keys], unresolvedIds }; } function assertRetractionIdentifierBudget( args: { ids?: string[]; dedupeKeys?: string[] }, // Named so the message matches the caller's actual argument surface — // `restore` accepts no ids at all, and telling its caller to "supply ids" // sends them after an argument the mutation would reject. accepts = "ids or dedupeKeys", ): void { const total = (args.ids?.length ?? 0) + (args.dedupeKeys?.length ?? 0); if (total === 0) { throw new Error(`intelHistory: at least one of ${accepts} is required`); } if (total > INTEL_HISTORY_MAX_RETRACT_IDENTIFIERS) { throw new Error( `intelHistory: at most ${INTEL_HISTORY_MAX_RETRACT_IDENTIFIERS} identifiers per call, got ${total}`, ); } } /** * Remove specific events from the history and keep them out (#5743). * * The store is agent-facing and durable for 180 days, so a single poisoned or * factually wrong feed item is retrievable for far longer than the live * snapshot that produced it. This is the supported way to take one back * without a hand-run Convex console operation. * * Deletion alone would not hold. `append` treats "no row with this dedupeKey" * as "never seen", and the producing feed keeps serving the item, so a bare * delete is reversed by the next seed tick. Each retraction therefore writes a * tombstone on the `dedupeKey` as well, which `append` consults. * * Scoped deliberately narrowly — explicit ids and/or dedupe keys, nothing * pattern-shaped. A retraction erases evidence from an intelligence archive; * "delete everything matching this substring" is the wrong amount of power to * hand a shared relay secret, and the identifiers are cheap to enumerate from * a search result. */ export const retract = internalMutation({ args: { ids: v.optional(v.array(v.string())), dedupeKeys: v.optional(v.array(v.string())), reason: v.string(), }, handler: async (ctx, args) => { assertRetractionIdentifierBudget(args); const reason = args.reason.trim(); if (!reason) { throw new Error("intelHistory.retract: reason is required"); } await touchAppendLock(ctx); const { keys, unresolvedIds } = await resolveRetractionKeys(ctx, args); // Fail loudly rather than report a successful no-op. A document id is the // handle an operator copies out of a search result, and a row that was // already pruned — or already retracted — no longer resolves to one. With // no key to tombstone the loop below does nothing, and a 200 saying // `deleted: 0, tombstoned: 0` reads as "already clean" when the truth is // "nothing was suppressed and the next seed tick will re-add it". The // recovery is in the message because the operator needs it right there. if (keys.length === 0) { throw new Error( `intelHistory.retract: no identifier resolved to a dedupeKey — nothing was tombstoned. ` + `Unresolved ids: ${unresolvedIds.join(", ")}. A deleted or pruned row cannot be ` + `retracted by id; re-run with --dedupe-key to suppress the identity itself.`, ); } const retractedAt = Date.now(); let deleted = 0; let tombstoned = 0; let refreshed = 0; for (const dedupeKey of keys) { const [rows, existingTombstone] = await Promise.all([ // `collect` over the dedupe index rather than `first`: the index is // meant to hold one row per key, and if an old bug ever put two there // a retraction that removed only one would leave the poisoned text // live while reporting success. ctx.db .query("intelHistory") .withIndex("by_dedupeKey", (q) => q.eq("dedupeKey", dedupeKey)) .collect(), ctx.db .query("intelHistoryRetractions") .withIndex("by_dedupeKey", (q) => q.eq("dedupeKey", dedupeKey)) .first(), ]); for (const row of rows) { await ctx.db.delete(row._id); deleted += 1; } if (existingTombstone) { // Re-retracting restarts the tombstone's own 180-day clock. An // operator repeating the call is telling us the item is still live // upstream, which is exactly when expiry would be premature. await ctx.db.patch(existingTombstone._id, { retractedAt, reason }); refreshed += 1; } else { await ctx.db.insert("intelHistoryRetractions", { dedupeKey, retractedAt, reason, }); tombstoned += 1; } } // Deliberately loud, and deliberately naming the identities rather than // counting them. Retraction removes intelligence from an archive on the // authority of a shared secret, and the tombstone row — the only other // record of what was taken out — is deletable by that same credential via // `restore`. A count cannot answer "which records were removed, and why?" // six weeks later; the arrays can, and both are already bounded (100 keys // x 256 chars) by the caller's identifier budget. console.info( JSON.stringify({ breadcrumb: "intel_history_retracted", keys, deleted, tombstoned, refreshed, unresolvedIds, reason, }), ); return { deleted, tombstoned, refreshed, keys, unresolvedIds }; }, }); /** * Lift a retraction (#5743). Removes the tombstones so the producing seeder * may store the event again. * * It does NOT resurrect the deleted rows — their embeddings are gone and this * mutation has no way to recompute one. If the event is still inside the * seeder's live window it reappears on the next tick; if it is not, the * retraction is effectively permanent for that row, which is the honest * outcome and the reason `retract` demands a stated reason up front. */ export const restore = internalMutation({ args: { dedupeKeys: v.array(v.string()), }, handler: async (ctx, args) => { assertRetractionIdentifierBudget({ dedupeKeys: args.dedupeKeys }, "dedupeKeys"); await touchAppendLock(ctx); // De-duplicated like `retract` does. Without it a repeated key deletes its // tombstone on the first pass and then — read-your-writes inside the // mutation — finds nothing on the second, so the same call reports the key // as both removed AND never-retracted. const dedupeKeys = [...new Set(args.dedupeKeys)]; let removed = 0; const notRetracted: string[] = []; for (const dedupeKey of dedupeKeys) { // The schema intentionally does not claim uniqueness for this index. // Remove every matching tombstone so legacy or manually introduced // duplicates cannot leave the key partially retracted. const tombstones = await ctx.db .query("intelHistoryRetractions") .withIndex("by_dedupeKey", (q) => q.eq("dedupeKey", dedupeKey)) .collect(); if (tombstones.length === 0) { notRetracted.push(dedupeKey); continue; } for (const tombstone of tombstones) { await ctx.db.delete(tombstone._id); removed += 1; } } // Names the keys for the same reason `retract` does: this is the operation // that destroys the other record of what was retracted. console.info( JSON.stringify({ breadcrumb: "intel_history_restored", dedupeKeys, removed, notRetracted, }), ); return { removed, notRetracted }; }, }); /** * Review read over the tombstones, newest retraction first. * * Retractions are invisible by construction — the whole point is that the row * is gone and stays gone — so without this the only record of what was taken * out is a log line that ages out long before the tombstone does. */ export const listRetractions = internalQuery({ args: { limit: v.optional(v.number()) }, handler: async (ctx, args) => { const limit = clamp( args.limit ?? INTEL_HISTORY_MAX_RETRACTION_PAGE, 1, INTEL_HISTORY_MAX_RETRACTION_PAGE, ); // Over-fetch by one to distinguish "that is all of them" from "the page // filled". Without the flag an operator checking whether a key is still // suppressed reads a truncated page as an exhaustive one and concludes the // retraction was lifted — the same `partial` contract `timeline` and // `search` already return, for the same reason. const rows = await ctx.db .query("intelHistoryRetractions") .withIndex("by_retractedAt") .order("desc") .take(limit + 1); return { retractions: rows.slice(0, limit).map((row) => ({ dedupeKey: row.dedupeKey, retractedAt: row.retractedAt, reason: row.reason, })), partial: rows.length > limit, }; }, }); /** * Chronological read. At least one of `domain` / `country` is required: the * only other option would be a table-wide scan ordered by `occurredAt`, which * has no index and no bounded cost. The REST layer rejects the unscoped case * too, so this throw is the backstop, not the user-facing error. * * `domain` wins the index when both are given, matching `search` below so the * two read paths agree; `country` is then post-filtered from a window * over-fetched by POST_FILTER_OVERFETCH. When that candidate window is full, * `partial` tells callers it may contain more country matches beyond the * bounded scan; they must not treat a short or empty page as exhaustive. */ export const timeline = internalQuery({ args: { domain: v.optional(v.string()), country: v.optional(v.string()), from: v.optional(v.number()), to: v.optional(v.number()), limit: v.optional(v.number()), }, handler: async (ctx, args) => { if (args.domain === undefined && args.country === undefined) { throw new Error( "intelHistory.timeline: at least one of domain or country is required", ); } const limit = clamp(args.limit ?? TIMELINE_DEFAULT_LIMIT, 1, TIMELINE_MAX_LIMIT); const { from, to } = args; // Over-fetch only when something is left to filter after the index range. const needsPostFilter = args.domain !== undefined && args.country !== undefined; const scanLimit = needsPostFilter ? Math.min(limit * POST_FILTER_OVERFETCH, TIMELINE_MAX_SCAN) : limit; // Both indexes are (scopeField, occurredAt), so the occurredAt bounds are // part of the index range rather than a post-read filter. The duplicated // branches are the price of that: `withIndex` types the range builder // against the specific index, so one shared helper cannot serve both. const docs = args.domain !== undefined ? await ctx.db .query("intelHistory") .withIndex("by_domain_occurredAt", (q) => { const scoped = q.eq("domain", args.domain as string); if (from !== undefined && to !== undefined) { return scoped.gte("occurredAt", from).lte("occurredAt", to); } if (from !== undefined) return scoped.gte("occurredAt", from); if (to !== undefined) return scoped.lte("occurredAt", to); return scoped; }) .order("desc") .take(scanLimit) : await ctx.db .query("intelHistory") .withIndex("by_country_occurredAt", (q) => { const scoped = q.eq("country", args.country as string); if (from !== undefined && to !== undefined) { return scoped.gte("occurredAt", from).lte("occurredAt", to); } if (from !== undefined) return scoped.gte("occurredAt", from); if (to !== undefined) return scoped.lte("occurredAt", to); return scoped; }) .order("desc") .take(scanLimit); const matched = needsPostFilter ? docs.filter((doc) => doc.country === args.country) : docs; return { records: matched.slice(0, limit).map(projectRecord), partial: needsPostFilter && docs.length === scanLimit, }; }, }); /** * Hydrate vector-search hits. `ctx.vectorSearch` returns ids and scores only, * so the action reads the documents back through this query. * * Input order is preserved (the caller's order IS the relevance order) and * ids that no longer resolve are dropped rather than throwing — a prune tick * between the search and the hydration is normal, not an error. */ export const getByIds = internalQuery({ args: { ids: v.array(v.id("intelHistory")) }, handler: async (ctx, args) => { if (args.ids.length > VECTOR_SEARCH_MAX_LIMIT) { throw new Error( `intelHistory.getByIds: at most ${VECTOR_SEARCH_MAX_LIMIT} ids per call, got ${args.ids.length}`, ); } // Independent reads on a user-facing search path — issue them together // rather than serially. Promise.all resolves in input order, which IS the // relevance order the caller depends on. const docs = await Promise.all(args.ids.map((id) => ctx.db.get(id))); return docs.filter((doc) => doc !== null).map((doc) => projectRecord(doc)); }, }); /** * Semantic read. Runs as an action because `ctx.vectorSearch` is only * available there. * * FILTER PUSHDOWN — Convex's vector filter builder exposes `eq` and `or` and * nothing else, so `domain AND country` cannot be expressed. `domain` is * pushed into the index (matching `timeline`'s index preference so the two * read paths agree) and `country` is post-filtered, as is the `occurredAt` * range, which the vector index cannot express at all. Both post-filters run * after the top-k cut, so the query over-fetches by POST_FILTER_OVERFETCH * whenever one applies. If that bounded candidate window fills, `partial` * marks the response as potentially incomplete rather than presenting a * false-empty or short result as exhaustive. See the final note in * convex/schema.ts. */ export const search = internalAction({ args: { embedding: v.array(v.float64()), domain: v.optional(v.string()), country: v.optional(v.string()), from: v.optional(v.number()), to: v.optional(v.number()), limit: v.optional(v.number()), minScore: v.optional(v.number()), }, // The return type is annotated, not inferred. `getByIds` lives in THIS // module, so `internal.intelHistory.getByIds` resolves through the module's // own type — inferring the handler's return would be circular (TS7022/7023), // and TypeScript resolves such cycles by degrading the whole `internal` // surface to `any`, which silently un-types every OTHER module's // self-referential `ctx.runQuery` too. handler: async ( ctx, args, ): Promise<{ records: IntelHistorySearchRecord[]; partial: boolean }> => { assertEmbedding(args.embedding, "intelHistory.search"); const limit = clamp(args.limit ?? SEARCH_DEFAULT_LIMIT, 1, SEARCH_MAX_LIMIT); const pushedDown = args.domain !== undefined ? "domain" : args.country !== undefined ? "country" : null; const needsPostFilter = (pushedDown === "domain" && args.country !== undefined) || args.from !== undefined || args.to !== undefined; const vectorLimit = Math.min( needsPostFilter ? limit * POST_FILTER_OVERFETCH : limit, VECTOR_SEARCH_MAX_LIMIT, ); const hits = await ctx.vectorSearch("intelHistory", "by_embedding", { vector: args.embedding, limit: vectorLimit, ...(pushedDown === "domain" ? { filter: (q) => q.eq("domain", args.domain as string) } : pushedDown === "country" ? { filter: (q) => q.eq("country", args.country as string) } : {}), }); if (hits.length === 0) return { records: [], partial: false }; const scoreById = new Map(hits.map((hit) => [hit._id as string, hit._score])); const hydrated: IntelHistoryRecord[] = await ctx.runQuery( internal.intelHistory.getByIds, { ids: hits.map((hit) => hit._id) }, ); const records = hydrated .filter((rec) => { if (pushedDown === "domain" && args.country !== undefined) { if (rec.country !== args.country) return false; } if (args.from !== undefined && rec.occurredAt < args.from) return false; if (args.to !== undefined && rec.occurredAt > args.to) return false; return true; }) .filter((rec) => { const score = scoreById.get(rec.id) ?? 0; return args.minScore === undefined || score >= args.minScore; }) .slice(0, limit) .map((rec) => ({ ...rec, _score: scoreById.get(rec.id) ?? 0 })); return { records, partial: needsPostFilter && hits.length === vectorLimit }; }, }); /** * Retention prune. Ages rows out by `ingestedAt`, not `occurredAt`, so a * backfill of genuinely old events survives its first night. * * Modelled on `pruneApiPlanLimitData` (convex/apiPlanLimitNotices.ts): the * per-run batch keeps one invocation inside Convex's write limit, and a full * batch means more aged rows remain, so the mutation reschedules itself * immediately — carrying the SAME resolved `now` so the cutoff stays fixed * across the chain. Without the self-drain, one daily run would delete at * most `batch` rows and the first prune 180 days after launch would take * weeks to clear. */ export const prune = internalMutation({ args: { now: v.optional(v.number()), retentionMs: v.optional(v.number()), // Per-run delete cap. Optional so tests can drive the drain without // seeding PRUNE_BATCH rows. limit: v.optional(v.number()), }, handler: async (ctx, args) => { const now = args.now ?? Date.now(); const cutoff = now - (args.retentionMs ?? RETENTION_MS); // Floor of 1, and non-finite falls back to the default rather than // propagating: a batch of 0 makes `take(0)` return [] and the // `stale.length >= batch` check below read 0 >= 0 as "a full batch", so the // mutation would reschedule itself forever, deleting nothing. The cron // passes {}, but this is operator-callable and `limit: 0` is a natural // thing to type when probing the drain. const requestedBatch = args.limit; const batch = Number.isFinite(requestedBatch) ? Math.max(1, Math.floor(requestedBatch as number)) : PRUNE_BATCH; const stale = await ctx.db .query("intelHistory") .withIndex("by_ingestedAt", (q) => q.lt("ingestedAt", cutoff)) .take(batch); for (const doc of stale) { await ctx.db.delete(doc._id); } // Tombstones age out on the same clock (#5743), measured from when the // operator retracted rather than when the event happened: a retraction // only has to outlive the upstream item's presence in the seeder's live // window, and one retention window past the operator's action is a // generous bound on that. They drain in the same pass rather than under // their own cron — there are a handful of them, hand-created, and a second // scheduled function for that volume is cost with no signal. The row is // tiny next to a 512-float history row, so the batch budget is unaffected. // Tombstones use the FULL retention window no matter what the caller // passed. `retentionMs` is an operator-callable override inherited from // the apiPlanLimit prune this was modelled on, where shortening it merely // deletes old rows early — recoverable, and nothing depends on them. Here // it is a security control: `prune({ retentionMs: 0 })` would set the // cutoff to `now`, drain EVERY tombstone, and hand every retracted // identity back to the producing feed on its next tick, silently undoing // every retraction anyone had ever made. Clamping to the floor lets the // override only ever LENGTHEN suppression, never shorten it — the same // shape as the `limit` floor above, for the same reason. const tombstoneCutoff = now - Math.max(args.retentionMs ?? RETENTION_MS, RETENTION_MS); const staleRetractions = await ctx.db .query("intelHistoryRetractions") .withIndex("by_retractedAt", (q) => q.lt("retractedAt", tombstoneCutoff)) .take(batch); for (const doc of staleRetractions) { await ctx.db.delete(doc._id); } // Either table filling its batch means more aged rows remain. Each pass // deletes at least one row from whichever table is still full, so the // self-drain terminates. const rescheduled = stale.length >= batch || staleRetractions.length >= batch; if (rescheduled) { await ctx.scheduler.runAfter(0, internal.intelHistory.prune, { now, retentionMs: args.retentionMs, limit: batch, }); } return { deleted: stale.length, deletedRetractions: staleRetractions.length, rescheduled, }; }, });