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#[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 #[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 #[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 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 #[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}