diff --git a/PROJECT_STATE.md b/PROJECT_STATE.md index c3b179d8..3cad6963 100644 --- a/PROJECT_STATE.md +++ b/PROJECT_STATE.md @@ -1,6 +1,6 @@ # ThothII — Project State -> Starting-point snapshot for new sessions. Last updated: 2026-08-08 (Task 13 verification audit). +> Starting-point snapshot for new sessions. Last updated: 2026-08-08 (final review fix round 1). > Point a fresh session here ("read PROJECT_STATE.md") before substantial work. ## Internal Qdrant + Ollama semantic infrastructure — LIVE 2026-08-08 @@ -14,6 +14,12 @@ - **Semantic contract.** Internal semantic indexing is fixed to `qwen3-embedding:0.6b`, `1024` dimensions, and cosine distance. Schema-v3 descriptors are operational; schema-v1/v2 descriptors remain `migration_required` until an explicit reviewed migration writes schema version 3. One workspace owns one Qdrant collection, and schema, Evidence, and Memory records coexist inside that collection with payload `kind` separation. +- **Final review runtime barriers.** Operational routes, retained session pins, and runtime + rendering now require schema version 3 before resolving bindings, readiness, diagnostics, or + Pi. Session admission verifies the exact internal Qdrant collection (dimensions, cosine + distance, and required keyword payload indexes) before Ollama and before manifest persistence. + The Qdrant adapter binds every search/list/delete filter to its constructed workspace identity + and rejects conflicting caller namespaces. - **Boundary and persistence.** Only DWH and LLM remain external runtime application endpoints. There are no active external vector or embedding endpoint instructions, bindings, or secrets in the supported operator manuals. Qdrant remains a derived but persistent semantic index: the @@ -26,7 +32,9 @@ --confirm-project ` requires the exact repeated project confirmation, validates manifest and archive safety before stopping `qdrant`, stages rollback content, restores in place, and restarts `qdrant` only if it was previously running. Restore does not migrate legacy workspace - descriptors, rename collections, or repair a semantic-index incompatibility. + descriptors, rename collections, or repair a semantic-index incompatibility. Backup and restore + share one atomic Docker-daemon lock per Compose project/Qdrant volume; contenders fail before + volume resolution, and cleanup removes the lock only when its ownership labels still match. - **Verification recorded for Task 13 final audit.** On Apple M4 Pro (`Darwin 25.5.0`, Docker Server `29.6.2 linux/arm64`), harness pytest passed **819 passed / 4 deselected**; backend Vitest passed **464/464** plus TypeScript and build; @@ -59,6 +67,14 @@ SQL, L2 legacy fixtures, gitignored task notes, and historical reference notes. No active schema-v3 operator manual or supported runtime deployment path retains external vector or embedding endpoint coupling. +- **Final review fix verification.** Backend Vitest passed **477/477** plus TypeScript and build; + harness pytest passed **824 passed / 4 deselected** with the existing 74 warnings; touched Python + files are Ruff-clean. The deterministic backup/restore safety test proves lock ownership, + backup–backup and backup–restore contention, rollback, and cleanup. Internal semantic Compose + and no-deployment-coupling contracts pass. A fresh one-shot unified deployment smoke reached its + pre-existing `thothctl` bad-candidate rollback scenario and failed there before backup/restore; + its exact resource-cleanup proof passed, so this run is recorded as a limitation, not as a + successful revalidation of the earlier Task 13 unified-smoke result. ## Historical snapshots and archived reference notes diff --git a/backend/src/routes/sessions.ts b/backend/src/routes/sessions.ts index 785473e2..d0da061b 100644 --- a/backend/src/routes/sessions.ts +++ b/backend/src/routes/sessions.ts @@ -8,7 +8,7 @@ import type { PrincipalContext } from "../auth/principal.js"; import type { ReadinessManager } from "../runtime/readiness-manager.js"; import type { ListModelsFn } from "./meta.js"; import type { WorkspaceRegistry } from "../workspaces/registry.js"; -import type { WorkspaceDescriptor } from "../workspaces/schema.js"; +import { validateOperationalWorkspace, type WorkspaceDescriptor } from "../workspaces/schema.js"; import type { MaintenanceBarrier } from "../runtime/maintenance-gate.js"; const BOOTSTRAP_FAILURE_MESSAGE = @@ -108,7 +108,11 @@ export function sessionRoutes( const isNotFound = (error: unknown) => /not found|non trovata|inesistente|404/i.test(error instanceof Error ? error.message : String(error)); - type LocatedSession = { manifest: any; workspaceConfigPath: string }; + type LocatedSession = { + manifest: any; + workspaceConfigPath: string; + workspace?: WorkspaceDescriptor; + }; const workspaceRevisionUnavailable = () => Object.assign( new Error("workspace revision unavailable"), { code: "workspace_revision_unavailable" }, @@ -168,7 +172,12 @@ export function sessionRoutes( if (!saved.workspace_id || !saved.workspace_revision) return located; try { const pinned = await d.workspaceRegistry.readPinned(saved.workspace_id, saved.workspace_revision); - return { ...located, workspaceConfigPath: pinned.workspaceConfigPath ?? (pinned as any).revision?.snapshotPath }; + const workspace = validateOperationalWorkspace(pinned.workspace); + return { + ...located, + workspace, + workspaceConfigPath: pinned.workspaceConfigPath ?? (pinned as any).revision?.snapshotPath, + }; } catch { throw workspaceRevisionUnavailable(); } @@ -319,6 +328,7 @@ export function sessionRoutes( let workspaceConfigPath: string | undefined; let workspaceId: string | undefined; let workspaceRevision: string | undefined; + let workspaceDescriptor: WorkspaceDescriptor | undefined; let allowedModels: readonly string[] | undefined; if (requestedWorkspaceId) { try { @@ -344,6 +354,7 @@ export function sessionRoutes( workspaceConfigPath = resolved.revision.snapshotPath; workspaceId = resolved.revision.id; workspaceRevision = resolved.revision.commit; + workspaceDescriptor = resolved.workspace; allowedModels = resolved.workspace.llm_policy.allowed; } catch { return reply.code(409).send({ @@ -362,8 +373,13 @@ export function sessionRoutes( // runtime owned by this principal, while runtimes belonging to other users remain intact. // Optional chaining preserves the deliberately narrow manager stubs used by route tests. for (const id of d.mgr.teardownForPrincipal?.(principal) ?? []) boundRuntimes.delete(id); - const ensure = await d.readiness.ensure(workspaceConfigPath ?? "", principal); - if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE }); + const ensure = await d.readiness.ensure( + workspaceConfigPath ?? "", principal, workspaceDescriptor, + ); + if (!ensure.ok) return reply.code(503).send({ + error: READINESS_FAILURE_MESSAGE, + ...(ensure.code ? { code: ensure.code } : {}), + }); // Local-only: verify the DWH is reachable BEFORE creating the session, so a dropped // VPN surfaces as an up-front alert instead of a session that spawns Pi and then dies // in bootstrap retrieval. `code` lets the client show a specific message. @@ -546,7 +562,12 @@ export function sessionRoutes( workspace_id?: string; workspace_revision?: string; }; let workspaceConfigPath: string; - try { workspaceConfigPath = (await resolveSessionWorkspace(located)).workspaceConfigPath; } + let workspaceDescriptor: WorkspaceDescriptor | undefined; + try { + const resolved = await resolveSessionWorkspace(located); + workspaceConfigPath = resolved.workspaceConfigPath; + workspaceDescriptor = resolved.workspace; + } catch { return unavailableWorkspaceReply(reply); } try { settings = await d.getSettings(principal); } catch { return storageFailure(reply); } // This check belongs inside the per-session lock: a preceding cold Resume may have @@ -558,8 +579,13 @@ export function sessionRoutes( return reply.code(200).send({ id, alreadyActive: true }); } } - const ensure = await d.readiness.ensure(workspaceConfigPath ?? "", principal); - if (!ensure.ok) return reply.code(503).send({ error: READINESS_FAILURE_MESSAGE }); + const ensure = await d.readiness.ensure( + workspaceConfigPath ?? "", principal, workspaceDescriptor, + ); + if (!ensure.ok) return reply.code(503).send({ + error: READINESS_FAILURE_MESSAGE, + ...(ensure.code ? { code: ensure.code } : {}), + }); const options = { provider: saved?.provider, model: saved?.model, diff --git a/backend/src/routes/workspaces.ts b/backend/src/routes/workspaces.ts index 5f727fd4..3ce75baa 100644 --- a/backend/src/routes/workspaces.ts +++ b/backend/src/routes/workspaces.ts @@ -19,6 +19,7 @@ import { parseWorkspaceYaml, serializeWorkspaceYaml, validateCanonicalWorkspace, + validateOperationalWorkspace, type CanonicalWorkspace, type WorkspaceDescriptor, } from "../workspaces/schema.js"; @@ -327,9 +328,22 @@ export function workspaceRoutes(app: FastifyInstance, deps: WorkspaceRoutesDeps) app.post("/workspaces/:id/test", async (request, reply) => { try { const { id } = z.object({ id: workspaceId }).parse(request.params); - const { workspace } = await deps.registry.read(id); - const bindings = resolveRuntimeBindings(workspace, process.env, deps.config.secretRoots); - return await deps.diagnose(workspace, bindings, { writeProbe: false }); + const { workspace, revision } = await deps.registry.read(id); + if (revision.state !== "operational") { + throw new WorkspaceRegistryError( + "workspace_not_activatable", "Workspace requires explicit migration", + ); + } + let operational: CanonicalWorkspace; + try { + operational = validateOperationalWorkspace(workspace); + } catch { + throw new WorkspaceRegistryError( + "workspace_not_activatable", "Workspace requires explicit migration", + ); + } + const bindings = resolveRuntimeBindings(operational, process.env, deps.config.secretRoots); + return await deps.diagnose(operational, bindings, { writeProbe: false }); } catch (error) { return errorReply(reply, error); } diff --git a/backend/src/runtime/readiness-manager.ts b/backend/src/runtime/readiness-manager.ts index 1c49494e..2b0963f0 100644 --- a/backend/src/runtime/readiness-manager.ts +++ b/backend/src/runtime/readiness-manager.ts @@ -1,9 +1,16 @@ -import type { OllamaEnsureResult, ThtRunner } from "../tht/tht-runner.js"; +import type { + OllamaEnsureResult, + SemanticReadinessCode, + ThtRunner, +} from "../tht/tht-runner.js"; import type { PrincipalContext } from "../auth/principal.js"; +import type { WorkspaceDescriptor } from "../workspaces/schema.js"; + +export type ReadinessResult = OllamaEnsureResult & { code?: SemanticReadinessCode }; interface ReadyEntry { expiresAt: number; - result: OllamaEnsureResult; + result: ReadinessResult; } /** @@ -11,7 +18,7 @@ interface ReadyEntry { * Failures are deliberately not cached so a submit can retry after a transient outage. */ export class ReadinessManager { - private inFlight = new Map>(); + private inFlight = new Map>(); private ready = new Map(); constructor( @@ -21,7 +28,11 @@ export class ReadinessManager { private now: () => number = Date.now, ) {} - ensure(workspace = "", principal?: PrincipalContext): Promise { + ensure( + workspace = "", + principal?: PrincipalContext, + descriptor?: WorkspaceDescriptor, + ): Promise { const key = `${principal?.issuer ?? ""}\0${principal?.subject ?? ""}\0${workspace}`; const cached = this.ready.get(key); if (cached && cached.expiresAt > this.now()) return Promise.resolve(cached.result); @@ -32,7 +43,20 @@ export class ReadinessManager { const runner = principal && typeof (this.tht as any).withPrincipal === "function" ? this.tht.withPrincipal(principal) : this.tht; - const pending = runner.ollamaEnsure(workspace, this.timeoutSec) + const pending = (async (): Promise => { + try { + if (descriptor) { + const qdrant = await runner.qdrantEnsure(descriptor, this.timeoutSec); + if (!qdrant.ok) return qdrant; + } + const ollama = await runner.ollamaEnsure(workspace, this.timeoutSec); + return ollama.ok + ? ollama + : { ...ollama, code: "workspace_not_activatable" }; + } catch { + return { ok: false, code: "workspace_not_activatable" }; + } + })() .then((result) => { if (result.ok) { this.ready.set(key, { result, expiresAt: this.now() + this.ttlMs }); diff --git a/backend/src/tht/tht-runner.ts b/backend/src/tht/tht-runner.ts index fbc64b83..43be047c 100644 --- a/backend/src/tht/tht-runner.ts +++ b/backend/src/tht/tht-runner.ts @@ -15,7 +15,11 @@ import { type RuntimePaths, type SemanticRuntimeConfig, } from "../workspaces/runtime-renderer.js"; -import { parseWorkspaceYaml } from "../workspaces/schema.js"; +import { + parseWorkspaceYaml, + validateOperationalWorkspace, + type WorkspaceDescriptor, +} from "../workspaces/schema.js"; export interface ThtConfig extends SecretBundleConfig { thtBin: string; @@ -25,6 +29,7 @@ export interface ThtConfig extends SecretBundleConfig { runtimeSnapshotRoot?: string; secretRoots?: readonly string[]; semanticRuntime: SemanticRuntimeConfig; + qdrantRequest?: typeof fetch; } export interface RuntimeConfigLease { @@ -64,6 +69,24 @@ export interface OllamaEnsureResult { model_name?: string; } +export type SemanticReadinessCode = "workspace_not_activatable" | "semantic_index_incompatible"; + +export interface QdrantEnsureResult { + ok: boolean; + code?: SemanticReadinessCode; +} + +const REQUIRED_QDRANT_PAYLOAD_INDEXES = [ + "content_hash", + "document_id", + "kind", + "record_key", + "record_kind", + "vector_generation", + "workspace_id", + "workspace_revision", +] as const; + interface RuntimeSnapshot { path: string; dev: number; @@ -136,7 +159,7 @@ export class ThtRunner { if (before.dev !== after.dev || before.ino !== after.ino || before.size !== after.size) { throw new Error("workspace snapshot changed while reading"); } - const workspace = parseWorkspaceYaml(source); + const workspace = validateOperationalWorkspace(parseWorkspaceYaml(source)); if (workspace.workspace.id !== identity.workspaceId) { throw new Error("workspace snapshot identity does not match its path"); } @@ -537,4 +560,52 @@ export class ThtRunner { error: parsed?.error ?? (stderr.trim() || `tht ollama ensure exit ${code}`), }; } + + async qdrantEnsure( + workspace: WorkspaceDescriptor, + timeoutSec: number, + ): Promise { + let descriptor; + try { + descriptor = validateOperationalWorkspace(workspace); + } catch { + return { ok: false, code: "workspace_not_activatable" }; + } + const collection = descriptor.semantic_index.vector_store; + 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" }; + } catch { + return { ok: false, code: "workspace_not_activatable" }; + } finally { + clearTimeout(timer); + controller.abort(); + } + } } diff --git a/backend/src/workspaces/registry.ts b/backend/src/workspaces/registry.ts index 9f3b3c90..84b45380 100644 --- a/backend/src/workspaces/registry.ts +++ b/backend/src/workspaces/registry.ts @@ -13,6 +13,7 @@ import { isCanonicalWorkspace, parseWorkspaceYaml, serializeWorkspaceYaml, + validateOperationalWorkspace, type CanonicalWorkspace, type WorkspaceDescriptor, } from "./schema.js"; @@ -216,7 +217,9 @@ export class WorkspaceRegistry { } let workspace: WorkspaceDescriptor; try { - workspace = parseWorkspaceYaml(await readFile(revision.snapshotPath, "utf8")); + workspace = validateOperationalWorkspace( + parseWorkspaceYaml(await readFile(revision.snapshotPath, "utf8")), + ); } catch (error) { throw workspaceError(error); } @@ -259,7 +262,10 @@ export class WorkspaceRegistry { const snapshotPath = this.snapshotPath(safeCommit(commit), id); try { const source = await readFile(snapshotPath, "utf8"); - return { workspace: parseWorkspaceYaml(source), workspaceConfigPath: snapshotPath }; + return { + workspace: validateOperationalWorkspace(parseWorkspaceYaml(source)), + workspaceConfigPath: snapshotPath, + }; } catch (error) { throw workspaceError(error); } diff --git a/backend/test/readiness-manager.test.ts b/backend/test/readiness-manager.test.ts index be0cb039..5eb22fcc 100644 --- a/backend/test/readiness-manager.test.ts +++ b/backend/test/readiness-manager.test.ts @@ -1,6 +1,21 @@ import { expect, test } from "vitest"; import { ReadinessManager } from "../src/runtime/readiness-manager.js"; +const workspace = { + workspace: { schema_version: 3, id: "psd", name: "PSD", language: "it" }, + dwh: { + engine: "postgres", database: "warehouse", schema: "public", + supported_transports: ["postgres_direct"], + }, + semantic_index: { + vector_store: { engine: "qdrant", collection: "psd", dimensions: 1024, distance: "cosine" }, + embedding: { + provider: "ollama_internal", model: "qwen3-embedding:0.6b", dimensions: 1024, + }, + }, + llm_policy: { allowed: ["zai/glm-5.2"] }, +} as const; + function deferred() { let resolve!: (value: T) => void; const promise = new Promise((r) => { resolve = r; }); @@ -60,3 +75,31 @@ test("readiness does not cache failed results", async () => { await expect(readiness.ensure("psd")).resolves.toMatchObject({ ok: true }); expect(calls).toBe(2); }); + +test("readiness checks Qdrant before Ollama and skips Ollama on semantic incompatibility", async () => { + let ollamaCalls = 0; + const tht = { + qdrantEnsure: async () => ({ ok: false, code: "semantic_index_incompatible" }), + ollamaEnsure: async () => { ollamaCalls += 1; return { ok: true }; }, + } as any; + const readiness = new ReadinessManager(tht, 60); + + await expect(readiness.ensure("/registry/psd.yaml", undefined, workspace as any)).resolves.toEqual({ + ok: false, + code: "semantic_index_incompatible", + }); + expect(ollamaCalls).toBe(0); +}); + +test("readiness returns a sanitized activation code when a semantic probe throws", async () => { + const tht = { + qdrantEnsure: async () => { throw new Error("dial http://qdrant:6333/private"); }, + ollamaEnsure: async () => ({ ok: true }), + } as any; + const readiness = new ReadinessManager(tht, 60); + + await expect(readiness.ensure("/registry/psd.yaml", undefined, workspace as any)).resolves.toEqual({ + ok: false, + code: "workspace_not_activatable", + }); +}); diff --git a/backend/test/routes-sessions.test.ts b/backend/test/routes-sessions.test.ts index a1461109..451590e4 100644 --- a/backend/test/routes-sessions.test.ts +++ b/backend/test/routes-sessions.test.ts @@ -14,17 +14,34 @@ import { validateDeclarativePiConfig } from "../src/pi/managed-config.js"; const FAKE = path.resolve("../harness/tests/fake_pi/fake_pi_rpc.mjs"); const SCRIPT = path.resolve("../harness/tests/fake_pi/scripts/f1_disambiguation.json"); +function operationalWorkspace(id = "default") { + return { + workspace: { schema_version: 3, id, name: id, language: "en" }, + dwh: { + engine: "postgres", database: "warehouse", schema: "public", + supported_transports: ["postgres_direct"], + }, + semantic_index: { + vector_store: { + engine: "qdrant", collection: id, dimensions: 1024, distance: "cosine", + }, + embedding: { + provider: "ollama_internal", model: "qwen3-embedding:0.6b", dimensions: 1024, + }, + }, + llm_policy: { + allowed: ["zai/glm-5.2", "deepseek/deepseek-v4-pro", "local-qwen/qwen3.6-35b-a3b"], + }, + } as const; +} + const defaultWorkspaceRegistry = { list: vi.fn(async () => [{ id: "default", commit: "e".repeat(40), blob: "f".repeat(40), snapshotPath: `/data/workspace-registry/snapshots/${"e".repeat(40)}/default.yaml`, state: "operational", }]), read: vi.fn(async (id: string) => ({ - workspace: { - llm_policy: { - allowed: ["zai/glm-5.2", "deepseek/deepseek-v4-pro", "local-qwen/qwen3.6-35b-a3b"], - }, - }, + workspace: operationalWorkspace(id), revision: { id, commit: "e".repeat(40), blob: "f".repeat(40), snapshotPath: `/data/workspace-registry/snapshots/${"e".repeat(40)}/${id}.yaml`, state: "operational", @@ -33,9 +50,13 @@ const defaultWorkspaceRegistry = { }; function buildApp(config: Parameters[0], deps: Record = {}) { + const thtRunner = deps.thtRunner + ? { qdrantEnsure: async () => ({ ok: true }), ...(deps.thtRunner as object) } + : undefined; return buildRealApp(config, { workspaceRuntimeSupport: () => true, ...deps, + ...(thtRunner ? { thtRunner } : {}), workspaceRegistry: { ...defaultWorkspaceRegistry, ...(deps.workspaceRegistry as object | undefined) }, } as any); } @@ -663,7 +684,7 @@ test("session lifecycle locates a B session when installation default is A", asy ], readPinned: vi.fn(async (id: string, revision: string) => { expect([id, revision]).toEqual(["b-workspace", "c".repeat(40)]); - return { workspace: { llm_policy: { allowed: ["zai/glm-5.2"] } }, workspaceConfigPath: bPinnedPath }; + return { workspace: operationalWorkspace(id), workspaceConfigPath: bPinnedPath }; }), } as any, }); @@ -932,7 +953,7 @@ test("POST /sessions/:id/resume uses the manifest's retained workspace revision" getSettings: () => ({ workspace: "legacy" }) as any, workspaceRegistry: { readPinned: vi.fn(async () => ({ - workspace: { workspace: { id: "psd-clinical" } }, + workspace: operationalWorkspace("psd-clinical"), revision: { id: "psd-clinical", commit: "a".repeat(40), blob: "b".repeat(40), snapshotPath: "/data/workspace-registry/snapshots/aaaaaaaa/psd-clinical.yaml", state: "operational", @@ -968,6 +989,61 @@ test("POST /sessions/:id/resume returns a sanitized error when its retained revi expect(response.json()).toMatchObject({ code: "workspace_revision_unavailable" }); }); +test("POST /sessions/:id/resume rejects a pinned schema-v2 workspace before readiness or runtime", async () => { + const readiness = vi.fn(async () => ({ ok: true })); + const reopenSession = vi.fn(async () => {}); + const acquireWorkspaceRuntime = vi.fn(); + const createFor = vi.fn(); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + thtRunner: { + sessionShow: async () => ({ + status: "open", archived: false, + workspace_id: "psd-clinical", workspace_revision: "a".repeat(40), + }), + reopenSession, + acquireWorkspaceRuntime, + } as any, + readiness: { ensure: readiness } as any, + mgr: { get: () => undefined, createFor } as any, + getSettings: () => ({ workspace: "legacy" }) as any, + workspaceRegistry: { + readPinned: vi.fn(async () => ({ + workspace: { + workspace: { schema_version: 2, id: "psd-clinical", name: "PSD", language: "it" }, + dwh: { + engine: "postgres", database: "warehouse", schema: "public", + supported_transports: ["rest_api"], + }, + semantic_index: { + vector_store: { + engine: "pgvector", database: "warehouse", schema: "vectors", + collection: "documents", dimensions: 768, distance: "cosine", + supported_transports: ["rest_api"], + }, + embedding: { + provider: "ollama_compatible", model: "nomic-embed-text", dimensions: 768, + }, + }, + llm_policy: { allowed: ["zai/glm-5.2"] }, + }, + workspaceConfigPath: `/data/workspace-registry/snapshots/${"a".repeat(40)}/psd-clinical.yaml`, + })), + } as any, + }); + + const response = await app.inject({ method: "POST", url: "/sessions/pinned-v2/resume" }); + + expect(response.statusCode).toBe(409); + expect(response.json()).toEqual({ + code: "workspace_revision_unavailable", + error: "Session workspace configuration is unavailable. Check configuration and try again.", + }); + expect(readiness).not.toHaveBeenCalled(); + expect(reopenSession).not.toHaveBeenCalled(); + expect(acquireWorkspaceRuntime).not.toHaveBeenCalled(); + expect(createFor).not.toHaveBeenCalled(); +}); + test("a pruned pin blocks Resume but not active or mutation lifecycle routes", async () => { const activePath = "/registry/snapshots/aaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaaa/b-workspace.yaml"; const prunedError = "cannot read /registry/snapshots/secret-pruned-revision/b-workspace.yaml"; @@ -2340,11 +2416,35 @@ test("POST /sessions readiness failure returns one fixed public message without expect(res.statusCode).toBe(503); expect(res.json()).toEqual({ error: "Session services are not ready. Check configuration and connectivity, then try again.", + code: "workspace_not_activatable", }); expect(res.body).not.toMatch(/secret\.invalid|DO_NOT_LEAK|\/srv\/private\/model-key/); expect(createdCalled).toBe(false); }); +test.each(["semantic_index_incompatible", "workspace_not_activatable"] as const)( + "POST /sessions does not persist when Qdrant readiness returns %s", + async (code) => { + const sessionNew = vi.fn(async () => ({ id: "must-not-exist" })); + const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { + thtRunner: { sessionNew } as any, + readiness: { ensure: async () => ({ ok: false, code }) } as any, + getSettings: () => ({ workspace: "psd" }) as any, + }); + + const response = await app.inject({ + method: "POST", url: "/sessions", payload: { question: "q" }, + }); + + expect(response.statusCode).toBe(503); + expect(response.json()).toEqual({ + error: "Session services are not ready. Check configuration and connectivity, then try again.", + code, + }); + expect(sessionNew).not.toHaveBeenCalled(); + }, +); + test("POST /sessions returns storage 503 before creating a Pi runtime when session persistence fails", async () => { let piCreated = false; const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { @@ -2365,8 +2465,10 @@ test("POST /sessions returns storage 503 before creating a Pi runtime when sessi test("POST /sessions proceeds when ollamaEnsure succeeds", async () => { let ensureWs: string | undefined; + const qdrantEnsure = vi.fn(async () => ({ ok: true })); const app = buildApp(loadConfig({ THT_HARNESS_DIR: "../harness" }), { thtRunner: { + qdrantEnsure, ollamaEnsure: async (ws: string) => { ensureWs = ws; return { ok: true }; }, searchPack: async () => {}, sessionNew: async () => ({ id: "s1" }), @@ -2376,6 +2478,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(ensureWs).toContain(`/snapshots/${"e".repeat(40)}/psd.yaml`); }); @@ -2574,6 +2677,7 @@ test("POST /sessions/:id/resume readiness failure returns the same fixed public expect(res.statusCode).toBe(503); expect(res.json()).toEqual({ error: "Session services are not ready. Check configuration and connectivity, then try again.", + code: "workspace_not_activatable", }); expect(res.body).not.toMatch(/secret\.invalid|DO_NOT_LEAK|\/srv\/private\/resume-key/); }); diff --git a/backend/test/routes-workspaces.test.ts b/backend/test/routes-workspaces.test.ts index d422be5b..481a86fa 100644 --- a/backend/test/routes-workspaces.test.ts +++ b/backend/test/routes-workspaces.test.ts @@ -238,7 +238,7 @@ test("validates a canonical workspace and runs the injected installation diagnos expect(diagnose).not.toHaveBeenCalled(); }); -test("runs the injected installation diagnostic for a migration-required v2 workspace when legacy bindings resolve", async () => { +test("rejects a migration-required v2 workspace before resolving semantic diagnostics", async () => { const diagnose = vi.fn(async () => ({ activatable: false, diagnostics: [{ level: "error" as const, code: "binding_missing" as const, field: "THT_WS_PSD_CLINICAL_VECTOR_BASE_URL", message: "Installation binding is missing or invalid." }], @@ -257,13 +257,12 @@ test("runs the injected installation diagnostic for a migration-required v2 work try { const testResult = await app.inject({ method: "POST", url: "/workspaces/psd-clinical/test", payload: {} }); - expect(testResult.statusCode).toBe(200); - expect(testResult.json()).toMatchObject({ activatable: false, diagnostics: [{ code: "binding_missing" }] }); - expect(diagnose).toHaveBeenCalledWith(workspaceV2, expect.objectContaining({ - dwh: expect.objectContaining({ transport: "rest_api", missing: [] }), - vector: expect.objectContaining({ transport: "rest_api", missing: [] }), - embedding: expect.objectContaining({ transport: "rest_api", missing: [] }), - }), { writeProbe: false }); + expect(testResult.statusCode).toBe(400); + expect(testResult.json()).toEqual({ + code: "workspace_not_activatable", + message: "Workspace cannot be activated on this installation.", + }); + expect(diagnose).not.toHaveBeenCalled(); } finally { process.env = originalEnv; } diff --git a/backend/test/tht-qdrant-readiness.test.ts b/backend/test/tht-qdrant-readiness.test.ts new file mode 100644 index 00000000..6eaa0a83 --- /dev/null +++ b/backend/test/tht-qdrant-readiness.test.ts @@ -0,0 +1,100 @@ +import { expect, test, vi } from "vitest"; +import { ThtRunner } from "../src/tht/tht-runner.js"; +import type { CanonicalWorkspace } from "../src/workspaces/schema.js"; + +const keywordIndexes = [ + "content_hash", "document_id", "kind", "record_key", "record_kind", + "vector_generation", "workspace_id", "workspace_revision", +]; + +const workspace: CanonicalWorkspace = { + workspace: { schema_version: 3, id: "psd", name: "PSD", language: "it" }, + dwh: { + engine: "postgres", database: "warehouse", schema: "public", + supported_transports: ["postgres_direct"], + }, + semantic_index: { + vector_store: { engine: "qdrant", collection: "psd", dimensions: 1024, distance: "cosine" }, + embedding: { + provider: "ollama_internal", model: "qwen3-embedding:0.6b", dimensions: 1024, + }, + }, + llm_policy: { allowed: ["zai/glm-5.2"] }, +}; + +function runner(request: (...args: any[]) => Promise) { + return new ThtRunner({ + thtBin: "tht", + harnessDir: "/harness", + configPath: "config/tht.yaml", + semanticRuntime: { + internalQdrantUrl: "http://qdrant:6333", + internalEmbeddingUrl: "http://embedding:11434", + internalEmbeddingModel: "qwen3-embedding:0.6b", + internalEmbeddingDimensions: 1024, + }, + qdrantRequest: request, + }); +} + +function response(status: number, body: unknown) { + return { + ok: status >= 200 && status < 300, + status, + json: async () => body, + }; +} + +function collection(overrides: Record = {}) { + return { + result: { + config: { params: { vectors: { size: 1024, distance: "Cosine" } } }, + payload_schema: Object.fromEntries(keywordIndexes.map((field) => [field, { data_type: "keyword" }])), + ...overrides, + }, + }; +} + +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 }); + 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) }); +}); + +test("Qdrant readiness classifies a missing collection as semantic incompatibility", async () => { + const request = vi.fn(async () => response(404, { status: "error", detail: "secret" })); + + await expect(runner(request).qdrantEnsure(workspace, 3)).resolves.toEqual({ + ok: false, + code: "semantic_index_incompatible", + }); +}); + +test.each([ + ["dimensions", collection({ + config: { params: { vectors: { size: 768, distance: "Cosine" } } }, + })], + ["distance", collection({ + config: { params: { vectors: { size: 1024, distance: "Dot" } } }, + })], + ["payload indexes", collection({ payload_schema: { workspace_id: { data_type: "keyword" } } })], +])("Qdrant readiness rejects incompatible %s", async (_label, body) => { + const request = vi.fn(async () => response(200, body)); + + await expect(runner(request).qdrantEnsure(workspace, 3)).resolves.toEqual({ + ok: false, + code: "semantic_index_incompatible", + }); +}); + +test("Qdrant readiness sanitizes unreachable internal service failures", async () => { + const request = vi.fn(async () => { throw new Error("connect http://qdrant:6333/private"); }); + + await expect(runner(request).qdrantEnsure(workspace, 3)).resolves.toEqual({ + ok: false, + code: "workspace_not_activatable", + }); +}); diff --git a/backend/test/workspace-registry.test.ts b/backend/test/workspace-registry.test.ts index 5835bbce..0111b9bd 100644 --- a/backend/test/workspace-registry.test.ts +++ b/backend/test/workspace-registry.test.ts @@ -662,6 +662,9 @@ test("lists a schema v2 descriptor as migration_required and refuses to acquire await expect(registry.acquireSessionRevision("psd-clinical")).rejects.toMatchObject({ code: "workspace_invalid", }); + await expect(registry.readPinned("psd-clinical", remote.initialCommit)).rejects.toMatchObject({ + code: "workspace_invalid", + }); }); test("lists operational descriptors retained after their workspace was removed from the active revision", async () => { diff --git a/backend/test/workspace-runtime-handoff.test.ts b/backend/test/workspace-runtime-handoff.test.ts index 63c7eebe..73b1febe 100644 --- a/backend/test/workspace-runtime-handoff.test.ts +++ b/backend/test/workspace-runtime-handoff.test.ts @@ -42,6 +42,18 @@ llm_policy: allowed: [zai/glm-5.2] `; +const migrationRequiredWorkspace = canonicalWorkspace + .replace("schema_version: 3", "schema_version: 2") + .replace( + " engine: qdrant\n collection: psd-clinical", + " engine: pgvector\n database: analytics\n schema: vectors\n collection: documents", + ) + .replace(" dimensions: 1024", " dimensions: 768") + .replace(" distance: cosine", " distance: cosine\n supported_transports: [rest_api]") + .replace(" provider: ollama_internal", " provider: ollama_compatible") + .replace(" model: qwen3-embedding:0.6b", " model: nomic-embed-text") + .replace(" dimensions: 1024", " dimensions: 768"); + afterEach(() => { vi.unstubAllEnvs(); roots.splice(0).forEach((root) => rmSync(root, { recursive: true, force: true })); @@ -51,7 +63,7 @@ async function git(cwd: string, args: string[]): Promise { return (await runFile("git", args, { cwd })).stdout.trim(); } -async function fixture() { +async function fixture(workspaceSource = canonicalWorkspace) { const root = mkdtempSync(join(tmpdir(), "tht-runtime-handoff-")); roots.push(root); const remote = join(root, "remote.git"); @@ -65,7 +77,7 @@ async function fixture() { await git(source, ["config", "user.name", "Runtime Handoff Test"]); await git(source, ["config", "user.email", "runtime-handoff@example.invalid"]); mkdirSync(join(source, "workspaces")); - writeFileSync(join(source, "workspaces", "psd-clinical.yaml"), canonicalWorkspace); + writeFileSync(join(source, "workspaces", "psd-clinical.yaml"), workspaceSource); await git(source, ["add", "workspaces/psd-clinical.yaml"]); await git(source, ["commit", "-m", "Canonical workspace"]); await git(source, ["remote", "add", "origin", remote]); @@ -160,6 +172,15 @@ test("separate runtime leases hand off one stable logical workspace identity", a } }); +test("ThtRunner refuses to render a migration-required registry snapshot", async () => { + const f = await fixture(migrationRequiredWorkspace); + const runner = runnerFor(f); + + expect(() => runner.acquireWorkspaceRuntime(f.revision.snapshotPath)).toThrow( + "Workspace descriptor requires explicit migration to schema version 3", + ); +}); + test("local GET sessions mine uses the real canonical handoff and returns an empty inventory", async () => { const f = await fixture(); const app = buildApp(loadConfig({ diff --git a/harness/tests/test_qdrant_vector_store.py b/harness/tests/test_qdrant_vector_store.py index 5da91c5b..a729f611 100644 --- a/harness/tests/test_qdrant_vector_store.py +++ b/harness/tests/test_qdrant_vector_store.py @@ -212,6 +212,27 @@ def test_upsert_refuses_collection_dimension_or_distance_mismatch_without_recrea assert creates == [] +def test_health_fails_when_the_bound_collection_is_missing(): + fake = FakeQdrantHttp() + + health = _store(fake).health() + + assert health.ok is False + assert health.read_reachable is False + assert health.write_reachable is False + assert "missing" in (health.detail or "").lower() + + +def test_health_fails_when_required_payload_indexes_are_missing_without_creating_them(): + fake = FakeQdrantHttp() + fake.collection = {"vectors": {"size": 1024, "distance": "Cosine"}} + + health = _store(fake).health() + + assert health.ok is False + assert fake.payload_indexes == set() + + @pytest.mark.parametrize( ("record", "semantic_kind"), [ @@ -307,6 +328,90 @@ def test_existing_hashes_health_and_exact_generation_inventory_and_delete(): assert health.dimension_compatible is True +def test_metadata_search_rejects_a_workspace_id_different_from_the_bound_adapter(): + fake = FakeQdrantHttp() + store = _store(fake) + generation = "gen:" + "1" * 32 + store.upsert("evidence", [ + _write_record( + f"demo:{generation}:chunk:1", + "evidence", + metadata={ + "workspace_id": "demo", "vector_generation": generation, + "document_id": "doc:shared", + }, + ), + ]) + foreign = next(iter(fake.points.values())).copy() + foreign["id"] = point_id("other", "evidence", f"other:{generation}:chunk:1") + foreign["payload"] = { + **foreign["payload"], + "workspace_id": "other", + "record_key": f"other:{generation}:chunk:1", + "ref": "ref:foreign", + "title": "foreign", + "content": "foreign", + } + fake.points[foreign["id"]] = foreign + + with pytest.raises(VectorStoreError, match="workspace namespace does not match"): + store.search( + ["evidence"], [0.2] * 1024, limit=5, kinds=["evidence"], + metadata_filter={ + "workspace_id": "other", + "vector_generation": generation, + "document_ids": ["doc:shared"], + }, + ) + + +def test_generation_inventory_rejects_a_workspace_id_different_from_the_bound_adapter(): + fake = FakeQdrantHttp() + store = _store(fake) + + with pytest.raises(VectorStoreError, match="workspace namespace does not match"): + store.list_evidence_generations("evidence", "other") + assert not any(call[1].endswith("/points/scroll") for call in fake.calls) + + +def test_generation_delete_cannot_mutate_foreign_workspace_or_non_evidence_points(): + fake = FakeQdrantHttp() + store = _store(fake) + generation = "gen:" + "1" * 32 + store.upsert("evidence", [ + _write_record( + f"demo:{generation}:chunk:1", + "evidence", + metadata={ + "workspace_id": "demo", "vector_generation": generation, + "document_id": "doc:demo", + }, + ), + ]) + demo = next(iter(fake.points.values())) + foreign = demo.copy() + foreign["id"] = point_id("other", "evidence", f"other:{generation}:chunk:1") + foreign["payload"] = { + **demo["payload"], "workspace_id": "other", + "record_key": f"other:{generation}:chunk:1", + } + fake.points[foreign["id"]] = foreign + memory = demo.copy() + memory["id"] = point_id("other", "memory", "memory:foreign") + memory["payload"] = { + **demo["payload"], "workspace_id": "other", "kind": "memory", + "record_kind": "memory", "record_key": "memory:foreign", + } + fake.points[memory["id"]] = memory + before = set(fake.points) + + with pytest.raises(VectorStoreError, match="workspace namespace does not match"): + store.delete_generation("evidence", generation, "other") + + assert set(fake.points) == before + assert not any(call[1].endswith("/points/delete?wait=true") for call in fake.calls) + + def test_delete_kinds_is_workspace_scoped_and_preserves_other_semantic_kinds(): fake = FakeQdrantHttp() store = _store(fake) diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index 7d880d82..d3073f15 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -94,14 +94,11 @@ class QdrantVectorStore: expected_dimension=self._expected_dimension, ) - dimensions = () - compatible = None - if info is not None: - dimension = info["config"]["params"]["vectors"]["size"] - dimensions = (dimension,) - compatible = ( - None if self._expected_dimension is None else dimensions == (self._expected_dimension,) - ) + dimension = info["config"]["params"]["vectors"]["size"] + dimensions = (dimension,) + compatible = ( + None if self._expected_dimension is None else dimensions == (self._expected_dimension,) + ) return VectorHealth( ok=compatible is not False, read_configured=True, @@ -142,12 +139,11 @@ class QdrantVectorStore: or not isinstance(workspace_id, str) ): raise VectorStoreError("Invalid vector metadata filter") - filter_must = [ - {"key": "workspace_id", "match": {"value": workspace_id}}, - {"key": "record_kind", "match": {"any": allowed_record_kinds}}, + self._require_bound_workspace(workspace_id) + filter_must.extend([ {"key": "vector_generation", "match": {"value": generation}}, {"key": "document_id", "match": {"any": document_ids}}, - ] + ]) response = self._call( "POST", f"/collections/{self._collection}/points/query", @@ -232,27 +228,19 @@ class QdrantVectorStore: raise VectorStoreError("Only exact Evidence generations may be deleted") if _WORKSPACE.fullmatch(workspace_id) is None: raise VectorStoreError("Invalid Evidence workspace namespace") + self._require_bound_workspace(workspace_id) + must = [ + *self._workspace_filter(), + {"key": "record_kind", "match": {"any": ["evidence"]}}, + {"key": "vector_generation", "match": {"value": generation}}, + ] before = len( - self._scroll( - [ - {"key": "workspace_id", "match": {"value": workspace_id}}, - {"key": "record_kind", "match": {"any": ["evidence"]}}, - {"key": "vector_generation", "match": {"value": generation}}, - ] - ) + self._scroll(must) ) self._call( "POST", f"/collections/{self._collection}/points/delete?wait=true", - { - "filter": { - "must": [ - {"key": "workspace_id", "match": {"value": workspace_id}}, - {"key": "record_kind", "match": {"any": ["evidence"]}}, - {"key": "vector_generation", "match": {"value": generation}}, - ] - } - }, + {"filter": {"must": must}}, ) return before @@ -261,9 +249,10 @@ class QdrantVectorStore: raise VectorStoreError("Only exact Evidence generations may be listed") if _WORKSPACE.fullmatch(workspace_id) is None: raise VectorStoreError("Invalid Evidence workspace namespace") + self._require_bound_workspace(workspace_id) points = self._scroll( [ - {"key": "workspace_id", "match": {"value": workspace_id}}, + *self._workspace_filter(), {"key": "record_kind", "match": {"any": ["evidence"]}}, ] ) @@ -279,6 +268,10 @@ class QdrantVectorStore: def _workspace_filter(self) -> list[dict]: return [{"key": "workspace_id", "match": {"value": self._workspace_id}}] + def _require_bound_workspace(self, workspace_id: str) -> None: + if workspace_id != self._workspace_id: + raise VectorStoreError("Evidence workspace namespace does not match bound workspace") + def _allowed_record_kinds( self, collections: list[str], kinds: list[str] | None ) -> list[str]: @@ -303,7 +296,7 @@ class QdrantVectorStore: response = self._call("GET", f"/collections/{self._collection}", None, allow_missing=True) if response is None: if not strict: - return None + raise VectorStoreError("Qdrant collection is missing") self._call( "PUT", f"/collections/{self._collection}", @@ -329,6 +322,8 @@ class QdrantVectorStore: raise VectorStoreError("Qdrant collection configuration mismatch") for field_name in _KEYWORD_INDEXES: if field_name not in result.get("payload_schema", {}): + if not strict: + raise VectorStoreError("Qdrant collection payload indexes mismatch") self._call( "PUT", f"/collections/{self._collection}/index", diff --git a/scripts/lib/vector-operation-lock.sh b/scripts/lib/vector-operation-lock.sh new file mode 100644 index 00000000..a148b237 --- /dev/null +++ b/scripts/lib/vector-operation-lock.sh @@ -0,0 +1,48 @@ +#!/bin/sh + +# Shared, daemon-scoped lock for operations that stop or replace one Compose Qdrant volume. +# The stopped container name is the atomic primitive; ownership labels prevent a cleanup trap +# from deleting a lock that it did not create. +vector_operation_lock_init() { + operation_lock_name="${expected_volume_name}-operation-lock" + operation_lock_owner="${project_name}-$$-$(date -u +%Y%m%dT%H%M%S)" + operation_lock_acquired=0 +} + +acquire_vector_operation_lock() { + if docker create \ + --name "$operation_lock_name" \ + --label "com.thothii.qdrant-operation-owner=$operation_lock_owner" \ + --label "com.thothii.qdrant-operation-project=$project_name" \ + --label "com.thothii.qdrant-operation-volume=$expected_volume_name" \ + "$helper_image" /bin/true >/dev/null 2>&1; then + operation_lock_acquired=1 + return 0 + fi + + if docker inspect "$operation_lock_name" >/dev/null 2>&1; then + echo "Qdrant operation already in progress for $expected_volume_name" >&2 + else + echo "Unable to acquire Qdrant operation lock for $expected_volume_name" >&2 + fi + return 2 +} + +release_vector_operation_lock() { + [ "${operation_lock_acquired:-0}" -eq 1 ] || return 0 + metadata=$(docker inspect --format \ + '{{ index .Config.Labels "com.thothii.qdrant-operation-owner" }} {{ index .Config.Labels "com.thothii.qdrant-operation-volume" }}' \ + "$operation_lock_name" 2>/dev/null || true) + set -- $metadata + if [ "${1-}" != "$operation_lock_owner" ] || [ "${2-}" != "$expected_volume_name" ]; then + echo "Qdrant operation lock ownership changed; refusing to remove it" >&2 + operation_lock_acquired=0 + return 1 + fi + if ! docker rm -f "$operation_lock_name" >/dev/null; then + echo "Unable to release Qdrant operation lock for $expected_volume_name" >&2 + operation_lock_acquired=0 + return 1 + fi + operation_lock_acquired=0 +} diff --git a/scripts/test-vector-backup-restore-safety.sh b/scripts/test-vector-backup-restore-safety.sh index 7c0d947b..d173d224 100755 --- a/scripts/test-vector-backup-restore-safety.sh +++ b/scripts/test-vector-backup-restore-safety.sh @@ -3,7 +3,17 @@ set -eu cd "$(dirname "$0")/.." tmp=$(mktemp -d) -trap 'rm -rf "$tmp"' EXIT HUP INT TERM +holder_pid= +holder_release="$tmp/holder-release" +cleanup() { + touch "$holder_release" + if [ -n "${holder_pid:-}" ]; then + kill "$holder_pid" 2>/dev/null || true + wait "$holder_pid" 2>/dev/null || true + fi + rm -rf "$tmp" +} +trap cleanup EXIT HUP INT TERM fakebin="$tmp/bin" mkdir "$fakebin" @@ -28,6 +38,15 @@ run_backup() { volume_name=$1 backup_dir=$2 output_name=$3 + if [ -n "${RUN_BLOCK_READY:-}" ]; then + : >"$RUN_BLOCK_READY" + attempts=0 + while [ ! -e "${RUN_BLOCK_RELEASE:?}" ]; do + attempts=$((attempts + 1)) + [ "$attempts" -lt 400 ] || exit 24 + sleep 0.05 + done + fi staging=$(mktemp -d "${TMPDIR:-/tmp}/fake-qdrant-backup.XXXXXX") mkdir -p "$staging/payload" cat >"$staging/manifest.env" </dev/null || exit 1 + printf '%s\n' "$lock_owner" >"$state/owner" + printf '%s\n' "$lock_volume" >"$state/volume" + printf '%s\n' "$lock_name" + exit 0 +fi + +if [ "$1" = inspect ]; then + for lock_name in "$@"; do :; done + state="$lock_state_root/$lock_name" + [ -d "$state" ] || exit 1 + printf '%s %s\n' "$(cat "$state/owner")" "$(cat "$state/volume")" + exit 0 +fi + +if [ "$1" = rm ]; then + for lock_name in "$@"; do :; done + state="$lock_state_root/$lock_name" + [ -d "$state" ] || exit 1 + rm -rf "$state" + exit 0 +fi + run_restore() { volume_name=$1 backup_dir=$2 @@ -219,6 +285,7 @@ PY backup_output="$tmp/qdrant-backup.tar" docker_log="$tmp/docker-backup.log" +lock_name="${project}_qdrant-data-operation-lock" PATH="$fakebin:$PATH" DOCKER_LOG="$docker_log" PROJECT_NAME="$project" VOLUME_ROOT="$volume_root" \ HELPER_IMAGE="qdrant/qdrant:v1.18.2@sha256:75eab8c4ba42096724fdcfde8b4de0b5713d529dde32f285a1f86fdcb2c9e50c" \ ./scripts/vector-backup.sh --project-name "$project" --output "$backup_output" >/dev/null @@ -231,6 +298,15 @@ tar -xOf "$backup_output" manifest.env | grep -qx "volume_name=$volume_name" grep -q "run --rm --mount type=volume,src=$volume_name,dst=/qdrant-data,readonly --mount type=bind,src=$tmp,dst=/backup" "$docker_log" grep -q "compose --project-name $project stop qdrant" "$docker_log" grep -q "compose --project-name $project start qdrant" "$docker_log" +grep -q "create --name $lock_name" "$docker_log" +grep -q "rm -f $lock_name" "$docker_log" +test ! -e "$tmp/locks/$lock_name" +lock_line=$(grep -n "create --name $lock_name" "$docker_log" | sed -n '1s/:.*//p') +volume_line=$(grep -n '^volume ls ' "$docker_log" | sed -n '1s/:.*//p') +[ "$lock_line" -lt "$volume_line" ] || { + echo "backup resolved volume state before acquiring the operation lock" >&2 + exit 1 +} if grep -q "volume inspect --format {{ .Mountpoint }}" "$docker_log"; then echo "backup consulted Docker mountpoints" >&2 exit 1 @@ -240,6 +316,61 @@ if grep -q "prune" "$docker_log"; then exit 1 fi +mkdir -p "$tmp/locks/$lock_name" +printf '%s\n' foreign-owner >"$tmp/locks/$lock_name/owner" +printf '%s\n' "$volume_name" >"$tmp/locks/$lock_name/volume" +foreign_lock_log="$tmp/foreign-lock.log" +if PATH="$fakebin:$PATH" DOCKER_LOG="$foreign_lock_log" PROJECT_NAME="$project" VOLUME_ROOT="$volume_root" \ +HELPER_IMAGE="qdrant/qdrant:v1.18.2@sha256:75eab8c4ba42096724fdcfde8b4de0b5713d529dde32f285a1f86fdcb2c9e50c" \ + ./scripts/vector-backup.sh --project-name "$project" --output "$tmp/foreign-lock.tar" \ + >"$tmp/foreign-lock.out" 2>"$tmp/foreign-lock.err"; then + echo "backup ignored a foreign operation lock" >&2 + exit 1 +fi +grep -q 'operation already in progress' "$tmp/foreign-lock.err" +test "$(cat "$tmp/locks/$lock_name/owner")" = foreign-owner +rm -rf "$tmp/locks/$lock_name" + +holder_ready="$tmp/holder-ready" +holder_output="$tmp/holder.tar" +holder_log="$tmp/holder.log" +PATH="$fakebin:$PATH" DOCKER_LOG="$holder_log" PROJECT_NAME="$project" VOLUME_ROOT="$volume_root" \ +RUN_BLOCK_READY="$holder_ready" RUN_BLOCK_RELEASE="$holder_release" \ +HELPER_IMAGE="qdrant/qdrant:v1.18.2@sha256:75eab8c4ba42096724fdcfde8b4de0b5713d529dde32f285a1f86fdcb2c9e50c" \ + ./scripts/vector-backup.sh --project-name "$project" --output "$holder_output" >/dev/null 2>"$tmp/holder.err" & +holder_pid=$! +attempts=0 +while [ ! -e "$holder_ready" ]; do + attempts=$((attempts + 1)) + [ "$attempts" -lt 400 ] || { echo "lock holder did not reach copy phase" >&2; exit 1; } + sleep 0.05 +done + +for contender in backup restore; do + contender_log="$tmp/contender-$contender.log" + if [ "$contender" = backup ]; then + set -- ./scripts/vector-backup.sh --project-name "$project" --output "$tmp/contender.tar" + else + set -- ./scripts/vector-restore.sh --project-name "$project" --input "$backup_output" --confirm-project "$project" + fi + if PATH="$fakebin:$PATH" DOCKER_LOG="$contender_log" PROJECT_NAME="$project" VOLUME_ROOT="$volume_root" \ + HELPER_IMAGE="qdrant/qdrant:v1.18.2@sha256:75eab8c4ba42096724fdcfde8b4de0b5713d529dde32f285a1f86fdcb2c9e50c" \ + "$@" >"$tmp/contender-$contender.out" 2>"$tmp/contender-$contender.err"; then + echo "$contender interleaved with an active qdrant operation" >&2 + exit 1 + fi + grep -q 'operation already in progress' "$tmp/contender-$contender.err" + if grep -Eq '^volume (ls|inspect)|^compose .* (stop|start) qdrant|^run ' "$contender_log"; then + echo "$contender touched qdrant state after lock contention" >&2 + exit 1 + fi +done +touch "$holder_release" +wait "$holder_pid" +holder_pid= +test -s "$holder_output" +test ! -e "$tmp/locks/$lock_name" + existing="$tmp/existing.tar" printf '%s' sentinel >"$existing" if PATH="$fakebin:$PATH" DOCKER_LOG="$tmp/existing.log" PROJECT_NAME="$project" VOLUME_ROOT="$volume_root" \ @@ -392,5 +523,6 @@ fi test "$(cat "$volume_dir/collections/demo/state.json")" = rollback-source grep -q "compose --project-name $project stop qdrant" "$rollback_log" grep -q "compose --project-name $project start qdrant" "$rollback_log" +test ! -e "$tmp/locks/$lock_name" echo "qdrant backup/restore archive validation, scoped helper execution, and rollback safety passed." diff --git a/scripts/vector-backup.sh b/scripts/vector-backup.sh index 47c3f52c..3efcd18d 100755 --- a/scripts/vector-backup.sh +++ b/scripts/vector-backup.sh @@ -28,6 +28,9 @@ output_name=$(basename "$output") [ -d "$output_dir" ] || { echo "backup destination directory does not exist" >&2; exit 2; } expected_volume_name="${project_name}_${volume_role}" +script_dir=$(CDPATH= cd "$(dirname "$0")" && pwd) +. "$script_dir/lib/vector-operation-lock.sh" +vector_operation_lock_init resolve_volume() { names=$(docker volume ls \ @@ -55,24 +58,32 @@ validate_volume_metadata() { [ "${3-}" = "$volume_role" ] || { echo "volume metadata role label mismatch" >&2; exit 2; } } -volume_name=$(resolve_volume) -validate_volume_metadata "$volume_name" - -running_container=$(docker compose --project-name "$project_name" ps --status running -q qdrant) restart_qdrant=0 temporary_output= cleanup() { status=$? + trap - EXIT HUP INT TERM if [ -n "${temporary_output:-}" ] && [ -e "${temporary_output:-}" ]; then rm -f "$temporary_output" fi if [ "$restart_qdrant" -eq 1 ]; then - docker compose --project-name "$project_name" start qdrant >/dev/null + if ! docker compose --project-name "$project_name" start qdrant >/dev/null; then + echo "Unable to restart Qdrant after backup" >&2 + [ "$status" -ne 0 ] || status=1 + fi + fi + if ! release_vector_operation_lock; then + [ "$status" -ne 0 ] || status=1 fi exit "$status" } trap cleanup EXIT HUP INT TERM +acquire_vector_operation_lock +volume_name=$(resolve_volume) +validate_volume_metadata "$volume_name" +running_container=$(docker compose --project-name "$project_name" ps --status running -q qdrant) + if [ -n "$running_container" ]; then docker compose --project-name "$project_name" stop qdrant >/dev/null restart_qdrant=1 diff --git a/scripts/vector-restore.sh b/scripts/vector-restore.sh index 6aae4f6a..5ee3a27b 100755 --- a/scripts/vector-restore.sh +++ b/scripts/vector-restore.sh @@ -30,15 +30,26 @@ done [ -r "$input" ] || { echo "backup input is not readable" >&2; exit 2; } expected_volume_name="${project_name}_${volume_role}" +script_dir=$(CDPATH= cd "$(dirname "$0")" && pwd) +. "$script_dir/lib/vector-operation-lock.sh" +vector_operation_lock_init private_archive_dir= private_archive_path= +restart_qdrant=0 cleanup() { status=$? + trap - EXIT HUP INT TERM if [ -n "${private_archive_dir:-}" ] && [ -d "${private_archive_dir:-}" ]; then rm -rf "$private_archive_dir" fi if [ "${restart_qdrant:-0}" -eq 1 ]; then - docker compose --project-name "$project_name" start qdrant >/dev/null + if ! docker compose --project-name "$project_name" start qdrant >/dev/null; then + echo "Unable to restart Qdrant after restore" >&2 + [ "$status" -ne 0 ] || status=1 + fi + fi + if ! release_vector_operation_lock; then + [ "$status" -ne 0 ] || status=1 fi exit "$status" } @@ -219,11 +230,11 @@ validate_archive_paths validate_archive_types validate_archive_manifest +acquire_vector_operation_lock volume_name=$(resolve_volume) validate_volume_metadata "$volume_name" running_container=$(docker compose --project-name "$project_name" ps --status running -q qdrant) -restart_qdrant=0 if [ -n "$running_container" ]; then docker compose --project-name "$project_name" stop qdrant >/dev/null