Files
Aura/fieldtests/research-breakout-deepdive/t1c_residency_scaling.rs
T
Brummel f66036bd61 feat(aura-engine): Source::resident_records — probe ring residency via the trait
The residency probe was an inherent method on the concrete M1FieldSource, so a Box<dyn Source> consumer could not read it without downcasting (the adjacent friction noted in #95). Hoist it onto the Source trait as resident_records(&self) -> Option<usize> with a default of None, mirroring bounds()'s Option-for-unknown idiom: None = "does not report" (the eager VecSource, which has no streaming ring), distinct from Some(0) = "reports, holds none" (e.g. exhausted).

M1FieldSource overrides it to return Some(chunk len). Callers updated (streaming_seam, the two deep-dive fieldtest bins). Tests: a hermetic VecSource default-is-None check, and a &dyn Source probe in the gated residency test proving M1FieldSource residency is readable without a downcast.

refs #95
2026-06-17 23:32:50 +02:00

79 lines
2.9 KiB
Rust

// Task 1 (deep dive 2) — the smoking gun: does PROCESS residency scale with
// window length, contradicting the docs' "a 20-year window streams in the same
// memory as a one-day one" (C12 cycle-0041 realization)?
//
// Drives ONE M1FieldSource (Close) to exhaustion over increasingly long windows,
// reading VmHWM after each. The aura-level resident_records stays <= CHUNK_SIZE
// (the docs' asserted predicate) the whole time; the question is whether VmHWM
// (the OS peak RSS) is flat or grows with the window. Fresh DataServer per window
// so a cache inside it cannot mask the effect. PUBLIC SURFACE ONLY.
use std::sync::Arc;
use chrono::TimeZone;
use data_server::{DataServer, DEFAULT_DATA_PATH};
use aura_engine::Source;
use aura_ingest::{M1Field, M1FieldSource};
const SYMBOL: &str = "GER40";
fn vmhwm_mib() -> u64 {
let status = std::fs::read_to_string("/proc/self/status").unwrap_or_default();
for line in status.lines() {
if let Some(rest) = line.strip_prefix("VmHWM:") {
let kib: u64 = rest.split_whitespace().next().and_then(|s| s.parse().ok()).unwrap_or(0);
return kib / 1024;
}
}
0
}
fn ms(y0: i32, m0: u32, y1: i32, m1: u32) -> (i64, i64) {
let from = chrono::Utc.with_ymd_and_hms(y0, m0, 1, 0, 0, 0).unwrap().timestamp_millis();
let (ny, nm) = if m1 == 12 { (y1 + 1, 1) } else { (y1, m1 + 1) };
let next = chrono::Utc.with_ymd_and_hms(ny, nm, 1, 0, 0, 0).unwrap().timestamp_millis();
(from, next - 1)
}
fn drive(label: &str, from_ms: i64, to_ms: i64) {
// Fresh server each time: no shared chunk cache carrying over.
let server = Arc::new(DataServer::new(DEFAULT_DATA_PATH));
let mut s = M1FieldSource::open(&server, SYMBOL, Some(from_ms), Some(to_ms), M1Field::Close)
.expect("window overlaps data");
let mut total = 0u64;
let mut peak_resident = 0usize;
while s.peek().is_some() {
peak_resident = peak_resident.max(s.resident_records().unwrap_or(0));
Source::next(&mut s);
total += 1;
}
println!(
"{label:<16} records={:<9} aura-resident<=1024? {:<5} | VmHWM={} MiB",
total,
peak_resident <= 1024,
vmhwm_mib(),
);
}
fn main() {
if !Arc::new(DataServer::new(DEFAULT_DATA_PATH)).has_symbol(SYMBOL) {
println!("skip: no {SYMBOL} data");
return;
}
println!("=== Task 1 deep dive: does VmHWM scale with window length? ===");
println!("(the docs promise it does NOT — O(one chunk) residency)\n");
println!("baseline VmHWM = {} MiB\n", vmhwm_mib());
let (f, t) = ms(2024, 9, 2024, 9);
drive("1 month", f, t);
let (f, t) = ms(2024, 1, 2024, 6);
drive("6 months", f, t);
let (f, t) = ms(2024, 1, 2024, 12);
drive("1 year", f, t);
let (f, t) = ms(2021, 1, 2024, 12);
drive("4 years", f, t);
let (f, t) = ms(2014, 8, 2024, 12);
drive("~10 years", f, t);
}