feat(evidence): add Qdrant BM25 vector in place

This commit is contained in:
2026-08-24 18:01:46 +02:00
parent ae0976a4aa
commit 0e9add09a9
11 changed files with 320 additions and 14 deletions
+2 -2
View File
@@ -20,7 +20,7 @@ import {
validateOperationalWorkspace,
type WorkspaceDescriptor,
} from "../workspaces/schema.js";
import { reconcileCollection } from "../workspaces/qdrant-collection.js";
import { reconcileCollection, type CollectionMode } from "../workspaces/qdrant-collection.js";
import type { WorkspaceSecretStore } from "../workspaces/secret-store.js";
export interface ThtConfig extends SecretBundleConfig {
@@ -575,7 +575,7 @@ export class ThtRunner {
async qdrantEnsure(
workspace: WorkspaceDescriptor,
timeoutSec: number,
mode: "self_heal" | "require_existing" = "require_existing",
mode: CollectionMode = "require_existing",
): Promise<QdrantEnsureResult> {
let descriptor;
try {
+4
View File
@@ -307,6 +307,10 @@ function createProductionService(): WorkspacePreprocessingService {
const result = await runner.qdrantEnsure(workspace, 30);
return result.ok ? { ok: true as const } : { ok: false as const, code: result.code ?? "workspace_not_activatable" };
},
evidencePreflight: async (workspace) => {
const result = await runner.qdrantEnsure(workspace, 30, "evidence_maintenance");
return result.ok ? { ok: true as const } : { ok: false as const, code: result.code ?? "workspace_not_activatable" };
},
});
}
@@ -14,6 +14,7 @@ export interface EvidencePreprocessingDependencies {
runStage(argv: string[]): Promise<Record<string, unknown>>;
persistJob(): void;
semanticPreflight(): Promise<{ ok: true } | { ok: false; code: SemanticFailureCode }>;
evidencePreflight(): Promise<{ ok: true } | { ok: false; code: SemanticFailureCode }>;
requireRunId(value: unknown): string;
numberRecord(value: unknown): Record<string, number> | undefined;
}
@@ -138,7 +139,7 @@ export async function preprocessEvidence(
}
const policy = evidencePolicy(request.evidence, request.httpPrivateHostAllowlist);
if (policy) return policy;
const semantic = await deps.semanticPreflight();
const semantic = await deps.evidencePreflight();
if (!semantic.ok) {
return { status: "failed", code: semantic.code, runId: request.job.runId };
}
@@ -73,6 +73,9 @@ export interface WorkspacePreprocessingServiceDeps {
semanticPreflight(workspace: WorkspaceDescriptor): Promise<
{ ok: true } | { ok: false; code: "workspace_not_activatable" | "semantic_index_incompatible" }
>;
evidencePreflight(workspace: WorkspaceDescriptor): Promise<
{ ok: true } | { ok: false; code: "workspace_not_activatable" | "semantic_index_incompatible" }
>;
httpPrivateHostAllowlist?: readonly string[];
}
@@ -410,6 +413,7 @@ export class WorkspacePreprocessingService {
runStage: async (argv) => await this.runJsonStage(scope.runtime, argv),
persistJob: () => this.state(scope.runtime.workspaceId).writeJob(scope.job),
semanticPreflight: async () => await this.deps.semanticPreflight(scope.runtime.workspace),
evidencePreflight: async () => await this.deps.evidencePreflight(scope.runtime.workspace),
requireRunId: (value) => this.requireRunId(value),
numberRecord: (value) => this.numberRecord(value),
};
+53 -8
View File
@@ -3,12 +3,12 @@ export const QDRANT_REQUIRED_INDEXES = Object.freeze([
"record_kind", "vector_generation", "workspace_id", "workspace_revision",
]);
export type CollectionMode = "self_heal" | "require_existing";
export type CollectionMode = "self_heal" | "require_existing" | "evidence_maintenance";
export interface CollectionCheck {
ok: boolean;
code?: "semantic_index_incompatible" | "workspace_not_activatable";
state?: "ready" | "created" | "repaired";
state?: "ready" | "created" | "repaired" | "upgraded";
}
export interface ReconcileCollectionOptions {
@@ -52,12 +52,43 @@ async function createCollection(opts: ReconcileCollectionOptions, request: typeo
const res = await request(qdrantUrl(opts.baseUrl, `/collections/${encodeURIComponent(opts.collection)}`), {
method: "PUT",
headers: { "content-type": "application/json" },
body: JSON.stringify({ vectors: { size: opts.dimensions, distance: qdrantDistance(opts.distance) } }),
body: JSON.stringify({
vectors: { size: opts.dimensions, distance: qdrantDistance(opts.distance) },
...(opts.mode === "evidence_maintenance" ? { sparse_vectors: { bm25: { modifier: "idf" } } } : {}),
}),
signal: opts.signal,
});
if (!res.ok && res.status !== 409) throw new Error("qdrant collection creation failed");
}
type EvidenceSparseCompatibility = "compatible" | "upgradeable" | "incompatible";
function evidenceSparseCompatibility(info: any): EvidenceSparseCompatibility {
const sparseVectors = info?.config?.params?.sparse_vectors;
if (sparseVectors === undefined) return "upgradeable";
if (!sparseVectors || typeof sparseVectors !== "object" || Array.isArray(sparseVectors)) {
return "incompatible";
}
const bm25 = sparseVectors.bm25;
if (bm25 === undefined) return "upgradeable";
return typeof bm25 === "object" && bm25 !== null && !Array.isArray(bm25)
&& typeof bm25.modifier === "string" && bm25.modifier.toLowerCase() === "idf"
? "compatible" : "incompatible";
}
async function createBm25Vector(opts: ReconcileCollectionOptions, request: typeof fetch): Promise<void> {
const res = await request(
qdrantUrl(opts.baseUrl, `/collections/${encodeURIComponent(opts.collection)}/vectors/bm25`),
{
method: "PUT",
headers: { "content-type": "application/json" },
body: JSON.stringify({ sparse: { modifier: "idf" } }),
signal: opts.signal,
},
);
if (!res.ok && res.status !== 409) throw new Error("qdrant BM25 vector creation failed");
}
async function createIndex(opts: ReconcileCollectionOptions, field: string, request: typeof fetch): Promise<void> {
const res = await request(qdrantUrl(opts.baseUrl, `/collections/${encodeURIComponent(opts.collection)}/index`), {
method: "PUT",
@@ -68,13 +99,14 @@ async function createIndex(opts: ReconcileCollectionOptions, field: string, requ
if (!res.ok && res.status !== 409) throw new Error("qdrant index creation failed");
}
/** Reconcile a Qdrant collection: self-heal creates missing collections/indexes; require_existing
* only validates and refuses incompatible contracts (never mutates). */
/** Reconcile a Qdrant collection. Only Evidence maintenance may add the BM25 sparse vector;
* session admission remains limited to the existing dense/index self-heal behavior. */
export async function reconcileCollection(opts: ReconcileCollectionOptions): Promise<CollectionCheck> {
const request = opts.request ?? fetch;
const allowsMutation = opts.mode === "self_heal" || opts.mode === "evidence_maintenance";
let info = await collectionInfo(opts, request);
if (info === undefined) {
if (opts.mode !== "self_heal") return { ok: false, code: "semantic_index_incompatible" };
if (!allowsMutation) return { ok: false, code: "semantic_index_incompatible" };
await createCollection(opts, request);
// Tolerate an already-compatible concurrent creator: re-read the final state.
info = await collectionInfo(opts, request);
@@ -83,9 +115,22 @@ export async function reconcileCollection(opts: ReconcileCollectionOptions): Pro
if (!vectorCompatibility(info, opts)) {
return { ok: false, code: "semantic_index_incompatible" };
}
let bm25Added = false;
if (opts.mode === "evidence_maintenance") {
const sparse = evidenceSparseCompatibility(info);
if (sparse === "incompatible") return { ok: false, code: "semantic_index_incompatible" };
if (sparse === "upgradeable") {
await createBm25Vector(opts, request);
info = await collectionInfo(opts, request);
if (!vectorCompatibility(info, opts) || evidenceSparseCompatibility(info) !== "compatible") {
return { ok: false, code: "semantic_index_incompatible" };
}
bm25Added = true;
}
}
const missing = await missingIndexes(opts, info);
if (missing.length > 0) {
if (opts.mode !== "self_heal") return { ok: false, code: "semantic_index_incompatible" };
if (!allowsMutation) return { ok: false, code: "semantic_index_incompatible" };
for (const field of missing) await createIndex(opts, field, request);
// Qdrant payload indexes become visible asynchronously: poll until the
// contract is complete or a bounded deadline passes (fail closed).
@@ -101,5 +146,5 @@ export async function reconcileCollection(opts: ReconcileCollectionOptions): Pro
}
return { ok: false, code: "semantic_index_incompatible" };
}
return { ok: true, state: "ready" };
return { ok: true, state: bm25Added ? "upgraded" : "ready" };
}
+84
View File
@@ -28,6 +28,90 @@ const compatible = (size = 1024, distance = "Cosine", schema = payloadSchema) =>
payload_schema: schema,
});
test("Evidence maintenance adds an absent BM25 vector without changing the dense contract", async () => {
const info = compatible();
const requests: Array<{ url: string; init?: any }> = [];
const request = async (url: string, init?: any) => {
requests.push({ url, init });
if (init?.method === "PUT" && /\/vectors\/bm25$/.test(url)) {
expect(JSON.parse(String(init.body))).toEqual({ sparse: { modifier: "idf" } });
(info.config.params as any).sparse_vectors = { bm25: { modifier: "idf" } };
return { status: 200, ok: true, json: async () => ({}) } as any;
}
return { status: 200, ok: true, json: async () => ({ result: info }) } as any;
};
const result = await reconcileCollection({
baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine",
mode: "evidence_maintenance", request,
});
expect(result).toEqual({ ok: true, state: "upgraded" });
expect(info.config.params.vectors).toEqual({ size: 1024, distance: "Cosine" });
expect(requests.filter(({ init }) => init?.method === "PUT")).toHaveLength(1);
expect(requests[1]?.url).toBe("http://qdrant:6333/collections/c/vectors/bm25");
});
test("Evidence maintenance creates a missing collection with both required vector contracts", async () => {
let info: any;
let createdBody: any;
const request = async (url: string, init?: any) => {
if (init?.method === "PUT") {
createdBody = JSON.parse(String(init.body));
info = { config: { params: createdBody }, payload_schema: payloadSchema };
return { status: 200, ok: true, json: async () => ({}) } as any;
}
if (info === undefined) return { status: 404, ok: false, json: async () => ({}) } as any;
return { status: 200, ok: true, json: async () => ({ result: info }) } as any;
};
const result = await reconcileCollection({
baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine",
mode: "evidence_maintenance", request,
});
expect(result).toEqual({ ok: true, state: "ready" });
expect(createdBody).toEqual({
vectors: { size: 1024, distance: "Cosine" },
sparse_vectors: { bm25: { modifier: "idf" } },
});
});
test("Evidence maintenance refuses an incompatible BM25 definition without mutating", async () => {
const info = compatible();
(info.config.params as any).sparse_vectors = { bm25: { modifier: "none" } };
const requests: Array<{ url: string; init?: any }> = [];
const request = async (url: string, init?: any) => {
requests.push({ url, init });
return { status: 200, ok: true, json: async () => ({ result: info }) } as any;
};
const result = await reconcileCollection({
baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine",
mode: "evidence_maintenance", request,
});
expect(result).toEqual({ ok: false, code: "semantic_index_incompatible" });
expect(requests.filter(({ init }) => init?.method === "PUT")).toEqual([]);
});
test("ordinary session reconciliation does not add BM25", async () => {
const info = compatible();
const requests: Array<{ url: string; init?: any }> = [];
const request = async (url: string, init?: any) => {
requests.push({ url, init });
return { status: 200, ok: true, json: async () => ({ result: info }) } as any;
};
const result = await reconcileCollection({
baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine",
mode: "self_heal", request,
});
expect(result).toEqual({ ok: true, state: "ready" });
expect(requests.filter(({ url }) => /\/vectors\/bm25$/.test(url))).toEqual([]);
});
test("self-heal creates a missing compatible collection", async () => {
const r = await reconcileCollection({
baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine",
@@ -161,6 +161,7 @@ function fixture(workspace = baseWorkspace) {
runChild,
listSessions: async () => [],
semanticPreflight: async () => ({ ok: true }),
evidencePreflight: async () => ({ ok: true }),
});
return { dataRoot, runChild, requests, service };
}
@@ -283,6 +284,7 @@ test("index schema fails closed when semantic preflight refuses the collection",
runChild,
listSessions: async () => [],
semanticPreflight: async () => ({ ok: false, code: "semantic_index_incompatible" }),
evidencePreflight: async () => ({ ok: true }),
});
const result = await service.indexSchema({ workspaceId: "psd-clinical" });
@@ -309,6 +311,7 @@ test("filesystem Evidence proceeds after materialization and private HTTP hosts
runChild: vi.fn(),
listSessions: async () => [],
semanticPreflight: async () => ({ ok: true }),
evidencePreflight: async () => ({ ok: true }),
httpPrivateHostAllowlist: ["metadata.internal"],
});
@@ -513,6 +516,7 @@ test("vector rebuild recreates the full collection contract including keyword in
runChild: vi.fn(),
listSessions: async () => [],
semanticPreflight: async () => ({ ok: true }),
evidencePreflight: async () => ({ ok: true }),
});
// replace global fetch used by vectorRebuild/reconcileCollection
const original = globalThis.fetch;
@@ -40,11 +40,13 @@ function dependencies(payload: Record<string, unknown> = {}): EvidencePreprocess
runStage: ReturnType<typeof vi.fn>;
persistJob: ReturnType<typeof vi.fn>;
semanticPreflight: ReturnType<typeof vi.fn>;
evidencePreflight: ReturnType<typeof vi.fn>;
} {
return {
runStage: vi.fn(async () => payload),
persistJob: vi.fn(),
semanticPreflight: vi.fn(async () => ({ ok: true as const })),
evidencePreflight: vi.fn(async () => ({ ok: true as const })),
requireRunId(value) {
if (typeof value !== "string" || !/^[0-9a-f]{32}$/.test(value)) {
throw new Error("child run id is invalid");
@@ -60,7 +62,7 @@ function dependencies(payload: Record<string, unknown> = {}): EvidencePreprocess
};
}
test("owns the standalone Evidence stage argv and mutation order", async () => {
test("Evidence maintenance preflights the additive BM25 contract before starting its stage", async () => {
const state = job({ childRuns: { evidence: "b".repeat(32) } });
const deps = dependencies({ run_id: "c".repeat(32), counts: { added: 2 } });
@@ -69,7 +71,8 @@ test("owns the standalone Evidence stage argv and mutation order", async () => {
deps,
);
expect(deps.semanticPreflight).toHaveBeenCalledOnce();
expect(deps.evidencePreflight).toHaveBeenCalledOnce();
expect(deps.semanticPreflight).not.toHaveBeenCalled();
expect(deps.runStage).toHaveBeenCalledWith([
"preprocess", "evidence", "--resume", "b".repeat(32), "--json", "-c", "/dev/fd/3",
]);