1use copybook_error::{Error, ErrorCode, ErrorContext, Result};
3use std::io::{BufReader, Read, Write};
4use tracing::{debug, warn};
5
6pub const BDW_HEADER_LEN: usize = 4;
8
9pub const BDW_MAX_BLOCK_LEN: usize = 32760;
11
12pub const VB_MAX_RECORD_LEN: usize = BDW_MAX_BLOCK_LEN - BDW_HEADER_LEN;
14
15#[derive(Debug, Clone, Copy, PartialEq, Eq)]
21pub struct BdwHeader {
22 bytes: [u8; BDW_HEADER_LEN],
23}
24
25impl BdwHeader {
26 #[must_use]
28 #[inline]
29 pub const fn from_bytes(bytes: [u8; BDW_HEADER_LEN]) -> Self {
30 Self { bytes }
31 }
32
33 #[inline]
39 #[must_use = "Handle the Result or propagate the error"]
40 pub fn from_block_len(block_len: usize) -> Result<Self> {
41 if !(BDW_HEADER_LEN..=BDW_MAX_BLOCK_LEN).contains(&block_len) {
42 return Err(Error::new(
43 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
44 format!(
45 "BDW block length out of bounds: {block_len} bytes (valid {BDW_HEADER_LEN}..={BDW_MAX_BLOCK_LEN})"
46 ),
47 ));
48 }
49 let len = u16::try_from(block_len).map_err(|_| {
50 Error::new(
51 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
52 format!("BDW block length out of bounds: {block_len} bytes"),
53 )
54 })?;
55 let len_bytes = len.to_be_bytes();
56 Ok(Self {
57 bytes: [len_bytes[0], len_bytes[1], 0, 0],
58 })
59 }
60
61 #[must_use]
63 #[inline]
64 pub const fn bytes(self) -> [u8; BDW_HEADER_LEN] {
65 self.bytes
66 }
67
68 #[must_use]
70 #[inline]
71 pub const fn length(self) -> u16 {
72 u16::from_be_bytes([self.bytes[0], self.bytes[1]])
73 }
74
75 #[must_use]
77 #[inline]
78 pub const fn reserved(self) -> u16 {
79 u16::from_be_bytes([self.bytes[2], self.bytes[3]])
80 }
81}
82
83#[derive(Debug, Clone, PartialEq, Eq)]
85pub struct VbRecord {
86 pub payload: Vec<u8>,
88 pub rdw: [u8; 4],
90 pub rdw_reserved: u16,
92 pub block_index: u64,
94 pub record_index_in_block: u32,
96 pub block_offset: u64,
98 pub physical_offset: u64,
100 pub block_len: usize,
102}
103
104#[derive(Debug)]
110pub struct VbBlockReader<R: Read> {
111 input: BufReader<R>,
112 strict_mode: bool,
113 block_index: u64,
114 record_count: u64,
115 physical_offset: u64,
116 block_remaining: usize,
117 block_offset: u64,
118 block_len: usize,
119 block_record_index: u32,
120 in_block: bool,
121}
122
123impl<R: Read> VbBlockReader<R> {
124 #[inline]
126 #[must_use]
127 pub fn new(input: R, strict_mode: bool) -> Self {
128 Self {
129 input: BufReader::with_capacity(BDW_MAX_BLOCK_LEN, input),
130 strict_mode,
131 block_index: 0,
132 record_count: 0,
133 physical_offset: 0,
134 block_remaining: 0,
135 block_offset: 0,
136 block_len: 0,
137 block_record_index: 0,
138 in_block: false,
139 }
140 }
141
142 #[inline]
148 #[must_use = "Handle the Result or propagate the error"]
149 pub fn read_record(&mut self) -> Result<Option<VbRecord>> {
150 loop {
151 if !self.in_block {
152 if !self.open_block()? {
153 return Ok(None);
154 }
155 if self.block_remaining == 0 {
156 self.in_block = false;
158 self.block_index += 1;
159 continue;
160 }
161 }
162 return self.read_block_record().map(Some);
163 }
164 }
165
166 fn open_block(&mut self) -> Result<bool> {
172 let mut first = [0u8; 1];
173 let read = self.input.read(&mut first).map_err(|error| {
174 Error::new(
175 ErrorCode::CBKR201_RDW_READ_ERROR,
176 format!("I/O error probing VB input: {error}"),
177 )
178 })?;
179 if read == 0 {
180 debug!("Reached EOF after {} VB blocks", self.block_index);
181 return Ok(false);
182 }
183 let mut header = [0u8; BDW_HEADER_LEN];
184 header[0] = first[0];
185 let mut have = 1usize;
186 while have < BDW_HEADER_LEN {
187 match self.input.read(&mut header[have..]) {
188 Ok(0) => break,
189 Ok(read_now) => have += read_now,
190 Err(error) => {
191 return Err(Error::new(
192 ErrorCode::CBKR201_RDW_READ_ERROR,
193 format!("I/O error reading BDW header: {error}"),
194 )
195 .with_context(self.block_context("Unable to read BDW header")));
196 }
197 }
198 }
199 if have < BDW_HEADER_LEN {
200 return self.short_block_header(have);
201 }
202 let parsed = BdwHeader::from_bytes(header);
203 let block_len = usize::from(parsed.length());
204 if !(BDW_HEADER_LEN..=BDW_MAX_BLOCK_LEN).contains(&block_len) {
205 return Err(Error::new(
206 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
207 format!(
208 "BDW block length out of bounds: {block_len} bytes (valid {BDW_HEADER_LEN}..={BDW_MAX_BLOCK_LEN})"
209 ),
210 )
211 .with_context(self.block_context("Invalid BDW block length")));
212 }
213 self.check_bdw_reserved(parsed.reserved())?;
214 self.block_offset = self.physical_offset;
215 self.physical_offset += BDW_HEADER_LEN as u64;
216 self.block_len = block_len;
217 self.block_remaining = block_len - BDW_HEADER_LEN;
218 self.block_record_index = 0;
219 self.in_block = true;
220 debug!("Opened VB block {}: length={block_len}", self.block_index);
221 Ok(true)
222 }
223
224 fn short_block_header(&mut self, available: usize) -> Result<bool> {
226 if self.strict_mode {
227 return Err(Error::new(
228 ErrorCode::CBKF223_BDW_UNDERFLOW,
229 format!("Truncated BDW header: {available} trailing bytes cannot form a block"),
230 )
231 .with_context(self.block_context("Trailing bytes after final VB block")));
232 }
233 warn!("Ignoring {available} trailing bytes after final VB block");
234 Ok(false)
235 }
236
237 fn check_bdw_reserved(&self, reserved: u16) -> Result<()> {
238 if reserved == 0 {
239 return Ok(());
240 }
241 let error = Error::new(
242 ErrorCode::CBKF225_BDW_RESERVED_NONZERO,
243 format!("BDW reserved bytes are non-zero: {reserved:04X}"),
244 )
245 .with_context(ErrorContext {
246 record_index: None,
247 field_path: None,
248 byte_offset: Some(self.physical_offset + 2),
249 line_number: None,
250 details: Some(format!("Expected 0000, got {reserved:04X}")),
251 });
252 if self.strict_mode {
253 return Err(error);
254 }
255 warn!(
256 "BDW reserved bytes non-zero (block {}): {:04X}",
257 self.block_index, reserved
258 );
259 Ok(())
260 }
261
262 fn read_block_record(&mut self) -> Result<VbRecord> {
263 let rdw_offset = self.physical_offset;
264 if self.block_remaining < 4 {
268 return Err(Error::new(
269 ErrorCode::CBKF223_BDW_UNDERFLOW,
270 format!(
271 "VB block {} ends with {} stray bytes; cannot form an RDW header",
272 self.block_index, self.block_remaining
273 ),
274 )
275 .with_context(self.record_context("Block ends mid-RDW header")));
276 }
277 let mut rdw = [0u8; 4];
278 self.input.read_exact(&mut rdw).map_err(|_| {
279 Error::new(
280 ErrorCode::CBKF223_BDW_UNDERFLOW,
281 format!(
282 "Truncated RDW inside VB block {}: {} block bytes remain",
283 self.block_index, self.block_remaining
284 ),
285 )
286 .with_context(self.record_context("Block ends mid-RDW header"))
287 })?;
288 let record_len = usize::from(u16::from_be_bytes([rdw[0], rdw[1]]));
289 if !(4..=VB_MAX_RECORD_LEN).contains(&record_len) {
290 return Err(Error::new(
291 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
292 format!("RDW length out of bounds inside VB block: {record_len} bytes"),
293 )
294 .with_context(self.record_context("Invalid nested RDW length")));
295 }
296 if record_len > self.block_remaining {
297 return Err(Error::new(
298 ErrorCode::CBKF224_RDW_BEYOND_BLOCK,
299 format!(
300 "RDW of {record_len} bytes escapes VB block {} ({} bytes remain)",
301 self.block_index, self.block_remaining
302 ),
303 )
304 .with_context(self.record_context("Nested RDW beyond block end")));
305 }
306 let payload_len = record_len - 4;
307 let mut payload = vec![0u8; payload_len];
308 self.input.read_exact(&mut payload).map_err(|_| {
309 Error::new(
310 ErrorCode::CBKF223_BDW_UNDERFLOW,
311 format!("Truncated RDW payload inside VB block {}", self.block_index),
312 )
313 .with_context(self.record_context("Block ends mid-RDW payload"))
314 })?;
315 self.block_remaining -= record_len;
316 self.physical_offset += record_len as u64;
317 let record = VbRecord {
318 payload,
319 rdw,
320 rdw_reserved: u16::from_be_bytes([rdw[2], rdw[3]]),
321 block_index: self.block_index,
322 record_index_in_block: self.block_record_index,
323 block_offset: self.block_offset,
324 physical_offset: rdw_offset,
325 block_len: self.block_len,
326 };
327 self.block_record_index += 1;
328 self.record_count += 1;
329 if self.block_remaining == 0 {
330 self.in_block = false;
331 self.block_index += 1;
332 }
333 Ok(record)
334 }
335
336 fn block_context(&self, details: impl Into<String>) -> ErrorContext {
337 ErrorContext {
338 record_index: None,
339 field_path: None,
340 byte_offset: Some(self.physical_offset),
341 line_number: None,
342 details: Some(details.into()),
343 }
344 }
345
346 fn record_context(&self, details: impl Into<String>) -> ErrorContext {
347 ErrorContext {
348 record_index: Some(self.record_count + 1),
349 field_path: None,
350 byte_offset: Some(self.physical_offset),
351 line_number: None,
352 details: Some(details.into()),
353 }
354 }
355
356 #[inline]
358 #[must_use]
359 pub const fn block_count(&self) -> u64 {
360 self.block_index
361 }
362
363 #[inline]
365 #[must_use]
366 pub const fn record_count(&self) -> u64 {
367 self.record_count
368 }
369
370 #[inline]
372 #[must_use]
373 pub const fn physical_bytes(&self) -> u64 {
374 self.physical_offset
375 }
376}
377
378#[derive(Debug)]
383pub struct VbBlockWriter<W: Write> {
384 output: W,
385 pending: Vec<u8>,
386 block_count: u64,
387 record_count: u64,
388 physical_bytes: u64,
389}
390
391impl<W: Write> VbBlockWriter<W> {
392 #[inline]
394 #[must_use]
395 pub fn new(output: W) -> Self {
396 Self {
397 output,
398 pending: Vec::new(),
399 block_count: 0,
400 record_count: 0,
401 physical_bytes: 0,
402 }
403 }
404
405 #[inline]
411 #[must_use = "Handle the Result or propagate the error"]
412 pub fn write_record_from_payload(&mut self, payload: &[u8], reserved: u16) -> Result<()> {
413 let record_len = payload
414 .len()
415 .checked_add(4)
416 .filter(|len| *len <= VB_MAX_RECORD_LEN)
417 .ok_or_else(|| {
418 Error::new(
419 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
420 format!(
421 "VB record too large: {} payload bytes exceed maximum of {VB_MAX_RECORD_LEN}",
422 payload.len()
423 ),
424 )
425 })?;
426 if !self.pending.is_empty()
427 && self.pending.len() + BDW_HEADER_LEN + record_len > BDW_MAX_BLOCK_LEN
428 {
429 self.flush_block()?;
430 }
431 let len_bytes = u16::try_from(record_len)
432 .map_err(|_| {
433 Error::new(
434 ErrorCode::CBKF222_BDW_LENGTH_INVALID,
435 format!("VB record length out of range: {record_len} bytes"),
436 )
437 })?
438 .to_be_bytes();
439 let reserved_bytes = reserved.to_be_bytes();
440 self.pending.extend_from_slice(&[
441 len_bytes[0],
442 len_bytes[1],
443 reserved_bytes[0],
444 reserved_bytes[1],
445 ]);
446 self.pending.extend_from_slice(payload);
447 self.record_count += 1;
448 Ok(())
449 }
450
451 #[inline]
456 #[must_use = "Handle the Result or propagate the error"]
457 pub fn finish(&mut self) -> Result<()> {
458 if !self.pending.is_empty() {
459 self.flush_block()?;
460 }
461 self.output.flush().map_err(|error| {
462 Error::new(
463 ErrorCode::CBKR202_RDW_WRITE_ERROR,
464 format!("I/O error finishing VB blocks: {error}"),
465 )
466 })
467 }
468
469 fn flush_block(&mut self) -> Result<()> {
470 let block_len = self.pending.len() + BDW_HEADER_LEN;
471 let header = BdwHeader::from_block_len(block_len)?;
472 self.output.write_all(&header.bytes()).map_err(|error| {
473 Error::new(
474 ErrorCode::CBKR202_RDW_WRITE_ERROR,
475 format!("I/O error writing VB block header: {error}"),
476 )
477 })?;
478 self.output.write_all(&self.pending).map_err(|error| {
479 Error::new(
480 ErrorCode::CBKR202_RDW_WRITE_ERROR,
481 format!("I/O error writing VB block body: {error}"),
482 )
483 })?;
484 self.physical_bytes += block_len as u64;
485 self.block_count += 1;
486 self.pending.clear();
487 Ok(())
488 }
489
490 #[inline]
492 #[must_use]
493 pub const fn block_count(&self) -> u64 {
494 self.block_count
495 }
496
497 #[inline]
499 #[must_use]
500 pub const fn record_count(&self) -> u64 {
501 self.record_count
502 }
503
504 #[inline]
506 #[must_use]
507 pub const fn physical_bytes(&self) -> u64 {
508 self.physical_bytes
509 }
510}