diff --git a/.github/workflows/reusable-drift-check.yml b/.github/workflows/reusable-drift-check.yml index a86aa52..ea9a6f4 100644 --- a/.github/workflows/reusable-drift-check.yml +++ b/.github/workflows/reusable-drift-check.yml @@ -308,7 +308,7 @@ jobs: uses: actions/checkout@v6 with: repository: QuantStrategyLab/AIAuditBridge - ref: 2351c987d6df12d48355aaa5803db8b089476310 + ref: 60bd64a2ae059a082614181eeb845b46df395523 path: external/AIAuditBridge - name: Dual-review critical drift diff --git a/docs/research_promotion_resume.zh-CN.md b/docs/research_promotion_resume.zh-CN.md new file mode 100644 index 0000000..74b7982 --- /dev/null +++ b/docs/research_promotion_resume.zh-CN.md @@ -0,0 +1,88 @@ +# 冻结研究输入的本地恢复 + +本轮在既有 research promotion ticket 上保存本地阶段进度。没有增加队列、排班或交易入口。基线为 QPK `d12309d`。 + +## 调用与归属 + +策略仓先验证真实输入、冻结窗口及 runner,再调用现有 +`run_actionable_research_promotion`,增加两个参数: + +- `ticket_dir`:调度者提供的、重启后仍保留的票据目录。 +- `research_identity`:包含下列五个非空字符串的 mapping。 + +| 字段 | 必须绑定的实际材料 | +| --- | --- | +| `code_revision` | 实际执行的冻结策略代码版本 | +| `input_revision` | 已验证的冻结输入根,包括开发、WFA、锁定 OOS 的窗口和来源 | +| `param_space_revision` | 基线参数、固定参数和允许搜索的参数空间 | +| `cost_model_revision` | 完整费用与执行假设,包括最低佣金、成交参与率、lot 等适用设置 | +| `validator_revision` | 实际严格 validator、promotion plan 和阈值版本 | + +复用已有 manifest/config 的根或真实版本号。QPK 不会把调用者填写的字符串当作输入已验证的证明;这些材料须在策略仓构造回调时验证并固定。修改其中任一材料必须修改对应 revision。 + +回调接口保持:`optimize(drift, budget) -> OptimizationProposal`、 +`enforce_backtest_gates(proposal) -> PromotionBacktestRun | Mapping`、 +`record_shadow(proposal) -> Mapping`。 +可选 `diagnose(drift, budget) -> Mapping` 返回既有 +`optimization_needed` 决定;AI 意见仍不能替代任何严格门或人工决定。 +可选 `pull_console(ticket_id)` 使用原 QRT 读取接口。 + +可选 `admit_new_research(ticket_dir: Path, created_at: str) -> bool` 由实际调度者提供,正常周期与 runner 均向下传递。仅在持有同一目录锁且确认该身份没有票据之后、首次保存和 AI 调用之前执行;只有字面值 `True` 放行,异常返回 `research_admission_unavailable`,其他值返回 `new_research_not_admitted`,两者都不创建票据或调用阶段。`created_at` 是本次新票据实际持久化的同一个 UTC 时间。已有票据恢复不再次准入;回调不可重复获取该目录锁。AAB 可据已有票据的 `created_at` 实现每日新实验限额,坏记录和未知计数必须拒绝放行;QPK 不另建计数存储,也不在此实现调度者的每日政策。 + +正常 `run_auto_pilot_cycle(..., research_identity=...)` 已接入此入口,以 +`store.local_root/research_promotion_tickets` 保存进度。未提供冻结身份或持久目录时,自动研究在 AI 前返回 `research_identity_unavailable`。 +正常周期发现尚未开始同步的 awaiting checkpoint 时,使用 `resume_delivery_only` 路径:必须重新匹配同一冻结身份且票据已经存在,只可回收决定或继续首次同步;不得借此创建新票据或运行 AI/回测。其他 pending profile 仍保持跳过。 +未传持久化参数的显式 runner 调用保留原接口,不能声称具有跨 run 去重。 + +QPK 管单次实验的本地进度、严格门与人工决定。AAB 管 AI job、额度准入和实际实验 dispatch;不同 GitHub 临时 runner 或不同目录不会共享 QPK 的锁。跨主机单重任务、每日候选上限和真实任务工件交接须由调度者保证,不由本补丁宣称完成。 + +## 去重与中断 + +去重身份包含五个 revision、target/domain、观察日期和 source revision、drift score、基线标识和硬预算。监测器的冷却、升级标记以及可选诊断回调的有无不会生成第二个相同实验。 + +持久化入口仍要求新鲜、可行动的 drift,绑定优化器、严格 backtest 和 shadow 回调;预算最多 25 次搜索、4 个参数,且必须 paired shadow。缺绑定、过期输入或身份缺失时不会调用 AI。 + +本地 ticket 增加 `research_progress`,包含身份及各阶段结果;该字段不进入 `to_dict()` 的控制台候选,也不发送给 QRT。保存继续使用临时文件、fsync 和原子替换。完整 proposal(含参数、费用、validation identity)和已完成 backtest/shadow 摘要可恢复;输入原始行和凭据不应作为回调结果传入。 + +| 保存状态 | 下次行为 | +| --- | --- | +| 某阶段尚未开始 | 从该阶段继续 | +| 某阶段已完成 | 复用结果;严格 WFA/OOS 门仍重新检查后才到 shadow | +| 阶段为 `running` 或 `unknown` | `research_outcome_unknown`,不再次调用该阶段;需要拥有原作业证据的操作方只读查明结果 | +| 明确额度延期 | 保存 `retry_at`;到期前不调用 AI,缺时间或首次收到时已经过期的时间不自动重试 | +| awaiting ticket,尚未开始同步 | 可以继续首次同步,原 sync adapter 仍先 GET 再决定是否 POST | +| 同步已开始但未确认 | 仅 GET 回收;不自动重复 POST,不能把失败当作不存在 | +| awaiting ticket 已有人工决定 | 原身份/参数/证据/权限校验后记录意图 | +| 本地 terminal ticket | 同一身份不重新诊断、优化或提交;新冻结证据可以产生新实验 | + +### 长周期 paired shadow 的只读续收 + +新增可选 `read_pending_shadow(proposal) -> Mapping`,通过正常周期和 actionable runner 传到 saved 入口。首次 `record_shadow` 仍只执行一次;后续只使用这个专用读取回调,不能重复首次可能创建外部观察的操作。AAB 的实际回调读取受控本地 observation,并校验冻结 policy、候选/参数/source、完整 forward 窗口及当期材料有效性。 + +首次记录或后续读取可以返回 `status=pending`、`passed=false`、`no_order=true`、`live_authority_granted=false` 及未来的有限 Unix `retry_at`。票据保持已有非终态 `shadow_recorded`,本地 shadow 阶段为 pending,对外返回 deferred。到期前零读取;缺失、非有限或回调返回时已过期的 retry_at 保存为 null,不自动轮询。回调明确失败则 PARKED,异常或进程在调用中停止则 unknown,不再次调用。后续仍 pending 时,只更新明确的新期限;不会重跑诊断、优化、严格回测或首次记录。 + +完成返回 `{"status": "complete", "observation": ...}`。observation 使用既有 `PairedShadowObservation` 或 `collect_paired_shadow_for_promotion` 输入 mapping,包括 policy、当前/前一 receipt 和 paired evidence。QPK 校验 policy 的 profile/domain 和观察数量,并调用原 paired validator 验证本次及前序关联;修复 adapter 二次校验漏传 previous 材料的问题,缺少前序仍拒绝。裸 `passed=true` 不能替代这些材料。完整窗口的真实性、当前可用性以及具体候选参数/source 绑定仍由 AAB 的真实 policy reader 验证,单条合法 receipt 不能证明窗口完成。 + +新实验仍须通过 7 日 drift 有效期。过期观察只可定位同身份的已有票据:此前要求的诊断、优化、严格回测必须已完成,原 proposal 预算与严格门重新校验通过,且 shadow 为明确 pending,才能续收。刚保存了完整 shadow 结果但尚未来得及组装人工候选的中断,也只复用这个已验证结果。以上两类恢复均不会重调 AI、优化、backtest 或首次 `record_shadow`。身份改变、阶段缺失、unknown、新票据、未来/缺失日期仍拒绝;不通过刷新原 as_of 创建替代实验。正常 auto-pilot、actionable runner 与 saved 入口均保留此限制。 + +已经到达 awaiting 的同身份票据,也能在原 drift 过期后只读回收控制台决定:要求上述阶段和 shadow 全部完成、原严格门有效、票据候选参数匹配保存的 proposal。此路径只 GET,不再次 POST,不读取 shadow 或运行前置研究。未知提交先读取原票据;接受/拒绝后返回已有终态,不因时间过去而变成新实验。 + +`research_in_progress` 表示同一目录已被另一进程持有,本次返回 `deferred`;不会等待或发起模型请求。研究、自动决定回收及本地 `research-promotion-decide` CLI 使用同一目录锁。损坏的独立票据不会阻断其他身份;损坏的当前身份不会被覆盖成新研究。 + +输出沿用 `status/reason/ticket/console_synced`,增加 `research_key`(本地去重键)、`ticket_path`、`resumed`,额度延期时有 `retry_at`。`console_synced=true` 只来自本次同步回调明确确认;恢复 GET 的结果单列 `reconciliation`。任何接受仍只有 intent,`live_authority_granted=false`。 + +## 验证边界 + +新增离线回归实际覆盖完成阶段重用、阶段间/阶段中中断、完整 proposal 保存恢复、输入/代码/参数/费用/validator 变更、坏票据隔离、额度延期、未知 POST 只读恢复、首次提交前恢复及人工决定。独立子进程验证目录锁竞争和进程退出后的释放;Windows 分支尚未在 Windows 运行。 + +使用合成证据和注入回调验证控制流,不等于真实 WFA/OOS、paired forward shadow、模型调用、QRT 生产写入、部署或交易验证。没有修改 frozen artifacts,也没有调用真实 provider。 + +最终本地结果:新增恢复测试 53 项;恢复/cycle/runner/人工回收/codex integration/CLI/AI provider/freshness/paired shadow/strict orchestrator 共 11 个测试文件,271 tests + 37 subtests 通过;Ruff 与 `git diff --check` 通过。新增准入 10 项覆盖拒绝/异常、锁内执行、UTC 时间一致及已有票据不重复准入,实际先 RED 后 GREEN。测试使用 `/usr/local/bin/python3`、`PYTHONPATH=src`、`PYTHONDONTWRITEBYTECODE=1`、禁用 pytest 外部插件和 cache,并阻断真实 socket 连接。 + +本次还将唯一 reusable drift workflow 的 AAB checkout 固定到 `60bd64a2ae059a082614181eeb845b46df395523`,采用已发布的研究主审 Codex-only、OIDC 和 `promotion_review` 边界。该版本主审显式传 `allowed_providers=["codex"]`,无需额外 provider 环境覆盖。精确引用断言先 RED 后 GREEN,工作流测试 1 项及 actionlint 通过;排班、采集和 validator 未修改。此处仅记录源码采用,实际发布和业务周期由部署方核实。 + +组合测试显式使用 `TZ=UTC`:既有 freshness 测试的一例采用本地 `date.today()`,生产观察时钟采用 UTC;本机跨日时,不指定 TZ 会先触发 future 拒绝,而非该测试预期的缺 source 拒绝。本轮保留生产 UTC 语义,没有修改该旧测试。 + +长周期观察增量在已发布分支 `f1e190a` 上新增 27 项离线测试:17 项先 RED,补正常 auto-pilot 和两处保存中断时再复现 3 项 RED;独立审查发现直接策略入口的过期 awaiting 回收被阻断,再用接受/拒绝 2 项 RED 修复,并保留缺阶段/候选错配拒绝。两条真实结构的合成 receipt 验证了多日材料传递,缺少前序、篡改 receipt、裸 passed、未满窗口、live 声明均不通过;没有真实观察数据。新 API 在 actual actionable runner 与正常 auto-pilot 上均验证,CN 的透传采用和 AAB 真实文件 reader 另由对应仓验证。 + +最终在相同隔离环境、阻断 socket 后,12 个相关测试文件为 298 tests + 37 subtests 通过;Ruff 和 `git diff --check` 通过。没有运行模型、真实 shadow、控制台生产写入或交易。 diff --git a/src/quant_platform_kit/strategy_lifecycle/cli.py b/src/quant_platform_kit/strategy_lifecycle/cli.py index f824c2a..f64123f 100644 --- a/src/quant_platform_kit/strategy_lifecycle/cli.py +++ b/src/quant_platform_kit/strategy_lifecycle/cli.py @@ -305,19 +305,10 @@ def _run_research_promotion_decide(args: argparse.Namespace) -> int: raise ValueError( "paper requires --paper-supported for a real broker paper/sim account" ) - load_ticket = _load_callable( - "quant_platform_kit.strategy_lifecycle.research_promotion_cycle", - "load_research_promotion_ticket", - ) apply_decision = _load_callable( "quant_platform_kit.strategy_lifecycle.research_promotion_cycle", - "apply_human_promotion_decision", + "decide_saved_research_promotion_ticket", ) - save_ticket = _load_callable( - "quant_platform_kit.strategy_lifecycle.research_promotion_cycle", - "save_research_promotion_ticket", - ) - ticket = load_ticket(args.ticket) confirmation = None if args.decision == "accept": confirmation = { @@ -326,13 +317,12 @@ def _run_research_promotion_decide(args: argparse.Namespace) -> int: "risk_profile": args.risk_profile, } decided = apply_decision( - ticket, + args.ticket, decision=args.decision, confirmation=confirmation, paper_supported=bool(args.paper_supported), + output_path=args.output, ) - output = args.output or args.ticket - save_ticket(decided, output) _print( f"[research-promotion-decide] ticket={decided.ticket_id} " f"state={decided.state.value} live_authority_granted={decided.live_authority_granted}" diff --git a/src/quant_platform_kit/strategy_lifecycle/codex_integration.py b/src/quant_platform_kit/strategy_lifecycle/codex_integration.py index f471d51..6b96077 100644 --- a/src/quant_platform_kit/strategy_lifecycle/codex_integration.py +++ b/src/quant_platform_kit/strategy_lifecycle/codex_integration.py @@ -489,21 +489,22 @@ def _process_optimization_decision( enforce_backtest_gates: Callable | None = None, record_shadow: Callable | None = None, sync_console: Callable | None = None, + pull_console: Callable | None = None, + research_identity: Mapping[str, str] | None = None, + resume_delivery_only: bool = False, + admit_new_research: Callable[[Path, str], bool] | None = None, + read_pending_shadow: Callable | None = None, ) -> dict[str, object]: """Codex diagnosis → bounded Python research → strict gates → human queue.""" from quant_platform_kit.strategy_lifecycle.promotion_actionable_runner import ( run_actionable_research_promotion, ) - from quant_platform_kit.strategy_lifecycle.research_promotion_cycle import ( - ResearchPromotionTicket, save_research_promotion_ticket, - ) - entry: dict[str, object] = { "strategy": drift.strategy_profile, "drift_status": drift.status.value, "execution_authorized": False, "live_authority_granted": False, } freshness_reason = _drift_freshness_reason(drift) - if freshness_reason: + if freshness_reason and not (freshness_reason == "observation_stale" and callable(read_pending_shadow) and not dry_run): return {**entry, "reason": freshness_reason, "research_promotion_state": "parked"} if drift.status not in (DriftStatus.REVIEW, DriftStatus.CRITICAL) or drift.alert_suppressed: return {**entry, "reason": "drift_not_actionable"} @@ -519,14 +520,17 @@ def _process_optimization_decision( return {**entry, "reason": "observation_source_unavailable", "research_promotion_state": "parked"} context = AiOptimizationContext(strategy_profile=drift.strategy_profile, domain=drift.domain, drift=drift, snapshot=snapshot) - decision = call_ai_optimization_decision(context, dry_run=dry_run) - entry["ai_decision"] = decision if dry_run: - return {**entry, "note": "Dry run: no AI call or research execution"} - if decision.get("reason") == "codex_research_deferred": - return {**entry, "reason": "codex_research_deferred", "research_promotion_state": "deferred", "retry_at": decision.get("retry_at")} - if decision.get("optimization_needed") is not True: - return {**entry, "reason": "codex_did_not_recommend_research", "research_promotion_state": "parked"} + return {**entry, "ai_decision": call_ai_optimization_decision(context, dry_run=True), + "note": "Dry run: no AI call or research execution"} + ticket_dir = getattr(store, "local_root", None) + if not isinstance(ticket_dir, (str, Path)) or research_identity is None: + return {**entry, "reason": "research_identity_unavailable", "research_promotion_state": "parked"} + + def diagnose(*_): + decision = call_ai_optimization_decision(context, dry_run=False) + entry["ai_decision"] = decision + return decision def bounded_optimize(active_drift, budget): from quant_platform_kit.strategy_lifecycle.param_optimizer import run_optimization @@ -540,7 +544,7 @@ def reviewed_backtest_gates(proposal): if verdict.verdict != "approve": # Retain existing deterministic risk/parameter checks. No LLM # escalation may replace them or spend an API fallback budget. - raise ValueError("deterministic_proposal_review_failed") + return {"status": "FAIL", "reason": "deterministic_proposal_review_failed"} return enforce_backtest_gates(proposal) try: @@ -551,20 +555,25 @@ def reviewed_backtest_gates(proposal): optimize=optimize if optimize is not None else bounded_optimize, enforce_backtest_gates=reviewed_backtest_gates, record_shadow=record_shadow, sync_console=sync_console, + research_identity=research_identity, ticket_dir=Path(ticket_dir) / "research_promotion_tickets", + diagnose=diagnose, pull_console=pull_console, + resume_delivery_only=resume_delivery_only, + admit_new_research=admit_new_research, + read_pending_shadow=read_pending_shadow, ) entry["research_promotion_state"] = summary["status"] entry["console_synced"] = summary.get("console_synced") + entry["reason"] = summary["reason"] + for key in ("research_key", "resumed", "retry_at"): + if key in summary: + entry[key] = summary[key] if "ticket" not in summary: return {**entry, "reason": summary["reason"]} - ticket = ResearchPromotionTicket.from_dict(summary["ticket"]) - entry["research_promotion_ticket_id"] = ticket.ticket_id - entry["requires_human_approval"] = ticket.state.value == "awaiting_human" - entry["research_promotion_notes"] = list(ticket.notes) - ticket_dir = getattr(store, "local_root", None) - if ticket_dir is not None: - path = Path(ticket_dir) / "research_promotion_tickets" / f"{ticket.ticket_id}.json" - save_research_promotion_ticket(ticket, path) - entry["research_promotion_ticket_path"] = str(path) + ticket = summary["ticket"] + entry["research_promotion_ticket_id"] = ticket["ticket_id"] + entry["requires_human_approval"] = ticket["state"] == "awaiting_human" + entry["research_promotion_notes"] = ticket["notes"] + entry["research_promotion_ticket_path"] = summary["ticket_path"] except Exception: # No provider details or automatic retries, including uncertain writes. entry["research_error"] = "research_cycle_failed" @@ -583,6 +592,9 @@ def run_auto_pilot_cycle( record_shadow: Callable | None = None, sync_console: Callable | None = None, pull_console: Callable | None = None, + research_identity: Mapping[str, str] | None = None, + admit_new_research: Callable[[Path, str], bool] | None = None, + read_pending_shadow: Callable | None = None, ) -> dict[str, Any]: """Run one drift-triggered cycle with candidate-bound research callbacks. @@ -602,6 +614,7 @@ def run_auto_pilot_cycle( # Recovery does not re-run optimizers, POST tickets or apply human intent. # Historical terminal tickets do not block a future independent observation. pending_profiles: set[str] = set() + undelivered_profiles: set[str] = set() ticket_root = getattr(store, "local_root", None) if not dry_run and isinstance(ticket_root, (str, Path)): from quant_platform_kit.strategy_lifecycle.research_promotion_cycle import ( @@ -622,6 +635,8 @@ def run_auto_pilot_cycle( result["state"] == "awaiting_human" or result["status"] == "updated" ) and result["strategy_profile"]: pending_profiles.add(result["strategy_profile"]) + if result["state"] == "awaiting_human" and result.get("delivery_unstarted") is True: + undelivered_profiles.add(result["strategy_profile"]) # Phase 1 snapshots = _run_monitor_phase(domain, store) @@ -641,7 +656,8 @@ def run_auto_pilot_cycle( for drift in alert_drifts: if drift.status not in (DriftStatus.REVIEW, DriftStatus.CRITICAL): continue - if drift.strategy_profile in pending_profiles: + resume_delivery_only = drift.strategy_profile in pending_profiles + if resume_delivery_only and drift.strategy_profile not in undelivered_profiles: summary["actions"].append({ "strategy_profile": drift.strategy_profile, "action": "skipped", "reason": "saved_research_ticket_pending", @@ -653,6 +669,10 @@ def run_auto_pilot_cycle( drift, store, dry_run, optimize=optimize, enforce_backtest_gates=enforce_backtest_gates, record_shadow=record_shadow, sync_console=sync_console, + pull_console=pull_console, research_identity=research_identity, + resume_delivery_only=resume_delivery_only, + admit_new_research=admit_new_research, + read_pending_shadow=read_pending_shadow, ) ) diff --git a/src/quant_platform_kit/strategy_lifecycle/paired_shadow_adapter.py b/src/quant_platform_kit/strategy_lifecycle/paired_shadow_adapter.py index f0c25e5..aab1a1d 100644 --- a/src/quant_platform_kit/strategy_lifecycle/paired_shadow_adapter.py +++ b/src/quant_platform_kit/strategy_lifecycle/paired_shadow_adapter.py @@ -83,6 +83,8 @@ def collect_paired_shadow_for_promotion( evidence, policy=payload["policy"], forward_observation_receipt=payload["forward_observation_receipt"], + previous_evidence=payload.get("previous_evidence"), + previous_forward_observation_receipt=payload.get("previous_forward_observation_receipt"), ) record["adapter"] = "paired_shadow_adapter.v1" record["live_authority_granted"] = False diff --git a/src/quant_platform_kit/strategy_lifecycle/promotion_actionable_runner.py b/src/quant_platform_kit/strategy_lifecycle/promotion_actionable_runner.py index 00c8657..04129ec 100644 --- a/src/quant_platform_kit/strategy_lifecycle/promotion_actionable_runner.py +++ b/src/quant_platform_kit/strategy_lifecycle/promotion_actionable_runner.py @@ -10,6 +10,7 @@ import json from collections.abc import Callable, Mapping, Sequence from datetime import date +from pathlib import Path from typing import Any from quant_platform_kit.strategy_lifecycle.contracts import DriftResult, DriftStatus @@ -23,6 +24,7 @@ ResearchPromotionTicket, make_console_research_promotion_sync, run_research_promotion_cycle, + run_saved_research_promotion_cycle, ) @@ -69,6 +71,13 @@ def run_actionable_research_promotion( record_shadow: Callable[[Any], Mapping[str, Any]] | None = None, enforce_backtest_gates: Callable[[Any], Any] | None = None, sync_console: Callable[[ResearchPromotionTicket], bool] | None = None, + pull_console: Callable | None = None, + research_identity: Mapping[str, str] | None = None, + ticket_dir: str | Path | None = None, + diagnose: Callable | None = None, + resume_delivery_only: bool = False, + admit_new_research: Callable[[Path, str], bool] | None = None, + read_pending_shadow: Callable | None = None, cycle: Callable[..., ResearchPromotionTicket] | None = None, ) -> dict[str, Any]: """Run promotion once for REVIEW/CRITICAL drift and require paired shadow. @@ -114,7 +123,12 @@ def run_actionable_research_promotion( if not from_store: health["source_revision"] = source_revision - if not health["actionable"]: + # Stale observations may identify an already-computed pending ticket. The + # saved entry alone decides whether its read-only tail can continue. + stale_checkpoint = (health.get("reason") == "observation_stale" + and ticket_dir is not None and research_identity is not None + and health.get("risk_status") in {"review", "critical"}) + if not health["actionable"] and not stale_checkpoint: return { **health, "drift_status": health["status"], @@ -146,7 +160,7 @@ def run_actionable_research_promotion( domain=domain, as_of=resolved_as_of, drift_score=float(health["score"]), - status=DriftStatus(str(health["status"])), + status=DriftStatus(str(health["risk_status"] if stale_checkpoint else health["status"])), source_revision=health.get("source_revision") or "", baseline_artifact_id=health.get("baseline_artifact_id"), baseline_param_set_id=health.get("baseline_param_set_id"), @@ -171,6 +185,23 @@ def deliver_to_console(ticket: ResearchPromotionTicket) -> bool: console_synced = False return console_synced + if (ticket_dir is not None or research_identity is not None or diagnose is not None + or resume_delivery_only or admit_new_research is not None or read_pending_shadow is not None): + if ticket_dir is None or research_identity is None or cycle is not None: + return {**health, "status": "parked", "reason": "research_identity_unavailable", + "console_synced": None} + saved = run_saved_research_promotion_cycle( + drift, research_identity=research_identity, ticket_dir=ticket_dir, + optimize=optimize or _bounded_optimize, enforce_backtest_gates=enforce_backtest_gates, + record_shadow=record_shadow, diagnose=diagnose, budget=budget, + sync_console=deliver_to_console, pull_console=pull_console, + evaluation_date=evaluation_date, max_age_days=max_age_days, + resume_delivery_only=resume_delivery_only, + admit_new_research=admit_new_research, + read_pending_shadow=read_pending_shadow, + ) + return {**health, **saved} + ticket = (cycle or run_research_promotion_cycle)( drift, optimize=optimize or _bounded_optimize, 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 a152ae6..47c25ef 100644 --- a/src/quant_platform_kit/strategy_lifecycle/research_promotion_cycle.py +++ b/src/quant_platform_kit/strategy_lifecycle/research_promotion_cycle.py @@ -8,10 +8,13 @@ from __future__ import annotations import calendar +import hashlib import json +import math import uuid from collections.abc import Callable, Mapping -from dataclasses import asdict, dataclass, field +from contextlib import contextmanager +from dataclasses import asdict, dataclass, field, replace from datetime import date, datetime, timedelta, timezone from enum import Enum from pathlib import Path @@ -157,9 +160,13 @@ class ResearchPromotionTicket: confirmation_execution_mode: str = "" confirmation_risk_profile: str = "" notes: tuple[str, ...] = () + # Local checkpoint only. QRT's candidate contract must not carry job state. + research_progress: Mapping[str, Any] = field(default_factory=dict, repr=False) - def to_dict(self) -> dict[str, Any]: + def to_dict(self, *, include_progress: bool = False) -> dict[str, Any]: payload = asdict(self) + if not include_progress: + payload.pop("research_progress") payload["state"] = self.state.value payload["notes"] = list(self.notes) payload["proposed_params"] = dict(self.proposed_params) @@ -202,6 +209,7 @@ def from_dict(cls, raw: Mapping[str, Any]) -> ResearchPromotionTicket: ), confirmation_risk_profile=str(raw.get("confirmation_risk_profile") or ""), notes=tuple(str(item) for item in (raw.get("notes") or ())), + research_progress=dict(raw.get("research_progress") or {}), ) @@ -400,6 +408,8 @@ def shadow_record_from_paired_evidence( *, policy: Any | None = None, forward_observation_receipt: Mapping[str, Any] | None = None, + previous_evidence: Mapping[str, Any] | None = None, + previous_forward_observation_receipt: Mapping[str, Any] | None = None, ) -> dict[str, Any]: """Validate paired-shadow evidence into a non-live cycle shadow record.""" from quant_platform_kit.strategy_lifecycle.paired_shadow_evidence import ( @@ -411,6 +421,8 @@ def shadow_record_from_paired_evidence( evidence, policy=policy, forward_observation_receipt=forward_observation_receipt, + previous_evidence=previous_evidence, + previous_forward_observation_receipt=previous_forward_observation_receipt, ) if validated.get("live_authority_granted") is True: raise ValueError("paired shadow evidence must not grant live authority") @@ -801,7 +813,7 @@ def apply_console_research_promotion_decision( if datetime.fromisoformat(remote_decided_at.replace("Z", "+00:00")).tzinfo is None: raise ValueError("console decision time requires timezone") decided = apply_human_promotion_decision( - ResearchPromotionTicket.from_dict(ticket.to_dict()), + ResearchPromotionTicket.from_dict(ticket.to_dict(include_progress=True)), decision=decision, confirmation=confirmation, paper_supported=paper_supported, @@ -820,6 +832,28 @@ def reconcile_saved_research_promotion_ticket( pull_console: Callable[[str], Mapping[str, Any] | None] | None = None, output_path: str | Path | None = None, domain: str | None = None, +) -> dict[str, Any]: + """Serialize decision recovery with research and manual decisions.""" + result = {"ticket_id": None, "strategy_profile": None, "state": None, + "status": "unavailable", "reason": "local_ticket_unavailable", + "live_authority_granted": False} + try: + with _research_directory_lock(Path(ticket_path).parent) as acquired: + if not acquired: + return {**result, "status": "deferred", "reason": "research_in_progress"} + return _reconcile_saved_research_promotion_ticket( + ticket_path, pull_console=pull_console, output_path=output_path, domain=domain, + ) + except (OSError, TypeError, ValueError): + return result + + +def _reconcile_saved_research_promotion_ticket( + ticket_path: str | Path, + *, + pull_console: Callable[[str], Mapping[str, Any] | None] | None = None, + output_path: str | Path | None = None, + domain: str | None = None, ) -> dict[str, Any]: """Recover one saved human decision without starting research or applying it. @@ -843,6 +877,10 @@ def reconcile_saved_research_promotion_ticket( return {**result, "status": "already_terminal", "reason": "local_ticket_terminal"} if ticket.state is not ResearchPromotionState.AWAITING_HUMAN: return {**result, "status": "skipped", "reason": "ticket_not_awaiting_human"} + result["delivery_unstarted"] = ( + isinstance(ticket.research_progress.get("identity"), Mapping) + and "console_delivery" not in ticket.research_progress + ) try: remote = (pull_console or make_console_research_promotion_pull())(ticket.ticket_id) except Exception: @@ -895,6 +933,11 @@ def _attach_shadow_or_park( kind = _shadow_kind(shadow) ticket.shadow_evidence_kind = kind ticket.shadow_passed = bool(shadow.get("passed", False)) + if (shadow.get("status") == "pending" and shadow.get("passed") is False + and shadow.get("no_order") is True and shadow.get("live_authority_granted") is False): + ticket.state = ResearchPromotionState.SHADOW_RECORDED + ticket.notes = ticket.notes + ("paired_shadow_observation_pending",) + return ticket digest = str(shadow.get("paired_shadow_evidence_sha256") or "").strip() if digest: ticket.notes = ticket.notes + (f"paired_shadow_evidence_sha256={digest}",) @@ -956,6 +999,10 @@ def run_research_promotion_cycle( ticket.state = ResearchPromotionState.BOUNDED_REOPT proposal = optimize(drift, budget) + if (proposal.strategy_profile != drift.strategy_profile or proposal.domain != drift.domain): + ticket.state = ResearchPromotionState.PARKED + ticket.notes = ("proposal_target_mismatch",) + return ticket ok, reason = enforce_optimization_budget(proposal, budget) ticket.search_iterations = int(proposal.search_iterations) ticket.proposed_params = dict(proposal.proposed_params or {}) @@ -1119,7 +1166,7 @@ def save_research_promotion_ticket( target = Path(path) target.parent.mkdir(parents=True, exist_ok=True) - payload = json.dumps(ticket.to_dict(), indent=2, sort_keys=True, allow_nan=False) + "\n" + payload = json.dumps(ticket.to_dict(include_progress=True), indent=2, sort_keys=True, allow_nan=False) + "\n" temporary: Path | None = None try: with tempfile.NamedTemporaryFile( @@ -1147,6 +1194,417 @@ def load_research_promotion_ticket(path: str | Path) -> ResearchPromotionTicket: return ticket +@contextmanager +def _research_directory_lock(directory: Path): + """One process at a time for this shared directory, never a global queue.""" + import os + + directory.mkdir(parents=True, exist_ok=True) + with (directory / ".research.lock").open("a+b") as stream: + if os.name == "nt": + import msvcrt + + if stream.tell() == 0: + stream.write(b"0") + stream.flush() + stream.seek(0) + try: + msvcrt.locking(stream.fileno(), msvcrt.LK_NBLCK, 1) + except OSError: + yield False + return + try: + yield True + finally: + stream.seek(0) + msvcrt.locking(stream.fileno(), msvcrt.LK_UNLCK, 1) + else: + import fcntl + + try: + fcntl.flock(stream.fileno(), fcntl.LOCK_EX | fcntl.LOCK_NB) + except BlockingIOError: + yield False + return + try: + yield True + finally: + fcntl.flock(stream.fileno(), fcntl.LOCK_UN) + + +def _saved_proposal(raw: Mapping[str, Any]) -> OptimizationProposal: + from quant_platform_kit.strategy_lifecycle.contracts import BacktestValidationIdentity + from quant_platform_kit.strategy_lifecycle.performance_store import _backtest_from_dict + + values = dict(raw) + for key in ("current_metrics", "proposed_metrics"): + if values.get(key) is not None: + stored = values[key] + validation = stored.get("validation_identity") + if validation is not None: + validation = dict(validation) + for date_key in ("train_start", "train_end", "test_start", "test_end", + "locked_oos_start", "locked_oos_end"): + validation[date_key] = date.fromisoformat(validation[date_key]) if validation[date_key] else None + validation = BacktestValidationIdentity(**validation) + values[key] = replace(_backtest_from_dict(stored), validation_identity=validation, + cost_inputs=dict(stored.get("cost_inputs") or {})) + for key in ("winning_dimensions", "regressing_dimensions"): + values[key] = tuple(values.get(key) or ()) + proposal = OptimizationProposal(**values) + if _canonical_ticket_value(proposal.to_dict()) != _canonical_ticket_value(raw): + raise ValueError("saved_proposal_invalid") + return proposal + + +def run_saved_research_promotion_cycle( + drift: DriftResult, + *, + research_identity: Mapping[str, str], + ticket_dir: str | Path, + optimize: Callable, + enforce_backtest_gates: Callable, + record_shadow: Callable, + diagnose: Callable | None = None, + sync_console: Callable | None = None, + pull_console: Callable | None = None, + budget: ResearchPromotionBudget | None = None, + evaluation_date: date | str | None = None, + max_age_days: int = 7, + resume_delivery_only: bool = False, + admit_new_research: Callable[[Path, str], bool] | None = None, + read_pending_shadow: Callable[[OptimizationProposal], Mapping[str, Any]] | None = None, +) -> dict[str, Any]: + """Resume a bound experiment using the existing ticket's local progress. + + Revisions refer to the caller's validated, frozen input, code, parameter + space (including baseline/fixed parameters), cost model and validator. They + are not evidence validation themselves. Only processes sharing ticket_dir + share the nonblocking lock; the dispatcher owns cross-host admission. + + Each side-effecting stage is saved as running BEFORE invocation. Completed + results are reused; running/unknown outcomes never get automatically called + again. A definite diagnosis deferral may run again only after retry_at. + Console writes are attempted once and thereafter recovered only by GET. + Optional admission runs under this lock only for a new ticket; its directory + and UTC created_at match the eventual ticket. It must not acquire this lock + again or maintain a second counter. Only literal True permits creation. + Explicit pending shadow may use read_pending_shadow after its deadline. + A stale original observation only permits this already-computed, identity- + matched tail; it never permits new AI, optimization or backtest calls. + """ + from quant_platform_kit.strategy_lifecycle.production_drift_health_probe import ( + probe_production_drift_health, + ) + + base = {"status": "parked", "reason": "research_input_invalid", "resumed": False, + "console_synced": None, "live_authority_granted": False} + if not all(callable(fn) for fn in (optimize, enforce_backtest_gates, record_shadow)): + return {**base, "reason": "research_bindings_unavailable"} + fields = {"code_revision", "input_revision", "param_space_revision", + "cost_model_revision", "validator_revision"} + if (not isinstance(research_identity, Mapping) or set(research_identity) != fields + or any(not isinstance(v, str) or not v.strip() for v in research_identity.values()) + or not isinstance(drift.source_revision, str) or not drift.source_revision.strip()): + return {**base, "reason": "research_identity_unavailable"} + budget = budget or ResearchPromotionBudget(require_paired_shadow=True) + if (type(budget.max_search_iterations) is not int or not 1 <= budget.max_search_iterations <= 25 + or type(budget.max_param_keys) is not int or not 1 <= budget.max_param_keys <= 4 + or budget.require_paired_shadow is not True or budget.allow_live_enablement): + return {**base, "reason": "research_budget_invalid"} + try: + health = probe_production_drift_health( + strategy_profile=drift.strategy_profile, domain=drift.domain, + as_of=drift.as_of, drift_score=drift.drift_score, + evaluation_date=evaluation_date, max_age_days=max_age_days, + ) + except (ValueError, TypeError, OverflowError): + return base + stale = health.get("reason") == "observation_stale" + if ((not health["actionable"] and not stale) + or drift.status not in _ACTIVE_DRIFT or drift.alert_suppressed): + return {**base, "reason": health.get("reason", "drift_not_actionable")} + observation = drift.to_dict() + identity = {"revisions": dict(research_identity), + "drift": {key: observation[key] for key in ( + "strategy_profile", "domain", "as_of", "source_revision", "drift_score", + "baseline_param_set_id", "baseline_param_version", "baseline_artifact_id")}, + "budget": asdict(budget)} + try: + canonical = _canonical_ticket_value(identity) + research_key = hashlib.sha256(canonical.encode()).hexdigest() + directory = Path(ticket_dir) + except (ValueError, TypeError): + return base + ticket_id = f"rpt_{research_key}" + path = directory / f"{ticket_id}.json" + base.update(research_key=research_key, ticket_path=str(path)) + + try: + with _research_directory_lock(directory) as acquired: + if not acquired: + return {**base, "status": "deferred", "reason": "research_in_progress"} + if path.exists(): + ticket = load_research_promotion_ticket(path) + progress = dict(ticket.research_progress) + if (ticket.ticket_id != ticket_id or ticket.strategy_profile != drift.strategy_profile + or ticket.domain != drift.domain + or _canonical_ticket_value(progress.get("identity")) != canonical): + return {**base, "reason": "research_checkpoint_mismatch"} + base["resumed"] = True + else: + if stale: + return {**base, "reason": "observation_stale"} + if resume_delivery_only: + return {**base, "reason": "saved_research_ticket_pending"} + now = _now_iso() + if admit_new_research is not None: + try: + admitted = admit_new_research(directory, now) + except Exception: + return {**base, "reason": "research_admission_unavailable"} + if admitted is not True: + return {**base, "status": "deferred", "reason": "new_research_not_admitted"} + progress = {"identity": identity, "stages": {}, "diagnosis_required": diagnose is not None} + ticket = ResearchPromotionTicket( + ticket_id=ticket_id, strategy_profile=drift.strategy_profile, + domain=drift.domain, state=ResearchPromotionState.BOUNDED_REOPT, + drift_status=drift.status.value, drift_score=drift.drift_score, + created_at=now, updated_at=now, budget=asdict(budget), + research_progress=progress, + ) + save_research_promotion_ticket(ticket, path) + + def output(reason: str, *, status: str | None = None) -> dict[str, Any]: + return {**base, "status": status or ticket.state.value, "reason": reason, + "ticket": ticket.to_dict()} + + def deliver_saved_ticket() -> None: + if (ticket.state != ResearchPromotionState.AWAITING_HUMAN or sync_console is None + or progress.get("console_delivery") is not None): + return + progress["console_delivery"] = "running" + save_research_promotion_ticket(ticket, path) + try: + candidate = ResearchPromotionTicket.from_dict(ticket.to_dict()) + confirmed = sync_console(candidate) is True + base["console_synced"] = ( + confirmed and _canonical_ticket_value(candidate.to_dict()) + == _canonical_ticket_value(ticket.to_dict()) + ) + except Exception: + base["console_synced"] = False + progress["console_delivery"] = "confirmed" if base["console_synced"] else "unconfirmed" + save_research_promotion_ticket(ticket, path) + + if resume_delivery_only and ticket.state != ResearchPromotionState.AWAITING_HUMAN: + return output("saved_research_ticket_pending", status="parked") + stages = progress["stages"] + if not isinstance(stages, dict): + return output("research_checkpoint_invalid") + if any(value.get("status") in {"running", "unknown"} for value in stages.values()): + return output("research_outcome_unknown", status="parked") + if ticket.state in _TERMINAL: + return output("saved_research_ticket_terminal") + shadow_pending = stages.get("shadow", {}).get("status") == "pending" + shadow_completed = (ticket.state == ResearchPromotionState.SHADOW_RECORDED + and stages.get("shadow", {}).get("status") == "completed" + and stages["shadow"].get("result", {}).get("status") == "complete") + awaiting_recovery = stale and ticket.state == ResearchPromotionState.AWAITING_HUMAN + if shadow_pending or shadow_completed or awaiting_recovery: + # A pending observation never permits incomplete earlier work + # to restart, even while the original drift is still fresh. + required = ("optimize", "backtest") + (("diagnose",) if progress.get("diagnosis_required") else ()) + if (ticket.state not in {ResearchPromotionState.SHADOW_RECORDED, ResearchPromotionState.AWAITING_HUMAN} + or any(stages.get(name, {}).get("status") != "completed" for name in required)): + return output("research_checkpoint_invalid", status="parked") + proposal = _saved_proposal(stages["optimize"]["result"]) + if awaiting_recovery: + recorded = stages.get("shadow", {}) + shadow = recorded.get("result", {}) + if (recorded.get("status") != "completed" or shadow.get("passed") is not True + or not _is_paired_shadow_kind(_shadow_kind(shadow)) + or shadow.get("live_authority_granted") is True or ticket.shadow_passed is not True + or _canonical_ticket_value(ticket.proposed_params) != _canonical_ticket_value(proposal.proposed_params)): + 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" + 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)): + return output("research_checkpoint_invalid", status="parked") + if shadow_pending: + pending_result = stages["shadow"].get("result", {}) + if (pending_result.get("status") != "pending" or pending_result.get("passed") is not False + or pending_result.get("no_order") is not True or pending_result.get("live_authority_granted") is not False): + return output("research_checkpoint_invalid", status="parked") + retry = pending_result.get("retry_at") + if (type(retry) not in (int, float) or not math.isfinite(retry) + or retry > datetime.now(timezone.utc).timestamp()): + return {**output("paired_shadow_observation_pending", status="deferred"), "retry_at": retry} + if not callable(read_pending_shadow): + return output("shadow_reader_unavailable", status="deferred") + elif stale and not (shadow_completed or awaiting_recovery): + return output("observation_stale", status="parked") + if ticket.state == ResearchPromotionState.AWAITING_HUMAN: + reconciliation = _reconcile_saved_research_promotion_ticket(path, pull_console=pull_console) + ticket = load_research_promotion_ticket(path) + progress = ticket.research_progress + if not stale: + deliver_saved_ticket() + return {**output("saved_research_ticket_reused"), "reconciliation": reconciliation} + + def stage(name: str, callback: Callable, encode: Callable = lambda value: value): + existing = stages.get(name) + if existing is not None: + if existing.get("status") == "completed": + return existing["result"] + if not (name == "shadow" and shadow_pending and existing.get("status") == "pending"): + raise ValueError("research_checkpoint_invalid") + stages[name] = {"status": "running"} + save_research_promotion_ticket(ticket, path) + try: + result = json.loads(_canonical_ticket_value(encode(callback()))) + except Exception: + stages[name] = {"status": "unknown"} + save_research_promotion_ticket(ticket, path) + raise ValueError("research_outcome_unknown") from None + status = "pending" if name == "shadow" and result.get("status") == "pending" else "completed" + stages[name] = {"status": status, "result": result} + if name == "shadow" and result.get("status") in {"pending", "complete"}: + # Save the resumable state with the result, not in a later + # write after the caller may already have stopped. + ticket.state = ResearchPromotionState.SHADOW_RECORDED + save_research_promotion_ticket(ticket, path) + return result + + if progress.get("diagnosis_required") is True: + if diagnose is None: + return output("research_bindings_unavailable", status="parked") + old = stages.get("diagnose", {}) + if old.get("status") == "deferred": + retry = old.get("retry_at") + if type(retry) not in (int, float) or retry > datetime.now(timezone.utc).timestamp(): + return {**output("codex_research_deferred", status="deferred"), "retry_at": retry} + del stages["diagnose"] + def checked_diagnosis(): + value = dict(diagnose(drift, budget)) + if value.get("reason") == "codex_research_deferred": + retry = value.get("retry_at") + # A reset already in the past is not a usable admission + # deadline. Preserve deferral without polling every run. + if type(retry) not in (int, float) or retry <= datetime.now(timezone.utc).timestamp(): + value["retry_at"] = None + return value + + decision = stage("diagnose", checked_diagnosis) + if not isinstance(decision, Mapping): + raise ValueError("research_checkpoint_invalid") + if decision.get("reason") == "codex_research_deferred": + retry = decision.get("retry_at") + stages["diagnose"] = {"status": "deferred", "retry_at": retry} + save_research_promotion_ticket(ticket, path) + return {**output("codex_research_deferred", status="deferred"), "retry_at": retry} + if decision.get("optimization_needed") is not True: + ticket.state = ResearchPromotionState.PARKED + ticket.notes = ("codex_did_not_recommend_research",) + save_research_promotion_ticket(ticket, path) + return output("codex_did_not_recommend_research") + + def cached_optimize(*_): + raw = stage("optimize", lambda: optimize(drift, budget), lambda value: value.to_dict()) + return _saved_proposal(raw) + + def shadow_result(proposal): + # The first callback may create an external observation. Only + # the dedicated, caller-owned read-only callback can be polled. + callback = read_pending_shadow if shadow_pending else record_shadow + + def encode(value): + value = dict(value) + if value.get("live_authority_granted") is True: + raise ValueError("shadow_live_authority_refused") + if value.get("status") == "pending": + if (value.get("passed") is not False or value.get("no_order") is not True + or value.get("live_authority_granted") is not False): + raise ValueError("invalid_shadow_pending") + retry = value.get("retry_at") + if (type(retry) not in (int, float) or not math.isfinite(retry) + or retry <= datetime.now(timezone.utc).timestamp()): + retry = None + return {"status": "pending", "retry_at": retry, "evidence_kind": "paired_shadow_pending", + "passed": False, "no_order": True, "live_authority_granted": False} + if value.get("status") == "complete": + from quant_platform_kit.strategy_lifecycle.paired_shadow_adapter import ( + PairedShadowObservation, collect_paired_shadow_for_promotion, + ) + observation = value["observation"] + policy = observation.policy if isinstance(observation, PairedShadowObservation) else observation["policy"] + receipt = (observation.forward_observation_receipt if isinstance(observation, PairedShadowObservation) + else observation["forward_observation_receipt"]) + if (policy.strategy_profile != proposal.strategy_profile or policy.domain != proposal.domain + or receipt["observation_index"] < policy.required_trading_sessions): + raise ValueError("shadow_window_not_complete") + # The owner also validates the frozen candidate/params/ + # source and full window. A bare passed flag is not proof. + return {**collect_paired_shadow_for_promotion(observation), "status": "complete"} + if shadow_pending and value.get("passed") is not False: + raise ValueError("shadow_completion_evidence_required") + return value + + return stage("shadow", lambda: callback(proposal), encode) + + result = run_research_promotion_cycle( + drift, budget=budget, ticket_id=ticket_id, optimize=cached_optimize, + enforce_backtest_gates=lambda proposal: stage("backtest", lambda: enforce_backtest_gates(proposal), + lambda value: _promotion_backtest_evidence_mapping(value) if value is not None else None), + record_shadow=shadow_result, + ) + if any(value.get("status") in {"running", "unknown"} for value in stages.values()): + return output("research_outcome_unknown", status="parked") + result.created_at = ticket.created_at + result.research_progress = progress + ticket = result + save_research_promotion_ticket(ticket, path) + deliver_saved_ticket() + if stages.get("shadow", {}).get("status") == "pending": + return {**output("paired_shadow_observation_pending", status="deferred"), + "retry_at": stages["shadow"]["result"]["retry_at"]} + return output("promotion_cycle_completed") + except Exception: + # Persisted running stage is deliberately left as unknown. Never include + # exceptions, input rows, credentials or a provider's response text. + try: + saved = load_research_promotion_ticket(path) + if any(value.get("status") in {"running", "unknown"} + for value in saved.research_progress.get("stages", {}).values()): + return {**base, "reason": "research_outcome_unknown"} + except Exception: + pass + return {**base, "reason": "research_checkpoint_unavailable"} + + +def decide_saved_research_promotion_ticket( + ticket_path: str | Path, + *, + decision: str, + confirmation: PromotionConfirmation | Mapping[str, Any] | None = None, + paper_supported: bool = False, + output_path: str | Path | None = None, +) -> ResearchPromotionTicket: + """Record explicit local human intent under the same research-directory lock.""" + with _research_directory_lock(Path(ticket_path).parent) as acquired: + if not acquired: + raise ValueError("research_in_progress") + ticket = load_research_promotion_ticket(ticket_path) + decided = apply_human_promotion_decision( + ticket, decision=decision, confirmation=confirmation, paper_supported=paper_supported, + ) + save_research_promotion_ticket(decided, output_path or ticket_path) + return decided + + __all__ = [ "DEFAULT_SUGGESTED_RISK_PROFILE", "EXECUTION_MODES", @@ -1158,6 +1616,7 @@ def load_research_promotion_ticket(path: str | Path) -> ResearchPromotionTicket: "apply_console_research_promotion_decision", "apply_human_promotion_decision", "build_human_promotion_notification", + "decide_saved_research_promotion_ticket", "enforce_optimization_budget", "enforce_promotion_backtest_gates", "load_research_promotion_ticket", @@ -1167,6 +1626,7 @@ def load_research_promotion_ticket(path: str | Path) -> ResearchPromotionTicket: "open_awaiting_human_ticket", "reconcile_saved_research_promotion_ticket", "run_research_promotion_cycle", + "run_saved_research_promotion_cycle", "save_research_promotion_ticket", "shadow_record_from_paired_evidence", "validate_promotion_confirmation", diff --git a/tests/test_lifecycle_codex_integration.py b/tests/test_lifecycle_codex_integration.py index 05f0cfe..17956da 100644 --- a/tests/test_lifecycle_codex_integration.py +++ b/tests/test_lifecycle_codex_integration.py @@ -21,6 +21,7 @@ def _fixed_probe_clock(): create_github_issue, ) from quant_platform_kit.strategy_lifecycle.contracts import DriftResult, DriftStatus, StrategyPerformanceSnapshot +from tests.test_research_promotion_resume import IDENTITY def test_drift_phase_excludes_suppressed_results_from_automation() -> None: @@ -141,7 +142,7 @@ def sync(ticket): patch(prefix + "call_ai_optimization_decision", side_effect=decision)): result = run_auto_pilot_cycle("us_equity", store=store, create_issues=False, optimize=optimize, enforce_backtest_gates=gates, - record_shadow=shadow, sync_console=sync) + record_shadow=shadow, sync_console=sync, research_identity=IDENTITY) action = result["actions"][0] assert events == (["codex", "optimize", "gates", "shadow", "console"] if gate_passes and risk_passes else ["codex", "optimize", "gates"] if risk_passes else ["codex", "optimize"]) @@ -195,19 +196,21 @@ def test_optimization_decision_is_codex_only_with_no_paid_fallback(success, outp assert result["optimization_needed"] is needed -def test_codex_deferral_is_pending_instead_of_a_negative_research_recommendation(): +def test_codex_deferral_is_pending_instead_of_a_negative_research_recommendation(tmp_path): + from datetime import datetime, timezone from quant_platform_kit.strategy_lifecycle.codex_integration import _process_optimization_decision - store = Mock() + retry_at = datetime.now(timezone.utc).timestamp() + 3600 + store = Mock(local_root=tmp_path) store.load_latest_snapshot.return_value = StrategyPerformanceSnapshot( strategy_profile="demo_strategy", domain="us_equity", platform="test", as_of=date(2026, 9, 7), source_revision="source-v1") optimize = Mock() with patch("quant_platform_kit.strategy_lifecycle.ai_provider.AiServiceClient") as factory: factory.return_value.execute.return_value = SimpleNamespace(success=False, provider="codex", output="", - raw={"status": "deferred", "retry_at": 9000}) + raw={"status": "deferred", "retry_at": retry_at}) result = _process_optimization_decision(_critical(), store, False, - optimize=optimize, enforce_backtest_gates=Mock(), record_shadow=Mock()) + optimize=optimize, enforce_backtest_gates=Mock(), record_shadow=Mock(), research_identity=IDENTITY) assert result["research_promotion_state"] == "deferred" - assert result["retry_at"] == 9000 + assert result["retry_at"] == retry_at optimize.assert_not_called() diff --git a/tests/test_research_promotion_reconciliation.py b/tests/test_research_promotion_reconciliation.py index c5e981c..43a4807 100644 --- a/tests/test_research_promotion_reconciliation.py +++ b/tests/test_research_promotion_reconciliation.py @@ -225,7 +225,7 @@ def test_reconciliation_save_failure_keeps_original_and_can_recover(tmp_path): assert result["status"] == "unavailable" assert result["reason"] == "local_ticket_save_failed" assert path.read_bytes() == before - assert list(tmp_path.iterdir()) == [path] + assert set(tmp_path.iterdir()) == {path, tmp_path / ".research.lock"} assert cycle.reconcile_saved_research_promotion_ticket(path, pull_console=pull)["status"] == "updated" diff --git a/tests/test_research_promotion_resume.py b/tests/test_research_promotion_resume.py new file mode 100644 index 0000000..3478ac9 --- /dev/null +++ b/tests/test_research_promotion_resume.py @@ -0,0 +1,524 @@ +"""Durable research stages, using synthetic evidence and no external services.""" + +from dataclasses import replace +from datetime import date, datetime, timezone +from unittest.mock import Mock, patch + +import pytest + +from quant_platform_kit.strategy_lifecycle import research_promotion_cycle as cycle +from tests.test_research_promotion_cycle import _drift, _proposal, _promotion_backtest_evidence +from tests.test_research_promotion_reconciliation import remote_decision + + +IDENTITY = { + "code_revision": "frozen-code-v1", + "input_revision": "frozen-input-v1", + "param_space_revision": "bounded-space-v1", + "cost_model_revision": "cost-v1", + "validator_revision": "strict-validator-v1", +} + + +def invoke(tmp_path, **overrides): + kwargs = dict( + research_identity=IDENTITY, ticket_dir=tmp_path, + optimize=Mock(return_value=_proposal()), + enforce_backtest_gates=Mock(return_value=_promotion_backtest_evidence()), + record_shadow=Mock(return_value={"evidence_kind": "paired_shadow", "passed": True, + "paired_shadow_evidence_sha256": "a" * 64}), + sync_console=Mock(return_value=True), pull_console=Mock(return_value=None), + evaluation_date=date(2026, 9, 8), + ) + kwargs.update(overrides) + drift = kwargs.pop("drift", replace(_drift(), source_revision="observation-v1")) + return cycle.run_saved_research_promotion_cycle(drift, **kwargs) + + +def test_same_frozen_input_does_not_repeat_completed_stages_or_console_post(tmp_path): + optimize, gates, shadow, sync = (Mock(return_value=value) for value in + (_proposal(), _promotion_backtest_evidence(), + {"evidence_kind": "paired_shadow", "passed": True}, True)) + first = invoke(tmp_path, optimize=optimize, enforce_backtest_gates=gates, + record_shadow=shadow, sync_console=sync) + second = invoke(tmp_path, optimize=optimize, enforce_backtest_gates=gates, + record_shadow=shadow, sync_console=sync) + assert first["status"] == second["status"] == "awaiting_human" + assert first["ticket"]["ticket_id"] == second["ticket"]["ticket_id"] + assert second["resumed"] is True + assert [fn.call_count for fn in (optimize, gates, shadow, sync)] == [1, 1, 1, 1] + assert "research_progress" not in first["ticket"] # local state is never QRT material + + +@pytest.mark.parametrize("field", list(IDENTITY)) +def test_changed_binding_starts_a_new_experiment(tmp_path, field): + first = invoke(tmp_path) + second = invoke(tmp_path, research_identity={**IDENTITY, field: "changed-v2"}) + assert first["research_key"] != second["research_key"] + assert first["ticket"]["ticket_id"] != second["ticket"]["ticket_id"] + + +def test_restart_after_saved_proposal_resumes_gate_without_optimizer(tmp_path): + save = cycle.save_research_promotion_ticket + + def stop_after_proposal(ticket, path): + result = save(ticket, path) + if ticket.research_progress.get("stages", {}).get("optimize", {}).get("status") == "completed": + raise SystemExit("synthetic process stop between stages") + return result + + optimize = Mock(return_value=_proposal()) + gates = Mock(return_value=_promotion_backtest_evidence()) + with patch.object(cycle, "save_research_promotion_ticket", side_effect=stop_after_proposal): + with pytest.raises(SystemExit): + invoke(tmp_path, optimize=optimize, enforce_backtest_gates=gates) + result = invoke(tmp_path, optimize=optimize, enforce_backtest_gates=gates) + assert result["status"] == "awaiting_human" + assert optimize.call_count == gates.call_count == 1 + + +@pytest.mark.parametrize("stage", ["optimize", "enforce_backtest_gates", "record_shadow", "diagnose"]) +def test_unknown_stage_outcome_is_not_retried(tmp_path, stage): + callback = Mock(side_effect=TimeoutError("private data must not escape")) + first = invoke(tmp_path, **{stage: callback}) + second = invoke(tmp_path, **{stage: callback}) + assert first["reason"] == second["reason"] == "research_outcome_unknown" + assert callback.call_count == 1 + assert "private data" not in str(first) + str(second) + + +def test_unknown_console_write_recovers_only_by_get(tmp_path): + sync = Mock(side_effect=TimeoutError("sensitive")) + first = invoke(tmp_path, sync_console=sync) + local = cycle.load_research_promotion_ticket(first["ticket_path"]) + pull = Mock(return_value=remote_decision(local, "reject")) + second = invoke(tmp_path, sync_console=sync, pull_console=pull) + third = invoke(tmp_path, sync_console=sync, pull_console=pull) + assert second["status"] == third["status"] == "human_rejected" + assert sync.call_count == pull.call_count == 1 + assert second["ticket"]["live_authority_granted"] is False + assert cycle.load_research_promotion_ticket(first["ticket_path"]).research_progress + + +def test_same_evidence_rejected_candidate_does_not_restart_ai(tmp_path): + diagnose = Mock(return_value={"optimization_needed": False}) + first = invoke(tmp_path, diagnose=diagnose) + second = invoke(tmp_path, diagnose=diagnose) + assert first["status"] == second["status"] == "parked" + assert diagnose.call_count == 1 + + +@pytest.mark.parametrize("overrides", [ + {"enforce_backtest_gates": None}, {"record_shadow": None}, + {"research_identity": {}}, {"evaluation_date": date(2026, 9, 20)}, + {"drift": replace(_drift(), source_revision="")}, +]) +def test_invalid_or_stale_input_and_missing_bindings_reject_before_ai(tmp_path, overrides): + diagnose, optimize = Mock(), Mock() + result = invoke(tmp_path, diagnose=diagnose, optimize=optimize, **overrides) + assert result["status"] == "parked" + diagnose.assert_not_called() + optimize.assert_not_called() + + +def test_wrong_proposal_target_never_reaches_gate_or_shadow(tmp_path): + gates, shadow = Mock(), Mock() + result = invoke(tmp_path, optimize=Mock(return_value=replace(_proposal(), domain="cn_equity")), + enforce_backtest_gates=gates, record_shadow=shadow) + assert result["status"] == "parked" + gates.assert_not_called() + shadow.assert_not_called() + + +def test_shared_directory_limits_concurrency_before_ai(tmp_path): + nested = [] + + def optimize(*_): + nested.append(invoke(tmp_path, research_identity={**IDENTITY, "input_revision": "next"})) + return _proposal() + + invoke(tmp_path, optimize=optimize) + assert nested[0]["status"] == "deferred" + assert nested[0]["reason"] == "research_in_progress" + + +def test_actionable_runner_uses_saved_cycle(tmp_path): + from quant_platform_kit.strategy_lifecycle.promotion_actionable_runner import run_actionable_research_promotion + + optimize = Mock(return_value=_proposal()) + kwargs = dict(strategy_profile="demo_strategy", domain="us_equity", as_of="2026-09-07", + drift_score=0.8, source_revision="observation-v1", evaluation_date="2026-09-08", + optimize=optimize, enforce_backtest_gates=Mock(return_value=_promotion_backtest_evidence()), + record_shadow=Mock(return_value={"evidence_kind": "paired_shadow", "passed": True}), + sync_console=Mock(return_value=True), pull_console=Mock(return_value=None), + ticket_dir=tmp_path, research_identity=IDENTITY) + first = run_actionable_research_promotion(**kwargs) + second = run_actionable_research_promotion(**kwargs) + assert first["status"] == second["status"] == "awaiting_human" + assert optimize.call_count == 1 + + +def test_normal_cycle_reuses_declined_diagnosis_across_runs(tmp_path): + from quant_platform_kit.strategy_lifecycle.codex_integration import run_auto_pilot_cycle + from quant_platform_kit.strategy_lifecycle.contracts import StrategyPerformanceSnapshot + + store = Mock(local_root=tmp_path) + store.load_latest_snapshot.return_value = StrategyPerformanceSnapshot( + strategy_profile="demo_strategy", domain="us_equity", platform="test", as_of=date(2026, 9, 7), + source_revision="observation-v1") + prefix = "quant_platform_kit.strategy_lifecycle.codex_integration." + with patch(prefix + "_run_monitor_phase", return_value=[]), \ + patch(prefix + "_run_drift_phase", return_value=([_drift()], [_drift()])), \ + patch(prefix + "_drift_freshness_reason", return_value=None), \ + patch(prefix + "call_ai_optimization_decision", return_value={"optimization_needed": False}) as ai: + for _ in range(2): + result = run_auto_pilot_cycle("us_equity", store=store, create_issues=False, + research_identity=IDENTITY, optimize=Mock(), enforce_backtest_gates=Mock(), + record_shadow=Mock(), pull_console=Mock(return_value=None)) + assert result["actions"][0]["research_promotion_state"] == "parked" + assert ai.call_count == 1 + + +def test_automatic_cycle_missing_frozen_identity_does_not_call_ai(tmp_path): + from quant_platform_kit.strategy_lifecycle.codex_integration import _process_optimization_decision + from quant_platform_kit.strategy_lifecycle.contracts import StrategyPerformanceSnapshot + + store = Mock(local_root=tmp_path) + store.load_latest_snapshot.return_value = StrategyPerformanceSnapshot( + strategy_profile="demo_strategy", domain="us_equity", platform="test", as_of=date(2026, 9, 7), + source_revision="observation-v1") + prefix = "quant_platform_kit.strategy_lifecycle.codex_integration." + with patch(prefix + "_drift_freshness_reason", return_value=None), \ + patch(prefix + "call_ai_optimization_decision", return_value={"optimization_needed": False}) as ai: + result = _process_optimization_decision(_drift(), store, False, + optimize=Mock(), enforce_backtest_gates=Mock(), record_shadow=Mock()) + assert result["reason"] == "research_identity_unavailable" + ai.assert_not_called() + + +def test_monitor_bookkeeping_does_not_change_frozen_research_identity(tmp_path): + optimize = Mock(return_value=_proposal()) + drift = replace(_drift(), source_revision="observation-v1") + first = invoke(tmp_path, drift=drift, optimize=optimize) + second = invoke(tmp_path, drift=replace(drift, escalated=True, cooldown_active=True), optimize=optimize) + assert first["research_key"] == second["research_key"] + assert optimize.call_count == 1 + + +def test_full_proposal_survives_checkpoint_without_losing_cost_or_validation_identity(tmp_path): + from quant_platform_kit.strategy_lifecycle.contracts import BacktestResult, BacktestValidationIdentity + + metrics = BacktestResult(strategy_profile="demo_strategy", domain="us_equity", param_set_id="candidate", + params={"a": 2}, start_date=date(2025, 1, 1), end_date=date(2025, 12, 31), + cost_inputs={"commission_bps": 1.0, "minimum_commission": 5.0}, + validation_identity=BacktestValidationIdentity(protocol="purged_walk_forward.v1", fold_id="fold1", + fold_role="walk_forward", train_start=date(2022, 1, 1), train_end=date(2022, 12, 31), + test_start=date(2023, 1, 2), test_end=date(2023, 6, 30), locked_oos_start=date(2024, 1, 1), + locked_oos_end=date(2025, 1, 1), purge_days=1, embargo_days=1)) + proposal = replace(_proposal(), current_metrics=metrics, proposed_metrics=metrics, + search_iterations=7, winning_dimensions=("sharpe",), walk_forward_passed=True) + seen = [] + + def gate(restored): + seen.append(restored) + assert restored == proposal + return _promotion_backtest_evidence() + + result = invoke(tmp_path, optimize=Mock(return_value=proposal), enforce_backtest_gates=gate) + assert result["status"] == "awaiting_human" + assert len(seen) == 1 + + +def test_deferral_waits_until_admitted_retry_time_then_rechecks_once(tmp_path): + retry = datetime(2026, 9, 8, 12, tzinfo=timezone.utc).timestamp() + diagnose = Mock(side_effect=[{"optimization_needed": False, "reason": "codex_research_deferred", "retry_at": retry}, + {"optimization_needed": False}]) + with patch.object(cycle, "datetime") as clock: + clock.now.return_value = datetime(2026, 9, 8, 11, tzinfo=timezone.utc) + first = invoke(tmp_path, diagnose=diagnose) + second = invoke(tmp_path, diagnose=diagnose) + assert first["status"] == second["status"] == "deferred" + assert diagnose.call_count == 1 + clock.now.return_value = datetime(2026, 9, 8, 12, tzinfo=timezone.utc) + third = invoke(tmp_path, diagnose=diagnose) + assert third["status"] == "parked" + assert diagnose.call_count == 2 + + +def test_damaged_checkpoint_does_not_rerun_or_block_a_different_identity(tmp_path): + import json + + first = invoke(tmp_path) + path = cycle.Path(first["ticket_path"]) + raw = json.loads(path.read_text()) + raw["research_progress"]["identity"]["revisions"]["input_revision"] = "mismatch" + path.write_text(json.dumps(raw)) + optimize = Mock() + bad = invoke(tmp_path, optimize=optimize) + assert bad["reason"] == "research_checkpoint_mismatch" + optimize.assert_not_called() + assert invoke(tmp_path, research_identity={**IDENTITY, "input_revision": "next"})["status"] == "awaiting_human" + + +def test_console_decision_recovery_does_not_race_an_active_directory_owner(tmp_path): + first = invoke(tmp_path) + path = cycle.Path(first["ticket_path"]) + local = cycle.load_research_promotion_ticket(path) + pull = Mock(return_value=remote_decision(local, "reject")) + before = path.read_bytes() + with cycle._research_directory_lock(tmp_path): + result = cycle.reconcile_saved_research_promotion_ticket(path, pull_console=pull) + assert result["status"] == "deferred" + assert result["reason"] == "research_in_progress" + assert path.read_bytes() == before + pull.assert_not_called() + assert cycle.reconcile_saved_research_promotion_ticket(path, pull_console=pull)["status"] == "updated" + + +def test_manual_decision_cli_uses_the_same_directory_lock(tmp_path, capsys): + from quant_platform_kit.strategy_lifecycle import cli + + first = invoke(tmp_path) + path = cycle.Path(first["ticket_path"]) + before = path.read_bytes() + with cycle._research_directory_lock(tmp_path): + assert cli.main(["research-promotion-decide", "--ticket", str(path), "--decision", "reject"]) == 1 + assert path.read_bytes() == before + assert "research_in_progress" in capsys.readouterr().err + assert cli.main(["research-promotion-decide", "--ticket", str(path), "--decision", "reject"]) == 0 + assert cycle.load_research_promotion_ticket(path).research_progress + + +def test_process_stop_during_stage_leaves_unknown_not_an_implicit_retry(tmp_path): + optimize = Mock(side_effect=SystemExit("synthetic process interruption")) + with pytest.raises(SystemExit): + invoke(tmp_path, optimize=optimize) + result = invoke(tmp_path, optimize=optimize) + assert result["reason"] == "research_outcome_unknown" + assert optimize.call_count == 1 + + +def test_checkpointed_strict_gate_is_revalidated_before_shadow(tmp_path): + import json + + save = cycle.save_research_promotion_ticket + path = None + + def stop_after_gate(ticket, target): + nonlocal path + path = target + result = save(ticket, target) + if ticket.research_progress.get("stages", {}).get("backtest", {}).get("status") == "completed": + raise SystemExit("synthetic process stop") + return result + + optimize, gates, shadow = Mock(return_value=_proposal()), Mock(return_value=_promotion_backtest_evidence()), Mock() + with patch.object(cycle, "save_research_promotion_ticket", side_effect=stop_after_gate): + with pytest.raises(SystemExit): + invoke(tmp_path, optimize=optimize, enforce_backtest_gates=gates, record_shadow=shadow) + raw = json.loads(path.read_text()) + raw["research_progress"]["stages"]["backtest"]["result"]["locked_independent_oos"]["reused_for_selection"] = True + path.write_text(json.dumps(raw)) + result = invoke(tmp_path, optimize=optimize, enforce_backtest_gates=gates, record_shadow=shadow) + assert result["status"] == "parked" + assert "locked_independent_oos_failed" in result["ticket"]["notes"] + assert optimize.call_count == gates.call_count == 1 + shadow.assert_not_called() + + +def test_operating_system_lock_blocks_another_process_and_releases_on_exit(tmp_path): + import os + import subprocess + import sys + + script = "from pathlib import Path; from quant_platform_kit.strategy_lifecycle.research_promotion_cycle import _research_directory_lock\nimport sys\nwith _research_directory_lock(Path(sys.argv[1])) as owned:\n print(owned, flush=True)\n sys.stdin.readline()\n" + env = {"PATH": os.environ.get("PATH", ""), "PYTHONPATH": str(cycle.Path(cycle.__file__).parents[2]), + "PYTHONDONTWRITEBYTECODE": "1"} + proc = subprocess.Popen([sys.executable, "-c", script, str(tmp_path)], stdin=subprocess.PIPE, + stdout=subprocess.PIPE, stderr=subprocess.PIPE, text=True, env=env) + try: + assert proc.stdout.readline().strip() == "True" + diagnose = Mock() + result = invoke(tmp_path, diagnose=diagnose) + assert result["reason"] == "research_in_progress" + diagnose.assert_not_called() + finally: + proc.communicate("\n", timeout=10) + assert invoke(tmp_path)["status"] == "awaiting_human" + + +def test_restart_before_first_console_attempt_delivers_the_saved_ticket_once(tmp_path): + save = cycle.save_research_promotion_ticket + sync = Mock(return_value=True) + + def stop_before_delivery(ticket, path): + result = save(ticket, path) + if ticket.state is cycle.ResearchPromotionState.AWAITING_HUMAN: + raise SystemExit("synthetic process stop before first delivery") + return result + + with patch.object(cycle, "save_research_promotion_ticket", side_effect=stop_before_delivery): + with pytest.raises(SystemExit): + invoke(tmp_path, sync_console=sync) + sync.assert_not_called() + optimize = Mock() + result = invoke(tmp_path, sync_console=sync, optimize=optimize) + assert result["status"] == "awaiting_human" + assert result["console_synced"] is True + sync.assert_called_once() + optimize.assert_not_called() + + +def test_optional_diagnosis_does_not_create_a_second_experiment_for_same_evidence(tmp_path): + first = invoke(tmp_path) + diagnose, optimize = Mock(), Mock() + second = invoke(tmp_path, diagnose=diagnose, optimize=optimize) + assert first["research_key"] == second["research_key"] + diagnose.assert_not_called() + optimize.assert_not_called() + + +@pytest.mark.parametrize("field,value", [("live_authority_granted", True), ("proposed_params", {"a": 999})]) +def test_console_callback_cannot_mutate_saved_authority_or_candidate(tmp_path, field, value): + def bad_sync(ticket): + setattr(ticket, field, value) + return True + + result = invoke(tmp_path, sync_console=bad_sync) + assert result["ticket"]["live_authority_granted"] is False + assert result["ticket"]["proposed_params"] == _proposal().proposed_params + assert result["console_synced"] is False + saved = cycle.load_research_promotion_ticket(result["ticket_path"]) + assert saved.live_authority_granted is False + + +@pytest.mark.parametrize("retry", [0, -1, 1000]) +def test_already_expired_deferral_does_not_create_a_retry_loop(tmp_path, retry): + diagnose = Mock(return_value={"optimization_needed": False, "reason": "codex_research_deferred", + "retry_at": retry}) + first = invoke(tmp_path, diagnose=diagnose) + second = invoke(tmp_path, diagnose=diagnose) + assert first["status"] == second["status"] == "deferred" + assert first["retry_at"] is None + assert diagnose.call_count == 1 + + +@pytest.mark.parametrize("delivery,identity_changes,expected_syncs", [ + (None, False, 1), (None, True, 0), ("running", False, 0), ("unconfirmed", False, 0), +]) +def test_normal_auto_pilot_resumes_only_matching_unstarted_console_delivery( + tmp_path, delivery, identity_changes, expected_syncs, +): + from types import SimpleNamespace + from quant_platform_kit.strategy_lifecycle.codex_integration import run_auto_pilot_cycle + from quant_platform_kit.strategy_lifecycle.contracts import StrategyPerformanceSnapshot + from quant_platform_kit.strategy_lifecycle.promotion_actionable_runner import run_actionable_research_promotion + + directory = tmp_path / "research_promotion_tickets" + save = cycle.save_research_promotion_ticket + before_sync = Mock() + + def stop_before_delivery(ticket, path): + if ticket.state is cycle.ResearchPromotionState.AWAITING_HUMAN: + if delivery is not None: + ticket.research_progress["console_delivery"] = delivery + save(ticket, path) + raise SystemExit("synthetic crash before delivery acknowledgement") + return save(ticket, path) + + with patch.object(cycle, "save_research_promotion_ticket", side_effect=stop_before_delivery): + with pytest.raises(SystemExit): + run_actionable_research_promotion( + strategy_profile="demo_strategy", domain="us_equity", as_of="2026-09-07", + drift_score=0.8, source_revision="observation-v1", evaluation_date="2026-09-08", + research_identity=IDENTITY, ticket_dir=directory, optimize=Mock(return_value=_proposal()), + enforce_backtest_gates=Mock(return_value=_promotion_backtest_evidence()), + record_shadow=Mock(return_value={"evidence_kind": "paired_shadow", "passed": True}), + sync_console=before_sync, + ) + before_sync.assert_not_called() + original_files = list(directory.glob("*.json")) + store = SimpleNamespace(local_root=tmp_path, load_latest_snapshot=lambda *_: StrategyPerformanceSnapshot( + strategy_profile="demo_strategy", domain="us_equity", platform="test", as_of=date(2026, 9, 7), + source_revision="observation-v1")) + sync, optimize, gates, shadow = Mock(return_value=True), Mock(), Mock(), Mock() + prefix = "quant_platform_kit.strategy_lifecycle.codex_integration." + identity = {**IDENTITY, "input_revision": "new-input"} if identity_changes else IDENTITY + with patch(prefix + "_run_monitor_phase", return_value=[]), \ + patch(prefix + "_run_drift_phase", return_value=([_drift()], [_drift()])), \ + patch(prefix + "call_ai_optimization_decision") as ai: + run_auto_pilot_cycle("us_equity", store=store, create_issues=False, research_identity=identity, + optimize=optimize, enforce_backtest_gates=gates, record_shadow=shadow, + sync_console=sync, pull_console=Mock(return_value=None)) + assert sync.call_count == expected_syncs + ai.assert_not_called() + optimize.assert_not_called() + gates.assert_not_called() + shadow.assert_not_called() + assert list(directory.glob("*.json")) == original_files + + +@pytest.mark.parametrize("decision", [False, None, 1, "true", RuntimeError("private admission detail")]) +def test_new_experiment_admission_rejects_before_ticket_or_model(tmp_path, decision): + admission = Mock(side_effect=decision) if isinstance(decision, Exception) else Mock(return_value=decision) + diagnose, optimize = Mock(), Mock() + result = invoke(tmp_path, admit_new_research=admission, diagnose=diagnose, optimize=optimize) + assert result["status"] in {"parked", "deferred"} + assert result["reason"] in {"new_research_not_admitted", "research_admission_unavailable"} + assert not list(tmp_path.glob("*.json")) + assert "private admission detail" not in str(result) + admission.assert_called_once() + diagnose.assert_not_called() + optimize.assert_not_called() + + +def test_admission_runs_under_directory_lock_and_binds_the_persisted_utc_clock(tmp_path): + admitted_times = [] + + def admit(directory, created_at): + assert directory == tmp_path + assert not list(directory.glob("*.json")) + with cycle._research_directory_lock(directory) as acquired: + assert acquired is False + assert datetime.fromisoformat(created_at).utcoffset().total_seconds() == 0 + admitted_times.append(created_at) + return True + + result = invoke(tmp_path, admit_new_research=admit) + assert result["ticket"]["created_at"] == admitted_times[0] + + +@pytest.mark.parametrize("outcome", ["awaiting", "terminal", "unknown"]) +def test_saved_experiments_never_require_a_second_new_experiment_admission(tmp_path, outcome): + admission = Mock(return_value=True) + kwargs = {"admit_new_research": admission} + if outcome == "terminal": + kwargs["diagnose"] = Mock(return_value={"optimization_needed": False}) + if outcome == "unknown": + kwargs["optimize"] = Mock(side_effect=TimeoutError("synthetic unknown")) + first = invoke(tmp_path, **kwargs) + admission.return_value = False + second = invoke(tmp_path, **kwargs) + assert first["research_key"] == second["research_key"] + assert admission.call_count == 1 + + +def test_actual_auto_pilot_passes_new_admission_before_ai(tmp_path): + from types import SimpleNamespace + from quant_platform_kit.strategy_lifecycle.codex_integration import run_auto_pilot_cycle + from quant_platform_kit.strategy_lifecycle.contracts import StrategyPerformanceSnapshot + + store = SimpleNamespace(local_root=tmp_path, load_latest_snapshot=lambda *_: StrategyPerformanceSnapshot( + strategy_profile="demo_strategy", domain="us_equity", platform="test", as_of=date(2026, 9, 7), + source_revision="observation-v1")) + admission = Mock(return_value=False) + prefix = "quant_platform_kit.strategy_lifecycle.codex_integration." + with patch(prefix + "_run_monitor_phase", return_value=[]), \ + patch(prefix + "_run_drift_phase", return_value=([_drift()], [_drift()])), \ + patch(prefix + "call_ai_optimization_decision") as ai: + result = run_auto_pilot_cycle("us_equity", store=store, create_issues=False, research_identity=IDENTITY, + optimize=Mock(), enforce_backtest_gates=Mock(), record_shadow=Mock(), admit_new_research=admission) + assert result["actions"][0]["reason"] == "new_research_not_admitted" + assert admission.call_count == 1 + ai.assert_not_called() + assert not list((tmp_path / "research_promotion_tickets").glob("*.json")) diff --git a/tests/test_research_promotion_shadow_resume.py b/tests/test_research_promotion_shadow_resume.py new file mode 100644 index 0000000..61cacc9 --- /dev/null +++ b/tests/test_research_promotion_shadow_resume.py @@ -0,0 +1,281 @@ +"""Synthetic local observations; no model, broker, real shadow or console.""" + +from datetime import date, datetime, timezone +from unittest.mock import Mock + +import pytest + +from quant_platform_kit.strategy_lifecycle import research_promotion_cycle as cycle +from quant_platform_kit.strategy_lifecycle.promotion_actionable_runner import run_actionable_research_promotion +from tests.test_research_promotion_resume import IDENTITY +from tests.test_research_promotion_cycle import _proposal, _promotion_backtest_evidence +from tests.test_paired_shadow_evidence import _dependencies, _leg, _policy + + +class Clock(datetime): + current = datetime(2026, 9, 8, tzinfo=timezone.utc) + + @classmethod + def now(cls, tz=None): + return cls.current.astimezone(tz) + + +@pytest.fixture(autouse=True) +def clock(monkeypatch): + Clock.current = datetime(2026, 9, 8, tzinfo=timezone.utc) + monkeypatch.setattr(cycle, "datetime", Clock) + + +def pending(**changes): + return {"status": "pending", "passed": False, "no_order": True, + "live_authority_granted": False, + "retry_at": datetime(2026, 9, 20, tzinfo=timezone.utc).timestamp(), **changes} + + +def complete(): + from quant_platform_kit.strategy_lifecycle.forward_observation_receipt import build_forward_observation_receipt + from quant_platform_kit.strategy_lifecycle.paired_shadow_evidence import build_paired_shadow_evidence + + policy = _policy(strategy_profile="demo_strategy", required_trading_sessions=2, + review_milestones=(1,), observation_start_session="2026-09-07") + first = build_forward_observation_receipt(policy=policy, observation_session="2026-09-07", + observation_index=1, dependency_digests=_dependencies(), evidence_modes=policy.non_live_evidence_modes) + earlier = build_paired_shadow_evidence(policy=policy, forward_observation_receipt=first, + baseline_id="baseline", observed_at="2026-09-07T20:00:00Z", input_snapshot_sha256="a" * 64, + candidate=_leg("candidate"), baseline=_leg("baseline")) + second = build_forward_observation_receipt(policy=policy, observation_session="2026-09-08", + observation_index=2, dependency_digests=_dependencies(), evidence_modes=policy.non_live_evidence_modes, + previous_receipt=first) + return {"status": "complete", "observation": dict(policy=policy, forward_observation_receipt=second, + baseline_id="baseline", observed_at="2026-09-08T20:00:00Z", input_snapshot_sha256="a" * 64, + candidate=_leg("candidate"), baseline=_leg("baseline"), previous_evidence=earlier, + previous_forward_observation_receipt=first)} + + +def job(tmp_path, **changes): + kwargs = dict(strategy_profile="demo_strategy", domain="us_equity", as_of="2026-09-07", + drift_score=.8, source_revision="observation-v1", evaluation_date="2026-09-08", + research_identity=IDENTITY, ticket_dir=tmp_path, + diagnose=Mock(return_value={"optimization_needed": True}), optimize=Mock(return_value=_proposal()), + enforce_backtest_gates=Mock(return_value=_promotion_backtest_evidence()), + record_shadow=Mock(return_value=pending()), read_pending_shadow=Mock(return_value=complete()), + sync_console=Mock(return_value=True), pull_console=Mock(return_value=None)) + return {**kwargs, **changes} + + +def test_actual_runner_resumes_only_local_shadow_after_original_drift_expires(tmp_path): + args = job(tmp_path) + first = run_actionable_research_promotion(**args) + assert first["status"] == "deferred" + assert first["ticket"]["state"] == "shadow_recorded" + assert first["ticket"]["shadow_passed"] is False + args["sync_console"].assert_not_called() + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + second = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert second["status"] == "awaiting_human" + assert second["research_key"] == first["research_key"] + assert second["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")] == [1] * 6 + + +@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))) + first = run_actionable_research_promotion(**args) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + second = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert first["status"] == second["status"] == "deferred" + args["read_pending_shadow"].assert_not_called() + assert args["record_shadow"].call_count == 1 + + +def test_not_due_shadow_does_not_poll_or_repeat_admission(tmp_path): + args = job(tmp_path, admit_new_research=Mock(return_value=True)) + run_actionable_research_promotion(**args) + run_actionable_research_promotion(**args) + args["read_pending_shadow"].assert_not_called() + assert args["admit_new_research"].call_count == 1 + + +@pytest.mark.parametrize("case", ["new", "identity", "incomplete", "gate_failed"]) +def test_stale_input_cannot_start_or_resume_unfinished_research(tmp_path, case): + args = job(tmp_path) + if case != "new": + first = run_actionable_research_promotion(**args) + if case == "identity": + args["research_identity"] = {**IDENTITY, "input_revision": "changed"} + else: + ticket = cycle.load_research_promotion_ticket(first["ticket_path"]) + stages = ticket.research_progress["stages"] + if case == "incomplete": + del stages["optimize"] + else: + stages["backtest"]["result"]["status"] = "FAIL" + cycle.save_research_promotion_ticket(ticket, first["ticket_path"]) + before = [args[key].call_count for key in ("diagnose", "optimize", "enforce_backtest_gates", "record_shadow")] + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert result["status"] == "parked" + assert [args[key].call_count for key in ("diagnose", "optimize", "enforce_backtest_gates", "record_shadow")] == before + args["read_pending_shadow"].assert_not_called() + + +@pytest.mark.parametrize("case", ["exception", "invalid_receipt", "missing_previous", "bare_passed", "live", "incomplete_window"]) +def test_unknown_or_invalid_shadow_completion_cannot_reach_console_or_retry(tmp_path, case): + reader = Mock(return_value=complete()) + if case == "exception": + reader.side_effect = TimeoutError("private material must not escape") + elif case == "invalid_receipt": + reader.return_value["observation"]["forward_observation_receipt"]["receipt_sha256"] = "0" * 64 + elif case == "missing_previous": + del reader.return_value["observation"]["previous_evidence"] + del reader.return_value["observation"]["previous_forward_observation_receipt"] + elif case == "bare_passed": + reader.return_value = {"status": "complete", "evidence_kind": "paired_shadow", "passed": True} + elif case == "live": + reader.return_value = pending(live_authority_granted=True) + else: + from dataclasses import replace + value = reader.return_value["observation"] + value["policy"] = replace(value["policy"], required_trading_sessions=63) + args = job(tmp_path, read_pending_shadow=reader) + run_actionable_research_promotion(**args) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + for _ in range(2): + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert result["status"] == "parked" + assert "private material" not in str(result) + assert reader.call_count == 1 + args["sync_console"].assert_not_called() + + +def test_shadow_read_uses_existing_directory_lock(tmp_path): + args = job(tmp_path) + run_actionable_research_promotion(**args) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + with cycle._research_directory_lock(tmp_path) as acquired: + assert acquired + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert result["reason"] == "research_in_progress" + args["read_pending_shadow"].assert_not_called() + + +def test_pending_reader_can_wait_again_without_restarting_the_experiment(tmp_path): + next_retry = datetime(2026, 9, 28, tzinfo=timezone.utc).timestamp() + reader = Mock(side_effect=[pending(retry_at=next_retry), complete()]) + args = job(tmp_path, read_pending_shadow=reader) + run_actionable_research_promotion(**args) + for day, calls in [(21, 1), (22, 1), (29, 2)]: + Clock.current = datetime(2026, 9, day, tzinfo=timezone.utc) + result = run_actionable_research_promotion(**{**args, "evaluation_date": f"2026-09-{day}"}) + assert result["status"] == ("awaiting_human" if day == 29 else "deferred") + assert reader.call_count == calls + assert args["record_shadow"].call_count == args["optimize"].call_count == 1 + + +def test_missing_reader_never_repeats_first_callback_and_explicit_failure_stays_parked(tmp_path): + args = job(tmp_path, read_pending_shadow=None) + run_actionable_research_promotion(**args) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert result["reason"] == "shadow_reader_unavailable" + assert args["record_shadow"].call_count == 1 + reader = Mock(return_value={"status": "failed", "passed": False, "evidence_kind": "paired_shadow", + "no_order": True, "live_authority_granted": False}) + for _ in range(2): + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21", "read_pending_shadow": reader}) + assert result["status"] == "parked" + assert reader.call_count == 1 + args["sync_console"].assert_not_called() + + +@pytest.mark.parametrize("stage_status", ["pending", "completed"]) +def test_interrupt_after_shadow_result_save_resumes_without_repeating_callback(tmp_path, monkeypatch, stage_status): + args = job(tmp_path) + save = cycle.save_research_promotion_ticket + if stage_status == "completed": + run_actionable_research_promotion(**args) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + args["evaluation_date"] = "2026-09-21" + + def stop(ticket, path): + save(ticket, path) + if ticket.research_progress.get("stages", {}).get("shadow", {}).get("status") == stage_status: + raise SystemExit("synthetic checkpoint boundary") + + monkeypatch.setattr(cycle, "save_research_promotion_ticket", stop) + with pytest.raises(SystemExit): + run_actionable_research_promotion(**args) + monkeypatch.setattr(cycle, "save_research_promotion_ticket", save) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + assert result["status"] == "awaiting_human" + assert args["read_pending_shadow"].call_count == args["record_shadow"].call_count == 1 + assert args["optimize"].call_count == args["enforce_backtest_gates"].call_count == 1 + + +def test_normal_auto_pilot_passes_reader_through_expired_observation(tmp_path, monkeypatch): + from quant_platform_kit.strategy_lifecycle import codex_integration as codex, production_drift_health_probe as probe + from quant_platform_kit.strategy_lifecycle.contracts import StrategyPerformanceSnapshot + from tests.test_research_promotion_cycle import _drift + + monkeypatch.setattr(probe, "datetime", Clock) + monkeypatch.setattr(codex, "_run_monitor_phase", Mock(return_value=[])) + monkeypatch.setattr(codex, "_run_drift_phase", Mock(return_value=([_drift()], [_drift()]))) + ai = Mock(return_value={"optimization_needed": True}) + monkeypatch.setattr(codex, "call_ai_optimization_decision", ai) + from quant_platform_kit.strategy_lifecycle import ai_reviewer + monkeypatch.setattr(ai_reviewer, "review_proposal", Mock(return_value=Mock(verdict="approve", to_dict=lambda: {}))) + store = Mock(local_root=tmp_path) + store.load_latest_snapshot.return_value = StrategyPerformanceSnapshot(strategy_profile="demo_strategy", domain="us_equity", + platform="test", as_of=date(2026, 9, 7), source_revision="observation-v1") + args = job(tmp_path) + kwargs = {key: args[key] for key in ("research_identity", "optimize", "enforce_backtest_gates", "record_shadow", + "read_pending_shadow", "sync_console", "pull_console")} + first = codex.run_auto_pilot_cycle("us_equity", store=store, create_issues=False, **kwargs) + assert first["actions"][0]["research_promotion_state"] == "deferred" + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + second = codex.run_auto_pilot_cycle("us_equity", store=store, create_issues=False, **kwargs) + assert second["actions"][0]["research_promotion_state"] == "awaiting_human" + assert ai.call_count == args["optimize"].call_count == args["enforce_backtest_gates"].call_count == 1 + assert args["read_pending_shadow"].call_count == 1 + + +@pytest.mark.parametrize("decision", ["accept", "reject"]) +def test_actual_caller_recovers_unknown_console_delivery_and_human_decision_after_sixty_days(tmp_path, decision): + from tests.test_research_promotion_reconciliation import remote_decision + + args = job(tmp_path, sync_console=Mock(side_effect=TimeoutError("private uncertain write"))) + run_actionable_research_promotion(**args) + 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" and ready["console_synced"] is False + 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): + recovered = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-11-21"}) + assert recovered["status"] == ("human_accepted" if decision == "accept" else "human_rejected") + assert recovered["ticket"]["live_authority_granted"] is False + assert args["pull_console"].call_count == args["sync_console"].call_count == 1 + assert [args[key].call_count for key in ("diagnose", "optimize", "enforce_backtest_gates", + "record_shadow", "read_pending_shadow")] == [1] * 5 + + +@pytest.mark.parametrize("case", ["missing_stage", "changed_candidate"]) +def test_stale_awaiting_without_matching_completed_checkpoint_cannot_read_console(tmp_path, case): + args = job(tmp_path) + run_actionable_research_promotion(**args) + Clock.current = datetime(2026, 9, 21, tzinfo=timezone.utc) + ready = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-09-21"}) + local = cycle.load_research_promotion_ticket(ready["ticket_path"]) + if case == "missing_stage": + del local.research_progress["stages"]["backtest"] + else: + local.proposed_params = {"a": 999} + cycle.save_research_promotion_ticket(local, ready["ticket_path"]) + result = run_actionable_research_promotion(**{**args, "evaluation_date": "2026-11-21"}) + assert result["status"] == "parked" + args["pull_console"].assert_not_called() + assert args["sync_console"].call_count == 1 diff --git a/tests/test_reusable_drift_workflow.py b/tests/test_reusable_drift_workflow.py index 3938039..4396ac5 100644 --- a/tests/test_reusable_drift_workflow.py +++ b/tests/test_reusable_drift_workflow.py @@ -61,7 +61,7 @@ def test_reusable_drift_workflow_enforces_lifecycle_preflight() -> None: assert "create_issues_for_domain" in workflow assert 'CODEX_AUDIT_SERVICE_URL: ${{ secrets.codex_audit_service_url }}' in workflow assert 'AI_GATEWAY_SERVICE_URL: ${{ inputs.ai_gateway_service_url }}' in workflow - assert 'ref: 2351c987d6df12d48355aaa5803db8b089476310' in workflow + assert 'ref: 60bd64a2ae059a082614181eeb845b46df395523' in workflow assert workflow.count('GH_TOKEN: ${{ github.token }}') >= 2 assert "emit_parked_record" in workflow assert '"schema": "qsl.drift_dual_review_availability.v1"' in workflow