Artifex/control_plane/trading_studio/indicators/engine.py
2026-08-18 02:00:53 +07:00

131 lines
5.1 KiB
Python

"""Formula-free execution boundary. It refuses rather than silently falling back."""
from __future__ import annotations
from dataclasses import dataclass
from enum import StrEnum
from typing import Any
import numpy as np
from .definitions import IndicatorVariant
from .feature_catalog import FeatureUnavailableError, require_validated_features
from .registry import HISTORICAL_REGISTRY, HS22_REQUIRED_V1, IndicatorRegistry
from .schema import HS22State, parse_hs22
class RefusalCode(StrEnum):
FORMULA_UNAVAILABLE = "formula_unavailable"
UNKNOWN_INDICATOR = "unknown_indicator"
UNSUPPORTED_VARIANT = "unsupported_variant"
UNREGISTERED = "unregistered"
@dataclass(frozen=True, slots=True)
class RefusalState:
code: RefusalCode
message: str
all_nan_behavior: bool = False
class IndicatorEngine:
def __init__(self, registry: IndicatorRegistry = HISTORICAL_REGISTRY) -> None:
self.registry = registry
def resolve(self, variant: IndicatorVariant) -> RefusalState:
definition = self.registry.definition(variant.indicator_id)
if definition is None:
return RefusalState(
RefusalCode.UNKNOWN_INDICATOR, "unknown ID returns all-NaN; no fallback", True
)
if definition.status.value == "unregistered":
return RefusalState(
RefusalCode.UNREGISTERED, "indicator is not registered for execution", True
)
if self.registry.variant(variant.indicator_id, variant.period, variant.p1) is None:
return RefusalState(
RefusalCode.UNSUPPORTED_VARIANT, "variant is outside the registry", False
)
return RefusalState(
RefusalCode.FORMULA_UNAVAILABLE, "indicator formulas are not implemented", False
)
def compute(self, variant: IndicatorVariant, length: int) -> RefusalState:
# This API intentionally does not emit a numeric fallback array.
if length < 0:
raise ValueError("length must be non-negative")
return self.resolve(variant)
def _arrays(close: Any, high: Any, low: Any, volume: Any) -> tuple[np.ndarray, ...]:
arrays = tuple(np.asarray(item, dtype=np.float64) for item in (close, high, low, volume))
if not arrays[0].ndim == 1 or any(item.shape != arrays[0].shape for item in arrays[1:]):
raise ValueError("HS22 OHLCV inputs must be equally sized one-dimensional arrays")
return arrays
def evaluate_indicator(
variant: IndicatorVariant, close: Any, high: Any, low: Any, volume: Any
) -> np.ndarray:
"""Evaluate exactly one frozen Cohort001 role triple; all others fail closed."""
try:
require_validated_features([variant])
except FeatureUnavailableError as error:
raise ValueError(str(error)) from error
from .historical_formulae import compute_indicator
close, high, low, volume = _arrays(close, high, low, volume)
return compute_indicator(
variant.indicator_id, close, high, low, volume, variant.period, variant.p1
)
def signal_color(values: Any) -> np.ndarray:
values = np.asarray(values, dtype=np.float64)
colors = np.zeros(len(values), dtype=np.int8)
previous = 0
for index in range(1, len(values)):
if np.isnan(values[index]) or np.isnan(values[index - 1]):
previous = 0
elif values[index] > values[index - 1]:
previous = 1
elif values[index] < values[index - 1]:
previous = -1
colors[index] = previous
return colors
def compute_hs22_state(
state: HS22State | object, close: Any, high: Any, low: Any, volume: Any
) -> dict[str, np.ndarray]:
"""Port of historical ``paper_replay.compute_combo_state`` without a reference import."""
if not isinstance(state, HS22State):
state = parse_hs22(state, HS22_REQUIRED_V1)
close, high, low, volume = _arrays(close, high, low, volume)
trend = evaluate_indicator(state.trend, close, high, low, volume)
signal = evaluate_indicator(state.signal, close, high, low, volume)
trigger = evaluate_indicator(state.trigger, close, high, low, volume)
confirm = evaluate_indicator(state.confirm, close, high, low, volume)
vol = evaluate_indicator(state.volatility, close, high, low, volume)
return {"trend": trend, "signal": signal, "trigger": trigger, "confirm": confirm, "vol": vol,
"trend_color": signal_color(trend), "signal_color": signal_color(signal)}
def analyze_corpus_coverage(
path: str, registry: IndicatorRegistry = HISTORICAL_REGISTRY
) -> dict[str, object]:
"""Inspect a parquet corpus when an optional parquet reader is installed."""
try:
import pyarrow.parquet as parquet # type: ignore[import-not-found]
except ImportError:
return {"available": False, "reason": "pyarrow is not installed"}
table = parquet.read_table(path)
fields = set(table.column_names)
required = {"open", "high", "low", "close", "volume"}
return {
"available": True,
"rows": table.num_rows,
"fields": sorted(fields),
"required_fields_present": required <= fields,
"registered_variants": len(registry.variants),
}