Просмотр исходного кода

框架:复用同一冻结快照的语义验收结果

当重试产生完全相同的 Artifact 与 Evidence 快照时,复用已经完成的语义 Validation,避免重复调用 Validator 模型;保留新的 Validation 记录和来源追踪,并补充失败 verdict 复用测试。
SamLee 5 часов назад
Родитель
Сommit
b1f187184f

+ 70 - 0
agent/agent/orchestration/coordinator.py

@@ -1506,6 +1506,13 @@ class TaskCoordinator:
             )
             if preflight_result is not None:
                 validator_result = preflight_result
+            elif (
+                reused := await self._reuse_semantic_validation(
+                    validation_context,
+                    validation,
+                )
+            ) is not None:
+                validator_result = reused
             elif validation.validation_plan.mode == ValidationMode.DETERMINISTIC:
                 validator_result = await asyncio.wait_for(
                     self._run_deterministic_validation(
@@ -1586,6 +1593,69 @@ class TaskCoordinator:
             failure=validation.failure,
         )
 
+    async def _reuse_semantic_validation(
+        self,
+        context: ValidationContext,
+        validation: ValidationReport,
+    ) -> Optional[ValidatorRunResult]:
+        """Reuse one completed semantic verdict for the same frozen Task snapshot."""
+
+        if validation.validation_plan.mode is not ValidationMode.AGENT:
+            return None
+        ledger = await self.task_store.load(context.root_trace_id)
+        for validation_id in reversed(context.task.validation_ids):
+            prior = ledger.validations.get(validation_id)
+            if (
+                prior is None
+                or prior.validation_id == validation.validation_id
+                or prior.status is not ValidationRunStatus.COMPLETED
+                or prior.verdict is None
+                or prior.spec_version != context.attempt.spec_version
+                or prior.validation_plan.mode is not ValidationMode.AGENT
+                or prior.summary.startswith("Deterministic preflight rejected")
+            ):
+                continue
+            prior_attempt = ledger.attempts.get(prior.attempt_id)
+            if prior_attempt is None or prior_attempt.snapshot_id is None:
+                continue
+            prior_snapshot = await self.artifact_store.get(
+                context.root_trace_id,
+                prior_attempt.snapshot_id,
+            )
+            if (
+                prior_snapshot.artifact_refs != context.snapshot.artifact_refs
+                or prior_snapshot.evidence_refs != context.snapshot.evidence_refs
+            ):
+                continue
+            await self.submit_validation(
+                {
+                    "role": AgentRole.VALIDATOR.value,
+                    "root_trace_id": context.root_trace_id,
+                    "task_id": context.task.task_id,
+                    "attempt_id": context.attempt.attempt_id,
+                    "validation_id": validation.validation_id,
+                    "snapshot_id": context.snapshot.snapshot_id,
+                    "trace_id": validation.validator_trace_id,
+                    "tool_call_id": f"reuse:{prior.validation_id}",
+                    "operation_id": validation.operation_id,
+                    "execution_epoch": validation.execution_epoch,
+                },
+                prior.verdict,
+                tuple(prior.criterion_results),
+                f"semantic_validation_reused_from={prior.validation_id};{prior.summary}"[:500],
+                tuple(prior.evidence_refs),
+                tuple(prior.unverified_claims),
+                tuple(prior.risks),
+                prior.recommendation,
+            )
+            return ValidatorRunResult(
+                trace_id=validation.validator_trace_id,
+                status="completed",
+                summary="semantic_validation_reused",
+                execution_stats=ExecutionStats(total_tokens=0, total_cost=0.0),
+            )
+        return None
+
     async def _run_validation_preflight(
         self,
         context: ValidationContext,

+ 42 - 0
agent/tests/test_orchestration_validation_policy.py

@@ -19,6 +19,7 @@ from agent.orchestration.models import (
     ArtifactRef,
     ArtifactSnapshot,
     AttemptSubmission,
+    DecisionAction,
     TaskAttempt,
     TaskRecord,
     TaskSpec,
@@ -217,6 +218,47 @@ async def test_coordinator_freezes_default_agent_plan(tmp_path):
     assert ValidationPlan.from_dict(report.validation_plan) == report.validation_plan
 
 
+@pytest.mark.asyncio
+async def test_same_frozen_snapshot_reuses_one_semantic_validation(tmp_path):
+    class SameArtifactExecutor(FakeExecutor):
+        async def run_worker(self, context):
+            self.worker_calls += 1
+            await self.coordinator.submit_attempt(
+                {
+                    **context,
+                    "role": AgentRole.WORKER.value,
+                    "trace_id": context["worker_trace_id"],
+                    "tool_call_id": f"submit-{context['attempt_id']}",
+                },
+                AttemptSubmission(
+                    summary="same content",
+                    artifact_refs=[ArtifactRef(uri="memory://same", version="1")],
+                ),
+            )
+            return WorkerRunResult(context["worker_trace_id"], "completed")
+
+    executor = SameArtifactExecutor([ValidationVerdict.FAILED])
+    coordinator, store, _ = await make_coordinator(tmp_path, executor)
+    task_id = await create_task(coordinator, "semantic validation cache")
+
+    first = (await coordinator.dispatch_tasks("root", [task_id]))[0]
+    await coordinator.decide_task(
+        "root",
+        task_id,
+        first.validation_id,
+        DecisionAction.RETRY,
+        {"reason": "retry unchanged content"},
+        "retry-same-snapshot",
+    )
+    second = (await coordinator.dispatch_tasks("root", [task_id]))[0]
+    report = (await store.load("root")).validations[second.validation_id]
+
+    assert executor.validator_calls == 1
+    assert report.verdict is ValidationVerdict.FAILED
+    assert report.summary.startswith("semantic_validation_reused_from=")
+    assert report.execution_stats.total_tokens == 0
+
+
 @pytest.mark.asyncio
 @pytest.mark.parametrize("verdict", [ValidationVerdict.PASSED, ValidationVerdict.FAILED])
 async def test_coordinator_executes_deterministic_plan_without_agent_validator(tmp_path, verdict):