sz-orm-core 7.6.0

Core ORM engine: Model trait, ActiveRecord, QueryBuilder, Pool, Transaction, migration, and SQL dialect abstraction
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
796
797
798
799
800
801
802
803
804
805
806
807
808
809
810
811
//! 零拷贝结果集传输深化(v6.8.0 TASK W1-6)
//!
//! 在结果集传输路径上引入 `bytes::Bytes` 切片引用传递,
//! 减少不必要的数据拷贝。提供零拷贝命中/回退拷贝统计。

use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;

use bytes::Bytes;

use crate::value_borrowed::{BorrowedRowData, BorrowedValue};
use crate::Value;

/// 零拷贝统计
#[derive(Debug, Default)]
pub struct ZeroCopyStats {
    /// 零拷贝命中次数
    zero_copy_hits: AtomicU64,
    /// 回退拷贝次数
    fallback_copies: AtomicU64,
    /// v7.3.0 任务 1.3:峰值 RSS(字节)
    peak_rss_bytes: AtomicU64,
    /// v7.3.0 任务 1.3:分配次数
    allocation_count: AtomicU64,
    /// v7.4.0:DecimalBytes 类型未命中次数
    decimal_bytes_misses: AtomicU64,
    /// v7.4.0:JsonBytes 类型未命中次数
    json_bytes_misses: AtomicU64,
    /// v7.4.0:BytesRef 类型未命中次数
    bytes_ref_misses: AtomicU64,
    /// v7.4.0:DateTimeInt 类型未命中次数
    datetime_int_misses: AtomicU64,
    /// v7.6.0:ArrayRef5.0:ArrayRef 借用命中次数(零拷贝)
    array_ref_hits: AtomicU64,
    /// v7.6.0:ObjectRef 借用命中次数(零拷贝)
    object_ref_hits: AtomicU64,
    /// v7.6.0:Array 堆分配回退次数
    array_heap_fallbacks: AtomicU64,
    /// v7.6.0:Object 堆分配回退次数
    object_heap_fallbacks: AtomicU64,
}

impl ZeroCopyStats {
    pub fn new() -> Self {
        Self::default()
    }

    pub fn record_zero_copy_hit(&self) {
        self.zero_copy_hits.fetch_add(1, Ordering::Relaxed);
    }

    pub fn record_fallback_copy(&self) {
        self.fallback_copies.fetch_add(1, Ordering::Relaxed);
    }

    pub fn zero_copy_hits(&self) -> u64 {
        self.zero_copy_hits.load(Ordering::Relaxed)
    }

    pub fn fallback_copies(&self) -> u64 {
        self.fallback_copies.load(Ordering::Relaxed)
    }

    pub fn total(&self) -> u64 {
        self.zero_copy_hits() + self.fallback_copies()
    }

    pub fn hit_rate(&self) -> f64 {
        let total = self.total();
        if total == 0 {
            0.0
        } else {
            self.zero_copy_hits() as f64 / total as f64
        }
    }

    /// v7.3.0 任务 1.3:记录峰值 RSS(取最大值)
    pub fn record_rss(&self, rss_bytes: u64) {
        let mut current = self.peak_rss_bytes.load(Ordering::Relaxed);
        while rss_bytes > current {
            match self.peak_rss_bytes.compare_exchange_weak(
                current,
                rss_bytes,
                Ordering::Relaxed,
                Ordering::Relaxed,
            ) {
                Ok(_) => break,
                Err(actual) => current = actual,
            }
        }
    }

    /// v7.3.0 任务 1.3:记录一次分配
    pub fn record_allocation(&self) {
        self.allocation_count.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.3.0 任务 1.3:峰值 RSS(字节)
    pub fn peak_rss_bytes(&self) -> u64 {
        self.peak_rss_bytes.load(Ordering::Relaxed)
    }

    /// v7.3.0 任务 1.3:分配次数
    pub fn allocation_count(&self) -> u64 {
        self.allocation_count.load(Ordering::Relaxed)
    }

    /// v7.3.0 任务 1.3:RSS 降低百分比(对比基准)
    pub fn rss_reduction_pct(&self, baseline_rss: u64) -> f64 {
        if baseline_rss == 0 {
            return 0.0;
        }
        let current = self.peak_rss_bytes();
        if current >= baseline_rss {
            0.0
        } else {
            (baseline_rss - current) as f64 / baseline_rss as f64 * 100.0
        }
    }

    /// v7.4.0:记录 DecimalBytes 类型未命中
    pub fn record_decimal_bytes_miss(&self) {
        self.decimal_bytes_misses.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.4.0:记录 JsonBytes 类型未命中
    pub fn record_json_bytes_miss(&self) {
        self.json_bytes_misses.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.4.0:记录 BytesRef 类型未命中
    pub fn record_bytes_ref_miss(&self) {
        self.bytes_ref_misses.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.4.0:记录 DateTimeInt 类型未命中
    pub fn record_datetime_int_miss(&self) {
        self.datetime_int_misses.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.4.0:DecimalBytes 未命中次数
    pub fn decimal_bytes_misses(&self) -> u64 {
        self.decimal_bytes_misses.load(Ordering::Relaxed)
    }

    /// v7.4.0:JsonBytes 未命中次数
    pub fn json_bytes_misses(&self) -> u64 {
        self.json_bytes_misses.load(Ordering::Relaxed)
    }

    /// v7.4.0:BytesRef 未命中次数
    pub fn bytes_ref_misses(&self) -> u64 {
        self.bytes_ref_misses.load(Ordering::Relaxed)
    }

    /// v7.4.0:DateTimeInt 未命中次数
    pub fn datetime_int_misses(&self) -> u64 {
        self.datetime_int_misses.load(Ordering::Relaxed)
    }

    /// v7.4.0:按类型未命中总次数
    pub fn total_type_misses(&self) -> u64 {
        self.decimal_bytes_misses()
            + self.json_bytes_misses()
            + self.bytes_ref_misses()
            + self.datetime_int_misses()
    }

    /// v7.6.0:记录 ArrayRef 借用命中(零拷贝)
    pub fn record_array_ref_hit(&self) {
        self.array_ref_hits.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.6.0:记录 ObjectRef 借用命中(零拷贝)
    pub fn record_object_ref_hit(&self) {
        self.object_ref_hits.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.6.0:记录 Array 堆分配回退
    pub fn record_array_heap_fallback(&self) {
        self.array_heap_fallbacks.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.6.0:记录 Object 堆分配回退
    pub fn record_object_heap_fallback(&self) {
        self.object_heap_fallbacks.fetch_add(1, Ordering::Relaxed);
    }

    /// v7.6.0:ArrayRef 借用命中次数
    pub fn array_ref_hits(&self) -> u64 {
        self.array_ref_hits.load(Ordering::Relaxed)
    }

    /// v7.6.0:ObjectRef 借用命中次数
    pub fn object_ref_hits(&self) -> u64 {
        self.object_ref_hits.load(Ordering::Relaxed)
    }

    /// v7.6.0:Array 堆分配回退次数
    pub fn array_heap_fallbacks(&self) -> u64 {
        self.array_heap_fallbacks.load(Ordering::Relaxed)
    }

    /// v7.6.0:Object 堆分配回退次数
    pub fn object_heap_fallbacks(&self) -> u64 {
        self.object_heap_fallbacks.load(Ordering::Relaxed)
    }

    /// v7.6.0:堆分配降低率(0.0 ~ 1.0)
    ///
    /// 计算公式:`ref_hits / (ref_hits + heap_fallbacks)`
    /// 其中 `ref_hits = array_ref_hits + object_ref_hits`,
    /// `heap_fallbacks = array_heap_fallbacks + object_heap_fallbacks`。
    pub fn heap_reduction_rate(&self) -> f64 {
        let ref_hits = self.array_ref_hits() + self.object_ref_hits();
        let heap_fallbacks = self.array_heap_fallbacks() + self.object_heap_fallbacks();
        let total = ref_hits + heap_fallbacks;
        if total == 0 {
            0.0
        } else {
            ref_hits as f64 / total as f64
        }
    }
}

impl Clone for ZeroCopyStats {
    fn clone(&self) -> Self {
        Self {
            zero_copy_hits: AtomicU64::new(self.zero_copy_hits()),
            fallback_copies: AtomicU64::new(self.fallback_copies()),
            peak_rss_bytes: AtomicU64::new(self.peak_rss_bytes()),
            allocation_count: AtomicU64::new(self.allocation_count()),
            decimal_bytes_misses: AtomicU64::new(self.decimal_bytes_misses()),
            json_bytes_misses: AtomicU64::new(self.json_bytes_misses()),
            bytes_ref_misses: AtomicU64::new(self.bytes_ref_misses()),
            datetime_int_misses: AtomicU64::new(self.datetime_int_misses()),
            array_ref_hits: AtomicU64::new(self.array_ref_hits()),
            object_ref_hits: AtomicU64::new(self.object_ref_hits()),
            array_heap_fallbacks: AtomicU64::new(self.array_heap_fallbacks()),
            object_heap_fallbacks: AtomicU64::new(self.object_heap_fallbacks()),
        }
    }
}

/// 零拷贝行(使用 Bytes 引用计数切片)
#[derive(Debug, Clone)]
pub struct ZeroCopyRow {
    /// 列名
    pub columns: Arc<Vec<String>>,
    /// 行数据(Bytes 引用计数,clone 廉价)
    pub data: Bytes,
    /// 列偏移(每列的起始位置和长度)
    pub offsets: Vec<(usize, usize)>,
}

impl ZeroCopyRow {
    pub fn get(&self, col: &str) -> Option<&[u8]> {
        let idx = self.columns.iter().position(|c| c == col)?;
        let (start, len) = self.offsets[idx];
        Some(&self.data[start..start + len])
    }

    pub fn column_count(&self) -> usize {
        self.columns.len()
    }
}

/// 零拷贝行流
pub struct ZeroCopyRowStream {
    rows: Vec<ZeroCopyRow>,
    index: usize,
    stats: Arc<ZeroCopyStats>,
}

impl ZeroCopyRowStream {
    pub fn next_row(&mut self) -> Option<&ZeroCopyRow> {
        if self.index < self.rows.len() {
            let row = &self.rows[self.index];
            self.index += 1;
            Some(row)
        } else {
            None
        }
    }

    pub fn remaining(&self) -> usize {
        self.rows.len() - self.index
    }

    pub fn total_rows(&self) -> usize {
        self.rows.len()
    }

    pub fn stats(&self) -> &ZeroCopyStats {
        &self.stats
    }
}

impl Iterator for ZeroCopyRowStream {
    type Item = ZeroCopyRow;

    fn next(&mut self) -> Option<Self::Item> {
        if self.index < self.rows.len() {
            let row = self.rows[self.index].clone();
            self.index += 1;
            Some(row)
        } else {
            None
        }
    }
}

/// 零拷贝管线
pub struct ZeroCopyPipeline {
    stats: Arc<ZeroCopyStats>,
}

impl ZeroCopyPipeline {
    pub fn new() -> Self {
        Self {
            stats: Arc::new(ZeroCopyStats::new()),
        }
    }

    /// 从行数据创建零拷贝行流
    ///
    /// 将 `Vec<HashMap<String, Value>>` 转换为零拷贝行流,
    /// 字符串/字节数据使用 `Bytes` 引用计数切片。
    ///
    /// # 关于"零拷贝"语义
    ///
    /// 此处的"零拷贝"是指将多行数据序列化到连续 `Bytes` buffer 后,
    /// 通过引用计数切片(`Bytes::slice`)共享同一内存区域,
    /// **减少跨行/跨列的引用计数开销**,而非完全无拷贝。
    /// `Value::String` / `Value::Bytes` 分支中 `extend_from_slice`
    /// 仍有一次从源数据到连续 buffer 的拷贝。
    pub fn stream_rows(
        &self,
        rows: Vec<std::collections::HashMap<String, Value>>,
        columns: &[String],
    ) -> ZeroCopyRowStream {
        let col_count = columns.len();
        let col_arc: Arc<Vec<String>> = Arc::new(columns.to_vec());
        let mut zero_copy_rows = Vec::with_capacity(rows.len());

        for row in rows {
            let mut buffer = Vec::new();
            let mut offsets = Vec::with_capacity(col_count);

            for col in columns.iter() {
                let value = row.get(col);
                let start = buffer.len();
                match value {
                    Some(Value::String(s)) => {
                        buffer.extend_from_slice(s.as_bytes());
                        self.stats.record_zero_copy_hit();
                    }
                    Some(Value::Bytes(b)) => {
                        buffer.extend_from_slice(b);
                        self.stats.record_zero_copy_hit();
                    }
                    Some(v) => {
                        let formatted = format!("{}", v);
                        buffer.extend_from_slice(formatted.as_bytes());
                        self.stats.record_fallback_copy();
                    }
                    None => {
                        self.stats.record_fallback_copy();
                    }
                }
                let len = buffer.len() - start;
                offsets.push((start, len));
            }

            zero_copy_rows.push(ZeroCopyRow {
                columns: col_arc.clone(),
                data: Bytes::from(buffer),
                offsets,
            });
        }

        ZeroCopyRowStream {
            rows: zero_copy_rows,
            index: 0,
            stats: self.stats.clone(),
        }
    }

    /// 从 BorrowedRowData 创建零拷贝行流(直接引用,零拷贝)
    pub fn stream_borrowed(
        &self,
        rows: Vec<BorrowedRowData<'_>>,
        columns: &[String],
    ) -> ZeroCopyRowStream {
        let col_arc: Arc<Vec<String>> = Arc::new(columns.to_vec());
        let mut zero_copy_rows = Vec::with_capacity(rows.len());

        for row in rows {
            let mut buffer = Vec::new();
            let mut offsets = Vec::with_capacity(columns.len());

            for col in columns.iter() {
                let start = buffer.len();
                if let Some(val) = row.get(col) {
                    match val {
                        BorrowedValue::String(s) => {
                            buffer.extend_from_slice(s.as_bytes());
                        }
                        BorrowedValue::Bytes(b) => {
                            buffer.extend_from_slice(b.as_ref());
                        }
                        _ => {
                            let owned = val.to_owned_value();
                            let formatted = format!("{}", owned);
                            buffer.extend_from_slice(formatted.as_bytes());
                            self.stats.record_fallback_copy();
                        }
                    }
                    self.stats.record_zero_copy_hit();
                } else {
                    self.stats.record_fallback_copy();
                }
                let len = buffer.len() - start;
                offsets.push((start, len));
            }

            zero_copy_rows.push(ZeroCopyRow {
                columns: col_arc.clone(),
                data: Bytes::from(buffer),
                offsets,
            });
        }

        ZeroCopyRowStream {
            rows: zero_copy_rows,
            index: 0,
            stats: self.stats.clone(),
        }
    }

    pub fn stats(&self) -> &ZeroCopyStats {
        &self.stats
    }
}

impl Default for ZeroCopyPipeline {
    fn default() -> Self {
        Self::new()
    }
}

// ============================================================================
// v7.3.0 任务 1.3:零拷贝类型注册表
// ============================================================================

/// 零拷贝类型标识(v7.3.0)
///
/// 标识支持零拷贝序列化的内置类型。
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub enum ZeroCopyTypeId {
    /// i32
    I32,
    /// i64
    I64,
    /// f32
    F32,
    /// f64
    F64,
    /// bool
    Bool,
    /// String
    String,
    /// Bytes
    Bytes,
}

/// 零拷贝类型注册表(v7.3.0)
///
/// 注册/查询类型是否支持零拷贝序列化。内置类型默认支持,
/// 可通过 `register` 扩展(标记自定义类型支持零拷贝)。
#[derive(Debug, Clone)]
pub struct ZeroCopyTypeRegistry {
    supported: std::collections::HashSet<ZeroCopyTypeId>,
}

impl Default for ZeroCopyTypeRegistry {
    fn default() -> Self {
        Self::with_builtins()
    }
}

impl ZeroCopyTypeRegistry {
    /// 创建包含所有内置类型的注册表
    pub fn with_builtins() -> Self {
        let mut supported = std::collections::HashSet::new();
        supported.insert(ZeroCopyTypeId::I32);
        supported.insert(ZeroCopyTypeId::I64);
        supported.insert(ZeroCopyTypeId::F32);
        supported.insert(ZeroCopyTypeId::F64);
        supported.insert(ZeroCopyTypeId::Bool);
        supported.insert(ZeroCopyTypeId::String);
        supported.insert(ZeroCopyTypeId::Bytes);
        Self { supported }
    }

    /// 创建空注册表
    pub fn empty() -> Self {
        Self {
            supported: std::collections::HashSet::new(),
        }
    }

    /// 注册类型支持零拷贝
    pub fn register(&mut self, type_id: ZeroCopyTypeId) {
        self.supported.insert(type_id);
    }

    /// 查询类型是否支持零拷贝
    pub fn is_supported(&self, type_id: ZeroCopyTypeId) -> bool {
        self.supported.contains(&type_id)
    }

    /// 已注册类型数
    pub fn len(&self) -> usize {
        self.supported.len()
    }

    /// 是否为空
    pub fn is_empty(&self) -> bool {
        self.supported.is_empty()
    }
}

impl ZeroCopyPipeline {
    /// v7.3.0 任务 1.3:使用类型注册表尝试零拷贝解析
    ///
    /// 对行中每列判断 Value 类型是否在注册表中支持零拷贝。
    /// 全部列支持时返回 `Some(ZeroCopyRow)` 并记录 `zero_copy_hit`;
    /// 任一列不支持时返回 `None`,记录 `fallback_copy` + 分配次数,
    /// 触发回退传统 clone 路径。
    pub fn try_parse_with_registry(
        &self,
        row: &std::collections::HashMap<String, Value>,
        columns: &[String],
        registry: &ZeroCopyTypeRegistry,
    ) -> Option<ZeroCopyRow> {
        let col_arc: Arc<Vec<String>> = Arc::new(columns.to_vec());
        let mut buffer = Vec::new();
        let mut offsets = Vec::with_capacity(columns.len());
        let mut all_supported = true;

        for col in columns.iter() {
            let start = buffer.len();
            if let Some(value) = row.get(col) {
                let type_id = value_to_type_id(value);
                if registry.is_supported(type_id) {
                    match value {
                        Value::String(s) => buffer.extend_from_slice(s.as_bytes()),
                        Value::Bytes(b) => buffer.extend_from_slice(b),
                        _ => {
                            let formatted = format!("{}", value);
                            buffer.extend_from_slice(formatted.as_bytes());
                        }
                    }
                    self.stats.record_zero_copy_hit();
                } else {
                    all_supported = false;
                    let formatted = format!("{}", value);
                    buffer.extend_from_slice(formatted.as_bytes());
                    self.stats.record_fallback_copy();
                    self.stats.record_allocation();
                }
            } else {
                self.stats.record_fallback_copy();
                self.stats.record_allocation();
            }
            let len = buffer.len() - start;
            offsets.push((start, len));
        }

        if all_supported {
            Some(ZeroCopyRow {
                columns: col_arc,
                data: Bytes::from(buffer),
                offsets,
            })
        } else {
            None
        }
    }
}

/// 将 Value 映射到 ZeroCopyTypeId
fn value_to_type_id(value: &Value) -> ZeroCopyTypeId {
    match value {
        Value::I32(_) => ZeroCopyTypeId::I32,
        Value::I64(_) => ZeroCopyTypeId::I64,
        Value::F32(_) => ZeroCopyTypeId::F32,
        Value::F64(_) => ZeroCopyTypeId::F64,
        Value::Bool(_) => ZeroCopyTypeId::Bool,
        Value::String(_) => ZeroCopyTypeId::String,
        Value::Bytes(_) => ZeroCopyTypeId::Bytes,
        _ => ZeroCopyTypeId::String, // 其他类型回退为 String
    }
}

#[cfg(test)]
mod tests {
    use super::*;
    use std::collections::HashMap;

    fn make_rows(n: usize, cols: &[&str]) -> Vec<HashMap<String, Value>> {
        (0..n)
            .map(|i| {
                let mut row = HashMap::new();
                for col in cols {
                    row.insert(col.to_string(), Value::String(format!("val_{}_{}", i, col)));
                }
                row
            })
            .collect()
    }

    #[test]
    fn stream_rows_preserves_count() {
        let pipeline = ZeroCopyPipeline::new();
        let columns = vec!["id".to_string(), "name".to_string()];
        let rows = make_rows(100, &["id", "name"]);
        let stream = pipeline.stream_rows(rows, &columns);
        assert_eq!(stream.total_rows(), 100);
    }

    #[test]
    fn stream_rows_data_accessible() {
        let pipeline = ZeroCopyPipeline::new();
        let columns = vec!["name".to_string()];
        let mut rows = Vec::new();
        let mut row = HashMap::new();
        row.insert("name".to_string(), Value::String("Alice".into()));
        rows.push(row);
        let mut stream = pipeline.stream_rows(rows, &columns);
        let first = stream.next().unwrap();
        assert_eq!(first.get("name").unwrap(), b"Alice");
    }

    #[test]
    fn stats_track_zero_copy_hits() {
        let pipeline = ZeroCopyPipeline::new();
        let columns = vec!["name".to_string()];
        let rows = make_rows(10, &["name"]);
        let _stream = pipeline.stream_rows(rows, &columns);
        assert!(pipeline.stats().zero_copy_hits() > 0);
    }

    #[test]
    fn stats_track_fallback_copies() {
        let pipeline = ZeroCopyPipeline::new();
        let columns = vec!["num".to_string()];
        let mut rows = Vec::new();
        let mut row = HashMap::new();
        row.insert("num".to_string(), Value::I64(42));
        rows.push(row);
        let _stream = pipeline.stream_rows(rows, &columns);
        assert!(pipeline.stats().fallback_copies() > 0);
    }

    #[test]
    fn empty_rows_stream() {
        let pipeline = ZeroCopyPipeline::new();
        let columns = vec!["id".to_string()];
        let stream = pipeline.stream_rows(Vec::new(), &columns);
        assert_eq!(stream.total_rows(), 0);
        assert_eq!(stream.remaining(), 0);
    }

    #[test]
    fn iterator_interface() {
        let pipeline = ZeroCopyPipeline::new();
        let columns = vec!["id".to_string()];
        let rows = make_rows(5, &["id"]);
        let stream = pipeline.stream_rows(rows, &columns);
        let collected: Vec<_> = stream.collect();
        assert_eq!(collected.len(), 5);
    }

    #[test]
    fn hit_rate_calculation() {
        let stats = ZeroCopyStats::new();
        for _ in 0..7 {
            stats.record_zero_copy_hit();
        }
        for _ in 0..3 {
            stats.record_fallback_copy();
        }
        assert!((stats.hit_rate() - 0.7).abs() < 0.001);
    }

    #[test]
    fn large_resultset_streaming() {
        let pipeline = ZeroCopyPipeline::new();
        let columns: Vec<String> = (0..20).map(|i| format!("col_{}", i)).collect();
        let rows: Vec<HashMap<String, Value>> = (0..10000)
            .map(|i| {
                let mut row = HashMap::new();
                for j in 0..20 {
                    row.insert(
                        format!("col_{}", j),
                        Value::String(format!("val_{}_{}", i, j)),
                    );
                }
                row
            })
            .collect();
        let stream = pipeline.stream_rows(rows, &columns);
        assert_eq!(stream.total_rows(), 10000);
        assert!(pipeline.stats().zero_copy_hits() > 0);
    }

    // ========================================================================
    // v7.6.0 任务 1.3:ZeroCopyStats 扩展测试
    // ========================================================================

    #[test]
    fn test_array_ref_hits_tracking() {
        let stats = ZeroCopyStats::new();
        stats.record_array_ref_hit();
        stats.record_array_ref_hit();
        stats.record_array_ref_hit();
        assert_eq!(stats.array_ref_hits(), 3);
    }

    #[test]
    fn test_object_ref_hits_tracking() {
        let stats = ZeroCopyStats::new();
        stats.record_object_ref_hit();
        stats.record_object_ref_hit();
        assert_eq!(stats.object_ref_hits(), 2);
    }

    #[test]
    fn test_array_heap_fallbacks_tracking() {
        let stats = ZeroCopyStats::new();
        stats.record_array_heap_fallback();
        assert_eq!(stats.array_heap_fallbacks(), 1);
    }

    #[test]
    fn test_object_heap_fallbacks_tracking() {
        let stats = ZeroCopyStats::new();
        stats.record_object_heap_fallback();
        stats.record_object_heap_fallback();
        assert_eq!(stats.object_heap_fallbacks(), 2);
    }

    #[test]
    fn test_heap_reduction_rate_empty() {
        let stats = ZeroCopyStats::new();
        assert_eq!(stats.heap_reduction_rate(), 0.0);
    }

    #[test]
    fn test_heap_reduction_rate_all_hits() {
        let stats = ZeroCopyStats::new();
        for _ in 0..8 {
            stats.record_array_ref_hit();
        }
        for _ in 0..2 {
            stats.record_object_ref_hit();
        }
        assert!((stats.heap_reduction_rate() - 1.0).abs() < 1e-9);
    }

    #[test]
    fn test_heap_reduction_rate_mixed() {
        let stats = ZeroCopyStats::new();
        for _ in 0..8 {
            stats.record_array_ref_hit();
        }
        for _ in 0..2 {
            stats.record_array_heap_fallback();
        }
        let rate = stats.heap_reduction_rate();
        assert!((rate - 0.8).abs() < 1e-9, "rate={}", rate);
    }

    #[test]
    fn test_heap_reduction_rate_above_80_pct() {
        let stats = ZeroCopyStats::new();
        for _ in 0..85 {
            stats.record_array_ref_hit();
        }
        for _ in 0..15 {
            stats.record_array_heap_fallback();
        }
        assert!(stats.heap_reduction_rate() >= 0.80);
    }

    #[test]
    fn test_zero_copy_stats_clone_preserves_v760_fields() {
        let stats = ZeroCopyStats::new();
        stats.record_array_ref_hit();
        stats.record_object_ref_hit();
        stats.record_array_heap_fallback();
        stats.record_object_heap_fallback();
        let cloned = stats.clone();
        assert_eq!(cloned.array_ref_hits(), 1);
        assert_eq!(cloned.object_ref_hits(), 1);
        assert_eq!(cloned.array_heap_fallbacks(), 1);
        assert_eq!(cloned.object_heap_fallbacks(), 1);
    }
}