fix: upsert in bounded chunks, recreate Qdrant indexes on rebuild, larger maintenance tmpfs
This commit is contained in:
@@ -11,6 +11,7 @@ import {
|
|||||||
} from "./preprocessing-state.js";
|
} from "./preprocessing-state.js";
|
||||||
import type { DeterministicRuntimeConfigLease } from "./runtime-config-lease.js";
|
import type { DeterministicRuntimeConfigLease } from "./runtime-config-lease.js";
|
||||||
import { readAnnotationsSync } from "./annotations-sync.js";
|
import { readAnnotationsSync } from "./annotations-sync.js";
|
||||||
|
import { reconcileCollection } from "./qdrant-collection.js";
|
||||||
|
|
||||||
export interface WorkspaceOperationResult {
|
export interface WorkspaceOperationResult {
|
||||||
schemaVersion: 1;
|
schemaVersion: 1;
|
||||||
@@ -150,12 +151,15 @@ export class WorkspacePreprocessingService {
|
|||||||
const q = `${runtime.configLease.semanticQdrantUrl}/collections/${encodeURIComponent(collection)}`;
|
const q = `${runtime.configLease.semanticQdrantUrl}/collections/${encodeURIComponent(collection)}`;
|
||||||
const del = await fetch(q, { method: "DELETE" });
|
const del = await fetch(q, { method: "DELETE" });
|
||||||
if (!del.ok && del.status !== 404) return baseResult(runtime, "vector rebuild", "failed", "semantic_index_incompatible", { warnings: ["collection delete failed"] });
|
if (!del.ok && del.status !== 404) return baseResult(runtime, "vector rebuild", "failed", "semantic_index_incompatible", { warnings: ["collection delete failed"] });
|
||||||
const put = await fetch(q, {
|
// Recreate the complete contract (dimensions + distance + the 8 required keyword indexes).
|
||||||
method: "PUT",
|
const recreated = await reconcileCollection({
|
||||||
headers: { "content-type": "application/json" },
|
baseUrl: runtime.configLease.semanticQdrantUrl,
|
||||||
body: JSON.stringify({ vectors: { size: runtime.workspace.semantic_index.vector_store.dimensions, distance: runtime.workspace.semantic_index.vector_store.distance } }),
|
collection,
|
||||||
|
dimensions: runtime.workspace.semantic_index.vector_store.dimensions,
|
||||||
|
distance: runtime.workspace.semantic_index.vector_store.distance,
|
||||||
|
mode: "self_heal",
|
||||||
});
|
});
|
||||||
if (!put.ok) return baseResult(runtime, "vector rebuild", "failed", "semantic_index_incompatible", { warnings: ["collection recreate failed"] });
|
if (!recreated.ok) return baseResult(runtime, "vector rebuild", "failed", "semantic_index_incompatible", { warnings: ["collection recreate failed"] });
|
||||||
return baseResult(runtime, "vector rebuild", "succeeded", "ok", { warnings: [`recreated collection=${collection}`] });
|
return baseResult(runtime, "vector rebuild", "succeeded", "ok", { warnings: [`recreated collection=${collection}`] });
|
||||||
}
|
}
|
||||||
async inspect(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
|
async inspect(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
|
||||||
|
|||||||
@@ -123,6 +123,7 @@ function runtime(workspace = baseWorkspace, workspaceId = workspace.workspace.id
|
|||||||
catalogBlob: "c".repeat(40),
|
catalogBlob: "c".repeat(40),
|
||||||
configDigest: "sha256:config",
|
configDigest: "sha256:config",
|
||||||
bindingDigest: "sha256:bindings",
|
bindingDigest: "sha256:bindings",
|
||||||
|
semanticQdrantUrl: "http://qdrant:6333",
|
||||||
effectiveConfig: {
|
effectiveConfig: {
|
||||||
schemaVersion: 1,
|
schemaVersion: 1,
|
||||||
dwh: {
|
dwh: {
|
||||||
@@ -491,3 +492,36 @@ test("full runs continue after schema accept only when the accepted blob matches
|
|||||||
const stillBlocked = await second.service.run({ workspaceId: "psd-clinical", resumeRunId: secondRunId });
|
const stillBlocked = await second.service.run({ workspaceId: "psd-clinical", resumeRunId: secondRunId });
|
||||||
expect(stillBlocked).toMatchObject({ status: "blocked", code: "manual_review_required" });
|
expect(stillBlocked).toMatchObject({ status: "blocked", code: "manual_review_required" });
|
||||||
});
|
});
|
||||||
|
|
||||||
|
test("vector rebuild recreates the full collection contract including keyword indexes", async () => {
|
||||||
|
const f = fixture();
|
||||||
|
// mock fetch: DELETE ok, then reconcileCollection self-heals create + indexes (real fetch in deps)
|
||||||
|
const calls: string[] = [];
|
||||||
|
const fakeFetch = async (url: string, init?: any) => {
|
||||||
|
calls.push(`${init?.method ?? "GET"} ${url}`);
|
||||||
|
if ((init?.method ?? "GET") === "DELETE") return new Response("", { status: 200 });
|
||||||
|
if (url.endsWith("/collections/psd-clinical") && init?.method === "PUT") return new Response("", { status: 200 });
|
||||||
|
if (url.endsWith("/collections/psd-clinical") && init?.method === "GET") {
|
||||||
|
return new Response(JSON.stringify({ result: { config: { params: { vectors: { size: 1024, distance: "Cosine" } } }, payload_schema: { content_hash: { data_type: "keyword" }, document_id: { data_type: "keyword" }, kind: { data_type: "keyword" }, record_key: { data_type: "keyword" }, record_kind: { data_type: "keyword" }, vector_generation: { data_type: "keyword" }, workspace_id: { data_type: "keyword" }, workspace_revision: { data_type: "keyword" } } } }), { status: 200 });
|
||||||
|
}
|
||||||
|
if (url.endsWith("/collections/psd-clinical/index") && init?.method === "PUT") return new Response("", { status: 200 });
|
||||||
|
return new Response(JSON.stringify({ result: {} }), { status: 200 });
|
||||||
|
};
|
||||||
|
const service = new WorkspacePreprocessingService({
|
||||||
|
dataRoot: f.dataRoot,
|
||||||
|
acquireActiveRuntime: async () => runtime(baseWorkspace),
|
||||||
|
runChild: vi.fn(),
|
||||||
|
listSessions: async () => [],
|
||||||
|
semanticPreflight: async () => ({ ok: true }),
|
||||||
|
});
|
||||||
|
// replace global fetch used by vectorRebuild/reconcileCollection
|
||||||
|
const original = globalThis.fetch;
|
||||||
|
globalThis.fetch = fakeFetch as any;
|
||||||
|
try {
|
||||||
|
const result = await service.vectorRebuild({ workspaceId: "psd-clinical", collection: "psd-clinical", confirm: "psd-clinical", destroy: true });
|
||||||
|
expect(result).toMatchObject({ status: "succeeded", code: "ok" });
|
||||||
|
} finally {
|
||||||
|
globalThis.fetch = original;
|
||||||
|
}
|
||||||
|
expect(calls.some((c) => c.startsWith("DELETE "))).toBe(true);
|
||||||
|
});
|
||||||
|
|||||||
+2
-2
@@ -99,8 +99,8 @@ services:
|
|||||||
user: "10001:10001"
|
user: "10001:10001"
|
||||||
read_only: true
|
read_only: true
|
||||||
tmpfs:
|
tmpfs:
|
||||||
- /tmp:rw,noexec,nosuid,nodev,size=64m,mode=1777
|
- /tmp:rw,noexec,nosuid,nodev,size=1g,mode=1777
|
||||||
- /var/tmp:rw,noexec,nosuid,nodev,size=32m,mode=1777
|
- /var/tmp:rw,noexec,nosuid,nodev,size=128m,mode=1777
|
||||||
cap_drop:
|
cap_drop:
|
||||||
- ALL
|
- ALL
|
||||||
security_opt:
|
security_opt:
|
||||||
|
|||||||
@@ -25,6 +25,7 @@ from tht.vectorstore.store import VectorHit, hit_from_metadata
|
|||||||
_GENERATION = re.compile(r"gen:[0-9a-f]{32}")
|
_GENERATION = re.compile(r"gen:[0-9a-f]{32}")
|
||||||
_WORKSPACE = re.compile(r"[a-z][a-z0-9_-]{0,63}")
|
_WORKSPACE = re.compile(r"[a-z][a-z0-9_-]{0,63}")
|
||||||
_KEYWORD_INDEXES = (
|
_KEYWORD_INDEXES = (
|
||||||
|
|
||||||
"content_hash",
|
"content_hash",
|
||||||
"document_id",
|
"document_id",
|
||||||
"kind",
|
"kind",
|
||||||
@@ -35,6 +36,8 @@ _KEYWORD_INDEXES = (
|
|||||||
"workspace_revision",
|
"workspace_revision",
|
||||||
)
|
)
|
||||||
|
|
||||||
|
UPSERT_BATCH_SIZE = 256
|
||||||
|
|
||||||
|
|
||||||
def point_id(workspace_id: str, kind: str, record_key: str, workspace_revision: str | None = None) -> str:
|
def point_id(workspace_id: str, kind: str, record_key: str, workspace_revision: str | None = None) -> str:
|
||||||
# P3: schema/Evidence points are revision-scoped; memory/solved remain workspace-wide.
|
# P3: schema/Evidence points are revision-scoped; memory/solved remain workspace-wide.
|
||||||
@@ -215,11 +218,14 @@ class QdrantVectorStore:
|
|||||||
),
|
),
|
||||||
}
|
}
|
||||||
)
|
)
|
||||||
self._call(
|
# Qdrant rejects request bodies larger than its JSON limit (32 MiB by default).
|
||||||
"PUT",
|
# A large schema/Evidence corpus therefore must be upserted in bounded chunks.
|
||||||
f"/collections/{self._collection}/points?wait=true",
|
for start in range(0, len(points), UPSERT_BATCH_SIZE):
|
||||||
{"points": points},
|
self._call(
|
||||||
)
|
"PUT",
|
||||||
|
f"/collections/{self._collection}/points?wait=true",
|
||||||
|
{"points": points[start:start + UPSERT_BATCH_SIZE]},
|
||||||
|
)
|
||||||
return len(records)
|
return len(records)
|
||||||
|
|
||||||
def delete_kinds(self, collection: str, kinds: list[str]) -> int:
|
def delete_kinds(self, collection: str, kinds: list[str]) -> int:
|
||||||
|
|||||||
Reference in New Issue
Block a user