fix: close internal semantic review gaps

This commit is contained in:
2026-08-08 23:27:58 +02:00
parent 43d8063922
commit 2c4534d968
18 changed files with 806 additions and 77 deletions
+18 -2
View File
@@ -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 <name>` 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
+34 -8
View File
@@ -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,
+17 -3
View File
@@ -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);
}
+29 -5
View File
@@ -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<string, Promise<OllamaEnsureResult>>();
private inFlight = new Map<string, Promise<ReadinessResult>>();
private ready = new Map<string, ReadyEntry>();
constructor(
@@ -21,7 +28,11 @@ export class ReadinessManager {
private now: () => number = Date.now,
) {}
ensure(workspace = "", principal?: PrincipalContext): Promise<OllamaEnsureResult> {
ensure(
workspace = "",
principal?: PrincipalContext,
descriptor?: WorkspaceDescriptor,
): Promise<ReadinessResult> {
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<ReadinessResult> => {
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 });
+73 -2
View File
@@ -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<QdrantEnsureResult> {
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();
}
}
}
+8 -2
View File
@@ -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);
}
+43
View File
@@ -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<T>() {
let resolve!: (value: T) => void;
const promise = new Promise<T>((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",
});
});
+111 -7
View File
@@ -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<typeof buildRealApp>[0], deps: Record<string, unknown> = {}) {
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/);
});
+7 -8
View File
@@ -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;
}
+100
View File
@@ -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<any>) {
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<string, unknown> = {}) {
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",
});
});
+3
View File
@@ -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 () => {
+23 -2
View File
@@ -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<string> {
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({
+105
View File
@@ -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)
+25 -30
View File
@@ -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",
+48
View File
@@ -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
}
+133 -1
View File
@@ -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" <<EOF
@@ -36,12 +55,59 @@ project_name=${PROJECT_NAME:?}
volume_name=$volume_name
volume_role=qdrant-data
helper_image=${HELPER_IMAGE:?}
created_utc=2026-08-08T17:35:36Z
EOF
cp -R "$(volume_dir_for "$volume_name")"/. "$staging/payload"/
tar -C "$staging" -cf "$backup_dir/$output_name" manifest.env payload
rm -rf "$staging"
}
lock_state_root=${LOCK_ROOT:-$(dirname "${VOLUME_ROOT:?}")/locks}
mkdir -p "$lock_state_root"
if [ "$1" = create ]; then
shift
lock_name=
lock_owner=
lock_volume=
while [ "$#" -gt 0 ]; do
case "$1" in
--name) lock_name=$2; shift 2 ;;
--label)
case "$2" in
com.thothii.qdrant-operation-owner=*) lock_owner=${2#*=} ;;
com.thothii.qdrant-operation-volume=*) lock_volume=${2#*=} ;;
esac
shift 2
;;
*) shift ;;
esac
done
[ -n "$lock_name" ] && [ -n "$lock_owner" ] && [ -n "$lock_volume" ] || exit 25
state="$lock_state_root/$lock_name"
mkdir "$state" 2>/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."
+16 -5
View File
@@ -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
+13 -2
View File
@@ -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