Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
46 changes: 46 additions & 0 deletions apps/app/src/lib/embedding/embedding.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -6,6 +6,7 @@ const queryMock = vi.fn();
const infoMock = vi.fn();
const deleteMock = vi.fn();
const rangeMock = vi.fn();
const fetchMock = vi.fn();

vi.mock('@upstash/vector', () => ({
Index: vi.fn().mockImplementation(() => ({
Expand All @@ -14,6 +15,7 @@ vi.mock('@upstash/vector', () => ({
info: infoMock,
delete: deleteMock,
range: rangeMock,
fetch: fetchMock,
})),
}));

Expand All @@ -34,6 +36,7 @@ import {
upsertEntityEmbeddings,
findSimilarTasks,
cosineToUnitScore,
computeEntityContentHash,
waitForIndexed,
pruneOrphanTaskVectors,
} from './index';
Expand All @@ -44,6 +47,11 @@ beforeEach(() => {
infoMock.mockReset();
deleteMock.mockReset();
rangeMock.mockReset();
fetchMock.mockReset();
// Default: every requested vector is present in Upstash, so the dedup guard's
// existence check confirms hash-matched entities can be safely skipped. Tests
// exercising the missing-vector desync override this to return null per id.
fetchMock.mockImplementation(async (ids: string[]) => ids.map((id) => ({ id })));
process.env.UPSTASH_VECTOR_REST_URL = 'https://test.upstash.io';
process.env.UPSTASH_VECTOR_REST_TOKEN = 'test-token';
});
Expand Down Expand Up @@ -168,6 +176,44 @@ describe('upsertEntityEmbeddings', () => {
expect(second.appliedHashes.map((h) => h.id).sort()).toEqual(['tsk_a', 'tsk_c']);
expect(second.skippedCount).toBe(1);
});

it('re-embeds a hash-matched entity whose vector is missing from Upstash (CS-681)', async () => {
// Desync: the stored hash matches the current content (dedup guard would
// normally skip), but the vector is GONE from Upstash (index re-provision /
// eviction / partial delete). Skipping on the hash alone never rewrites the
// vector, so findSimilarTasks enumerates zero candidates and "Draft plan &
// suggest links" returns 0/0 forever — the hash never changes so it recurs
// on every run. The guard must detect the absent vector and re-embed it.
const text = 'unchanged task content';
const hash = computeEntityContentHash({ text });

// tsk_present still has its vector; tsk_missing's is gone.
fetchMock.mockImplementation(async (ids: string[]) =>
ids.map((id) => (id === 'task_org_1_tsk_missing' ? null : { id })),
);

const result = await upsertEntityEmbeddings({
organizationId: 'org_1',
kind: 'task',
entities: [
{ id: 'tsk_present', text },
{ id: 'tsk_missing', text },
],
existingHashes: new Map([
['tsk_present', hash],
['tsk_missing', hash],
]),
});

// Only the vectorless task is re-embedded + re-upserted; the present one
// stays skipped (the hash optimization still holds when the vector exists).
expect(upsertMock).toHaveBeenCalledTimes(1);
expect(upsertMock).toHaveBeenCalledWith(
expect.objectContaining({ id: 'task_org_1_tsk_missing' }),
);
expect(result.appliedHashes).toEqual([{ id: 'tsk_missing', hash }]);
expect(result.skippedCount).toBe(1);
});
});

describe('cosineToUnitScore', () => {
Expand Down
72 changes: 67 additions & 5 deletions apps/app/src/lib/embedding/index.ts
Original file line number Diff line number Diff line change
Expand Up @@ -65,6 +65,9 @@ const TASK_VECTOR_PAGE_SIZE = 1000;
// Backstop so a misbehaving cursor can't loop forever. 500 pages × 1000 = 500k
// task vectors, far beyond any real org — hitting it signals a bug, not scale.
const TASK_VECTOR_MAX_PAGES = 500;
// Upstash caps `fetch` at 1000 ids per request — batch the existence check for
// hash-matched entities (see `findEntitiesWithoutVectors`) accordingly.
const VECTOR_FETCH_BATCH_SIZE = 1000;

let cachedIndex: Index | null = null;

Expand Down Expand Up @@ -127,6 +130,49 @@ export function computeEntityContentHash({
.digest('hex');
}

/**
* Of a set of entities whose stored hash matches the current content (so the
* dedup guard in `upsertEntityEmbeddings` would otherwise skip them), return
* those whose vector is actually ABSENT from Upstash and therefore must be
* re-embedded.
*
* A matching `embeddingHash` proves the content is unchanged — NOT that the
* vector still exists. Vectors can vanish independently of the DB hash (index
* re-provision, cache eviction, a partial earlier delete), leaving a row
* hash-cached with no vector. Skipping those on the hash alone never rewrites
* the missing vector, so `findSimilarTasks` enumerates zero candidates and
* "Draft plan & suggest links" returns 0/0 on every run — recurring forever
* because the hash never changes (CS-681). Earlier fixes pruned ORPHAN vectors;
* this heals the inverse desync (a live row with a MISSING vector).
*
* `fetch` returns one entry per id aligned to the input, or `null` when the id
* is absent — existence is all we need, so `includeVectors` is omitted to keep
* the payload tiny.
*/
async function findEntitiesWithoutVectors({
kind,
organizationId,
candidates,
}: {
kind: EntityKind;
organizationId: string;
candidates: Array<{ entity: EntityInput; hash: string }>;
}): Promise<Array<{ entity: EntityInput; hash: string }>> {
if (candidates.length === 0) return [];
const index = getIndex();
const missing: Array<{ entity: EntityInput; hash: string }> = [];
for (let i = 0; i < candidates.length; i += VECTOR_FETCH_BATCH_SIZE) {
const batch = candidates.slice(i, i + VECTOR_FETCH_BATCH_SIZE);
const fetched = await index.fetch(
batch.map(({ entity }) => embeddingId(kind, organizationId, entity.id)),
);
fetched.forEach((row, j) => {
if (!row) missing.push(batch[j]);
});
}
return missing;
}

/**
* Upsert per-org entity embeddings into Upstash Vector.
*
Expand All @@ -143,17 +189,33 @@ export async function upsertEntityEmbeddings({
const valid = entities.filter((e) => e.text.trim().length > 0);
if (valid.length === 0) return { appliedHashes: [], skippedCount: 0 };

// Skip entities whose stored hash matches the current content hash
// text + model + dims + department haven't changed, so the existing
// vector is still authoritative. Saves the OpenAI embed AND the Upstash
// upsert (the two non-trivial costs in this path).
// A stored hash matching the current content hash means text + model + dims
// + department are unchanged, so the existing vector is still authoritative —
// skipping saves the OpenAI embed AND the Upstash upsert (the two non-trivial
// costs in this path).
const withHashes = valid.map((entity) => ({
entity,
hash: computeEntityContentHash({ text: entity.text, department: entity.department }),
}));
const toEmbed = withHashes.filter(
const changed = withHashes.filter(
({ entity, hash }) => existingHashes?.get(entity.id) !== hash,
);
// The hash optimization is only sound while a matching hash implies the
// vector still exists. When a vector has vanished but its hash lingers
// (index re-provision, eviction, a partial earlier delete), skipping on the
// hash alone leaves that entity permanently vectorless — the desync that made
// treatment-plan suggestions return 0/0 on every run (CS-681). So verify the
// hash-matched entities' vectors are actually present, and re-embed any gone.
const unchanged = withHashes.filter(
({ entity, hash }) => existingHashes?.get(entity.id) === hash,
);
const missingVectors = await findEntitiesWithoutVectors({
kind,
organizationId,
candidates: unchanged,
});

const toEmbed = [...changed, ...missingVectors];
const skippedCount = withHashes.length - toEmbed.length;
if (toEmbed.length === 0) {
return { appliedHashes: [], skippedCount };
Expand Down
Loading