From e056c19e6214254a9e3b2390e24c389920b84e95 Mon Sep 17 00:00:00 2001 From: mptyl Date: Wed, 12 Aug 2026 20:00:14 +0200 Subject: [PATCH] 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 --- PROJECT_STATE.md | 25 ++ backend/scripts/p4-acceptance.mjs | 403 ++++++++++++++++++ backend/scripts/p4-acceptance.test.mjs | 53 +++ backend/src/runtime/readiness-manager.ts | 2 +- backend/src/tht/tht-runner.ts | 43 +- backend/src/workspace-maintenance.ts | 13 +- .../src/workspaces/preprocessing-service.ts | 32 ++ backend/src/workspaces/qdrant-collection.ts | 105 +++++ .../src/workspaces/runtime-config-lease.ts | 5 + backend/src/workspaces/runtime-renderer.ts | 2 +- backend/test/qdrant-collection.test.ts | 88 ++++ backend/test/routes-sessions.test.ts | 2 +- backend/test/tht-qdrant-readiness.test.ts | 2 +- backend/test/workspace-maintenance.test.ts | 29 ++ docs/contracts/workspace-preprocessing-cli.md | 26 ++ docs/testing/p2-p6-manual-verification.md | 22 + scripts/p4-acceptance.sh | 59 +++ scripts/test-p4-acceptance.sh | 8 + .../internal/workspaceops/operations.go | 87 ++++ .../internal/workspaceops/operations_test.go | 67 +++ 20 files changed, 1041 insertions(+), 32 deletions(-) create mode 100644 backend/scripts/p4-acceptance.mjs create mode 100644 backend/scripts/p4-acceptance.test.mjs create mode 100644 backend/src/workspaces/qdrant-collection.ts create mode 100644 backend/test/qdrant-collection.test.ts create mode 100755 scripts/p4-acceptance.sh create mode 100755 scripts/test-p4-acceptance.sh diff --git a/PROJECT_STATE.md b/PROJECT_STATE.md index 31f7239c..9f693f5b 100644 --- a/PROJECT_STATE.md +++ b/PROJECT_STATE.md @@ -31,6 +31,31 @@ `3b0726472e15c157…`. - **Manual gate:** P3 walkthrough in `docs/testing/p2-p6-manual-verification.md`; decision **PENDING**. P4 starts only after an explicit new authorization. +### P4 Qdrant collection lifecycle — implementation complete, automated PASS, manual PENDING (2026-08-12) + +- **Scope:** P4 (PRD D4): one shared TypeScript collection manager owns the Qdrant collection + and payload-index contract; session admission self-heals a missing collection (1024/cosine + + the 8 required keyword payload indexes) and adds missing indexes, but never mutates an + incompatible collection (`semantic_index_incompatible`); the operator path keeps + `require_existing` semantics. +- **Host CLI:** `thothctl workspace vector inspect` (read-only contract report) and + `thothctl workspace vector rebuild --workspace --collection --confirm --destroy` + (guarded delete/recreate of only the descriptor-owned collection, with durable state before + deletion and verification after recreation; mismatched confirmation or missing `--destroy` + → exit 2). +- **Key files:** `backend/src/workspaces/qdrant-collection.ts` (+test), `backend/src/tht/tht-runner.ts` + (`qdrantEnsure` self-heal for admission; default `require_existing` elsewhere), + `backend/src/workspaces/runtime-config-lease.ts` (lease exposes `semanticQdrantUrl`), + `backend/src/workspace-maintenance.ts` + `preprocessing-service.ts` (`vector-inspect`/`vector-rebuild` + operator commands), `tools/thothctl/internal/workspaceops/operations.go` (+tests). +- **Automated acceptance:** PASS 11/11 (run `p4-3a001f83fae22fe72056dc52e5ff63b5`, + report `.artifacts/p4-integration/...` retained via `--keep`): preflight, clean_state, ownership, + qdrant_up, self_heal_create_missing, self_heal_repairs_missing_index, incompatible_refused, + require_existing_refused, rebuild_recreates_contract, secret_scan, cleanup_confinement. +- **Gates:** backend 666/666 + tsc clean; Go build+test 9/9; p4 runner unit tests 3/3; harness + 841 passed (only the two pre-existing debt failures unchanged). +- **Manual acceptance:** PENDING — walkthrough section P4 in `docs/testing/p2-p6-manual-verification.md`. + ### P2 host preprocessing CLI — implementation complete, automated PASS, manual PENDING (2026-08-11) diff --git a/backend/scripts/p4-acceptance.mjs b/backend/scripts/p4-acceptance.mjs new file mode 100644 index 00000000..e2a328e2 --- /dev/null +++ b/backend/scripts/p4-acceptance.mjs @@ -0,0 +1,403 @@ +#!/usr/bin/env node +// P4 automated integration acceptance: Qdrant collection lifecycle (self-heal + guarded rebuild). +import { createHash, randomBytes } from "node:crypto"; +import { execFile, execFileSync } from "node:child_process"; +import { promisify } from "node:util"; +import { fileURLToPath } from "node:url"; +import { existsSync, lstatSync, mkdirSync, readFileSync, readdirSync, realpathSync, rmSync, statSync, writeFileSync } from "node:fs"; +import { mkdir, readFile, rm, writeFile } from "node:fs/promises"; +import { basename, dirname, isAbsolute, join, relative, resolve, sep } from "node:path"; +import { createServer as createNetServer } from "node:net"; +import process from "node:process"; + +import { stringify as yamlStringify } from "yaml"; + +import { buildSafeEnvironment, deriveOverall, scanSecrets } from "./p1-acceptance.mjs"; + +const execFileAsync = promisify(execFile); +const modulePath = fileURLToPath(import.meta.url); +const defaultRepositoryRoot = realpathSync(resolve(dirname(modulePath), "../..")); +const RUN_ID = /^p4-[0-9a-f]{32}$/; +const HEX64 = /^[0-9a-f]{64}$/; +const QDRANT_IMAGE = "qdrant/qdrant:v1.18.2"; +export const CHECK_IDS = Object.freeze([ + "preflight", + "clean_state", + "ownership", + "qdrant_up", + "self_heal_create_missing", + "self_heal_repairs_missing_index", + "incompatible_refused", + "require_existing_refused", + "rebuild_recreates_contract", + "secret_scan", + "cleanup_confinement", +]); +const TOPOLOGY = ["installation", "fixtures", "logs", "qdrant-volumes"]; +const MAX_REPORT_JSON_BYTES = 64 * 1024; +const MAX_REPORT_MD_BYTES = 32 * 1024; + + +function resolveSystemExecutable(name) { + for (const candidate of [`/usr/bin/${name}`, `/bin/${name}`, `/opt/homebrew/bin/${name}`, `/usr/local/bin/${name}`, `/usr/local/sbin/${name}`]) { + try { + const resolved = realpathSync(candidate); + if (statSync(resolved).isFile()) return resolved; + } catch { /* continue */ } + } + throw new Error(`required executable ${name} is unavailable`); +} +const DOCKER_BIN = (() => { try { return resolveSystemExecutable("docker"); } catch { return "docker"; } })(); + +function nowIso() { return new Date().toISOString(); } +function sha256(value) { return createHash("sha256").update(value).digest("hex"); } +function assert(condition, message) { if (!condition) throw new Error(message); } +function sleep(ms) { return new Promise((resolve) => setTimeout(resolve, ms)); } + +function canonicalRoot(repositoryRoot = defaultRepositoryRoot) { + return realpathSync(repositoryRoot); +} +export function canonicalIntegrationBase(repositoryRoot = defaultRepositoryRoot) { + return join(canonicalRoot(repositoryRoot), ".artifacts", "p4-integration"); +} +export function validateRunRoot(repositoryRoot, runRoot, runId) { + if (!RUN_ID.test(runId)) throw new Error("invalid owned run id"); + const base = canonicalIntegrationBase(repositoryRoot); + const lexical = resolve(runRoot); + if (dirname(lexical) !== base || basename(lexical) !== runId) throw new Error("run root is not a direct integration child"); + return lexical; +} +function validateNoSymlinkAncestors(repositoryRoot, target) { + const repo = canonicalRoot(repositoryRoot); + const rel = relative(repo, target); + if (rel.startsWith("..") || isAbsolute(rel)) throw new Error("target escapes the repository"); + let cursor = repo; + for (const part of rel.split(sep)) { + cursor = join(cursor, part); + if (existsSync(cursor) && lstatSyncIsSymlink(cursor)) throw new Error(`symlink ancestor: ${cursor}`); + } +} +function lstatSyncIsSymlink(path) { return lstatSync(path).isSymbolicLink(); } + +export function createOwnedRun(repositoryRoot, nonce = randomBytes(16).toString("hex")) { + const runId = `p4-${nonce}`; + if (!RUN_ID.test(runId)) throw new Error("invalid run id"); + const base = canonicalIntegrationBase(repositoryRoot); + mkdirSync(base, { recursive: true }); + const runRoot = join(base, runId); + validateNoSymlinkAncestors(repositoryRoot, runRoot); + mkdirSync(join(runRoot, "installation"), { recursive: true }); + mkdirSync(join(runRoot, "fixtures"), { recursive: true }); + mkdirSync(join(runRoot, "logs"), { recursive: true }); + mkdirSync(join(runRoot, "qdrant-volumes"), { recursive: true }); + const marker = { runId, createdAt: nowIso(), repositoryRoot: canonicalRoot(repositoryRoot), sha256: "" }; + marker.sha256 = sha256(JSON.stringify(marker) + "\n"); + writeFileSync(join(runRoot, "run.json"), JSON.stringify(marker, null, 2) + "\n", { mode: 0o600 }); + return { runId, runRoot }; +} + +export function cleanupOwnedRun(repositoryRoot, runRoot, runId) { + const validated = validateRunRoot(repositoryRoot, runRoot, runId); + const base = canonicalIntegrationBase(repositoryRoot); + for (const sibling of readdirSync(base)) { + if (sibling.startsWith("p4-") && sibling !== runId) throw new Error("refusing cleanup with sibling p4 runs present"); + } + rmSync(validated, { recursive: true, force: true }); +} + + +function result(checkId, ok, detail, cause) { + const message = cause ? `${String(detail)} :: ${String(cause)}` : String(detail); + return { checkId, status: ok ? "PASS" : "FAIL", ok: !!ok, detail: ok ? "PASS" : message.slice(0, 500) }; +} + +function execCapture(command, args, options = {}) { + const spawned = execFileSync(command, args, { encoding: "utf8", maxBuffer: 64 * 1024 * 1024, ...options }); + return String(spawned ?? ""); +} + +async function waitForQdrant(baseUrl, timeoutMs = 120000) { + const deadline = Date.now() + timeoutMs; + while (Date.now() < deadline) { + try { + const res = await fetch(`${baseUrl}/readyz`, { signal: AbortSignal.timeout(3000) }); + if (res.ok) return true; + } catch { /* retry */ } + await sleep(1500); + } + throw new Error("qdrant did not become ready"); +} + +async function qdrantGet(baseUrl, path) { + const res = await fetch(`${baseUrl}${path}`); + if (!res.ok) throw new Error(`qdrant GET ${path} -> ${res.status}`); + return (await res.json()).result; +} +async function qdrantPut(baseUrl, path, body) { + const payload = { ...body }; + if (payload.vectors && typeof payload.vectors.distance === "string" && payload.vectors.distance.length > 0) { + payload.vectors = { ...payload.vectors, distance: payload.vectors.distance.charAt(0).toUpperCase() + payload.vectors.distance.slice(1) }; + } + const res = await fetch(`${baseUrl}${path}`, { + method: "PUT", + headers: { "content-type": "application/json" }, + body: JSON.stringify(payload), + }); + if (!res.ok && res.status !== 409) throw new Error(`qdrant PUT ${path} -> ${res.status}`); + return res.ok || res.status === 409; +} +async function qdrantDelete(baseUrl, path) { + const res = await fetch(`${baseUrl}${path}`, { method: "DELETE" }); + if (!res.ok && res.status !== 404) throw new Error(`qdrant DELETE ${path} -> ${res.status}`); +} + +function contractOk(info, dimensions, distance) { + const vectors = info?.config?.params?.vectors; + const schema = info?.payload_schema; + const required = ["content_hash","document_id","kind","record_key","record_kind","vector_generation","workspace_id","workspace_revision"]; + if (!vectors || vectors.size !== dimensions || String(vectors.distance).toLowerCase() !== distance) return false; + if (!schema || typeof schema !== "object") return false; + return required.every((field) => schema[field]?.data_type === "keyword"); +} + +async function runIntegration(repositoryRoot, runRoot, runId, qdrantBaseUrl) { + const checks = []; + const record = (checkId, fn) => checks.push(async () => { + try { return result(checkId, await fn()); } + catch (error) { return result(checkId, false, error.message, error.cause?.message ?? error.code); } + }); + const ctx = { run: { root: runRoot, id: runId }, repo: repositoryRoot }; + + record("preflight", async () => { + execCapture(DOCKER_BIN, ["version", "--format", "{{.Server.Version}}"]); + execCapture("node", ["--version"]); + execCapture("npm", ["--version"]); + return true; + }); + + record("clean_state", async () => { + const base = canonicalIntegrationBase(repositoryRoot); + const leftovers = readdirSync(base).filter((entry) => entry.startsWith("p4-") && entry !== runId); + if (leftovers.length > 0) throw new Error(`leftover p4 runs: ${leftovers.join(", ")}`); + return true; + }); + + record("ownership", async () => { + const marker = JSON.parse(await readFile(join(runRoot, "run.json"), "utf8")); + if (marker.runId !== runId) throw new Error("run marker mismatch"); + return true; + }); + + const containerName = `p4acc-qdrant-${runId.slice(3, 11)}`; + let started = false; + const startQdrant = async () => { + await execFileAsync(DOCKER_BIN, ["rm", "-f", containerName], { stdio: "ignore" }).catch(() => {}); + const hostPort = await freePort(); + try { + await execFileAsync(DOCKER_BIN, ["run", "-d", "--name", containerName, + "-p", `127.0.0.1:${hostPort}:6333`, "-v", `${containerName}-vol:/qdrant/storage`, + "--restart", "no", QDRANT_IMAGE], { stdio: "ignore" }); + } catch (error) { + const detail = error.stderr ?? error.message; + throw new Error(`docker run qdrant failed: ${String(detail).slice(0, 300)}`); + } + started = true; + return `http://127.0.0.1:${hostPort}`; + }; + const stopQdrant = async () => { + if (!started) return; + try { + const logs = await execFileAsync(DOCKER_BIN, ["logs", containerName]); + const insp = await execFileAsync(DOCKER_BIN, ["inspect", "--format", "{{.State.Status}} exit={{.State.ExitCode}} oom={{.State.OOMKilled}}", containerName]).catch(() => ({ stdout: "inspect failed" })); + await writeFile(join(runRoot, "qdrant.log"), `INSPECT: ${String(insp.stdout).trim()}\n` + String(logs.stdout).slice(-3000) + "\n---STDERR---\n" + String(logs.stderr).slice(-3000)); + } catch { /* best effort */ } + await execFileAsync(DOCKER_BIN, ["rm", "-f", containerName], { stdio: "ignore" }).catch(() => {}); + await execFileAsync(DOCKER_BIN, ["volume", "rm", "-f", `${containerName}-vol`], { stdio: "ignore" }).catch(() => {}); + }; + + + +function freePort() { + return new Promise((resolve, reject) => { + const server = createNetServer(); + server.unref(); + server.on("error", reject); + server.listen(0, "127.0.0.1", () => { + const port = server.address().port; + server.close(() => resolve(port)); + }); + }); +} + +async function dockerPortRetry(containerName, attempts = 20) { + for (let attempt = 0; attempt < attempts; attempt += 1) { + try { + const inspect = await execFileAsync(DOCKER_BIN, ["port", containerName, "6333"]); + const line = String(inspect.stdout).trim(); + const hostPort = line.split("\n")[0].split(":")[1]; + if (hostPort) return `http://127.0.0.1:${hostPort}`; + } catch { /* transient */ } + await sleep(1000); + } + throw new Error(`docker port ${containerName} did not resolve`); +} + + let manager; + try { + const qdrantUrl = await startQdrant(); + await waitForQdrant(qdrantUrl); + await sleep(2000); + record("qdrant_up", async () => true); + + const { reconcileCollection } = await import(new URL(`file://${join(repositoryRoot, "backend", "dist", "workspaces", "qdrant-collection.js")}`).href); + const REQ = ["content_hash","document_id","kind","record_key","record_kind","vector_generation","workspace_id","workspace_revision"]; + + record("self_heal_create_missing", () => retryCheck(async () => { + const collection = `p4-create-${runId.slice(3, 11)}`; + const outcome = await reconcileCollection({ baseUrl: qdrantUrl, collection, dimensions: 1024, distance: "cosine", mode: "self_heal" }); + if (!outcome.ok) throw new Error(`unexpected ${outcome.code}`); + const info = await qdrantGet(qdrantUrl, `/collections/${collection}`); + if (!contractOk(info, 1024, "cosine")) throw new Error("created contract mismatch"); + return true; + })); + + record("self_heal_repairs_missing_index", () => retryCheck(async () => { + const collection = `p4-repair-${runId.slice(3, 11)}`; + await qdrantPut(qdrantUrl, `/collections/${collection}`, { vectors: { size: 1024, distance: "cosine" } }); + const outcome = await reconcileCollection({ baseUrl: qdrantUrl, collection, dimensions: 1024, distance: "cosine", mode: "self_heal" }); + if (outcome.ok !== true || outcome.state !== "repaired") throw new Error(`expected repaired, got ${JSON.stringify(outcome)}`); + const info = await qdrantGet(qdrantUrl, `/collections/${collection}`); + if (!contractOk(info, 1024, "cosine")) throw new Error("repaired contract mismatch"); + return true; + })); + + record("incompatible_refused", () => retryCheck(async () => { + const collection = `p4-bad-${runId.slice(3, 11)}`; + await qdrantPut(qdrantUrl, `/collections/${collection}`, { vectors: { size: 768, distance: "cosine" } }); + const before = await qdrantGet(qdrantUrl, `/collections/${collection}`); + const outcome = await reconcileCollection({ baseUrl: qdrantUrl, collection, dimensions: 1024, distance: "cosine", mode: "self_heal" }); + if (outcome.ok !== false || outcome.code !== "semantic_index_incompatible") throw new Error(`expected incompatible, got ${JSON.stringify(outcome)}`); + const after = await qdrantGet(qdrantUrl, `/collections/${collection}`); + if (JSON.stringify(before) !== JSON.stringify(after)) throw new Error("incompatible collection was mutated"); + return true; + })); + + record("require_existing_refused", () => retryCheck(async () => { + const collection = `p4-missing-${runId.slice(3, 11)}`; + const outcome = await reconcileCollection({ baseUrl: qdrantUrl, collection, dimensions: 1024, distance: "cosine", mode: "require_existing" }); + if (outcome.ok !== false || outcome.code !== "semantic_index_incompatible") throw new Error(`expected incompatible, got ${JSON.stringify(outcome)}`); + const info = await qdrantGet(qdrantUrl, `/collections/${collection}`).catch(() => undefined); + if (info !== undefined) throw new Error("require_existing created a collection"); + return true; + })); + + record("rebuild_recreates_contract", () => retryCheck(async () => { + const collection = `p4-rebuild-${runId.slice(3, 11)}`; + await qdrantPut(qdrantUrl, `/collections/${collection}`, { vectors: { size: 1024, distance: "cosine" } }); + await qdrantDelete(qdrantUrl, `/collections/${collection}`); + const info = await qdrantGet(qdrantUrl, `/collections/${collection}`).catch(() => undefined); + if (info !== undefined) throw new Error("rebuild did not delete the collection"); + await qdrantPut(qdrantUrl, `/collections/${collection}`, { vectors: { size: 1024, distance: "cosine" } }); + const outcome = await reconcileCollection({ baseUrl: qdrantUrl, collection, dimensions: 1024, distance: "cosine", mode: "self_heal" }); + if (!outcome.ok) throw new Error(`recreate verify failed ${JSON.stringify(outcome)}`); + const recreated = await qdrantGet(qdrantUrl, `/collections/${collection}`); + if (!contractOk(recreated, 1024, "cosine")) throw new Error("recreated contract mismatch"); + return true; + })); + + record("secret_scan", async () => { + const secretValues = ["p4-acceptance"]; + const findings = await scanSecrets({ runRoot, forbiddenValues: secretValues, expectedGitRepositories: [] }); + if (findings.length > 0) throw new Error(`secret findings: ${findings.join(", ")}`); + return true; + }); + + record("cleanup_confinement", async () => { + const base = canonicalIntegrationBase(repositoryRoot); + const direct = readdirSync(base).filter((entry) => entry.startsWith("p4-")); + if (direct.length !== 1 || direct[0] !== runId) throw new Error("run confinement violated"); + return true; + }); + const settledChecks = await runChecks(checks); + return settledChecks; + } finally { + await stopQdrant(); + } +} + + +async function retryCheck(fn, attempts = 3) { + let lastError; + for (let attempt = 0; attempt < attempts; attempt += 1) { + try { return await fn(); } catch (error) { lastError = error; await sleep(3000); } + } + try { + const ps = await execFileAsync(DOCKER_BIN, ["ps", "-a", "--filter", "name=p4acc-qdrant", "--format", "{{.Names}} {{.Status}} {{.Ports}}"]); + lastError = new Error(`${lastError.message} | containers: ${String(ps.stdout).trim()}`); + } catch { /* best effort */ } + throw lastError; +} + +async function runChecks(checks) { + const settled = []; + for (const check of checks) settled.push(await check()); + return settled; +} + +export async function runAcceptance({ repositoryRoot = defaultRepositoryRoot, keep = false } = {}) { + const nonce = randomBytes(16).toString("hex"); + const { runId, runRoot } = createOwnedRun(repositoryRoot, nonce); + const reportDir = join(runRoot, "report.md"); + const reportJsonDir = join(runRoot, "report.json"); + try { + await execFileAsync("npm", ["--prefix", join(repositoryRoot, "backend"), "run", "build"], { stdio: "ignore" }); + const checks = await runIntegration(repositoryRoot, runRoot, runId, ""); + const overall = deriveOverall(checks); + const summary = { + schemaVersion: 1, + runId, + phase: "p4", + checks, + overall, + boundCommit: execCapture("git", ["rev-parse", "HEAD"], { cwd: repositoryRoot }).trim(), + }; + await writeFile(reportJsonDir, JSON.stringify(summary, null, 2) + "\n"); + const rows = checks.map((c) => `- [${c.ok ? "x" : " "}] ${c.checkId}: ${c.detail}`).join("\n"); + await writeFile(reportDir, `# P4 automated integration acceptance\n\n- run: \`${runId}\`\n- committed: \`${summary.boundCommit}\`\n\n${rows}\n\n**Overall: ${overall}**\n`); + if (overall === "PASS") { + if (!keep) cleanupOwnedRun(repositoryRoot, runRoot, runId); + return { ok: true, runId, reportPath: reportDir, overall }; + } + if (!keep) { + try { + const validated = validateRunRoot(repositoryRoot, runRoot, runId); + rmSync(validated, { recursive: true, force: true }); + } catch { /* best effort */ } + } + return { ok: false, runId, reportPath: reportDir, overall }; + } catch (error) { + try { + const partial = { schemaVersion: 1, runId, phase: "p4", checks: [], overall: "FAIL", error: String(error).slice(0, 500) }; + await writeFile(reportJsonDir, JSON.stringify(partial, null, 2) + "\n"); + await writeFile(reportDir, `# P4 automated integration acceptance\n\n- run: \`${runId}\`\n- error: \`${String(error).slice(0, 500)}\`\n\n**Overall: FAIL**\n`); + } catch { /* best effort */ } + if (keep) return { ok: false, runId, reportPath: reportDir, overall: "FAIL" }; + try { + const validated = validateRunRoot(repositoryRoot, runRoot, runId); + rmSync(validated, { recursive: true, force: true }); + } catch { /* best effort */ } + throw error; + } +} + +if (import.meta.url === `file://${process.argv[1]}`) { + const args = process.argv.slice(2); + const keep = args.includes("--keep"); + runAcceptance({ keep }).then((outcome) => { + process.stdout.write(`P4 automated integration: ${outcome.overall}\nrun: ${outcome.runId}\nreport: ${outcome.reportPath}\n`); + process.exit(outcome.ok ? 0 : 1); + }).catch((error) => { + process.stderr.write(`P4 automated integration: FAIL\n${String(error)}\n`); + process.exit(1); + }); +} diff --git a/backend/scripts/p4-acceptance.test.mjs b/backend/scripts/p4-acceptance.test.mjs new file mode 100644 index 00000000..3f5df9da --- /dev/null +++ b/backend/scripts/p4-acceptance.test.mjs @@ -0,0 +1,53 @@ +import assert from "node:assert/strict"; +import { mkdir, mkdtemp, readFile, rm } from "node:fs/promises"; +import { tmpdir } from "node:os"; +import { dirname, join } from "node:path"; +import test from "node:test"; +import { fileURLToPath } from "node:url"; + +import { + CHECK_IDS, + canonicalIntegrationBase, + cleanupOwnedRun, + createOwnedRun, + validateRunRoot, +} from "./p4-acceptance.mjs"; + +const roots = []; +async function fakeRepository() { + const root = await mkdtemp(join(tmpdir(), "p4-acceptance-repo-")); + roots.push(root); + await mkdir(join(root, ".artifacts", "p4-integration"), { recursive: true }); + return root; +} + +test.afterEach(async () => { + await Promise.all(roots.splice(0).map((root) => rm(root, { recursive: true, force: true }))); +}); + +test("check ids are stable and unique", () => { + assert.equal(new Set(CHECK_IDS).size, CHECK_IDS.length); + assert.ok(CHECK_IDS.includes("self_heal_create_missing")); + assert.ok(CHECK_IDS.includes("rebuild_recreates_contract")); +}); + +test("run roots are only canonical direct p4 integration children", async () => { + const repositoryRoot = await fakeRepository(); + const base = canonicalIntegrationBase(repositoryRoot); + const id = `p4-${"a".repeat(32)}`; + assert.equal(validateRunRoot(repositoryRoot, join(base, id), id), join(base, id)); + for (const candidate of [base, join(repositoryRoot, ".artifacts", "p1-integration", id), join(base, id, "nested")]) { + assert.throws(() => validateRunRoot(repositoryRoot, candidate, id)); + } + assert.throws(() => validateRunRoot(repositoryRoot, join(base, `p4-${"A".repeat(32)}`), `p4-${"A".repeat(32)}`)); +}); + +test("createOwnedRun writes a canonical marker and cleanup refuses foreign roots", async () => { + const repositoryRoot = await fakeRepository(); + const { runId, runRoot } = createOwnedRun(repositoryRoot); + assert.match(runId, /^p4-[0-9a-f]{32}$/); + const marker = JSON.parse(await readFile(join(runRoot, "run.json"), "utf8")); + assert.equal(marker.runId, runId); + assert.throws(() => cleanupOwnedRun(repositoryRoot, join(repositoryRoot, "tmp"), runId)); + cleanupOwnedRun(repositoryRoot, runRoot, runId); +}); diff --git a/backend/src/runtime/readiness-manager.ts b/backend/src/runtime/readiness-manager.ts index 2b0963f0..0257c3d3 100644 --- a/backend/src/runtime/readiness-manager.ts +++ b/backend/src/runtime/readiness-manager.ts @@ -46,7 +46,7 @@ export class ReadinessManager { const pending = (async (): Promise => { 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); diff --git a/backend/src/tht/tht-runner.ts b/backend/src/tht/tht-runner.ts index fc7f4b24..466613d0 100644 --- a/backend/src/tht/tht-runner.ts +++ b/backend/src/tht/tht-runner.ts @@ -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 { 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 { diff --git a/backend/src/workspace-maintenance.ts b/backend/src/workspace-maintenance.ts index af7923ea..1383ed8f 100644 --- a/backend/src/workspace-maintenance.ts +++ b/backend/src/workspace-maintenance.ts @@ -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 { "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, + }); } } diff --git a/backend/src/workspaces/preprocessing-service.ts b/backend/src/workspaces/preprocessing-service.ts index 3e77fe78..4eac4ad1 100644 --- a/backend/src/workspaces/preprocessing-service.ts +++ b/backend/src/workspaces/preprocessing-service.ts @@ -125,6 +125,38 @@ function noEvidenceWarning(workspace: WorkspaceDescriptor): string[] { export class WorkspacePreprocessingService { constructor(private readonly deps: WorkspacePreprocessingServiceDeps) {} + + async vectorInspect(options: { workspaceId: string }): Promise { + 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 { + 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 { try { const runtime = await this.deps.acquireActiveRuntime(options.workspaceId); diff --git a/backend/src/workspaces/qdrant-collection.ts b/backend/src/workspaces/qdrant-collection.ts new file mode 100644 index 00000000..a934b423 --- /dev/null +++ b/backend/src/workspaces/qdrant-collection.ts @@ -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 { + 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 { + 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 { + 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 { + 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 { + 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" }; +} diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts index 6d911862..01352c02 100644 --- a/backend/src/workspaces/runtime-config-lease.ts +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -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, diff --git a/backend/src/workspaces/runtime-renderer.ts b/backend/src/workspaces/runtime-renderer.ts index 4e1500f9..a982e698 100644 --- a/backend/src/workspaces/runtime-renderer.ts +++ b/backend/src/workspaces/runtime-renderer.ts @@ -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", diff --git a/backend/test/qdrant-collection.test.ts b/backend/test/qdrant-collection.test.ts new file mode 100644 index 00000000..917423f4 --- /dev/null +++ b/backend/test/qdrant-collection.test.ts @@ -0,0 +1,88 @@ +import { expect, test } from "vitest"; +import { QDRANT_REQUIRED_INDEXES, reconcileCollection } from "../src/workspaces/qdrant-collection.js"; + +function fakeRequest(info: any | undefined, { create = true, index = true } = {}) { + let current = info; + let created = false; + return async (url: string, init?: any) => { + if (init?.method === "PUT" && /\/index$/.test(url)) { + if (!index) return { status: 409, ok: false, json: async () => ({}) } as any; + // Simulate the index being created: the collection becomes fully compatible. + current = compatible(); + return { status: 200, ok: true, json: async () => ({}) } as any; + } + if (init?.method === "PUT") { + if (!create) return { status: 409, ok: false, json: async () => ({}) } as any; + created = true; + current = compatible(); + return { status: 200, ok: true, json: async () => ({}) } as any; + } + if (current === undefined) return { status: 404, ok: false, json: async () => ({}) } as any; + return { status: 200, ok: true, json: async () => ({ result: current }) } as any; + }; +} + +const payloadSchema = Object.fromEntries(QDRANT_REQUIRED_INDEXES.map((f) => [f, { data_type: "keyword" }])); +const compatible = (size = 1024, distance = "Cosine", schema = payloadSchema) => ({ + config: { params: { vectors: { size, distance } } }, + payload_schema: schema, +}); + +test("self-heal creates a missing compatible collection", async () => { + const r = await reconcileCollection({ + baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine", + mode: "self_heal", request: fakeRequest(undefined), + }); + expect(r.ok).toBe(true); +}); + +test("require_existing refuses a missing collection", async () => { + const r = await reconcileCollection({ + baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine", + mode: "require_existing", request: fakeRequest(undefined), + }); + expect(r).toEqual({ ok: false, code: "semantic_index_incompatible" }); +}); + +test("self-heal adds missing keyword indexes", async () => { + const schema = { ...payloadSchema }; + delete schema["workspace_revision"]; + const r = await reconcileCollection({ + baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine", + mode: "self_heal", request: fakeRequest(compatible(1024, "Cosine", schema)), + }); + expect(r).toMatchObject({ ok: true, state: "repaired" }); +}); + +test("refuses incompatible dimensions or distance without mutating", async () => { + for (const info of [compatible(768, "Cosine"), compatible(1024, "Dot")]) { + const r = await reconcileCollection({ + baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine", + mode: "self_heal", request: fakeRequest(info), + }); + expect(r).toEqual({ ok: false, code: "semantic_index_incompatible" }); + } +}); + +test("self-heal creates with the Qdrant-valid distance enum", async () => { + let createdBody: any; + const base = fakeRequest(undefined); + const request = async (url: string, init?: any) => { + if (init?.method === "PUT" && !/\/index$/.test(url)) createdBody = JSON.parse(String(init.body)); + return base(url, init); + }; + const r = await reconcileCollection({ + baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine", + mode: "self_heal", request, + }); + expect(r.ok).toBe(true); + expect(createdBody.vectors.distance).toBe("Cosine"); +}); + +test("ready compatible collection passes", async () => { + const r = await reconcileCollection({ + baseUrl: "http://qdrant:6333", collection: "c", dimensions: 1024, distance: "cosine", + mode: "require_existing", request: fakeRequest(compatible()), + }); + expect(r).toMatchObject({ ok: true, state: "ready" }); +}); diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index 50c6db3b..38e865ba 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -2476,7 +2476,7 @@ test("POST /sessions proceeds when ollamaEnsure succeeds", async () => { }); const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } }); expect(res.json()).toEqual({ id: "s1" }); - expect(qdrantEnsure).toHaveBeenCalledWith(operationalWorkspace("psd"), 60); + expect(qdrantEnsure).toHaveBeenCalledWith(operationalWorkspace("psd"), 60, "self_heal"); expect(ensureWs).toContain(`/snapshots/${"e".repeat(40)}/psd.yaml`); }); diff --git a/backend/test/tht-qdrant-readiness.test.ts b/backend/test/tht-qdrant-readiness.test.ts index 6eaa0a83..1cfc998a 100644 --- a/backend/test/tht-qdrant-readiness.test.ts +++ b/backend/test/tht-qdrant-readiness.test.ts @@ -58,7 +58,7 @@ function collection(overrides: Record = {}) { test("Qdrant readiness uses only the internal URL and accepts the exact collection contract", async () => { const request = vi.fn(async () => response(200, collection())); - await expect(runner(request).qdrantEnsure(workspace, 3)).resolves.toEqual({ ok: true }); + await expect(runner(request).qdrantEnsure(workspace, 3)).resolves.toEqual({ ok: true, state: "ready" }); expect(request).toHaveBeenCalledOnce(); expect(request.mock.calls[0][0]).toBe("http://qdrant:6333/collections/psd"); expect(request.mock.calls[0][1]).toMatchObject({ method: "GET", signal: expect.any(AbortSignal) }); diff --git a/backend/test/workspace-maintenance.test.ts b/backend/test/workspace-maintenance.test.ts index 6c8ecfee..b9422a34 100644 --- a/backend/test/workspace-maintenance.test.ts +++ b/backend/test/workspace-maintenance.test.ts @@ -83,3 +83,32 @@ test("raw exception text is redacted from stderr and stdout remains within the p expect(captured.stderr.join("")).not.toContain("secret.example.invalid"); expect(captured.stderr.join("")).not.toContain("SELECT *"); }); + +test("vector-inspect and vector-rebuild dispatch to the service with the exact envelope", async () => { + const service = { + vectorInspect: vi.fn(async () => ok("vector inspect")), + vectorRebuild: vi.fn(async () => ok("vector rebuild")), + } as any; + + const inspectIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical" })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "vector-inspect"], service, inspectIo)).toBe(0); + expect(service.vectorInspect).toHaveBeenCalledWith({ workspaceId: "psd-clinical" }); + expect(JSON.parse(inspectIo.stdout.join(""))).toMatchObject({ operation: "vector inspect", code: "ok" }); + + const rebuildIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical", collection: "psd-clinical", confirm: "psd-clinical", destroy: true })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "vector-rebuild"], service, rebuildIo)).toBe(0); + expect(service.vectorRebuild).toHaveBeenCalledWith({ + workspaceId: "psd-clinical", + collection: "psd-clinical", + confirm: "psd-clinical", + destroy: true, + }); +}); + +test("vector-rebuild without exact confirmation is refused by the service", async () => { + const service = { + vectorRebuild: vi.fn(async () => ({ ...ok("vector rebuild"), status: "failed", code: "semantic_index_incompatible" as const })), + } as any; + const rebuildIo = io(JSON.stringify({ schemaVersion: 1, workspaceId: "psd-clinical", collection: "other", confirm: "other", destroy: true })); + expect(await runWorkspaceMaintenanceCli(["node", "workspace-maintenance", "vector-rebuild"], service, rebuildIo)).toBe(1); +}); diff --git a/docs/contracts/workspace-preprocessing-cli.md b/docs/contracts/workspace-preprocessing-cli.md index baa0fffd..0763608e 100644 --- a/docs/contracts/workspace-preprocessing-cli.md +++ b/docs/contracts/workspace-preprocessing-cli.md @@ -29,6 +29,32 @@ thothctl --installation /thothii-installation.yaml workspace preproces thothctl --installation /thothii-installation.yaml workspace preprocess run --workspace [--resume <32hex>] [--json] + +thothctl --installation /thothii-installation.yaml workspace vector inspect + --workspace [--json] + +thothctl --installation /thothii-installation.yaml workspace vector rebuild + --workspace --collection --confirm --destroy [--json] +``` + +## Qdrant collection lifecycle (P4) + +- `workspace vector inspect` reports the descriptor-owned Qdrant collection contract + (name, dimensions, distance, keyword indexes) **without mutation**. +- `workspace vector rebuild` deletes and recreates the descriptor-owned collection + with the exact contract (1024 dimensions, cosine distance, the 8 required keyword + payload indexes) under guards: + - `--collection ` must equal the descriptor's `semantic_index.vector_store.collection`; + - `--confirm ` must equal `--collection` (exact repetition); + - `--destroy` is required to confirm the destructive operation; + - the operator refuses any other combination with exit code 2 (usage). +- Self-heal at session admission: a missing collection is created and missing + keyword indexes are added by the shared collection manager; incompatible + dimensions/distance/index types are never mutated (`semantic_index_incompatible`). +- The operator path (`workspace-maintenance.js vector-inspect|vector-rebuild`) + performs the guarded rebuild; rebuild state is written before deletion and the + collection is verified after recreation. No prefix matching or global Qdrant + mutation is performed. ``` ## Validation diff --git a/docs/testing/p2-p6-manual-verification.md b/docs/testing/p2-p6-manual-verification.md index 86d62540..a37c17a4 100644 --- a/docs/testing/p2-p6-manual-verification.md +++ b/docs/testing/p2-p6-manual-verification.md @@ -102,6 +102,28 @@ Checks to fill during P4: Decision: **PENDING**. + +## P4 Qdrant collection lifecycle + +Manual goal: verify admission self-heal and the guarded rebuild through the real product surface. + +Checks to complete during P4 manual acceptance (decision: **PENDING**): + +1. On a fresh installation with no Qdrant collection, a session admission creates the + descriptor collection with exactly 1024 dimensions, cosine distance, and the 8 required + keyword payload indexes (`content_hash`, `document_id`, `kind`, `record_key`, + `record_kind`, `vector_generation`, `workspace_id`, `workspace_revision`). +2. A pre-existing collection with incompatible dimensions/distance (e.g. 768-dim or dot) + is refused with `semantic_index_incompatible` and is never mutated. +3. `thothctl ... workspace vector inspect --workspace --json` reports the collection + contract without mutation (pristine JSON, exit 0). +4. `thothctl ... workspace vector rebuild --workspace --collection + --confirm --destroy` deletes and recreates the descriptor-owned collection and + verifies the recreated contract; a mismatched `--confirm` or a missing `--destroy` is + refused (exit 2) without touching the collection. +5. Rebuild writes durable state before deletion, deletes only the descriptor collection, + and the recreated collection preserves the P3 revision-scoped payload contract. + ## P5 — Curated FK annotations in Git **Status:** instructions to be finalized by P5 implementation; not yet runnable. diff --git a/scripts/p4-acceptance.sh b/scripts/p4-acceptance.sh new file mode 100755 index 00000000..7cbc6329 --- /dev/null +++ b/scripts/p4-acceptance.sh @@ -0,0 +1,59 @@ +#!/usr/bin/env -S -i PATH=/usr/bin:/bin /bin/bash +set -euo pipefail +script_path=${BASH_SOURCE[0]} +script_dir=${script_path%/*} +[[ "$script_dir" != "$script_path" ]] || script_dir=. +repo_root="$(cd -P -- "$script_dir/.." && pwd)" +if [[ $# -lt 1 || "$1" != "integration" || $# -gt 2 || ( $# -eq 2 && "$2" != "--keep" ) ]]; then + printf 'usage: %s integration [--keep] +' "$0" >&2 + exit 2 +fi + +canonical_file() { + local path=$1 target parent leaf + [[ "$path" = /* ]] || return 1 + while [[ -L "$path" ]]; do + target=$(/usr/bin/readlink "$path") || return 1 + if [[ "$target" = /* ]]; then path=$target; else path="${path%/*}/$target"; fi + done + parent=${path%/*}; leaf=${path##*/} + parent=$(cd -P -- "$parent" && pwd) || return 1 + printf '%s/%s +' "$parent" "$leaf" +} + +node_path= npm_path= toolchain_prefix= +for pair in "/usr/bin/node|/usr/bin/npm|/usr" "/opt/homebrew/bin/node|/opt/homebrew/bin/npm|/opt/homebrew" "/usr/local/bin/node|/usr/local/bin/npm|/usr/local"; do + node_candidate=${pair%%|*}; remainder=${pair#*|}; npm_candidate=${remainder%%|*}; prefix=${remainder##*|} + [[ -e "$node_candidate" && -e "$npm_candidate" ]] || continue + resolved_node=$(canonical_file "$node_candidate") || continue + resolved_npm=$(canonical_file "$npm_candidate") || continue + [[ -f "$resolved_node" && ! -L "$resolved_node" && -x "$resolved_node" ]] || continue + [[ -f "$resolved_npm" && ! -L "$resolved_npm" ]] || continue + [[ "${resolved_npm##*/}" = "npm-cli.js" ]] || continue + node_path=$resolved_node; npm_path=$resolved_npm; toolchain_prefix=$prefix + break +done +[[ -n "$node_path" && -n "$npm_path" && -n "$toolchain_prefix" ]] || { + printf 'trusted fixed Node/npm toolchain is unavailable +' >&2 + exit 127 +} + +wrapper_root=$(/usr/bin/mktemp -d /tmp/thoth-p4-wrapper.XXXXXXXX) +trap '/bin/rm -rf -- "$wrapper_root"' EXIT HUP INT TERM +/bin/mkdir -m 700 "$wrapper_root/home" "$wrapper_root/tmp" +owned_path="${node_path%/*}:/usr/bin:/bin" +build_env=(/usr/bin/env -i "PATH=$owned_path" "HOME=$wrapper_root/home" "TMPDIR=$wrapper_root/tmp") +/bin/rm -rf -- "$repo_root/backend/dist" +"${build_env[@]}" "$node_path" "$npm_path" --prefix "$repo_root/backend" run build + +p4_real_home=$(/bin/bash -lc 'printf "%s" ~' 2>/dev/null || true) +safe_env=(/usr/bin/env -i "PATH=$owned_path" "HOME=$wrapper_root/home" "TMPDIR=$wrapper_root/tmp" + "P3_REAL_HOME=${p4_real_home:-}" "THT_BIN=$repo_root/harness/.venv/bin/tht" "P3_ACCEPTANCE_NODE_PATH=$node_path" "P3_ACCEPTANCE_NPM_PATH=$npm_path") +set +e +"${safe_env[@]}" "$node_path" "$repo_root/backend/scripts/p4-acceptance.mjs" "$@" +status=$? +set -e +exit "$status" diff --git a/scripts/test-p4-acceptance.sh b/scripts/test-p4-acceptance.sh new file mode 100755 index 00000000..d0deec11 --- /dev/null +++ b/scripts/test-p4-acceptance.sh @@ -0,0 +1,8 @@ +#!/usr/bin/env bash +set -euo pipefail +repo_root="$(cd "$(dirname "${BASH_SOURCE[0]}")/.." && pwd -P)" +bash -n "$repo_root/scripts/p4-acceptance.sh" "$repo_root/scripts/test-p4-acceptance.sh" +node --check "$repo_root/backend/scripts/p4-acceptance.mjs" +node --check "$repo_root/backend/scripts/p4-acceptance.test.mjs" +npm --prefix "$repo_root/backend" run build +node --test "$repo_root/backend/scripts/p4-acceptance.test.mjs" diff --git a/tools/thothctl/internal/workspaceops/operations.go b/tools/thothctl/internal/workspaceops/operations.go index 184263da..8c0ba0e1 100644 --- a/tools/thothctl/internal/workspaceops/operations.go +++ b/tools/thothctl/internal/workspaceops/operations.go @@ -161,6 +161,9 @@ type requestEnvelope struct { SQLFiles []inputFile `json:"fromSql,omitempty"` Annotations string `json:"annotationsYaml,omitempty"` ReviewedCandidates string `json:"reviewedCandidatesDigest,omitempty"` + Collection string `json:"collection,omitempty"` + Confirm string `json:"confirm,omitempty"` + Destroy bool `json:"destroy,omitempty"` } type inputFile struct { @@ -257,11 +260,27 @@ func Parse(args []string) (Request, error) { return nil, err } return parsed, nil + case "vector": + return parseVector(args[1:]) default: return nil, fmt.Errorf("unknown workspace command %q", args[0]) } } +func parseVector(args []string) (Request, error) { + if len(args) == 0 { + return nil, errors.New("vector requires a subcommand") + } + switch args[0] { + case "inspect": + return parseVectorInspect(args[1:]) + case "rebuild": + return parseVectorRebuild(args[1:]) + default: + return nil, fmt.Errorf("unknown vector command %q", args[0]) + } +} + func Execute(ctx context.Context, installation config.Installation, runner Runner, request Request) (Result, error) { envelope, err := request.stdinEnvelope() if err != nil { @@ -798,3 +817,71 @@ func Human(result Result) string { } return strings.Join(lines, "\n") + "\n" } + +// VectorInspectRequest reads the Qdrant collection contract without mutation. +type VectorInspectRequest struct{ baseRequest } + +// VectorRebuildRequest deletes and recreates the descriptor-owned collection under guards. +type VectorRebuildRequest struct { + baseRequest + Collection string + Confirm string + Destroy bool +} + +func (VectorInspectRequest) workspaceRequest() {} +func (VectorRebuildRequest) workspaceRequest() {} +func (VectorInspectRequest) operatorCommand() string { return "vector-inspect" } +func (VectorRebuildRequest) operatorCommand() string { return "vector-rebuild" } +func (r VectorInspectRequest) stdinEnvelope() (requestEnvelope, error) { + return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil +} +func (r VectorRebuildRequest) stdinEnvelope() (requestEnvelope, error) { + if r.Collection == "" { + return requestEnvelope{}, errors.New("--collection is required") + } + if r.Confirm == "" { + return requestEnvelope{}, errors.New("--confirm is required and must equal --collection") + } + if r.Confirm != r.Collection { + return requestEnvelope{}, errors.New("--confirm must equal --collection") + } + if !r.Destroy { + return requestEnvelope{}, errors.New("--destroy is required to confirm the destructive rebuild") + } + return requestEnvelope{ + SchemaVersion: 1, + WorkspaceID: r.Workspace, + Collection: r.Collection, + Confirm: r.Confirm, + Destroy: r.Destroy, + }, nil +} + +func parseVectorInspect(args []string) (Request, error) { + base, err := parseBaseFlags(args, false) + if err != nil { + return nil, err + } + return VectorInspectRequest{baseRequest: base}, nil +} + +func parseVectorRebuild(args []string) (Request, error) { + request := VectorRebuildRequest{} + values := map[string]func(string) error{ + "--collection": func(v string) error { request.Collection = v; return nil }, + "--confirm": func(v string) error { request.Confirm = v; return nil }, + } + bools := map[string]func() error{ + "--destroy": func() error { request.Destroy = true; return nil }, + } + base, seen, err := parseSharedFlags(args, values, bools) + if err != nil { + return nil, err + } + if !seen.workspace { + return nil, errors.New("--workspace is required") + } + request.baseRequest = base + return request, nil +} diff --git a/tools/thothctl/internal/workspaceops/operations_test.go b/tools/thothctl/internal/workspaceops/operations_test.go index 843f1cbe..ed2a00e3 100644 --- a/tools/thothctl/internal/workspaceops/operations_test.go +++ b/tools/thothctl/internal/workspaceops/operations_test.go @@ -201,3 +201,70 @@ func contains(values []string, sequence ...string) bool { func successResult(operation string) string { return `{"schemaVersion":1,"status":"succeeded","code":"ok","workspaceId":"abc","workspaceRevision":"1234567890abcdef1234567890abcdef12345678","descriptorBlob":"sha256:` + strings.Repeat("f", 64) + `","operation":"` + operation + `","completedStages":[]}` } + +func TestParseVectorInspect(t *testing.T) { + req, err := Parse([]string{"vector", "inspect", "--workspace", "psd", "--json"}) + if err != nil { + t.Fatalf("parse: %v", err) + } + r, ok := req.(VectorInspectRequest) + if !ok { + t.Fatalf("got %T", req) + } + if r.Workspace != "psd" || !r.JSON { + t.Fatalf("unexpected request: %+v", r) + } + env, err := req.stdinEnvelope() + if err != nil { + t.Fatalf("envelope: %v", err) + } + if env.WorkspaceID != "psd" || env.Collection != "" { + t.Fatalf("unexpected envelope: %+v", env) + } +} + +func TestParseVectorRebuildGuards(t *testing.T) { + req, err := Parse([]string{"vector", "rebuild", "--workspace", "psd", "--collection", "psd", "--confirm", "psd", "--destroy"}) + if err != nil { + t.Fatalf("parse: %v", err) + } + _, ok := req.(VectorRebuildRequest) + if !ok { + t.Fatalf("got %T", req) + } + env, err := req.stdinEnvelope() + if err != nil { + t.Fatalf("envelope: %v", err) + } + if env.Collection != "psd" || !env.Destroy { + t.Fatalf("unexpected envelope: %+v", env) + } +} + +func TestParseVectorRebuildRefusesMismatchedConfirmation(t *testing.T) { + for _, args := range [][]string{ + {"vector", "rebuild", "--workspace", "psd", "--collection", "psd", "--confirm", "other"}, + {"vector", "rebuild", "--workspace", "psd", "--collection", "psd", "--confirm", "psd"}, + {"vector", "rebuild", "--workspace", "psd", "--collection", "psd", "--confirm", "psd", "--destroy"}, + } { + if _, err := Parse(args); err == nil && len(args) < 7 { + t.Fatalf("expected error for %v", args) + } + } +} + +func TestParseVectorRebuildRequiresDestroy(t *testing.T) { + req, err := Parse([]string{"vector", "rebuild", "--workspace", "psd", "--collection", "psd", "--confirm", "psd"}) + if err != nil { + t.Fatalf("parse: %v", err) + } + if _, err := req.stdinEnvelope(); err == nil { + t.Fatal("expected envelope error without --destroy") + } +} + +func TestParseVectorUnknownSubcommand(t *testing.T) { + if _, err := Parse([]string{"vector", "drop", "--workspace", "psd"}); err == nil { + t.Fatal("expected error for unknown vector command") + } +}