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:
@@ -31,6 +31,31 @@
|
|||||||
`3b0726472e15c157…`.
|
`3b0726472e15c157…`.
|
||||||
- **Manual gate:** P3 walkthrough in `docs/testing/p2-p6-manual-verification.md`; decision
|
- **Manual gate:** P3 walkthrough in `docs/testing/p2-p6-manual-verification.md`; decision
|
||||||
**PENDING**. P4 starts only after an explicit new authorization.
|
**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 <id> --collection <name> --confirm <name> --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)
|
### P2 host preprocessing CLI — implementation complete, automated PASS, manual PENDING (2026-08-11)
|
||||||
|
|
||||||
|
|||||||
@@ -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);
|
||||||
|
});
|
||||||
|
}
|
||||||
@@ -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);
|
||||||
|
});
|
||||||
@@ -46,7 +46,7 @@ export class ReadinessManager {
|
|||||||
const pending = (async (): Promise<ReadinessResult> => {
|
const pending = (async (): Promise<ReadinessResult> => {
|
||||||
try {
|
try {
|
||||||
if (descriptor) {
|
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;
|
if (!qdrant.ok) return qdrant;
|
||||||
}
|
}
|
||||||
const ollama = await runner.ollamaEnsure(workspace, this.timeoutSec);
|
const ollama = await runner.ollamaEnsure(workspace, this.timeoutSec);
|
||||||
|
|||||||
@@ -10,6 +10,7 @@ import { clearPrincipalEnvironment, principalEnvironment, type PrincipalContext
|
|||||||
import { secretValue, type SecretBundleConfig } from "../config/secret-bundle.js";
|
import { secretValue, type SecretBundleConfig } from "../config/secret-bundle.js";
|
||||||
import { renderWorkspaceRuntimeFromSnapshotPath } from "../workspaces/runtime-config-lease.js";
|
import { renderWorkspaceRuntimeFromSnapshotPath } from "../workspaces/runtime-config-lease.js";
|
||||||
import {
|
import {
|
||||||
|
DEFAULT_SEMANTIC_RUNTIME,
|
||||||
type RuntimeInstallationOverlay,
|
type RuntimeInstallationOverlay,
|
||||||
type RuntimePaths,
|
type RuntimePaths,
|
||||||
type SemanticRuntimeConfig,
|
type SemanticRuntimeConfig,
|
||||||
@@ -19,6 +20,7 @@ import {
|
|||||||
validateOperationalWorkspace,
|
validateOperationalWorkspace,
|
||||||
type WorkspaceDescriptor,
|
type WorkspaceDescriptor,
|
||||||
} from "../workspaces/schema.js";
|
} from "../workspaces/schema.js";
|
||||||
|
import { reconcileCollection } from "../workspaces/qdrant-collection.js";
|
||||||
|
|
||||||
export interface ThtConfig extends SecretBundleConfig {
|
export interface ThtConfig extends SecretBundleConfig {
|
||||||
thtBin: string;
|
thtBin: string;
|
||||||
@@ -29,6 +31,8 @@ export interface ThtConfig extends SecretBundleConfig {
|
|||||||
secretRoots?: readonly string[];
|
secretRoots?: readonly string[];
|
||||||
semanticRuntime: SemanticRuntimeConfig;
|
semanticRuntime: SemanticRuntimeConfig;
|
||||||
qdrantRequest?: typeof fetch;
|
qdrantRequest?: typeof fetch;
|
||||||
|
/** "self_heal" for session admission (create missing collections/indexes), default "require_existing". */
|
||||||
|
qdrantCollectionMode?: "self_heal" | "require_existing";
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface RuntimeConfigLease {
|
export interface RuntimeConfigLease {
|
||||||
@@ -216,7 +220,7 @@ export class ThtRunner {
|
|||||||
throw new Error("registry workspace runtime requires an absolute data root");
|
throw new Error("registry workspace runtime requires an absolute data root");
|
||||||
})(),
|
})(),
|
||||||
secretRoots: this.cfg.secretRoots ?? [],
|
secretRoots: this.cfg.secretRoots ?? [],
|
||||||
semanticRuntime: this.cfg.semanticRuntime,
|
semanticRuntime: this.cfg.semanticRuntime ?? DEFAULT_SEMANTIC_RUNTIME,
|
||||||
});
|
});
|
||||||
const path = this.createRuntimeSnapshot(rendered.renderedConfig);
|
const path = this.createRuntimeSnapshot(rendered.renderedConfig);
|
||||||
let released = false;
|
let released = false;
|
||||||
@@ -561,6 +565,7 @@ export class ThtRunner {
|
|||||||
async qdrantEnsure(
|
async qdrantEnsure(
|
||||||
workspace: WorkspaceDescriptor,
|
workspace: WorkspaceDescriptor,
|
||||||
timeoutSec: number,
|
timeoutSec: number,
|
||||||
|
mode: "self_heal" | "require_existing" = "require_existing",
|
||||||
): Promise<QdrantEnsureResult> {
|
): Promise<QdrantEnsureResult> {
|
||||||
let descriptor;
|
let descriptor;
|
||||||
try {
|
try {
|
||||||
@@ -572,32 +577,16 @@ export class ThtRunner {
|
|||||||
const controller = new AbortController();
|
const controller = new AbortController();
|
||||||
const timer = setTimeout(() => controller.abort(), Math.max(1, timeoutSec) * 1000);
|
const timer = setTimeout(() => controller.abort(), Math.max(1, timeoutSec) * 1000);
|
||||||
try {
|
try {
|
||||||
const url = new URL(
|
const checked = await reconcileCollection({
|
||||||
`/collections/${encodeURIComponent(collection.collection)}`,
|
baseUrl: this.cfg.semanticRuntime.internalQdrantUrl,
|
||||||
this.cfg.semanticRuntime.internalQdrantUrl,
|
collection: collection.collection,
|
||||||
);
|
dimensions: collection.dimensions,
|
||||||
const request = this.cfg.qdrantRequest ?? fetch;
|
distance: collection.distance,
|
||||||
const response = await request(url.toString(), { method: "GET", signal: controller.signal });
|
mode,
|
||||||
if (response.status === 404) {
|
request: this.cfg.qdrantRequest ?? fetch,
|
||||||
return { ok: false, code: "semantic_index_incompatible" };
|
signal: controller.signal,
|
||||||
}
|
});
|
||||||
if (!response.ok) return { ok: false, code: "workspace_not_activatable" };
|
return checked;
|
||||||
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" };
|
|
||||||
} catch {
|
} catch {
|
||||||
return { ok: false, code: "workspace_not_activatable" };
|
return { ok: false, code: "workspace_not_activatable" };
|
||||||
} finally {
|
} finally {
|
||||||
|
|||||||
@@ -18,7 +18,7 @@ export interface WorkspaceMaintenanceIo {
|
|||||||
writeStderr(value: string): void;
|
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(
|
function failureResult(
|
||||||
operation: string,
|
operation: string,
|
||||||
@@ -69,6 +69,8 @@ function parseRequest(command: string, stdin: string): Record<string, unknown> {
|
|||||||
"index-schema": ["schemaVersion", "workspaceId", "resumeRunId"],
|
"index-schema": ["schemaVersion", "workspaceId", "resumeRunId"],
|
||||||
"preprocess-evidence": ["schemaVersion", "workspaceId", "dryRun", "resumeRunId"],
|
"preprocess-evidence": ["schemaVersion", "workspaceId", "dryRun", "resumeRunId"],
|
||||||
"preprocess-run": ["schemaVersion", "workspaceId", "resumeRunId"],
|
"preprocess-run": ["schemaVersion", "workspaceId", "resumeRunId"],
|
||||||
|
"vector-inspect": ["schemaVersion", "workspaceId"],
|
||||||
|
"vector-rebuild": ["schemaVersion", "workspaceId", "collection", "confirm", "destroy"],
|
||||||
};
|
};
|
||||||
const allowed = allowedByCommand[command];
|
const allowed = allowedByCommand[command];
|
||||||
if (!allowed) throw new Error("unknown command");
|
if (!allowed) throw new Error("unknown command");
|
||||||
@@ -121,6 +123,15 @@ async function dispatch(command: Command, service: WorkspacePreprocessingService
|
|||||||
workspaceId: request.workspaceId as string,
|
workspaceId: request.workspaceId as string,
|
||||||
resumeRunId: request.resumeRunId as string | undefined,
|
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 {
|
export class WorkspacePreprocessingService {
|
||||||
constructor(private readonly deps: WorkspacePreprocessingServiceDeps) {}
|
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> {
|
async inspect(options: { workspaceId: string }): Promise<WorkspaceOperationResult> {
|
||||||
try {
|
try {
|
||||||
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
|
const runtime = await this.deps.acquireActiveRuntime(options.workspaceId);
|
||||||
|
|||||||
@@ -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;
|
bindings: RuntimeBindings;
|
||||||
bindingDigest: string;
|
bindingDigest: string;
|
||||||
renderedConfig: string;
|
renderedConfig: string;
|
||||||
|
semanticQdrantUrl: string;
|
||||||
}
|
}
|
||||||
|
|
||||||
export interface ActiveRenderedWorkspaceRuntime extends RenderedWorkspaceRuntime {
|
export interface ActiveRenderedWorkspaceRuntime extends RenderedWorkspaceRuntime {
|
||||||
@@ -71,6 +72,7 @@ export interface DeterministicRuntimeConfigLease extends RuntimeConfigLease {
|
|||||||
catalogBlob: string;
|
catalogBlob: string;
|
||||||
configDigest: string;
|
configDigest: string;
|
||||||
bindingDigest: string;
|
bindingDigest: string;
|
||||||
|
semanticQdrantUrl: string;
|
||||||
effectiveConfig: CanonicalEffectiveConfig;
|
effectiveConfig: CanonicalEffectiveConfig;
|
||||||
effectiveConfigIdentity: string;
|
effectiveConfigIdentity: string;
|
||||||
configFingerprint: string;
|
configFingerprint: string;
|
||||||
@@ -310,6 +312,7 @@ function renderWorkspaceRuntimeFromWorkspace(options: {
|
|||||||
installationOverlay: overlay,
|
installationOverlay: overlay,
|
||||||
bindings,
|
bindings,
|
||||||
bindingDigest: stableBindingDigest(bindings),
|
bindingDigest: stableBindingDigest(bindings),
|
||||||
|
semanticQdrantUrl: options.semanticRuntime.internalQdrantUrl,
|
||||||
renderedConfig: renderRuntimeConfig(
|
renderedConfig: renderRuntimeConfig(
|
||||||
options.workspace,
|
options.workspace,
|
||||||
bindings,
|
bindings,
|
||||||
@@ -508,6 +511,7 @@ export async function publishDeterministicRuntimeConfigLease(options: {
|
|||||||
catalogBlob: rendered.catalogBlob,
|
catalogBlob: rendered.catalogBlob,
|
||||||
configDigest,
|
configDigest,
|
||||||
bindingDigest: rendered.bindingDigest,
|
bindingDigest: rendered.bindingDigest,
|
||||||
|
semanticQdrantUrl: rendered.semanticQdrantUrl,
|
||||||
effectiveConfig,
|
effectiveConfig,
|
||||||
effectiveConfigIdentity: effectiveConfigIdentityValue,
|
effectiveConfigIdentity: effectiveConfigIdentityValue,
|
||||||
configFingerprint: configFingerprintValue,
|
configFingerprint: configFingerprintValue,
|
||||||
@@ -570,6 +574,7 @@ export async function publishDeterministicRuntimeConfigLease(options: {
|
|||||||
catalogBlob: rendered.catalogBlob,
|
catalogBlob: rendered.catalogBlob,
|
||||||
configDigest,
|
configDigest,
|
||||||
bindingDigest: rendered.bindingDigest,
|
bindingDigest: rendered.bindingDigest,
|
||||||
|
semanticQdrantUrl: rendered.semanticQdrantUrl,
|
||||||
effectiveConfig,
|
effectiveConfig,
|
||||||
effectiveConfigIdentity: effectiveConfigIdentityValue,
|
effectiveConfigIdentity: effectiveConfigIdentityValue,
|
||||||
configFingerprint: configFingerprintValue,
|
configFingerprint: configFingerprintValue,
|
||||||
|
|||||||
@@ -34,7 +34,7 @@ export interface SemanticRuntimeConfig {
|
|||||||
internalEmbeddingDimensions: number;
|
internalEmbeddingDimensions: number;
|
||||||
}
|
}
|
||||||
|
|
||||||
const DEFAULT_SEMANTIC_RUNTIME: SemanticRuntimeConfig = {
|
export const DEFAULT_SEMANTIC_RUNTIME: SemanticRuntimeConfig = {
|
||||||
internalQdrantUrl: "http://qdrant:6333",
|
internalQdrantUrl: "http://qdrant:6333",
|
||||||
internalEmbeddingUrl: "http://embedding:11434",
|
internalEmbeddingUrl: "http://embedding:11434",
|
||||||
internalEmbeddingModel: "qwen3-embedding:0.6b",
|
internalEmbeddingModel: "qwen3-embedding:0.6b",
|
||||||
|
|||||||
@@ -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" });
|
||||||
|
});
|
||||||
@@ -2476,7 +2476,7 @@ test("POST /sessions proceeds when ollamaEnsure succeeds", async () => {
|
|||||||
});
|
});
|
||||||
const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
const res = await app.inject({ method: "POST", url: "/sessions", payload: { question: "q" } });
|
||||||
expect(res.json()).toEqual({ id: "s1" });
|
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`);
|
expect(ensureWs).toContain(`/snapshots/${"e".repeat(40)}/psd.yaml`);
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -58,7 +58,7 @@ function collection(overrides: Record<string, unknown> = {}) {
|
|||||||
test("Qdrant readiness uses only the internal URL and accepts the exact collection contract", async () => {
|
test("Qdrant readiness uses only the internal URL and accepts the exact collection contract", async () => {
|
||||||
const request = vi.fn(async () => response(200, collection()));
|
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).toHaveBeenCalledOnce();
|
||||||
expect(request.mock.calls[0][0]).toBe("http://qdrant:6333/collections/psd");
|
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) });
|
expect(request.mock.calls[0][1]).toMatchObject({ method: "GET", signal: expect.any(AbortSignal) });
|
||||||
|
|||||||
@@ -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("secret.example.invalid");
|
||||||
expect(captured.stderr.join("")).not.toContain("SELECT *");
|
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);
|
||||||
|
});
|
||||||
|
|||||||
@@ -29,6 +29,32 @@ thothctl --installation <absolute>/thothii-installation.yaml workspace preproces
|
|||||||
|
|
||||||
thothctl --installation <absolute>/thothii-installation.yaml workspace preprocess run
|
thothctl --installation <absolute>/thothii-installation.yaml workspace preprocess run
|
||||||
--workspace <id> [--resume <32hex>] [--json]
|
--workspace <id> [--resume <32hex>] [--json]
|
||||||
|
|
||||||
|
thothctl --installation <absolute>/thothii-installation.yaml workspace vector inspect
|
||||||
|
--workspace <id> [--json]
|
||||||
|
|
||||||
|
thothctl --installation <absolute>/thothii-installation.yaml workspace vector rebuild
|
||||||
|
--workspace <id> --collection <name> --confirm <name> --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 <name>` must equal the descriptor's `semantic_index.vector_store.collection`;
|
||||||
|
- `--confirm <name>` 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
|
## Validation
|
||||||
|
|||||||
@@ -102,6 +102,28 @@ Checks to fill during P4:
|
|||||||
|
|
||||||
Decision: **PENDING**.
|
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 <id> --json` reports the collection
|
||||||
|
contract without mutation (pristine JSON, exit 0).
|
||||||
|
4. `thothctl ... workspace vector rebuild --workspace <id> --collection <name>
|
||||||
|
--confirm <name> --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
|
## P5 — Curated FK annotations in Git
|
||||||
|
|
||||||
**Status:** instructions to be finalized by P5 implementation; not yet runnable.
|
**Status:** instructions to be finalized by P5 implementation; not yet runnable.
|
||||||
|
|||||||
Executable
+59
@@ -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"
|
||||||
Executable
+8
@@ -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"
|
||||||
@@ -161,6 +161,9 @@ type requestEnvelope struct {
|
|||||||
SQLFiles []inputFile `json:"fromSql,omitempty"`
|
SQLFiles []inputFile `json:"fromSql,omitempty"`
|
||||||
Annotations string `json:"annotationsYaml,omitempty"`
|
Annotations string `json:"annotationsYaml,omitempty"`
|
||||||
ReviewedCandidates string `json:"reviewedCandidatesDigest,omitempty"`
|
ReviewedCandidates string `json:"reviewedCandidatesDigest,omitempty"`
|
||||||
|
Collection string `json:"collection,omitempty"`
|
||||||
|
Confirm string `json:"confirm,omitempty"`
|
||||||
|
Destroy bool `json:"destroy,omitempty"`
|
||||||
}
|
}
|
||||||
|
|
||||||
type inputFile struct {
|
type inputFile struct {
|
||||||
@@ -257,11 +260,27 @@ func Parse(args []string) (Request, error) {
|
|||||||
return nil, err
|
return nil, err
|
||||||
}
|
}
|
||||||
return parsed, nil
|
return parsed, nil
|
||||||
|
case "vector":
|
||||||
|
return parseVector(args[1:])
|
||||||
default:
|
default:
|
||||||
return nil, fmt.Errorf("unknown workspace command %q", args[0])
|
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) {
|
func Execute(ctx context.Context, installation config.Installation, runner Runner, request Request) (Result, error) {
|
||||||
envelope, err := request.stdinEnvelope()
|
envelope, err := request.stdinEnvelope()
|
||||||
if err != nil {
|
if err != nil {
|
||||||
@@ -798,3 +817,71 @@ func Human(result Result) string {
|
|||||||
}
|
}
|
||||||
return strings.Join(lines, "\n") + "\n"
|
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
|
||||||
|
}
|
||||||
|
|||||||
@@ -201,3 +201,70 @@ func contains(values []string, sequence ...string) bool {
|
|||||||
func successResult(operation string) string {
|
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":[]}`
|
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")
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user