use std::collections::BTreeMap;
use std::fs;
use std::path::{Path, PathBuf};
use arrow_array::RecordBatch;
use optionstratlib::OptionStyle;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use serde::{Deserialize, Serialize};
use serde_json::Value;
use crate::error::{BundleError, ConfigError};
mod tables;
mod timeline;
mod validate;
pub use timeline::{Playback, PlaybackSpeed, TimelineCursor};
pub use validate::{BundleDivergence, ORACLE_ABS_TOL, ORACLE_REL_TOL, compare_bundles};
pub(crate) use validate::parse_contract_id;
pub const CONTRACT_ID_FORMAT: &str = "v1:{UNDERLYING}:{expiration_ns}:{strike_cents}:{C|P}";
pub const CONTRACT_ID_VERSION_PREFIX: &str = "v1";
pub const CONTRACT_ID_UNDERLYING_PATTERN: &str = "^[A-Z0-9._]{1,32}$";
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
pub struct BundleManifest {
pub schema: String,
pub run_id: String,
pub created_utc: String,
pub code_version: String,
pub lockfile_sha256: String,
pub seed: u64,
pub config: Value,
pub strategy: Value,
pub data_source: Value,
pub metrics: Value,
pub row_counts: BTreeMap<String, u64>,
}
impl BundleManifest {
pub fn capital_config(&self) -> Result<CapitalConfig, serde_json::Error> {
CapitalConfig::deserialize(&self.config)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
pub struct CapitalConfig {
pub initial_capital: u64,
}
impl CapitalConfig {
pub fn capital_cents(&self) -> Result<i64, BundleError> {
i64::try_from(self.initial_capital).map_err(|_| {
BundleError::Invariant(format!(
"config.initial_capital {} exceeds the i64 cents domain",
self.initial_capital
))
})
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[repr(u8)]
pub enum PositionSide {
Long,
Short,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[repr(u8)]
pub enum ExecMode {
Naive,
Realistic,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct Fill {
pub step: u32,
pub ts_ns: i64,
pub strategy_run_id: String,
pub trade_id: u64,
pub position_id: u64,
pub order_id: u64,
pub fill_seq: u32,
pub underlying: String,
pub expiration_ns: i64,
pub contract_id: String,
pub strike_cents: u64,
#[serde(with = "option_style_serde")]
pub style: OptionStyle,
pub side: PositionSide,
pub quantity: u32,
pub price_cents: u64,
pub fees_cents: u64,
pub slippage_cents: i64,
pub mode: ExecMode,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct EquityPoint {
pub step: u32,
pub ts_ns: i64,
pub cash_cents: i64,
pub position_value_cents: i64,
pub equity_cents: i64,
pub drawdown: f64,
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct PositionRow {
pub step: u32,
pub ts_ns: i64,
pub position_id: u64,
pub trade_id: u64,
pub contract_id: String,
pub side: PositionSide,
pub quantity: u32,
pub avg_price_cents: u64,
pub mark_cents: u64,
pub unrealized_cents: i64,
pub stale_mark: bool,
pub exit_reason: Option<String>,
pub open_at_end: bool,
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
#[serde(deny_unknown_fields)]
pub struct GreeksAttribution {
pub step: u32,
pub ts_ns: i64,
pub theta_pnl_cents: i64,
pub delta_pnl_cents: i64,
pub vega_pnl_cents: i64,
pub spread_capture_cents: i64,
pub fees_cents: u64,
pub residual_cents: i64,
}
#[derive(Debug, Clone, PartialEq)]
pub struct LoadedBundle {
pub manifest: BundleManifest,
pub fills: Vec<Fill>,
pub equity: Vec<EquityPoint>,
pub positions: Vec<PositionRow>,
pub greeks: Vec<GreeksAttribution>,
}
const TABLES: [(&str, &str); 4] = [
("fills.parquet", "fills"),
("equity_curve.parquet", "equity_curve"),
("positions.parquet", "positions"),
("greeks_attribution.parquet", "greeks_attribution"),
];
const MANIFEST_FILE: &str = "manifest.json";
pub const SUPPORTED_SCHEMA: &str = "ironcondor.bundle.v1";
pub const MAX_MANIFEST_BYTES: u64 = 8 * 1024 * 1024;
pub const MAX_TABLE_BYTES: u64 = 512 * 1024 * 1024;
pub const MAX_TABLE_ROWS: u64 = 5_000_000;
pub const MAX_WORKING_SET: u64 = 2 * 1024 * 1024 * 1024;
pub const MAX_BATCH_ROWS: usize = 65_536;
pub const MAX_BATCH_BYTES: u64 = 256 * 1024 * 1024;
pub const DECODED_OVERHEAD_PERMILLE: u64 = 1_500;
pub const MAX_EXPANSION_RATIO: u64 = 20;
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct ResourceCeilings {
pub max_manifest_bytes: u64,
pub max_table_bytes: u64,
pub max_table_rows: u64,
pub max_working_set: u64,
pub max_batch_rows: usize,
pub max_batch_bytes: u64,
pub decoded_overhead_permille: u64,
pub max_expansion_ratio: u64,
}
impl Default for ResourceCeilings {
fn default() -> Self {
Self {
max_manifest_bytes: MAX_MANIFEST_BYTES,
max_table_bytes: MAX_TABLE_BYTES,
max_table_rows: MAX_TABLE_ROWS,
max_working_set: MAX_WORKING_SET,
max_batch_rows: MAX_BATCH_ROWS,
max_batch_bytes: MAX_BATCH_BYTES,
decoded_overhead_permille: DECODED_OVERHEAD_PERMILLE,
max_expansion_ratio: MAX_EXPANSION_RATIO,
}
}
}
impl ResourceCeilings {
pub fn validate(&self) -> Result<(), ConfigError> {
let invalid = |field: &str, reason: String| ConfigError::InvalidValue {
field: field.to_owned(),
reason,
};
if self.max_manifest_bytes == 0 {
return Err(invalid(
"replay.max_manifest_bytes",
"must be > 0".to_owned(),
));
}
if self.max_table_bytes == 0 {
return Err(invalid("replay.max_table_bytes", "must be > 0".to_owned()));
}
if self.max_table_rows == 0 {
return Err(invalid("replay.max_table_rows", "must be > 0".to_owned()));
}
if self.max_working_set == 0 {
return Err(invalid("replay.max_working_set", "must be > 0".to_owned()));
}
if self.max_batch_rows == 0 {
return Err(invalid("replay.max_batch_rows", "must be > 0".to_owned()));
}
if self.max_batch_bytes == 0 {
return Err(invalid("replay.max_batch_bytes", "must be > 0".to_owned()));
}
if self.max_batch_bytes > self.max_working_set {
return Err(invalid(
"replay.max_batch_bytes",
format!(
"per-batch cap {} must not exceed the working-set ceiling {}",
self.max_batch_bytes, self.max_working_set
),
));
}
if self.decoded_overhead_permille < 1_000 {
return Err(invalid(
"replay.decoded_overhead_permille",
format!(
"must be >= 1000 (1.0x); got {}",
self.decoded_overhead_permille
),
));
}
if self.max_expansion_ratio == 0 {
return Err(invalid(
"replay.max_expansion_ratio",
"must be >= 1".to_owned(),
));
}
Ok(())
}
}
#[cold]
#[inline(never)]
fn too_large(detail: String) -> BundleError {
BundleError::TooLarge(detail)
}
#[cold]
#[inline(never)]
fn invariant(detail: String) -> BundleError {
BundleError::Invariant(detail)
}
#[cold]
#[inline(never)]
fn io_err(detail: String) -> BundleError {
BundleError::Io(detail)
}
#[cold]
#[inline(never)]
fn parquet_err(detail: String) -> BundleError {
BundleError::Parquet(detail)
}
fn catch_decode_panic<T>(
file: &str,
op: impl FnOnce() -> Result<T, BundleError>,
) -> Result<T, BundleError> {
crate::terminal::contained(op).unwrap_or_else(|| {
Err(parquet_err(format!(
"{file}: upstream Parquet/Arrow decoder panicked on malformed input"
)))
})
}
#[cold]
#[inline(never)]
fn missing_table(detail: String) -> BundleError {
BundleError::MissingTable(detail)
}
const MAX_SCHEMA_TAG_CHARS: usize = 64;
fn clamp_schema_tag(tag: String) -> String {
if tag.chars().count() <= MAX_SCHEMA_TAG_CHARS {
return tag;
}
let mut clamped: String = tag.chars().take(MAX_SCHEMA_TAG_CHARS).collect();
clamped.push('…');
clamped
}
#[cold]
#[inline(never)]
fn unsupported_schema(tag: String) -> BundleError {
BundleError::UnsupportedSchema(clamp_schema_tag(tag))
}
const MAX_ECHO_CHARS: usize = 64;
fn clamp_echo(value: &str) -> String {
if value.chars().count() <= MAX_ECHO_CHARS {
return value.to_owned();
}
let mut clamped: String = value.chars().take(MAX_ECHO_CHARS).collect();
clamped.push('…');
clamped
}
#[inline]
fn footer_i64_to_u64(raw: i64, what: &str) -> Result<u64, BundleError> {
u64::try_from(raw).map_err(|_| too_large(format!("{what}: negative footer value {raw}")))
}
#[inline]
fn apply_overhead(bytes: u64, permille: u64) -> Result<u64, BundleError> {
bytes
.checked_mul(permille)
.map(|scaled| scaled / 1_000)
.ok_or_else(|| {
too_large(format!(
"decoded size {bytes} overflows the overhead multiplier"
))
})
}
#[derive(Debug, Clone, Copy)]
struct WorkingSetBudget {
used: u64,
max_working_set: u64,
max_batch_bytes: u64,
}
impl WorkingSetBudget {
fn new(ceilings: &ResourceCeilings) -> Self {
Self {
used: 0,
max_working_set: ceilings.max_working_set,
max_batch_bytes: ceilings.max_batch_bytes,
}
}
fn account(&mut self, batch_bytes: u64) -> Result<(), BundleError> {
if batch_bytes > self.max_batch_bytes {
return Err(too_large(format!(
"decoded batch {batch_bytes} B exceeds per-batch cap {} B",
self.max_batch_bytes
)));
}
let next = self
.used
.checked_add(batch_bytes)
.ok_or_else(|| too_large("cumulative working set overflowed u64".to_owned()))?;
if next > self.max_working_set {
return Err(too_large(format!(
"cumulative working set {next} B would exceed ceiling {} B",
self.max_working_set
)));
}
self.used = next;
Ok(())
}
#[inline]
fn used(&self) -> u64 {
self.used
}
}
fn validate_row_counts_shape(row_counts: &BTreeMap<String, u64>) -> Result<(), BundleError> {
for (_file, key) in TABLES {
if !row_counts.contains_key(key) {
return Err(invariant(format!(
"row_counts missing required key `{key}`"
)));
}
}
if row_counts.len() != TABLES.len() {
return Err(invariant(format!(
"row_counts carries {} keys; exactly the {} table names are required",
row_counts.len(),
TABLES.len()
)));
}
Ok(())
}
#[derive(Debug, Clone)]
pub struct BundleReader {
root: PathBuf,
manifest: BundleManifest,
ceilings: ResourceCeilings,
}
impl BundleReader {
pub fn open(root: impl Into<PathBuf>) -> Result<Self, BundleError> {
Self::open_with_ceilings(root, ResourceCeilings::default())
}
pub fn open_with_ceilings(
root: impl Into<PathBuf>,
ceilings: ResourceCeilings,
) -> Result<Self, BundleError> {
ceilings.validate()?;
let root = root.into();
let meta = fs::metadata(&root).map_err(|e| {
io_err(format!(
"cannot access bundle directory {}: {e}",
root.display()
))
})?;
if !meta.is_dir() {
return Err(io_err(format!(
"bundle root is not a directory: {}",
root.display()
)));
}
let manifest_path = root.join(MANIFEST_FILE);
let manifest_len = match fs::metadata(&manifest_path) {
Ok(meta) => meta.len(),
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return Err(missing_table(MANIFEST_FILE.to_owned()));
}
Err(e) => return Err(io_err(format!("stat {MANIFEST_FILE}: {e}"))),
};
if manifest_len > ceilings.max_manifest_bytes {
return Err(too_large(format!(
"{MANIFEST_FILE} is {manifest_len} B; exceeds manifest ceiling {} B",
ceilings.max_manifest_bytes
)));
}
let manifest_bytes = match fs::read(&manifest_path) {
Ok(bytes) => bytes,
Err(e) if e.kind() == std::io::ErrorKind::NotFound => {
return Err(missing_table(MANIFEST_FILE.to_owned()));
}
Err(e) => return Err(io_err(format!("read {MANIFEST_FILE}: {e}"))),
};
let manifest: BundleManifest = serde_json::from_slice(&manifest_bytes)
.map_err(|e| invariant(format!("malformed {MANIFEST_FILE}: {e}")))?;
if manifest.schema != SUPPORTED_SCHEMA {
return Err(unsupported_schema(manifest.schema));
}
validate_row_counts_shape(&manifest.row_counts)?;
for (file, _key) in TABLES {
if !root.join(file).is_file() {
return Err(missing_table(file.to_owned()));
}
}
Ok(Self {
root,
manifest,
ceilings,
})
}
#[must_use]
pub fn root(&self) -> &Path {
&self.root
}
#[must_use]
pub fn manifest(&self) -> &BundleManifest {
&self.manifest
}
#[must_use]
pub fn ceilings(&self) -> &ResourceCeilings {
&self.ceilings
}
pub fn load(&self) -> Result<LoadedBundle, BundleError> {
self.load_cancellable(&|| false)
}
pub fn load_cancellable(
&self,
cancelled: &dyn Fn() -> bool,
) -> Result<LoadedBundle, BundleError> {
if cancelled() {
return Err(BundleError::Cancelled);
}
let mut budget = WorkingSetBudget::new(&self.ceilings);
let mut fills: Vec<Fill> = Vec::new();
if cancelled() {
return Err(BundleError::Cancelled);
}
self.scan_table("fills.parquet", "fills", &mut budget, cancelled, &mut |b| {
tables::read_fills(b, &mut fills)
})?;
let mut equity: Vec<EquityPoint> = Vec::new();
if cancelled() {
return Err(BundleError::Cancelled);
}
self.scan_table(
"equity_curve.parquet",
"equity_curve",
&mut budget,
cancelled,
&mut |b| tables::read_equity(b, &mut equity),
)?;
let mut positions: Vec<PositionRow> = Vec::new();
if cancelled() {
return Err(BundleError::Cancelled);
}
self.scan_table(
"positions.parquet",
"positions",
&mut budget,
cancelled,
&mut |b| tables::read_positions(b, &mut positions),
)?;
let mut greeks: Vec<GreeksAttribution> = Vec::new();
if cancelled() {
return Err(BundleError::Cancelled);
}
self.scan_table(
"greeks_attribution.parquet",
"greeks_attribution",
&mut budget,
cancelled,
&mut |b| tables::read_greeks(b, &mut greeks),
)?;
if cancelled() {
return Err(BundleError::Cancelled);
}
let loaded = LoadedBundle {
manifest: self.manifest.clone(),
fills,
equity,
positions,
greeks,
};
validate::run_validation_chain(&loaded)?;
Ok(loaded)
}
fn scan_table(
&self,
file: &str,
key: &str,
budget: &mut WorkingSetBudget,
cancelled: &dyn Fn() -> bool,
decode: &mut dyn FnMut(&RecordBatch) -> Result<u64, BundleError>,
) -> Result<(), BundleError> {
let path = self.root.join(file);
let file_bytes = fs::metadata(&path)
.map_err(|e| match e.kind() {
std::io::ErrorKind::NotFound => missing_table(file.to_owned()),
_ => io_err(format!("stat {file}: {e}")),
})?
.len();
if file_bytes > self.ceilings.max_table_bytes {
return Err(too_large(format!(
"{file} is {file_bytes} B; exceeds per-file ceiling {} B",
self.ceilings.max_table_bytes
)));
}
let handle = fs::File::open(&path).map_err(|e| match e.kind() {
std::io::ErrorKind::NotFound => missing_table(file.to_owned()),
_ => io_err(format!("open {file}: {e}")),
})?;
let builder = catch_decode_panic(file, move || {
ParquetRecordBatchReaderBuilder::try_new(handle)
.map_err(|e| parquet_err(format!("{file}: {e}")))
})?;
let (footer_rows, uncompressed, compressed) = {
let metadata = builder.metadata();
let footer_rows = footer_i64_to_u64(
metadata.file_metadata().num_rows(),
&format!("{file} footer row count"),
)?;
let mut uncompressed: u64 = 0;
let mut compressed: u64 = 0;
for rg in metadata.row_groups() {
uncompressed = uncompressed
.checked_add(footer_i64_to_u64(
rg.total_byte_size(),
&format!("{file} row-group uncompressed size"),
)?)
.ok_or_else(|| {
too_large(format!("{file}: footer uncompressed size overflowed u64"))
})?;
compressed = compressed
.checked_add(footer_i64_to_u64(
rg.compressed_size(),
&format!("{file} row-group compressed size"),
)?)
.ok_or_else(|| {
too_large(format!("{file}: footer compressed size overflowed u64"))
})?;
}
(footer_rows, uncompressed, compressed)
};
if footer_rows > self.ceilings.max_table_rows {
return Err(too_large(format!(
"{file} footer declares {footer_rows} rows; exceeds per-table ceiling {}",
self.ceilings.max_table_rows
)));
}
let declared = self
.manifest
.row_counts
.get(key)
.copied()
.ok_or_else(|| invariant(format!("row_counts missing key `{key}`")))?;
if declared != footer_rows {
return Err(invariant(format!(
"{file}: row_counts says {declared} but the Parquet footer says {footer_rows}"
)));
}
let on_disk = file_bytes.max(compressed);
let Some(bomb_limit) = on_disk.checked_mul(self.ceilings.max_expansion_ratio) else {
return Err(too_large(format!(
"{file}: on-disk size {on_disk} B is too large to bound at {}x - \
rejected before decode",
self.ceilings.max_expansion_ratio
)));
};
if uncompressed > bomb_limit {
return Err(too_large(format!(
"{file}: uncompressed {uncompressed} B exceeds {}x its on-disk size — \
decompression bomb rejected",
self.ceilings.max_expansion_ratio
)));
}
let estimate = apply_overhead(uncompressed, self.ceilings.decoded_overhead_permille)?;
if budget
.used()
.checked_add(estimate)
.is_none_or(|total| total > self.ceilings.max_working_set)
{
return Err(too_large(format!(
"{file}: estimated decoded working set would exceed ceiling {} B",
self.ceilings.max_working_set
)));
}
let mut reader = catch_decode_panic(file, || {
builder
.with_batch_size(self.ceilings.max_batch_rows)
.build()
.map_err(|e| parquet_err(format!("{file}: {e}")))
})?;
let mut decoded_rows: usize = 0;
loop {
if cancelled() {
return Err(BundleError::Cancelled);
}
let next = catch_decode_panic(file, || {
reader
.next()
.transpose()
.map_err(|e| parquet_err(format!("{file}: {e}")))
})?;
let Some(batch) = next else { break };
let batch_bytes = u64::try_from(batch.get_array_memory_size())
.map_err(|_| too_large(format!("{file}: decoded batch size exceeds u64")))?;
let measured = apply_overhead(batch_bytes, self.ceilings.decoded_overhead_permille)?;
budget.account(measured)?;
let rows = batch.num_rows();
let retained = decode(&batch)?;
budget.account(retained)?;
decoded_rows = decoded_rows
.checked_add(rows)
.ok_or_else(|| too_large(format!("{file}: decoded row count overflowed usize")))?;
}
let decoded_rows = u64::try_from(decoded_rows)
.map_err(|_| too_large(format!("{file}: decoded row count exceeds u64")))?;
if decoded_rows != footer_rows {
return Err(invariant(format!(
"{file}: decoded {decoded_rows} rows but the Parquet footer declares {footer_rows}"
)));
}
Ok(())
}
}
mod option_style_serde {
use optionstratlib::OptionStyle;
use serde::de::Error;
use serde::{Deserialize, Deserializer, Serializer};
const VARIANTS: &[&str] = &["call", "put"];
pub(super) fn serialize<S>(style: &OptionStyle, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
serializer.serialize_str(match style {
OptionStyle::Call => "call",
OptionStyle::Put => "put",
})
}
pub(super) fn deserialize<'de, D>(deserializer: D) -> Result<OptionStyle, D::Error>
where
D: Deserializer<'de>,
{
let raw = String::deserialize(deserializer)?;
match raw.as_str() {
"call" => Ok(OptionStyle::Call),
"put" => Ok(OptionStyle::Put),
_ => Err(D::Error::unknown_variant(&raw, VARIANTS)),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
#[track_caller]
fn from_json<T: for<'de> Deserialize<'de>>(json: &str) -> T {
match serde_json::from_str::<T>(json) {
Ok(value) => value,
Err(e) => panic!("expected `{json}` to deserialize: {e}"),
}
}
#[track_caller]
fn to_json<T: Serialize>(value: &T) -> String {
match serde_json::to_string(value) {
Ok(s) => s,
Err(e) => panic!("serialize failed: {e}"),
}
}
#[track_caller]
fn to_value<T: Serialize>(value: &T) -> Value {
match serde_json::to_value(value) {
Ok(v) => v,
Err(e) => panic!("to_value failed: {e}"),
}
}
fn well_formed_manifest_json() -> String {
r#"{
"schema": "ironcondor.bundle.v1",
"run_id": "run-abc123",
"created_utc": "2026-07-16T00:00:00Z",
"code_version": "0.3.0",
"lockfile_sha256": "deadbeef",
"seed": 42,
"config": {
"initial_capital": 1000000,
"mode": "realistic",
"fees": { "per_contract_cents": 65, "per_order_cents": 100 },
"slippage": { "model": "none" }
},
"strategy": { "kind": "iron_condor", "params": {} },
"data_source": { "kind": "parquet", "path": "tape.parquet", "sha256": "cafe" },
"metrics": { "sharpe": 1.5, "nested": { "total_pnl_cents": 12345 } },
"row_counts": {
"fills": 4, "equity_curve": 10, "positions": 8, "greeks_attribution": 10
}
}"#
.to_owned()
}
fn real_ironcondor_manifest_json() -> String {
r#"{
"code_version": "0.5.0",
"config": {
"data_source": {
"kind": "parquet",
"path": "conformance_condor.parquet",
"sha256": ""
},
"fees": {
"per_contract_cents": 65,
"per_order_cents": 100
},
"initial_capital": 10000000,
"limits": {
"max_contracts_per_snapshot": 50000,
"max_decompressed_bytes": 8589934592,
"max_file_bytes": 4294967296,
"max_manifest_bytes": 16777216,
"max_rows_per_table": 100000000,
"max_steps": 100000,
"max_string_len": 65536,
"max_total_bytes": 2147483648
},
"liquidity_profile": {
"decay": "1",
"depth_levels": 2,
"touch_size": {
"contracts": 3,
"kind": "flat"
}
},
"marketable_cap_ticks": 10,
"mode": "realistic",
"output_dir": "conformance-out",
"overwrite": false,
"seed": 7,
"slippage": {
"model": "none"
}
},
"created_utc": "1970-01-01T00:00:00+00:00",
"data_source": {
"kind": "parquet",
"path": "conformance_condor.parquet",
"sha256": ""
},
"lockfile_sha256": "7757b4ea1d081633037b2f1795252453baca654186b607bb8c85f14e0009740b",
"metrics": {
"custom_metrics": {
"max_drawdown_cents": "19625",
"net_premium_cents": "1136000",
"realized_pnl_cents": "-7425"
},
"general_performance": {
"sharpe_ratio": "-0.577350269189625731058868041",
"total_return": "-0.0019625"
}
},
"row_counts": {
"equity_curve": 5,
"fills": 10,
"greeks_attribution": 5,
"positions": 18
},
"run_id": "c4cd155fc156f3ed1488661d9cfd448af59b216c962ea4716f76685f92b21459",
"schema": "ironcondor.bundle.v1",
"seed": 7,
"strategy": {
"IronCondor": {
"close_fee": 65,
"long_call_strike": 520000,
"long_put_strike": 480000,
"open_fee": 65,
"quantity": 5,
"short_call_strike": 510000,
"short_put_strike": 490000,
"underlying": "SPX",
"underlying_price": 500000
}
}
}"#
.to_owned()
}
#[test]
fn test_bundle_manifest_roundtrips_well_formed() {
let manifest: BundleManifest = from_json(&well_formed_manifest_json());
assert_eq!(manifest.schema, "ironcondor.bundle.v1");
assert_eq!(manifest.run_id, "run-abc123");
assert_eq!(manifest.seed, 42);
assert_eq!(manifest.row_counts.get("fills"), Some(&4));
assert_eq!(manifest.row_counts.get("equity_curve"), Some(&10));
let reparsed: BundleManifest = from_json(&to_json(&manifest));
assert_eq!(manifest, reparsed);
}
#[test]
fn test_bundle_manifest_tolerates_unknown_fields() {
let json = well_formed_manifest_json().replace(
"\"seed\": 42,",
"\"seed\": 42, \"future_added_field\": {\"x\": 1},",
);
let manifest: BundleManifest = from_json(&json);
assert_eq!(manifest.schema, "ironcondor.bundle.v1");
}
#[test]
fn test_bundle_manifest_missing_required_field_errors() {
let json = well_formed_manifest_json().replace("\"seed\": 42,", "");
assert!(
serde_json::from_str::<BundleManifest>(&json).is_err(),
"a manifest missing `seed` must fail to deserialize"
);
}
#[test]
fn test_bundle_manifest_metrics_stays_opaque() {
let manifest: BundleManifest = from_json(&well_formed_manifest_json());
let expected = serde_json::json!({
"sharpe": 1.5,
"nested": { "total_pnl_cents": 12345 }
});
assert_eq!(manifest.metrics, expected);
}
#[test]
fn test_capital_config_extracts_initial_capital_ignoring_other_fields() {
let manifest: BundleManifest = from_json(&well_formed_manifest_json());
let capital = match manifest.capital_config() {
Ok(c) => c,
Err(e) => panic!("capital_config should project initial_capital: {e}"),
};
assert_eq!(capital.initial_capital, 1_000_000);
}
#[test]
fn test_capital_config_reads_real_ironcondor_manifest() {
let manifest: BundleManifest = from_json(&real_ironcondor_manifest_json());
assert_eq!(manifest.schema, "ironcondor.bundle.v1");
assert_eq!(manifest.seed, 7);
assert_eq!(manifest.row_counts.get("positions"), Some(&18));
let capital = match manifest.capital_config() {
Ok(c) => c,
Err(e) => panic!("real manifest must project initial_capital: {e}"),
};
assert_eq!(capital.initial_capital, 10_000_000);
match capital.capital_cents() {
Ok(cents) => assert_eq!(cents, 10_000_000),
Err(e) => panic!("checked capital narrowing must succeed: {e}"),
}
}
#[test]
fn test_capital_config_reports_absent_initial_capital() {
let config = serde_json::json!({ "mode": "naive", "fees": { "per_order_cents": 5 } });
assert!(
CapitalConfig::deserialize(&config).is_err(),
"an absent initial_capital must not default to 0"
);
}
#[test]
fn test_capital_config_rejects_legacy_capital_cents_shape() {
let config = serde_json::json!({ "capital_cents": 1_000_000 });
assert!(
CapitalConfig::deserialize(&config).is_err(),
"the legacy capital_cents-only shape must not project as initial_capital"
);
}
#[test]
fn test_capital_config_initial_capital_is_unsigned() {
let config = serde_json::json!({ "initial_capital": -25 });
assert!(
CapitalConfig::deserialize(&config).is_err(),
"a negative initial_capital must fail against the unsigned wire type"
);
}
#[test]
fn test_capital_config_capital_cents_checked_conversion_overflows() {
let config = serde_json::json!({ "initial_capital": u64::MAX });
let capital = match CapitalConfig::deserialize(&config) {
Ok(c) => c,
Err(e) => panic!("u64::MAX initial_capital should deserialize: {e}"),
};
assert_eq!(capital.initial_capital, u64::MAX);
match capital.capital_cents() {
Err(BundleError::Invariant(_)) => {}
other => panic!("overflowing capital narrowing should be Invariant, got {other:?}"),
}
}
#[test]
fn test_position_side_wire_strings_roundtrip() {
assert_eq!(to_json(&PositionSide::Long), "\"long\"");
assert_eq!(to_json(&PositionSide::Short), "\"short\"");
assert_eq!(from_json::<PositionSide>("\"long\""), PositionSide::Long);
assert_eq!(from_json::<PositionSide>("\"short\""), PositionSide::Short);
}
#[test]
fn test_position_side_rejects_unknown_string() {
assert!(serde_json::from_str::<PositionSide>("\"sideways\"").is_err());
assert!(serde_json::from_str::<PositionSide>("\"Long\"").is_err());
}
#[test]
fn test_exec_mode_wire_strings_roundtrip() {
assert_eq!(to_json(&ExecMode::Naive), "\"naive\"");
assert_eq!(to_json(&ExecMode::Realistic), "\"realistic\"");
assert_eq!(from_json::<ExecMode>("\"naive\""), ExecMode::Naive);
assert_eq!(from_json::<ExecMode>("\"realistic\""), ExecMode::Realistic);
}
#[test]
fn test_exec_mode_rejects_unknown_string() {
assert!(serde_json::from_str::<ExecMode>("\"paper\"").is_err());
}
#[test]
fn test_option_style_wire_strings_roundtrip_via_fill() {
let fill: Fill = from_json(&fill_json("call", "long", "realistic"));
assert_eq!(fill.style, OptionStyle::Call);
assert!(to_json(&fill).contains("\"style\":\"call\""));
let put: Fill = from_json(&fill_json("put", "short", "naive"));
assert_eq!(put.style, OptionStyle::Put);
assert!(to_json(&put).contains("\"style\":\"put\""));
}
#[test]
fn test_option_style_rejects_unknown_wire_string() {
assert!(serde_json::from_str::<Fill>(&fill_json("Call", "long", "naive")).is_err());
assert!(serde_json::from_str::<Fill>(&fill_json("american", "long", "naive")).is_err());
}
fn fill_json(style: &str, side: &str, mode: &str) -> String {
format!(
r#"{{
"step": 3,
"ts_ns": 1700000000000000000,
"strategy_run_id": "run-abc123",
"trade_id": 7,
"position_id": 11,
"order_id": 21,
"fill_seq": 0,
"underlying": "BTC",
"expiration_ns": 1735286400000000000,
"contract_id": "v1:BTC:1735286400000000000:6000000:C",
"strike_cents": 6000000,
"style": "{style}",
"side": "{side}",
"quantity": 2,
"price_cents": 12500,
"fees_cents": 30,
"slippage_cents": -15,
"mode": "{mode}"
}}"#
)
}
#[test]
fn test_fill_roundtrips() {
let fill: Fill = from_json(&fill_json("call", "long", "realistic"));
assert_eq!(fill.step, 3);
assert_eq!(fill.strike_cents, 6_000_000);
assert_eq!(fill.fees_cents, 30);
assert_eq!(fill.slippage_cents, -15);
assert_eq!(fill.side, PositionSide::Long);
assert_eq!(fill.mode, ExecMode::Realistic);
let reparsed: Fill = from_json(&to_json(&fill));
assert_eq!(fill, reparsed);
}
#[test]
fn test_fill_strict_rejects_unknown_field() {
let json = fill_json("call", "long", "naive")
.replace("\"step\": 3,", "\"step\": 3, \"unexpected\": 1,");
assert!(
serde_json::from_str::<Fill>(&json).is_err(),
"a row with an unknown field must be rejected (strict posture)"
);
}
#[test]
fn test_fill_missing_required_field_errors() {
let json = fill_json("call", "long", "naive").replace("\"quantity\": 2,", "");
assert!(serde_json::from_str::<Fill>(&json).is_err());
}
#[test]
fn test_equity_point_roundtrips() {
let json = r#"{
"step": 5,
"ts_ns": 1700000000000000000,
"cash_cents": 990000,
"position_value_cents": -1500,
"equity_cents": 988500,
"drawdown": -0.015
}"#;
let point: EquityPoint = from_json(json);
assert_eq!(point.equity_cents, 988_500);
assert_eq!(
point.cash_cents + point.position_value_cents,
point.equity_cents
);
assert!((point.drawdown - (-0.015)).abs() < 1e-12);
let reparsed: EquityPoint = from_json(&to_json(&point));
assert_eq!(point, reparsed);
}
#[test]
fn test_position_row_open_leg_has_null_exit_reason() {
let json = r#"{
"step": 4,
"ts_ns": 1700000000000000000,
"position_id": 11,
"trade_id": 7,
"contract_id": "v1:BTC:1735286400000000000:6000000:C",
"side": "short",
"quantity": 1,
"avg_price_cents": 12000,
"mark_cents": 11800,
"unrealized_cents": 200,
"stale_mark": false,
"exit_reason": null,
"open_at_end": true
}"#;
let row: PositionRow = from_json(json);
assert_eq!(row.side, PositionSide::Short);
assert_eq!(row.exit_reason, None);
assert!(row.open_at_end);
assert!(!row.stale_mark);
let reparsed: PositionRow = from_json(&to_json(&row));
assert_eq!(row, reparsed);
}
#[test]
fn test_position_row_terminal_carries_exit_reason() {
let json = r#"{
"step": 9,
"ts_ns": 1700000000000000000,
"position_id": 11,
"trade_id": 7,
"contract_id": "v1:BTC:1735286400000000000:6000000:C",
"side": "short",
"quantity": 1,
"avg_price_cents": 12000,
"mark_cents": 11800,
"unrealized_cents": 200,
"stale_mark": true,
"exit_reason": "profit_target",
"open_at_end": false
}"#;
let row: PositionRow = from_json(json);
assert_eq!(row.exit_reason.as_deref(), Some("profit_target"));
assert!(row.stale_mark);
}
#[test]
fn test_greeks_attribution_roundtrips() {
let json = r#"{
"step": 5,
"ts_ns": 1700000000000000000,
"theta_pnl_cents": 40,
"delta_pnl_cents": -120,
"vega_pnl_cents": 15,
"spread_capture_cents": 10,
"fees_cents": 30,
"residual_cents": 1
}"#;
let row: GreeksAttribution = from_json(json);
assert_eq!(row.theta_pnl_cents, 40);
assert_eq!(row.spread_capture_cents, 10);
assert_eq!(row.fees_cents, 30);
let reparsed: GreeksAttribution = from_json(&to_json(&row));
assert_eq!(row, reparsed);
}
#[test]
fn test_contract_id_grammar_constants_are_fixed() {
assert_eq!(
CONTRACT_ID_FORMAT,
"v1:{UNDERLYING}:{expiration_ns}:{strike_cents}:{C|P}"
);
assert_eq!(CONTRACT_ID_VERSION_PREFIX, "v1");
assert_eq!(CONTRACT_ID_UNDERLYING_PATTERN, "^[A-Z0-9._]{1,32}$");
assert!(CONTRACT_ID_FORMAT.starts_with(CONTRACT_ID_VERSION_PREFIX));
}
#[test]
fn test_run_id_is_a_plain_opaque_string() {
let manifest: BundleManifest = from_json(&well_formed_manifest_json());
let id: &str = &manifest.run_id;
assert_eq!(id, "run-abc123");
}
#[test]
fn test_only_pub_f64_field_is_drawdown() {
let src = include_str!("mod.rs");
let mut f64_field_lines = 0_usize;
for line in src.lines() {
let trimmed = line.trim_start();
if trimmed.starts_with("pub ") && trimmed.contains(": f64") {
assert!(
trimmed.contains("drawdown"),
"unexpected f64 field (money must be integer cents): `{trimmed}`"
);
f64_field_lines += 1;
}
}
assert_eq!(
f64_field_lines, 1,
"exactly one public f64 field (EquityPoint::drawdown) is expected"
);
}
#[test]
fn test_fill_serializes_to_json_object() {
let fill: Fill = from_json(&fill_json("put", "short", "realistic"));
assert!(to_value(&fill).is_object());
}
fn temp_bundle_dir() -> tempfile::TempDir {
match tempfile::tempdir() {
Ok(d) => d,
Err(e) => panic!("create tempdir: {e}"),
}
}
fn write_parquet(path: &Path, num_rows: usize) {
use std::sync::Arc;
use arrow_array::{ArrayRef, Int32Array, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
use parquet::arrow::ArrowWriter;
let schema = Arc::new(Schema::new(vec![Field::new(
"step",
DataType::Int32,
false,
)]));
let steps: Vec<i32> = (0..num_rows)
.map(|i| i32::try_from(i).unwrap_or(i32::MAX))
.collect();
let column: ArrayRef = Arc::new(Int32Array::from(steps));
let batch = match RecordBatch::try_new(Arc::clone(&schema), vec![column]) {
Ok(b) => b,
Err(e) => panic!("build record batch: {e}"),
};
let file = match std::fs::File::create(path) {
Ok(f) => f,
Err(e) => panic!("create {}: {e}", path.display()),
};
let mut writer = match ArrowWriter::try_new(file, schema, None) {
Ok(w) => w,
Err(e) => panic!("arrow writer: {e}"),
};
if let Err(e) = writer.write(&batch) {
panic!("write batch: {e}");
}
if let Err(e) = writer.close() {
panic!("close writer: {e}");
}
}
fn manifest_json_with(schema: &str, counts: [u64; 4], extra: bool) -> String {
let [c_fills, c_eq, c_pos, c_greeks] = counts;
let extra_field = if extra {
"\"future_field\": {\"x\": 1},"
} else {
""
};
format!(
r#"{{ "schema": "{schema}", "run_id": "run-abc123",
"created_utc": "2026-07-16T00:00:00Z", "code_version": "0.3.0",
"lockfile_sha256": "deadbeef", "seed": 42, {extra_field}
"config": {{ "initial_capital": 1000000, "mode": "realistic" }},
"strategy": {{ "kind": "iron_condor" }},
"data_source": {{ "kind": "parquet", "path": "tape.parquet", "sha256": "cafe" }},
"metrics": {{ "sharpe": 1.5 }},
"row_counts": {{ "fills": {c_fills}, "equity_curve": {c_eq},
"positions": {c_pos}, "greeks_attribution": {c_greeks} }} }}"#
)
}
fn write_bundle(
dir: &Path,
schema: &str,
file_rows: [usize; 4],
counts: [u64; 4],
extra: bool,
) {
if let Err(e) = std::fs::write(
dir.join(MANIFEST_FILE),
manifest_json_with(schema, counts, extra),
) {
panic!("write manifest: {e}");
}
let [(f0, _), (f1, _), (f2, _), (f3, _)] = TABLES;
let [r0, r1, r2, r3] = file_rows;
write_parquet(&dir.join(f0), r0);
write_parquet(&dir.join(f1), r1);
write_parquet(&dir.join(f2), r2);
write_parquet(&dir.join(f3), r3);
}
fn fills_batch(n: usize) -> RecordBatch {
use std::sync::Arc;
use arrow_array::{ArrayRef, Int32Array, Int64Array, StringArray};
use arrow_schema::{DataType, Field, Schema};
let steps: Vec<i32> = (0..n)
.map(|i| i32::try_from(i).unwrap_or(i32::MAX))
.collect();
let ts: Vec<i64> = steps
.iter()
.map(|&s| 1_700_000_000_000_000_000 + i64::from(s))
.collect();
let trade: Vec<i64> = steps.iter().map(|&s| 100 + i64::from(s)).collect();
let pos: Vec<i64> = steps.iter().map(|&s| 200 + i64::from(s)).collect();
let order: Vec<i64> = steps.iter().map(|&s| i64::from(s)).collect();
let schema = Arc::new(Schema::new(vec![
Field::new("step", DataType::Int32, false),
Field::new("ts_ns", DataType::Int64, false),
Field::new("strategy_run_id", DataType::Utf8, false),
Field::new("trade_id", DataType::Int64, false),
Field::new("position_id", DataType::Int64, false),
Field::new("order_id", DataType::Int64, false),
Field::new("fill_seq", DataType::Int32, false),
Field::new("underlying", DataType::Utf8, false),
Field::new("expiration_ns", DataType::Int64, false),
Field::new("contract_id", DataType::Utf8, false),
Field::new("strike_cents", DataType::Int64, false),
Field::new("style", DataType::Utf8, false),
Field::new("side", DataType::Utf8, false),
Field::new("quantity", DataType::Int32, false),
Field::new("price_cents", DataType::Int64, false),
Field::new("fees_cents", DataType::Int64, false),
Field::new("slippage_cents", DataType::Int64, false),
Field::new("mode", DataType::Utf8, false),
]));
let cols: Vec<ArrayRef> = vec![
Arc::new(Int32Array::from(steps)),
Arc::new(Int64Array::from(ts)),
Arc::new(StringArray::from(vec!["run-abc123"; n])),
Arc::new(Int64Array::from(trade)),
Arc::new(Int64Array::from(pos)),
Arc::new(Int64Array::from(order)),
Arc::new(Int32Array::from(vec![0_i32; n])),
Arc::new(StringArray::from(vec!["BTC"; n])),
Arc::new(Int64Array::from(vec![1_735_286_400_000_000_000_i64; n])),
Arc::new(StringArray::from(vec![
"v1:BTC:1735286400000000000:6000000:C";
n
])),
Arc::new(Int64Array::from(vec![6_000_000_i64; n])),
Arc::new(StringArray::from(vec!["call"; n])),
Arc::new(StringArray::from(vec!["long"; n])),
Arc::new(Int32Array::from(vec![1_i32; n])),
Arc::new(Int64Array::from(vec![12_500_i64; n])),
Arc::new(Int64Array::from(vec![30_i64; n])),
Arc::new(Int64Array::from(vec![-15_i64; n])),
Arc::new(StringArray::from(vec!["realistic"; n])),
];
match RecordBatch::try_new(schema, cols) {
Ok(b) => b,
Err(e) => panic!("build fills batch: {e}"),
}
}
fn fills_batch_extra_column(n: usize) -> RecordBatch {
use std::sync::Arc;
use arrow_array::{ArrayRef, Int64Array};
use arrow_schema::{DataType, Field, Schema};
let base = fills_batch(n);
let mut fields: Vec<Field> = base
.schema()
.fields()
.iter()
.map(|f| f.as_ref().clone())
.collect();
fields.push(Field::new("future_metric_cents", DataType::Int64, false));
let mut cols: Vec<ArrayRef> = base.columns().to_vec();
cols.push(Arc::new(Int64Array::from(vec![0_i64; n])));
match RecordBatch::try_new(Arc::new(Schema::new(fields)), cols) {
Ok(b) => b,
Err(e) => panic!("build extra-column fills batch: {e}"),
}
}
fn equity_batch(n: usize) -> RecordBatch {
use std::sync::Arc;
use arrow_array::{ArrayRef, Float64Array, Int32Array, Int64Array};
use arrow_schema::{DataType, Field, Schema};
let steps: Vec<i32> = (0..n)
.map(|i| i32::try_from(i).unwrap_or(i32::MAX))
.collect();
let ts: Vec<i64> = steps
.iter()
.map(|&s| 1_700_000_000_000_000_000 + i64::from(s))
.collect();
let cash: Vec<i64> = steps.iter().map(|&s| 990_000 + i64::from(s)).collect();
let equity: Vec<i64> = steps.iter().map(|&s| 988_500 + i64::from(s)).collect();
let schema = Arc::new(Schema::new(vec![
Field::new("step", DataType::Int32, false),
Field::new("ts_ns", DataType::Int64, false),
Field::new("cash_cents", DataType::Int64, false),
Field::new("position_value_cents", DataType::Int64, false),
Field::new("equity_cents", DataType::Int64, false),
Field::new("drawdown", DataType::Float64, false),
]));
let cols: Vec<ArrayRef> = vec![
Arc::new(Int32Array::from(steps)),
Arc::new(Int64Array::from(ts)),
Arc::new(Int64Array::from(cash)),
Arc::new(Int64Array::from(vec![-1_500_i64; n])),
Arc::new(Int64Array::from(equity)),
Arc::new(Float64Array::from(vec![-0.015_f64; n])),
];
match RecordBatch::try_new(schema, cols) {
Ok(b) => b,
Err(e) => panic!("build equity batch: {e}"),
}
}
fn positions_batch(n: usize) -> RecordBatch {
use std::sync::Arc;
use arrow_array::{ArrayRef, BooleanArray, Int32Array, Int64Array, StringArray};
use arrow_schema::{DataType, Field, Schema};
let steps: Vec<i32> = (0..n)
.map(|i| i32::try_from(i).unwrap_or(i32::MAX))
.collect();
let ts: Vec<i64> = steps
.iter()
.map(|&s| 1_700_000_000_000_000_000 + i64::from(s))
.collect();
let pos: Vec<i64> = steps.iter().map(|&s| i64::from(s)).collect();
let exit: Vec<Option<&str>> = steps
.iter()
.map(|&s| if s == 0 { Some("expiry") } else { None })
.collect();
let open_at_end: Vec<bool> = steps
.iter()
.map(|&s| usize::try_from(s).unwrap_or(usize::MAX) == n - 1)
.collect();
let schema = Arc::new(Schema::new(vec![
Field::new("step", DataType::Int32, false),
Field::new("ts_ns", DataType::Int64, false),
Field::new("position_id", DataType::Int64, false),
Field::new("trade_id", DataType::Int64, false),
Field::new("contract_id", DataType::Utf8, false),
Field::new("side", DataType::Utf8, false),
Field::new("quantity", DataType::Int32, false),
Field::new("avg_price_cents", DataType::Int64, false),
Field::new("mark_cents", DataType::Int64, false),
Field::new("unrealized_cents", DataType::Int64, false),
Field::new("stale_mark", DataType::Boolean, false),
Field::new("exit_reason", DataType::Utf8, true),
Field::new("open_at_end", DataType::Boolean, false),
]));
let cols: Vec<ArrayRef> = vec![
Arc::new(Int32Array::from(steps)),
Arc::new(Int64Array::from(ts)),
Arc::new(Int64Array::from(pos)),
Arc::new(Int64Array::from(vec![7_i64; n])),
Arc::new(StringArray::from(vec![
"v1:BTC:1735286400000000000:6000000:C";
n
])),
Arc::new(StringArray::from(vec!["short"; n])),
Arc::new(Int32Array::from(vec![1_i32; n])),
Arc::new(Int64Array::from(vec![12_000_i64; n])),
Arc::new(Int64Array::from(vec![11_800_i64; n])),
Arc::new(Int64Array::from(vec![200_i64; n])),
Arc::new(BooleanArray::from(vec![false; n])),
Arc::new(StringArray::from(exit)),
Arc::new(BooleanArray::from(open_at_end)),
];
match RecordBatch::try_new(schema, cols) {
Ok(b) => b,
Err(e) => panic!("build positions batch: {e}"),
}
}
fn greeks_batch(n: usize) -> RecordBatch {
use std::sync::Arc;
use arrow_array::{ArrayRef, Int32Array, Int64Array};
use arrow_schema::{DataType, Field, Schema};
let steps: Vec<i32> = (0..n)
.map(|i| i32::try_from(i).unwrap_or(i32::MAX))
.collect();
let ts: Vec<i64> = steps
.iter()
.map(|&s| 1_700_000_000_000_000_000 + i64::from(s))
.collect();
const CAPITAL: i64 = 1_000_000;
const BASE_TERMS: i64 = 40 - 120 + 15 + 10 - 30;
let equity_of = |s: i32| -> i64 { 988_500 + i64::from(s) };
let residual: Vec<i64> = steps
.iter()
.map(|&s| {
let step_pnl = if s == 0 {
equity_of(0) - CAPITAL
} else {
equity_of(s) - equity_of(s - 1)
};
step_pnl - BASE_TERMS
})
.collect();
let schema = Arc::new(Schema::new(vec![
Field::new("step", DataType::Int32, false),
Field::new("ts_ns", DataType::Int64, false),
Field::new("theta_pnl_cents", DataType::Int64, false),
Field::new("delta_pnl_cents", DataType::Int64, false),
Field::new("vega_pnl_cents", DataType::Int64, false),
Field::new("spread_capture_cents", DataType::Int64, false),
Field::new("fees_cents", DataType::Int64, false),
Field::new("residual_cents", DataType::Int64, false),
]));
let cols: Vec<ArrayRef> = vec![
Arc::new(Int32Array::from(steps)),
Arc::new(Int64Array::from(ts)),
Arc::new(Int64Array::from(vec![40_i64; n])),
Arc::new(Int64Array::from(vec![-120_i64; n])),
Arc::new(Int64Array::from(vec![15_i64; n])),
Arc::new(Int64Array::from(vec![10_i64; n])),
Arc::new(Int64Array::from(vec![30_i64; n])),
Arc::new(Int64Array::from(residual)),
];
match RecordBatch::try_new(schema, cols) {
Ok(b) => b,
Err(e) => panic!("build greeks batch: {e}"),
}
}
fn write_record_batch(path: &Path, batch: &RecordBatch, compressed: bool) {
use parquet::arrow::ArrowWriter;
use parquet::basic::{Compression, ZstdLevel};
use parquet::file::properties::WriterProperties;
let file = match std::fs::File::create(path) {
Ok(f) => f,
Err(e) => panic!("create {}: {e}", path.display()),
};
let props = if compressed {
Some(
WriterProperties::builder()
.set_compression(Compression::ZSTD(ZstdLevel::default()))
.build(),
)
} else {
None
};
let mut writer = match ArrowWriter::try_new(file, batch.schema(), props) {
Ok(w) => w,
Err(e) => panic!("arrow writer: {e}"),
};
if let Err(e) = writer.write(batch) {
panic!("write batch: {e}");
}
if let Err(e) = writer.close() {
panic!("close writer: {e}");
}
}
fn write_bomb_parquet(path: &Path, num_rows: usize) {
use std::sync::Arc;
use arrow_array::{ArrayRef, Int64Array, RecordBatch};
use arrow_schema::{DataType, Field, Schema};
use parquet::arrow::ArrowWriter;
use parquet::basic::{Compression, ZstdLevel};
use parquet::file::properties::WriterProperties;
let schema = Arc::new(Schema::new(vec![Field::new(
"step",
DataType::Int64,
false,
)]));
let col: ArrayRef = Arc::new(Int64Array::from(vec![0_i64; num_rows]));
let batch = match RecordBatch::try_new(Arc::clone(&schema), vec![col]) {
Ok(b) => b,
Err(e) => panic!("build bomb batch: {e}"),
};
let file = match std::fs::File::create(path) {
Ok(f) => f,
Err(e) => panic!("create {}: {e}", path.display()),
};
let props = WriterProperties::builder()
.set_compression(Compression::ZSTD(ZstdLevel::default()))
.set_dictionary_enabled(false)
.build();
let mut writer = match ArrowWriter::try_new(file, schema, Some(props)) {
Ok(w) => w,
Err(e) => panic!("arrow writer: {e}"),
};
if let Err(e) = writer.write(&batch) {
panic!("write bomb batch: {e}");
}
if let Err(e) = writer.close() {
panic!("close bomb writer: {e}");
}
}
fn write_full_bundle(
dir: &Path,
schema: &str,
file_rows: [usize; 4],
counts: [u64; 4],
extra: bool,
compressed: bool,
) {
if let Err(e) = std::fs::write(
dir.join(MANIFEST_FILE),
manifest_json_with(schema, counts, extra),
) {
panic!("write manifest: {e}");
}
let [(f0, _), (f1, _), (f2, _), (f3, _)] = TABLES;
let [r0, r1, r2, r3] = file_rows;
write_record_batch(&dir.join(f0), &fills_batch(r0), compressed);
write_record_batch(&dir.join(f1), &equity_batch(r1), compressed);
write_record_batch(&dir.join(f2), &positions_batch(r2), compressed);
write_record_batch(&dir.join(f3), &greeks_batch(r3), compressed);
}
#[track_caller]
fn open_ok(root: &Path, ceilings: ResourceCeilings) -> BundleReader {
match BundleReader::open_with_ceilings(root, ceilings) {
Ok(r) => r,
Err(e) => panic!("open_with_ceilings should succeed: {e}"),
}
}
fn dir_snapshot(dir: &Path) -> Vec<(String, u64)> {
let mut out = Vec::new();
let entries = match std::fs::read_dir(dir) {
Ok(e) => e,
Err(e) => panic!("read_dir: {e}"),
};
for entry in entries {
let entry = match entry {
Ok(e) => e,
Err(e) => panic!("dir entry: {e}"),
};
let name = entry.file_name().to_string_lossy().into_owned();
let len = match entry.metadata() {
Ok(m) => m.len(),
Err(e) => panic!("metadata: {e}"),
};
out.push((name, len));
}
out.sort();
out
}
#[test]
fn test_open_and_load_happy_path() {
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[4, 10, 8, 10],
[4, 10, 8, 10],
false,
false,
);
let reader = match BundleReader::open(dir.path()) {
Ok(r) => r,
Err(e) => panic!("a well-formed bundle should open: {e}"),
};
assert_eq!(reader.manifest().schema, SUPPORTED_SCHEMA);
let loaded = match reader.load() {
Ok(l) => l,
Err(e) => panic!("a well-formed bundle should load under default ceilings: {e}"),
};
assert_eq!(loaded.manifest.run_id, "run-abc123");
assert_eq!(loaded.fills.len(), 4);
assert_eq!(loaded.equity.len(), 10);
assert_eq!(loaded.positions.len(), 8);
assert_eq!(loaded.greeks.len(), 10);
assert!(
loaded
.fills
.is_sorted_by_key(|f| (f.step, f.order_id, f.fill_seq))
);
assert!(loaded.equity.is_sorted_by_key(|e| e.step));
assert!(
loaded
.positions
.is_sorted_by_key(|p| (p.step, p.position_id))
);
assert!(loaded.greeks.is_sorted_by_key(|g| g.step));
}
#[test]
fn test_open_missing_directory_is_io() {
match BundleReader::open("/no/such/bundle/dir-xyz") {
Err(BundleError::Io(_)) => {}
other => panic!("a missing bundle dir must be a typed Io error, got {other:?}"),
}
}
#[test]
fn test_open_missing_manifest_is_missing_table() {
let dir = temp_bundle_dir();
for (file, _key) in TABLES {
write_parquet(&dir.path().join(file), 1);
}
match BundleReader::open(dir.path()) {
Err(BundleError::MissingTable(t)) => assert_eq!(t, MANIFEST_FILE),
other => panic!("a missing manifest must be MissingTable, got {other:?}"),
}
}
#[test]
fn test_open_bad_schema_is_unsupported() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
"ironcondor.bundle.v2",
[1, 1, 1, 1],
[1, 1, 1, 1],
false,
);
match BundleReader::open(dir.path()) {
Err(BundleError::UnsupportedSchema(tag)) => assert_eq!(tag, "ironcondor.bundle.v2"),
other => panic!("a major-incompatible tag must be UnsupportedSchema, got {other:?}"),
}
}
#[test]
fn test_open_unknown_manifest_field_still_opens() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[1, 1, 1, 1],
[1, 1, 1, 1],
true,
);
assert!(
BundleReader::open(dir.path()).is_ok(),
"a newer-minor extra manifest field must still open (permissive)"
);
}
#[test]
fn test_schema_freeze_pin_is_exactly_v1() {
assert_eq!(SUPPORTED_SCHEMA, "ironcondor.bundle.v1");
}
#[test]
fn test_schema_gate_accepts_v1_rejects_fabricated_tags() {
let ok_dir = temp_bundle_dir();
write_bundle(
ok_dir.path(),
SUPPORTED_SCHEMA,
[1, 1, 1, 1],
[1, 1, 1, 1],
false,
);
assert!(
BundleReader::open(ok_dir.path()).is_ok(),
"the frozen tag `{SUPPORTED_SCHEMA}` must open"
);
for bogus in [
"ironcondor.bundle.v2",
"ironcondor.bundle.v0",
"IRONCONDOR.BUNDLE.V1",
"totally.bogus.format",
"",
] {
let dir = temp_bundle_dir();
write_bundle(dir.path(), bogus, [1, 1, 1, 1], [1, 1, 1, 1], false);
match BundleReader::open(dir.path()) {
Err(BundleError::UnsupportedSchema(tag)) => assert_eq!(
tag, bogus,
"the offending tag is echoed verbatim (short enough to be unclamped)"
),
other => {
panic!("a non-`v1` tag `{bogus}` must be UnsupportedSchema, got {other:?}")
}
}
}
}
#[test]
fn test_extra_optional_column_still_opens_and_loads() {
let dir = temp_bundle_dir();
if let Err(e) = std::fs::write(
dir.path().join(MANIFEST_FILE),
manifest_json_with(SUPPORTED_SCHEMA, [3, 3, 3, 3], false),
) {
panic!("write manifest: {e}");
}
let [(f0, _), (f1, _), (f2, _), (f3, _)] = TABLES;
write_record_batch(&dir.path().join(f0), &fills_batch_extra_column(3), false);
write_record_batch(&dir.path().join(f1), &equity_batch(3), false);
write_record_batch(&dir.path().join(f2), &positions_batch(3), false);
write_record_batch(&dir.path().join(f3), &greeks_batch(3), false);
let reader = match BundleReader::open(dir.path()) {
Ok(r) => r,
Err(e) => panic!("a newer-minor extra column must still open: {e}"),
};
match reader.load() {
Ok(loaded) => {
assert_eq!(loaded.fills.len(), 3, "the known columns still decode");
assert_eq!(loaded.equity.len(), 3);
}
Err(e) => panic!("a newer-minor extra column must still load: {e}"),
}
}
#[test]
fn test_open_missing_table_is_missing_table() {
let dir = temp_bundle_dir();
if let Err(e) = std::fs::write(
dir.path().join(MANIFEST_FILE),
manifest_json_with(SUPPORTED_SCHEMA, [1, 1, 1, 1], false),
) {
panic!("write manifest: {e}");
}
write_parquet(&dir.path().join("fills.parquet"), 1);
write_parquet(&dir.path().join("equity_curve.parquet"), 1);
write_parquet(&dir.path().join("positions.parquet"), 1);
match BundleReader::open(dir.path()) {
Err(BundleError::MissingTable(t)) => assert_eq!(t, "greeks_attribution.parquet"),
other => panic!("a missing table must be MissingTable, got {other:?}"),
}
}
#[test]
fn test_open_wrong_keyed_row_counts_is_invariant() {
let dir = temp_bundle_dir();
let manifest = r#"{ "schema": "ironcondor.bundle.v1", "run_id": "r",
"created_utc": "t", "code_version": "v", "lockfile_sha256": "s", "seed": 1,
"config": {}, "strategy": {}, "data_source": {}, "metrics": {},
"row_counts": { "fills": 1, "equity_curve": 1, "positions": 1 } }"#;
if let Err(e) = std::fs::write(dir.path().join(MANIFEST_FILE), manifest) {
panic!("write manifest: {e}");
}
for (file, _key) in TABLES {
write_parquet(&dir.path().join(file), 1);
}
match BundleReader::open(dir.path()) {
Err(BundleError::Invariant(_)) => {}
other => panic!("a wrong-keyed row_counts must be Invariant, got {other:?}"),
}
}
#[test]
fn test_ceiling1_oversized_file_is_too_large_pre_open() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[16, 16, 16, 16],
[16, 16, 16, 16],
false,
);
let reader = open_ok(
dir.path(),
ResourceCeilings {
max_table_bytes: 1,
..ResourceCeilings::default()
},
);
match reader.load() {
Err(BundleError::TooLarge(_)) => {}
other => panic!("an oversized file must be TooLarge, got {other:?}"),
}
}
#[test]
fn test_ceiling2_footer_rowcount_is_too_large_pre_decode() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[10, 10, 10, 10],
[10, 10, 10, 10],
false,
);
let reader = open_ok(
dir.path(),
ResourceCeilings {
max_table_rows: 5,
..ResourceCeilings::default()
},
);
match reader.load() {
Err(BundleError::TooLarge(_)) => {}
other => panic!("a footer row count over the ceiling must be TooLarge, got {other:?}"),
}
}
#[test]
fn test_ceiling2_footer_rowcount_mismatch_is_invariant() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[7, 10, 8, 10],
[10, 10, 8, 10],
false,
);
let reader = open_ok(dir.path(), ResourceCeilings::default());
match reader.load() {
Err(BundleError::Invariant(_)) => {}
other => panic!("a footer/row_counts mismatch must be Invariant, got {other:?}"),
}
}
#[test]
fn test_ceiling3_tiny_working_set_is_too_large() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[64, 64, 64, 64],
[64, 64, 64, 64],
false,
);
let reader = open_ok(
dir.path(),
ResourceCeilings {
max_working_set: 8,
max_batch_bytes: 8,
..ResourceCeilings::default()
},
);
match reader.load() {
Err(BundleError::TooLarge(_)) => {}
other => panic!("a working set over the ceiling must be TooLarge, got {other:?}"),
}
}
#[test]
fn test_working_set_budget_stops_before_crossing_ceiling() {
let ceilings = ResourceCeilings {
max_working_set: 100,
max_batch_bytes: 60,
..ResourceCeilings::default()
};
let mut budget = WorkingSetBudget::new(&ceilings);
assert!(budget.account(40).is_ok());
assert!(budget.account(40).is_ok());
assert_eq!(budget.used(), 80);
match budget.account(40) {
Err(BundleError::TooLarge(_)) => {}
other => {
panic!("the batch that would exceed the ceiling must be TooLarge, got {other:?}")
}
}
assert!(
budget.used() < ceilings.max_working_set,
"cumulative bytes must stay strictly under the ceiling even on the reject path"
);
assert_eq!(budget.used(), 80, "a rejected batch is never committed");
}
#[test]
fn test_working_set_budget_rejects_oversized_batch() {
let ceilings = ResourceCeilings {
max_working_set: 1_000,
max_batch_bytes: 50,
..ResourceCeilings::default()
};
let mut budget = WorkingSetBudget::new(&ceilings);
match budget.account(51) {
Err(BundleError::TooLarge(_)) => {}
other => panic!("a batch over the per-batch cap must be TooLarge, got {other:?}"),
}
assert_eq!(budget.used(), 0, "a rejected batch is never committed");
}
fn write_bomb_bundle(dir: &Path, fills_rows: usize) {
let counts = [u64::try_from(fills_rows).unwrap_or(u64::MAX), 1, 1, 1];
if let Err(e) = std::fs::write(
dir.join(MANIFEST_FILE),
manifest_json_with(SUPPORTED_SCHEMA, counts, false),
) {
panic!("write manifest: {e}");
}
write_bomb_parquet(&dir.join("fills.parquet"), fills_rows);
write_parquet(&dir.join("equity_curve.parquet"), 1);
write_parquet(&dir.join("positions.parquet"), 1);
write_parquet(&dir.join("greeks_attribution.parquet"), 1);
}
#[track_caller]
fn scan_fills_probing_decoder(reader: &BundleReader) -> (Result<(), BundleError>, bool, u64) {
use std::cell::Cell;
let mut budget = WorkingSetBudget::new(reader.ceilings());
let invoked = Cell::new(false);
let result = reader.scan_table(
"fills.parquet",
"fills",
&mut budget,
&|| false,
&mut |_batch| {
invoked.set(true);
Ok(0)
},
);
(result, invoked.get(), budget.used())
}
#[test]
fn test_decompression_bomb_rejects_before_invoking_the_decoder() {
let dir = temp_bundle_dir();
write_bomb_bundle(dir.path(), 8_000);
let reader = open_ok(dir.path(), ResourceCeilings::default());
let (result, invoked, used) = scan_fills_probing_decoder(&reader);
match result {
Err(BundleError::TooLarge(detail)) => assert!(
detail.contains("decompression bomb"),
"the bomb must be rejected by the pre-decode bomb check: {detail}"
),
other => panic!("a decompression bomb must be TooLarge, got {other:?}"),
}
assert!(
!invoked,
"the bomb must be rejected BEFORE the decoder materializes any batch"
);
assert_eq!(
used, 0,
"no working set is committed on the pre-decode reject"
);
}
#[test]
fn test_oversized_footer_rows_reject_before_invoking_the_decoder() {
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[8, 8, 8, 8],
[8, 8, 8, 8],
false,
false,
);
let reader = open_ok(
dir.path(),
ResourceCeilings {
max_table_rows: 4,
..ResourceCeilings::default()
},
);
let (result, invoked, used) = scan_fills_probing_decoder(&reader);
assert!(
matches!(result, Err(BundleError::TooLarge(_))),
"an over-ceiling footer must be TooLarge, got {result:?}"
);
assert!(!invoked, "ceiling 2 rejects before the decoder runs");
assert_eq!(
used, 0,
"no working set is committed on the pre-decode reject"
);
}
#[test]
fn test_rowcount_lie_rejects_before_invoking_the_decoder() {
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[4, 10, 8, 10],
[10, 10, 8, 10],
false,
false,
);
let reader = open_ok(dir.path(), ResourceCeilings::default());
let (result, invoked, used) = scan_fills_probing_decoder(&reader);
assert!(
matches!(result, Err(BundleError::Invariant(_))),
"a footer/row_counts lie must be Invariant, got {result:?}"
);
assert!(
!invoked,
"the row_counts cross-check rejects before the decoder runs — the hint \
never sizes an allocation"
);
assert_eq!(
used, 0,
"no working set is committed on the pre-decode reject"
);
}
#[test]
fn test_checked_footer_conversion_rejects_negative() {
match footer_i64_to_u64(-1, "footer row count") {
Err(BundleError::TooLarge(_)) => {}
other => panic!("a negative footer value must be TooLarge, got {other:?}"),
}
match footer_i64_to_u64(i64::MIN, "footer size") {
Err(BundleError::TooLarge(_)) => {}
other => panic!("i64::MIN must be TooLarge, got {other:?}"),
}
assert_eq!(footer_i64_to_u64(42, "x").ok(), Some(42));
}
#[test]
fn test_apply_overhead_overflow_is_too_large() {
match apply_overhead(u64::MAX, DECODED_OVERHEAD_PERMILLE) {
Err(BundleError::TooLarge(_)) => {}
other => panic!("an overflowing overhead multiply must be TooLarge, got {other:?}"),
}
assert_eq!(apply_overhead(1_000, 1_500).ok(), Some(1_500));
}
#[test]
fn test_load_pre_cancelled_returns_promptly() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[4, 10, 8, 10],
[4, 10, 8, 10],
false,
);
let reader = open_ok(dir.path(), ResourceCeilings::default());
match reader.load_cancellable(&|| true) {
Err(BundleError::Cancelled) => {}
other => panic!("a pre-cancelled load must be Cancelled, got {other:?}"),
}
}
#[test]
fn test_load_cancels_mid_decode_at_batch_boundary() {
use std::cell::Cell;
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[8, 8, 8, 8],
[8, 8, 8, 8],
false,
false,
);
let reader = open_ok(
dir.path(),
ResourceCeilings {
max_batch_rows: 1,
..ResourceCeilings::default()
},
);
let polls = Cell::new(0_u32);
let probe = || {
let n = polls.get().saturating_add(1);
polls.set(n);
n > 3
};
match reader.load_cancellable(&probe) {
Err(BundleError::Cancelled) => {}
other => panic!("a mid-decode cancellation must be Cancelled, got {other:?}"),
}
assert!(
polls.get() <= 6,
"cancellation must abort promptly, got {} polls",
polls.get()
);
}
#[test]
fn test_load_is_read_only() {
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[4, 10, 8, 10],
[4, 10, 8, 10],
false,
false,
);
let before = dir_snapshot(dir.path());
let reader = open_ok(dir.path(), ResourceCeilings::default());
let _ = reader.load();
let after = dir_snapshot(dir.path());
assert_eq!(
before, after,
"load must not add, remove, or modify any bundle file"
);
}
#[test]
fn test_resource_ceilings_default_validates() {
assert!(
ResourceCeilings::default().validate().is_ok(),
"the documented defaults must be valid"
);
}
#[test]
fn test_resource_ceilings_out_of_range_is_config_error() {
let zero = ResourceCeilings {
max_table_bytes: 0,
..ResourceCeilings::default()
};
match zero.validate() {
Err(ConfigError::InvalidValue { field, .. }) => {
assert_eq!(field, "replay.max_table_bytes")
}
other => panic!("a zero ceiling must be a ConfigError, got {other:?}"),
}
let bad_overhead = ResourceCeilings {
decoded_overhead_permille: 900,
..ResourceCeilings::default()
};
assert!(
matches!(
bad_overhead.validate(),
Err(ConfigError::InvalidValue { .. })
),
"an overhead below 1.0x must be rejected"
);
let bad_batch = ResourceCeilings {
max_batch_bytes: MAX_WORKING_SET + 1,
..ResourceCeilings::default()
};
assert!(
matches!(bad_batch.validate(), Err(ConfigError::InvalidValue { .. })),
"a per-batch cap above the working-set ceiling must be rejected"
);
let bad_manifest = ResourceCeilings {
max_manifest_bytes: 0,
..ResourceCeilings::default()
};
match bad_manifest.validate() {
Err(ConfigError::InvalidValue { field, .. }) => {
assert_eq!(field, "replay.max_manifest_bytes")
}
other => panic!("a zero manifest ceiling must be a ConfigError, got {other:?}"),
}
}
#[test]
fn test_ceiling_oversized_manifest_rejects_pre_read() {
let dir = temp_bundle_dir();
for (file, _key) in TABLES {
write_parquet(&dir.path().join(file), 1);
}
let junk = "x".repeat(4096);
if let Err(e) = std::fs::write(dir.path().join(MANIFEST_FILE), &junk) {
panic!("write manifest: {e}");
}
let ceilings = ResourceCeilings {
max_manifest_bytes: 64,
..ResourceCeilings::default()
};
match BundleReader::open_with_ceilings(dir.path(), ceilings) {
Err(BundleError::TooLarge(detail)) => assert!(
detail.contains(MANIFEST_FILE),
"the oversized artifact must be named the manifest: {detail}"
),
other => {
panic!("an oversized manifest must be TooLarge pre-read, got {other:?}")
}
}
}
#[test]
fn test_normal_manifest_opens_under_a_finite_manifest_ceiling() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[1, 1, 1, 1],
[1, 1, 1, 1],
false,
);
let ceilings = ResourceCeilings {
max_manifest_bytes: 8 * 1024,
..ResourceCeilings::default()
};
assert!(
BundleReader::open_with_ceilings(dir.path(), ceilings).is_ok(),
"a normal manifest must open under a finite manifest ceiling"
);
}
#[test]
fn test_open_with_invalid_ceilings_is_typed_error_not_silent_open() {
let dir = temp_bundle_dir();
write_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[1, 1, 1, 1],
[1, 1, 1, 1],
false,
);
let bad = ResourceCeilings {
max_batch_rows: 0,
..ResourceCeilings::default()
};
match BundleReader::open_with_ceilings(dir.path(), bad) {
Err(BundleError::Config(ConfigError::InvalidValue { field, .. })) => {
assert_eq!(field, "replay.max_batch_rows")
}
other => panic!(
"an invalid ceiling on the enforcement path must be a typed \
BundleError::Config, never a silent open, got {other:?}"
),
}
}
#[test]
fn test_unsupported_schema_tag_is_clamped() {
let junk = "z".repeat(10 * 1024);
let err = unsupported_schema(junk);
match &err {
BundleError::UnsupportedSchema(tag) => assert!(
tag.chars().count() <= MAX_SCHEMA_TAG_CHARS + 1,
"the tag must be clamped (<= {} chars + ellipsis), got {}",
MAX_SCHEMA_TAG_CHARS,
tag.chars().count()
),
other => panic!("expected UnsupportedSchema, got {other:?}"),
}
let rendered = err.to_string();
let bound = "unsupported schema: ".chars().count() + MAX_SCHEMA_TAG_CHARS + 1;
assert!(
rendered.chars().count() <= bound,
"rendered message must be bounded (<= {bound}), got {}",
rendered.chars().count()
);
assert!(rendered.starts_with("unsupported schema: "));
}
#[test]
fn test_catch_decode_panic_maps_an_upstream_panic_to_a_typed_parquet_error() {
match catch_decode_panic::<()>("greeks_attribution.parquet", || {
panic!("arrow-ipc get_data_type")
}) {
Err(BundleError::Parquet(detail)) => {
assert!(
detail.starts_with("greeks_attribution.parquet"),
"the typed error names the table: {detail}"
);
assert!(
detail.contains("panicked") && !detail.contains("get_data_type"),
"the payload is dropped, never interpolated: {detail}"
);
}
other => {
panic!("a contained decoder panic must be a typed Parquet error, got {other:?}")
}
}
}
#[test]
fn test_catch_decode_panic_passes_a_normal_result_through_unchanged() {
assert!(matches!(
catch_decode_panic("fills.parquet", || Ok::<u8, BundleError>(7)),
Ok(7)
));
match catch_decode_panic::<()>("fills.parquet", || Err(BundleError::Cancelled)) {
Err(BundleError::Cancelled) => {}
other => panic!("a normal `Err` must pass through, got {other:?}"),
}
}
#[test]
fn test_clamp_schema_tag_never_panics_on_multibyte_utf8() {
let multibyte = "🚀".repeat(200);
let clamped = clamp_schema_tag(multibyte);
assert!(clamped.chars().count() <= MAX_SCHEMA_TAG_CHARS + 1);
assert_eq!(
clamp_schema_tag("ironcondor.bundle.v9".to_owned()),
"ironcondor.bundle.v9"
);
}
#[test]
fn test_load_decodes_zstd_compressed_tables() {
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[3, 3, 3, 3],
[3, 3, 3, 3],
false,
true,
);
let reader = open_ok(dir.path(), ResourceCeilings::default());
match reader.load() {
Ok(loaded) => {
assert_eq!(loaded.manifest.schema, SUPPORTED_SCHEMA);
assert_eq!(loaded.fills.len(), 3);
assert_eq!(loaded.equity.len(), 3);
assert_eq!(loaded.positions.len(), 3);
assert_eq!(loaded.greeks.len(), 3);
}
Err(e) => {
panic!("a ZSTD-compressed bundle must decode under the enabled codec: {e}")
}
}
}
#[test]
fn test_load_returns_typed_rows_in_file_order() {
let dir = temp_bundle_dir();
write_full_bundle(
dir.path(),
SUPPORTED_SCHEMA,
[3, 3, 3, 3],
[3, 3, 3, 3],
false,
false,
);
let reader = open_ok(dir.path(), ResourceCeilings::default());
let loaded = match reader.load() {
Ok(l) => l,
Err(e) => panic!("a well-formed full bundle should load: {e}"),
};
let base_ts = 1_700_000_000_000_000_000_i64;
let cid = "v1:BTC:1735286400000000000:6000000:C";
let expected_fills: Vec<Fill> = (0..3_u32)
.map(|s| Fill {
step: s,
ts_ns: base_ts + i64::from(s),
strategy_run_id: "run-abc123".to_owned(),
trade_id: 100 + u64::from(s),
position_id: 200 + u64::from(s),
order_id: u64::from(s),
fill_seq: 0,
underlying: "BTC".to_owned(),
expiration_ns: 1_735_286_400_000_000_000,
contract_id: cid.to_owned(),
strike_cents: 6_000_000,
style: OptionStyle::Call,
side: PositionSide::Long,
quantity: 1,
price_cents: 12_500,
fees_cents: 30,
slippage_cents: -15,
mode: ExecMode::Realistic,
})
.collect();
let expected_equity: Vec<EquityPoint> = (0..3_u32)
.map(|s| EquityPoint {
step: s,
ts_ns: base_ts + i64::from(s),
cash_cents: 990_000 + i64::from(s),
position_value_cents: -1_500,
equity_cents: 988_500 + i64::from(s),
drawdown: -0.015,
})
.collect();
let expected_positions: Vec<PositionRow> = (0..3_u32)
.map(|s| PositionRow {
step: s,
ts_ns: base_ts + i64::from(s),
position_id: u64::from(s),
trade_id: 7,
contract_id: cid.to_owned(),
side: PositionSide::Short,
quantity: 1,
avg_price_cents: 12_000,
mark_cents: 11_800,
unrealized_cents: 200,
stale_mark: false,
exit_reason: if s == 0 {
Some("expiry".to_owned())
} else {
None
},
open_at_end: s == 2,
})
.collect();
let expected_greeks: Vec<GreeksAttribution> = (0..3_u32)
.map(|s| GreeksAttribution {
step: s,
ts_ns: base_ts + i64::from(s),
theta_pnl_cents: 40,
delta_pnl_cents: -120,
vega_pnl_cents: 15,
spread_capture_cents: 10,
fees_cents: 30,
residual_cents: if s == 0 { -11_415 } else { 86 },
})
.collect();
assert_eq!(
loaded.fills, expected_fills,
"fills must round-trip verbatim in file order"
);
assert_eq!(
loaded.equity, expected_equity,
"equity must round-trip verbatim in file order"
);
assert_eq!(
loaded.positions, expected_positions,
"positions must round-trip verbatim in file order"
);
assert_eq!(
loaded.greeks, expected_greeks,
"greeks must round-trip verbatim in file order"
);
}
fn fills_batch_with_dict_contract_id(n: usize, value: &str) -> RecordBatch {
use std::sync::Arc;
use arrow_array::types::Int32Type;
use arrow_array::{ArrayRef, DictionaryArray, Int32Array, StringArray};
use arrow_schema::{DataType, Field, Schema};
let base = fills_batch(n);
let keys = Int32Array::from(vec![0_i32; n]);
let values = Arc::new(StringArray::from(vec![value])) as ArrayRef;
let dict: ArrayRef = match DictionaryArray::<Int32Type>::try_new(keys, values) {
Ok(d) => Arc::new(d),
Err(e) => panic!("build dictionary array: {e}"),
};
let dict_type = DataType::Dictionary(Box::new(DataType::Int32), Box::new(DataType::Utf8));
let mut fields: Vec<Field> = base
.schema()
.fields()
.iter()
.map(|f| f.as_ref().clone())
.collect();
let mut cols: Vec<ArrayRef> = base.columns().to_vec();
let idx = match base.schema().index_of("contract_id") {
Ok(i) => i,
Err(e) => panic!("fills batch missing contract_id: {e}"),
};
if let Some(slot) = fields.get_mut(idx) {
*slot = Field::new("contract_id", dict_type, false);
}
if let Some(slot) = cols.get_mut(idx) {
*slot = dict;
}
match RecordBatch::try_new(Arc::new(Schema::new(fields)), cols) {
Ok(b) => b,
Err(e) => panic!("build dict fills batch: {e}"),
}
}
#[test]
fn test_working_set_budget_trips_on_retained_dictionary_bytes() {
let n = 64;
let long = "L".repeat(1024);
let batch = fills_batch_with_dict_contract_id(n, &long);
let arrow_bytes = u64::try_from(batch.get_array_memory_size()).unwrap_or(u64::MAX);
let batch_estimate = match apply_overhead(arrow_bytes, DECODED_OVERHEAD_PERMILLE) {
Ok(v) => v,
Err(e) => panic!("overhead multiply: {e}"),
};
let mut out = Vec::new();
let retained = match tables::read_fills(&batch, &mut out) {
Ok(r) => r,
Err(e) => panic!("read_fills should succeed: {e}"),
};
assert_eq!(out.len(), n);
assert!(
retained > batch_estimate,
"the per-row copied dictionary string must make retained {retained} exceed the \
Arrow batch estimate {batch_estimate}"
);
let ceilings = ResourceCeilings {
max_working_set: batch_estimate + retained - 1,
max_batch_bytes: retained,
..ResourceCeilings::default()
};
let mut budget = WorkingSetBudget::new(&ceilings);
match budget.account(batch_estimate) {
Ok(()) => {}
Err(e) => panic!("the transient batch estimate alone must fit under the ceiling: {e}"),
}
match budget.account(retained) {
Err(BundleError::TooLarge(_)) => {}
other => panic!("the retained bytes must trip the working-set ceiling, got {other:?}"),
}
assert!(
budget.used() < ceilings.max_working_set,
"a rejected retained batch is never committed; used stays under the ceiling"
);
}
}