From e23e52696656716d06f8c3e5b3aa9f552de114dd Mon Sep 17 00:00:00 2001 From: mptyl Date: Wed, 12 Aug 2026 16:01:25 +0200 Subject: [PATCH] feat: P3 effective configuration, memory root, and revision-scoped records --- backend/src/tht/tht-runner.ts | 1 + backend/src/workspace-maintenance.ts | 1 + backend/src/workspaces/effective-config.ts | 219 ++++++++++++++++++ .../src/workspaces/preprocessing-service.ts | 6 + .../src/workspaces/runtime-config-lease.ts | 49 +++- backend/src/workspaces/runtime-renderer.ts | 2 + backend/test/effective-config.test.ts | 197 ++++++++++++++++ .../workspace-preprocessing-service.test.ts | 20 +- .../workspace-runtime-config-lease.test.ts | 110 ++++++++- .../test/workspace-runtime-renderer.test.ts | 1 + docs/contracts/tht-dwh.md | 128 ++++++++++ docs/install/local-workspace-registry.md | 9 + docs/install/server-workspace-registry.md | 9 + docs/testing/p2-p6-manual-verification.md | 33 +-- harness/tests/test_effective_config_p3.py | 90 +++++++ harness/tests/test_p3_dwh_binding.py | 62 +++++ harness/tests/test_qdrant_vector_store.py | 3 +- harness/tht/adapters/vector/qdrant.py | 26 ++- harness/tht/cli/memory_cmd.py | 67 ++++++ harness/tht/config.py | 72 ++++++ harness/tht/jobs/dwh_pipeline.py | 47 ++-- 21 files changed, 1108 insertions(+), 44 deletions(-) create mode 100644 backend/src/workspaces/effective-config.ts create mode 100644 backend/test/effective-config.test.ts create mode 100644 docs/contracts/tht-dwh.md create mode 100644 harness/tests/test_effective_config_p3.py create mode 100644 harness/tests/test_p3_dwh_binding.py diff --git a/backend/src/tht/tht-runner.ts b/backend/src/tht/tht-runner.ts index 8aa57c8b..fc7f4b24 100644 --- a/backend/src/tht/tht-runner.ts +++ b/backend/src/tht/tht-runner.ts @@ -180,6 +180,7 @@ export class ThtRunner { sessions: join(root, "sessions"), artifacts: join(root, "artifacts"), indexes: join(root, "indexes"), + memory: join(root, "memory"), }; } diff --git a/backend/src/workspace-maintenance.ts b/backend/src/workspace-maintenance.ts index 20b8e01d..af7923ea 100644 --- a/backend/src/workspace-maintenance.ts +++ b/backend/src/workspace-maintenance.ts @@ -40,6 +40,7 @@ function failureResult( const STATE_ERROR_CODES: Record = { preprocessing_resume_mismatch: "preprocessing_resume_mismatch", preprocessing_conflict: "preprocessing_conflict", + effective_config_mismatch: "effective_config_mismatch", }; function boundedJson(result: WorkspaceOperationResult): string { diff --git a/backend/src/workspaces/effective-config.ts b/backend/src/workspaces/effective-config.ts new file mode 100644 index 00000000..29e06bf1 --- /dev/null +++ b/backend/src/workspaces/effective-config.ts @@ -0,0 +1,219 @@ +import { createHash } from "node:crypto"; +import { normalize } from "node:path"; + +export interface CanonicalEffectiveConfig { + schemaVersion: 1; + dwh: CanonicalDwhConfig; + vector: CanonicalVectorConfig; + embedding: CanonicalEmbeddingConfig; + roots: CanonicalRootsConfig; +} + +export interface CanonicalDwhConfig { + engine: "postgres"; + database: string; + schema: string; + transport: "postgres_direct" | "rest_api" | "ssh_tunnel"; + host?: string; + port?: number; + baseUrl?: string; + user?: string; +} + +export interface CanonicalVectorConfig { + collection: string; + dimensions: number; + distance: string; +} + +export interface CanonicalEmbeddingConfig { + model: string; + dimensions: number; +} + +export interface CanonicalRootsConfig { + artifacts: string; + indexes: string; +} + +function asRecord(value: unknown): Record | undefined { + if (typeof value === "object" && value !== null && !Array.isArray(value)) { + return value as Record; + } + return undefined; +} + +function requireString(value: Record, key: string): string { + const candidate = value[key]; + if (typeof candidate !== "string" || candidate.length === 0) { + throw new TypeError(`effective config is missing ${key}`); + } + return candidate; +} + +function optionalString(value: Record, key: string): string | undefined { + const candidate = value[key]; + if (candidate === undefined || candidate === null) return undefined; + if (typeof candidate !== "string") return undefined; + return candidate; +} + +function optionalNumber(value: Record, key: string): number | undefined { + const candidate = value[key]; + if (candidate === undefined || candidate === null) return undefined; + if (typeof candidate !== "number" || Number.isNaN(candidate)) return undefined; + return candidate; +} + +function requireNumber(value: Record, key: string): number { + const candidate = value[key]; + if (typeof candidate !== "number" || Number.isNaN(candidate)) { + throw new TypeError(`effective config is missing numeric ${key}`); + } + return candidate; +} + +function normalizeRoot(value: string): string { + return normalize(value); +} + +function dwhTransport(rendered: Record): CanonicalDwhConfig["transport"] { + const dwh = asRecord(rendered.dwh); + const type = dwh ? requireString(dwh, "type") : undefined; + if (type === "postgres_direct") return "postgres_direct"; + if (type === "thoth_rest") return "rest_api"; + if (type === "ssh_tunnel") return "ssh_tunnel"; + throw new TypeError(`effective config has unsupported dwh transport ${type}`); +} + +function buildDwhConfig(rendered: Record): CanonicalDwhConfig { + const transport = dwhTransport(rendered); + const databaseRecord = asRecord(rendered.database) ?? asRecord(asRecord(asRecord(rendered.dwh)?.connection)?.database); + if (!databaseRecord) { + throw new TypeError("effective config is missing database identity"); + } + const engine = "postgres"; + const database = requireString(databaseRecord, "database"); + const schema = requireString(databaseRecord, "schema"); + const dwh: Record = { engine, database, schema, transport }; + + if (transport === "postgres_direct") { + const host = optionalString(databaseRecord, "host"); + const port = optionalNumber(databaseRecord, "port"); + const user = optionalString(databaseRecord, "user"); + if (host !== undefined) dwh.host = host; + if (port !== undefined) dwh.port = port; + if (user !== undefined) dwh.user = user; + } else if (transport === "rest_api") { + const rest = asRecord(rendered.rest) ?? asRecord(asRecord(asRecord(rendered.dwh)?.endpoint)); + const baseUrl = rest ? optionalString(rest, "base_url") : undefined; + if (baseUrl !== undefined) dwh.baseUrl = baseUrl; + } + + return dwh as unknown as CanonicalDwhConfig; +} + +function buildVectorConfig(rendered: Record): CanonicalVectorConfig { + const resources = asRecord(rendered.resources); + const vector = resources ? asRecord(resources.vector) : undefined; + if (!vector) { + throw new TypeError("effective config is missing vector resources"); + } + const collection = requireString(vector, "collection"); + const semanticIndex = asRecord(rendered.semantic_index); + const vectorStore = semanticIndex ? asRecord(semanticIndex.vector_store) : undefined; + const dimensions = vectorStore + ? requireNumber(vectorStore, "dimensions") + : requireNumber(vector, "dimensions"); + const distance = vectorStore + ? requireString(vectorStore, "distance") + : (optionalString(vector, "distance") ?? "cosine"); + return { collection, dimensions, distance }; +} + +function buildEmbeddingConfig(rendered: Record): CanonicalEmbeddingConfig { + const resources = asRecord(rendered.resources); + const embeddings = resources ? asRecord(resources.embeddings) : undefined; + if (!embeddings) { + throw new TypeError("effective config is missing embedding resources"); + } + return { + model: requireString(embeddings, "model"), + dimensions: requireNumber(embeddings, "dimensions"), + }; +} + +function buildRootsConfig(rendered: Record): CanonicalRootsConfig { + const roots = asRecord(rendered.roots) ?? asRecord(rendered.paths); + if (!roots) { + throw new TypeError("effective config is missing roots"); + } + return { + artifacts: normalizeRoot(requireString(roots, "artifacts")), + indexes: normalizeRoot(requireString(roots, "indexes")), + }; +} + +/** + * Build the versioned, non-secret effective DWH/preprocessing configuration from a + * rendered runtime configuration object. The result contains only the fields that + * affect DWH generation identity; credentials, runtime identity, session storage, + * evidence, and service endpoints are excluded. + */ +export function buildCanonicalEffectiveConfig(renderedConfig: unknown): CanonicalEffectiveConfig { + const rendered = asRecord(renderedConfig); + if (!rendered) { + throw new TypeError("effective config requires a rendered configuration object"); + } + return { + schemaVersion: 1, + dwh: buildDwhConfig(rendered), + vector: buildVectorConfig(rendered), + embedding: buildEmbeddingConfig(rendered), + roots: buildRootsConfig(rendered), + }; +} + +function sha256(value: string | Buffer): string { + return `sha256:${createHash("sha256").update(value).digest("hex")}`; +} + +/** + * Serialize the canonical effective config to a deterministic JSON string with the + * fixed key order defined by the shared contract. No whitespace is included. + */ +export function canonicalEffectiveConfigJson(config: CanonicalEffectiveConfig): string { + const ordered: Record = { schemaVersion: config.schemaVersion }; + ordered.dwh = { ...config.dwh }; + ordered.vector = { ...config.vector }; + ordered.embedding = { ...config.embedding }; + ordered.roots = { ...config.roots }; + return JSON.stringify(ordered); +} + +/** + * Return the stable logical config-source identity for a workspace revision. + * This is `workspace://@v1:`. + */ +export function effectiveConfigIdentity(workspaceId: string, renderedConfig: unknown): string { + const canonical = buildCanonicalEffectiveConfig(renderedConfig); + const digest = createHash("sha256").update(canonicalEffectiveConfigJson(canonical)).digest("hex"); + return `workspace://${workspaceId}@v1:${digest}`; +} + +/** + * Return the config fingerprint: `sha256:` + the SHA-256 of the canonical effective + * config JSON bytes. + */ +export function configFingerprint(renderedConfig: unknown): string { + const canonical = buildCanonicalEffectiveConfig(renderedConfig); + return sha256(canonicalEffectiveConfigJson(canonical)); +} + +/** + * Return the input fingerprint: `sha256:` + the SHA-256 of the logical config-source + * identity string. + */ +export function inputFingerprint(workspaceId: string, renderedConfig: unknown): string { + return sha256(effectiveConfigIdentity(workspaceId, renderedConfig)); +} diff --git a/backend/src/workspaces/preprocessing-service.ts b/backend/src/workspaces/preprocessing-service.ts index 87015df1..3e77fe78 100644 --- a/backend/src/workspaces/preprocessing-service.ts +++ b/backend/src/workspaces/preprocessing-service.ts @@ -30,6 +30,9 @@ export interface WorkspaceOperationResult { completedStages: string[]; counts?: Record; artifactIdentities?: Array<{ kind: string; digest: string }>; + effectiveConfigIdentity?: string; + configFingerprint?: string; + inputFingerprint?: string; /** Suggested FK annotations YAML for the operator to write to --output (schema suggest-fks). */ suggestedFksYaml?: string; warnings?: string[]; @@ -92,6 +95,9 @@ function baseResult( descriptorBlob: runtime.descriptorBlob, operation, completedStages: [], + effectiveConfigIdentity: runtime.configLease.effectiveConfigIdentity, + configFingerprint: runtime.configLease.configFingerprint, + inputFingerprint: runtime.configLease.inputFingerprint, ...extra, }; } diff --git a/backend/src/workspaces/runtime-config-lease.ts b/backend/src/workspaces/runtime-config-lease.ts index af126c9c..6d911862 100644 --- a/backend/src/workspaces/runtime-config-lease.ts +++ b/backend/src/workspaces/runtime-config-lease.ts @@ -18,6 +18,14 @@ import { import { dirname, isAbsolute, join, relative, resolve } from "node:path"; import { readFile as readFileAsync } from "node:fs/promises"; import { parse, parseAllDocuments, stringify } from "yaml"; +import { + buildCanonicalEffectiveConfig, + canonicalEffectiveConfigJson, + configFingerprint, + effectiveConfigIdentity, + inputFingerprint, + type CanonicalEffectiveConfig, +} from "./effective-config.js"; import { resolveRuntimeBindings, type RuntimeBindings } from "./bindings.js"; import { GitWorkspaceRepository } from "./git-repository.js"; import { WorkspaceRegistry } from "./registry.js"; @@ -63,6 +71,10 @@ export interface DeterministicRuntimeConfigLease extends RuntimeConfigLease { catalogBlob: string; configDigest: string; bindingDigest: string; + effectiveConfig: CanonicalEffectiveConfig; + effectiveConfigIdentity: string; + configFingerprint: string; + inputFingerprint: string; } export class RuntimeConfigLeaseError extends Error { @@ -84,6 +96,9 @@ interface PublishedRuntimeConfigManifest { catalogBlob: string; configDigest: string; bindingDigest: string; + effectiveConfigIdentity: string; + configFingerprint: string; + inputFingerprint: string; path: string; file: { dev: number; @@ -230,6 +245,7 @@ function runtimePaths(dataRoot: string, workspaceId: string): RuntimePaths { sessions: join(root, "sessions"), artifacts: join(root, "artifacts"), indexes: join(root, "indexes"), + memory: join(root, "memory"), }; } @@ -388,6 +404,9 @@ function decodePublishedRuntimeConfigManifest(value: unknown): PublishedRuntimeC || typeof manifest.catalogBlob !== "string" || typeof manifest.configDigest !== "string" || typeof manifest.bindingDigest !== "string" + || typeof manifest.effectiveConfigIdentity !== "string" + || typeof manifest.configFingerprint !== "string" + || typeof manifest.inputFingerprint !== "string" || typeof manifest.path !== "string" || !file || typeof file.dev !== "number" @@ -413,6 +432,13 @@ export async function publishDeterministicRuntimeConfigLease(options: { }): Promise { const rendered = await renderActiveWorkspaceRuntime(options); const publishedConfig = applyCollectionLifecycle(rendered.renderedConfig, "require_existing"); + const renderedConfigObject = parse(rendered.renderedConfig) as Record; + const effectiveConfig = buildCanonicalEffectiveConfig(renderedConfigObject); + const effectiveConfigIdentityValue = effectiveConfigIdentity(rendered.workspaceId, renderedConfigObject); + const configFingerprintValue = configFingerprint(renderedConfigObject); + const inputFingerprintValue = inputFingerprint(rendered.workspaceId, renderedConfigObject); + const identitySuffix = inputFingerprintValue.slice(7, 23); + const preprocessingRoot = ensureTrustedDirectory(join( options.dataRoot, "sessions", @@ -421,8 +447,8 @@ export async function publishDeterministicRuntimeConfigLease(options: { )); const configDirectory = ensureTrustedDirectory(join(preprocessingRoot, "runtime-config")); const manifestDirectory = ensureTrustedDirectory(join(preprocessingRoot, "runtime-config-manifests")); - const path = join(configDirectory, `${rendered.workspaceRevision}.yaml`); - const manifestPath = join(manifestDirectory, `${rendered.workspaceRevision}.json`); + const path = join(configDirectory, `${rendered.workspaceRevision}-${identitySuffix}.yaml`); + const manifestPath = join(manifestDirectory, `${rendered.workspaceRevision}-${identitySuffix}.json`); const configDigest = sha256(publishedConfig); const verifyPublished = (): PublishedRuntimeConfigManifest | undefined => { @@ -456,6 +482,9 @@ export async function publishDeterministicRuntimeConfigLease(options: { || manifest.catalogBlob !== rendered.catalogBlob || manifest.configDigest !== configDigest || manifest.bindingDigest !== rendered.bindingDigest + || manifest.effectiveConfigIdentity !== effectiveConfigIdentityValue + || manifest.configFingerprint !== configFingerprintValue + || manifest.inputFingerprint !== inputFingerprintValue || manifest.path !== path || manifest.file.dev !== configStat!.dev || manifest.file.ino !== configStat!.ino @@ -479,6 +508,10 @@ export async function publishDeterministicRuntimeConfigLease(options: { catalogBlob: rendered.catalogBlob, configDigest, bindingDigest: rendered.bindingDigest, + effectiveConfig, + effectiveConfigIdentity: effectiveConfigIdentityValue, + configFingerprint: configFingerprintValue, + inputFingerprint: inputFingerprintValue, release: () => undefined, }; } @@ -505,6 +538,9 @@ export async function publishDeterministicRuntimeConfigLease(options: { catalogBlob: rendered.catalogBlob, configDigest, bindingDigest: rendered.bindingDigest, + effectiveConfigIdentity: effectiveConfigIdentityValue, + configFingerprint: configFingerprintValue, + inputFingerprint: inputFingerprintValue, path, file: { dev: Number(publishedStat.dev), @@ -518,7 +554,8 @@ export async function publishDeterministicRuntimeConfigLease(options: { readStrictJson(manifestPath); } catch (error) { if ((error as NodeJS.ErrnoException).code === "ENOENT") { - writeAtomicFile(manifestPath, `${JSON.stringify(manifest)}\n`, 0o600); + writeAtomicFile(manifestPath, `${JSON.stringify(manifest)} +`, 0o600); } else { throw error; } @@ -533,6 +570,10 @@ export async function publishDeterministicRuntimeConfigLease(options: { catalogBlob: rendered.catalogBlob, configDigest, bindingDigest: rendered.bindingDigest, + effectiveConfig, + effectiveConfigIdentity: effectiveConfigIdentityValue, + configFingerprint: configFingerprintValue, + inputFingerprint: inputFingerprintValue, release: () => undefined, }; -} +} \ No newline at end of file diff --git a/backend/src/workspaces/runtime-renderer.ts b/backend/src/workspaces/runtime-renderer.ts index f6499f89..4e1500f9 100644 --- a/backend/src/workspaces/runtime-renderer.ts +++ b/backend/src/workspaces/runtime-renderer.ts @@ -9,6 +9,7 @@ export interface RuntimePaths { sessions: string; artifacts: string; indexes: string; + memory: string; } export interface RuntimeIdentity { @@ -255,6 +256,7 @@ export function renderRuntimeConfig( ...(installation.profile === undefined ? {} : { profile: installation.profile }), language: descriptor.workspace.language, database, + semantic_index: descriptor.semantic_index, resources: { vector: { engine: "qdrant", diff --git a/backend/test/effective-config.test.ts b/backend/test/effective-config.test.ts new file mode 100644 index 00000000..ea60a30b --- /dev/null +++ b/backend/test/effective-config.test.ts @@ -0,0 +1,197 @@ +import { createHash } from "node:crypto"; +import { expect, test } from "vitest"; +import { + buildCanonicalEffectiveConfig, + canonicalEffectiveConfigJson, + configFingerprint, + effectiveConfigIdentity, + inputFingerprint, +} from "../src/workspaces/effective-config.js"; + +const semanticRuntime = { + internalQdrantUrl: "http://qdrant:6333", + internalEmbeddingUrl: "http://embedding:11434", + internalEmbeddingModel: "qwen3-embedding:0.6b", + internalEmbeddingDimensions: 1024, +}; + +function directRendered(): Record { + return { + runtime_identity: { + workspace_id: "psd-clinical", + workspace_revision: "a".repeat(40), + source_identity: "workspace://psd-clinical", + }, + session_storage: { mode: "local" }, + profile: "server", + language: "it", + database: { + host: "dwh.internal", + port: 5432, + database: "postgres", + schema: "datawarehouse", + user: "thoth_reader", + password_file: "/run/secrets/dwh-password", + ssl_ca_file: "/run/secrets/dwh-ca.pem", + transport: "direct", + }, + dwh: { type: "postgres_direct" }, + resources: { + vector: { + engine: "qdrant", + base_url: "http://qdrant:6333", + collection: "psd-clinical", + dimensions: 1024, + distance: "cosine", + collection_lifecycle: "require_existing", + }, + embeddings: { + provider: "ollama_internal", + base_url: "http://embedding:11434", + model: "qwen3-embedding:0.6b", + dimensions: 1024, + }, + }, + roots: { + sessions: "/data/sessions/psd-clinical/sessions", + artifacts: "/data/sessions/psd-clinical/artifacts", + indexes: "/data/sessions/psd-clinical/indexes", + }, + paths: { + sessions: "/data/sessions/psd-clinical/sessions", + artifacts: "/data/sessions/psd-clinical/artifacts", + indexes: "/data/sessions/psd-clinical/indexes", + memory: "/data/sessions/psd-clinical/memory", + }, + evidence: { + sources: [{ + type: "filesystem", + root: "/srv/registry/snapshots/rev/psd-clinical/evidence", + patterns: ["**/*.md"], + max_bytes: 10_485_760, + }], + }, + }; +} + +function restRendered(): Record { + return { + ...directRendered(), + database: { + host: "localhost", + port: 5432, + database: "postgres", + schema: "datawarehouse", + user: "rest", + password: "", + transport: "rest", + }, + dwh: { + type: "thoth_rest", + database: { database: "postgres", schema: "datawarehouse" }, + endpoint: { + base_url: "https://dwh.example.test", + api_key_file: "/run/secrets/dwh-api-key", + }, + }, + rest: { + base_url: "https://dwh.example.test", + api_key_file: "/run/secrets/dwh-api-key", + }, + }; +} + +const directCanonical = + `{"schemaVersion":1,"dwh":{` + + `"engine":"postgres","database":"postgres","schema":"datawarehouse",` + + `"transport":"postgres_direct","host":"dwh.internal","port":5432,"user":"thoth_reader"},` + + `"vector":{"collection":"psd-clinical","dimensions":1024,"distance":"cosine"},` + + `"embedding":{"model":"qwen3-embedding:0.6b","dimensions":1024},` + + `"roots":{"artifacts":"/data/sessions/psd-clinical/artifacts",` + + `"indexes":"/data/sessions/psd-clinical/indexes"}}`; + +test("canonical effective config is deterministic and contains the expected key order", () => { + const rendered = directRendered(); + const canonical = buildCanonicalEffectiveConfig(rendered); + expect(canonicalEffectiveConfigJson(canonical)).toBe(directCanonical); + expect(buildCanonicalEffectiveConfig(rendered)).toEqual(canonical); +}); + +test("canonical effective config excludes secrets, evidence, session storage, and runtime identity", () => { + const json = canonicalEffectiveConfigJson(buildCanonicalEffectiveConfig(directRendered())); + expect(json).not.toContain("password_file"); + expect(json).not.toContain("ssl_ca_file"); + expect(json).not.toContain("session_storage"); + expect(json).not.toContain("runtime_identity"); + expect(json).not.toContain("evidence"); + expect(json).not.toContain("sources"); + expect(json).not.toContain("collection_lifecycle"); + expect(json).not.toContain("base_url"); // vector/embedding service URLs are not identity + expect(json).not.toContain("memory"); + expect(json).not.toContain('"sessions"'); +}); + +test("REST transport canonicalizes to transport rest_api with baseUrl and no host/port", () => { + const canonical = buildCanonicalEffectiveConfig(restRendered()); + const json = canonicalEffectiveConfigJson(canonical); + expect(json).toContain(`"transport":"rest_api"`); + expect(json).toContain(`"baseUrl":"https://dwh.example.test"`); + expect(json).not.toContain(`"host":`); + expect(json).not.toContain(`"port":`); + expect(json).not.toContain("api_key_file"); +}); + +test("identity and fingerprint helpers produce stable prefixed hex values", () => { + const rendered = directRendered(); + const identity = effectiveConfigIdentity("psd-clinical", rendered); + const cfg = configFingerprint(rendered); + const input = inputFingerprint("psd-clinical", rendered); + + expect(identity).toMatch(/^workspace:\/\/psd-clinical@v1:[0-9a-f]{64}$/); + expect(cfg).toBe("sha256:" + createHash("sha256").update(directCanonical).digest("hex")); + expect(input).toBe("sha256:" + createHash("sha256").update(identity).digest("hex")); + expect(input).not.toBe(cfg); +}); + +test("content-only or session_storage changes keep the same effective config identity", () => { + const base = directRendered(); + const identityBefore = effectiveConfigIdentity("psd-clinical", base); + const fingerprintBefore = configFingerprint(base); + + const contentOnly = { + ...base, + runtime_identity: { + ...base.runtime_identity, + workspace_revision: "b".repeat(40), + }, + session_storage: { mode: "remote", url: "http://example.test" }, + evidence: { + sources: [{ + type: "filesystem", + root: "/srv/registry/snapshots/other/psd-clinical/evidence", + patterns: ["**/*.txt"], + max_bytes: 999, + }], + }, + }; + + expect(effectiveConfigIdentity("psd-clinical", contentOnly)).toBe(identityBefore); + expect(configFingerprint(contentOnly)).toBe(fingerprintBefore); +}); + +test("DWH-affecting changes alter the effective config identity", () => { + const base = directRendered(); + const identityBefore = effectiveConfigIdentity("psd-clinical", base); + + const changedHost = { ...base, database: { ...(base.database as object), host: "dwh-two.internal" } }; + expect(effectiveConfigIdentity("psd-clinical", changedHost)).not.toBe(identityBefore); + + const changedDatabase = { ...base, database: { ...(base.database as object), database: "analytics" } }; + expect(effectiveConfigIdentity("psd-clinical", changedDatabase)).not.toBe(identityBefore); + + const changedCollection = { ...base, resources: { ...base.resources, vector: { ...(base.resources as Record).vector, collection: "other" } } }; + expect(effectiveConfigIdentity("psd-clinical", changedCollection)).not.toBe(identityBefore); + + const changedTransport = restRendered(); + expect(effectiveConfigIdentity("psd-clinical", changedTransport)).not.toBe(identityBefore); +}); diff --git a/backend/test/workspace-preprocessing-service.test.ts b/backend/test/workspace-preprocessing-service.test.ts index 66471a93..9a7ef4c3 100644 --- a/backend/test/workspace-preprocessing-service.test.ts +++ b/backend/test/workspace-preprocessing-service.test.ts @@ -115,13 +115,31 @@ function runtime(workspace = baseWorkspace, workspaceId = workspace.workspace.id descriptorBlob: "b".repeat(40), catalogBlob: "c".repeat(40), configLease: { - path: `/data/sessions/${workspaceId}/preprocessing/runtime-config/${"a".repeat(40)}.yaml`, + path: `/data/sessions/${workspaceId}/preprocessing/runtime-config/${"a".repeat(40)}-identitysuffix.yaml`, workspaceId, workspaceRevision: "a".repeat(40), descriptorBlob: "b".repeat(40), catalogBlob: "c".repeat(40), configDigest: "sha256:config", bindingDigest: "sha256:bindings", + effectiveConfig: { + schemaVersion: 1, + dwh: { + engine: "postgres", + database: "analytics", + schema: "mart", + transport: "postgres_direct", + host: "dwh.internal", + port: 5432, + user: "reader", + }, + vector: { collection: workspaceId, dimensions: 1024, distance: "cosine" }, + embedding: { model: "qwen3-embedding:0.6b", dimensions: 1024 }, + roots: { artifacts: "/data/artifacts", indexes: "/data/indexes" }, + }, + effectiveConfigIdentity: "workspace://psd-clinical@v1:" + "d".repeat(64), + configFingerprint: "sha256:" + "e".repeat(64), + inputFingerprint: "sha256:" + "f".repeat(64), release: () => undefined, }, }; diff --git a/backend/test/workspace-runtime-config-lease.test.ts b/backend/test/workspace-runtime-config-lease.test.ts index ab491cac..1512797e 100644 --- a/backend/test/workspace-runtime-config-lease.test.ts +++ b/backend/test/workspace-runtime-config-lease.test.ts @@ -14,7 +14,13 @@ import { tmpdir } from "node:os"; import { join } from "node:path"; import { promisify } from "node:util"; import { afterEach, expect, test, vi } from "vitest"; +import { parse } from "yaml"; import { WorkspaceRegistry } from "../src/workspaces/registry.js"; +import { + buildCanonicalEffectiveConfig, + canonicalEffectiveConfigJson, + effectiveConfigIdentity, +} from "../src/workspaces/effective-config.js"; import type { WorkspaceRegistryConfig } from "../src/workspaces/types.js"; import { publishDeterministicRuntimeConfigLease, @@ -123,6 +129,7 @@ evidence: return { dataRoot, harnessDir, + source, registry, registryConfig, revision, @@ -163,7 +170,7 @@ test("active workspace rendering is byte-identical to direct snapshot rendering" expect(active.catalogBlob).toMatch(/^sha256:[0-9a-f]{64}$/); }); -test("deterministic operator leases publish one revision-bound protected config and refuse changed same-revision bytes", async () => { +test("deterministic operator leases are keyed by logical identity and stable across calls", async () => { const f = await fixture(); const first = await publishDeterministicRuntimeConfigLease({ workspaceId: "psd-clinical", @@ -186,6 +193,7 @@ test("deterministic operator leases publish one revision-bound protected config semanticRuntime, }); + const suffix = first.inputFingerprint.slice(7, 23); expect(second.path).toBe(first.path); expect(first.path).toBe(join( f.dataRoot, @@ -193,15 +201,34 @@ test("deterministic operator leases publish one revision-bound protected config "psd-clinical", "preprocessing", "runtime-config", - `${f.revision.commit}.yaml`, + `${f.revision.commit}-${suffix}.yaml`, )); expect(statSync(first.path).mode & 0o777).toBe(0o400); expect(statSync(first.manifestPath).mode & 0o777).toBe(0o600); expect(readFileSync(first.path, "utf8")).toContain("collection_lifecycle: require_existing"); + expect(readFileSync(first.path, "utf8")).toContain("memory:"); expect(existsSync(first.manifestPath)).toBe(true); + expect(first.effectiveConfigIdentity).toMatch(/^workspace:\/\/psd-clinical@v1:[0-9a-f]{64}$/); + expect(first.configFingerprint).toMatch(/^sha256:[0-9a-f]{64}$/); + expect(first.inputFingerprint).toMatch(/^sha256:[0-9a-f]{64}$/); + expect(first.inputFingerprint).not.toBe(first.configFingerprint); + const manifest = JSON.parse(readFileSync(first.manifestPath, "utf8")); + expect(manifest).toMatchObject({ + schemaVersion: 1, + workspaceId: "psd-clinical", + workspaceRevision: f.revision.commit, + descriptorBlob: first.descriptorBlob, + catalogBlob: first.catalogBlob, + configDigest: first.configDigest, + bindingDigest: first.bindingDigest, + effectiveConfigIdentity: first.effectiveConfigIdentity, + configFingerprint: first.configFingerprint, + inputFingerprint: first.inputFingerprint, + path: first.path, + }); vi.stubEnv("THT_WS_PSD_CLINICAL_DWH_HOST", "warehouse-two.internal"); - await expect(publishDeterministicRuntimeConfigLease({ + const changed = await publishDeterministicRuntimeConfigLease({ workspaceId: "psd-clinical", registry: f.registry, registryConfig: f.registryConfig, @@ -210,7 +237,11 @@ test("deterministic operator leases publish one revision-bound protected config dataRoot: f.dataRoot, secretRoots: f.registryConfig.secretRoots, semanticRuntime, - })).rejects.toMatchObject({ code: "effective_config_mismatch" }); + }); + expect(changed.path).not.toBe(first.path); + expect(changed.inputFingerprint).not.toBe(first.inputFingerprint); + expect(changed.configFingerprint).not.toBe(first.configFingerprint); + expect(readFileSync(changed.path, "utf8")).toContain("warehouse-two.internal"); }); test("runtime rendering rejects untrusted snapshot paths and symlinks", async () => { @@ -237,3 +268,74 @@ test("runtime rendering rejects untrusted snapshot paths and symlinks", async () semanticRuntime, })).toThrow(/trusted runtime snapshot/i); }); + + +test("operator lease and session snapshot produce byte-identical effective DWH bindings", async () => { + const f = await fixture(); + const session = renderWorkspaceRuntimeFromSnapshotPath({ + snapshotPath: f.revision.snapshotPath, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + const lease = await publishDeterministicRuntimeConfigLease({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + + const sessionCanonical = canonicalEffectiveConfigJson(buildCanonicalEffectiveConfig(parse(session.renderedConfig))); + const operatorCanonical = canonicalEffectiveConfigJson(lease.effectiveConfig); + expect(operatorCanonical).toBe(sessionCanonical); + expect(lease.effectiveConfigIdentity).toBe( + effectiveConfigIdentity("psd-clinical", parse(session.renderedConfig)), + ); +}); + +test("a content-only Evidence commit keeps the same effective config identity with a new revision lease", async () => { + const f = await fixture(); + const firstLease = await publishDeterministicRuntimeConfigLease({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + const firstIdentity = firstLease.effectiveConfigIdentity; + + writeFileSync( + join(f.source, "psd-clinical", "evidence", "guide.md"), + "# updated content only\n", + ); + await git(f.source, ["add", "psd-clinical/evidence/guide.md"]); + await git(f.source, ["commit", "-m", "Evidence content only"]); + await git(f.source, ["push", "origin", "main"]); + await f.registry.pull(); + const current = (await f.registry.list())[0]; + + const secondLease = await publishDeterministicRuntimeConfigLease({ + workspaceId: "psd-clinical", + registry: f.registry, + registryConfig: f.registryConfig, + harnessDir: f.harnessDir, + configPath: "config/tht.yaml", + dataRoot: f.dataRoot, + secretRoots: f.registryConfig.secretRoots, + semanticRuntime, + }); + + expect(secondLease.workspaceRevision).toBe(current.commit); + expect(secondLease.workspaceRevision).not.toBe(firstLease.workspaceRevision); + expect(secondLease.effectiveConfigIdentity).toBe(firstIdentity); + expect(secondLease.path).not.toBe(firstLease.path); +}); diff --git a/backend/test/workspace-runtime-renderer.test.ts b/backend/test/workspace-runtime-renderer.test.ts index 6fc85ec3..8a151f38 100644 --- a/backend/test/workspace-runtime-renderer.test.ts +++ b/backend/test/workspace-runtime-renderer.test.ts @@ -40,6 +40,7 @@ const paths: RuntimePaths = { sessions: "/data/workspaces/psd-clinical/sessions", artifacts: "/data/workspaces/psd-clinical/artifacts", indexes: "/data/workspaces/psd-clinical/indexes", + memory: "/data/workspaces/psd-clinical/memory", }; const semanticRuntime: SemanticRuntimeConfig = { internalQdrantUrl: "http://qdrant:6333", diff --git a/docs/contracts/tht-dwh.md b/docs/contracts/tht-dwh.md new file mode 100644 index 00000000..ab1b59e9 --- /dev/null +++ b/docs/contracts/tht-dwh.md @@ -0,0 +1,128 @@ +# `.tht-dwh` — DWH generations, `OWNER.json`, ACTIVE, and fingerprints + +> Operator contract. P3 makes the effective DWH/preprocessing configuration reproducible and +> versioned across the operator CLI and the application sessions, and documents what `.tht-dwh` +> is so operators can reason about why a rerun is instant or why it takes minutes. + +## What `.tht-dwh` is + +`.tht-dwh` is the workspace-local directory that stores the **prepared snapshots of the data +warehouse structure** (the catalog `physical.yaml` plus the LSH hashes used for fuzzy search). +ThothII does not re-read the whole database for every question: it prepares it once, stores the +result here, and reuses it. The directory lives under the workspace runtime root, for example: + +```text +/data/sessions//.tht-dwh/ +``` + +## Immutable generations + +Each preparation run produces a **generation**: an immutable directory containing the catalog and +the LSH artifacts for one exact "effective configuration" (see fingerprints below). Generations +are never modified in place; a new run writes a new generation, and an `ACTIVE` pointer selects +which generation the workspace currently uses. Keeping the old generations makes rollback and +diagnosis safe. + +## `OWNER.json` + +Every generation root contains an `OWNER.json` that records who owns it: + +```json +{ + "workspace_id": "", + "config_fingerprint": "sha256:<64 hex>", + "input_fingerprint": "sha256:<64 hex>" +} +``` + +- `config_fingerprint` is the digest of the **canonical effective configuration** (see below). +- `input_fingerprint` is the digest of the **logical configuration identity**. + +Before reusing a generation, the harness compares the current canonical identity with the one in +`OWNER.json`. If they differ, the generation is **refused** (never silently reused) and a new one +is produced. This is what protects ThothII from using artifacts prepared for a different database, +endpoint, user, schema, or index contract. + +The reader is compatible with the historical schema-v1 `OWNER.json` (same three keys, `sha256:` +values) so existing installations keep working; new writes use the versioned computation. There is +no automatic in-place reinterpretation: operators regenerate explicitly when a root is old. + +## The canonical effective configuration and the logical identity + +The **canonical effective configuration** is the non-secret subset of the rendered runtime +configuration that determines whether a prepared DWH generation is still valid: + +```json +{ + "schemaVersion": 1, + "dwh": { + "engine": "postgres", + "database": "", + "schema": "", + "transport": "postgres_direct | rest_api | ...", + "host": "", + "port": 5432, + "baseUrl": "", + "user": "" + }, + "vector": { "collection": "", "dimensions": 1024, "distance": "cosine" }, + "embedding": { "model": "", "dimensions": 1024 }, + "roots": { "artifacts": "", "indexes": "" } +} +``` + +Deliberately **excluded** (their change must not invalidate a DWH generation): + +- `session_storage` and `runtime_identity` (a content-only Git commit or an Evidence-only change + must not force a full database re-introspection); +- Evidence source/policy (P6 materialization and Evidence preprocessing are separate); +- memory, search, and execution settings; +- **all credentials** (passwords, API keys, signed URLs, and secret-file paths). + +The **logical configuration identity** is: + +```text +workspace://@v1: +``` + +It is the same for the operator CLI and for application sessions, because both derive it from the +same rendered configuration. That is the guarantee that the work prepared by `thothctl` is exactly +what the sessions will consume. + +## Why a rerun can be instant or take minutes + +- Same canonical identity (e.g., only Evidence files changed) → the generation is reused → the + DWH step is `unchanged` and fast. +- Changed canonical identity (different database, address, user, schema, collection, model, or + artifact/index roots) → the old generation is refused → ThothII re-introspects and writes a new + generation → the step takes as long as the first preparation. + +## Safe migration, regeneration, and recovery + +- **Migration**: existing schema-v1 `OWNER.json` roots are readable; to switch them to the + versioned identity, run a normal regeneration (explicit `--refresh`/new run). No automatic + in-place rewrite. +- **Regeneration**: a new run produces a new immutable generation and moves `ACTIVE`; the previous + generations remain for rollback. +- **Recovery**: if the active generation is corrupt or owned by another configuration, ThothII + fails closed (never mixes artifacts) and tells the operator to regenerate; the old generations + are still available for inspection. + +## Memory root + +P3 also gives each workspace an explicit **workspace-global memory root**: + +```text +/data/sessions//memory/ +``` + +All memory commands, locks, the canonical JSONL registry, and the Qdrant projection use this root +when present. A guarded migration copies and verifies exactly one legacy canonical JSONL from the +old `artifacts/memory` location under the workspace lock and rebuilds the projection; conflicting +legacy registries fail closed. There is no in-place reinterpretation. + +## Revision-scoped search records + +Schema and Evidence records in the Qdrant collection include the pinned `workspace_revision`, so +searches never mix descriptions or documents from different versions of the workspace. Memory and +solved-question records remain workspace-wide on purpose. diff --git a/docs/install/local-workspace-registry.md b/docs/install/local-workspace-registry.md index 1fea6358..2008bad9 100644 --- a/docs/install/local-workspace-registry.md +++ b/docs/install/local-workspace-registry.md @@ -30,6 +30,15 @@ per `docs/contracts/workspace-preprocessing-cli.md` and the P2 walkthrough in `workspace-maintenance` Compose service; it never starts a backend/Pi/frontend listener and never attaches Git credentials. + +## Effective configuration and `.tht-dwh` (P3) + +Prepared DWH generations are reusable and safe: `thothctl` and the application derive the same +canonical effective configuration and logical identity, so prepared work is reused when nothing +relevant changed and refused when the database/endpoint/identity changed. See +`docs/contracts/tht-dwh.md` for generations, `OWNER.json`, `ACTIVE`, fingerprints, migration and +recovery. A content-only or Evidence-only change never forces a full re-introspection. + ## Prerequisites - macOS: Docker Desktop, Git, and sufficient volume disk space. Git Credential Manager is useful diff --git a/docs/install/server-workspace-registry.md b/docs/install/server-workspace-registry.md index c96911d9..01a9c040 100644 --- a/docs/install/server-workspace-registry.md +++ b/docs/install/server-workspace-registry.md @@ -25,6 +25,15 @@ per `docs/contracts/workspace-preprocessing-cli.md` and the P2 walkthrough in `workspace-maintenance` Compose service; it never starts a backend/Pi/frontend listener and never attaches Git credentials. + +## Effective configuration and `.tht-dwh` (P3) + +Prepared DWH generations are reusable and safe: `thothctl` and the application derive the same +canonical effective configuration and logical identity, so prepared work is reused when nothing +relevant changed and refused when the database/endpoint/identity changed. See +`docs/contracts/tht-dwh.md` for generations, `OWNER.json`, `ACTIVE`, fingerprints, migration and +recovery. A content-only or Evidence-only change never forces a full re-introspection. + ## Service account, storage, and firewall Create a dedicated host service account and an operator root such as `/srv/thothii`. The core diff --git a/docs/testing/p2-p6-manual-verification.md b/docs/testing/p2-p6-manual-verification.md index dca9a92b..86d62540 100644 --- a/docs/testing/p2-p6-manual-verification.md +++ b/docs/testing/p2-p6-manual-verification.md @@ -56,25 +56,32 @@ Checks: Decision: **PENDING** (independent manual gate; automation never records PASS). -## P3 — Effective config and `.tht-dwh` +## P3 — Effective configuration and `.tht-dwh` -**Status:** instructions to be finalized by P3 implementation; not yet runnable. +**Status:** P3 implementation complete; automated integration PASS; manual acceptance PENDING. -Manual goal: compare operator and session effective DWH identities, inspect `OWNER.json` and -`ACTIVE` without exposing secrets, prove safe reuse after a content-only revision, and prove -fail-closed behavior after a DWH-affecting change. +Manual goal: prove that the operator CLI and application sessions derive the same effective +configuration, that a content-only revision reuses the prepared DWH generation (fast, `unchanged`), +that a DWH-affecting change fails closed and regenerates, that the workspace memory migration is +safe, and that search records are revision-scoped. See `docs/contracts/tht-dwh.md`. -Checks to fill during P3: +Checks: -1. canonical fingerprint comparison; -2. stable logical config-source identity; -3. schema-v1 ownership compatibility/migration; -4. DWH cache reuse across equivalent revisions; -5. revision-scoped schema/Evidence state; -6. mismatch rejection and recovery. +1. run `thothctl ... workspace preprocess dwh` twice with only an Evidence/content change between + them: the second run reports `unchanged` and does not re-introspect; +2. change a DWH-affecting field (host/port/database/schema/user/collection) in the descriptor, + push, pull: the next run refuses the old generation and regenerates, with a clear + `effective_config_mismatch`-style outcome and no mixed artifacts; +3. inspect `.tht-dwh` generations: immutable directories, `OWNER.json` with the canonical + fingerprints, `ACTIVE` pointer; old generations still present; +4. memory: after the guarded migration the workspace uses + `/sessions//memory/`; the JSONL registry and Qdrant projection are + rebuilt and consistent; a conflicting legacy registry fails closed; +5. search records: schema/Evidence points carry the pinned `workspace_revision`; memory/solved + records remain workspace-wide; +6. documentation: `docs/contracts/tht-dwh.md` matches the observed behavior. Decision: **PENDING**. - ## P4 — Qdrant bootstrap and guarded rebuild **Status:** instructions to be finalized by P4 implementation; not yet runnable. diff --git a/harness/tests/test_effective_config_p3.py b/harness/tests/test_effective_config_p3.py new file mode 100644 index 00000000..7e61bc26 --- /dev/null +++ b/harness/tests/test_effective_config_p3.py @@ -0,0 +1,90 @@ + +import pytest + +from tht.config import ( + canonical_effective_config_json, + effective_config_fingerprint, + effective_config_identity, + effective_config_input_fingerprint, +) + + +def _write_cfg(tmp_path, raw): + import yaml as _yaml + + from tht.config import load_config + path = tmp_path / "config.yaml" + path.write_text(_yaml.safe_dump(raw), encoding="utf-8") + cfg = load_config(path) + cfg._workspace_id = raw.get("workspace", {}).get("id", "psd") + cfg._config_source = f"workspace://{cfg._workspace_id}" + return cfg + + +def _cfg(tmp_path, *, transport="thoth_rest", base_url="http://dwh.example.invalid", collection="psd", model="qwen3-embedding:0.6b"): + return { + "schemaVersion": 1, + "workspace": {"schema_version": 3, "id": "psd", "name": "PSD", "language": "it"}, + "dwh": { + "type": transport, + "database": {"database": "warehouse", "schema": "dw"}, + "endpoint": {"base_url": base_url, "api_key": "secret"}, + } if transport == "thoth_rest" else { + "type": "postgres_direct", + "connection": {"host": "h", "port": 5432, "database": "warehouse", "schema": "dw", "user": "reader", "password": "secret"}, + }, + "vectors": {"type": "qdrant", "base_url": "http://qdrant:6333", "collection": collection, "collection_lifecycle": "self_heal"}, + "embeddings": {"provider": "ollama_internal", "base_url": "http://embedding:11434", "model": model, "dimensions": 1024}, + "roots": {"artifacts": str(tmp_path / "artifacts"), "indexes": str(tmp_path / "indexes")}, + "paths": {"artifacts": str(tmp_path / "artifacts"), "indexes": str(tmp_path / "indexes"), "sessions": str(tmp_path / "sessions")}, + } + + +@pytest.fixture() +def cfg(tmp_path): + return _write_cfg(tmp_path, _cfg(tmp_path)) + + +def test_canonical_json_is_deterministic_and_key_ordered(cfg): + doc = canonical_effective_config_json(cfg) + assert doc == canonical_effective_config_json(cfg) + keys = list(__import__("json").loads(doc).keys()) + assert keys == ["schemaVersion", "dwh", "vector", "embedding", "roots"] + + +def test_canonical_excludes_credentials_and_evidence(cfg): + doc = canonical_effective_config_json(cfg) + assert "secret" not in doc + assert "password" not in doc + assert "evidence" not in doc + assert "session_storage" not in doc + assert "runtime_identity" not in doc + + +def test_identity_and_fingerprints_format(cfg): + ident = effective_config_identity("psd", cfg) + assert ident == "workspace://psd@v1:" + effective_config_fingerprint(cfg)[len("sha256:"):] + assert effective_config_input_fingerprint("psd", cfg).startswith("sha256:") + assert len(effective_config_fingerprint(cfg)) == len("sha256:") + 64 + + +def test_content_only_change_keeps_identity(tmp_path): + base = _write_cfg(tmp_path, _cfg(tmp_path)) + changed = _cfg(tmp_path) + # A non-DWH-affecting change (collection lifecycle policy) must not alter the identity; + # the canonical document excludes it. Evidence changes are proven end-to-end by acceptance. + changed["vectors"]["collection_lifecycle"] = "require_existing" + changed_cfg = _write_cfg(tmp_path, changed) + assert effective_config_identity("psd", base) == effective_config_identity("psd", changed_cfg) + + +def test_dwh_affecting_change_alters_identity(tmp_path): + base = _write_cfg(tmp_path, _cfg(tmp_path)) + other = _write_cfg(tmp_path, _cfg(tmp_path, base_url="http://other.example.invalid")) + assert effective_config_identity("psd", base) != effective_config_identity("psd", other) + + +def test_transport_change_alters_identity(tmp_path): + rest = _write_cfg(tmp_path, _cfg(tmp_path, transport="thoth_rest")) + direct = _write_cfg(tmp_path, _cfg(tmp_path, transport="postgres_direct")) + assert effective_config_identity("psd", rest) != effective_config_identity("psd", direct) diff --git a/harness/tests/test_p3_dwh_binding.py b/harness/tests/test_p3_dwh_binding.py new file mode 100644 index 00000000..e1a7a6c9 --- /dev/null +++ b/harness/tests/test_p3_dwh_binding.py @@ -0,0 +1,62 @@ + +import pytest + + +def _cfg(tmp_path, **overrides): + base = { + "dwh": { + "type": "thoth_rest", + "database": {"database": "warehouse", "schema": "dw"}, + "endpoint": {"base_url": "http://dwh.example.invalid", "api_key": "secret"}, + }, + "vectors": {"type": "qdrant", "base_url": "http://qdrant:6333", "collection": "psd", "collection_lifecycle": "self_heal"}, + "embeddings": {"provider": "ollama_internal", "base_url": "http://embedding:11434", "model": "qwen3-embedding:0.6b", "dimensions": 1024}, + "roots": {"artifacts": str(tmp_path / "artifacts"), "indexes": str(tmp_path / "indexes")}, + "paths": {"artifacts": str(tmp_path / "artifacts"), "indexes": str(tmp_path / "indexes"), "sessions": str(tmp_path / "sessions")}, + } + base.update(overrides) + return base + + +def _write_cfg(tmp_path, raw): + import yaml as _yaml + + from tht.config import load_config + path = tmp_path / "config.yaml" + path.write_text(_yaml.safe_dump(raw), encoding="utf-8") + cfg = load_config(path) + cfg._workspace_id = "psd" + cfg._config_source = "workspace://psd" + return cfg + + +@pytest.fixture() +def binding(tmp_path): + from tht.jobs.dwh_pipeline import config_dwh_binding + return config_dwh_binding(_write_cfg(tmp_path, _cfg(tmp_path))) + + +def test_binding_has_versioned_fingerprints(binding): + assert binding["workspace_id"] == "psd" + assert binding["config_fingerprint"].startswith("sha256:") + assert binding["input_fingerprint"].startswith("sha256:") + assert len(binding["config_fingerprint"]) == 71 + + +def test_content_only_change_keeps_binding(tmp_path): + from tht.jobs.dwh_pipeline import config_dwh_binding + cfg1 = _write_cfg(tmp_path, _cfg(tmp_path)) + import copy + raw = copy.deepcopy(_cfg(tmp_path)) + raw["vectors"]["collection_lifecycle"] = "require_existing" + cfg3 = _write_cfg(tmp_path, raw) + assert config_dwh_binding(cfg1)["config_fingerprint"] == config_dwh_binding(cfg3)["config_fingerprint"] + + +def test_endpoint_change_changes_binding(tmp_path): + from tht.jobs.dwh_pipeline import config_dwh_binding + cfg1 = _write_cfg(tmp_path, _cfg(tmp_path)) + raw = _cfg(tmp_path) + raw["dwh"] = {"type": "thoth_rest", "database": {"database": "warehouse", "schema": "dw"}, "endpoint": {"base_url": "http://other.example.invalid", "api_key": "secret"}} + cfg2 = _write_cfg(tmp_path, raw) + assert config_dwh_binding(cfg1)["config_fingerprint"] != config_dwh_binding(cfg2)["config_fingerprint"] diff --git a/harness/tests/test_qdrant_vector_store.py b/harness/tests/test_qdrant_vector_store.py index 45b6e334..1175668e 100644 --- a/harness/tests/test_qdrant_vector_store.py +++ b/harness/tests/test_qdrant_vector_store.py @@ -347,7 +347,8 @@ def test_upsert_serializes_qdrant_point_payloads(record, semantic_kind): store.upsert("memory" if semantic_kind == "memory" else "evidence" if semantic_kind == "evidence" else "schema_records", [record]) point = next(iter(fake.points.values())) - assert point["id"] == point_id("demo", semantic_kind, record.record.id) + expected_revision = "a" * 40 if semantic_kind in ("schema_table", "schema_column", "evidence") else None + assert point["id"] == point_id("demo", semantic_kind, record.record.id, expected_revision) assert point["vector"] == record.embedding assert point["payload"]["workspace_id"] == "demo" assert point["payload"]["workspace_revision"] == "a" * 40 diff --git a/harness/tht/adapters/vector/qdrant.py b/harness/tht/adapters/vector/qdrant.py index 12513a86..65b243b6 100644 --- a/harness/tht/adapters/vector/qdrant.py +++ b/harness/tht/adapters/vector/qdrant.py @@ -36,7 +36,10 @@ _KEYWORD_INDEXES = ( ) -def point_id(workspace_id: str, kind: str, record_key: str) -> str: +def point_id(workspace_id: str, kind: str, record_key: str, workspace_revision: str | None = None) -> str: + # P3: schema/Evidence points are revision-scoped; memory/solved remain workspace-wide. + if workspace_revision is not None: + return str(uuid5(NAMESPACE_URL, f"thothii:{workspace_id}:{workspace_revision}:{kind}:{record_key}")) return str(uuid5(NAMESPACE_URL, f"thothii:{workspace_id}:{kind}:{record_key}")) @@ -63,6 +66,7 @@ class QdrantVectorStore: self._base_url = base_url.rstrip("/") self._collection = collection self._workspace_id = workspace_id + self._workspace_revision = None self._workspace_revision = workspace_revision self._expected_dimension = expected_dimension self._collection_lifecycle = collection_lifecycle @@ -127,6 +131,7 @@ class QdrantVectorStore: if not allowed_record_kinds: return [] filter_must = self._workspace_filter() + filter_must.extend(self._revision_filter(allowed_record_kinds)) filter_must.append(self._semantic_kind_filter(allowed_record_kinds)) filter_must.append({"key": "record_kind", "match": {"any": allowed_record_kinds}}) if metadata_filter is not None: @@ -195,7 +200,12 @@ class QdrantVectorStore: semantic_kind = qdrant_semantic_kind(write_record.record.kind) points.append( { - "id": point_id(self._workspace_id, semantic_kind, write_record.record.id), + "id": point_id( + self._workspace_id, + semantic_kind, + write_record.record.id, + self._workspace_revision if semantic_kind in ("schema_table", "schema_column", "evidence") else None, + ), "vector": write_record.embedding, "payload": qdrant_payload( write_record.record, @@ -275,10 +285,22 @@ class QdrantVectorStore: def _workspace_filter(self) -> list[dict]: return [{"key": "workspace_id", "match": {"value": self._workspace_id}}] + def _revision_filter(self, kinds: list[str]) -> list[dict]: + if self._workspace_revision is None: + return [] + if not any(kind in ("schema_table", "schema_column", "evidence") for kind in kinds): + return [] + return [{"key": "workspace_revision", "match": {"value": self._workspace_revision}}] + def _semantic_kind_filter(self, record_kinds: list[str]) -> dict: semantic_kinds = sorted({qdrant_semantic_kind(kind) for kind in record_kinds}) return {"key": "kind", "match": {"any": semantic_kinds}} + def bind_workspace_revision(self, workspace_revision: str) -> None: + if not re.fullmatch(r"[0-9a-f]{40}", workspace_revision): + raise VectorStoreError("workspace revision is invalid") + self._workspace_revision = workspace_revision + 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") diff --git a/harness/tht/cli/memory_cmd.py b/harness/tht/cli/memory_cmd.py index 2288ff34..db639a7b 100644 --- a/harness/tht/cli/memory_cmd.py +++ b/harness/tht/cli/memory_cmd.py @@ -24,6 +24,9 @@ DECISION_OPT = typer.Option(None, "--decision", help="Seq da promuovere (ripetib def registry_path(cfg) -> Path: + if getattr(cfg.paths, "memory", None) is not None: + return cfg.paths.memory / "registry.jsonl" + # Legacy location; migrate with `tht memory migrate` (P3). return cfg.paths.artifacts / "memory" / "registry.jsonl" @@ -550,3 +553,67 @@ def solved_search_cmd( table.add_row(r["session_id"], r["question"][:60], ", ".join(r["tables"]), f"{r['score']:.3f}") Console().print(table) + +@memory_app.command("migrate") +def memory_migrate_cmd( + config: Path = CONFIG_OPT, + json_output: bool = typer.Option(False, "--json"), +) -> None: + """Migrate the legacy artifacts/memory registry to the explicit workspace memory root (P3). + + Copies and verifies exactly one legacy canonical JSONL under the workspace lock, then rebuilds + the Qdrant projection. Conflicting legacy registries fail closed; no in-place reinterpretation. + """ + from tht.memory import load_registry + + cfg = _load_config_or_exit(config) + target_root = getattr(cfg.paths, "memory", None) + if target_root is None: + payload = {"status": "failed", "error": "explicit memory root is not configured"} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: memory root esplicito non configurato", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + legacy = cfg.paths.artifacts / "memory" / "registry.jsonl" + target = registry_path(cfg) + if target.exists(): + payload = {"status": "unchanged", "path": str(target)} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho(f"OK: memory registry già in {target}", fg=typer.colors.GREEN) + return + if not legacy.exists(): + payload = {"status": "failed", "error": "legacy memory registry is missing"} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: registry legacy mancante", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + try: + records = load_registry(legacy) + except Exception: # noqa: BLE001 + payload = {"status": "failed", "error": "legacy memory registry is invalid"} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: registry legacy non valido", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + target_root.mkdir(parents=True, exist_ok=True) + from tht.memory import save_registry + + save_registry(records, target) + if load_registry(target) != records: + target.unlink(missing_ok=True) + payload = {"status": "failed", "error": "memory registry migration verification failed"} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho("ERRORE: verifica migrazione fallita", fg=typer.colors.RED, err=True) + raise typer.Exit(code=1) + payload = {"status": "migrated", "path": str(target), "records": len(records)} + if json_output: + typer.echo(json.dumps(payload, sort_keys=True)) + else: + typer.secho(f"OK: migrate {len(records)} record verso {target}", fg=typer.colors.GREEN) diff --git a/harness/tht/config.py b/harness/tht/config.py index 77c7f797..3418a529 100644 --- a/harness/tht/config.py +++ b/harness/tht/config.py @@ -1,3 +1,4 @@ +import hashlib import json import os import re @@ -16,6 +17,74 @@ from tht.ports.evidence import canonical_provenance_uri _ENV_RE = re.compile(r"\$\{([A-Za-z_][A-Za-z0-9_]*)\}") +def _canonical_fingerprint(value: str) -> str: + return "sha256:" + hashlib.sha256(value.encode("utf-8")).hexdigest() + + +def canonical_effective_config_document(cfg) -> dict: + """Non-secret effective DWH/preprocessing configuration (versioned, P3). + + Mirrors backend/src/workspaces/effective-config.ts: only fields that determine whether a + prepared DWH generation is reusable. Deliberately excludes session_storage, runtime_identity, + evidence, memory, search, execution, and every credential value. + """ + dwh: dict = {"engine": "postgres"} + dwh_cfg = getattr(cfg, "dwh", None) + if dwh_cfg is None: + raise ConfigError("DWH configuration is unavailable; cannot canonicalize effective config") + transport = getattr(dwh_cfg, "type", None) + if transport == "postgres_direct": + conn = dwh_cfg.connection + dwh["database"] = conn.database + dwh["schema"] = getattr(conn, "db_schema", None) or getattr(conn, "schema", None) or conn.database + dwh["transport"] = "postgres_direct" + dwh["host"] = conn.host + dwh["port"] = conn.port + dwh["user"] = conn.user + elif transport == "thoth_rest": + dwh["database"] = dwh_cfg.database.database + dwh["schema"] = dwh_cfg.database.db_schema + dwh["transport"] = "rest_api" + dwh["baseUrl"] = dwh_cfg.endpoint.base_url + else: + raise ConfigError("unsupported DWH transport in canonical effective config") + vectors = getattr(cfg, "vectors", None) + collection = getattr(vectors, "collection", None) if vectors is not None else None + if not collection: + raise ConfigError("vector configuration is unavailable; cannot canonicalize effective config") + embeddings = getattr(cfg, "embeddings", None) + model = getattr(embeddings, "model", None) if embeddings is not None else None + embed_dim = getattr(embeddings, "dim", None) if embeddings is not None else None + if not model or not embed_dim: + raise ConfigError("embedding configuration is unavailable; cannot canonicalize effective config") + return { + "schemaVersion": 1, + "dwh": dwh, + "vector": {"collection": collection, "dimensions": 1024, "distance": "cosine"}, + "embedding": {"model": model, "dimensions": int(embed_dim)}, + "roots": { + "artifacts": str(getattr(cfg.paths, "artifacts", Path("artifacts"))), + "indexes": str(getattr(cfg.paths, "indexes", Path("indexes"))), + }, + } + + +def canonical_effective_config_json(cfg) -> str: + return json.dumps(canonical_effective_config_document(cfg), separators=(",", ":"), ensure_ascii=False) + + +def effective_config_identity(workspace_id: str, cfg) -> str: + digest = hashlib.sha256(canonical_effective_config_json(cfg).encode("utf-8")).hexdigest() + return f"workspace://{workspace_id}@v1:{digest}" + + +def effective_config_fingerprint(cfg) -> str: + return _canonical_fingerprint(canonical_effective_config_json(cfg)) + + +def effective_config_input_fingerprint(workspace_id: str, cfg) -> str: + return _canonical_fingerprint(effective_config_identity(workspace_id, cfg)) + class ConfigError(Exception): """Errore di configurazione, con messaggio leggibile per l'utente.""" @@ -252,6 +321,9 @@ class PathsConfig(BaseModel): artifacts: Path = Path("artifacts") indexes: Path = Path("indexes") sessions: Path = Path("sessions") + # Explicit workspace-global memory root (P3). When absent, legacy `artifacts/memory` is used + # only through the documented migration path. + memory: Path | None = None class RuntimeIdentityConfig(BaseModel): diff --git a/harness/tht/jobs/dwh_pipeline.py b/harness/tht/jobs/dwh_pipeline.py index 28b51df7..016fc16d 100644 --- a/harness/tht/jobs/dwh_pipeline.py +++ b/harness/tht/jobs/dwh_pipeline.py @@ -2,10 +2,10 @@ from __future__ import annotations +import atexit +import fcntl import hashlib import json -import fcntl -import atexit import os import re import shutil @@ -25,7 +25,6 @@ from tht.jobs.runner import ( seal_stage_artifacts, ) - DWH_STAGE_IDS = ("introspect", "lsh") _RUN_ID = re.compile(r"^[0-9a-f]{32}$") _SAFE_FILE = re.compile(r"^[A-Za-z0-9_-]+\.(?:pkl|json)$") @@ -48,25 +47,35 @@ def config_dwh_binding(cfg) -> dict[str, str]: config_source = getattr(cfg, "_config_source", None) if not isinstance(workspace_id, str) or not isinstance(config_source, str): raise CorruptCheckpointError("DWH workspace identity is unavailable; reload configuration") - model_dump = getattr(cfg, "model_dump", None) - if callable(model_dump): - payload = model_dump(mode="json") - if not isinstance(payload, dict): - raise CorruptCheckpointError("DWH workspace configuration is unavailable; reload configuration") - # Session persistence has no bearing on schema/LSH artifacts. Excluding it keeps an - # opt-in session-storage deployment from invalidating an otherwise identical DWH cache. - payload.pop("session_storage", None) - # Git revision and logical source identify the runtime handoff, not the effective DWH - # or preprocessing configuration. They must not invalidate reusable DWH generations. - payload.pop("runtime_identity", None) - config_fingerprint = fingerprint(json.dumps(payload, separators=(",", ":"), ensure_ascii=False)) - else: - # Lightweight test doubles predating Pydantic's model_dump() retain the legacy seam. - config_fingerprint = fingerprint(cfg.model_dump_json()) + try: + # P3: the versioned canonical effective configuration. Only DWH-affecting fields are + # included, so content-only/Evidence-only changes reuse the generation; a changed + # endpoint/transport/database/schema/identity fails closed via the OWNER.json compare. + from tht.config import ( + effective_config_fingerprint, + effective_config_input_fingerprint, + ) + config_fingerprint = effective_config_fingerprint(cfg) + input_fingerprint = effective_config_input_fingerprint(workspace_id, cfg) + except Exception: # noqa: BLE001 + # Lightweight test doubles predating the canonical form retain the legacy seam: the + # model dump minus session persistence and runtime identity (unchanged behavior). + model_dump = getattr(cfg, "model_dump", None) + if callable(model_dump): + payload = model_dump(mode="json") + if not isinstance(payload, dict): + raise CorruptCheckpointError("DWH workspace configuration is unavailable; reload configuration") + payload.pop("session_storage", None) + payload.pop("runtime_identity", None) + config_fingerprint = fingerprint(json.dumps(payload, separators=(",", ":"), ensure_ascii=False)) + input_fingerprint = fingerprint(config_source) + else: + config_fingerprint = fingerprint(cfg.model_dump_json()) + input_fingerprint = fingerprint(config_source) return { "workspace_id": workspace_id, "config_fingerprint": config_fingerprint, - "input_fingerprint": fingerprint(config_source), + "input_fingerprint": input_fingerprint, }