diff --git a/backend/src/workspaces/preprocessing-service.ts b/backend/src/workspaces/preprocessing-service.ts index 4dc94704..094ba766 100644 --- a/backend/src/workspaces/preprocessing-service.ts +++ b/backend/src/workspaces/preprocessing-service.ts @@ -244,23 +244,9 @@ export class WorkspacePreprocessingService { if (request.reviewed_candidates_digest !== reviewedCandidatesDigest || typeof request.annotations_digest !== "string") { return baseResult(runtime, "schema check", "failed", "annotation_invalid", { runId }); } - const annotationsDigest = request.annotations_digest; - const review = state.writeFkReview(runId, { - reviewedCandidatesDigest, - annotationsDigest, - workspaceRevision: runtime.workspaceRevision, - }); - const job = state.readJob(runId); - if (!job.completedStages.includes("fk_review")) { - job.reviewDigest = review.digest; - job.completedStages.push("fk_review"); - state.writeJob(job); - } - return baseResult(runtime, "schema check", "succeeded", "ok", { - runId, - completedStages: [...job.completedStages], - artifactIdentities: [{ kind: "fk_review", digest: review.digest }], - }); + // P5 supersedes the host-file FK review: schema check is read-only validation and never + // records a review. Only `schema accept` records a human review for the curated Git blob. + return baseResult(runtime, "schema check", "succeeded", "ok", { runId }); } async acceptSchema(options: { workspaceId: string; runId: string; yes?: boolean }): Promise { @@ -392,10 +378,12 @@ export class WorkspacePreprocessingService { } const candidate = this.state(scope.runtime.workspaceId).readFkCandidates(scope.job.runId); if (candidate && !scope.job.completedStages.includes("fk_review")) { - // The candidate content digest is authoritative: a human review accepted for ANY run - // carrying the exact same candidate digest counts as the review checkpoint for this run. + // P5: continuation requires a review accepted for this candidate whose accepted blob digest + // equals the current revision's synced annotations. A revision change (or a missing curated + // blob) therefore records a new review checkpoint instead of silently reusing the old one. const accepted = this.findAcceptedReviewForDigest(scope.runtime.workspaceId, candidate.digest); - if (!accepted) { + const currentDigest = this.currentAnnotationsDigest(scope.runtime); + if (accepted === undefined || currentDigest === undefined || accepted.annotationsDigest !== currentDigest) { return baseResult(scope.runtime, "preprocess run", "blocked", "manual_review_required", { runId: scope.job.runId, childRuns: { ...scope.job.childRuns }, @@ -562,6 +550,10 @@ export class WorkspacePreprocessingService { } } + private currentAnnotationsDigest(runtime: ActiveRuntime): string | undefined { + return readAnnotationsSync(this.deps.dataRoot, runtime.workspaceId, runtime.workspaceRevision)?.contentDigest; + } + private findAcceptedReviewForDigest( workspaceId: string, digestValue: string, diff --git a/backend/test/workspace-preprocessing-service.test.ts b/backend/test/workspace-preprocessing-service.test.ts index 38328887..18c2ff11 100644 --- a/backend/test/workspace-preprocessing-service.test.ts +++ b/backend/test/workspace-preprocessing-service.test.ts @@ -220,7 +220,7 @@ test("schema suggest-fks publishes a candidate artifact and blocks full runs for expect(f.runChild.mock.calls.map(([request]) => (request as ChildProcessRequest).argv[0])).toEqual(["preprocess", "schema"]); }); -test("schema check requires the exact candidate digest, stages annotations via temp file, and persists the review", async () => { +test("schema check requires the exact candidate digest and stages annotations via a temp file without recording a review", async () => { const f = fixture(); f.runChild.mockResolvedValueOnce({ exitCode: 0, @@ -269,7 +269,7 @@ test("schema check requires the exact candidate digest, stages annotations via t expect(checked).toMatchObject({ status: "succeeded", code: "ok" }); expect(() => readFileSync(stagedPath, "utf8")).toThrow(); const state = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); - expect(state.readFkReview(suggest.runId!)?.reviewedCandidatesDigest).toBe(reviewedDigest); + expect(state.readFkReview(suggest.runId!)).toBeUndefined(); }); test("index schema fails closed when semantic preflight refuses the collection", async () => { @@ -422,3 +422,67 @@ test("schema accept fails closed without --yes, for an unknown run, or with no s .resolves.toMatchObject({ status: "failed", code: "annotation_invalid" }); expect(f.runChild).toHaveBeenCalledTimes(1); }); + +test("full runs continue after schema accept only when the accepted blob matches the current revision", async () => { + const f = fixture(); + const revision = "a".repeat(40); + f.runChild + .mockResolvedValueOnce({ exitCode: 0, stdout: JSON.stringify({ status: "succeeded", run_id: "d".repeat(32) }), stderr: "" }) + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", candidate_count: 1, candidate_digest: "sha256:" + "e".repeat(64), candidate_yaml: "tables: {}\n" }), + stderr: "", + }); + + const blocked = await f.service.run({ workspaceId: "psd-clinical" }); + expect(blocked).toMatchObject({ status: "blocked", code: "manual_review_required" }); + const runId = blocked.runId!; + const state = new PreprocessingStateStore({ dataRoot: f.dataRoot, workspaceId: "psd-clinical" }); + const candidateDigest = state.readFkCandidates(runId)!.digest; + + const synced = syncAnnotations({ + dataRoot: f.dataRoot, + workspaceId: "psd-clinical", + commit: revision, + blobId: "b".repeat(40), + contents: Buffer.from("tables: {}\n"), + }); + + // Accept writes the review; the resume then passes the gate and reaches schema indexing. + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", orphan_count: 0, annotations_digest: synced.contentDigest, reviewed_candidates_digest: candidateDigest }), + stderr: "", + }); + await f.service.acceptSchema({ workspaceId: "psd-clinical", runId, yes: true }); + + f.runChild.mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", counts: { added: 1, updated: 0, deleted: 0, unchanged: 0 } }), + stderr: "", + }); + const resumed = await f.service.run({ workspaceId: "psd-clinical", resumeRunId: runId }); + expect(resumed).toMatchObject({ status: "succeeded", code: "ok" }); + + // A review whose accepted blob digest no longer matches the current revision stays blocked. + const second = fixture(); + second.runChild + .mockResolvedValueOnce({ exitCode: 0, stdout: JSON.stringify({ status: "succeeded", run_id: "d".repeat(32) }), stderr: "" }) + .mockResolvedValueOnce({ + exitCode: 0, + stdout: JSON.stringify({ status: "succeeded", candidate_count: 1, candidate_digest: "sha256:" + "e".repeat(64), candidate_yaml: "tables: {}\n" }), + stderr: "", + }); + const secondBlocked = await second.service.run({ workspaceId: "psd-clinical" }); + const secondRunId = secondBlocked.runId!; + const secondState = new PreprocessingStateStore({ dataRoot: second.dataRoot, workspaceId: "psd-clinical" }); + const secondCandidate = secondState.readFkCandidates(secondRunId)!.digest; + secondState.writeFkReview(secondRunId, { + reviewedCandidatesDigest: secondCandidate, + annotationsDigest: "sha256:" + "0".repeat(64), + workspaceRevision: revision, + blobId: "b".repeat(40), + }); + const stillBlocked = await second.service.run({ workspaceId: "psd-clinical", resumeRunId: secondRunId }); + expect(stillBlocked).toMatchObject({ status: "blocked", code: "manual_review_required" }); +});