"""Causal anchor-time historical-reference features for forecasting.""" from __future__ import annotations from collections import deque from dataclasses import dataclass import numpy as np REFERENCE_WINDOWS = (1, 5, 20) WINDOW_FEATURES = ( "distance_to_peak_bps", "distance_to_bottom_bps", "range_width_bps", "range_position", "available", ) REFERENCE_FEATURE_NAMES = tuple( f"prior_{window}_session_{feature}" for window in REFERENCE_WINDOWS for feature in WINDOW_FEATURES ) SESSION_FEATURE_NAMES = ( "distance_to_session_open_bps", "distance_to_session_max_close_so_far_bps", "distance_to_session_min_close_so_far_bps", "session_range_position_so_far", "session_history_available", ) GENERIC_WINDOW_FEATURES = ( "net_return_bps", "rms_observed_return_bps", "range_width_bps", "terminal_drawdown_bps", "available", ) GENERIC_FEATURE_NAMES = tuple( f"prior_{window}_session_{feature}" for window in REFERENCE_WINDOWS for feature in GENERIC_WINDOW_FEATURES ) VISIBLE_FEATURE_NAMES = ( "visible_128_distance_to_peak_bps", "visible_128_distance_to_bottom_bps", "visible_128_range_width_bps", "visible_128_range_position", "visible_128_history_available", ) REFERENCE_TREATMENTS = ( "P0", "Pmask", "P1", "Pmulti", "Pgeneric", "P128", "Pstale", "Psession", ) @dataclass(frozen=True) class ReferenceFeatureBundle: reference_values: np.ndarray generic_values: np.ndarray stale_reference_values: np.ndarray visible_values: np.ndarray session_values: np.ndarray relative_log_close: np.ndarray def __post_init__(self) -> None: row_count = len(self.relative_log_close) if self.reference_values.shape != ( row_count, len(REFERENCE_FEATURE_NAMES), ): raise ValueError("reference feature shape is invalid") if self.generic_values.shape != ( row_count, len(GENERIC_FEATURE_NAMES), ): raise ValueError("generic feature shape is invalid") if self.stale_reference_values.shape != ( row_count, len(REFERENCE_FEATURE_NAMES), ): raise ValueError("stale reference feature shape is invalid") if self.visible_values.shape != ( row_count, len(VISIBLE_FEATURE_NAMES), ): raise ValueError("visible feature shape is invalid") if self.session_values.shape != ( row_count, len(SESSION_FEATURE_NAMES), ): raise ValueError("session feature shape is invalid") if not ( np.isfinite(self.reference_values).all() and np.isfinite(self.generic_values).all() and np.isfinite(self.stale_reference_values).all() and np.isfinite(self.visible_values).all() and np.isfinite(self.session_values).all() and np.isfinite(self.relative_log_close).all() ): raise ValueError("reference feature bundle contains nonfinite values") def _relative_log_close(features: dict[str, np.ndarray]) -> np.ndarray: values = np.asarray(features["X"]) if values.ndim != 2 or values.shape[1] < 2: raise ValueError("token feature matrix is invalid") returns = values[:, 0].astype(np.float64, copy=False) observed = values[:, 1] > 0.5 if not np.isfinite(returns).all(): raise ValueError("close-return stream contains nonfinite values") if np.any((~observed) & (returns != 0.0)): raise ValueError("unobserved rows contain nonzero close returns") return np.cumsum(returns, dtype=np.float64) def _session_runs(session_date: np.ndarray) -> tuple[np.ndarray, np.ndarray]: dates = np.asarray(session_date) if dates.ndim != 1 or not len(dates): raise ValueError("session dates must be a nonempty vector") changes = np.flatnonzero(dates[1:] != dates[:-1]) + 1 starts = np.concatenate(([0], changes)).astype(np.int64) stops = np.concatenate((changes, [len(dates)])).astype(np.int64) labels = dates[starts] if np.any(labels[1:] <= labels[:-1]): raise ValueError("session dates are not strictly increasing") return starts, stops def _geometry( anchor: np.ndarray, peak: float, bottom: float, ) -> np.ndarray: width = float(peak - bottom) denominator = max(width, np.finfo(np.float64).eps) return np.column_stack( ( 10_000.0 * (anchor - peak), 10_000.0 * (anchor - bottom), np.full(len(anchor), 10_000.0 * width), (anchor - bottom) / denominator, ) ) def _generic_summary(values: np.ndarray) -> np.ndarray: observed_returns = np.diff(values) rms_return = ( float(np.sqrt(np.mean(np.square(observed_returns)))) if len(observed_returns) else 0.0 ) peak = float(np.max(values)) bottom = float(np.min(values)) return np.array( ( 10_000.0 * float(values[-1] - values[0]), 10_000.0 * rms_return, 10_000.0 * (peak - bottom), 10_000.0 * float(values[-1] - peak), ), dtype=np.float64, ) def _visible_geometry( relative_close: np.ndarray, observed: np.ndarray, *, context_length: int = 128, ) -> np.ndarray: result = np.zeros( (len(relative_close), len(VISIBLE_FEATURE_NAMES)), dtype=np.float64, ) maximum: deque[int] = deque() minimum: deque[int] = deque() for index, value in enumerate(relative_close): first = index + 1 - context_length while maximum and maximum[0] < first: maximum.popleft() while minimum and minimum[0] < first: minimum.popleft() if not observed[index]: continue while maximum and relative_close[maximum[-1]] <= value: maximum.pop() while minimum and relative_close[minimum[-1]] >= value: minimum.pop() maximum.append(index) minimum.append(index) peak = float(relative_close[maximum[0]]) bottom = float(relative_close[minimum[0]]) result[index, :4] = _geometry( np.asarray([value]), peak, bottom, )[0] result[index, 4] = 1.0 return result def materialize_reference_features( features: dict[str, np.ndarray], ) -> ReferenceFeatureBundle: """Build close-only features available at each completed minute anchor.""" relative_close = _relative_log_close(features) dates = np.asarray(features["session_date"]) minute = np.asarray(features["minute_of_session"]) if len(dates) != len(relative_close) or len(minute) != len(relative_close): raise ValueError("feature row metadata has inconsistent lengths") observed = np.asarray(features["X"])[:, 1] > 0.5 starts, stops = _session_runs(dates) session_peak = np.full(len(starts), np.nan, dtype=np.float64) session_bottom = np.full(len(starts), np.nan, dtype=np.float64) session_available = np.zeros(len(starts), dtype=np.bool_) for index, (start, stop) in enumerate(zip(starts, stops, strict=True)): local = observed[start:stop] if np.any(local): values = relative_close[start:stop][local] session_peak[index] = float(np.max(values)) session_bottom[index] = float(np.min(values)) session_available[index] = True reference = np.zeros( (len(relative_close), len(REFERENCE_FEATURE_NAMES)), dtype=np.float64, ) generic = np.zeros( (len(relative_close), len(GENERIC_FEATURE_NAMES)), dtype=np.float64, ) stale_reference = np.zeros_like(reference) visible = _visible_geometry(relative_close, observed) session = np.zeros( (len(relative_close), len(SESSION_FEATURE_NAMES)), dtype=np.float64, ) for session_index, (start, stop) in enumerate( zip(starts, stops, strict=True) ): local_observed = observed[start:stop] if np.any(local_observed): rows = np.arange(start, stop)[local_observed] values = relative_close[rows] running_peak = np.maximum.accumulate(values) running_bottom = np.minimum.accumulate(values) opening = values[0] width = running_peak - running_bottom session[rows, :4] = np.column_stack( ( 10_000.0 * (values - opening), 10_000.0 * (values - running_peak), 10_000.0 * (values - running_bottom), (values - running_bottom) / np.maximum(width, np.finfo(np.float64).eps), ) ) session[rows, 4] = 1.0 for window_index, window in enumerate(REFERENCE_WINDOWS): first = session_index - window if first < 0: continue window_available = session_available[first:session_index] if not np.all(window_available) or not np.any(local_observed): continue peak = float(np.max(session_peak[first:session_index])) bottom = float(np.min(session_bottom[first:session_index])) rows = np.arange(start, stop)[local_observed] offset = window_index * len(WINDOW_FEATURES) reference[rows, offset : offset + 4] = _geometry( relative_close[rows], peak, bottom, ) reference[rows, offset + 4] = 1.0 source_start = starts[first] source_stop = stops[session_index - 1] source_observed = observed[source_start:source_stop] source_values = relative_close[source_start:source_stop][ source_observed ] generic[rows, offset : offset + 4] = _generic_summary( source_values ) generic[rows, offset + 4] = 1.0 stale_first = first - 1 stale_stop = session_index - 1 if stale_first < 0: continue stale_available = session_available[stale_first:stale_stop] if not np.all(stale_available): continue stale_peak = float( np.max(session_peak[stale_first:stale_stop]) ) stale_bottom = float( np.min(session_bottom[stale_first:stale_stop]) ) stale_reference[rows, offset : offset + 4] = _geometry( relative_close[rows], stale_peak, stale_bottom, ) stale_reference[rows, offset + 4] = 1.0 return ReferenceFeatureBundle( reference_values=reference.astype(np.float32), generic_values=generic.astype(np.float32), stale_reference_values=stale_reference.astype(np.float32), visible_values=visible.astype(np.float32), session_values=session.astype(np.float32), relative_log_close=relative_close, ) def treatment_values( bundle: ReferenceFeatureBundle, treatment: str, ) -> tuple[np.ndarray, tuple[str, ...]]: if treatment not in REFERENCE_TREATMENTS: raise ValueError(f"unknown reference treatment: {treatment}") if treatment == "Psession": return bundle.session_values.copy(), SESSION_FEATURE_NAMES if treatment == "P128": return bundle.visible_values.copy(), VISIBLE_FEATURE_NAMES if treatment == "Pgeneric": return bundle.generic_values.copy(), GENERIC_FEATURE_NAMES if treatment == "Pstale": return ( bundle.stale_reference_values.copy(), REFERENCE_FEATURE_NAMES, ) if treatment == "P1": width = len(WINDOW_FEATURES) return ( bundle.reference_values[:, :width].copy(), REFERENCE_FEATURE_NAMES[:width], ) values = np.zeros_like(bundle.reference_values) if treatment == "Pmask": mask_indices = np.arange( len(WINDOW_FEATURES) - 1, len(REFERENCE_FEATURE_NAMES), len(WINDOW_FEATURES), ) values[:, mask_indices] = bundle.reference_values[:, mask_indices] elif treatment == "Pmulti": values[:] = bundle.reference_values return values, REFERENCE_FEATURE_NAMES __all__ = [ "GENERIC_FEATURE_NAMES", "REFERENCE_FEATURE_NAMES", "REFERENCE_TREATMENTS", "REFERENCE_WINDOWS", "SESSION_FEATURE_NAMES", "VISIBLE_FEATURE_NAMES", "ReferenceFeatureBundle", "materialize_reference_features", "treatment_values", ]