diff --git a/README.zh-CN.md b/README.zh-CN.md index 0981066d..74a74d25 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -59,6 +59,8 @@ quant-lifecycle dashboard --format all - [`docs/strategy_lifecycle_policy.zh-CN.md`](docs/strategy_lifecycle_policy.zh-CN.md) - [`docs/strategy_portfolio_action_matrix.zh-CN.md`](docs/strategy_portfolio_action_matrix.zh-CN.md) - [`docs/evidence_package_template.zh-CN.md`](docs/evidence_package_template.zh-CN.md) +- [`docs/cross_asset_research_driver.zh-CN.md`](docs/cross_asset_research_driver.zh-CN.md) +- [`docs/cross_asset_forward_risk_driver.zh-CN.md`](docs/cross_asset_forward_risk_driver.zh-CN.md) ## 云服务抽象层 diff --git a/docs/cross_asset_forward_risk_driver.zh-CN.md b/docs/cross_asset_forward_risk_driver.zh-CN.md new file mode 100644 index 00000000..75ab36ed --- /dev/null +++ b/docs/cross_asset_forward_risk_driver.zh-CN.md @@ -0,0 +1,51 @@ +# 跨资产 P4/P5 观察适配契约 + +## 当前架构理解 + +P1–P3 research driver 负责绑定输入、冻结配置和历史/OOS 证据;P4/P5 的时间属性和账户/组合属性不同,不应塞回同一个 artifact。`forward_risk_terminal.v1` 因此作为第二个纯观察 envelope: + +```text +P1–P3 terminal digest +→ P4 shadow/forward(或已有 paper 结果)身份 +→ P5 portfolio RiskSnapshot 身份 +→ READY / DEFERRED / PARKED terminal artifact +``` + +## 权限边界 + +该适配器固定: + +- `no_order=true`; +- `permission_effect=none`; +- `broker_dependency=false`; +- 不启动 paper、shadow 或 live runtime; +- 不抓行情、不运行策略、不计算仓位、不连接 broker; +- `READY` 只描述证据齐全,不授予 shadow、paper 或 live 权限。 + +P4 的 `mode=paper` 仅表示消费了其他平台已经产生并验证的模拟结果;QPK 适配器本身仍不连接 paper broker。没有 paper 能力的平台直接使用 `mode=shadow`,不会因此被阻塞。 + +## P4/P5 READY 条件 + +| 阶段 | artifact schema | 额外约束 | +|---|---|---| +| P4 | `forward_observation.v1` | P1–P3 terminal 必须 READY;candidate 必须一致;artifact 在 terminal 生成时未过期 | +| P5 | `portfolio_risk_snapshot.v1` | P4 必须 READY;candidate 必须一致;RiskSnapshot 身份未过期 | + +缺少正常上游时输出 `DEFERRED`;格式错误、过期、candidate 不一致或越过阶段依赖时输出 `PARKED`。每个合法 P1–P3 terminal 都可生成一个 P4/P5 terminal artifact,不能以 workflow 绿色代替终态文件。 + +## 现有样板评估 + +- `ShadowValidator` 已能读取近期 performance snapshot 并比较候选,但尚未定义跨资产不可变 P4 artifact;各策略 producer 需要后续输出 `forward_observation.v1`。 +- `RiskSnapshot` 已有 fail-closed、expiry、Kelly fraction、风险预算和熔断状态校验,适合作为 P5 producer;本契约只绑定其不可变身份,不复制风险算法。 +- `lifecycle_matrix_runtime` 已能只读聚合 P4/P5 terminal 状态,可在 producer 接线后继续复用。 +- paper 不是通用前置条件。支持 paper 的平台可提供 P4 paper observation;不支持的平台以 no-order shadow/forward observation 完成 P4。 + +## 低风险迁移 + +1. 所有资产先输出 P1–P3 terminal;缺证据时如实 `DEFERRED/PARKED`。 +2. 各策略将现有 shadow 日报适配为 `forward_observation.v1`,不改变策略计算。 +3. 将现有 `RiskSnapshot` 序列化、摘要后注册为 `portfolio_risk_snapshot.v1`。 +4. workflow 的 `always()` 终态步骤输出 P4/P5 envelope。 +5. 只读 matrix 消费这些 artifact;live/runtime 继续使用独立 authority 和 Risk Gate。 + +不推荐新增第二套 scheduler、通用事件总线或让此模块直接运行 paper/broker;这些会把证据适配层与执行层重新耦合。 diff --git a/docs/cross_asset_integration_plan.zh-CN.md b/docs/cross_asset_integration_plan.zh-CN.md index 49c0637f..09b2d6bd 100644 --- a/docs/cross_asset_integration_plan.zh-CN.md +++ b/docs/cross_asset_integration_plan.zh-CN.md @@ -12,3 +12,12 @@ 当前覆盖:CN equity 8 条、HK equity 3 条、crypto 5 条;详见 `docs/registry/cross_asset_strategy_inventory.json`。 + +P1–P3 生产者统一写出 `research_driver_terminal.v1`,包括正常、等待和失败 +分支;具体字段、终态推导和 no-order 边界见 +[`cross_asset_research_driver.zh-CN.md`](cross_asset_research_driver.zh-CN.md)。 + +P4/P5 使用独立的 `forward_risk_terminal.v1` 观察契约绑定 P1–P3 terminal +digest、forward/shadow(或已有 paper)证据与组合 RiskSnapshot;该契约固定 +no-order、无 broker 依赖且不产生权限。详见 +[`cross_asset_forward_risk_driver.zh-CN.md`](cross_asset_forward_risk_driver.zh-CN.md)。 diff --git a/docs/cross_asset_research_driver.zh-CN.md b/docs/cross_asset_research_driver.zh-CN.md new file mode 100644 index 00000000..906b8baa --- /dev/null +++ b/docs/cross_asset_research_driver.zh-CN.md @@ -0,0 +1,66 @@ +# 跨资产 Research Driver 契约 + +## 目标与边界 + +`research_driver_terminal.v1` 为 US、CN、HK、Crypto 提供相同的 P1–P3 研究终态。它只汇总已经生成并验证过的 artifact 身份,不抓取数据、不执行回测、不读取策略 catalog,也不接入 broker。 + +固定边界: + +```text +P1 research input manifest +→ P2 strategy config freeze +→ P3 strategy evidence package +→ READY / DEFERRED / PARKED terminal artifact +``` + +每次有合法的 run/strategy/candidate 身份时都必须生成 terminal artifact: + +- 三个阶段均有有效 artifact:`READY`; +- artifact 尚未生成或正常等待上游:`DEFERRED`; +- artifact 格式错误、摘要错误或越过上游依赖:`PARKED`; +- 任意结果始终为 `no_order=true`、`permission_effect=none`。 + +## 证据要求 + +| 阶段 | READY 所需 artifact | 说明 | +|---|---|---| +| P1 | `research_input_manifest.v1` | 由现有 P1 validator 校验后提供身份与 SHA-256 | +| P2 | `strategy_config_freeze.v1` | 冻结候选配置的身份与 SHA-256 | +| P3 | `strategy_evidence_package.v2` | 由现有 P3 validator 校验后提供身份与 SHA-256 | + +`catalog_status`、`runtime_enabled`、inventory 条目或网页显示都不是证据。契约不接受这些字段,并固定输出 `catalog_status_used_as_evidence=false`。 + +## 生产者用法 + +```python +from quant_platform_kit.strategy_lifecycle import ( + build_ready_research_stage, + build_research_driver_terminal_artifact, +) + +terminal = build_research_driver_terminal_artifact( + run_id="daily-2026-08-24", + generated_at="2026-08-24T22:00:00+08:00", + strategy_id="example_strategy", + candidate_id="candidate-001", + domain="cn_equity", + p1_input=build_ready_research_stage( + "p1_input", artifact_id="manifest-001", artifact_sha256="a" * 64 + ), + p2_freeze=build_ready_research_stage( + "p2_freeze", artifact_id="freeze-001", artifact_sha256="b" * 64 + ), + p3_evidence=None, +) +assert terminal["terminal_status"] == "DEFERRED" +assert terminal["no_order"] is True +``` + +生产 workflow 应在正常、跳过、依赖缺失和失败分支的 finally/always 步骤中写出该 JSON。这个契约不会自动把 `READY` 提升为 shadow、paper 或 live;后续生命周期消费者仍需独立验证 artifact 与权限。 + +## 迁移风险 + +- 旧策略可先只写 `DEFERRED`,不需要伪造 P1/P2/P3 完成状态。 +- domain 仅允许 `us_equity`、`cn_equity`、`hk_equity`、`crypto`。 +- 该模块不替代 P1/P3 原始 validator;它只绑定验证后的不可变身份。 +- 不允许用 catalog/inventory 状态填充 READY,也不允许从该 artifact 推导交易权限。 diff --git a/src/quant_platform_kit/schemas/forward-risk-terminal.v1.schema.json b/src/quant_platform_kit/schemas/forward-risk-terminal.v1.schema.json new file mode 100644 index 00000000..cca5f79a --- /dev/null +++ b/src/quant_platform_kit/schemas/forward-risk-terminal.v1.schema.json @@ -0,0 +1,142 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://quantstrategylab.dev/schemas/forward-risk-terminal.v1.schema.json", + "title": "forward_risk_terminal.v1", + "description": "Terminal, cross-asset P4/P5 observation artifact. It has no broker or execution authority.", + "type": "object", + "additionalProperties": false, + "required": [ + "schema_version", + "terminal", + "terminal_status", + "generated_at", + "run_id", + "strategy_id", + "candidate_id", + "domain", + "research_terminal_status", + "research_terminal_sha256", + "no_order", + "permission_effect", + "broker_dependency", + "stages" + ], + "properties": { + "schema_version": {"const": "forward_risk_terminal.v1"}, + "terminal": {"const": true}, + "terminal_status": {"$ref": "#/$defs/terminalStatus"}, + "generated_at": {"type": "string", "format": "date-time"}, + "run_id": {"$ref": "#/$defs/nonEmptyString"}, + "strategy_id": {"$ref": "#/$defs/nonEmptyString"}, + "candidate_id": {"$ref": "#/$defs/nonEmptyString"}, + "domain": {"enum": ["us_equity", "cn_equity", "hk_equity", "crypto"]}, + "research_terminal_status": {"$ref": "#/$defs/terminalStatus"}, + "research_terminal_sha256": {"$ref": "#/$defs/sha256"}, + "no_order": {"const": true}, + "permission_effect": {"const": "none"}, + "broker_dependency": {"const": false}, + "stages": { + "type": "object", + "additionalProperties": false, + "required": ["p4_forward", "p5_risk"], + "properties": { + "p4_forward": { + "$ref": "#/$defs/stageRecord", + "properties": { + "stage": {"const": "P4"}, + "mode": {"enum": ["shadow", "paper"]}, + "artifact": { + "oneOf": [ + {"type": "null"}, + { + "$ref": "#/$defs/artifactIdentity", + "properties": {"schema_version": {"const": "forward_observation.v1"}} + } + ] + } + } + }, + "p5_risk": { + "$ref": "#/$defs/stageRecord", + "properties": { + "stage": {"const": "P5"}, + "mode": {"const": "portfolio_risk"}, + "artifact": { + "oneOf": [ + {"type": "null"}, + { + "$ref": "#/$defs/artifactIdentity", + "properties": {"schema_version": {"const": "portfolio_risk_snapshot.v1"}} + } + ] + } + } + } + } + } + }, + "$defs": { + "nonEmptyString": {"type": "string", "minLength": 1, "pattern": ".*\\S.*"}, + "sha256": {"type": "string", "pattern": "^[0-9a-f]{64}$"}, + "reasonCode": {"type": "string", "pattern": "^[a-z][a-z0-9_]{0,127}$"}, + "terminalStatus": {"enum": ["READY", "DEFERRED", "PARKED"]}, + "artifactIdentity": { + "type": "object", + "additionalProperties": false, + "required": [ + "artifact_id", + "schema_version", + "sha256", + "candidate_id", + "observed_at", + "expires_at" + ], + "properties": { + "artifact_id": {"$ref": "#/$defs/nonEmptyString"}, + "schema_version": {"type": "string"}, + "sha256": {"$ref": "#/$defs/sha256"}, + "candidate_id": {"$ref": "#/$defs/nonEmptyString"}, + "observed_at": {"type": "string", "format": "date-time"}, + "expires_at": {"type": "string", "format": "date-time"} + } + }, + "stageRecord": { + "type": "object", + "additionalProperties": false, + "required": ["stage", "status", "mode", "artifact", "reason_codes"], + "properties": { + "stage": {"enum": ["P4", "P5"]}, + "status": {"$ref": "#/$defs/terminalStatus"}, + "mode": {"enum": ["shadow", "paper", "portfolio_risk"]}, + "artifact": { + "oneOf": [ + {"type": "null"}, + {"$ref": "#/$defs/artifactIdentity"} + ] + }, + "reason_codes": { + "type": "array", + "uniqueItems": true, + "items": {"$ref": "#/$defs/reasonCode"} + } + }, + "allOf": [ + { + "if": {"properties": {"status": {"const": "READY"}}, "required": ["status"]}, + "then": { + "properties": { + "artifact": {"type": "object"}, + "reason_codes": {"maxItems": 0} + } + }, + "else": { + "properties": { + "artifact": {"type": "null"}, + "reason_codes": {"minItems": 1} + } + } + } + ] + } + } +} diff --git a/src/quant_platform_kit/schemas/research-driver-terminal.v1.schema.json b/src/quant_platform_kit/schemas/research-driver-terminal.v1.schema.json new file mode 100644 index 00000000..10ee6d50 --- /dev/null +++ b/src/quant_platform_kit/schemas/research-driver-terminal.v1.schema.json @@ -0,0 +1,134 @@ +{ + "$schema": "https://json-schema.org/draft/2020-12/schema", + "$id": "https://quantstrategylab.dev/schemas/research-driver-terminal.v1.schema.json", + "title": "research_driver_terminal.v1", + "description": "Terminal, cross-asset P1-P3 research artifact. It never grants execution authority.", + "type": "object", + "additionalProperties": false, + "required": [ + "schema_version", + "terminal", + "terminal_status", + "generated_at", + "run_id", + "strategy_id", + "candidate_id", + "domain", + "no_order", + "permission_effect", + "catalog_status_used_as_evidence", + "stages" + ], + "properties": { + "schema_version": {"const": "research_driver_terminal.v1"}, + "terminal": {"const": true}, + "terminal_status": {"enum": ["READY", "DEFERRED", "PARKED"]}, + "generated_at": {"type": "string", "format": "date-time"}, + "run_id": {"$ref": "#/$defs/nonEmptyString"}, + "strategy_id": {"$ref": "#/$defs/nonEmptyString"}, + "candidate_id": {"$ref": "#/$defs/nonEmptyString"}, + "domain": {"enum": ["us_equity", "cn_equity", "hk_equity", "crypto"]}, + "no_order": {"const": true}, + "permission_effect": {"const": "none"}, + "catalog_status_used_as_evidence": {"const": false}, + "stages": { + "type": "object", + "additionalProperties": false, + "required": ["p1_input", "p2_freeze", "p3_evidence"], + "properties": { + "p1_input": { + "$ref": "#/$defs/stageRecord", + "properties": { + "stage": {"const": "P1"}, + "artifact": { + "oneOf": [ + {"type": "null"}, + { + "$ref": "#/$defs/artifactIdentity", + "properties": {"schema_version": {"const": "research_input_manifest.v1"}} + } + ] + } + } + }, + "p2_freeze": { + "$ref": "#/$defs/stageRecord", + "properties": { + "stage": {"const": "P2"}, + "artifact": { + "oneOf": [ + {"type": "null"}, + { + "$ref": "#/$defs/artifactIdentity", + "properties": {"schema_version": {"const": "strategy_config_freeze.v1"}} + } + ] + } + } + }, + "p3_evidence": { + "$ref": "#/$defs/stageRecord", + "properties": { + "stage": {"const": "P3"}, + "artifact": { + "oneOf": [ + {"type": "null"}, + { + "$ref": "#/$defs/artifactIdentity", + "properties": {"schema_version": {"const": "strategy_evidence_package.v2"}} + } + ] + } + } + } + } + } + }, + "$defs": { + "nonEmptyString": {"type": "string", "minLength": 1, "pattern": ".*\\S.*"}, + "sha256": {"type": "string", "pattern": "^[0-9a-f]{64}$"}, + "reasonCode": {"type": "string", "pattern": "^[a-z][a-z0-9_]{0,127}$"}, + "artifactIdentity": { + "type": "object", + "additionalProperties": false, + "required": ["artifact_id", "schema_version", "sha256"], + "properties": { + "artifact_id": {"$ref": "#/$defs/nonEmptyString"}, + "schema_version": {"type": "string"}, + "sha256": {"$ref": "#/$defs/sha256"} + } + }, + "stageRecord": { + "type": "object", + "additionalProperties": false, + "required": ["stage", "status", "artifact", "reason_codes"], + "properties": { + "stage": {"enum": ["P1", "P2", "P3"]}, + "status": {"enum": ["READY", "DEFERRED", "PARKED"]}, + "artifact": {"oneOf": [{"type": "null"}, {"$ref": "#/$defs/artifactIdentity"}]}, + "reason_codes": { + "type": "array", + "uniqueItems": true, + "items": {"$ref": "#/$defs/reasonCode"} + } + }, + "allOf": [ + { + "if": {"properties": {"status": {"const": "READY"}}, "required": ["status"]}, + "then": { + "properties": { + "artifact": {"type": "object"}, + "reason_codes": {"maxItems": 0} + } + }, + "else": { + "properties": { + "artifact": {"type": "null"}, + "reason_codes": {"minItems": 1} + } + } + } + ] + } + } +} diff --git a/src/quant_platform_kit/strategy_lifecycle/__init__.py b/src/quant_platform_kit/strategy_lifecycle/__init__.py index e6cf2636..94301d81 100644 --- a/src/quant_platform_kit/strategy_lifecycle/__init__.py +++ b/src/quant_platform_kit/strategy_lifecycle/__init__.py @@ -35,6 +35,19 @@ read_evidence_package_v2_json, validate_evidence_package_v2, ) +from quant_platform_kit.strategy_lifecycle.forward_risk_driver import ( + FORWARD_RISK_SCHEMA_VERSION, + FORWARD_RISK_TERMINAL_STATUSES, + P4_OBSERVATION_MODES, + InvalidForwardRiskArtifact, + build_forward_risk_terminal_artifact, + build_nonready_forward_risk_stage, + build_ready_forward_observation_stage, + build_ready_portfolio_risk_stage, + canonical_forward_risk_terminal_bytes, + forward_risk_terminal_sha256, + validate_forward_risk_terminal_artifact, +) from quant_platform_kit.strategy_lifecycle.live_candidate_notifications import ( LiveCandidateNotificationEvent, build_live_candidate_notification, @@ -47,6 +60,18 @@ normalize_catalog_lifecycle_status, require_canonical_lifecycle_write, ) +from quant_platform_kit.strategy_lifecycle.research_driver import ( + InvalidResearchDriverArtifact, + RESEARCH_DRIVER_DOMAINS, + RESEARCH_DRIVER_SCHEMA_VERSION, + RESEARCH_DRIVER_TERMINAL_STATUSES, + build_nonready_research_stage, + build_ready_research_stage, + build_research_driver_terminal_artifact, + canonical_research_driver_terminal_bytes, + research_driver_terminal_sha256, + validate_research_driver_terminal_artifact, +) __all__ = [ "BacktestResult", @@ -66,22 +91,43 @@ "WindowPerformance", "EvidenceGateResult", "EvidencePackage", + "FORWARD_RISK_SCHEMA_VERSION", + "FORWARD_RISK_TERMINAL_STATUSES", + "P4_OBSERVATION_MODES", + "InvalidForwardRiskArtifact", "LiveCandidateNotificationEvent", "CANONICAL_LIFECYCLE_STATES", "LEGACY_CATALOG_STATUS_MAP", + "InvalidResearchDriverArtifact", + "RESEARCH_DRIVER_DOMAINS", + "RESEARCH_DRIVER_SCHEMA_VERSION", + "RESEARCH_DRIVER_TERMINAL_STATUSES", "load_evidence_package", "canonical_evidence_package_v2_bytes", "read_evidence_package_v2_json", "build_live_candidate_notification", + "build_forward_risk_terminal_artifact", + "build_nonready_forward_risk_stage", + "build_nonready_research_stage", + "build_ready_forward_observation_stage", + "build_ready_portfolio_risk_stage", + "build_ready_research_stage", + "build_research_driver_terminal_artifact", + "canonical_research_driver_terminal_bytes", + "canonical_forward_risk_terminal_bytes", "catalog_status_grants_execution_permission", "migrate_legacy_lifecycle_status", "normalize_catalog_lifecycle_status", "require_canonical_lifecycle_write", + "research_driver_terminal_sha256", + "forward_risk_terminal_sha256", "validate_evidence_package", "validate_evidence_package_file", "validate_evidence_package_v2", "validate_optimization_spec", "validate_research_spec", + "validate_research_driver_terminal_artifact", + "validate_forward_risk_terminal_artifact", "validate_strategy_spec", "validate_strategy_spec_file", ] diff --git a/src/quant_platform_kit/strategy_lifecycle/forward_risk_driver.py b/src/quant_platform_kit/strategy_lifecycle/forward_risk_driver.py new file mode 100644 index 00000000..d190e782 --- /dev/null +++ b/src/quant_platform_kit/strategy_lifecycle/forward_risk_driver.py @@ -0,0 +1,496 @@ +"""Cross-asset, no-order P4/P5 terminal observation contract. + +This adapter binds future/shadow observations and portfolio-risk snapshots to +an already terminal P1-P3 research run. It only validates immutable artifact +identities. It never starts a paper account, fetches market data, calls a +broker, changes a position, or grants lifecycle authority. +""" + +from __future__ import annotations + +from collections.abc import Mapping, Sequence +from datetime import datetime, timezone +from hashlib import sha256 +import json +import re +from typing import Any + +from .research_driver import ( + RESEARCH_DRIVER_DOMAINS, + research_driver_terminal_sha256, + validate_research_driver_terminal_artifact, +) + + +FORWARD_RISK_SCHEMA_VERSION = "forward_risk_terminal.v1" +FORWARD_RISK_TERMINAL_STATUSES = frozenset({"READY", "DEFERRED", "PARKED"}) +P4_OBSERVATION_MODES = frozenset({"shadow", "paper"}) + +_SHA256_RE = re.compile(r"^[0-9a-f]{64}$") +_REASON_CODE_RE = re.compile(r"^[a-z][a-z0-9_]{0,127}$") +_STAGE_FIELDS = frozenset({"stage", "status", "mode", "artifact", "reason_codes"}) +_BASE_ARTIFACT_FIELDS = frozenset( + { + "artifact_id", + "schema_version", + "sha256", + "candidate_id", + "observed_at", + "expires_at", + } +) +_TOP_LEVEL_FIELDS = frozenset( + { + "schema_version", + "terminal", + "terminal_status", + "generated_at", + "run_id", + "strategy_id", + "candidate_id", + "domain", + "research_terminal_status", + "research_terminal_sha256", + "no_order", + "permission_effect", + "broker_dependency", + "stages", + } +) + + +class InvalidForwardRiskArtifact(ValueError): + """Raised when a P4/P5 observation artifact fails closed validation.""" + + +def _invalid(message: str) -> None: + raise InvalidForwardRiskArtifact(message) + + +def _nonblank_string(value: object, field: str) -> str: + if not isinstance(value, str) or not value.strip(): + _invalid(f"{field} must be a non-empty string") + if any(ord(character) < 0x20 or ord(character) == 0x7F for character in value): + _invalid(f"{field} contains a control character") + return value.strip() + + +def _timestamp(value: object, field: str) -> tuple[str, datetime]: + text = _nonblank_string(value, field) + try: + parsed = datetime.fromisoformat(text.replace("Z", "+00:00")) + except ValueError: + _invalid(f"{field} must be an ISO-8601 timestamp") + if parsed.tzinfo is None or parsed.utcoffset() is None: + _invalid(f"{field} must include a timezone") + return text, parsed.astimezone(timezone.utc) + + +def _reason_codes(values: object, field: str, *, required: bool) -> list[str]: + if isinstance(values, (str, bytes)) or not isinstance(values, Sequence): + _invalid(f"{field} must be an array") + normalized = [_nonblank_string(value, field) for value in values] + if any(not _REASON_CODE_RE.fullmatch(value) for value in normalized): + _invalid(f"{field} contains an invalid reason code") + if normalized != sorted(set(normalized)): + _invalid(f"{field} must be sorted and unique") + if required and not normalized: + _invalid(f"{field} must explain a non-ready stage") + if not required and normalized: + _invalid(f"{field} must be empty for READY") + return normalized + + +def _artifact_identity( + value: object, + *, + field: str, + expected_schema_version: str, + expected_candidate_id: str | None = None, + generated_at: datetime | None = None, +) -> dict[str, str]: + if not isinstance(value, Mapping) or set(value) != _BASE_ARTIFACT_FIELDS: + _invalid(f"{field} must be a closed artifact identity") + digest = _nonblank_string(value.get("sha256"), f"{field}.sha256") + if not _SHA256_RE.fullmatch(digest): + _invalid(f"{field}.sha256 must be a lowercase SHA-256 digest") + schema_version = _nonblank_string( + value.get("schema_version"), f"{field}.schema_version" + ) + if schema_version != expected_schema_version: + _invalid(f"{field}.schema_version must equal {expected_schema_version}") + candidate_id = _nonblank_string( + value.get("candidate_id"), f"{field}.candidate_id" + ) + if expected_candidate_id is not None and candidate_id != expected_candidate_id: + _invalid(f"{field}.candidate_id does not match the research terminal") + observed_at, observed_time = _timestamp( + value.get("observed_at"), f"{field}.observed_at" + ) + expires_at, expiry_time = _timestamp( + value.get("expires_at"), f"{field}.expires_at" + ) + if expiry_time <= observed_time: + _invalid(f"{field}.expires_at must be later than observed_at") + if generated_at is not None and expiry_time <= generated_at: + _invalid(f"{field} is expired at generated_at") + return { + "artifact_id": _nonblank_string( + value.get("artifact_id"), f"{field}.artifact_id" + ), + "schema_version": schema_version, + "sha256": digest, + "candidate_id": candidate_id, + "observed_at": observed_at, + "expires_at": expires_at, + } + + +def _build_ready_stage( + *, + stage: str, + mode: str, + schema_version: str, + artifact_id: str, + artifact_sha256: str, + candidate_id: str, + observed_at: str, + expires_at: str, +) -> dict[str, Any]: + return { + "stage": stage, + "status": "READY", + "mode": mode, + "artifact": _artifact_identity( + { + "artifact_id": artifact_id, + "schema_version": schema_version, + "sha256": artifact_sha256, + "candidate_id": candidate_id, + "observed_at": observed_at, + "expires_at": expires_at, + }, + field=f"{stage.lower()}.artifact", + expected_schema_version=schema_version, + ), + "reason_codes": [], + } + + +def build_ready_forward_observation_stage( + *, + mode: str, + artifact_id: str, + artifact_sha256: str, + candidate_id: str, + observed_at: str, + expires_at: str, +) -> dict[str, Any]: + """Build a P4 identity for validated shadow or simulated-paper evidence.""" + + normalized_mode = _nonblank_string(mode, "mode").lower() + if normalized_mode not in P4_OBSERVATION_MODES: + _invalid("P4 mode must be shadow or paper") + return _build_ready_stage( + stage="P4", + mode=normalized_mode, + schema_version="forward_observation.v1", + artifact_id=artifact_id, + artifact_sha256=artifact_sha256, + candidate_id=candidate_id, + observed_at=observed_at, + expires_at=expires_at, + ) + + +def build_ready_portfolio_risk_stage( + *, + artifact_id: str, + artifact_sha256: str, + candidate_id: str, + observed_at: str, + expires_at: str, +) -> dict[str, Any]: + """Build a P5 identity for an already validated portfolio RiskSnapshot.""" + + return _build_ready_stage( + stage="P5", + mode="portfolio_risk", + schema_version="portfolio_risk_snapshot.v1", + artifact_id=artifact_id, + artifact_sha256=artifact_sha256, + candidate_id=candidate_id, + observed_at=observed_at, + expires_at=expires_at, + ) + + +def build_nonready_forward_risk_stage( + stage: str, *, status: str, reason_codes: Sequence[str] +) -> dict[str, Any]: + """Build a truthful P4/P5 DEFERRED or PARKED terminal stage.""" + + normalized_stage = _nonblank_string(stage, "stage").upper() + if normalized_stage not in {"P4", "P5"}: + _invalid("stage must be P4 or P5") + normalized_status = _nonblank_string(status, "status").upper() + if normalized_status not in {"DEFERRED", "PARKED"}: + _invalid("non-ready stage status must be DEFERRED or PARKED") + return { + "stage": normalized_stage, + "status": normalized_status, + "mode": "shadow" if normalized_stage == "P4" else "portfolio_risk", + "artifact": None, + "reason_codes": _reason_codes( + list(reason_codes), "reason_codes", required=True + ), + } + + +def _validate_stage( + value: object, + *, + stage: str, + candidate_id: str, + generated_at: datetime, +) -> dict[str, Any]: + field = "p4_forward" if stage == "P4" else "p5_risk" + if not isinstance(value, Mapping) or set(value) != _STAGE_FIELDS: + _invalid(f"{field} must be a closed stage record") + if value.get("stage") != stage: + _invalid(f"{field}.stage must equal {stage}") + status = _nonblank_string(value.get("status"), f"{field}.status").upper() + if status not in FORWARD_RISK_TERMINAL_STATUSES: + _invalid(f"{field}.status is unsupported") + mode = _nonblank_string(value.get("mode"), f"{field}.mode").lower() + allowed_modes = P4_OBSERVATION_MODES if stage == "P4" else {"portfolio_risk"} + if mode not in allowed_modes: + _invalid(f"{field}.mode is unsupported") + if status == "READY": + artifact = _artifact_identity( + value.get("artifact"), + field=f"{field}.artifact", + expected_schema_version=( + "forward_observation.v1" + if stage == "P4" + else "portfolio_risk_snapshot.v1" + ), + expected_candidate_id=candidate_id, + generated_at=generated_at, + ) + reasons = _reason_codes( + value.get("reason_codes"), f"{field}.reason_codes", required=False + ) + else: + if value.get("artifact") is not None: + _invalid(f"{field}.artifact must be null unless status is READY") + artifact = None + reasons = _reason_codes( + value.get("reason_codes"), f"{field}.reason_codes", required=True + ) + return { + "stage": stage, + "status": status, + "mode": mode, + "artifact": artifact, + "reason_codes": reasons, + } + + +def _normalize_stage( + value: object | None, + *, + stage: str, + candidate_id: str, + generated_at: datetime, +) -> dict[str, Any]: + if value is None: + return build_nonready_forward_risk_stage( + stage, + status="DEFERRED", + reason_codes=(("p4_observation_not_produced" if stage == "P4" else "p5_risk_not_produced"),), + ) + try: + return _validate_stage( + value, + stage=stage, + candidate_id=candidate_id, + generated_at=generated_at, + ) + except InvalidForwardRiskArtifact: + return build_nonready_forward_risk_stage( + stage, + status="PARKED", + reason_codes=(("p4_observation_invalid" if stage == "P4" else "p5_risk_invalid"),), + ) + + +def _terminal_status( + research_status: str, stages: Mapping[str, Mapping[str, Any]] +) -> str: + statuses = {research_status, *(str(value["status"]) for value in stages.values())} + if "PARKED" in statuses: + return "PARKED" + if statuses == {"READY"}: + return "READY" + return "DEFERRED" + + +def build_forward_risk_terminal_artifact( + *, + research_terminal: Mapping[str, Any], + generated_at: str, + p4_forward: Mapping[str, Any] | None = None, + p5_risk: Mapping[str, Any] | None = None, +) -> dict[str, Any]: + """Bind P4/P5 evidence to a validated P1-P3 terminal artifact.""" + + research = validate_research_driver_terminal_artifact(research_terminal) + generated_at_text, generated_time = _timestamp(generated_at, "generated_at") + candidate_id = str(research["candidate_id"]) + stages = { + "p4_forward": _normalize_stage( + p4_forward, + stage="P4", + candidate_id=candidate_id, + generated_at=generated_time, + ), + "p5_risk": _normalize_stage( + p5_risk, + stage="P5", + candidate_id=candidate_id, + generated_at=generated_time, + ), + } + if stages["p5_risk"]["status"] == "READY" and stages["p4_forward"]["status"] != "READY": + stages["p5_risk"] = build_nonready_forward_risk_stage( + "P5", status="PARKED", reason_codes=("p4_forward_not_ready",) + ) + if research["terminal_status"] != "READY": + for key, stage_name in (("p4_forward", "P4"), ("p5_risk", "P5")): + if stages[key]["status"] == "READY": + stages[key] = build_nonready_forward_risk_stage( + stage_name, + status="PARKED", + reason_codes=("research_terminal_not_ready",), + ) + artifact = { + "schema_version": FORWARD_RISK_SCHEMA_VERSION, + "terminal": True, + "terminal_status": _terminal_status(str(research["terminal_status"]), stages), + "generated_at": generated_at_text, + "run_id": str(research["run_id"]), + "strategy_id": str(research["strategy_id"]), + "candidate_id": candidate_id, + "domain": str(research["domain"]), + "research_terminal_status": str(research["terminal_status"]), + "research_terminal_sha256": research_driver_terminal_sha256(research), + "no_order": True, + "permission_effect": "none", + "broker_dependency": False, + "stages": stages, + } + return validate_forward_risk_terminal_artifact(artifact) + + +def validate_forward_risk_terminal_artifact( + artifact: Mapping[str, Any], +) -> dict[str, Any]: + """Validate a closed P4/P5 terminal envelope without side effects.""" + + if not isinstance(artifact, Mapping) or set(artifact) != _TOP_LEVEL_FIELDS: + _invalid("terminal artifact must be a closed object") + if artifact.get("schema_version") != FORWARD_RISK_SCHEMA_VERSION: + _invalid(f"schema_version must equal {FORWARD_RISK_SCHEMA_VERSION}") + if artifact.get("terminal") is not True: + _invalid("terminal must remain true") + if artifact.get("no_order") is not True: + _invalid("no_order must remain true") + if artifact.get("permission_effect") != "none": + _invalid("permission_effect must remain none") + if artifact.get("broker_dependency") is not False: + _invalid("broker_dependency must remain false") + generated_at, generated_time = _timestamp(artifact.get("generated_at"), "generated_at") + domain = _nonblank_string(artifact.get("domain"), "domain").lower() + if domain not in RESEARCH_DRIVER_DOMAINS: + _invalid("domain is unsupported") + candidate_id = _nonblank_string(artifact.get("candidate_id"), "candidate_id") + digest = _nonblank_string( + artifact.get("research_terminal_sha256"), "research_terminal_sha256" + ) + if not _SHA256_RE.fullmatch(digest): + _invalid("research_terminal_sha256 must be a lowercase SHA-256 digest") + research_status = _nonblank_string( + artifact.get("research_terminal_status"), "research_terminal_status" + ).upper() + if research_status not in FORWARD_RISK_TERMINAL_STATUSES: + _invalid("research_terminal_status is unsupported") + stages_value = artifact.get("stages") + if not isinstance(stages_value, Mapping) or set(stages_value) != {"p4_forward", "p5_risk"}: + _invalid("stages must contain exactly P4 and P5") + stages = { + "p4_forward": _validate_stage( + stages_value["p4_forward"], + stage="P4", + candidate_id=candidate_id, + generated_at=generated_time, + ), + "p5_risk": _validate_stage( + stages_value["p5_risk"], + stage="P5", + candidate_id=candidate_id, + generated_at=generated_time, + ), + } + if stages["p5_risk"]["status"] == "READY" and stages["p4_forward"]["status"] != "READY": + _invalid("P5 READY cannot bypass a non-ready P4") + if research_status != "READY" and any( + value["status"] == "READY" for value in stages.values() + ): + _invalid("P4/P5 READY cannot bypass a non-ready research terminal") + terminal_status = _nonblank_string( + artifact.get("terminal_status"), "terminal_status" + ).upper() + if terminal_status != _terminal_status(research_status, stages): + _invalid("terminal_status does not match upstream and stage results") + result = { + "schema_version": FORWARD_RISK_SCHEMA_VERSION, + "terminal": True, + "terminal_status": terminal_status, + "generated_at": generated_at, + "run_id": _nonblank_string(artifact.get("run_id"), "run_id"), + "strategy_id": _nonblank_string( + artifact.get("strategy_id"), "strategy_id" + ), + "candidate_id": candidate_id, + "domain": domain, + "research_terminal_status": research_status, + "research_terminal_sha256": digest, + "no_order": True, + "permission_effect": "none", + "broker_dependency": False, + "stages": stages, + } + return json.loads( + json.dumps(result, sort_keys=True, separators=(",", ":"), allow_nan=False) + ) + + +def canonical_forward_risk_terminal_bytes(artifact: Mapping[str, Any]) -> bytes: + """Return deterministic UTF-8 bytes for a validated P4/P5 artifact.""" + + validated = validate_forward_risk_terminal_artifact(artifact) + return json.dumps( + validated, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=False, + allow_nan=False, + ).encode("utf-8") + + +def forward_risk_terminal_sha256(artifact: Mapping[str, Any]) -> str: + """Return the SHA-256 digest of the canonical P4/P5 terminal bytes.""" + + return sha256(canonical_forward_risk_terminal_bytes(artifact)).hexdigest() + diff --git a/src/quant_platform_kit/strategy_lifecycle/research_driver.py b/src/quant_platform_kit/strategy_lifecycle/research_driver.py new file mode 100644 index 00000000..a70c437b --- /dev/null +++ b/src/quant_platform_kit/strategy_lifecycle/research_driver.py @@ -0,0 +1,373 @@ +"""Cross-asset, research-only P1-P3 terminal driver contract. + +The driver is deliberately a pure evidence-envelope builder. It does not +fetch data, run a backtest, inspect a catalog, call a broker, or grant a +lifecycle permission. Producers validate their P1/P2/P3 artifacts first and +pass only their immutable identities into this boundary. +""" + +from __future__ import annotations + +from collections.abc import Mapping, Sequence +from datetime import datetime +from hashlib import sha256 +import json +import re +from typing import Any + + +RESEARCH_DRIVER_SCHEMA_VERSION = "research_driver_terminal.v1" +RESEARCH_DRIVER_DOMAINS = frozenset( + {"us_equity", "cn_equity", "hk_equity", "crypto"} +) +RESEARCH_DRIVER_TERMINAL_STATUSES = frozenset( + {"READY", "DEFERRED", "PARKED"} +) + +_STAGE_SPECS = { + "p1_input": ("P1", "research_input_manifest.v1"), + "p2_freeze": ("P2", "strategy_config_freeze.v1"), + "p3_evidence": ("P3", "strategy_evidence_package.v2"), +} +_SHA256_RE = re.compile(r"^[0-9a-f]{64}$") +_REASON_CODE_RE = re.compile(r"^[a-z][a-z0-9_]{0,127}$") +_ARTIFACT_FIELDS = frozenset({"artifact_id", "schema_version", "sha256"}) +_STAGE_FIELDS = frozenset({"stage", "status", "artifact", "reason_codes"}) +_TOP_LEVEL_FIELDS = frozenset( + { + "schema_version", + "terminal", + "terminal_status", + "generated_at", + "run_id", + "strategy_id", + "candidate_id", + "domain", + "no_order", + "permission_effect", + "catalog_status_used_as_evidence", + "stages", + } +) + + +class InvalidResearchDriverArtifact(ValueError): + """Raised when a terminal research-driver artifact is not trustworthy.""" + + +def _invalid(message: str) -> None: + raise InvalidResearchDriverArtifact(message) + + +def _nonblank_string(value: object, field: str) -> str: + if not isinstance(value, str) or not value.strip(): + _invalid(f"{field} must be a non-empty string") + if any(ord(character) < 0x20 or ord(character) == 0x7F for character in value): + _invalid(f"{field} contains a control character") + return value.strip() + + +def _timezone_timestamp(value: object, field: str) -> str: + text = _nonblank_string(value, field) + try: + parsed = datetime.fromisoformat(text.replace("Z", "+00:00")) + except ValueError: + _invalid(f"{field} must be an ISO-8601 timestamp") + if parsed.tzinfo is None or parsed.utcoffset() is None: + _invalid(f"{field} must include a timezone") + return text + + +def _reason_codes(values: object, field: str, *, required: bool) -> list[str]: + if isinstance(values, (str, bytes)) or not isinstance(values, Sequence): + _invalid(f"{field} must be an array") + normalized: list[str] = [] + for value in values: + code = _nonblank_string(value, field) + if not _REASON_CODE_RE.fullmatch(code): + _invalid(f"{field} contains an invalid reason code") + normalized.append(code) + if normalized != sorted(set(normalized)): + _invalid(f"{field} must be sorted and unique") + if required and not normalized: + _invalid(f"{field} must explain a non-ready stage") + if not required and normalized: + _invalid(f"{field} must be empty for READY") + return normalized + + +def _artifact_identity( + value: object, *, field: str, expected_schema_version: str +) -> dict[str, str]: + if not isinstance(value, Mapping) or set(value) != _ARTIFACT_FIELDS: + _invalid(f"{field} must be a closed artifact identity") + artifact_id = _nonblank_string(value.get("artifact_id"), f"{field}.artifact_id") + schema_version = _nonblank_string( + value.get("schema_version"), f"{field}.schema_version" + ) + if schema_version != expected_schema_version: + _invalid( + f"{field}.schema_version must equal {expected_schema_version}" + ) + digest = _nonblank_string(value.get("sha256"), f"{field}.sha256") + if not _SHA256_RE.fullmatch(digest): + _invalid(f"{field}.sha256 must be a lowercase SHA-256 digest") + return { + "artifact_id": artifact_id, + "schema_version": schema_version, + "sha256": digest, + } + + +def build_ready_research_stage( + stage_name: str, *, artifact_id: str, artifact_sha256: str +) -> dict[str, Any]: + """Build a READY stage from an already validated immutable artifact.""" + + if stage_name not in _STAGE_SPECS: + _invalid(f"unsupported research stage: {stage_name!r}") + stage, schema_version = _STAGE_SPECS[stage_name] + record = { + "stage": stage, + "status": "READY", + "artifact": { + "artifact_id": artifact_id, + "schema_version": schema_version, + "sha256": artifact_sha256, + }, + "reason_codes": [], + } + return _validate_stage_record(stage_name, record) + + +def build_nonready_research_stage( + stage_name: str, *, status: str, reason_codes: Sequence[str] +) -> dict[str, Any]: + """Build a truthful DEFERRED or PARKED stage without claiming evidence.""" + + if stage_name not in _STAGE_SPECS: + _invalid(f"unsupported research stage: {stage_name!r}") + normalized_status = _nonblank_string(status, "status").upper() + if normalized_status not in {"DEFERRED", "PARKED"}: + _invalid("non-ready stage status must be DEFERRED or PARKED") + stage, _ = _STAGE_SPECS[stage_name] + record = { + "stage": stage, + "status": normalized_status, + "artifact": None, + "reason_codes": list(reason_codes), + } + return _validate_stage_record(stage_name, record) + + +def _validate_stage_record(stage_name: str, value: object) -> dict[str, Any]: + expected_stage, expected_schema_version = _STAGE_SPECS[stage_name] + if not isinstance(value, Mapping) or set(value) != _STAGE_FIELDS: + _invalid(f"{stage_name} must be a closed stage record") + stage = _nonblank_string(value.get("stage"), f"{stage_name}.stage") + if stage != expected_stage: + _invalid(f"{stage_name}.stage must equal {expected_stage}") + status = _nonblank_string(value.get("status"), f"{stage_name}.status").upper() + if status not in RESEARCH_DRIVER_TERMINAL_STATUSES: + _invalid(f"{stage_name}.status is unsupported") + if status == "READY": + artifact = _artifact_identity( + value.get("artifact"), + field=f"{stage_name}.artifact", + expected_schema_version=expected_schema_version, + ) + reasons = _reason_codes( + value.get("reason_codes"), f"{stage_name}.reason_codes", required=False + ) + else: + if value.get("artifact") is not None: + _invalid(f"{stage_name}.artifact must be null unless status is READY") + artifact = None + reasons = _reason_codes( + value.get("reason_codes"), f"{stage_name}.reason_codes", required=True + ) + return { + "stage": expected_stage, + "status": status, + "artifact": artifact, + "reason_codes": reasons, + } + + +def _normalize_stage(stage_name: str, value: object) -> dict[str, Any]: + if value is None: + return build_nonready_research_stage( + stage_name, + status="DEFERRED", + reason_codes=(f"{stage_name}_not_produced",), + ) + try: + return _validate_stage_record(stage_name, value) + except InvalidResearchDriverArtifact: + return build_nonready_research_stage( + stage_name, + status="PARKED", + reason_codes=(f"{stage_name}_invalid",), + ) + + +def _enforce_dependencies(stages: dict[str, dict[str, Any]]) -> None: + dependencies = { + "p2_freeze": "p1_input", + "p3_evidence": "p2_freeze", + } + for stage_name, dependency in dependencies.items(): + if ( + stages[stage_name]["status"] == "READY" + and stages[dependency]["status"] != "READY" + ): + stages[stage_name] = build_nonready_research_stage( + stage_name, + status="PARKED", + reason_codes=(f"{dependency}_not_ready",), + ) + + +def _terminal_status(stages: Mapping[str, Mapping[str, Any]]) -> str: + statuses = {str(stage["status"]) for stage in stages.values()} + if "PARKED" in statuses: + return "PARKED" + if statuses == {"READY"}: + return "READY" + return "DEFERRED" + + +def build_research_driver_terminal_artifact( + *, + run_id: str, + generated_at: str, + strategy_id: str, + candidate_id: str, + domain: str, + p1_input: Mapping[str, Any] | None = None, + p2_freeze: Mapping[str, Any] | None = None, + p3_evidence: Mapping[str, Any] | None = None, +) -> dict[str, Any]: + """Return one terminal, no-order P1-P3 artifact for every valid run identity. + + Missing stages become ``DEFERRED``. Malformed evidence or an impossible + dependency chain becomes ``PARKED``. Catalog state is intentionally not an + input and cannot contribute to readiness. + """ + + normalized_domain = _nonblank_string(domain, "domain").lower() + if normalized_domain not in RESEARCH_DRIVER_DOMAINS: + _invalid(f"unsupported research domain: {domain!r}") + stages = { + "p1_input": _normalize_stage("p1_input", p1_input), + "p2_freeze": _normalize_stage("p2_freeze", p2_freeze), + "p3_evidence": _normalize_stage("p3_evidence", p3_evidence), + } + _enforce_dependencies(stages) + artifact = { + "schema_version": RESEARCH_DRIVER_SCHEMA_VERSION, + "terminal": True, + "terminal_status": _terminal_status(stages), + "generated_at": _timezone_timestamp(generated_at, "generated_at"), + "run_id": _nonblank_string(run_id, "run_id"), + "strategy_id": _nonblank_string(strategy_id, "strategy_id"), + "candidate_id": _nonblank_string(candidate_id, "candidate_id"), + "domain": normalized_domain, + "no_order": True, + "permission_effect": "none", + "catalog_status_used_as_evidence": False, + "stages": stages, + } + return validate_research_driver_terminal_artifact(artifact) + + +def validate_research_driver_terminal_artifact( + artifact: Mapping[str, Any], +) -> dict[str, Any]: + """Validate a closed terminal artifact and return a detached JSON value.""" + + if not isinstance(artifact, Mapping) or set(artifact) != _TOP_LEVEL_FIELDS: + _invalid("terminal artifact must be a closed object") + if artifact.get("schema_version") != RESEARCH_DRIVER_SCHEMA_VERSION: + _invalid(f"schema_version must equal {RESEARCH_DRIVER_SCHEMA_VERSION}") + if artifact.get("terminal") is not True: + _invalid("terminal must remain true") + if artifact.get("no_order") is not True: + _invalid("no_order must remain true") + if artifact.get("permission_effect") != "none": + _invalid("permission_effect must remain none") + if artifact.get("catalog_status_used_as_evidence") is not False: + _invalid("catalog status cannot be used as evidence") + domain = _nonblank_string(artifact.get("domain"), "domain").lower() + if domain not in RESEARCH_DRIVER_DOMAINS: + _invalid("domain is unsupported") + stages_value = artifact.get("stages") + if not isinstance(stages_value, Mapping) or set(stages_value) != set(_STAGE_SPECS): + _invalid("stages must contain exactly P1, P2, and P3") + stages = { + stage_name: _validate_stage_record(stage_name, stages_value[stage_name]) + for stage_name in _STAGE_SPECS + } + dependency_copy = json.loads(json.dumps(stages)) + _enforce_dependencies(dependency_copy) + if dependency_copy != stages: + _invalid("a READY stage cannot bypass a non-ready dependency") + terminal_status = _nonblank_string( + artifact.get("terminal_status"), "terminal_status" + ).upper() + if terminal_status != _terminal_status(stages): + _invalid("terminal_status does not match stage results") + result = { + "schema_version": RESEARCH_DRIVER_SCHEMA_VERSION, + "terminal": True, + "terminal_status": terminal_status, + "generated_at": _timezone_timestamp(artifact.get("generated_at"), "generated_at"), + "run_id": _nonblank_string(artifact.get("run_id"), "run_id"), + "strategy_id": _nonblank_string( + artifact.get("strategy_id"), "strategy_id" + ), + "candidate_id": _nonblank_string( + artifact.get("candidate_id"), "candidate_id" + ), + "domain": domain, + "no_order": True, + "permission_effect": "none", + "catalog_status_used_as_evidence": False, + "stages": stages, + } + return json.loads( + json.dumps(result, sort_keys=True, separators=(",", ":"), allow_nan=False) + ) + + +def canonical_research_driver_terminal_bytes(artifact: Mapping[str, Any]) -> bytes: + """Return deterministic UTF-8 JSON bytes for a validated artifact.""" + + validated = validate_research_driver_terminal_artifact(artifact) + return json.dumps( + validated, + sort_keys=True, + separators=(",", ":"), + ensure_ascii=False, + allow_nan=False, + ).encode("utf-8") + + +def research_driver_terminal_sha256(artifact: Mapping[str, Any]) -> str: + """Return the digest of the canonical terminal artifact bytes.""" + + return sha256(canonical_research_driver_terminal_bytes(artifact)).hexdigest() + + +__all__ = [ + "InvalidResearchDriverArtifact", + "RESEARCH_DRIVER_DOMAINS", + "RESEARCH_DRIVER_SCHEMA_VERSION", + "RESEARCH_DRIVER_TERMINAL_STATUSES", + "build_nonready_research_stage", + "build_ready_research_stage", + "build_research_driver_terminal_artifact", + "canonical_research_driver_terminal_bytes", + "research_driver_terminal_sha256", + "validate_research_driver_terminal_artifact", +] diff --git a/tests/test_cross_asset_forward_risk_driver.py b/tests/test_cross_asset_forward_risk_driver.py new file mode 100644 index 00000000..7b5dcacd --- /dev/null +++ b/tests/test_cross_asset_forward_risk_driver.py @@ -0,0 +1,212 @@ +from __future__ import annotations + +import copy +import json +from pathlib import Path + +import pytest +from jsonschema import Draft202012Validator, FormatChecker + +from quant_platform_kit.strategy_lifecycle.forward_risk_driver import ( + InvalidForwardRiskArtifact, + build_forward_risk_terminal_artifact, + build_nonready_forward_risk_stage, + build_ready_forward_observation_stage, + build_ready_portfolio_risk_stage, + canonical_forward_risk_terminal_bytes, + forward_risk_terminal_sha256, + validate_forward_risk_terminal_artifact, +) +from quant_platform_kit.strategy_lifecycle.research_driver import ( + RESEARCH_DRIVER_DOMAINS, + build_nonready_research_stage, + build_ready_research_stage, + build_research_driver_terminal_artifact, +) + + +ROOT = Path(__file__).parents[1] +SCHEMA_PATH = ROOT / "src/quant_platform_kit/schemas/forward-risk-terminal.v1.schema.json" + + +def _research_terminal(domain: str = "us_equity", *, ready: bool = True): + p3 = ( + build_ready_research_stage( + "p3_evidence", artifact_id="evidence-001", artifact_sha256="c" * 64 + ) + if ready + else build_nonready_research_stage( + "p3_evidence", status="DEFERRED", reason_codes=("evidence_pending",) + ) + ) + return build_research_driver_terminal_artifact( + run_id="daily-20260824", + generated_at="2026-08-24T20:00:00+08:00", + strategy_id="strategy-001", + candidate_id="candidate-001", + domain=domain, + p1_input=build_ready_research_stage( + "p1_input", artifact_id="manifest-001", artifact_sha256="a" * 64 + ), + p2_freeze=build_ready_research_stage( + "p2_freeze", artifact_id="freeze-001", artifact_sha256="b" * 64 + ), + p3_evidence=p3, + ) + + +def _p4(mode: str = "shadow"): + return build_ready_forward_observation_stage( + mode=mode, + artifact_id="forward-001", + artifact_sha256="d" * 64, + candidate_id="candidate-001", + observed_at="2026-08-24T20:30:00+08:00", + expires_at="2026-08-26T20:30:00+08:00", + ) + + +def _p5(): + return build_ready_portfolio_risk_stage( + artifact_id="risk-001", + artifact_sha256="e" * 64, + candidate_id="candidate-001", + observed_at="2026-08-24T20:40:00+08:00", + expires_at="2026-08-25T20:40:00+08:00", + ) + + +def _build(domain: str = "us_equity", **overrides): + values = {"p4_forward": _p4(), "p5_risk": _p5()} + values.update(overrides) + return build_forward_risk_terminal_artifact( + research_terminal=_research_terminal(domain), + generated_at="2026-08-24T21:00:00+08:00", + **values, + ) + + +@pytest.mark.parametrize("domain", sorted(RESEARCH_DRIVER_DOMAINS)) +def test_all_asset_domains_share_one_ready_p4_p5_contract(domain): + artifact = _build(domain) + assert artifact["terminal_status"] == "READY" + assert artifact["domain"] == domain + assert artifact["no_order"] is True + assert artifact["permission_effect"] == "none" + assert artifact["broker_dependency"] is False + + +@pytest.mark.parametrize("mode", ["shadow", "paper"]) +def test_p4_can_describe_shadow_or_existing_paper_evidence_without_broker_access(mode): + artifact = _build(p4_forward=_p4(mode)) + assert artifact["stages"]["p4_forward"]["mode"] == mode + assert artifact["broker_dependency"] is False + + +def test_missing_observations_still_emit_a_deferred_terminal_artifact(): + artifact = _build(p4_forward=None, p5_risk=None) + assert artifact["terminal_status"] == "DEFERRED" + assert artifact["stages"]["p4_forward"]["status"] == "DEFERRED" + assert artifact["stages"]["p5_risk"]["status"] == "DEFERRED" + + +def test_p5_cannot_bypass_missing_p4(): + artifact = _build(p4_forward=None) + assert artifact["terminal_status"] == "PARKED" + assert artifact["stages"]["p5_risk"] == { + "stage": "P5", + "status": "PARKED", + "mode": "portfolio_risk", + "artifact": None, + "reason_codes": ["p4_forward_not_ready"], + } + + +def test_ready_p4_p5_cannot_bypass_nonready_p1_p3_terminal(): + artifact = build_forward_risk_terminal_artifact( + research_terminal=_research_terminal(ready=False), + generated_at="2026-08-24T21:00:00+08:00", + p4_forward=_p4(), + p5_risk=_p5(), + ) + assert artifact["terminal_status"] == "PARKED" + assert artifact["research_terminal_status"] == "DEFERRED" + assert artifact["stages"]["p4_forward"]["reason_codes"] == [ + "research_terminal_not_ready" + ] + + +def test_stale_or_wrong_candidate_observation_is_parked(): + stale = _p4() + stale["artifact"]["expires_at"] = "2026-08-24T20:59:59+08:00" + artifact = _build(p4_forward=stale) + assert artifact["terminal_status"] == "PARKED" + assert artifact["stages"]["p4_forward"]["reason_codes"] == [ + "p4_observation_invalid" + ] + + wrong_candidate = _p4() + wrong_candidate["artifact"]["candidate_id"] = "candidate-002" + artifact = _build(p4_forward=wrong_candidate) + assert artifact["stages"]["p4_forward"]["status"] == "PARKED" + + +def test_python_output_matches_closed_json_schema(): + schema = json.loads(SCHEMA_PATH.read_text(encoding="utf-8")) + validator = Draft202012Validator(schema, format_checker=FormatChecker()) + assert list(validator.iter_errors(_build("crypto"))) == [] + + +@pytest.mark.parametrize( + ("field", "value"), + [ + ("no_order", False), + ("permission_effect", "live"), + ("broker_dependency", True), + ("terminal", False), + ], +) +def test_authority_guards_fail_closed(field, value): + artifact = _build() + artifact[field] = value + with pytest.raises(InvalidForwardRiskArtifact): + validate_forward_risk_terminal_artifact(artifact) + + +def test_validator_rejects_tampered_status_and_dependency_chain(): + artifact = _build() + artifact["terminal_status"] = "DEFERRED" + with pytest.raises(InvalidForwardRiskArtifact, match="terminal_status"): + validate_forward_risk_terminal_artifact(artifact) + + artifact = _build() + artifact["stages"]["p4_forward"] = build_nonready_forward_risk_stage( + "P4", status="DEFERRED", reason_codes=("forward_pending",) + ) + artifact["terminal_status"] = "DEFERRED" + with pytest.raises(InvalidForwardRiskArtifact, match="P5 READY"): + validate_forward_risk_terminal_artifact(artifact) + + +def test_canonical_digest_and_copy_isolation_are_deterministic(): + artifact = _build("hk_equity") + reordered = json.loads(json.dumps(artifact)) + assert canonical_forward_risk_terminal_bytes(reordered) == ( + canonical_forward_risk_terminal_bytes(artifact) + ) + assert forward_risk_terminal_sha256(reordered) == ( + forward_risk_terminal_sha256(artifact) + ) + validated = validate_forward_risk_terminal_artifact(artifact) + validated["stages"]["p4_forward"]["artifact"]["artifact_id"] = "changed" + assert artifact["stages"]["p4_forward"]["artifact"]["artifact_id"] == "forward-001" + + +def test_research_terminal_digest_is_bound_and_rejects_bad_format(): + artifact = _build() + assert len(artifact["research_terminal_sha256"]) == 64 + tampered = copy.deepcopy(artifact) + tampered["research_terminal_sha256"] = "sha256:not-valid" + with pytest.raises(InvalidForwardRiskArtifact, match="SHA-256"): + validate_forward_risk_terminal_artifact(tampered) + diff --git a/tests/test_cross_asset_research_driver.py b/tests/test_cross_asset_research_driver.py new file mode 100644 index 00000000..06f57cec --- /dev/null +++ b/tests/test_cross_asset_research_driver.py @@ -0,0 +1,177 @@ +from __future__ import annotations + +import copy +import json +from pathlib import Path + +import pytest +from jsonschema import Draft202012Validator, FormatChecker + +from quant_platform_kit.strategy_lifecycle.research_driver import ( + InvalidResearchDriverArtifact, + RESEARCH_DRIVER_DOMAINS, + build_nonready_research_stage, + build_ready_research_stage, + build_research_driver_terminal_artifact, + canonical_research_driver_terminal_bytes, + research_driver_terminal_sha256, + validate_research_driver_terminal_artifact, +) + + +ROOT = Path(__file__).parents[1] +SCHEMA_PATH = ( + ROOT + / "src/quant_platform_kit/schemas/research-driver-terminal.v1.schema.json" +) + + +def _ready_stages() -> dict[str, dict[str, object]]: + return { + "p1_input": build_ready_research_stage( + "p1_input", artifact_id="manifest-001", artifact_sha256="a" * 64 + ), + "p2_freeze": build_ready_research_stage( + "p2_freeze", artifact_id="freeze-001", artifact_sha256="b" * 64 + ), + "p3_evidence": build_ready_research_stage( + "p3_evidence", artifact_id="evidence-001", artifact_sha256="c" * 64 + ), + } + + +def _build(domain: str = "us_equity", **overrides): + stages = _ready_stages() + stages.update(overrides) + return build_research_driver_terminal_artifact( + run_id="daily-20260824", + generated_at="2026-08-24T22:00:00+08:00", + strategy_id="strategy-001", + candidate_id="candidate-001", + domain=domain, + **stages, + ) + + +@pytest.mark.parametrize("domain", sorted(RESEARCH_DRIVER_DOMAINS)) +def test_all_supported_asset_domains_share_one_ready_contract(domain): + artifact = _build(domain) + assert artifact["terminal"] is True + assert artifact["terminal_status"] == "READY" + assert artifact["domain"] == domain + assert artifact["no_order"] is True + assert artifact["permission_effect"] == "none" + assert artifact["catalog_status_used_as_evidence"] is False + + +def test_missing_stage_still_emits_truthful_deferred_terminal_artifact(): + artifact = _build(p3_evidence=None) + assert artifact["terminal_status"] == "DEFERRED" + assert artifact["stages"]["p3_evidence"] == { + "stage": "P3", + "status": "DEFERRED", + "artifact": None, + "reason_codes": ["p3_evidence_not_produced"], + } + + +def test_malformed_evidence_is_parked_instead_of_crashing_or_claiming_ready(): + malformed = _ready_stages()["p1_input"] + malformed["artifact"]["sha256"] = "not-a-digest" + artifact = _build(p1_input=malformed) + assert artifact["terminal_status"] == "PARKED" + assert artifact["stages"]["p1_input"]["reason_codes"] == [ + "p1_input_invalid" + ] + assert artifact["stages"]["p2_freeze"]["status"] == "PARKED" + assert artifact["stages"]["p3_evidence"]["status"] == "PARKED" + + +def test_ready_downstream_stage_cannot_bypass_a_deferred_dependency(): + deferred = build_nonready_research_stage( + "p1_input", status="DEFERRED", reason_codes=("input_not_available",) + ) + artifact = _build(p1_input=deferred) + assert artifact["terminal_status"] == "PARKED" + assert artifact["stages"]["p2_freeze"]["reason_codes"] == [ + "p1_input_not_ready" + ] + assert artifact["stages"]["p3_evidence"]["reason_codes"] == [ + "p2_freeze_not_ready" + ] + + +def test_catalog_metadata_cannot_be_smuggled_in_as_stage_evidence(): + catalog_record = _ready_stages()["p1_input"] + catalog_record["catalog_status"] = "live_enabled" + artifact = _build(p1_input=catalog_record) + assert artifact["terminal_status"] == "PARKED" + assert artifact["stages"]["p1_input"]["artifact"] is None + + +def test_python_output_matches_closed_json_schema(): + schema = json.loads(SCHEMA_PATH.read_text(encoding="utf-8")) + validator = Draft202012Validator(schema, format_checker=FormatChecker()) + assert list(validator.iter_errors(_build("hk_equity"))) == [] + + +@pytest.mark.parametrize( + ("field", "value"), + [ + ("no_order", False), + ("permission_effect", "live"), + ("catalog_status_used_as_evidence", True), + ("terminal", False), + ], +) +def test_authority_and_terminal_guards_fail_closed(field, value): + artifact = _build() + artifact[field] = value + with pytest.raises(InvalidResearchDriverArtifact): + validate_research_driver_terminal_artifact(artifact) + + +def test_validator_rejects_tampered_terminal_status_and_dependency_chain(): + artifact = _build() + artifact["terminal_status"] = "DEFERRED" + with pytest.raises(InvalidResearchDriverArtifact, match="terminal_status"): + validate_research_driver_terminal_artifact(artifact) + + artifact = _build() + artifact["stages"]["p1_input"] = build_nonready_research_stage( + "p1_input", status="DEFERRED", reason_codes=("input_not_available",) + ) + artifact["terminal_status"] = "DEFERRED" + with pytest.raises(InvalidResearchDriverArtifact, match="bypass"): + validate_research_driver_terminal_artifact(artifact) + + +def test_canonical_bytes_digest_and_copy_isolation_are_deterministic(): + artifact = _build("crypto") + reordered = json.loads(json.dumps(artifact)) + assert canonical_research_driver_terminal_bytes(reordered) == ( + canonical_research_driver_terminal_bytes(artifact) + ) + assert research_driver_terminal_sha256(reordered) == ( + research_driver_terminal_sha256(artifact) + ) + validated = validate_research_driver_terminal_artifact(artifact) + validated["stages"]["p1_input"]["artifact"]["artifact_id"] = "changed" + assert artifact["stages"]["p1_input"]["artifact"]["artifact_id"] == ( + "manifest-001" + ) + + +def test_unknown_domain_and_invalid_generated_at_are_rejected(): + with pytest.raises(InvalidResearchDriverArtifact, match="unsupported"): + _build("forex") + stages = _ready_stages() + with pytest.raises(InvalidResearchDriverArtifact, match="timezone"): + build_research_driver_terminal_artifact( + run_id="run", + generated_at="2026-08-24T22:00:00", + strategy_id="strategy", + candidate_id="candidate", + domain="crypto", + **stages, + )