From 0459a6cd3e4b58497f2a602c54b262d18c2e6bbc Mon Sep 17 00:00:00 2001 From: mptyl Date: Thu, 13 Aug 2026 05:06:23 +0200 Subject: [PATCH] feat: operator schema accept command for curated FK review (P5) --- backend/src/workspace-maintenance.ts | 10 ++- backend/src/workspaces/annotations-sync.ts | 35 +++++++++ .../src/workspaces/preprocessing-service.ts | 57 +++++++++++++++ backend/src/workspaces/preprocessing-state.ts | 3 + .../workspace-preprocessing-service.test.ts | 73 +++++++++++++++++++ docs/contracts/workspace-preprocessing-cli.md | 26 +++++++ tools/thothctl/cmd/thothctl/main.go | 1 + .../internal/workspaceops/operations.go | 54 ++++++++++++++ .../internal/workspaceops/operations_test.go | 35 +++++++++ 9 files changed, 293 insertions(+), 1 deletion(-) diff --git a/backend/src/workspace-maintenance.ts b/backend/src/workspace-maintenance.ts index 1383ed8f..fdf26285 100644 --- a/backend/src/workspace-maintenance.ts +++ b/backend/src/workspace-maintenance.ts @@ -18,7 +18,7 @@ export interface WorkspaceMaintenanceIo { writeStderr(value: string): void; } -type Command = "inspect" | "preprocess-dwh" | "schema-suggest-fks" | "schema-check" | "index-schema" | "preprocess-evidence" | "preprocess-run" | "vector-inspect" | "vector-rebuild"; +type Command = "inspect" | "preprocess-dwh" | "schema-suggest-fks" | "schema-check" | "schema-accept" | "index-schema" | "preprocess-evidence" | "preprocess-run" | "vector-inspect" | "vector-rebuild"; function failureResult( operation: string, @@ -66,6 +66,7 @@ function parseRequest(command: string, stdin: string): Record { "preprocess-dwh": ["schemaVersion", "workspaceId", "resumeRunId"], "schema-suggest-fks": ["schemaVersion", "workspaceId", "fromSql", "assume", "resumeRunId"], "schema-check": ["schemaVersion", "workspaceId", "annotationsYaml", "reviewedCandidatesDigest"], + "schema-accept": ["schemaVersion", "workspaceId", "runId", "yes"], "index-schema": ["schemaVersion", "workspaceId", "resumeRunId"], "preprocess-evidence": ["schemaVersion", "workspaceId", "dryRun", "resumeRunId"], "preprocess-run": ["schemaVersion", "workspaceId", "resumeRunId"], @@ -107,6 +108,12 @@ async function dispatch(command: Command, service: WorkspacePreprocessingService annotationsYaml: request.annotationsYaml as string | undefined, reviewedCandidatesDigest: request.reviewedCandidatesDigest as string | undefined, }); + case "schema-accept": + return await service.acceptSchema({ + workspaceId: request.workspaceId as string, + runId: request.runId as string, + yes: request.yes === true, + }); case "index-schema": return await service.indexSchema({ workspaceId: request.workspaceId as string, @@ -171,6 +178,7 @@ export async function runWorkspaceMaintenanceCli( || message === "invalid workspace id"; return command in { inspect: true, "preprocess-dwh": true, "schema-suggest-fks": true, "schema-check": true, + "schema-accept": true, "index-schema": true, "preprocess-evidence": true, "preprocess-run": true, } ? (requestError ? 2 : 1) : 2; } diff --git a/backend/src/workspaces/annotations-sync.ts b/backend/src/workspaces/annotations-sync.ts index 1592308e..9eb8e136 100644 --- a/backend/src/workspaces/annotations-sync.ts +++ b/backend/src/workspaces/annotations-sync.ts @@ -172,3 +172,38 @@ export function syncAnnotations(input: AnnotationsSyncInput): AnnotationsSyncRes writeAtomicFile(manifestPath, `${JSON.stringify(manifest)}\n`, 0o600); return { path, manifestPath, contentDigest }; } + +export interface SyncedAnnotations { + blobId: string; + contents: Buffer; + contentDigest: string; + manifestPath: string; +} + +/** Read the synced annotations and verify their ownership manifest; undefined when not yet synced. */ +export function readAnnotationsSync( + dataRoot: string, + workspaceId: string, + commit: string, +): SyncedAnnotations | undefined { + const directory = join(annotationsSyncRoot(dataRoot, workspaceId, commit), "mschema"); + const path = join(directory, "annotations.yaml"); + const manifestPath = join(directory, "annotations.ownership.json"); + let contents: Buffer; + try { + contents = readTrustedFile(path); + } catch (error) { + if ((error as NodeJS.ErrnoException).code === "ENOENT") return undefined; + throw error; + } + const manifest = parseManifest(readTrustedFile(manifestPath).toString("utf8")); + const contentDigest = sha256(contents); + if ( + manifest.workspace !== workspaceId + || manifest.commit !== commit + || manifest.contentDigest !== contentDigest + ) { + throw new Error("annotations ownership manifest does not match the pinned revision"); + } + return { blobId: manifest.blobId, contents, contentDigest, manifestPath }; +} diff --git a/backend/src/workspaces/preprocessing-service.ts b/backend/src/workspaces/preprocessing-service.ts index 4eac4ad1..4dc94704 100644 --- a/backend/src/workspaces/preprocessing-service.ts +++ b/backend/src/workspaces/preprocessing-service.ts @@ -10,6 +10,7 @@ import { type SessionInventoryRow, } from "./preprocessing-state.js"; import type { DeterministicRuntimeConfigLease } from "./runtime-config-lease.js"; +import { readAnnotationsSync } from "./annotations-sync.js"; export interface WorkspaceOperationResult { schemaVersion: 1; @@ -262,6 +263,62 @@ export class WorkspacePreprocessingService { }); } + async acceptSchema(options: { workspaceId: string; runId: string; yes?: boolean }): Promise { + const runtime = await this.deps.acquireActiveRuntime(options.workspaceId); + const state = this.state(runtime.workspaceId); + if (options.yes !== true) { + return baseResult(runtime, "schema accept", "failed", "annotation_invalid", { + runId: options.runId, + warnings: ["accept requires --yes"], + }); + } + if (!/^[0-9a-f]{32}$/.test(options.runId)) { + return baseResult(runtime, "schema accept", "failed", "annotation_invalid"); + } + const candidate = state.readFkCandidates(options.runId); + if (candidate === undefined) { + return baseResult(runtime, "schema accept", "failed", "annotation_invalid", { + runId: options.runId, + warnings: ["candidate run is unavailable"], + }); + } + const synced = readAnnotationsSync(this.deps.dataRoot, runtime.workspaceId, runtime.workspaceRevision); + if (synced === undefined || synced.contents.toString("utf8").trim() === "") { + return baseResult(runtime, "schema accept", "failed", "annotation_invalid", { + runId: options.runId, + warnings: ["curated annotations are not synchronized"], + }); + } + // The harness parser validates the curated blob against the physical schema; the recorded + // candidate digest must round-trip and the blob digest must match the synced destination. + const payload = await this.runJsonStage(runtime, [ + "schema", "check", "--reviewed-candidates", candidate.digest, "--json", "-c", "/dev/fd/3", + ]); + if (payload.annotations_digest !== synced.contentDigest + || payload.reviewed_candidates_digest !== candidate.digest + || Number(payload.orphan_count ?? 0) !== 0) { + return baseResult(runtime, "schema accept", "failed", "annotation_invalid", { runId: options.runId }); + } + const review = state.writeFkReview(options.runId, { + reviewedCandidatesDigest: candidate.digest, + annotationsDigest: synced.contentDigest, + workspaceRevision: runtime.workspaceRevision, + blobId: synced.blobId, + }); + const job = state.readJob(options.runId); + job.reviewDigest = review.digest; + if (!job.completedStages.includes("fk_review")) job.completedStages.push("fk_review"); + state.writeJob(job); + return baseResult(runtime, "schema accept", "succeeded", "ok", { + runId: options.runId, + completedStages: [...job.completedStages], + artifactIdentities: [ + { kind: "fk_review", digest: review.digest }, + { kind: "annotations", digest: synced.contentDigest }, + ], + }); + } + async indexSchema(options: { workspaceId: string; resumeRunId?: string }): Promise { const scope = await this.startRun(options.workspaceId, "index-schema", options.resumeRunId); const semantic = await this.deps.semanticPreflight(scope.runtime.workspace); diff --git a/backend/src/workspaces/preprocessing-state.ts b/backend/src/workspaces/preprocessing-state.ts index d7d52f13..f950c4ef 100644 --- a/backend/src/workspaces/preprocessing-state.ts +++ b/backend/src/workspaces/preprocessing-state.ts @@ -68,6 +68,8 @@ export interface FkReviewRecord { reviewedCandidatesDigest: string; annotationsDigest: string; workspaceRevision: string; + /** Curated Git blob id accepted at review time (P5); absent for legacy host-file reviews. */ + blobId?: string; } function sha256(value: string | Buffer): string { @@ -171,6 +173,7 @@ function decodeReview(value: unknown): FkReviewRecord { typeof record.reviewedCandidatesDigest !== "string" || typeof record.annotationsDigest !== "string" || typeof record.workspaceRevision !== "string" + || (record.blobId !== undefined && typeof record.blobId !== "string") ) throw new Error("preprocessing review state is invalid"); return record as unknown as FkReviewRecord; } diff --git a/backend/test/workspace-preprocessing-service.test.ts b/backend/test/workspace-preprocessing-service.test.ts index 9a7ef4c3..38328887 100644 --- a/backend/test/workspace-preprocessing-service.test.ts +++ b/backend/test/workspace-preprocessing-service.test.ts @@ -4,6 +4,7 @@ import { join } from "node:path"; import { afterEach, expect, test, vi } from "vitest"; import { parseWorkspaceYaml } from "../src/workspaces/schema.js"; import { PreprocessingStateStore } from "../src/workspaces/preprocessing-state.js"; +import { syncAnnotations } from "../src/workspaces/annotations-sync.js"; import { WorkspacePreprocessingService, type ChildProcessRequest, @@ -349,3 +350,75 @@ test("full runs follow the explicit order and finish unchanged when no Evidence "vector index-schema", ]); }); + + +test("schema accept validates the synced Git blob and records the review", async () => { + const f = fixture(); + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + candidate_count: 1, + candidate_digest: "sha256:" + "e".repeat(64), + candidate_yaml: "tables: {}\n", + }), + stderr: "", + }); + const suggest = await f.service.suggestFks({ workspaceId: "psd-clinical" }); + const runId = suggest.runId!; + const candidateDigest = suggest.artifactIdentities![0]!.digest; + const synced = syncAnnotations({ + dataRoot: f.dataRoot, + workspaceId: "psd-clinical", + commit: "a".repeat(40), + blobId: "b".repeat(40), + contents: Buffer.from("tables: {}\n"), + }); + + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + orphan_count: 0, + annotations_digest: synced.contentDigest, + reviewed_candidates_digest: candidateDigest, + }), + stderr: "", + }); + + const result = await f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: true }); + expect(result).toMatchObject({ status: "succeeded", code: "ok", operation: "schema accept" }); + const state = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); + expect(state.readFkReview(runId)).toMatchObject({ + reviewedCandidatesDigest: candidateDigest, + annotationsDigest: synced.contentDigest, + workspaceRevision: "a".repeat(40), + blobId: "b".repeat(40), + }); +}); + +test("schema accept fails closed without --yes, for an unknown run, or with no synced annotations", async () => { + const f = fixture(); + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ + status: "succeeded", + candidate_count: 1, + candidate_digest: "sha256:" + "e".repeat(64), + candidate_yaml: "tables: {}\n", + }), + stderr: "", + }); + const suggest = await f.service.suggestFks({ workspaceId: "psd-clinical" }); + const runId = suggest.runId!; + + await expect(f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: false })) + .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); + + await expect(f.service.acceptSchema({ workspaceId: "psd-clinical", runId: "e".repeat(32), yes: true })) + .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); + + await expect(f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: true })) + .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); + expect(f.runChild).toHaveBeenCalledTimes(1); +}); diff --git a/docs/contracts/workspace-preprocessing-cli.md b/docs/contracts/workspace-preprocessing-cli.md index 0763608e..204300a9 100644 --- a/docs/contracts/workspace-preprocessing-cli.md +++ b/docs/contracts/workspace-preprocessing-cli.md @@ -21,6 +21,9 @@ thothctl --installation /thothii-installation.yaml workspace schema ch [--annotations --reviewed-candidates ] [--json] +thothctl --installation /thothii-installation.yaml workspace schema accept + --workspace --run <32hex> --yes [--json] + thothctl --installation /thothii-installation.yaml workspace index-schema --workspace [--json] @@ -57,6 +60,25 @@ thothctl --installation /thothii-installation.yaml workspace vector re mutation is performed. ``` +## Curated FK annotations (P5) + +- The canonical curated annotations file is `/schema/annotations.yaml`, a regular + Git blob at the same commit as the descriptor. Absence is compatible (empty canonical set + + warning); symlinks, trees/gitlinks, oversized (>16 MiB), non-UTF-8, and malformed objects are + refused at activation. +- Activation synchronizes the blob to the immutable revision-qualified root + `/data/sessions//revisions//artifacts/mschema/annotations.yaml` with a restrictive + mode and an adjacent ownership manifest `{ workspace, commit, blobId, contentDigest, + destination }`. Pinned runtimes resolve annotations from `paths.annotations_root`. +- `workspace schema accept --run --yes` is the only human FK review primitive: after + commit/push/pull, it reads the current synced blob, validates it with the harness parser against + the physical schema and the recorded candidate digest, and records + `{ reviewedCandidatesDigest, annotationsDigest, workspaceRevision, blobId }`. `--yes` is + required; an empty file, an unknown run, a malformed blob, or a non-matching candidate fails + closed without recording a review. `schema check` alone is not evidence of human review. +- `preprocess run` continues only with the exact accepted blob digest and a compatible reusable + DWH binding; otherwise it records a new review checkpoint. + ## Validation - `--installation` is mandatory and absolute. @@ -73,6 +95,10 @@ thothctl --installation /thothii-installation.yaml workspace vector re - `--annotations` and `--reviewed-candidates` are all-or-nothing; - annotations must be UTF-8, canonical, non-symlink, max 16 MiB; - `--reviewed-candidates` must match `sha256:<64 lowercase hex>`. +- `schema accept` + - `--run` is mandatory and must be 32 lowercase hex characters; + - `--yes` is mandatory and may be supplied once; + - `--annotations`/`--reviewed-candidates`/`--from-sql`/`--assume` are not accepted. - Unknown flags, passthrough separators, and shell fragments are rejected before Docker runs. ## Container boundary diff --git a/tools/thothctl/cmd/thothctl/main.go b/tools/thothctl/cmd/thothctl/main.go index 0a364ebb..509c0eb6 100644 --- a/tools/thothctl/cmd/thothctl/main.go +++ b/tools/thothctl/cmd/thothctl/main.go @@ -56,6 +56,7 @@ Commands: workspace preprocess dwh --workspace ID [--resume RUN] [--json] workspace schema suggest-fks --workspace ID [--from-sql FILE]... [--assume COLUMN=TABLE]... [--output FILE] [--json] workspace schema check --workspace ID [--annotations FILE --reviewed-candidates sha256:HEX] [--json] + workspace schema accept --workspace ID --run RUN --yes [--json] workspace index-schema --workspace ID [--json] workspace preprocess evidence --workspace ID [--dry-run] [--resume RUN] [--json] workspace preprocess run --workspace ID [--resume RUN] [--json] diff --git a/tools/thothctl/internal/workspaceops/operations.go b/tools/thothctl/internal/workspaceops/operations.go index 8c0ba0e1..7ff6f848 100644 --- a/tools/thothctl/internal/workspaceops/operations.go +++ b/tools/thothctl/internal/workspaceops/operations.go @@ -71,6 +71,12 @@ type CheckSchemaRequest struct { ReviewedCandidates string } +type AcceptSchemaRequest struct { + baseRequest + Run string + Yes bool +} + type IndexSchemaRequest struct{ baseRequest } type EvidenceRequest struct { @@ -88,6 +94,7 @@ func (InspectRequest) workspaceRequest() {} func (DwhRequest) workspaceRequest() {} func (SuggestFksRequest) workspaceRequest() {} func (CheckSchemaRequest) workspaceRequest() {} +func (AcceptSchemaRequest) workspaceRequest() {} func (IndexSchemaRequest) workspaceRequest() {} func (EvidenceRequest) workspaceRequest() {} func (RunRequest) workspaceRequest() {} @@ -98,6 +105,7 @@ func (SuggestFksRequest) operatorCommand() string { return "schema-suggest-fks" func (CheckSchemaRequest) operatorCommand() string { return "schema-check" } +func (AcceptSchemaRequest) operatorCommand() string { return "schema-accept" } func (IndexSchemaRequest) operatorCommand() string { return "index-schema" } func (EvidenceRequest) operatorCommand() string { return "preprocess-evidence" } func (RunRequest) operatorCommand() string { return "preprocess-run" } @@ -144,6 +152,10 @@ func (r IndexSchemaRequest) stdinEnvelope() (requestEnvelope, error) { return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace}, nil } +func (r AcceptSchemaRequest) stdinEnvelope() (requestEnvelope, error) { + return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, RunID: r.Run, Yes: r.Yes}, nil +} + func (r EvidenceRequest) stdinEnvelope() (requestEnvelope, error) { return requestEnvelope{SchemaVersion: 1, WorkspaceID: r.Workspace, Resume: r.Resume, DryRun: r.DryRun}, nil } @@ -164,6 +176,8 @@ type requestEnvelope struct { Collection string `json:"collection,omitempty"` Confirm string `json:"confirm,omitempty"` Destroy bool `json:"destroy,omitempty"` + RunID string `json:"runId,omitempty"` + Yes bool `json:"yes,omitempty"` } type inputFile struct { @@ -464,6 +478,8 @@ func parseSchema(args []string) (Request, error) { return parseSuggestFks(args[1:]) case "check": return parseSchemaCheck(args[1:]) + case "accept": + return parseSchemaAccept(args[1:]) default: return nil, fmt.Errorf("unknown workspace schema command %q", args[0]) } @@ -588,6 +604,44 @@ func parseSchemaCheck(args []string) (CheckSchemaRequest, error) { return request, nil } +func parseSchemaAccept(args []string) (AcceptSchemaRequest, error) { + request := AcceptSchemaRequest{} + base, seen, err := parseSharedFlags(args, map[string]func(string) error{ + "--run": func(value string) error { + if request.Run != "" { + return errors.New("--run may be supplied once") + } + if !runIDPattern.MatchString(value) { + return errors.New("--run must be 32 lowercase hex characters") + } + request.Run = value + return nil + }, + }, map[string]func() error{ + "--yes": func() error { + if request.Yes { + return errors.New("--yes may be supplied once") + } + request.Yes = true + return nil + }, + }) + if err != nil { + return AcceptSchemaRequest{}, err + } + if !seen.workspace { + return AcceptSchemaRequest{}, errors.New("--workspace is required") + } + if request.Run == "" { + return AcceptSchemaRequest{}, errors.New("--run is required") + } + if !request.Yes { + return AcceptSchemaRequest{}, errors.New("--yes is required") + } + request.baseRequest = base + return request, nil +} + func parseBaseFlags(args []string, allowDryRun bool) (baseRequest, error) { base, seen, err := parseSharedFlags(args, nil, nil) if err != nil { diff --git a/tools/thothctl/internal/workspaceops/operations_test.go b/tools/thothctl/internal/workspaceops/operations_test.go index ed2a00e3..d1fc5d5a 100644 --- a/tools/thothctl/internal/workspaceops/operations_test.go +++ b/tools/thothctl/internal/workspaceops/operations_test.go @@ -268,3 +268,38 @@ func TestParseVectorUnknownSubcommand(t *testing.T) { t.Fatal("expected error for unknown vector command") } } + +func TestParseSchemaAccept(t *testing.T) { + req, err := Parse([]string{"schema", "accept", "--workspace", "psd", "--run", strings.Repeat("d", 32), "--yes"}) + if err != nil { + t.Fatalf("parse: %v", err) + } + r, ok := req.(AcceptSchemaRequest) + if !ok { + t.Fatalf("got %T", req) + } + if r.Workspace != "psd" || r.Run != strings.Repeat("d", 32) || !r.Yes { + t.Fatalf("unexpected request: %+v", r) + } + env, err := req.stdinEnvelope() + if err != nil { + t.Fatalf("envelope: %v", err) + } + if env.WorkspaceID != "psd" || env.RunID != strings.Repeat("d", 32) || !env.Yes { + t.Fatalf("unexpected envelope: %+v", env) + } +} + +func TestParseSchemaAcceptRequiresRunAndYes(t *testing.T) { + for _, args := range [][]string{ + {"schema", "accept", "--workspace", "psd"}, + {"schema", "accept", "--workspace", "psd", "--yes"}, + {"schema", "accept", "--workspace", "psd", "--run", strings.Repeat("d", 32)}, + {"schema", "accept", "--workspace", "psd", "--run", "not-hex", "--yes"}, + {"schema", "accept", "--workspace", "psd", "--run", strings.Repeat("d", 32), "--yes", "--run", strings.Repeat("e", 32)}, + } { + if _, err := Parse(args); err == nil { + t.Fatalf("expected error for %v", args) + } + } +}