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