Files
quant_engine/src/quant_engine/governed_pipeline.py
T
2026-08-30 16:14:34 +08:00

541 lines
19 KiB
Python

"""Versioned, risk-gated, paper-only quantitative research vertical slice.
The module composes existing ``quant_engine`` calculations into a small set of
storage-neutral governance contracts. It never reads a database, calls a data
vendor, loads credentials, or routes an order to a broker.
"""
from __future__ import annotations
import hashlib
import json
import math
import re
from collections.abc import Mapping
from dataclasses import asdict, dataclass
from datetime import UTC, datetime
from enum import StrEnum
from types import MappingProxyType
import pandas as pd
from quant_engine.execution import ExecutionConfig
from quant_engine.research_pipeline import FactorBacktestResult, run_factor_backtest_research
_SHA256 = re.compile(r"^[0-9a-f]{64}$")
_GIT_SHA = re.compile(r"^[0-9a-f]{40}$")
__all__ = [
"DatasetSnapshot",
"FactorVersion",
"StrategyStage",
"StrategyVersion",
"BacktestRun",
"PortfolioTarget",
"RiskPolicy",
"RiskDecisionStatus",
"RiskDecision",
"PaperOrderIntent",
"GovernedFactorSliceResult",
"evaluate_portfolio_risk",
"create_paper_order_intent",
"run_governed_factor_slice",
]
def _required_text(value: str, name: str) -> str:
normalized = value.strip()
if not normalized:
raise ValueError(f"{name} must be non-empty")
return normalized
def _aware_utc(value: datetime, name: str) -> datetime:
if not isinstance(value, datetime) or value.tzinfo is None or value.utcoffset() is None:
raise ValueError(f"{name} must be timezone-aware")
return value.astimezone(UTC)
def _canonical_value(value: object) -> object:
if value is None or isinstance(value, str | bool | int):
return value
if isinstance(value, float):
if math.isnan(value):
return "NaN"
if math.isinf(value):
return "Infinity" if value > 0 else "-Infinity"
return value
if isinstance(value, datetime):
return _aware_utc(value, "datetime").isoformat()
if isinstance(value, StrEnum):
return value.value
if isinstance(value, Mapping):
return {
str(key): _canonical_value(item)
for key, item in sorted(value.items(), key=lambda pair: str(pair[0]))
}
if isinstance(value, list | tuple):
return [_canonical_value(item) for item in value]
raise TypeError(f"unsupported canonical value: {type(value).__name__}")
def _stable_id(prefix: str, payload: Mapping[str, object]) -> str:
encoded = json.dumps(
_canonical_value(payload),
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
allow_nan=False,
).encode("utf-8")
return f"{prefix}:{hashlib.sha256(encoded).hexdigest()}"
def _immutable_weights(values: Mapping[str, float]) -> Mapping[str, float]:
normalized: dict[str, float] = {}
for raw_asset, raw_weight in values.items():
asset = _required_text(str(raw_asset), "asset_id")
weight = float(raw_weight)
if not math.isfinite(weight) or weight < 0:
raise ValueError("portfolio weights must be finite and non-negative")
if asset in normalized:
raise ValueError(f"duplicate asset_id: {asset}")
normalized[asset] = weight
return MappingProxyType(dict(sorted(normalized.items())))
@dataclass(frozen=True, slots=True)
class DatasetSnapshot:
"""Point-in-time identity for caller-supplied research data."""
snapshot_id: str
schema_version: str
content_sha256: str
effective_at: datetime
available_at: datetime
ingested_at: datetime
def __post_init__(self) -> None:
object.__setattr__(self, "snapshot_id", _required_text(self.snapshot_id, "snapshot_id"))
object.__setattr__(
self, "schema_version", _required_text(self.schema_version, "schema_version")
)
digest = self.content_sha256.strip().lower()
if not _SHA256.fullmatch(digest):
raise ValueError("content_sha256 must be a lowercase SHA-256 digest")
object.__setattr__(self, "content_sha256", digest)
effective_at = _aware_utc(self.effective_at, "effective_at")
available_at = _aware_utc(self.available_at, "available_at")
ingested_at = _aware_utc(self.ingested_at, "ingested_at")
if not effective_at <= available_at <= ingested_at:
raise ValueError("timestamps must satisfy effective_at <= available_at <= ingested_at")
object.__setattr__(self, "effective_at", effective_at)
object.__setattr__(self, "available_at", available_at)
object.__setattr__(self, "ingested_at", ingested_at)
@dataclass(frozen=True, slots=True)
class FactorVersion:
"""Versioned factor definition identity without factor implementation duplication."""
factor_id: str
version: str
definition_sha256: str
dataset_schema_version: str
def __post_init__(self) -> None:
object.__setattr__(self, "factor_id", _required_text(self.factor_id, "factor_id"))
object.__setattr__(self, "version", _required_text(self.version, "version"))
object.__setattr__(
self,
"dataset_schema_version",
_required_text(self.dataset_schema_version, "dataset_schema_version"),
)
digest = self.definition_sha256.strip().lower()
if not _SHA256.fullmatch(digest):
raise ValueError("definition_sha256 must be a lowercase SHA-256 digest")
object.__setattr__(self, "definition_sha256", digest)
@property
def version_id(self) -> str:
return f"{self.factor_id}@{self.version}"
class StrategyStage(StrEnum):
DRAFT = "Draft"
RESEARCH = "Research"
VALIDATED = "Validated"
APPROVED = "Approved"
PAPER = "Paper"
LIVE = "Live"
PAUSED = "Paused"
RETIRED = "Retired"
@dataclass(frozen=True, slots=True)
class StrategyVersion:
"""Strategy lineage and lifecycle state used by the governed slice."""
strategy_id: str
version: str
factor_version_id: str
stage: StrategyStage
def __post_init__(self) -> None:
object.__setattr__(self, "strategy_id", _required_text(self.strategy_id, "strategy_id"))
object.__setattr__(self, "version", _required_text(self.version, "version"))
object.__setattr__(
self,
"factor_version_id",
_required_text(self.factor_version_id, "factor_version_id"),
)
if not isinstance(self.stage, StrategyStage):
raise TypeError("stage must be a StrategyStage")
@property
def version_id(self) -> str:
return f"{self.strategy_id}@{self.version}"
@dataclass(frozen=True, slots=True)
class BacktestRun:
run_id: str
dataset_snapshot_id: str
factor_version_id: str
strategy_version_id: str
code_revision: str
config_hash: str
created_at: datetime
@dataclass(frozen=True, slots=True)
class PortfolioTarget:
target_id: str
backtest_run_id: str
dataset_snapshot_id: str
weights: Mapping[str, float]
created_at: datetime
def __post_init__(self) -> None:
object.__setattr__(self, "weights", _immutable_weights(self.weights))
@property
def gross_exposure(self) -> float:
return float(sum(abs(weight) for weight in self.weights.values()))
@dataclass(frozen=True, slots=True)
class RiskPolicy:
policy_id: str
max_gross_exposure: float
max_single_asset_weight: float
max_positions: int
def __post_init__(self) -> None:
object.__setattr__(self, "policy_id", _required_text(self.policy_id, "policy_id"))
if not math.isfinite(self.max_gross_exposure) or self.max_gross_exposure <= 0:
raise ValueError("max_gross_exposure must be positive and finite")
if not math.isfinite(self.max_single_asset_weight) or not (
0 < self.max_single_asset_weight <= 1
):
raise ValueError("max_single_asset_weight must be in (0, 1]")
if (
isinstance(self.max_positions, bool)
or not isinstance(self.max_positions, int)
or self.max_positions <= 0
):
raise ValueError("max_positions must be a positive integer")
class RiskDecisionStatus(StrEnum):
APPROVED = "approved"
REJECTED = "rejected"
@dataclass(frozen=True, slots=True)
class RiskDecision:
decision_id: str
portfolio_target_id: str
policy_id: str
status: RiskDecisionStatus
reasons: tuple[str, ...]
created_at: datetime
@dataclass(frozen=True, slots=True, init=False)
class PaperOrderIntent:
intent_id: str
portfolio_target_id: str
risk_decision_id: str
environment: str
target_weights: Mapping[str, float]
created_at: datetime
def __init__(self, target: PortfolioTarget, decision: RiskDecision) -> None:
if decision.status is not RiskDecisionStatus.APPROVED:
raise ValueError("an approved risk decision is required")
if decision.portfolio_target_id != target.target_id:
raise ValueError("risk decision does not bind the supplied portfolio target")
target_weights = _immutable_weights(target.weights)
payload = {
"portfolio_target_id": target.target_id,
"risk_decision_id": decision.decision_id,
"environment": "paper",
"target_weights": target_weights,
"created_at": decision.created_at,
}
object.__setattr__(self, "intent_id", _stable_id("order-intent", payload))
object.__setattr__(self, "portfolio_target_id", target.target_id)
object.__setattr__(self, "risk_decision_id", decision.decision_id)
object.__setattr__(self, "environment", "paper")
object.__setattr__(self, "target_weights", target_weights)
object.__setattr__(self, "created_at", decision.created_at)
@dataclass(frozen=True, slots=True)
class GovernedFactorSliceResult:
dataset_snapshot: DatasetSnapshot
factor_version: FactorVersion
strategy_version: StrategyVersion
research_result: FactorBacktestResult
backtest_run: BacktestRun
portfolio_target: PortfolioTarget
risk_decision: RiskDecision
order_intent: PaperOrderIntent | None
def evaluate_portfolio_risk(
target: PortfolioTarget,
policy: RiskPolicy,
*,
created_at: datetime | None = None,
) -> RiskDecision:
"""Apply fail-closed paper risk limits to a versioned target portfolio."""
decision_time = target.created_at if created_at is None else _aware_utc(created_at, "created_at")
reasons: list[str] = []
if target.gross_exposure > policy.max_gross_exposure + 1e-12:
reasons.append(
f"gross exposure {target.gross_exposure:.12g} exceeds "
f"limit {policy.max_gross_exposure:.12g}"
)
positive_weights = [weight for weight in target.weights.values() if weight > 0]
largest_weight = max(positive_weights, default=0.0)
if largest_weight > policy.max_single_asset_weight + 1e-12:
reasons.append(
f"single-asset weight {largest_weight:.12g} exceeds "
f"limit {policy.max_single_asset_weight:.12g}"
)
if len(positive_weights) > policy.max_positions:
reasons.append(
f"position count {len(positive_weights)} exceeds limit {policy.max_positions}"
)
status = RiskDecisionStatus.REJECTED if reasons else RiskDecisionStatus.APPROVED
payload = {
"portfolio_target_id": target.target_id,
"policy_id": policy.policy_id,
"status": status,
"reasons": reasons,
"created_at": decision_time,
}
return RiskDecision(
decision_id=_stable_id("risk-decision", payload),
portfolio_target_id=target.target_id,
policy_id=policy.policy_id,
status=status,
reasons=tuple(reasons),
created_at=decision_time,
)
def create_paper_order_intent(
target: PortfolioTarget,
decision: RiskDecision,
) -> PaperOrderIntent:
"""Create a paper-only intent after an exact, approved risk decision."""
return PaperOrderIntent(target, decision)
def run_governed_factor_slice(
*,
factor_scores: pd.DataFrame,
execution_prices: pd.DataFrame,
valuation_prices: pd.DataFrame,
dataset_snapshot: DatasetSnapshot,
factor_version: FactorVersion,
strategy_version: StrategyVersion,
risk_policy: RiskPolicy,
code_revision: str,
created_at: datetime,
top_k: int,
execution_price_field: str,
valuation_price_field: str,
lag_sessions: int = 1,
gross_exposure: float = 1.0,
largest: bool = True,
initial_cash: float = 1_000_000.0,
execution_config: ExecutionConfig | None = None,
) -> GovernedFactorSliceResult:
"""Run snapshot → factor → strategy → backtest → target → risk → paper intent."""
if strategy_version.factor_version_id != factor_version.version_id:
raise ValueError("strategy factor lineage does not match factor_version")
if dataset_snapshot.schema_version != factor_version.dataset_schema_version:
raise ValueError("factor dataset schema does not match dataset snapshot")
if strategy_version.stage not in {StrategyStage.APPROVED, StrategyStage.PAPER}:
raise ValueError("strategy must be in Approved or Paper stage")
normalized_revision = code_revision.strip().lower()
if not _GIT_SHA.fullmatch(normalized_revision):
raise ValueError("code_revision must be a lowercase 40-character Git SHA")
run_time = _aware_utc(created_at, "created_at")
if dataset_snapshot.available_at > run_time or dataset_snapshot.ingested_at > run_time:
raise ValueError("dataset snapshot must be available before the research run")
if not factor_scores.empty:
latest_decision = pd.Timestamp(factor_scores.index.max())
latest_decision_date = latest_decision.date()
if latest_decision_date > run_time.date():
raise ValueError("factor_scores contain future decision dates")
config = ExecutionConfig() if execution_config is None else execution_config
if not isinstance(config, ExecutionConfig):
raise TypeError("execution_config must be an ExecutionConfig")
parameters: dict[str, object] = {
"top_k": top_k,
"execution_price_field": execution_price_field,
"valuation_price_field": valuation_price_field,
"lag_sessions": lag_sessions,
"gross_exposure": gross_exposure,
"largest": largest,
"initial_cash": initial_cash,
"execution_config": asdict(config),
}
config_hash = _stable_id("config", parameters).split(":", maxsplit=1)[1]
run_payload = {
"dataset_snapshot_id": dataset_snapshot.snapshot_id,
"dataset_content_sha256": dataset_snapshot.content_sha256,
"factor_version_id": factor_version.version_id,
"factor_definition_sha256": factor_version.definition_sha256,
"strategy_version_id": strategy_version.version_id,
"code_revision": normalized_revision,
"config_hash": config_hash,
"created_at": run_time,
}
backtest_run = BacktestRun(
run_id=_stable_id("backtest-run", run_payload),
dataset_snapshot_id=dataset_snapshot.snapshot_id,
factor_version_id=factor_version.version_id,
strategy_version_id=strategy_version.version_id,
code_revision=normalized_revision,
config_hash=config_hash,
created_at=run_time,
)
research_result = run_factor_backtest_research(
factor_scores,
execution_prices,
valuation_prices,
top_k=top_k,
execution_price_field=execution_price_field,
valuation_price_field=valuation_price_field,
lag_sessions=lag_sessions,
gross_exposure=gross_exposure,
largest=largest,
initial_cash=initial_cash,
config=config,
)
if research_result.schedule.decision_weights.empty:
raise ValueError("governed slice requires at least one target portfolio")
final_weights = {
str(asset): float(weight)
for asset, weight in research_result.schedule.decision_weights.iloc[-1].items()
}
target_payload = {
"backtest_run_id": backtest_run.run_id,
"dataset_snapshot_id": dataset_snapshot.snapshot_id,
"weights": final_weights,
"created_at": run_time,
}
portfolio_target = PortfolioTarget(
target_id=_stable_id("portfolio-target", target_payload),
backtest_run_id=backtest_run.run_id,
dataset_snapshot_id=dataset_snapshot.snapshot_id,
weights=final_weights,
created_at=run_time,
)
risk_decision = evaluate_portfolio_risk(
portfolio_target,
risk_policy,
created_at=run_time,
)
order_intent = (
create_paper_order_intent(portfolio_target, risk_decision)
if risk_decision.status is RiskDecisionStatus.APPROVED
else None
)
return GovernedFactorSliceResult(
dataset_snapshot=dataset_snapshot,
factor_version=factor_version,
strategy_version=strategy_version,
research_result=research_result,
backtest_run=backtest_run,
portfolio_target=portfolio_target,
risk_decision=risk_decision,
order_intent=order_intent,
)
def _demo() -> Mapping[str, object]:
dates = pd.date_range("2026-01-05", periods=4, freq="B")
scores = pd.DataFrame({"A": [2.0, 1.0], "B": [1.0, 2.0]}, index=dates[:2])
opens = pd.DataFrame({"A": [10.0, 10.0, 10.2, 10.3], "B": [20.0, 20.0, 20.5, 21.0]}, index=dates)
result = run_governed_factor_slice(
factor_scores=scores,
execution_prices=opens,
valuation_prices=opens * 1.01,
dataset_snapshot=DatasetSnapshot(
snapshot_id="dataset:architecture-smoke-v1",
schema_version="1.0.0",
content_sha256="a" * 64,
effective_at=datetime(2026, 1, 8, 7, tzinfo=UTC),
available_at=datetime(2026, 1, 8, 8, tzinfo=UTC),
ingested_at=datetime(2026, 1, 8, 8, 5, tzinfo=UTC),
),
factor_version=FactorVersion(
factor_id="factor:architecture-smoke",
version="1.0.0",
definition_sha256="b" * 64,
dataset_schema_version="1.0.0",
),
strategy_version=StrategyVersion(
strategy_id="strategy:architecture-smoke",
version="1.0.0",
factor_version_id="factor:architecture-smoke@1.0.0",
stage=StrategyStage.APPROVED,
),
risk_policy=RiskPolicy(
policy_id="risk:architecture-smoke@1.0.0",
max_gross_exposure=1.0,
max_single_asset_weight=0.6,
max_positions=10,
),
code_revision="c" * 40,
created_at=datetime(2026, 1, 9, 1, tzinfo=UTC),
top_k=2,
execution_price_field="open",
valuation_price_field="close",
execution_config=ExecutionConfig(
commission_bps=0,
stamp_tax_bps=0,
slippage_bps=0,
min_trade_amount=0,
),
)
return {
"backtest_run_id": result.backtest_run.run_id,
"portfolio_target_id": result.portfolio_target.target_id,
"risk_decision": result.risk_decision.status.value,
"order_intent_id": result.order_intent.intent_id if result.order_intent else None,
"environment": result.order_intent.environment if result.order_intent else None,
}
if __name__ == "__main__":
print(json.dumps(_demo(), ensure_ascii=False, sort_keys=True))