From ff3e18ee086dc4057c1fca055d80d467a5ecb181 Mon Sep 17 00:00:00 2001 From: Pigbibi <20649888+Pigbibi@users.noreply.github.com> Date: Wed, 9 Sep 2026 06:00:30 +0800 Subject: [PATCH] fix: preserve research progress across multi-day shadow observations Co-Authored-By: Codex --- .../paired_shadow_evidence.py | 86 +++++++++++-------- .../research_promotion_cycle.py | 5 +- tests/test_paired_shadow_adapter.py | 16 ++++ tests/test_paired_shadow_evidence.py | 81 +++++++++++++++++ .../test_research_promotion_shadow_resume.py | 50 +++++++++++ 5 files changed, 200 insertions(+), 38 deletions(-) diff --git a/src/quant_platform_kit/strategy_lifecycle/paired_shadow_evidence.py b/src/quant_platform_kit/strategy_lifecycle/paired_shadow_evidence.py index 03d78547..ef488cc6 100644 --- a/src/quant_platform_kit/strategy_lifecycle/paired_shadow_evidence.py +++ b/src/quant_platform_kit/strategy_lifecycle/paired_shadow_evidence.py @@ -231,22 +231,14 @@ def build_paired_shadow_evidence( ) -def validate_paired_shadow_evidence( +def _validate_paired_shadow_content( value: Mapping[str, object], *, policy: ForwardObservationPolicy | None = None, forward_observation_receipt: Mapping[str, object] | None = None, - previous_evidence: Mapping[str, object] | None = None, previous_forward_observation_receipt: Mapping[str, object] | None = None, ) -> dict[str, object]: - """Validate integrity and, when supplied, bind the evidence to P4 receipts. - - A standalone validation establishes only the artifact's closed schema and - digest. A trusting consumer must provide its frozen policy and matching - forward-observation receipt; a continuous consumer also provides both - predecessors. This prevents a recent production performance snapshot - from being re-labelled as paired shadow evidence. - """ + """Check one artifact's content and optional receipt binding, not ancestry.""" if not isinstance(value, Mapping) or set(value) != _TOP_LEVEL_FIELDS: _invalid("evidence must be a closed paired-shadow object") @@ -305,23 +297,61 @@ def validate_paired_shadow_evidence( elif previous_forward_observation_receipt is not None: _invalid("previous forward-observation receipt requires current receipt") + return { + "schema_version": PAIRED_SHADOW_EVIDENCE_SCHEMA_VERSION, + "evidence_kind": PAIRED_SHADOW_EVIDENCE_KIND, + "candidate_id": candidate_id, + "baseline_id": baseline_id, + "observation_session": session, + "observation_index": index, + "observed_at": timestamp, + "input_snapshot_sha256": input_digest, + "forward_observation_receipt_sha256": forward_digest, + "previous_paired_shadow_evidence_sha256": previous_digest, + "candidate": candidate, + "baseline": baseline, + "no_order": True, + "live_authority_granted": False, + "paired_shadow_evidence_sha256": claimed_digest, + } + + +def validate_paired_shadow_evidence( + value: Mapping[str, object], + *, + policy: ForwardObservationPolicy | None = None, + forward_observation_receipt: Mapping[str, object] | None = None, + previous_evidence: Mapping[str, object] | None = None, + previous_forward_observation_receipt: Mapping[str, object] | None = None, +) -> dict[str, object]: + """Validate one artifact and its actual adjacent evidence/receipt pair. + + Every non-first artifact requires both actual predecessors. A full-window + consumer must walk from the first record; validating a later adjacent pair + alone does not establish the earlier history or window completion. + """ + current = _validate_paired_shadow_content( + value, policy=policy, forward_observation_receipt=forward_observation_receipt, + previous_forward_observation_receipt=previous_forward_observation_receipt, + ) + previous_digest = current["previous_paired_shadow_evidence_sha256"] if previous_evidence is not None: if forward_observation_receipt is None or previous_forward_observation_receipt is None: _invalid("continuous evidence validation requires both forward-observation receipts") - previous = validate_paired_shadow_evidence( + previous = _validate_paired_shadow_content( previous_evidence, policy=policy, forward_observation_receipt=previous_forward_observation_receipt, ) if previous_digest != previous["paired_shadow_evidence_sha256"]: _invalid("previous_paired_shadow_evidence_sha256 does not match predecessor") - if index != int(previous["observation_index"]) + 1: + if current["observation_index"] != int(previous["observation_index"]) + 1: _invalid("observation_index does not increment from the predecessor") - if session <= str(previous["observation_session"]): + if current["observation_session"] <= str(previous["observation_session"]): _invalid("observation_session does not advance from the predecessor") - if timestamp <= str(previous["observed_at"]): + if current["observed_at"] <= str(previous["observed_at"]): _invalid("observed_at does not advance from the predecessor") - if candidate_id != previous["candidate_id"] or baseline_id != previous["baseline_id"]: + if current["candidate_id"] != previous["candidate_id"] or current["baseline_id"] != previous["baseline_id"]: _invalid("candidate or baseline identity changed within the evidence chain") previous_forward_digest = previous["forward_observation_receipt_sha256"] forward_previous_digest = previous_forward_observation_receipt.get( @@ -337,36 +367,20 @@ def validate_paired_shadow_evidence( elif previous_digest is not None: _invalid("previous paired-shadow digest requires predecessor evidence") - return { - "schema_version": PAIRED_SHADOW_EVIDENCE_SCHEMA_VERSION, - "evidence_kind": PAIRED_SHADOW_EVIDENCE_KIND, - "candidate_id": candidate_id, - "baseline_id": baseline_id, - "observation_session": session, - "observation_index": index, - "observed_at": timestamp, - "input_snapshot_sha256": input_digest, - "forward_observation_receipt_sha256": forward_digest, - "previous_paired_shadow_evidence_sha256": previous_digest, - "candidate": candidate, - "baseline": baseline, - "no_order": True, - "live_authority_granted": False, - "paired_shadow_evidence_sha256": claimed_digest, - } + return current def canonical_paired_shadow_evidence_bytes(value: Mapping[str, object]) -> bytes: - """Return canonical bytes for a schema-valid paired-shadow artifact.""" + """Serialize valid content; this proves neither ancestry nor window completion.""" - return _canonical_bytes(validate_paired_shadow_evidence(value)) + return _canonical_bytes(_validate_paired_shadow_content(value)) def paired_shadow_evidence_sha256(value: Mapping[str, object]) -> str: - """Return the deterministic SHA-256 identity of a valid artifact.""" + """Return a content identity, not evidence of ancestry or window completion.""" return str( - validate_paired_shadow_evidence(value)["paired_shadow_evidence_sha256"] + _validate_paired_shadow_content(value)["paired_shadow_evidence_sha256"] ) diff --git a/src/quant_platform_kit/strategy_lifecycle/research_promotion_cycle.py b/src/quant_platform_kit/strategy_lifecycle/research_promotion_cycle.py index 47c25ef7..70d8d935 100644 --- a/src/quant_platform_kit/strategy_lifecycle/research_promotion_cycle.py +++ b/src/quant_platform_kit/strategy_lifecycle/research_promotion_cycle.py @@ -40,6 +40,7 @@ class ResearchPromotionState(str, Enum): _ACTIVE_DRIFT = {DriftStatus.REVIEW, DriftStatus.CRITICAL} +_RESEARCH_CANDIDATE_RECOMMENDATIONS = frozenset({"promote", "needs_review", "research_candidate"}) _TERMINAL = { ResearchPromotionState.PARKED, ResearchPromotionState.HUMAN_ACCEPTED, @@ -1012,7 +1013,7 @@ def run_research_promotion_cycle( ticket.updated_at = _now_iso() return ticket - if proposal.recommendation not in {"promote", "needs_review", "research_candidate"}: + if proposal.recommendation not in _RESEARCH_CANDIDATE_RECOMMENDATIONS: ticket.state = ResearchPromotionState.PARKED ticket.notes = (f"recommendation={proposal.recommendation}",) ticket.updated_at = _now_iso() @@ -1429,7 +1430,7 @@ def deliver_saved_ticket() -> None: return output("research_checkpoint_invalid", status="parked") gates_ok, _ = enforce_promotion_backtest_gates(proposal, stages["backtest"]["result"]) budget_ok, _ = enforce_optimization_budget(proposal, budget) - if (not gates_ok or not budget_ok or proposal.recommendation != "promote" + if (not gates_ok or not budget_ok or proposal.recommendation not in _RESEARCH_CANDIDATE_RECOMMENDATIONS or proposal.strategy_profile != drift.strategy_profile or proposal.domain != drift.domain or (progress.get("diagnosis_required") and stages["diagnose"]["result"].get("optimization_needed") is not True)): diff --git a/tests/test_paired_shadow_adapter.py b/tests/test_paired_shadow_adapter.py index 8a523532..2cb9aa28 100644 --- a/tests/test_paired_shadow_adapter.py +++ b/tests/test_paired_shadow_adapter.py @@ -157,3 +157,19 @@ def _fake_build(**kwargs): # noqa: ARG001 monkeypatch.setattr(adapter, "build_paired_shadow_evidence", _fake_build) with pytest.raises(ValueError, match="live_authority_granted"): collect_paired_shadow_for_promotion(_observation()) + + +def test_collect_advances_three_actual_receipt_pairs() -> None: + previous = previous_receipt = None + for index, session in enumerate(("2026-08-26", "2026-08-27", "2026-08-28"), 1): + receipt = _forward_receipt(previous=previous_receipt, index=index, session=session) + record = collect_paired_shadow_for_promotion(PairedShadowObservation( + policy=_policy(), forward_observation_receipt=receipt, + baseline_id="soxl-v6-live-baseline", observed_at=session + "T20:00:00-04:00", + input_snapshot_sha256="a" * 64, candidate=_leg("candidate"), baseline=_leg("baseline"), + previous_evidence=previous, previous_forward_observation_receipt=previous_receipt, + )) + assert record["passed"] is True + assert record["no_order"] is True and record["live_authority_granted"] is False + assert record["evidence"]["observation_index"] == index + previous, previous_receipt = record["evidence"], receipt diff --git a/tests/test_paired_shadow_evidence.py b/tests/test_paired_shadow_evidence.py index dca9b064..3136fda7 100644 --- a/tests/test_paired_shadow_evidence.py +++ b/tests/test_paired_shadow_evidence.py @@ -2,6 +2,7 @@ import copy from datetime import date +from hashlib import sha256 import json import pytest @@ -242,3 +243,83 @@ def test_report_artifacts_embed_the_validated_evidence_without_runtime_changes() json.loads(serialized["artifacts"]["paired_shadow_evidence_json"]) == evidence ) + + +def _three_session_chain(): + chain = [] + for index, session in enumerate(("2026-08-26", "2026-08-27", "2026-08-28"), 1): + previous, previous_receipt = chain[-1] if chain else (None, None) + receipt = _forward_receipt(previous=previous_receipt, index=index, session=session) + evidence = _evidence( + forward_observation_receipt=receipt, + observed_at=session + "T20:00:00-04:00", + previous_evidence=previous, + previous_forward_observation_receipt=previous_receipt, + ) + chain.append((evidence, receipt)) + return chain + + +def test_three_sessions_validate_serialize_and_preserve_no_order_guards() -> None: + chain = _three_session_chain() + for index, (evidence, receipt) in enumerate(chain): + previous, previous_receipt = chain[index - 1] if index else (None, None) + assert validate_paired_shadow_evidence( + evidence, policy=_policy(), forward_observation_receipt=receipt, + previous_evidence=previous, previous_forward_observation_receipt=previous_receipt, + ) == evidence + assert json.loads(canonical_paired_shadow_evidence_bytes(evidence)) == evidence + assert paired_shadow_evidence_sha256(evidence) == evidence["paired_shadow_evidence_sha256"] + artifacts = build_paired_shadow_evidence_report_artifacts( + evidence, policy=_policy(), forward_observation_receipt=receipt, + previous_evidence=previous, previous_forward_observation_receipt=previous_receipt, + ) + assert json.loads(artifacts["paired_shadow_evidence_json"]) == evidence + assert artifacts["paired_shadow_evidence_no_order"] is True + assert artifacts["paired_shadow_evidence_live_authority_granted"] is False + assert evidence["observation_index"] < _policy().required_trading_sessions + + +@pytest.mark.parametrize("failure", ["missing", "tampered", "wrong_receipt", "wrong_candidate"]) +def test_three_session_chain_rejects_missing_or_invalid_predecessor(failure) -> None: + chain = _three_session_chain() + current, receipt = chain[-1] + previous, previous_receipt = copy.deepcopy(chain[-2]) + if failure == "missing": + previous = previous_receipt = None + elif failure == "tampered": + previous["candidate"]["cost"]["source"] = "changed" + elif failure == "wrong_receipt": + previous_receipt = chain[0][1] + else: + previous["candidate_id"] = "another-candidate" + previous["paired_shadow_evidence_sha256"] = sha256(json.dumps( + {k: v for k, v in previous.items() if k != "paired_shadow_evidence_sha256"}, + sort_keys=True, separators=(",", ":"), ensure_ascii=False, allow_nan=False, + ).encode()).hexdigest() + with pytest.raises(ValueError): + validate_paired_shadow_evidence( + current, policy=_policy(), forward_observation_receipt=receipt, + previous_evidence=previous, previous_forward_observation_receipt=previous_receipt, + ) + + +def test_full_chain_walk_rejects_changed_ancestor() -> None: + chain = _three_session_chain() + changed = copy.deepcopy(chain) + first = changed[0][0] + first["baseline_id"] = "changed-historical-baseline" + first["paired_shadow_evidence_sha256"] = sha256(json.dumps( + {k: v for k, v in first.items() if k != "paired_shadow_evidence_sha256"}, + sort_keys=True, separators=(",", ":"), ensure_ascii=False, allow_nan=False, + ).encode()).hexdigest() + # The full-window consumer walks from the first actual record, so a later + # internally valid pair cannot hide a changed ancestor. + previous = previous_receipt = None + with pytest.raises(ValueError, match="predecessor|identity changed"): + for evidence, receipt in changed: + previous = validate_paired_shadow_evidence( + evidence, policy=_policy(), forward_observation_receipt=receipt, + previous_evidence=previous, previous_forward_observation_receipt=previous_receipt, + ) + previous_receipt = receipt diff --git a/tests/test_research_promotion_shadow_resume.py b/tests/test_research_promotion_shadow_resume.py index 61cacc92..7abf0ef8 100644 --- a/tests/test_research_promotion_shadow_resume.py +++ b/tests/test_research_promotion_shadow_resume.py @@ -1,5 +1,6 @@ """Synthetic local observations; no model, broker, real shadow or console.""" +from dataclasses import replace from datetime import date, datetime, timezone from unittest.mock import Mock @@ -79,6 +80,55 @@ def test_actual_runner_resumes_only_local_shadow_after_original_drift_expires(tm "record_shadow", "read_pending_shadow", "sync_console")] == [1] * 6 +@pytest.mark.parametrize("recommendation", ["promote", "needs_review", "research_candidate"]) +@pytest.mark.parametrize("decision", ["accept", "reject"]) +def test_saved_candidate_recommendations_resume_through_the_same_human_gate(tmp_path, recommendation, decision): + from tests.test_research_promotion_reconciliation import remote_decision + + args = job(tmp_path, optimize=Mock(return_value=replace(_proposal(), recommendation=recommendation)), + sync_console=Mock(return_value=False)) + first = run_actionable_research_promotion(**args) + assert first["status"] == "deferred" + local = cycle.load_research_promotion_ticket(first["ticket_path"]) + assert local.research_progress["stages"]["optimize"]["result"]["recommendation"] == recommendation + args["sync_console"].assert_not_called() + + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + ready = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert ready["status"] == "awaiting_human" + assert ready["research_key"] == first["research_key"] + assert ready["ticket"]["ticket_id"] == first["ticket"]["ticket_id"] + local = cycle.load_research_promotion_ticket(ready["ticket_path"]) + args["pull_console"].return_value = remote_decision(local, decision) + Clock.current = datetime(2026, 11, 21, tzinfo=timezone.utc) + for _ in range(2): + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-11-21"}) + assert result["status"] == ("human_accepted" if decision == "accept" else "human_rejected") + assert result["ticket"]["live_authority_granted"] is False + assert [args[key].call_count for key in ("diagnose", "optimize", "enforce_backtest_gates", + "record_shadow", "read_pending_shadow", "sync_console", "pull_console")] == [1] * 7 + + +@pytest.mark.parametrize("mutation", ["recommendation", "backtest"]) +def test_candidate_recommendation_does_not_bypass_saved_strict_gates(tmp_path, mutation): + args = job(tmp_path, optimize=Mock(return_value=replace(_proposal(), recommendation="research_candidate"))) + first = run_actionable_research_promotion(**args) + assert first["status"] == "deferred" + local = cycle.load_research_promotion_ticket(first["ticket_path"]) + stages = local.research_progress["stages"] + if mutation == "recommendation": + stages["optimize"]["result"]["recommendation"] = "reject" + else: + stages["backtest"]["result"]["status"] = "FAIL" + cycle.save_research_promotion_ticket(local, first["ticket_path"]) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert result["status"] == "parked" and result["reason"] == "research_checkpoint_invalid" + assert [args[key].call_count for key in ("diagnose", "optimize", "enforce_backtest_gates", "record_shadow")] == [1] * 4 + args["read_pending_shadow"].assert_not_called() + args["sync_console"].assert_not_called() + + @pytest.mark.parametrize("retry", [None, True, float("nan"), float("inf"), 0]) def test_missing_or_invalid_pending_deadline_never_automatically_reads(tmp_path, retry): args = job(tmp_path, record_shadow=Mock(return_value=pending(retry_at=retry)))