Skip to main content

copybook_rdw/
reader.rs

1use crate::{
2    RDW_HEADER_LEN, RDWRecord, rdw_is_suspect_ascii_corruption, rdw_read_len, rdw_slice_body,
3    rdw_try_peek_len, rdw_validate_and_finish, schema_prefix::calculate_schema_fixed_prefix,
4};
5use copybook_core::Schema;
6use copybook_error::{Error, ErrorCode, ErrorContext, Result};
7use std::io::{BufRead, BufReader, Read};
8use tracing::{debug, warn};
9
10const RDW_READER_BUF_CAPACITY: usize = (u16::MAX as usize) + RDW_HEADER_LEN;
11
12/// RDW (Record Descriptor Word) record reader for variable-length records.
13#[derive(Debug)]
14pub struct RDWRecordReader<R: Read> {
15    input: BufReader<R>,
16    record_count: u64,
17    strict_mode: bool,
18}
19
20impl<R: Read> RDWRecordReader<R> {
21    /// Create a new RDW record reader.
22    #[inline]
23    #[must_use]
24    pub fn new(input: R, strict_mode: bool) -> Self {
25        Self {
26            input: BufReader::with_capacity(RDW_READER_BUF_CAPACITY, input),
27            record_count: 0,
28            strict_mode,
29        }
30    }
31
32    #[inline]
33    fn peek_header(&mut self) -> Result<Option<[u8; RDW_HEADER_LEN]>> {
34        let peek = rdw_try_peek_len(&mut self.input).map_err(|error| {
35            error.with_context(ErrorContext {
36                record_index: Some(self.record_count + 1),
37                field_path: None,
38                byte_offset: Some(0),
39                line_number: None,
40                details: Some("Unable to peek RDW header".to_string()),
41            })
42        })?;
43
44        if peek.is_none() {
45            return self.handle_short_or_empty_header();
46        }
47
48        let buf = self.input.fill_buf().map_err(|e| {
49            Error::new(
50                ErrorCode::CBKF104_RDW_SUSPECT_ASCII,
51                format!("I/O error reading RDW header: {e}"),
52            )
53            .with_context(ErrorContext {
54                record_index: Some(self.record_count + 1),
55                field_path: None,
56                byte_offset: Some(0),
57                line_number: None,
58                details: Some("Unable to read RDW header".to_string()),
59            })
60        })?;
61
62        if buf.len() < RDW_HEADER_LEN {
63            return self.handle_truncated_header();
64        }
65
66        Ok(Some([buf[0], buf[1], buf[2], buf[3]]))
67    }
68
69    fn handle_short_or_empty_header(&mut self) -> Result<Option<[u8; RDW_HEADER_LEN]>> {
70        let context = self.header_context("Unable to read RDW header");
71        let buf = self.input.fill_buf().map_err(|e| {
72            Error::new(
73                ErrorCode::CBKF104_RDW_SUSPECT_ASCII,
74                format!("I/O error reading RDW header: {e}"),
75            )
76            .with_context(context)
77        })?;
78
79        if buf.is_empty() {
80            debug!("Reached EOF after {} RDW records", self.record_count);
81            return Ok(None);
82        }
83
84        self.handle_truncated_header()
85    }
86
87    fn handle_truncated_header(&mut self) -> Result<Option<[u8; RDW_HEADER_LEN]>> {
88        if self.strict_mode {
89            return Err(Error::new(
90                ErrorCode::CBKF221_RDW_UNDERFLOW,
91                "Incomplete RDW header: expected 4 bytes".to_string(),
92            )
93            .with_context(self.header_context("File ends with incomplete RDW header")));
94        }
95
96        debug!(
97            "Reached EOF after {} RDW records (truncated header ignored)",
98            self.record_count
99        );
100        let context = self.header_context("Unable to read RDW header");
101        let remaining = self
102            .input
103            .fill_buf()
104            .map_err(|e| {
105                Error::new(
106                    ErrorCode::CBKF104_RDW_SUSPECT_ASCII,
107                    format!("I/O error reading RDW header: {e}"),
108                )
109                .with_context(context)
110            })?
111            .len();
112        self.input.consume(remaining);
113        Ok(None)
114    }
115
116    fn header_context(&self, details: impl Into<String>) -> ErrorContext {
117        ErrorContext {
118            record_index: Some(self.record_count + 1),
119            field_path: None,
120            byte_offset: Some(0),
121            line_number: None,
122            details: Some(details.into()),
123        }
124    }
125
126    /// Read the next RDW record.
127    ///
128    /// # Errors
129    /// Returns an error if the record cannot be read due to I/O errors or
130    /// framing issues.
131    #[inline]
132    #[must_use = "Handle the Result or propagate the error"]
133    pub fn read_record(&mut self) -> Result<Option<RDWRecord>> {
134        let Some(header) = self.peek_header()? else {
135            return Ok(None);
136        };
137
138        let length = match rdw_read_len(&mut self.input) {
139            Ok(len) => len,
140            Err(error) => {
141                return Err(error.with_context(ErrorContext {
142                    record_index: Some(self.record_count + 1),
143                    field_path: None,
144                    byte_offset: Some(0),
145                    line_number: None,
146                    details: Some("Unable to read RDW body length".to_string()),
147                }));
148            }
149        };
150
151        // Consume reserved bytes so the buffer now points at the body.
152        self.input.consume(2);
153        let reserved = u16::from_be_bytes([header[2], header[3]]);
154
155        self.record_count += 1;
156        debug!(
157            "Read RDW header for record {}: length={}, reserved={:04X}",
158            self.record_count,
159            u32::from(length),
160            reserved
161        );
162
163        self.validate_reserved(reserved)?;
164        self.reject_ascii_corruption(header)?;
165
166        if length == 0 {
167            debug!("Zero-length RDW record {}", self.record_count);
168            return Ok(Some(RDWRecord {
169                header,
170                payload: Vec::new(),
171            }));
172        }
173
174        let payload = self.read_payload(length)?;
175
176        debug!(
177            "Read RDW record {} payload: {} bytes",
178            self.record_count, length
179        );
180        Ok(Some(RDWRecord { header, payload }))
181    }
182
183    fn validate_reserved(&self, reserved: u16) -> Result<()> {
184        if reserved == 0 {
185            return Ok(());
186        }
187
188        let error = Error::new(
189            ErrorCode::CBKR211_RDW_RESERVED_NONZERO,
190            format!("RDW reserved bytes are non-zero: {reserved:04X}"),
191        )
192        .with_context(ErrorContext {
193            record_index: Some(self.record_count),
194            field_path: None,
195            byte_offset: Some(2),
196            line_number: None,
197            details: Some(format!("Expected 0000, got {reserved:04X}")),
198        });
199
200        if self.strict_mode {
201            return Err(error);
202        }
203
204        warn!(
205            "RDW reserved bytes non-zero (record {}): {:04X}",
206            self.record_count, reserved
207        );
208        Ok(())
209    }
210
211    fn reject_ascii_corruption(&self, header: [u8; RDW_HEADER_LEN]) -> Result<()> {
212        if !Self::is_suspect_ascii_corruption(header) {
213            return Ok(());
214        }
215
216        warn!(
217            "RDW appears to be ASCII-corrupted (record {}): {:02X} {:02X} {:02X} {:02X}",
218            self.record_count, header[0], header[1], header[2], header[3]
219        );
220
221        Err(Error::new(
222            ErrorCode::CBKF104_RDW_SUSPECT_ASCII,
223            format!(
224                "RDW appears to be ASCII-corrupted: {:02X} {:02X} {:02X} {:02X}",
225                header[0], header[1], header[2], header[3]
226            ),
227        )
228        .with_context(ErrorContext {
229            record_index: Some(self.record_count),
230            field_path: None,
231            byte_offset: Some(0),
232            line_number: None,
233            details: Some("Suspected ASCII transfer corruption".to_string()),
234        }))
235    }
236
237    fn read_payload(&mut self, length: u16) -> Result<Vec<u8>> {
238        let payload_len = usize::from(length);
239        let body_slice = match rdw_slice_body(&mut self.input, length) {
240            Ok(slice) => slice,
241            Err(error) => {
242                return Err(error.with_context(ErrorContext {
243                    record_index: Some(self.record_count),
244                    field_path: None,
245                    byte_offset: Some(4),
246                    line_number: None,
247                    details: Some("File ends with incomplete RDW payload".to_string()),
248                }));
249            }
250        };
251
252        let payload = rdw_validate_and_finish(body_slice).to_vec();
253        self.input.consume(payload_len);
254        Ok(payload)
255    }
256
257    /// Validate a zero-length record against schema requirements.
258    ///
259    /// # Errors
260    /// Returns an error when the schema requires non-zero fixed bytes.
261    #[inline]
262    #[must_use = "Handle the Result or propagate the error"]
263    pub fn validate_zero_length_record(&self, schema: &Schema) -> Result<()> {
264        let min_size = Self::calculate_schema_fixed_prefix(schema);
265
266        if min_size > 0 {
267            return Err(Error::new(
268                ErrorCode::CBKF221_RDW_UNDERFLOW,
269                format!("Zero-length RDW record invalid: schema requires minimum {min_size} bytes"),
270            )
271            .with_context(ErrorContext {
272                record_index: Some(self.record_count),
273                field_path: None,
274                byte_offset: None,
275                line_number: None,
276                details: Some("Zero-length record with non-zero schema prefix".to_string()),
277            }));
278        }
279
280        Ok(())
281    }
282
283    /// Number of RDW records consumed from the stream.
284    #[inline]
285    #[must_use]
286    pub fn record_count(&self) -> u64 {
287        self.record_count
288    }
289
290    #[inline]
291    fn calculate_schema_fixed_prefix(schema: &Schema) -> u32 {
292        calculate_schema_fixed_prefix(schema)
293    }
294
295    #[inline]
296    fn is_suspect_ascii_corruption(rdw_header: [u8; RDW_HEADER_LEN]) -> bool {
297        rdw_is_suspect_ascii_corruption(rdw_header)
298    }
299}