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