feat: operator schema accept command for curated FK review (P5)

This commit is contained in:
2026-08-13 05:06:23 +02:00
parent 5249798c03
commit 0459a6cd3e
9 changed files with 293 additions and 1 deletions
+9 -1
View File
@@ -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<string, unknown> {
"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;
}
@@ -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 };
}
@@ -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<WorkspaceOperationResult> {
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<WorkspaceOperationResult> {
const scope = await this.startRun(options.workspaceId, "index-schema", options.resumeRunId);
const semantic = await this.deps.semanticPreflight(scope.runtime.workspace);
@@ -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;
}
@@ -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);
});
@@ -21,6 +21,9 @@ thothctl --installation <absolute>/thothii-installation.yaml workspace schema ch
[--annotations <regular-file> --reviewed-candidates <sha256:hex>]
[--json]
thothctl --installation <absolute>/thothii-installation.yaml workspace schema accept
--workspace <id> --run <32hex> --yes [--json]
thothctl --installation <absolute>/thothii-installation.yaml workspace index-schema
--workspace <id> [--json]
@@ -57,6 +60,25 @@ thothctl --installation <absolute>/thothii-installation.yaml workspace vector re
mutation is performed.
```
## Curated FK annotations (P5)
- The canonical curated annotations file is `<workspace-id>/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/<id>/revisions/<commit>/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 <id> --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 <absolute>/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
+1
View File
@@ -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]
@@ -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 {
@@ -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)
}
}
}