feat: P4 qdrant collection lifecycle (self-heal + guarded rebuild)

- shared TS collection manager: self-heal creates missing collection (1024/cosine)
  and missing keyword payload indexes; never mutates incompatible contracts
  (semantic_index_incompatible); async index visibility polled with bounded deadline
- session admission (qdrantEnsure) uses the manager in self-heal mode; operator path
  keeps require_existing semantics
- runtime lease exposes semanticQdrantUrl to the operator
- operator commands vector-inspect/vector-rebuild with exact confirmation guards
- thothctl workspace vector inspect|rebuild (Go) with --collection/--confirm/--destroy
- p4 acceptance runner: real Qdrant (v1.18.2) lifecycle checks, 11/11 PASS
- docs: CLI contract, manual walkthrough P4 (PENDING), PROJECT_STATE
This commit is contained in:
2026-08-12 20:00:14 +02:00
parent 230a876314
commit e056c19e62
20 changed files with 1041 additions and 32 deletions
+1 -1
View File
@@ -46,7 +46,7 @@ export class ReadinessManager {
const pending = (async (): Promise<ReadinessResult> => {
try {
if (descriptor) {
const qdrant = await runner.qdrantEnsure(descriptor, this.timeoutSec);
const qdrant = await runner.qdrantEnsure(descriptor, this.timeoutSec, "self_heal");
if (!qdrant.ok) return qdrant;
}
const ollama = await runner.ollamaEnsure(workspace, this.timeoutSec);
+16 -27
View File
@@ -10,6 +10,7 @@ import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext
import { secretValue, type SecretBundleConfig } from "../config/secret-bundle.js";
import { renderWorkspaceRuntimeFromSnapshotPath } from "../workspaces/runtime-config-lease.js";
import {
DEFAULT_SEMANTIC_RUNTIME,
type RuntimeInstallationOverlay,
type RuntimePaths,
type SemanticRuntimeConfig,
@@ -19,6 +20,7 @@ import {
validateOperationalWorkspace,
type WorkspaceDescriptor,
} from "../workspaces/schema.js";
import { reconcileCollection } from "../workspaces/qdrant-collection.js";
export interface ThtConfig extends SecretBundleConfig {
thtBin: string;
@@ -29,6 +31,8 @@ export interface ThtConfig extends SecretBundleConfig {
secretRoots?: readonly string[];
semanticRuntime: SemanticRuntimeConfig;
qdrantRequest?: typeof fetch;
/** "self_heal" for session admission (create missing collections/indexes), default "require_existing". */
qdrantCollectionMode?: "self_heal" | "require_existing";
}
export interface RuntimeConfigLease {
@@ -216,7 +220,7 @@ export class ThtRunner {
throw new Error("registry workspace runtime requires an absolute data root");
})(),
secretRoots: this.cfg.secretRoots ?? [],
semanticRuntime: this.cfg.semanticRuntime,
semanticRuntime: this.cfg.semanticRuntime ?? DEFAULT_SEMANTIC_RUNTIME,
});
const path = this.createRuntimeSnapshot(rendered.renderedConfig);
let released = false;
@@ -561,6 +565,7 @@ export class ThtRunner {
async qdrantEnsure(
workspace: WorkspaceDescriptor,
timeoutSec: number,
mode: "self_heal" | "require_existing" = "require_existing",
): Promise<QdrantEnsureResult> {
let descriptor;
try {
@@ -572,32 +577,16 @@ export class ThtRunner {
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), Math.max(1, timeoutSec) * 1000);
try {
const url = new URL(
`/collections/${encodeURIComponent(collection.collection)}`,
this.cfg.semanticRuntime.internalQdrantUrl,
);
const request = this.cfg.qdrantRequest ?? fetch;
const response = await request(url.toString(), { method: "GET", signal: controller.signal });
if (response.status === 404) {
return { ok: false, code: "semantic_index_incompatible" };
}
if (!response.ok) return { ok: false, code: "workspace_not_activatable" };
const body = await response.json() as any;
const result = body?.result;
const vectors = result?.config?.params?.vectors;
const payloadSchema = result?.payload_schema;
const configurationMatches = vectors
&& vectors.size === collection.dimensions
&& typeof vectors.distance === "string"
&& vectors.distance.toLowerCase() === collection.distance;
const indexesMatch = payloadSchema
&& typeof payloadSchema === "object"
&& REQUIRED_QDRANT_PAYLOAD_INDEXES.every(
(field) => payloadSchema[field]?.data_type === "keyword",
);
return configurationMatches && indexesMatch
? { ok: true }
: { ok: false, code: "semantic_index_incompatible" };
const checked = await reconcileCollection({
baseUrl: this.cfg.semanticRuntime.internalQdrantUrl,
collection: collection.collection,
dimensions: collection.dimensions,
distance: collection.distance,
mode,
request: this.cfg.qdrantRequest ?? fetch,
signal: controller.signal,
});
return checked;
} catch {
return { ok: false, code: "workspace_not_activatable" };
} finally {
+12 -1
View File
@@ -18,7 +18,7 @@ export interface WorkspaceMaintenanceIo {
writeStderr(value: string): void;
}
type Command = "inspect" | "preprocess-dwh" | "schema-suggest-fks" | "schema-check" | "index-schema" | "preprocess-evidence" | "preprocess-run";
type Command = "inspect" | "preprocess-dwh" | "schema-suggest-fks" | "schema-check" | "index-schema" | "preprocess-evidence" | "preprocess-run" | "vector-inspect" | "vector-rebuild";
function failureResult(
operation: string,
@@ -69,6 +69,8 @@ function parseRequest(command: string, stdin: string): Record<string, unknown> {
"index-schema": ["schemaVersion", "workspaceId", "resumeRunId"],
"preprocess-evidence": ["schemaVersion", "workspaceId", "dryRun", "resumeRunId"],
"preprocess-run": ["schemaVersion", "workspaceId", "resumeRunId"],
"vector-inspect": ["schemaVersion", "workspaceId"],
"vector-rebuild": ["schemaVersion", "workspaceId", "collection", "confirm", "destroy"],
};
const allowed = allowedByCommand[command];
if (!allowed) throw new Error("unknown command");
@@ -121,6 +123,15 @@ async function dispatch(command: Command, service: WorkspacePreprocessingService
workspaceId: request.workspaceId as string,
resumeRunId: request.resumeRunId as string | undefined,
});
case "vector-inspect":
return await service.vectorInspect({ workspaceId: request.workspaceId as string });
case "vector-rebuild":
return await service.vectorRebuild({
workspaceId: request.workspaceId as string,
collection: request.collection as string | undefined,
confirm: request.confirm as string | undefined,
destroy: request.destroy === true,
});
}
}
@@ -125,6 +125,38 @@ function noEvidenceWarning(workspace: WorkspaceDescriptor): string[] {
export class WorkspacePreprocessingService {
constructor(private readonly deps: WorkspacePreprocessingServiceDeps) {}
async vectorInspect(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
const collection = runtime.workspace.semantic_index.vector_store.collection;
const res = await fetch(`${runtime.configLease.semanticQdrantUrl}/collections/${encodeURIComponent(collection)}`, { method: "GET" });
if (!res.ok) return baseResult(runtime, "vector inspect", "failed", "semantic_index_incompatible", { warnings: ["collection unavailable"] });
const body = await res.json() as any;
const info = body?.result;
const vectors = info?.config?.params?.vectors;
return baseResult(runtime, "vector inspect", "succeeded", "ok", {
counts: { dimensions: vectors?.size ?? 0 },
warnings: [`collection=${collection} distance=${vectors?.distance ?? "unknown"}`],
});
}
async vectorRebuild(options: { workspaceId: string; collection?: string; confirm?: string; destroy?: boolean }): Promise<WorkspaceOperationResult> {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
const collection = runtime.workspace.semantic_index.vector_store.collection;
if (options.collection !== collection || options.confirm !== collection || options.destroy !== true) {
return baseResult(runtime, "vector rebuild", "failed", "semantic_index_incompatible", { warnings: ["rebuild requires exact confirmation and --destroy"] });
}
const q = `${runtime.configLease.semanticQdrantUrl}/collections/${encodeURIComponent(collection)}`;
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"] });
const put = await fetch(q, {
method: "PUT",
headers: { "content-type": "application/json" },
body: JSON.stringify({ vectors: { size: runtime.workspace.semantic_index.vector_store.dimensions, distance: runtime.workspace.semantic_index.vector_store.distance } }),
});
if (!put.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}`] });
}
async inspect(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
try {
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
+105
View File
@@ -0,0 +1,105 @@
export const QDRANT_REQUIRED_INDEXES = Object.freeze([
"content_hash", "document_id", "kind", "record_key",
"record_kind", "vector_generation", "workspace_id", "workspace_revision",
]);
export type CollectionMode = "self_heal" | "require_existing";
export interface CollectionCheck {
ok: boolean;
code?: "semantic_index_incompatible" | "workspace_not_activatable";
state?: "ready" | "created" | "repaired";
}
export interface ReconcileCollectionOptions {
baseUrl: string;
collection: string;
dimensions: number;
distance: string;
mode: CollectionMode;
request?: typeof fetch;
signal?: AbortSignal;
}
function qdrantDistance(distance: string): string {
return distance.length === 0 ? distance : distance.charAt(0).toUpperCase() + distance.slice(1);
}
function qdrantUrl(baseUrl: string, path: string): string {
return new URL(path, baseUrl).toString();
}
async function collectionInfo(opts: ReconcileCollectionOptions, request: typeof fetch): Promise<unknown | undefined> {
const res = await request(qdrantUrl(opts.baseUrl, `/collections/${encodeURIComponent(opts.collection)}`), { method: "GET", signal: opts.signal });
if (res.status === 404) return undefined;
if (!res.ok) throw new Error("qdrant collection check failed");
return (await res.json() as any)?.result;
}
function vectorCompatibility(info: any, opts: ReconcileCollectionOptions): boolean {
const vectors = info?.config?.params?.vectors;
return Boolean(vectors && vectors.size === opts.dimensions && typeof vectors.distance === "string"
&& vectors.distance.toLowerCase() === opts.distance);
}
async function missingIndexes(opts: ReconcileCollectionOptions, info: any): Promise<string[]> {
const payloadSchema = info?.payload_schema;
if (!payloadSchema || typeof payloadSchema !== "object") return [...QDRANT_REQUIRED_INDEXES];
return QDRANT_REQUIRED_INDEXES.filter((field) => payloadSchema[field]?.data_type !== "keyword");
}
async function createCollection(opts: ReconcileCollectionOptions, request: typeof fetch): Promise<void> {
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) } }),
signal: opts.signal,
});
if (!res.ok && res.status !== 409) throw new Error("qdrant collection 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",
headers: { "content-type": "application/json" },
body: JSON.stringify({ field_name: field, field_schema: "keyword" }),
signal: opts.signal,
});
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). */
export async function reconcileCollection(opts: ReconcileCollectionOptions): Promise<CollectionCheck> {
const request = opts.request ?? fetch;
let info = await collectionInfo(opts, request);
if (info === undefined) {
if (opts.mode !== "self_heal") 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);
if (info === undefined) return { ok: false, code: "workspace_not_activatable" };
}
if (!vectorCompatibility(info, opts)) {
return { ok: false, code: "semantic_index_incompatible" };
}
const missing = await missingIndexes(opts, info);
if (missing.length > 0) {
if (opts.mode !== "self_heal") 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).
const deadline = Date.now() + 15000;
let current: any = info;
while (Date.now() < deadline) {
current = await collectionInfo(opts, request);
if (!vectorCompatibility(current, opts)) break;
if ((await missingIndexes(opts, current)).length === 0) {
return { ok: true, state: "repaired" };
}
await new Promise((resolve) => setTimeout(resolve, 500));
}
return { ok: false, code: "semantic_index_incompatible" };
}
return { ok: true, state: "ready" };
}
@@ -57,6 +57,7 @@ export interface RenderedWorkspaceRuntime {
bindings: RuntimeBindings;
bindingDigest: string;
renderedConfig: string;
semanticQdrantUrl: string;
}
export interface ActiveRenderedWorkspaceRuntime extends RenderedWorkspaceRuntime {
@@ -71,6 +72,7 @@ export interface DeterministicRuntimeConfigLease extends RuntimeConfigLease {
catalogBlob: string;
configDigest: string;
bindingDigest: string;
semanticQdrantUrl: string;
effectiveConfig: CanonicalEffectiveConfig;
effectiveConfigIdentity: string;
configFingerprint: string;
@@ -310,6 +312,7 @@ function renderWorkspaceRuntimeFromWorkspace(options: {
installationOverlay: overlay,
bindings,
bindingDigest: stableBindingDigest(bindings),
semanticQdrantUrl: options.semanticRuntime.internalQdrantUrl,
renderedConfig: renderRuntimeConfig(
options.workspace,
bindings,
@@ -508,6 +511,7 @@ export async function publishDeterministicRuntimeConfigLease(options: {
catalogBlob: rendered.catalogBlob,
configDigest,
bindingDigest: rendered.bindingDigest,
semanticQdrantUrl: rendered.semanticQdrantUrl,
effectiveConfig,
effectiveConfigIdentity: effectiveConfigIdentityValue,
configFingerprint: configFingerprintValue,
@@ -570,6 +574,7 @@ export async function publishDeterministicRuntimeConfigLease(options: {
catalogBlob: rendered.catalogBlob,
configDigest,
bindingDigest: rendered.bindingDigest,
semanticQdrantUrl: rendered.semanticQdrantUrl,
effectiveConfig,
effectiveConfigIdentity: effectiveConfigIdentityValue,
configFingerprint: configFingerprintValue,
+1 -1
View File
@@ -34,7 +34,7 @@ export interface SemanticRuntimeConfig {
internalEmbeddingDimensions: number;
}
const DEFAULT_SEMANTIC_RUNTIME: SemanticRuntimeConfig = {
export const DEFAULT_SEMANTIC_RUNTIME: SemanticRuntimeConfig = {
internalQdrantUrl: "http://qdrant:6333",
internalEmbeddingUrl: "http://embedding:11434",
internalEmbeddingModel: "qwen3-embedding:0.6b",