opendeviationbar-streaming 13.75.0

Real-time streaming engine for open deviation bar processing
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
690
691
692
693
694
695
696
697
698
699
700
701
702
703
704
705
706
707
708
709
710
711
712
713
714
715
716
717
718
719
720
721
722
723
724
725
726
727
728
729
730
731
732
733
734
735
736
737
738
739
740
741
742
743
744
745
746
747
748
749
750
751
752
753
754
755
756
757
758
759
760
761
762
763
764
765
766
767
768
769
770
771
772
773
774
775
776
777
778
779
780
781
782
783
784
785
786
787
788
789
790
791
792
793
794
795
//! Parquet dead-letter writer for failed ClickHouse flush batches.
//!
//! Phase 35: Dead-letter resilience -- no bar is ever silently lost.
//!
//! When `do_flush()` fails after all retries, or when `on_bar()` gets
//! backpressure (`TrySendError::Full`), bars are serialized to Parquet
//! files in `/tmp/opendeviationbar-dead-letter/` for later replay by
//! Python Charon (`kintsugi/dead_letter.py`).
//!
//! # File naming
//!
//! `{symbol}_{threshold}_{unix_timestamp}.parquet`
//!
//! # Schema compatibility
//!
//! Column names match `CORE_COLUMNS` from `row.rs` exactly. Python Charon
//! reads these via `polars.read_parquet()` and passes to `store_bars_batch()`.

use std::path::{Path, PathBuf};

use arrow_array::builder::{
    BooleanBuilder, Float64Builder, Int64Builder, StringBuilder, UInt32Builder, UInt8Builder,
};
use arrow_array::RecordBatch;
use arrow_schema::{DataType, Field, Schema};

use super::row::ClickHouseBarRow;

/// Directory for dead-letter Parquet files (matches Python `_DEAD_LETTER_DIR`).
pub const DEAD_LETTER_DIR: &str = "/tmp/opendeviationbar-dead-letter";

/// Errors from dead-letter operations.
#[derive(Debug)]
pub enum DeadLetterError {
    /// I/O error (directory creation, file write).
    Io(std::io::Error),
    /// Arrow/Parquet error.
    Arrow(String),
}

impl std::fmt::Display for DeadLetterError {
    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
        match self {
            DeadLetterError::Io(e) => write!(f, "dead-letter I/O error: {e}"),
            DeadLetterError::Arrow(e) => write!(f, "dead-letter Arrow error: {e}"),
        }
    }
}

impl From<std::io::Error> for DeadLetterError {
    fn from(e: std::io::Error) -> Self {
        DeadLetterError::Io(e)
    }
}

impl From<arrow_schema::ArrowError> for DeadLetterError {
    fn from(e: arrow_schema::ArrowError) -> Self {
        DeadLetterError::Arrow(e.to_string())
    }
}

impl From<parquet::errors::ParquetError> for DeadLetterError {
    fn from(e: parquet::errors::ParquetError) -> Self {
        DeadLetterError::Arrow(format!("Parquet error: {e}"))
    }
}

/// Build the Arrow schema for dead-letter Parquet files.
///
/// Column names match `CORE_COLUMNS` exactly. Nullable flags match the
/// `Option<T>` fields in `ClickHouseBarRow`.
pub fn dead_letter_schema() -> Schema {
    Schema::new(vec![
        // -- Cache key components --
        Field::new("symbol", DataType::Utf8, false),
        Field::new("threshold_decimal_bps", DataType::UInt32, false),
        // -- OHLCV --
        Field::new("close_time_us", DataType::Int64, false),
        Field::new("open_time_us", DataType::Int64, false),
        Field::new("open", DataType::Float64, false),
        Field::new("high", DataType::Float64, false),
        Field::new("low", DataType::Float64, false),
        Field::new("close", DataType::Float64, false),
        Field::new("volume", DataType::Float64, false),
        // -- Market microstructure --
        Field::new("vwap", DataType::Float64, false),
        Field::new("buy_volume", DataType::Float64, false),
        Field::new("sell_volume", DataType::Float64, false),
        Field::new("individual_trade_count", DataType::UInt32, false),
        Field::new("agg_record_count", DataType::UInt32, false),
        // -- Microstructure features --
        Field::new("duration_us", DataType::Int64, false),
        Field::new("ofi", DataType::Float64, false),
        Field::new("vwap_close_deviation", DataType::Float64, false),
        Field::new("price_impact", DataType::Float64, false),
        Field::new("kyle_lambda_proxy", DataType::Float64, false),
        Field::new("trade_intensity", DataType::Float64, false),
        Field::new("volume_per_trade", DataType::Float64, false),
        Field::new("aggression_ratio", DataType::Float64, false),
        Field::new("aggregation_density", DataType::Float64, false),
        Field::new("turnover_imbalance", DataType::Float64, false),
        // -- Ouroboros --
        Field::new("ouroboros_mode", DataType::Utf8, false),
        // -- Exchange session flags --
        Field::new("exchange_session_sydney", DataType::UInt8, false),
        Field::new("exchange_session_tokyo", DataType::UInt8, false),
        Field::new("exchange_session_london", DataType::UInt8, false),
        Field::new("exchange_session_newyork", DataType::UInt8, false),
        // -- Inter-bar features (all nullable) --
        Field::new("lookback_trade_count", DataType::UInt32, true),
        Field::new("lookback_ofi", DataType::Float64, true),
        Field::new("lookback_duration_us", DataType::Int64, true),
        Field::new("lookback_intensity", DataType::Float64, true),
        Field::new("lookback_vwap_raw", DataType::Float64, true),
        Field::new("lookback_vwap_position", DataType::Float64, true),
        Field::new("lookback_count_imbalance", DataType::Float64, true),
        Field::new("lookback_kyle_lambda", DataType::Float64, true),
        Field::new("lookback_burstiness", DataType::Float64, true),
        Field::new("lookback_volume_skew", DataType::Float64, true),
        Field::new("lookback_volume_kurt", DataType::Float64, true),
        Field::new("lookback_price_range", DataType::Float64, true),
        Field::new("lookback_kaufman_er", DataType::Float64, true),
        Field::new("lookback_garman_klass_vol", DataType::Float64, true),
        Field::new("lookback_hurst", DataType::Float64, true),
        Field::new("lookback_permutation_entropy", DataType::Float64, true),
        // -- Intra-bar features (all nullable) --
        Field::new("intra_bull_epoch_density", DataType::Float64, true),
        Field::new("intra_bear_epoch_density", DataType::Float64, true),
        Field::new("intra_bull_excess_gain", DataType::Float64, true),
        Field::new("intra_bear_excess_gain", DataType::Float64, true),
        Field::new("intra_bull_cv", DataType::Float64, true),
        Field::new("intra_bear_cv", DataType::Float64, true),
        Field::new("intra_max_drawdown", DataType::Float64, true),
        Field::new("intra_max_runup", DataType::Float64, true),
        Field::new("intra_trade_count", DataType::UInt32, true),
        Field::new("intra_ofi", DataType::Float64, true),
        Field::new("intra_duration_us", DataType::Int64, true),
        Field::new("intra_intensity", DataType::Float64, true),
        Field::new("intra_vwap_position", DataType::Float64, true),
        Field::new("intra_count_imbalance", DataType::Float64, true),
        Field::new("intra_kyle_lambda", DataType::Float64, true),
        Field::new("intra_burstiness", DataType::Float64, true),
        Field::new("intra_volume_skew", DataType::Float64, true),
        Field::new("intra_volume_kurt", DataType::Float64, true),
        Field::new("intra_kaufman_er", DataType::Float64, true),
        Field::new("intra_garman_klass_vol", DataType::Float64, true),
        Field::new("intra_hurst", DataType::Float64, true),
        Field::new("intra_permutation_entropy", DataType::Float64, true),
        // -- Gap awareness --
        Field::new("has_gap", DataType::Boolean, false),
        Field::new("gap_trade_count", DataType::Int64, false),
        Field::new("max_gap_duration_us", DataType::Int64, false),
        Field::new("is_exchange_gap", DataType::Boolean, false),
        // -- Trade ID range --
        Field::new("first_agg_trade_id", DataType::Int64, false),
        Field::new("last_agg_trade_id", DataType::Int64, false),
        // -- Bar flags --
        Field::new("is_orphan", DataType::UInt8, false),
        // -- Cache metadata --
        Field::new("cache_key", DataType::Utf8, false),
        Field::new("opendeviationbar_version", DataType::Utf8, false),
        Field::new("source_start_ts", DataType::Int64, false),
        Field::new("source_end_ts", DataType::Int64, false),
    ])
}

/// Public entry point for `rows_to_record_batch` (used by niffler tests).
pub fn rows_to_record_batch_public(
    rows: &[ClickHouseBarRow],
    schema: &Schema,
) -> Result<RecordBatch, DeadLetterError> {
    rows_to_record_batch(rows, schema)
}

/// Convert a slice of `ClickHouseBarRow` to an Arrow `RecordBatch`.
fn rows_to_record_batch(
    rows: &[ClickHouseBarRow],
    schema: &Schema,
) -> Result<RecordBatch, DeadLetterError> {
    let n = rows.len();

    // Non-nullable string builders
    let mut symbol_b = StringBuilder::with_capacity(n, n * 10);
    let mut ouroboros_mode_b = StringBuilder::with_capacity(n, n * 4);
    let mut cache_key_b = StringBuilder::with_capacity(n, n * 32);
    let mut version_b = StringBuilder::with_capacity(n, n * 10);

    // Non-nullable numeric builders
    let mut threshold_b = UInt32Builder::with_capacity(n);
    let mut close_time_b = Int64Builder::with_capacity(n);
    let mut open_time_b = Int64Builder::with_capacity(n);
    let mut open_b = Float64Builder::with_capacity(n);
    let mut high_b = Float64Builder::with_capacity(n);
    let mut low_b = Float64Builder::with_capacity(n);
    let mut close_b = Float64Builder::with_capacity(n);
    let mut volume_b = Float64Builder::with_capacity(n);
    let mut vwap_b = Float64Builder::with_capacity(n);
    let mut buy_vol_b = Float64Builder::with_capacity(n);
    let mut sell_vol_b = Float64Builder::with_capacity(n);
    let mut ind_count_b = UInt32Builder::with_capacity(n);
    let mut agg_count_b = UInt32Builder::with_capacity(n);
    let mut duration_b = Int64Builder::with_capacity(n);
    let mut ofi_b = Float64Builder::with_capacity(n);
    let mut vwap_dev_b = Float64Builder::with_capacity(n);
    let mut price_impact_b = Float64Builder::with_capacity(n);
    let mut kyle_b = Float64Builder::with_capacity(n);
    let mut intensity_b = Float64Builder::with_capacity(n);
    let mut vol_per_trade_b = Float64Builder::with_capacity(n);
    let mut aggression_b = Float64Builder::with_capacity(n);
    let mut agg_density_b = Float64Builder::with_capacity(n);
    let mut turnover_b = Float64Builder::with_capacity(n);

    // Exchange session flags (UInt8)
    let mut sess_sydney_b = UInt8Builder::with_capacity(n);
    let mut sess_tokyo_b = UInt8Builder::with_capacity(n);
    let mut sess_london_b = UInt8Builder::with_capacity(n);
    let mut sess_ny_b = UInt8Builder::with_capacity(n);

    // Inter-bar nullable builders
    let mut lb_trade_count_b = UInt32Builder::with_capacity(n);
    let mut lb_ofi_b = Float64Builder::with_capacity(n);
    let mut lb_duration_b = Int64Builder::with_capacity(n);
    let mut lb_intensity_b = Float64Builder::with_capacity(n);
    let mut lb_vwap_raw_b = Float64Builder::with_capacity(n);
    let mut lb_vwap_pos_b = Float64Builder::with_capacity(n);
    let mut lb_count_imb_b = Float64Builder::with_capacity(n);
    let mut lb_kyle_b = Float64Builder::with_capacity(n);
    let mut lb_burst_b = Float64Builder::with_capacity(n);
    let mut lb_vol_skew_b = Float64Builder::with_capacity(n);
    let mut lb_vol_kurt_b = Float64Builder::with_capacity(n);
    let mut lb_price_range_b = Float64Builder::with_capacity(n);
    let mut lb_kaufman_b = Float64Builder::with_capacity(n);
    let mut lb_gk_vol_b = Float64Builder::with_capacity(n);
    let mut lb_hurst_b = Float64Builder::with_capacity(n);
    let mut lb_perm_ent_b = Float64Builder::with_capacity(n);

    // Intra-bar nullable builders
    let mut intra_bull_density_b = Float64Builder::with_capacity(n);
    let mut intra_bear_density_b = Float64Builder::with_capacity(n);
    let mut intra_bull_gain_b = Float64Builder::with_capacity(n);
    let mut intra_bear_gain_b = Float64Builder::with_capacity(n);
    let mut intra_bull_cv_b = Float64Builder::with_capacity(n);
    let mut intra_bear_cv_b = Float64Builder::with_capacity(n);
    let mut intra_max_dd_b = Float64Builder::with_capacity(n);
    let mut intra_max_ru_b = Float64Builder::with_capacity(n);
    let mut intra_trade_count_b = UInt32Builder::with_capacity(n);
    let mut intra_ofi_b = Float64Builder::with_capacity(n);
    let mut intra_duration_b = Int64Builder::with_capacity(n);
    let mut intra_intensity_b = Float64Builder::with_capacity(n);
    let mut intra_vwap_pos_b = Float64Builder::with_capacity(n);
    let mut intra_count_imb_b = Float64Builder::with_capacity(n);
    let mut intra_kyle_b = Float64Builder::with_capacity(n);
    let mut intra_burst_b = Float64Builder::with_capacity(n);
    let mut intra_vol_skew_b = Float64Builder::with_capacity(n);
    let mut intra_vol_kurt_b = Float64Builder::with_capacity(n);
    let mut intra_kaufman_b = Float64Builder::with_capacity(n);
    let mut intra_gk_vol_b = Float64Builder::with_capacity(n);
    let mut intra_hurst_b = Float64Builder::with_capacity(n);
    let mut intra_perm_ent_b = Float64Builder::with_capacity(n);

    // Gap awareness
    let mut has_gap_b = BooleanBuilder::with_capacity(n);
    let mut gap_count_b = Int64Builder::with_capacity(n);
    let mut gap_dur_b = Int64Builder::with_capacity(n);
    let mut is_exch_gap_b = BooleanBuilder::with_capacity(n);

    // Trade IDs
    let mut first_tid_b = Int64Builder::with_capacity(n);
    let mut last_tid_b = Int64Builder::with_capacity(n);

    // Bar flags
    let mut is_orphan_b = UInt8Builder::with_capacity(n);

    // Cache metadata
    let mut src_start_b = Int64Builder::with_capacity(n);
    let mut src_end_b = Int64Builder::with_capacity(n);

    for r in rows {
        symbol_b.append_value(&r.symbol);
        threshold_b.append_value(r.threshold_decimal_bps);
        close_time_b.append_value(r.close_time_us);
        open_time_b.append_value(r.open_time_us);
        open_b.append_value(r.open);
        high_b.append_value(r.high);
        low_b.append_value(r.low);
        close_b.append_value(r.close);
        volume_b.append_value(r.volume);
        vwap_b.append_value(r.vwap);
        buy_vol_b.append_value(r.buy_volume);
        sell_vol_b.append_value(r.sell_volume);
        ind_count_b.append_value(r.individual_trade_count);
        agg_count_b.append_value(r.agg_record_count);
        duration_b.append_value(r.duration_us);
        ofi_b.append_value(r.ofi);
        vwap_dev_b.append_value(r.vwap_close_deviation);
        price_impact_b.append_value(r.price_impact);
        kyle_b.append_value(r.kyle_lambda_proxy);
        intensity_b.append_value(r.trade_intensity);
        vol_per_trade_b.append_value(r.volume_per_trade);
        aggression_b.append_value(r.aggression_ratio);
        agg_density_b.append_value(r.aggregation_density);
        turnover_b.append_value(r.turnover_imbalance);
        ouroboros_mode_b.append_value(&r.ouroboros_mode);

        sess_sydney_b.append_value(r.exchange_session_sydney);
        sess_tokyo_b.append_value(r.exchange_session_tokyo);
        sess_london_b.append_value(r.exchange_session_london);
        sess_ny_b.append_value(r.exchange_session_newyork);

        // Inter-bar (nullable)
        lb_trade_count_b.append_option(r.lookback_trade_count);
        lb_ofi_b.append_option(r.lookback_ofi);
        lb_duration_b.append_option(r.lookback_duration_us);
        lb_intensity_b.append_option(r.lookback_intensity);
        lb_vwap_raw_b.append_option(r.lookback_vwap_raw);
        lb_vwap_pos_b.append_option(r.lookback_vwap_position);
        lb_count_imb_b.append_option(r.lookback_count_imbalance);
        lb_kyle_b.append_option(r.lookback_kyle_lambda);
        lb_burst_b.append_option(r.lookback_burstiness);
        lb_vol_skew_b.append_option(r.lookback_volume_skew);
        lb_vol_kurt_b.append_option(r.lookback_volume_kurt);
        lb_price_range_b.append_option(r.lookback_price_range);
        lb_kaufman_b.append_option(r.lookback_kaufman_er);
        lb_gk_vol_b.append_option(r.lookback_garman_klass_vol);
        lb_hurst_b.append_option(r.lookback_hurst);
        lb_perm_ent_b.append_option(r.lookback_permutation_entropy);

        // Intra-bar (nullable)
        intra_bull_density_b.append_option(r.intra_bull_epoch_density);
        intra_bear_density_b.append_option(r.intra_bear_epoch_density);
        intra_bull_gain_b.append_option(r.intra_bull_excess_gain);
        intra_bear_gain_b.append_option(r.intra_bear_excess_gain);
        intra_bull_cv_b.append_option(r.intra_bull_cv);
        intra_bear_cv_b.append_option(r.intra_bear_cv);
        intra_max_dd_b.append_option(r.intra_max_drawdown);
        intra_max_ru_b.append_option(r.intra_max_runup);
        intra_trade_count_b.append_option(r.intra_trade_count);
        intra_ofi_b.append_option(r.intra_ofi);
        intra_duration_b.append_option(r.intra_duration_us);
        intra_intensity_b.append_option(r.intra_intensity);
        intra_vwap_pos_b.append_option(r.intra_vwap_position);
        intra_count_imb_b.append_option(r.intra_count_imbalance);
        intra_kyle_b.append_option(r.intra_kyle_lambda);
        intra_burst_b.append_option(r.intra_burstiness);
        intra_vol_skew_b.append_option(r.intra_volume_skew);
        intra_vol_kurt_b.append_option(r.intra_volume_kurt);
        intra_kaufman_b.append_option(r.intra_kaufman_er);
        intra_gk_vol_b.append_option(r.intra_garman_klass_vol);
        intra_hurst_b.append_option(r.intra_hurst);
        intra_perm_ent_b.append_option(r.intra_permutation_entropy);

        // Gap awareness
        has_gap_b.append_value(r.has_gap);
        gap_count_b.append_value(r.gap_trade_count);
        gap_dur_b.append_value(r.max_gap_duration_us);
        is_exch_gap_b.append_value(r.is_exchange_gap);

        // Trade IDs
        first_tid_b.append_value(r.first_agg_trade_id);
        last_tid_b.append_value(r.last_agg_trade_id);

        // Bar flags
        is_orphan_b.append_value(r.is_orphan);

        // Cache metadata
        cache_key_b.append_value(&r.cache_key);
        version_b.append_value(&r.opendeviationbar_version);
        src_start_b.append_value(r.source_start_ts);
        src_end_b.append_value(r.source_end_ts);
    }

    let columns: Vec<arrow_array::ArrayRef> = vec![
        std::sync::Arc::new(symbol_b.finish()),
        std::sync::Arc::new(threshold_b.finish()),
        std::sync::Arc::new(close_time_b.finish()),
        std::sync::Arc::new(open_time_b.finish()),
        std::sync::Arc::new(open_b.finish()),
        std::sync::Arc::new(high_b.finish()),
        std::sync::Arc::new(low_b.finish()),
        std::sync::Arc::new(close_b.finish()),
        std::sync::Arc::new(volume_b.finish()),
        std::sync::Arc::new(vwap_b.finish()),
        std::sync::Arc::new(buy_vol_b.finish()),
        std::sync::Arc::new(sell_vol_b.finish()),
        std::sync::Arc::new(ind_count_b.finish()),
        std::sync::Arc::new(agg_count_b.finish()),
        std::sync::Arc::new(duration_b.finish()),
        std::sync::Arc::new(ofi_b.finish()),
        std::sync::Arc::new(vwap_dev_b.finish()),
        std::sync::Arc::new(price_impact_b.finish()),
        std::sync::Arc::new(kyle_b.finish()),
        std::sync::Arc::new(intensity_b.finish()),
        std::sync::Arc::new(vol_per_trade_b.finish()),
        std::sync::Arc::new(aggression_b.finish()),
        std::sync::Arc::new(agg_density_b.finish()),
        std::sync::Arc::new(turnover_b.finish()),
        std::sync::Arc::new(ouroboros_mode_b.finish()),
        std::sync::Arc::new(sess_sydney_b.finish()),
        std::sync::Arc::new(sess_tokyo_b.finish()),
        std::sync::Arc::new(sess_london_b.finish()),
        std::sync::Arc::new(sess_ny_b.finish()),
        std::sync::Arc::new(lb_trade_count_b.finish()),
        std::sync::Arc::new(lb_ofi_b.finish()),
        std::sync::Arc::new(lb_duration_b.finish()),
        std::sync::Arc::new(lb_intensity_b.finish()),
        std::sync::Arc::new(lb_vwap_raw_b.finish()),
        std::sync::Arc::new(lb_vwap_pos_b.finish()),
        std::sync::Arc::new(lb_count_imb_b.finish()),
        std::sync::Arc::new(lb_kyle_b.finish()),
        std::sync::Arc::new(lb_burst_b.finish()),
        std::sync::Arc::new(lb_vol_skew_b.finish()),
        std::sync::Arc::new(lb_vol_kurt_b.finish()),
        std::sync::Arc::new(lb_price_range_b.finish()),
        std::sync::Arc::new(lb_kaufman_b.finish()),
        std::sync::Arc::new(lb_gk_vol_b.finish()),
        std::sync::Arc::new(lb_hurst_b.finish()),
        std::sync::Arc::new(lb_perm_ent_b.finish()),
        std::sync::Arc::new(intra_bull_density_b.finish()),
        std::sync::Arc::new(intra_bear_density_b.finish()),
        std::sync::Arc::new(intra_bull_gain_b.finish()),
        std::sync::Arc::new(intra_bear_gain_b.finish()),
        std::sync::Arc::new(intra_bull_cv_b.finish()),
        std::sync::Arc::new(intra_bear_cv_b.finish()),
        std::sync::Arc::new(intra_max_dd_b.finish()),
        std::sync::Arc::new(intra_max_ru_b.finish()),
        std::sync::Arc::new(intra_trade_count_b.finish()),
        std::sync::Arc::new(intra_ofi_b.finish()),
        std::sync::Arc::new(intra_duration_b.finish()),
        std::sync::Arc::new(intra_intensity_b.finish()),
        std::sync::Arc::new(intra_vwap_pos_b.finish()),
        std::sync::Arc::new(intra_count_imb_b.finish()),
        std::sync::Arc::new(intra_kyle_b.finish()),
        std::sync::Arc::new(intra_burst_b.finish()),
        std::sync::Arc::new(intra_vol_skew_b.finish()),
        std::sync::Arc::new(intra_vol_kurt_b.finish()),
        std::sync::Arc::new(intra_kaufman_b.finish()),
        std::sync::Arc::new(intra_gk_vol_b.finish()),
        std::sync::Arc::new(intra_hurst_b.finish()),
        std::sync::Arc::new(intra_perm_ent_b.finish()),
        std::sync::Arc::new(has_gap_b.finish()),
        std::sync::Arc::new(gap_count_b.finish()),
        std::sync::Arc::new(gap_dur_b.finish()),
        std::sync::Arc::new(is_exch_gap_b.finish()),
        std::sync::Arc::new(first_tid_b.finish()),
        std::sync::Arc::new(last_tid_b.finish()),
        std::sync::Arc::new(is_orphan_b.finish()),
        std::sync::Arc::new(cache_key_b.finish()),
        std::sync::Arc::new(version_b.finish()),
        std::sync::Arc::new(src_start_b.finish()),
        std::sync::Arc::new(src_end_b.finish()),
    ];

    let schema_ref = std::sync::Arc::new(schema.clone());
    RecordBatch::try_new(schema_ref, columns).map_err(|e| DeadLetterError::Arrow(e.to_string()))
}

/// Write failed rows to a Parquet dead-letter file.
///
/// Returns the path of the written file, or an empty `PathBuf` if `rows` is empty.
///
/// # File naming
///
/// `{symbol}_{threshold}_{unix_secs}.parquet` using `rows[0]` metadata.
///
/// # Compression
///
/// Uses ZSTD for compression (matches Python Parquet conventions).
pub fn write_dead_letter(rows: &[ClickHouseBarRow]) -> Result<PathBuf, DeadLetterError> {
    if rows.is_empty() {
        return Ok(PathBuf::new());
    }

    let dir = Path::new(DEAD_LETTER_DIR);
    std::fs::create_dir_all(dir)?;

    let symbol = &rows[0].symbol;
    let threshold = rows[0].threshold_decimal_bps;
    let timestamp = std::time::SystemTime::now()
        .duration_since(std::time::UNIX_EPOCH)
        .unwrap_or_default()
        .as_secs();

    let filename = format!("{symbol}_{threshold}_{timestamp}.parquet");
    let path = dir.join(&filename);

    let schema = dead_letter_schema();
    let batch = rows_to_record_batch(rows, &schema)?;

    let props = parquet::file::properties::WriterProperties::builder()
        .set_compression(parquet::basic::Compression::ZSTD(
            parquet::basic::ZstdLevel::try_new(3)?,
        ))
        .build();

    let file = std::fs::File::create(&path)?;
    let mut writer =
        parquet::arrow::ArrowWriter::try_new(file, std::sync::Arc::new(schema), Some(props))?;
    writer.write(&batch)?;
    writer.close()?;

    tracing::warn!(
        symbol,
        threshold,
        rows = rows.len(),
        path = %path.display(),
        "dead-lettered failed bars to Parquet"
    );

    Ok(path)
}

/// Write a single failed row to a Parquet dead-letter file.
///
/// Convenience wrapper around `write_dead_letter` for the backpressure path.
pub fn write_dead_letter_single(row: ClickHouseBarRow) -> Result<PathBuf, DeadLetterError> {
    write_dead_letter(&[row])
}

#[cfg(test)]
mod tests {
    use super::*;
    use super::super::row::{ClickHouseBarRow, CORE_COLUMNS};
    use opendeviationbar_core::fixed_point::FixedPoint;
    use opendeviationbar_core::OpenDeviationBar;
    use parquet::file::reader::FileReader;
    use std::sync::Arc;

    use crate::live_engine::CompletedBar;

    fn test_row(first_tid: i64, last_tid: i64) -> ClickHouseBarRow {
        let mut bar = OpenDeviationBar::default();
        bar.open = FixedPoint::from_str("50000.0").unwrap();
        bar.high = FixedPoint::from_str("50100.0").unwrap();
        bar.low = FixedPoint::from_str("49900.0").unwrap();
        bar.close = FixedPoint::from_str("50050.0").unwrap();
        bar.vwap = FixedPoint::from_str("50025.0").unwrap();
        bar.open_time = 1_700_000_000_000_000;
        bar.close_time = 1_700_000_100_000_000;
        bar.first_agg_trade_id = first_tid;
        bar.last_agg_trade_id = last_tid;
        bar.individual_trade_count = 100;
        bar.agg_record_count = 50;
        bar.duration_us = 100_000_000;
        bar.lookback_trade_count = Some(200);
        bar.lookback_ofi = Some(0.1);

        let completed = CompletedBar {
            symbol: Arc::from("BTCUSDT"),
            threshold_decimal_bps: 250,
            bar,
        };
        ClickHouseBarRow::from_completed_bar(&completed)
    }

    /// Use a unique temp dir per test to avoid test isolation issues.
    /// Uses PID + high-resolution timestamp for uniqueness across parallel nextest.
    fn test_dead_letter_dir() -> PathBuf {
        let pid = std::process::id();
        let nanos = std::time::SystemTime::now()
            .duration_since(std::time::UNIX_EPOCH)
            .unwrap_or_default()
            .as_nanos();
        let dir = std::env::temp_dir().join(format!(
            "opendeviationbar-dead-letter-test-{pid}-{nanos}"
        ));
        let _ = std::fs::remove_dir_all(&dir);
        std::fs::create_dir_all(&dir).unwrap();
        dir
    }

    #[test]
    fn test_dead_letter_write_3_rows() {
        let dir = test_dead_letter_dir();
        let rows = vec![test_row(1, 10), test_row(11, 20), test_row(21, 30)];

        let schema = dead_letter_schema();
        let batch = rows_to_record_batch(&rows, &schema).unwrap();

        // Write to test dir
        let path = dir.join("BTCUSDT_250_12345.parquet");
        let props = parquet::file::properties::WriterProperties::builder()
            .set_compression(parquet::basic::Compression::ZSTD(
                parquet::basic::ZstdLevel::try_new(3).unwrap(),
            ))
            .build();
        let file = std::fs::File::create(&path).unwrap();
        let mut writer = parquet::arrow::ArrowWriter::try_new(
            file,
            std::sync::Arc::new(schema),
            Some(props),
        )
        .unwrap();
        writer.write(&batch).unwrap();
        writer.close().unwrap();

        assert!(path.exists());

        // Read back and verify row count
        let reader = parquet::arrow::arrow_reader::ParquetRecordBatchReader::try_new(
            std::fs::File::open(&path).unwrap(),
            1024,
        )
        .unwrap();
        let batches: Vec<RecordBatch> = reader.into_iter().collect::<Result<_, _>>().unwrap();
        let total_rows: usize = batches.iter().map(RecordBatch::num_rows).sum();
        assert_eq!(total_rows, 3);

        let _ = std::fs::remove_dir_all(&dir);
    }

    #[test]
    fn test_dead_letter_write_single() {
        let dir = test_dead_letter_dir();
        let row = test_row(1, 10);

        let schema = dead_letter_schema();
        let batch = rows_to_record_batch(&[row], &schema).unwrap();

        let path = dir.join("BTCUSDT_250_single.parquet");
        let props = parquet::file::properties::WriterProperties::builder()
            .set_compression(parquet::basic::Compression::ZSTD(
                parquet::basic::ZstdLevel::try_new(3).unwrap(),
            ))
            .build();
        let file = std::fs::File::create(&path).unwrap();
        let mut writer = parquet::arrow::ArrowWriter::try_new(
            file,
            std::sync::Arc::new(schema),
            Some(props),
        )
        .unwrap();
        writer.write(&batch).unwrap();
        writer.close().unwrap();

        assert!(path.exists());
        let _ = std::fs::remove_dir_all(&dir);
    }

    #[test]
    fn test_dead_letter_schema_nullable_flags() {
        let schema = dead_letter_schema();
        let fields = schema.fields();

        // All lookback_* and intra_* should be nullable
        for field in fields.iter() {
            let name = field.name();
            if name.starts_with("lookback_") || name.starts_with("intra_") {
                assert!(
                    field.is_nullable(),
                    "Field {name} should be nullable"
                );
            }
        }

        // Core non-nullable fields
        let non_nullable_names = [
            "symbol",
            "threshold_decimal_bps",
            "open",
            "high",
            "low",
            "close",
            "volume",
            "has_gap",
            "first_agg_trade_id",
            "last_agg_trade_id",
        ];
        for name in &non_nullable_names {
            let field = schema.field_with_name(name).unwrap();
            assert!(
                !field.is_nullable(),
                "Field {name} should NOT be nullable"
            );
        }
    }

    #[test]
    fn test_dead_letter_column_count() {
        let schema = dead_letter_schema();
        assert_eq!(
            schema.fields().len(),
            CORE_COLUMNS.len(),
            "Schema column count must match CORE_COLUMNS ({})",
            CORE_COLUMNS.len()
        );
    }

    #[test]
    fn test_dead_letter_zstd_compression() {
        let dir = test_dead_letter_dir();
        let row = test_row(1, 10);

        let schema = dead_letter_schema();
        let batch = rows_to_record_batch(&[row], &schema).unwrap();

        let path = dir.join("BTCUSDT_250_zstd.parquet");
        let props = parquet::file::properties::WriterProperties::builder()
            .set_compression(parquet::basic::Compression::ZSTD(
                parquet::basic::ZstdLevel::try_new(3).unwrap(),
            ))
            .build();
        let file = std::fs::File::create(&path).unwrap();
        let mut writer = parquet::arrow::ArrowWriter::try_new(
            file,
            std::sync::Arc::new(schema),
            Some(props),
        )
        .unwrap();
        writer.write(&batch).unwrap();
        writer.close().unwrap();

        // Read back and check compression from metadata
        let file = std::fs::File::open(&path).unwrap();
        let reader = parquet::file::reader::SerializedFileReader::new(file).unwrap();
        let meta = reader.metadata();
        let row_group = meta.row_group(0);
        let col_meta = row_group.column(0);
        assert!(
            matches!(
                col_meta.compression(),
                parquet::basic::Compression::ZSTD(_)
            ),
            "Expected ZSTD compression, got {:?}",
            col_meta.compression()
        );

        let _ = std::fs::remove_dir_all(&dir);
    }

    #[test]
    fn test_dead_letter_empty_slice() {
        let result = write_dead_letter(&[]);
        assert!(result.is_ok());
        let path = result.unwrap();
        assert_eq!(path, PathBuf::new());
    }

    /// Write a dead-letter Parquet artifact for cross-language round-trip testing (RESIL-02).
    ///
    /// Writes a known ClickHouseBarRow to a fixed path so the Python test
    /// `tests/test_dead_letter_roundtrip.py` can read it with polars.read_parquet().
    #[test]
    fn test_write_roundtrip_artifact() {
        let artifact_dir = std::env::temp_dir().join("opendeviationbar-dead-letter-roundtrip");
        let _ = std::fs::remove_dir_all(&artifact_dir);
        std::fs::create_dir_all(&artifact_dir).unwrap();

        let mut row = test_row(1000, 1099);
        // Set some nullable fields to None to test null handling
        row.lookback_duration_us = None;
        row.lookback_intensity = None;
        row.intra_bull_epoch_density = None;
        row.intra_hurst = None;
        // Set some to specific values to verify round-trip
        row.lookback_trade_count = Some(500);
        row.lookback_ofi = Some(-0.42);

        let schema = dead_letter_schema();
        let batch = rows_to_record_batch(&[row], &schema).unwrap();

        let path = artifact_dir.join("BTCUSDT_250_9999999999.parquet");
        let props = parquet::file::properties::WriterProperties::builder()
            .set_compression(parquet::basic::Compression::ZSTD(
                parquet::basic::ZstdLevel::try_new(3).unwrap(),
            ))
            .build();
        let file = std::fs::File::create(&path).unwrap();
        let mut writer = parquet::arrow::ArrowWriter::try_new(
            file,
            std::sync::Arc::new(schema),
            Some(props),
        )
        .unwrap();
        writer.write(&batch).unwrap();
        writer.close().unwrap();

        assert!(path.exists());
        // Print path so Python test can find it
        eprintln!("ROUNDTRIP_ARTIFACT={}", path.display());
    }

    #[test]
    fn test_dead_letter_column_names_match_core_columns() {
        let schema = dead_letter_schema();
        for (i, field) in schema.fields().iter().enumerate() {
            assert_eq!(
                field.name(),
                CORE_COLUMNS[i],
                "Column {i} name mismatch: schema has '{}', CORE_COLUMNS has '{}'",
                field.name(),
                CORE_COLUMNS[i]
            );
        }
    }
}