From 6de74d103926d9090f056aeebe7be393ec381ea1 Mon Sep 17 00:00:00 2001 From: Anonymous Authors Date: Sat, 25 Jul 2026 05:59:42 -0500 Subject: Align five-stage pipeline with manuscript --- src/gap_pipeline/e2e.py | 3 +- src/gap_pipeline/kernel_models.py | 54 ++++++++----- src/gap_pipeline/kernel_prompts.py | 153 ++++++++++++++++++------------------- src/gap_pipeline/paper_pipeline.py | 90 ++++++++++++++-------- src/gap_pipeline/pipeline.py | 2 +- src/gap_pipeline/release.py | 21 ++--- 6 files changed, 185 insertions(+), 138 deletions(-) (limited to 'src/gap_pipeline') diff --git a/src/gap_pipeline/e2e.py b/src/gap_pipeline/e2e.py index 66101e7..e0859dc 100644 --- a/src/gap_pipeline/e2e.py +++ b/src/gap_pipeline/e2e.py @@ -26,7 +26,6 @@ def _review_accept() -> dict[str, str]: return { "verdict": "accept", "step_by_step_check": "n1 passes; n2 passes; n3 passes", - "replacement_check": "s1 satisfies its positivity guard", "blocking_issues": "", "patch_suggestion": "", } @@ -204,7 +203,7 @@ async def run_offline_smoke(work_root: Path) -> dict[str, Any]: f"{item_id}.stage3.replacement": { "changes": [ { - "slot_id": "s1", + "slot_id": "slot1", "source_node_id": "n1", "description": "positive square root at equality", "original_value": "1", diff --git a/src/gap_pipeline/kernel_models.py b/src/gap_pipeline/kernel_models.py index 32abf39..f835473 100644 --- a/src/gap_pipeline/kernel_models.py +++ b/src/gap_pipeline/kernel_models.py @@ -38,6 +38,11 @@ class ProofDAG(StrictModel): def node_ids(self) -> list[str]: return [node.node_id for node in self.nodes] + def leaf_node_ids(self) -> list[str]: + """Return source leaves: nodes with no prerequisite dependencies.""" + + return [node.node_id for node in self.nodes if not node.dependencies] + @model_validator(mode="after") def validate_graph(self) -> "ProofDAG": node_ids = self.node_ids() @@ -113,6 +118,8 @@ class ReplacementChange(StrictModel): ] if any(not value.strip() for value in text_fields): raise ValueError("replacement fields must be non-empty") + if not re.fullmatch(r"slot[1-9][0-9]*", self.slot_id): + raise ValueError("replacement slot IDs must be slot1, slot2, ...") if self.original_value.strip() == self.replacement_value.strip(): raise ValueError("replacement must differ from the original value") return self @@ -140,6 +147,34 @@ class ReplacementPlan(StrictModel): } if unknown: raise ValueError(f"replacement plan references unknown nodes {sorted(unknown)}") + leaves = set(dag.leaf_node_ids()) + non_leaf = { + change.source_node_id + for change in self.changes + if change.source_node_id not in leaves + } + if non_leaf: + raise ValueError( + "replacement plan must target source leaf nodes; " + f"received {sorted(non_leaf)}" + ) + return self + + def validate_repair_of( + self, + previous: "ReplacementPlan", + ) -> "ReplacementPlan": + current_targets = [ + (change.slot_id, change.source_node_id) for change in self.changes + ] + previous_targets = [ + (change.slot_id, change.source_node_id) + for change in previous.changes + ] + if current_targets != previous_targets: + raise ValueError( + "repair must preserve replacement slot IDs and source leaf nodes" + ) return self @@ -234,13 +269,11 @@ class CandidateBundle(StrictModel): class JudgeVerdict(StrictModel): verdict: Literal["accept", "reject"] step_by_step_check: str - replacement_check: str blocking_issues: str = "" patch_suggestion: str = "" @field_validator( "step_by_step_check", - "replacement_check", "blocking_issues", "patch_suggestion", mode="before", @@ -279,7 +312,6 @@ class JudgeVerdict(StrictModel): def validate_coverage( self, dag: ProofDAG, - replacement_plan: ReplacementPlan, ) -> "JudgeVerdict": missing_nodes = [ node_id @@ -290,20 +322,8 @@ class JudgeVerdict(StrictModel): ) is None ] - missing_slots = [ - change.slot_id - for change in replacement_plan.changes - if re.search( - rf"(?>> - -OFFICIAL SOLUTION: -<<<{solution}>>> - -SOURCE PROOF DAG: -{dag} - -METHOD PLAN: -{methods} - -REPLACEMENT PLAN: -{replacements} - -DIFFUSED PROOF: -{diffused} - -RENDERED VARIANT: -{variant}""" +JUDGE_SYSTEM = JUDGE_SYSTEM_PROMPT def _dump(value: object) -> str: @@ -201,18 +165,27 @@ def replacement_user( dag: ProofDAG, methods: MethodPlan, *, + previous_replacements: ReplacementPlan | None = None, feedback: str = "", ) -> str: - feedback_block = ( - f"Previous verification feedback to address:\n{feedback}" - if feedback - else "This is the initial replacement proposal." + repair_context = ( + "This is the initial replacement proposal." + if previous_replacements is None + else ( + "This is a repair pass. Preserve the previous slot IDs, source leaf " + "nodes, and accepted changes unless the verification feedback " + "identifies them as the blocking issue. Make only the smallest " + "necessary correction.\n\nPREVIOUS REPLACEMENT PLAN:\n" + f"{_dump(previous_replacements)}\n\nVERIFICATION FEEDBACK:\n" + f"{feedback or 'No additional textual feedback was supplied.'}" + ) ) return REPLACEMENT_USER.format( - feedback=feedback_block, + repair_context=repair_context, question=item.problem, solution=item.solution, dag=_dump(dag), + leaf_node_ids=json.dumps(dag.leaf_node_ids(), ensure_ascii=False), methods=_dump(methods), ) @@ -222,8 +195,23 @@ def diffusion_user( dag: ProofDAG, methods: MethodPlan, replacements: ReplacementPlan, + *, + previous_diffused: DiffusedProof | None = None, + feedback: str = "", ) -> str: + repair_context = ( + "This is the initial DAG diffusion." + if previous_diffused is None + else ( + "This is a repair pass. Apply only corrections required by the " + "verification feedback; preserve every unaffected node.\n\n" + f"PREVIOUS DIFFUSED PROOF:\n{_dump(previous_diffused)}\n\n" + f"VERIFICATION FEEDBACK:\n" + f"{feedback or 'No additional textual feedback was supplied.'}" + ) + ) return DIFFUSION_USER.format( + repair_context=repair_context, question=item.problem, dag=_dump(dag), methods=_dump(methods), @@ -234,8 +222,23 @@ def diffusion_user( def render_user( replacements: ReplacementPlan, diffused: DiffusedProof, + *, + previous_variant: object | None = None, + feedback: str = "", ) -> str: + repair_context = ( + "This is the initial rendering." + if previous_variant is None + else ( + "This is a repair pass. Preserve all unaffected wording and apply " + "only corrections required by the verification feedback.\n\n" + f"PREVIOUS RENDERED VARIANT:\n{_dump(previous_variant)}\n\n" + f"VERIFICATION FEEDBACK:\n" + f"{feedback or 'No additional textual feedback was supplied.'}" + ) + ) return RENDER_USER.format( + repair_context=repair_context, node_order=json.dumps( [node.node_id for node in diffused.nodes], ensure_ascii=False, @@ -247,30 +250,26 @@ def render_user( def judge_user( item: CanonicalItem, - dag: ProofDAG, methods: MethodPlan, replacements: ReplacementPlan, - diffused: DiffusedProof, variant: object, - *, - format_feedback: str = "", ) -> str: - return JUDGE_USER.format( - required_node_ids=json.dumps(dag.node_ids(), ensure_ascii=False), - required_slot_ids=json.dumps( - [change.slot_id for change in replacements.changes], + variant_payload = ( + variant.model_dump(mode="json") + if hasattr(variant, "model_dump") + else variant + ) + return JUDGE_USER_TEMPLATE.format( + original_problem=item.problem, + original_solution=item.solution, + method_labels=json.dumps( + [ + f"{node.node_id}: {node.method_label}" + for node in methods.nodes + ], ensure_ascii=False, ), - format_feedback=( - f"Previous report-format error: {format_feedback}" - if format_feedback - else "" - ), - question=item.problem, - solution=item.solution, - dag=_dump(dag), - methods=_dump(methods), - replacements=_dump(replacements), - diffused=_dump(diffused), - variant=_dump(variant), + slot_replacement=_dump(replacements), + candidate_problem=str(variant_payload["question"]), + candidate_proof=str(variant_payload["solution"]), ) diff --git a/src/gap_pipeline/paper_pipeline.py b/src/gap_pipeline/paper_pipeline.py index 84f34bf..c90f3f6 100644 --- a/src/gap_pipeline/paper_pipeline.py +++ b/src/gap_pipeline/paper_pipeline.py @@ -133,6 +133,7 @@ class PaperKernelPipeline: methods: MethodPlan, *, version: int, + previous_replacements: ReplacementPlan | None = None, feedback: str = "", ) -> ReplacementPlan: request_id = f"{item.item_id}.stage3.replacement.v{version:02d}" @@ -145,10 +146,13 @@ class PaperKernelPipeline: item, dag, methods, + previous_replacements=previous_replacements, feedback=feedback, ), ) ).validate_against(dag) + if previous_replacements is not None: + replacements.validate_repair_of(previous_replacements) self.store.write_stage( f"03_replacement_v{version:02d}", replacements, @@ -164,6 +168,8 @@ class PaperKernelPipeline: replacements: ReplacementPlan, *, version: int, + previous_diffused: DiffusedProof | None = None, + feedback: str = "", ) -> DiffusedProof: request_id = f"{item.item_id}.stage4.diffusion.v{version:02d}" diffused = DiffusedProof.model_validate( @@ -171,7 +177,14 @@ class PaperKernelPipeline: self.proposer, request_id=request_id, system_prompt=DIFFUSION_SYSTEM, - user_prompt=diffusion_user(item, dag, methods, replacements), + user_prompt=diffusion_user( + item, + dag, + methods, + replacements, + previous_diffused=previous_diffused, + feedback=feedback, + ), ) ).validate_against(dag, methods) self.store.write_stage( @@ -188,6 +201,8 @@ class PaperKernelPipeline: diffused: DiffusedProof, *, version: int, + previous_variant: RenderedVariant | None = None, + feedback: str = "", ) -> RenderedVariant: request_id = f"{self.store.item_id}.stage5.render.v{version:02d}" variant = RenderedVariant.model_validate( @@ -195,7 +210,12 @@ class PaperKernelPipeline: self.proposer, request_id=request_id, system_prompt=RENDER_SYSTEM, - user_prompt=render_user(replacements, diffused), + user_prompt=render_user( + replacements, + diffused, + previous_variant=previous_variant, + feedback=feedback, + ), ) ).validate_against(dag, diffused) self.store.write_stage( @@ -212,6 +232,7 @@ class PaperKernelPipeline: methods: MethodPlan, *, version: int, + previous_bundle: CandidateBundle | None = None, feedback: str = "", ) -> CandidateBundle: replacements = await self.generate_replacements( @@ -219,6 +240,11 @@ class PaperKernelPipeline: dag, methods, version=version, + previous_replacements=( + previous_bundle.replacement_plan + if previous_bundle is not None + else None + ), feedback=feedback, ) diffused = await self.diffuse_dag( @@ -227,18 +253,34 @@ class PaperKernelPipeline: methods, replacements, version=version, + previous_diffused=( + previous_bundle.diffused_proof + if previous_bundle is not None + else None + ), + feedback=feedback, ) variant = await self.render_variant( dag, replacements, diffused, version=version, + previous_variant=( + previous_bundle.variant if previous_bundle is not None else None + ), + feedback=feedback, ) - return CandidateBundle( + bundle = CandidateBundle( replacement_plan=replacements, diffused_proof=diffused, variant=variant, ) + if ( + previous_bundle is not None + and sha256_payload(bundle) == sha256_payload(previous_bundle) + ): + raise ValueError("repair pass returned an unchanged candidate bundle") + return bundle async def _judge_once( self, @@ -251,36 +293,21 @@ class PaperKernelPipeline: iteration: int, judge_id: int, ) -> JudgeVerdict: - format_feedback = "" - for attempt in range(1, 4): - request_id = ( - f"{item.item_id}.verify.t{iteration:02d}." - f"j{judge_id}.a{attempt}" - ) - verdict = JudgeVerdict.model_validate( - await self._call( - judge, - request_id=request_id, - system_prompt=JUDGE_SYSTEM, - user_prompt=judge_user( - item, - dag, - methods, - bundle.replacement_plan, - bundle.diffused_proof, - bundle.variant, - format_feedback=format_feedback, - ), - ) + request_id = f"{item.item_id}.verify.t{iteration:02d}.j{judge_id}" + verdict = JudgeVerdict.model_validate( + await self._call( + judge, + request_id=request_id, + system_prompt=JUDGE_SYSTEM, + user_prompt=judge_user( + item, + methods, + bundle.replacement_plan, + bundle.variant, + ), ) - try: - return verdict.validate_coverage(dag, bundle.replacement_plan) - except ValueError as exc: - format_feedback = str(exc) - raise ValueError( - f"judge {judge_id} failed coverage after three format attempts: " - f"{format_feedback}" ) + return verdict.validate_coverage(dag) @staticmethod def _feedback(verdicts: list[JudgeVerdict]) -> str: @@ -369,6 +396,7 @@ class PaperKernelPipeline: dag, methods, version=iteration + 1, + previous_bundle=bundle, feedback=self._feedback(verdicts), ) repaired_from_previous = True diff --git a/src/gap_pipeline/pipeline.py b/src/gap_pipeline/pipeline.py index 4a5332a..9d03525 100644 --- a/src/gap_pipeline/pipeline.py +++ b/src/gap_pipeline/pipeline.py @@ -1,4 +1,4 @@ -"""Historical two-call Putnam generator retained for provenance tests. +"""Consolidated two-call Putnam generation interface. The default manuscript-aligned implementation is ``paper_pipeline.py``. """ diff --git a/src/gap_pipeline/release.py b/src/gap_pipeline/release.py index 7df6566..5bf8d95 100644 --- a/src/gap_pipeline/release.py +++ b/src/gap_pipeline/release.py @@ -89,21 +89,22 @@ def export_release( continue candidate = kernel_payload["accepted_candidate"] + method_nodes = kernel_payload["method_plan"]["nodes"] + replacement_changes = kernel_payload["accepted_replacement_plan"]["changes"] variants["kernel_variant"] = { "question": candidate["question"], "solution": candidate["solution"], "_meta": { - "proof_dag": kernel_payload["proof_dag"], - "method_plan": kernel_payload["method_plan"], - "replacement_plan": kernel_payload["accepted_replacement_plan"], - "diffused_proof": kernel_payload["accepted_diffused_proof"], - "terminal_answer": candidate["terminal_answer"], - "accepted_candidate_sha256": kernel_payload[ - "accepted_candidate_sha256" - ], - "accepted_bundle_sha256": kernel_payload[ - "accepted_bundle_sha256" + "core_steps": [ + node["method_label"] for node in method_nodes ], + "mutable_slots": { + change["slot_id"]: { + "description": change["description"], + "original": change["original_value"], + } + for change in replacement_changes + }, }, } output_record = copy.deepcopy(record) -- cgit v1.2.3