feat(0071): stage1 breakout candidate — Donchian channel signal (#137)
A second Stage-1 R strategy candidate for the edge-research milestone, a structurally
different trend mechanic than the refuted SMA-momentum: a Donchian channel breakout.
Swaps only the signal leg of the stage1-r graph —
close -> Delay(1) -> {RollingMax,RollingMin}(N) -> {Gt,Gt} -> {Latch,Latch} ->
Sub = bias in {-1,0,+1} — keeping the vol-stop R definition unchanged, so it screens
under the identical R yardstick (clean A/B vs momentum).
- aura-std: new RollingMax / RollingMin nodes (sliding-window extremum via a monotonic
deque, O(1) amortized; mirror the Sma/Delay ring-buffer shape, warm-up skip-emit,
C7/C8). The +/-1 direction latch composes from two existing Latch nodes + Sub (no new
latch type).
- aura-cli: stage1_breakout_graph + a manual-grid stage1_breakout_sweep_family
(fully-bound graph per point, sidestepping parameter-ganging since one channel length
drives both rolling nodes); --strategy stage1-breakout + --channel grid flag; usage
strings updated.
- Causality (C2): one Delay(1) on close feeds both rolling nodes, so each channel covers
close[t-N..t-1] — the current bar is excluded. Pinned by a contrastive e2e test (a
strictly rising series up-breaks every warmed bar; a window including the current bar
never could) plus a +/-1-latch-hold test.
All existing goldens byte-identical (stage1-r/sma/momentum untouched); folded-no-trace
== raw-trace metrics; workspace tests + clippy green.
refs #137
This commit is contained in:
@@ -30,6 +30,8 @@ mod mul;
|
||||
mod position_management;
|
||||
mod recorder;
|
||||
mod resample;
|
||||
mod rolling_max;
|
||||
mod rolling_min;
|
||||
mod series_reducer;
|
||||
mod session;
|
||||
mod sim_broker;
|
||||
@@ -56,6 +58,8 @@ pub use position_management::{
|
||||
};
|
||||
pub use recorder::Recorder;
|
||||
pub use resample::Resample;
|
||||
pub use rolling_max::RollingMax;
|
||||
pub use rolling_min::RollingMin;
|
||||
pub use series_reducer::SeriesReducer;
|
||||
pub use session::Session;
|
||||
pub use sim_broker::SimBroker;
|
||||
|
||||
@@ -0,0 +1,222 @@
|
||||
//! `RollingMax` — maximum over the last `length` values of one f64 input.
|
||||
//!
|
||||
//! Sliding-window maximum maintained by a **descending monotonic deque** (the
|
||||
//! textbook O(1)-amortized sliding-window-max): each cycle pops from the back every
|
||||
//! element <= the new sample (they can never again be the window max), pushes the new
|
||||
//! one, and evicts the front once its stream index leaves the window. The front is
|
||||
//! always the current window max. A per-cycle O(N) re-scan would be too slow at the
|
||||
//! large `N` (channel lengths in the thousands of bars) the breakout screen sweeps.
|
||||
//!
|
||||
//! Like `Sma`/`Delay`, the window lives in node state, so `eval` reads only the newest
|
||||
//! sample and `lookbacks()` is `1`. Warm-up is **skip-emit**: `None` until `length`
|
||||
//! samples have passed (the window is not yet full). The deque is pre-sized to
|
||||
//! `length` so the hot path is allocation-free (C7). Per the "operator is topology"
|
||||
//! convention (see `gt.rs`), `RollingMax` and `RollingMin` are two node types, not one
|
||||
//! node with a max/min param.
|
||||
|
||||
use aura_core::{
|
||||
Cell, Ctx, FieldSpec, Firing, Node, NodeSchema, ParamSpec, PortSpec, PrimitiveBuilder,
|
||||
ScalarKind,
|
||||
};
|
||||
use std::collections::VecDeque;
|
||||
|
||||
/// Maximum over the last `length` values of one f64 input, maintained by a descending
|
||||
/// monotonic deque (front = current window max). `None` during warm-up.
|
||||
pub struct RollingMax {
|
||||
length: usize,
|
||||
// descending monotonic deque of (value, stream_index); front is the window max.
|
||||
// node-owned, bounded by `length` (C7), pre-sized so push_back does not allocate.
|
||||
deque: VecDeque<(f64, u64)>,
|
||||
seen: u64, // stream position of the next sample (drives front eviction)
|
||||
count: usize, // samples seen so far, capped at `length` — the warm-up gate
|
||||
out: [Cell; 1],
|
||||
}
|
||||
|
||||
impl RollingMax {
|
||||
/// Build a rolling max of window `length` (must be >= 1; mirror `Sma::new`).
|
||||
pub fn new(length: usize) -> Self {
|
||||
assert!(length >= 1, "RollingMax length must be >= 1");
|
||||
Self {
|
||||
length,
|
||||
deque: VecDeque::with_capacity(length),
|
||||
seen: 0,
|
||||
count: 0,
|
||||
out: [Cell::from_f64(0.0)],
|
||||
}
|
||||
}
|
||||
|
||||
/// The param-generic recipe for a blueprint primitive: declares `length` and builds
|
||||
/// through `RollingMax::new`.
|
||||
pub fn builder() -> PrimitiveBuilder {
|
||||
PrimitiveBuilder::new(
|
||||
"RollingMax",
|
||||
NodeSchema {
|
||||
inputs: vec![PortSpec { kind: ScalarKind::F64, firing: Firing::Any, name: "series".into() }],
|
||||
output: vec![FieldSpec { name: "value".into(), kind: ScalarKind::F64 }],
|
||||
params: vec![ParamSpec { name: "length".into(), kind: ScalarKind::I64 }],
|
||||
},
|
||||
|p| Box::new(RollingMax::new(p[0].i64() as usize)),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
impl Node for RollingMax {
|
||||
// The window lives in node state, so only the newest sample is read each cycle.
|
||||
fn lookbacks(&self) -> Vec<usize> {
|
||||
vec![1]
|
||||
}
|
||||
|
||||
fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Cell]> {
|
||||
let w = ctx.f64_in(0);
|
||||
if w.is_empty() {
|
||||
return None; // no sample yet
|
||||
}
|
||||
let x = w[0]; // index 0 = newest (financial indexing)
|
||||
let idx = self.seen;
|
||||
self.seen += 1;
|
||||
// descending invariant: a back element <= x can never again be the max — drop it.
|
||||
while let Some(&(v, _)) = self.deque.back() {
|
||||
if v <= x {
|
||||
self.deque.pop_back();
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
self.deque.push_back((x, idx));
|
||||
// evict the front once it leaves the window [idx-length+1, idx] (no subtraction:
|
||||
// front index `i` is out when `i + length <= idx`).
|
||||
while let Some(&(_, i)) = self.deque.front() {
|
||||
if i + self.length as u64 <= idx {
|
||||
self.deque.pop_front();
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
if self.count < self.length {
|
||||
self.count += 1;
|
||||
}
|
||||
if self.count < self.length {
|
||||
return None; // not yet warmed up
|
||||
}
|
||||
self.out[0] = Cell::from_f64(self.deque.front().expect("deque non-empty after push").0);
|
||||
Some(&self.out)
|
||||
}
|
||||
|
||||
fn label(&self) -> String {
|
||||
format!("RollingMax({})", self.length)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use aura_core::{AnyColumn, Scalar, Timestamp};
|
||||
|
||||
fn drive(node: &mut RollingMax, feed: &[f64]) -> Vec<Option<f64>> {
|
||||
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, 1)];
|
||||
let mut out = Vec::new();
|
||||
for &v in feed {
|
||||
inputs[0].push(Scalar::f64(v)).unwrap();
|
||||
out.push(node.eval(Ctx::new(&inputs, Timestamp(0))).map(|r| r[0].f64()));
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rolling_max_warms_up_then_tracks_the_window_max() {
|
||||
// max of [1,3,2], [3,2,5], [2,5,4] once warmed up; silent for the first two.
|
||||
let mut node = RollingMax::new(3);
|
||||
let got = drive(&mut node, &[1.0, 3.0, 2.0, 5.0, 4.0]);
|
||||
assert_eq!(got, vec![None, None, Some(3.0), Some(5.0), Some(5.0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rolling_max_evicts_the_expiring_window_maximum() {
|
||||
// the front-eviction path: the early big value must leave the window. length 2.
|
||||
// windows: [9,1]->9, [1,2]->2, [2,3]->3 (the 9 has expired, not stuck as max).
|
||||
let mut node = RollingMax::new(2);
|
||||
let got = drive(&mut node, &[9.0, 1.0, 2.0, 3.0]);
|
||||
assert_eq!(got, vec![None, Some(9.0), Some(2.0), Some(3.0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rolling_max_length_one_is_identity() {
|
||||
let mut node = RollingMax::new(1);
|
||||
assert_eq!(drive(&mut node, &[7.0, 9.0, 4.0]), vec![Some(7.0), Some(9.0), Some(4.0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deque_matches_naive_window_max() {
|
||||
// the monotonic deque must equal a naive O(N) per-window max at every warmed
|
||||
// cycle, across a long deterministic series (mirrors Sma's re-sum cross-check).
|
||||
let length = 37;
|
||||
let mut node = RollingMax::new(length);
|
||||
let series: Vec<f64> = (0..3_000u64)
|
||||
.map(|i| (i.wrapping_mul(1103515245).wrapping_add(12345) % 1000) as f64 * 0.5 - 250.0)
|
||||
.collect();
|
||||
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, 1)];
|
||||
for (t, &x) in series.iter().enumerate() {
|
||||
inputs[0].push(Scalar::f64(x)).unwrap();
|
||||
let got = node.eval(Ctx::new(&inputs, Timestamp(0))).map(|r| r[0].f64());
|
||||
if t + 1 >= length {
|
||||
let want = series[t + 1 - length..=t].iter().cloned().fold(f64::NEG_INFINITY, f64::max);
|
||||
assert_eq!(got, Some(want), "at t={t}");
|
||||
} else {
|
||||
assert_eq!(got, None, "silent until warmed up at t={t}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lookback_is_one_window_lives_in_node_state() {
|
||||
assert_eq!(RollingMax::new(20).lookbacks(), vec![1]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn length_must_be_at_least_one() {
|
||||
let r = std::panic::catch_unwind(|| RollingMax::new(0));
|
||||
assert!(r.is_err(), "RollingMax::new(0) must panic");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn label_carries_the_window() {
|
||||
assert_eq!(RollingMax::new(20).label(), "RollingMax(20)");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_declares_length_param_and_value_output() {
|
||||
// the blueprint seam later tasks bootstrap through: the param-generic recipe
|
||||
// declares a single I64 `length` knob and an F64 `value` output (mirrors Sma's
|
||||
// nodes_declare_expected_params).
|
||||
let schema = RollingMax::builder().schema().clone();
|
||||
assert_eq!(schema.params, vec![ParamSpec { name: "length".into(), kind: ScalarKind::I64 }]);
|
||||
assert_eq!(schema.output.len(), 1);
|
||||
assert_eq!(schema.output[0].name, "value");
|
||||
assert_eq!(schema.output[0].kind, ScalarKind::F64);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_input_slot_is_named_series() {
|
||||
// the named f64 input the breakout graph wires close into (mirrors Sma's
|
||||
// input_slot_is_named_series).
|
||||
assert_eq!(RollingMax::builder().schema().inputs[0].name, "series");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_bind_removes_length_from_param_space() {
|
||||
// a bound length reports an empty param surface; the open form keeps it (mirrors
|
||||
// Sma's bind_removes_slot_from_param_space).
|
||||
let bound = RollingMax::builder().named("channel").bind("length", Scalar::i64(20));
|
||||
assert!(bound.schema().params.is_empty());
|
||||
assert_eq!(RollingMax::builder().named("channel").params().len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_bound_node_builds_with_injected_length() {
|
||||
// built from an empty open slice, the bound builder yields a RollingMax(20) —
|
||||
// the `p[0].i64() as usize` build closure is exercised end-to-end (mirrors Sma's
|
||||
// bound_node_builds_with_injected_value).
|
||||
let node = RollingMax::builder().bind("length", Scalar::i64(20)).build(&[]);
|
||||
assert_eq!(node.label(), "RollingMax(20)");
|
||||
}
|
||||
}
|
||||
@@ -0,0 +1,208 @@
|
||||
//! `RollingMin` — minimum over the last `length` values of one f64 input.
|
||||
//!
|
||||
//! The mirror of `RollingMax`: an **ascending** monotonic deque (front = window min).
|
||||
//! Each cycle pops from the back every element >= the new sample, pushes the new one,
|
||||
//! and evicts the front once it leaves the window. See `rolling_max.rs` for the shared
|
||||
//! rationale (O(1) amortized, window in node state, warm-up skip-emit, C7).
|
||||
|
||||
use aura_core::{
|
||||
Cell, Ctx, FieldSpec, Firing, Node, NodeSchema, ParamSpec, PortSpec, PrimitiveBuilder,
|
||||
ScalarKind,
|
||||
};
|
||||
use std::collections::VecDeque;
|
||||
|
||||
/// Minimum over the last `length` values of one f64 input, maintained by an ascending
|
||||
/// monotonic deque (front = current window min). `None` during warm-up.
|
||||
pub struct RollingMin {
|
||||
length: usize,
|
||||
deque: VecDeque<(f64, u64)>, // ascending; front = window min
|
||||
seen: u64,
|
||||
count: usize,
|
||||
out: [Cell; 1],
|
||||
}
|
||||
|
||||
impl RollingMin {
|
||||
/// Build a rolling min of window `length` (must be >= 1).
|
||||
pub fn new(length: usize) -> Self {
|
||||
assert!(length >= 1, "RollingMin length must be >= 1");
|
||||
Self {
|
||||
length,
|
||||
deque: VecDeque::with_capacity(length),
|
||||
seen: 0,
|
||||
count: 0,
|
||||
out: [Cell::from_f64(0.0)],
|
||||
}
|
||||
}
|
||||
|
||||
/// The param-generic recipe for a blueprint primitive: declares `length` and builds
|
||||
/// through `RollingMin::new`.
|
||||
pub fn builder() -> PrimitiveBuilder {
|
||||
PrimitiveBuilder::new(
|
||||
"RollingMin",
|
||||
NodeSchema {
|
||||
inputs: vec![PortSpec { kind: ScalarKind::F64, firing: Firing::Any, name: "series".into() }],
|
||||
output: vec![FieldSpec { name: "value".into(), kind: ScalarKind::F64 }],
|
||||
params: vec![ParamSpec { name: "length".into(), kind: ScalarKind::I64 }],
|
||||
},
|
||||
|p| Box::new(RollingMin::new(p[0].i64() as usize)),
|
||||
)
|
||||
}
|
||||
}
|
||||
|
||||
impl Node for RollingMin {
|
||||
fn lookbacks(&self) -> Vec<usize> {
|
||||
vec![1]
|
||||
}
|
||||
|
||||
fn eval(&mut self, ctx: Ctx<'_>) -> Option<&[Cell]> {
|
||||
let w = ctx.f64_in(0);
|
||||
if w.is_empty() {
|
||||
return None;
|
||||
}
|
||||
let x = w[0];
|
||||
let idx = self.seen;
|
||||
self.seen += 1;
|
||||
// ascending invariant: a back element >= x can never again be the min — drop it.
|
||||
while let Some(&(v, _)) = self.deque.back() {
|
||||
if v >= x {
|
||||
self.deque.pop_back();
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
self.deque.push_back((x, idx));
|
||||
while let Some(&(_, i)) = self.deque.front() {
|
||||
if i + self.length as u64 <= idx {
|
||||
self.deque.pop_front();
|
||||
} else {
|
||||
break;
|
||||
}
|
||||
}
|
||||
if self.count < self.length {
|
||||
self.count += 1;
|
||||
}
|
||||
if self.count < self.length {
|
||||
return None;
|
||||
}
|
||||
self.out[0] = Cell::from_f64(self.deque.front().expect("deque non-empty after push").0);
|
||||
Some(&self.out)
|
||||
}
|
||||
|
||||
fn label(&self) -> String {
|
||||
format!("RollingMin({})", self.length)
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use aura_core::{AnyColumn, Scalar, Timestamp};
|
||||
|
||||
fn drive(node: &mut RollingMin, feed: &[f64]) -> Vec<Option<f64>> {
|
||||
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, 1)];
|
||||
let mut out = Vec::new();
|
||||
for &v in feed {
|
||||
inputs[0].push(Scalar::f64(v)).unwrap();
|
||||
out.push(node.eval(Ctx::new(&inputs, Timestamp(0))).map(|r| r[0].f64()));
|
||||
}
|
||||
out
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rolling_min_warms_up_then_tracks_the_window_min() {
|
||||
// min of [9,7,8], [7,8,5], [8,5,6]; silent for the first two.
|
||||
let mut node = RollingMin::new(3);
|
||||
let got = drive(&mut node, &[9.0, 7.0, 8.0, 5.0, 6.0]);
|
||||
assert_eq!(got, vec![None, None, Some(7.0), Some(5.0), Some(5.0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rolling_min_evicts_the_expiring_window_minimum() {
|
||||
// length 2; the early small value must leave the window: [1,9]->1, [9,8]->8, [8,7]->7.
|
||||
let mut node = RollingMin::new(2);
|
||||
let got = drive(&mut node, &[1.0, 9.0, 8.0, 7.0]);
|
||||
assert_eq!(got, vec![None, Some(1.0), Some(8.0), Some(7.0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn rolling_min_length_one_is_identity() {
|
||||
let mut node = RollingMin::new(1);
|
||||
assert_eq!(drive(&mut node, &[7.0, 9.0, 4.0]), vec![Some(7.0), Some(9.0), Some(4.0)]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn deque_matches_naive_window_min() {
|
||||
// the monotonic deque must equal a naive O(N) per-window min at every warmed
|
||||
// cycle, across a long deterministic series (mirrors RollingMax's cross-check).
|
||||
// u64 + wrapping arithmetic matches rolling_max.rs and avoids i32 overflow panics.
|
||||
let length = 37;
|
||||
let mut node = RollingMin::new(length);
|
||||
let series: Vec<f64> = (0..3_000u64)
|
||||
.map(|i| (i.wrapping_mul(1103515245).wrapping_add(12345) % 1000) as f64 * 0.5 - 250.0)
|
||||
.collect();
|
||||
let mut inputs = vec![AnyColumn::with_capacity(ScalarKind::F64, 1)];
|
||||
for (t, &x) in series.iter().enumerate() {
|
||||
inputs[0].push(Scalar::f64(x)).unwrap();
|
||||
let got = node.eval(Ctx::new(&inputs, Timestamp(0))).map(|r| r[0].f64());
|
||||
if t + 1 >= length {
|
||||
let want = series[t + 1 - length..=t].iter().cloned().fold(f64::INFINITY, f64::min);
|
||||
assert_eq!(got, Some(want), "at t={t}");
|
||||
} else {
|
||||
assert_eq!(got, None, "silent until warmed up at t={t}");
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn lookback_is_one_window_lives_in_node_state() {
|
||||
assert_eq!(RollingMin::new(20).lookbacks(), vec![1]);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn length_must_be_at_least_one() {
|
||||
let r = std::panic::catch_unwind(|| RollingMin::new(0));
|
||||
assert!(r.is_err(), "RollingMin::new(0) must panic");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn label_carries_the_window() {
|
||||
assert_eq!(RollingMin::new(20).label(), "RollingMin(20)");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_declares_length_param_and_value_output() {
|
||||
// the blueprint seam later tasks bootstrap through: the param-generic recipe
|
||||
// declares a single I64 `length` knob and an F64 `value` output (mirrors
|
||||
// RollingMax's builder_declares_length_param_and_value_output).
|
||||
let schema = RollingMin::builder().schema().clone();
|
||||
assert_eq!(schema.params, vec![ParamSpec { name: "length".into(), kind: ScalarKind::I64 }]);
|
||||
assert_eq!(schema.output.len(), 1);
|
||||
assert_eq!(schema.output[0].name, "value");
|
||||
assert_eq!(schema.output[0].kind, ScalarKind::F64);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_input_slot_is_named_series() {
|
||||
// the named f64 input the breakout graph wires close into (mirrors RollingMax's
|
||||
// builder_input_slot_is_named_series).
|
||||
assert_eq!(RollingMin::builder().schema().inputs[0].name, "series");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_bind_removes_length_from_param_space() {
|
||||
// a bound length reports an empty param surface; the open form keeps it (mirrors
|
||||
// RollingMax's builder_bind_removes_length_from_param_space).
|
||||
let bound = RollingMin::builder().named("channel").bind("length", Scalar::i64(20));
|
||||
assert!(bound.schema().params.is_empty());
|
||||
assert_eq!(RollingMin::builder().named("channel").params().len(), 1);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn builder_bound_node_builds_with_injected_length() {
|
||||
// built from an empty open slice, the bound builder yields a RollingMin(20) —
|
||||
// the `p[0].i64() as usize` build closure is exercised end-to-end (mirrors
|
||||
// RollingMax's builder_bound_node_builds_with_injected_length).
|
||||
let node = RollingMin::builder().bind("length", Scalar::i64(20)).build(&[]);
|
||||
assert_eq!(node.label(), "RollingMin(20)");
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user