feat: reserve one qdrant collection per workspace

This commit is contained in:
2026-08-08 17:03:51 +02:00
parent ba1d7b0e78
commit 76bc94d5da
6 changed files with 288 additions and 107 deletions
+10 -13
View File
@@ -11,7 +11,7 @@ import { renderWorkspaceDocs, serializeWorkspaceYaml, type CanonicalWorkspace }
const workspace: CanonicalWorkspace = {
workspace: {
schema_version: 2,
schema_version: 3,
id: "psd-clinical",
name: "Policlinico San Donato",
description: "Clinical analytics workspace",
@@ -25,18 +25,15 @@ const workspace: CanonicalWorkspace = {
},
semantic_index: {
vector_store: {
engine: "pgvector",
database: "warehouse",
schema: "vectors",
collection: "clinical_documents",
dimensions: 768,
engine: "qdrant",
collection: "psd-clinical",
dimensions: 1024,
distance: "cosine",
supported_transports: ["pgvector_direct"],
},
embedding: {
provider: "ollama_compatible",
model: "nomic-embed-text-v2-moe",
dimensions: 768,
provider: "ollama_internal",
model: "qwen3-embedding:0.6b",
dimensions: 1024,
},
},
llm_policy: { allowed: ["zai/glm-5.2"] },
@@ -184,9 +181,9 @@ test("validates a canonical workspace and runs the injected installation diagnos
expect(validate.statusCode).toBe(200);
expect(validate.json()).toMatchObject({ workspace });
expect(testResult.statusCode).toBe(200);
expect(testResult.json()).toMatchObject({ activatable: false, diagnostics: [{ code: "binding_missing" }] });
expect(diagnose).toHaveBeenCalledWith(workspace, expect.any(Object), { writeProbe: false });
expect(testResult.statusCode).toBe(400);
expect(testResult.json()).toMatchObject({ code: "workspace_invalid" });
expect(diagnose).not.toHaveBeenCalled();
});
test("returns a 409 field conflict instead of overwriting a changed workspace", async () => {
+108 -36
View File
@@ -13,7 +13,7 @@ import { parseWorkspaceYaml, type CanonicalWorkspace } from "../src/workspaces/s
import type { WorkspaceRegistryConfig } from "../src/workspaces/types.js";
const validYaml = `workspace:
schema_version: 2
schema_version: 3
id: psd-clinical
name: Policlinico San Donato
language: it
@@ -24,17 +24,14 @@ dwh:
supported_transports: [postgres_direct]
semantic_index:
vector_store:
engine: pgvector
database: postgres
schema: vectors
collection: clinical_documents
dimensions: 768
engine: qdrant
collection: psd-clinical
dimensions: 1024
distance: cosine
supported_transports: [pgvector_direct]
embedding:
provider: ollama_compatible
model: nomic-embed-text-v2-moe
dimensions: 768
provider: ollama_internal
model: qwen3-embedding:0.6b
dimensions: 1024
llm_policy:
allowed: [zai/glm-5.2]
`;
@@ -130,6 +127,18 @@ function withReversibleVectorProbe(source: string): string {
`);
}
function legacyV1Yaml(source = validYaml): string {
return source
.replace(" engine: qdrant\n", " engine: pgvector\n database: postgres\n schema: vectors\n")
.replace(" collection: psd-clinical\n", " collection: psd_clinical\n")
.replace(" dimensions: 1024", " dimensions: 768")
.replace(" provider: ollama_internal", " provider: ollama_compatible")
.replace(" model: qwen3-embedding:0.6b", " model: nomic-embed-text-v2-moe")
.replace(" dimensions: 1024", " dimensions: 768")
.replace("distance: cosine\n", "distance: cosine\n supported_transports: [pgvector_direct]\n")
.replace("schema_version: 3", "schema_version: 1");
}
const runFile = promisify(execFile);
const temporaryRoots: string[] = [];
@@ -168,6 +177,30 @@ async function fixture(workspaceSource = validYaml): Promise<{
return { root, remote, source, initialCommit: stdout.trim() };
}
async function multiWorkspaceFixture(workspaces: Record<string, string>): Promise<{
root: string; remote: string; source: string; initialCommit: string;
}> {
const root = mkdtempSync(join(tmpdir(), "thoth-workspace-registry-"));
temporaryRoots.push(root);
const remote = join(root, "remote.git");
const source = join(root, "source");
await git(root, ["init", "--bare", "--initial-branch=main", remote]);
mkdirSync(source);
await git(source, ["init", "--initial-branch=main"]);
await git(source, ["config", "user.name", "Workspace Registry Test"]);
await git(source, ["config", "user.email", "workspace-registry@example.invalid"]);
mkdirSync(join(source, "workspaces"));
for (const [id, workspaceSource] of Object.entries(workspaces)) {
writeFileSync(join(source, "workspaces", `${id}.yaml`), workspaceSource);
}
await git(source, ["add", "workspaces"]);
await git(source, ["commit", "-m", "Initial workspaces"]);
await git(source, ["remote", "add", "origin", remote]);
await git(source, ["push", "origin", "main"]);
const { stdout } = await runFile("git", ["rev-parse", "HEAD"], { cwd: source });
return { root, remote, source, initialCommit: stdout.trim() };
}
function config(
root: string,
remoteUrl: string,
@@ -195,6 +228,10 @@ function workspaceWith(
return {
...workspace,
workspace: { ...workspace.workspace, id, name: id, ...changes },
semantic_index: {
...workspace.semantic_index,
vector_store: { ...workspace.semantic_index.vector_store, collection: id },
},
};
}
@@ -325,10 +362,10 @@ test("reports stale publish conflicts with expected and actual revisions", async
await registry.bootstrap();
const initial = await registry.read("psd-clinical");
writeFileSync(join(remote.source, "workspaces", "psd-clinical.yaml"), validYaml.replace(
"model: nomic-embed-text-v2-moe", "model: mxbai-embed-large",
"schema: datawarehouse", "schema: analytics",
));
await git(remote.source, ["add", "workspaces/psd-clinical.yaml"]);
await git(remote.source, ["commit", "-m", "Change embedding model"]);
await git(remote.source, ["commit", "-m", "Change dwh schema"]);
await git(remote.source, ["push", "origin", "main"]);
const actualCommit = await gitOutput(remote.source, ["rev-parse", "HEAD"]);
const actualBlob = await gitOutput(remote.source, ["rev-parse", "HEAD:workspaces/psd-clinical.yaml"]);
@@ -340,18 +377,16 @@ test("reports stale publish conflicts with expected and actual revisions", async
baseBlob: initial.revision.blob,
})).rejects.toMatchObject({
code: "workspace_conflict",
fields: ["semantic_index.embedding.model"],
fields: ["dwh.schema"],
expected: { commit: initial.revision.commit, blob: initial.revision.blob },
actual: { commit: actualCommit, blob: actualBlob },
});
});
test.each([
["adds", withEmbeddingDiagnostic(withDwhRestTransport(validYaml)), withDwhRestAndEmbeddingDiagnostics(validYaml), "diagnostics.dwh_rest"],
["removes", withDwhRestAndEmbeddingDiagnostics(validYaml), withEmbeddingDiagnostic(withDwhRestTransport(validYaml)), "diagnostics.dwh_rest"],
["adds", withVectorMetadataDiagnostic(validYaml), withReversibleVectorProbe(validYaml), "diagnostics.vector_rest.reversible_probe"],
["removes", withReversibleVectorProbe(validYaml), withVectorMetadataDiagnostic(validYaml), "diagnostics.vector_rest.reversible_probe"],
])("reports an optional diagnostics branch when the registry %s it", async (_operation, baseSource, remoteSource, field) => {
["adds", validYaml, withDwhRestDiagnostic(validYaml)],
["removes", withDwhRestDiagnostic(validYaml), validYaml],
])("reports an optional diagnostics branch when the registry %s it", async (_operation, baseSource, remoteSource) => {
const remote = await fixture(baseSource);
const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote));
await registry.bootstrap();
@@ -368,7 +403,7 @@ test.each([
baseBlob: initial.revision.blob,
})).rejects.toMatchObject({
code: "workspace_conflict",
fields: [field],
fields: ["dwh.supported_transports", "diagnostics"],
});
});
@@ -416,10 +451,7 @@ test("resets an ahead checkout after a rejected push and retries publication", a
});
test("lists a v1 descriptor in migration-required state without rendering operational artifacts", async () => {
const legacyYaml = validYaml.replace(
" database: postgres\n schema: vectors\n",
"",
).replace("schema_version: 2", "schema_version: 1");
const legacyYaml = legacyV1Yaml();
const remote = await fixture(legacyYaml);
const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote));
@@ -435,10 +467,7 @@ test("lists a v1 descriptor in migration-required state without rendering operat
});
test("migrates a validated pre-state manifest and keeps its v1 workspace migration-gated", async () => {
const legacyYaml = validYaml.replace(
" database: postgres\n schema: vectors\n",
"",
).replace("schema_version: 2", "schema_version: 1");
const legacyYaml = legacyV1Yaml();
const remote = await fixture(legacyYaml);
const root = join(remote.root, "registry");
const firstRegistry = new WorkspaceRegistry(config(root, remote.remote));
@@ -462,10 +491,7 @@ test("migrates a validated pre-state manifest and keeps its v1 workspace migrati
});
test("finishes a pre-state active manifest migration after its snapshot was atomically updated", async () => {
const legacyYaml = validYaml.replace(
" database: postgres\n schema: vectors\n",
"",
).replace("schema_version: 2", "schema_version: 1");
const legacyYaml = legacyV1Yaml();
const remote = await fixture(legacyYaml);
const root = join(remote.root, "registry");
const registry = new WorkspaceRegistry(config(root, remote.remote));
@@ -485,10 +511,7 @@ test("finishes a pre-state active manifest migration after its snapshot was atom
});
test("rejects a corrupt pre-state manifest rather than accepting it during migration", async () => {
const legacyYaml = validYaml.replace(
" database: postgres\n schema: vectors\n",
"",
).replace("schema_version: 2", "schema_version: 1");
const legacyYaml = legacyV1Yaml();
const remote = await fixture(legacyYaml);
const root = join(remote.root, "registry");
const registry = new WorkspaceRegistry(config(root, remote.remote));
@@ -516,6 +539,45 @@ test("keeps the last valid snapshot when a pulled commit has invalid YAML", asyn
});
});
test("rejects duplicate schema v3 collection ownership and keeps the previous active snapshot", async () => {
const v3Yaml = validYaml;
const remote = await multiWorkspaceFixture({
"psd-clinical": v3Yaml,
"research-clinical": v3Yaml
.replace("id: psd-clinical", "id: research-clinical")
.replace("name: Policlinico San Donato", "name: Research Clinical")
.replace("collection: psd-clinical", "collection: research-clinical"),
});
const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote));
await registry.bootstrap();
writeFileSync(
join(remote.source, "workspaces", "research-clinical.yaml"),
v3Yaml
.replace("id: psd-clinical", "id: research-clinical")
.replace("name: Policlinico San Donato", "name: Research Clinical")
.replace("collection: psd-clinical", "collection: shared"),
);
writeFileSync(
join(remote.source, "workspaces", "psd-clinical.yaml"),
v3Yaml.replace("collection: psd-clinical", "collection: shared"),
);
await git(remote.source, ["add", "workspaces"]);
await git(remote.source, ["commit", "-m", "Duplicate collection ownership"]);
await git(remote.source, ["push", "origin", "main"]);
await expect(registry.pull()).rejects.toMatchObject({
code: "workspace_invalid",
message: "Workspace repository content is invalid",
});
await expect(registry.read("psd-clinical")).resolves.toMatchObject({
revision: { commit: remote.initialCommit },
});
await expect(registry.read("research-clinical")).resolves.toMatchObject({
revision: { commit: remote.initialCommit },
});
});
test("retains a historical snapshot while a resumable manifest still references its revision", async () => {
const remote = await fixture();
const root = join(remote.root, "registry");
@@ -567,6 +629,16 @@ test("a session revision lease survives stale retention scans until its manifest
expect(existsSync(registry.snapshotPath(remote.initialCommit, "psd-clinical"))).toBe(false);
});
test("does not acquire a session revision lease for a migration_required workspace", async () => {
const remote = await fixture(legacyV1Yaml());
const registry = new WorkspaceRegistry(config(join(remote.root, "registry"), remote.remote));
await registry.bootstrap();
await expect(registry.acquireSessionRevision("psd-clinical")).rejects.toMatchObject({
code: "workspace_invalid",
});
});
test("lists operational descriptors retained after their workspace was removed from the active revision", async () => {
const remote = await fixture();
const root = join(remote.root, "registry");
@@ -575,7 +647,7 @@ test("lists operational descriptors retained after their workspace was removed f
writeFileSync(join(remote.source, "workspaces", "archive-only.yaml"), validYaml.replace(
"id: psd-clinical", "id: archive-only",
));
).replace("collection: psd-clinical", "collection: archive-only"));
await git(remote.source, ["add", "workspaces/archive-only.yaml"]);
await git(remote.source, ["commit", "-m", "Add retained workspace"]);
await git(remote.source, ["push", "origin", "main"]);
+19 -12
View File
@@ -22,30 +22,37 @@ function readFixture(name: string): string {
}
test("migrates the current local PSD descriptor without copying secret values", () => {
const result = migrateLegacyWorkspace(readFixture("local.yaml"), { id: "local" });
const result = migrateLegacyWorkspace(readFixture("local.yaml"), { id: "local", collection: "local" });
expect(result.workspace.workspace).toMatchObject({ id: "local", schema_version: 1, language: "it" });
expect(result.state).toBe("migration_required");
expect(result.workspace.workspace).toMatchObject({ id: "local", schema_version: 3, language: "it" });
expect(JSON.stringify(result)).not.toMatch(/password:|api_key:|\$\{THT_/i);
});
test("keeps an incomplete legacy vector identity readable and explicitly migration-required", () => {
const result = migrateLegacyWorkspace(readFixture("tht.example.yaml"), { id: "example" });
test("migrates a legacy descriptor only with an explicit target collection into schema v3", () => {
const result = migrateLegacyWorkspace(readFixture("tht.example.yaml"), { id: "example", collection: "shared" });
const isOperationalWorkspace = (workspaceSchema as { isOperationalWorkspace?: unknown }).isOperationalWorkspace;
expect(result.state).toBe("migration_required");
expect(result.workspace.workspace.schema_version).toBe(1);
expect(parseWorkspaceYaml(result.source).workspace.schema_version).toBe(1);
expect(result.workspace.workspace.schema_version).toBe(3);
expect(parseWorkspaceYaml(result.source).workspace.schema_version).toBe(3);
expect(isOperationalWorkspace).toBeTypeOf("function");
expect((isOperationalWorkspace as (workspace: ReturnType<typeof parseWorkspaceYaml>) => boolean)(
parseWorkspaceYaml(result.source),
)).toBe(false);
)).toBe(true);
expect(parseWorkspaceYaml(result.source)).toMatchObject({
semantic_index: { vector_store: { engine: "qdrant", collection: "shared" } },
});
});
test("requires an explicit target collection for legacy migration", () => {
expect(() => migrateLegacyWorkspace(readFixture("local.yaml"), { id: "local" } as never)).toThrow(
/collection/i,
);
});
test("writes versioned repository artifacts atomically without replacing a prior migration", async () => {
const root = await mkdtemp(join(tmpdir(), "thoth-workspace-migrate-"));
temporaryRoots.push(root);
const migration = migrateLegacyWorkspace(readFixture("local.yaml"), { id: "local" });
const migration = migrateLegacyWorkspace(readFixture("local.yaml"), { id: "local", collection: "local" });
const destination = await writeMigratedWorkspace(migration, root);
@@ -61,11 +68,11 @@ test("CLI accepts an explicit valid ID when a legacy filename contains dots", as
const input = join(root, "psd.clinical.yaml");
writeFileSync(input, readFixture("local.yaml"));
await main(["--input", input, "--output", root, "--id", "psd-clinical"]);
await main(["--input", input, "--output", root, "--id", "psd-clinical", "--collection", "psd-clinical"]);
const destination = join(root, "workspaces", "psd-clinical.yaml");
expect(parseWorkspaceYaml(readFileSync(destination, "utf8"))).toMatchObject({
workspace: { id: "psd-clinical", schema_version: 1 },
workspace: { id: "psd-clinical", schema_version: 3 },
});
});