1use crate::{
2 compression::{Compression, Compressors, Zstd},
3 DataReader, NippyJar, NippyJarError, NippyJarHeader, RefRow,
4};
5use smallvec::SmallVec;
6use std::{ops::Range, sync::Arc};
7use zstd::bulk::Decompressor;
8
9type ValueRanges = SmallVec<[ValueRange; 4]>;
16
17#[derive(Clone)]
19pub struct NippyJarCursor<'a, H = ()> {
20 jar: &'a NippyJar<H>,
22 reader: Arc<DataReader>,
24 internal_buffer: Vec<u8>,
27 row: u64,
29}
30
31impl<H: NippyJarHeader> std::fmt::Debug for NippyJarCursor<'_, H> {
32 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
33 f.debug_struct("NippyJarCursor").field("config", &self.jar).finish_non_exhaustive()
34 }
35}
36
37impl<'a, H: NippyJarHeader> NippyJarCursor<'a, H> {
38 pub fn new(jar: &'a NippyJar<H>) -> Result<Self, NippyJarError> {
40 Ok(Self {
41 jar,
42 reader: Arc::new(jar.open_data_reader()?),
43 internal_buffer: Vec::new(),
44 row: 0,
45 })
46 }
47
48 pub const fn with_reader(
51 jar: &'a NippyJar<H>,
52 reader: Arc<DataReader>,
53 ) -> Result<Self, NippyJarError> {
54 Ok(Self { jar, reader, internal_buffer: Vec::new(), row: 0 })
55 }
56
57 pub const fn jar(&self) -> &NippyJar<H> {
59 self.jar
60 }
61
62 pub const fn row_index(&self) -> u64 {
64 self.row
65 }
66
67 pub const fn reset(&mut self) {
69 self.row = 0;
70 }
71
72 pub fn row_by_number(&mut self, row: usize) -> Result<Option<RefRow<'_>>, NippyJarError> {
74 self.row = row as u64;
75 self.next_row()
76 }
77
78 pub fn next_row(&mut self) -> Result<Option<RefRow<'_>>, NippyJarError> {
80 self.internal_buffer.clear();
81
82 if self.row as usize >= self.jar.rows {
83 return Ok(None)
85 }
86
87 let mut row = ValueRanges::with_capacity(self.jar.columns);
88
89 for column in 0..self.jar.columns {
91 self.read_value(column, &mut row)?;
92 }
93
94 self.row += 1;
95
96 Ok(Some(
97 row.into_iter()
98 .map(|v| match v {
99 ValueRange::Mmap(range) => self.reader.data(range),
100 ValueRange::Internal(range) => &self.internal_buffer[range],
101 })
102 .collect(),
103 ))
104 }
105
106 pub fn row_by_number_with_cols(
108 &mut self,
109 row: usize,
110 mask: usize,
111 ) -> Result<Option<RefRow<'_>>, NippyJarError> {
112 self.row = row as u64;
113 self.next_row_with_cols(mask)
114 }
115
116 pub fn next_row_with_cols(&mut self, mask: usize) -> Result<Option<RefRow<'_>>, NippyJarError> {
120 self.internal_buffer.clear();
121
122 if self.row as usize >= self.jar.rows {
123 return Ok(None)
125 }
126
127 let columns = self.jar.columns;
128 let mut row = ValueRanges::with_capacity(columns);
129
130 for column in 0..columns {
131 if mask & (1 << column) != 0 {
132 self.read_value(column, &mut row)?
133 }
134 }
135 self.row += 1;
136
137 Ok(Some(
138 row.into_iter()
139 .map(|v| match v {
140 ValueRange::Mmap(range) => self.reader.data(range),
141 ValueRange::Internal(range) => &self.internal_buffer[range],
142 })
143 .collect(),
144 ))
145 }
146
147 fn read_value(&mut self, column: usize, row: &mut ValueRanges) -> Result<(), NippyJarError> {
149 let offset_pos = self.row as usize * self.jar.columns + column;
151 let value_offset = self.reader.offset(offset_pos)? as usize;
152
153 let column_offset_range = if self.jar.rows * self.jar.columns == offset_pos + 1 {
154 value_offset..self.reader.size()
156 } else {
157 let next_value_offset = self.reader.offset(offset_pos + 1)? as usize;
158 value_offset..next_value_offset
159 };
160
161 if let Some(compression) = self.jar.compressor() {
162 if self.internal_buffer.capacity() < self.jar.max_row_size {
165 self.internal_buffer.reserve(self.jar.max_row_size - self.internal_buffer.len());
166 }
167
168 let from = self.internal_buffer.len();
169 match compression {
170 Compressors::Zstd(z) if z.use_dict => {
171 let dictionaries = z.dictionaries.as_ref().expect("dictionaries to exist")
175 [column]
176 .loaded()
177 .expect("dictionary to be loaded");
178 let mut decompressor = Decompressor::with_prepared_dictionary(dictionaries)?;
179 Zstd::decompress_with_dictionary(
180 self.reader.data(column_offset_range),
181 &mut self.internal_buffer,
182 &mut decompressor,
183 )?;
184 }
185 _ => {
186 compression.decompress_to(
188 self.reader.data(column_offset_range),
189 &mut self.internal_buffer,
190 )?;
191 }
192 }
193 let to = self.internal_buffer.len();
194
195 row.push(ValueRange::Internal(from..to));
196 } else {
197 row.push(ValueRange::Mmap(column_offset_range));
199 }
200
201 Ok(())
202 }
203}
204
205enum ValueRange {
208 Mmap(Range<usize>),
209 Internal(Range<usize>),
210}