1#![allow(clippy::missing_inline_in_public_items)]
11
12use crate::options::{DecodeOptions, EncodeOptions, RecordFormat, ZonedEncodingFormat};
13use crate::zoned_overpunch::ZeroSignPolicy;
14use base64::Engine;
15use copybook_core::{Error, ErrorCode, Result, Schema};
16use serde_json::Value;
17use std::convert::TryFrom;
18use std::io::{BufRead, BufReader, Read, Write};
19use std::sync::Arc;
20use tracing::info;
21
22mod envelope;
23mod run_summary;
24mod telemetry;
25mod warnings;
26
27use envelope::build_json_envelope;
28pub use run_summary::RunSummary;
29pub use warnings::increment_warning_counter;
30use warnings::{reset_warning_counter, warning_count};
31
32const MAX_WORKERS: usize = 64;
33
34#[inline]
44#[must_use = "Handle the Result or propagate the error"]
45pub fn decode_record(schema: &Schema, data: &[u8], options: &DecodeOptions) -> Result<Value> {
46 decode_record_with_raw_data(schema, data, options, None, 0)
47}
48
49#[inline]
85#[must_use = "Handle the Result or propagate the error"]
86pub fn decode_record_with_scratch(
87 schema: &Schema,
88 data: &[u8],
89 options: &DecodeOptions,
90 scratch: &mut crate::memory::ScratchBuffers,
91) -> Result<Value> {
92 decode_record_with_scratch_and_raw(schema, data, options, None, 0, scratch)
93}
94
95fn decode_record_with_scratch_and_raw(
97 schema: &Schema,
98 data: &[u8],
99 options: &DecodeOptions,
100 raw_data: Option<Vec<u8>>,
101 record_index: u64,
102 scratch: &mut crate::memory::ScratchBuffers,
103) -> Result<Value> {
104 use serde_json::Map;
105
106 let mut fields_map = Map::new();
107 let mut record_raw = None;
108 let mut encoding_acc = Vec::new();
109
110 if let Some(raw_bytes) = raw_data.filter(|_| {
111 matches!(
112 options.emit_raw,
113 crate::options::RawMode::Record | crate::options::RawMode::RecordRDW
114 )
115 }) {
116 record_raw = Some(base64::engine::general_purpose::STANDARD.encode(raw_bytes));
117 }
118
119 process_fields_recursive_with_scratch(
120 &schema.fields,
121 data,
122 &mut fields_map,
123 options,
124 scratch,
125 record_index,
126 &mut encoding_acc,
127 )?;
128
129 Ok(build_json_envelope(
130 fields_map,
131 schema,
132 options,
133 record_index,
134 data.len(),
135 record_raw,
136 encoding_acc,
137 ))
138}
139
140#[inline]
145#[must_use = "Handle the Result or propagate the error"]
146pub fn decode_record_with_raw_data(
147 schema: &Schema,
148 data: &[u8],
149 options: &DecodeOptions,
150 raw_data_with_header: Option<&[u8]>,
151 record_index: u64,
152) -> Result<Value> {
153 use crate::options::RawMode;
154 use serde_json::Map;
155
156 let mut fields_map = Map::new();
157 let mut scratch_buffers: Option<crate::memory::ScratchBuffers> = None;
158 let mut encoding_acc = Vec::new();
159
160 process_fields_recursive(
161 &schema.fields,
162 data,
163 &mut fields_map,
164 options,
165 &mut scratch_buffers,
166 record_index,
167 &mut encoding_acc,
168 )?;
169
170 let mut record_raw = None;
171 match options.emit_raw {
172 RawMode::Off | RawMode::Field => {}
173 RawMode::Record => {
174 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(data);
175 record_raw = Some(raw_b64);
176 }
177 RawMode::RecordRDW => {
178 if let Some(full_raw) = raw_data_with_header {
179 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(full_raw);
180 record_raw = Some(raw_b64);
181 } else {
182 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(data);
183 record_raw = Some(raw_b64);
184 }
185 }
186 }
187
188 Ok(build_json_envelope(
189 fields_map,
190 schema,
191 options,
192 record_index,
193 data.len(),
194 record_raw,
195 encoding_acc,
196 ))
197}
198
199fn process_fields_recursive(
204 fields: &[copybook_core::Field],
205 data: &[u8],
206 json_obj: &mut serde_json::Map<String, Value>,
207 options: &DecodeOptions,
208 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
209 record_index: u64,
210 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
211) -> Result<()> {
212 use copybook_core::FieldKind;
213
214 let total_fields = fields.len();
215 let mut deferred_group_views = Vec::new();
216
217 for (field_index, field) in fields.iter().enumerate() {
218 match (&field.kind, &field.occurs) {
219 (_, Some(occurs)) => {
220 process_array_field(
221 field,
222 occurs,
223 data,
224 json_obj,
225 options,
226 fields,
227 scratch_buffers,
228 record_index,
229 encoding_acc,
230 )?;
231 }
232 (FieldKind::Group, None) if field.level > 1 => {
233 let mut group_obj = serde_json::Map::new();
234 process_fields_recursive(
235 &field.children,
236 data,
237 &mut group_obj,
238 options,
239 scratch_buffers,
240 record_index,
241 encoding_acc,
242 )?;
243 if is_scalar_target_group_redefine(field, fields) {
244 let group_value = Value::Object(group_obj);
245 if let Value::Object(group_fields) = &group_value {
246 for (name, value) in group_fields {
247 json_obj.insert(name.clone(), value.clone());
248 }
249 }
250 deferred_group_views.push((field.name.clone(), group_value));
251 } else if field.redefines_of.is_none() {
252 json_obj.insert(field.name.clone(), Value::Object(group_obj));
253 }
254 }
255 (FieldKind::Group, None) => {
256 process_fields_recursive(
257 &field.children,
258 data,
259 json_obj,
260 options,
261 scratch_buffers,
262 record_index,
263 encoding_acc,
264 )?;
265 }
266 _ => {
267 process_scalar_field_standard(
268 field,
269 field_index,
270 total_fields,
271 data,
272 json_obj,
273 options,
274 scratch_buffers,
275 encoding_acc,
276 )?;
277 }
278 }
279 }
280
281 for (name, value) in deferred_group_views {
282 json_obj.insert(name, value);
283 }
284
285 Ok(())
286}
287
288fn process_fields_recursive_with_scratch(
291 fields: &[copybook_core::Field],
292 data: &[u8],
293 json_obj: &mut serde_json::Map<String, Value>,
294 options: &DecodeOptions,
295 scratch: &mut crate::memory::ScratchBuffers,
296 record_index: u64,
297 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
298) -> Result<()> {
299 use copybook_core::FieldKind;
300
301 let mut deferred_group_views = Vec::new();
302
303 for field in fields {
304 if is_filler_field(field) && !options.emit_filler {
305 continue;
306 }
307
308 match (&field.kind, &field.occurs) {
309 (_, Some(occurs)) => {
310 process_array_field_with_scratch(
311 field,
312 occurs,
313 data,
314 json_obj,
315 options,
316 fields,
317 scratch,
318 record_index,
319 encoding_acc,
320 )?;
321 }
322 (FieldKind::Group, None) if field.level > 1 => {
323 let mut group_obj = serde_json::Map::new();
324 process_fields_recursive_with_scratch(
325 &field.children,
326 data,
327 &mut group_obj,
328 options,
329 scratch,
330 record_index,
331 encoding_acc,
332 )?;
333 if is_scalar_target_group_redefine(field, fields) {
334 let group_value = Value::Object(group_obj);
335 if let Value::Object(group_fields) = &group_value {
336 for (name, value) in group_fields {
337 json_obj.insert(name.clone(), value.clone());
338 }
339 }
340 deferred_group_views.push((field.name.clone(), group_value));
341 } else if field.redefines_of.is_none() {
342 json_obj.insert(field.name.clone(), Value::Object(group_obj));
343 }
344 }
345 (FieldKind::Group, None) => {
346 process_fields_recursive_with_scratch(
347 &field.children,
348 data,
349 json_obj,
350 options,
351 scratch,
352 record_index,
353 encoding_acc,
354 )?;
355 }
356 _ => {
357 process_scalar_field_with_scratch(
358 field,
359 data,
360 json_obj,
361 options,
362 scratch,
363 encoding_acc,
364 )?;
365 }
366 }
367 }
368
369 for (name, value) in deferred_group_views {
370 json_obj.insert(name, value);
371 }
372
373 Ok(())
374}
375
376#[inline]
386#[allow(clippy::too_many_arguments)]
387fn process_scalar_field_standard(
388 field: ©book_core::Field,
389 field_index: usize,
390 total_fields: usize,
391 data: &[u8],
392 json_obj: &mut serde_json::Map<String, Value>,
393 options: &DecodeOptions,
394 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
395 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
396) -> Result<()> {
397 if matches!(field.kind, copybook_core::FieldKind::Renames { .. }) {
399 let Some(resolved) = &field.resolved_renames else {
400 return Err(Error::new(
401 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
402 format!(
403 "RENAMES field '{name}' has no resolved metadata",
404 name = field.name
405 ),
406 ));
407 };
408
409 let alias_start = resolved.offset as usize;
410 let alias_end = alias_start + resolved.length as usize;
411
412 if alias_end > data.len() {
413 return Err(Error::new(
414 ErrorCode::CBKD301_RECORD_TOO_SHORT,
415 format!(
416 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
417 name = field.name,
418 offset = resolved.offset,
419 length = resolved.length,
420 data_len = data.len()
421 ),
422 ));
423 }
424
425 let alias_data = &data[alias_start..alias_end];
426 let text = crate::charset::ebcdic_to_utf8(
427 alias_data,
428 options.codepage,
429 options.on_decode_unmappable,
430 )?;
431 json_obj.insert(field.name.clone(), Value::String(text));
432 return Ok(());
433 }
434
435 let field_start = field.offset as usize;
436 let mut field_end = field_start + field.len as usize;
437
438 if options.format == RecordFormat::RDW
439 && field_index + 1 == total_fields
440 && matches!(field.kind, copybook_core::FieldKind::Alphanum { .. })
441 && data.len() > field_end
442 {
443 field_end = data.len();
444 }
445
446 if field_start > data.len() {
447 return Err(Error::new(
448 ErrorCode::CBKD301_RECORD_TOO_SHORT,
449 format!(
450 "Field '{name}' starts beyond record boundary",
451 name = field.name
452 ),
453 ));
454 }
455
456 field_end = field_end.min(data.len());
457
458 if field_start >= field_end {
459 return Ok(());
460 }
461
462 let field_data = &data[field_start..field_end];
463 let value = decode_scalar_field_value_standard(field, field_data, options, scratch_buffers)?;
464
465 if options.preserve_zoned_encoding {
467 collect_zoned_encoding_info(field, field_data, options, encoding_acc);
468 }
469
470 json_obj.insert(field.name.clone(), value);
471
472 if matches!(options.emit_raw, crate::options::RawMode::Field) {
474 let raw_key = format!("{}_raw_b64", field.name);
475 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(field_data);
476 json_obj.insert(raw_key, Value::String(raw_b64));
477 }
478
479 Ok(())
480}
481
482#[inline]
486fn process_scalar_field_with_scratch(
487 field: ©book_core::Field,
488 data: &[u8],
489 json_obj: &mut serde_json::Map<String, Value>,
490 options: &DecodeOptions,
491 scratch: &mut crate::memory::ScratchBuffers,
492 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
493) -> Result<()> {
494 if matches!(field.kind, copybook_core::FieldKind::Renames { .. }) {
496 let Some(resolved) = &field.resolved_renames else {
497 return Err(Error::new(
498 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
499 format!(
500 "RENAMES field '{name}' has no resolved metadata",
501 name = field.name
502 ),
503 ));
504 };
505
506 let alias_start = resolved.offset as usize;
507 let alias_end = alias_start + resolved.length as usize;
508
509 if alias_end > data.len() {
510 return Err(Error::new(
511 ErrorCode::CBKD301_RECORD_TOO_SHORT,
512 format!(
513 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
514 name = field.name,
515 offset = resolved.offset,
516 length = resolved.length,
517 data_len = data.len()
518 ),
519 ));
520 }
521
522 let alias_data = &data[alias_start..alias_end];
523 let text = crate::charset::ebcdic_to_utf8(
524 alias_data,
525 options.codepage,
526 options.on_decode_unmappable,
527 )?;
528 json_obj.insert(field.name.clone(), Value::String(text));
529 return Ok(());
530 }
531
532 let field_start = field.offset as usize;
533 let mut field_end = field_start + field.len as usize;
534
535 if field_start > data.len() {
536 return Err(Error::new(
537 ErrorCode::CBKD301_RECORD_TOO_SHORT,
538 format!(
539 "Field '{name}' starts beyond record boundary",
540 name = field.name
541 ),
542 ));
543 }
544
545 if options.format == RecordFormat::RDW {
546 field_end = field_end.min(data.len());
547 }
548
549 if field_start >= field_end {
550 return Ok(());
551 }
552
553 if field_end > data.len() {
554 return Err(Error::new(
555 ErrorCode::CBKD301_RECORD_TOO_SHORT,
556 format!(
557 "Field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
558 name = field.name,
559 offset = field.offset,
560 length = field.len,
561 data_len = data.len()
562 ),
563 ));
564 }
565
566 let field_data = &data[field_start..field_end];
567 let value = decode_scalar_field_value_with_scratch(field, field_data, options, scratch)?;
568
569 if options.preserve_zoned_encoding {
571 collect_zoned_encoding_info(field, field_data, options, encoding_acc);
572 }
573
574 json_obj.insert(field.name.clone(), value);
575
576 if matches!(options.emit_raw, crate::options::RawMode::Field) {
578 let raw_key = format!("{}_raw_b64", field.name);
579 let raw_b64 = base64::engine::general_purpose::STANDARD.encode(field_data);
580 json_obj.insert(raw_key, Value::String(raw_b64));
581 }
582
583 Ok(())
584}
585
586#[allow(clippy::too_many_arguments)]
588fn process_array_field(
589 field: ©book_core::Field,
590 occurs: ©book_core::Occurs,
591 data: &[u8],
592 json_obj: &mut serde_json::Map<String, Value>,
593 options: &DecodeOptions,
594 all_fields: &[copybook_core::Field],
595 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
596 record_index: u64,
597 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
598) -> Result<()> {
599 use copybook_core::{FieldKind, Occurs};
600
601 let count = match occurs {
602 Occurs::Fixed { count } => *count,
603 Occurs::ODO {
604 min,
605 max,
606 counter_path,
607 } => {
608 let scratch = scratch_buffers.get_or_insert_with(crate::memory::ScratchBuffers::new);
610 let counter_value =
611 find_and_read_counter_field(counter_path, all_fields, data, options, scratch)?;
612
613 let counter_field = find_field_by_path(all_fields, counter_path)?;
614 let validation_context = crate::odo_redefines::OdoValidationContext {
615 field_path: field.path.clone(),
616 counter_path: counter_path.clone(),
617 record_index,
618 byte_offset: u64::from(counter_field.offset),
619 };
620 let validation = crate::odo_redefines::validate_odo_decode(
621 counter_value,
622 *min,
623 *max,
624 &validation_context,
625 options,
626 )?;
627
628 if let Some(warning) = validation.warning {
629 tracing::warn!("{}", warning);
630 increment_warning_counter();
631 }
632
633 validation.actual_count
634 }
635 };
636
637 let element_size = field.len as usize;
638 let array_start = field.offset as usize;
639 let total_array_size = element_size * count as usize;
640 let array_end = array_start + total_array_size;
641
642 if array_end > data.len() {
644 return Err(Error::new(
645 ErrorCode::CBKD301_RECORD_TOO_SHORT,
646 format!(
647 "Array '{}' requires {} bytes but only {} bytes available",
648 field.name,
649 total_array_size,
650 data.len().saturating_sub(array_start)
651 ),
652 ));
653 }
654
655 let mut array_values = Vec::new();
657 for i in 0..count {
658 let element_start = array_start + (i as usize * element_size);
659 let element_end = element_start + element_size;
660
661 let element_value = match &field.kind {
662 FieldKind::Group => {
663 let mut element_obj = serde_json::Map::new();
665 let element_base_offset = u32::try_from(element_start).map_err(|_| {
666 Error::new(
667 ErrorCode::CBKD301_RECORD_TOO_SHORT,
668 format!("Array element offset {element_start} exceeds supported range"),
669 )
670 })?;
671 let adjusted_children = adjust_field_offsets(&field.children, element_base_offset);
672 process_fields_recursive(
673 &adjusted_children,
674 data,
675 &mut element_obj,
676 options,
677 scratch_buffers,
678 record_index,
679 encoding_acc,
680 )?;
681 Value::Object(element_obj)
682 }
683 FieldKind::Condition { values } => condition_value(values, "CONDITION_ARRAY"),
684 _ => {
685 let element_data = &data[element_start..element_end];
686 let val = decode_scalar_field_value_standard(
687 field,
688 element_data,
689 options,
690 scratch_buffers,
691 )?;
692 if options.preserve_zoned_encoding {
693 collect_zoned_encoding_info(field, element_data, options, encoding_acc);
694 }
695 val
696 }
697 };
698
699 array_values.push(element_value);
700 }
701
702 json_obj.insert(field.name.clone(), Value::Array(array_values));
703 Ok(())
704}
705
706#[allow(clippy::too_many_arguments)]
708fn process_array_field_with_scratch(
709 field: ©book_core::Field,
710 occurs: ©book_core::Occurs,
711 data: &[u8],
712 json_obj: &mut serde_json::Map<String, Value>,
713 options: &DecodeOptions,
714 all_fields: &[copybook_core::Field],
715 scratch: &mut crate::memory::ScratchBuffers,
716 record_index: u64,
717 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
718) -> Result<()> {
719 use copybook_core::{FieldKind, Occurs};
720 use serde_json::Value;
721
722 let count = match occurs {
723 Occurs::Fixed { count } => *count,
724 Occurs::ODO {
725 min,
726 max,
727 counter_path,
728 } => {
729 let counter_value =
731 find_and_read_counter_field(counter_path, all_fields, data, options, scratch)?;
732
733 let counter_field = find_field_by_path(all_fields, counter_path)?;
734 let validation_context = crate::odo_redefines::OdoValidationContext {
735 field_path: field.path.clone(),
736 counter_path: counter_path.clone(),
737 record_index,
738 byte_offset: u64::from(counter_field.offset),
739 };
740 let validation = crate::odo_redefines::validate_odo_decode(
741 counter_value,
742 *min,
743 *max,
744 &validation_context,
745 options,
746 )?;
747
748 if let Some(warning) = validation.warning {
749 tracing::warn!("{}", warning);
750 increment_warning_counter();
751 }
752
753 validation.actual_count
754 }
755 };
756
757 let element_size = field.len as usize;
758 let array_start = field.offset as usize;
759 let total_array_size = element_size * count as usize;
760 let array_end = array_start + total_array_size;
761
762 if array_end > data.len() {
763 return Err(Error::new(
764 ErrorCode::CBKD301_RECORD_TOO_SHORT,
765 format!(
766 "Array field '{}' with {} elements at offset {} requires {} bytes but record has {}",
767 field.name,
768 count,
769 array_start,
770 total_array_size,
771 data.len() - array_start
772 ),
773 ));
774 }
775
776 let mut array_values = Vec::new();
777
778 for i in 0..count {
779 let element_offset = array_start + (i as usize * element_size);
780 let element_data = &data[element_offset..element_offset + element_size];
781
782 let element_value = match &field.kind {
783 FieldKind::Group => {
784 let mut group_obj = serde_json::Map::new();
786
787 let element_offset_u32 = u32::try_from(element_offset).map_err(|_| {
789 Error::new(
790 ErrorCode::CBKD301_RECORD_TOO_SHORT,
791 format!("Array element offset {element_offset} exceeds supported range"),
792 )
793 })?;
794
795 let mut element_field = field.clone();
796 element_field.offset = element_offset_u32;
797 element_field.occurs = None; process_fields_recursive_with_scratch(
800 &element_field.children,
801 data,
802 &mut group_obj,
803 options,
804 scratch,
805 record_index,
806 encoding_acc,
807 )?;
808 Value::Object(group_obj)
809 }
810 FieldKind::Condition { values } => condition_value(values, "CONDITION_ARRAY"),
811 _ => {
812 let val =
813 decode_scalar_field_value_with_scratch(field, element_data, options, scratch)?;
814 if options.preserve_zoned_encoding {
815 collect_zoned_encoding_info(field, element_data, options, encoding_acc);
816 }
817 val
818 }
819 };
820
821 array_values.push(element_value);
822 }
823
824 json_obj.insert(field.name.clone(), Value::Array(array_values));
825 Ok(())
826}
827
828fn find_and_read_counter_field(
830 counter_path: &str,
831 all_fields: &[copybook_core::Field],
832 data: &[u8],
833 options: &DecodeOptions,
834 scratch: &mut crate::memory::ScratchBuffers,
835) -> Result<u32> {
836 let counter_field = find_field_by_path(all_fields, counter_path)?;
838
839 let field_start = counter_field.offset as usize;
841 let field_end = field_start + counter_field.len as usize;
842
843 if field_end > data.len() {
844 return Err(Error::new(
845 ErrorCode::CBKD301_RECORD_TOO_SHORT,
846 format!("Counter field '{counter_path}' extends beyond record"),
847 ));
848 }
849
850 let field_data = &data[field_start..field_end];
851
852 match &counter_field.kind {
854 copybook_core::FieldKind::ZonedDecimal {
855 digits,
856 scale,
857 signed,
858 sign_separate,
859 } => {
860 let count = if let Some(sign_sep) = sign_separate {
861 let decimal = crate::numeric::decode_zoned_decimal_sign_separate(
862 field_data,
863 *digits,
864 *scale,
865 sign_sep,
866 options.codepage,
867 )?;
868 decimal_counter_to_u32(&decimal, counter_path)?
869 } else {
870 let decimal_str = crate::numeric::decode_zoned_decimal_to_string_with_scratch(
871 field_data,
872 *digits,
873 *scale,
874 *signed,
875 options.codepage,
876 counter_field.blank_when_zero,
877 scratch,
878 )?;
879 decimal_str.parse::<u32>().map_err(|_| {
880 Error::new(
881 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
882 format!("ODO counter '{counter_path}' has invalid value: {decimal_str}"),
883 )
884 })?
885 };
886
887 Ok(count)
888 }
889 copybook_core::FieldKind::BinaryInt { bits, signed } => {
890 let int_value = crate::numeric::decode_binary_int(field_data, *bits, *signed)?;
891 if int_value < 0 {
892 return Err(Error::new(
893 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
894 format!("ODO counter '{counter_path}' has negative value: {int_value}"),
895 ));
896 }
897 Ok(u32::try_from(int_value).map_err(|_| {
898 Error::new(
899 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
900 format!("ODO counter '{counter_path}' exceeds supported range: {int_value}"),
901 )
902 })?)
903 }
904 copybook_core::FieldKind::PackedDecimal {
905 digits,
906 scale,
907 signed,
908 } => {
909 let decimal_str = crate::numeric::decode_packed_decimal_to_string_with_scratch(
910 field_data, *digits, *scale, *signed, scratch,
911 )?;
912 let count = decimal_str.parse::<u32>().map_err(|_| {
913 Error::new(
914 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
915 format!("ODO counter '{counter_path}' has invalid value: {decimal_str}"),
916 )
917 })?;
918 Ok(count)
919 }
920 _ => Err(Error::new(
921 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
922 format!("ODO counter '{counter_path}' has unsupported type"),
923 )),
924 }
925}
926
927fn find_field_by_path<'a>(
929 fields: &'a [copybook_core::Field],
930 path: &str,
931) -> Result<&'a copybook_core::Field> {
932 for field in fields {
933 if field.path == path || field.name == path {
934 return Ok(field);
935 }
936 if let Ok(found) = find_field_by_path(&field.children, path) {
938 return Ok(found);
939 }
940 }
941
942 Err(Error::new(
943 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
944 format!("ODO counter field '{path}' not found"),
945 ))
946}
947
948fn adjust_field_offsets(
954 fields: &[copybook_core::Field],
955 base_offset: u32,
956) -> Vec<copybook_core::Field> {
957 fields
958 .iter()
959 .map(|field| {
960 let mut adjusted_field = field.clone();
961 adjusted_field.offset = base_offset;
962 if !adjusted_field.children.is_empty() {
963 adjusted_field.children =
964 adjust_field_offsets(&adjusted_field.children, base_offset);
965 }
966 adjusted_field
967 })
968 .collect()
969}
970
971#[inline]
973fn is_filler_field(field: ©book_core::Field) -> bool {
974 field.name.eq_ignore_ascii_case("FILLER") || field.name.starts_with("_filler_")
975}
976
977#[inline]
982fn collect_zoned_encoding_info(
983 field: ©book_core::Field,
984 field_data: &[u8],
985 options: &DecodeOptions,
986 encoding_acc: &mut Vec<(String, ZonedEncodingFormat)>,
987) {
988 if let copybook_core::FieldKind::ZonedDecimal { digits, signed, .. } = &field.kind
989 && let Ok((_, Some(info))) = crate::numeric::decode_zoned_decimal_with_encoding(
990 field_data,
991 *digits,
992 0, *signed,
994 options.codepage,
995 field.blank_when_zero,
996 true,
997 )
998 && !info.has_mixed_encoding
999 {
1000 encoding_acc.push((field.name.clone(), info.detected_format));
1001 }
1002}
1003
1004fn numeric_string_to_value(s: String, options: &DecodeOptions) -> Value {
1010 use crate::options::JsonNumberMode;
1011 match options.json_number_mode {
1012 JsonNumberMode::Lossless => Value::String(s),
1013 JsonNumberMode::Native => {
1014 if !s.contains('.') && !s.contains('e') && !s.contains('E') {
1016 if let Ok(n) = s.parse::<i64>() {
1017 return Value::Number(serde_json::Number::from(n));
1018 }
1019 if let Ok(n) = s.parse::<u64>() {
1020 return Value::Number(serde_json::Number::from(n));
1021 }
1022 }
1023 if let Ok(f) = s.parse::<f64>()
1025 && let Some(n) = serde_json::Number::from_f64(f)
1026 {
1027 return Value::Number(n);
1028 }
1029 Value::String(s)
1031 }
1032 }
1033}
1034
1035#[allow(clippy::too_many_lines)]
1037fn decode_scalar_field_value_standard(
1038 field: ©book_core::Field,
1039 field_data: &[u8],
1040 options: &DecodeOptions,
1041 scratch_buffers: &mut Option<crate::memory::ScratchBuffers>,
1042) -> Result<Value> {
1043 use copybook_core::FieldKind;
1044
1045 match &field.kind {
1046 FieldKind::Alphanum { .. } => {
1047 let text = crate::charset::ebcdic_to_utf8(
1048 field_data,
1049 options.codepage,
1050 options.on_decode_unmappable,
1051 )?;
1052 Ok(Value::String(text))
1053 }
1054 FieldKind::ZonedDecimal {
1055 digits,
1056 scale,
1057 signed,
1058 sign_separate,
1059 } => {
1060 if let Some(sign_sep) = sign_separate {
1061 let decimal = crate::numeric::decode_zoned_decimal_sign_separate(
1062 field_data,
1063 *digits,
1064 *scale,
1065 sign_sep,
1066 options.codepage,
1067 )?;
1068 Ok(zoned_decimal_to_json_value(
1069 &decimal,
1070 *digits,
1071 *scale,
1072 field.blank_when_zero,
1073 options,
1074 ))
1075 } else if options.preserve_zoned_encoding {
1076 let (decimal, _encoding_info) = crate::numeric::decode_zoned_decimal_with_encoding(
1078 field_data,
1079 *digits,
1080 *scale,
1081 *signed,
1082 options.codepage,
1083 field.blank_when_zero,
1084 true, )?;
1086
1087 Ok(zoned_decimal_to_json_value(
1090 &decimal,
1091 *digits,
1092 *scale,
1093 field.blank_when_zero,
1094 options,
1095 ))
1096 } else {
1097 let decimal = crate::numeric::decode_zoned_decimal(
1099 field_data,
1100 *digits,
1101 *scale,
1102 *signed,
1103 options.codepage,
1104 field.blank_when_zero,
1105 )?;
1106 Ok(zoned_decimal_to_json_value(
1107 &decimal,
1108 *digits,
1109 *scale,
1110 field.blank_when_zero,
1111 options,
1112 ))
1113 }
1114 }
1115 FieldKind::BinaryInt { bits, signed } => {
1116 let int_value = crate::numeric::decode_binary_int(field_data, *bits, *signed)?;
1117 let scratch = scratch_buffers.get_or_insert_with(crate::memory::ScratchBuffers::new);
1118 let formatted =
1119 crate::numeric::format_binary_int_to_string_with_scratch(int_value, scratch);
1120 Ok(numeric_string_to_value(formatted, options))
1121 }
1122 FieldKind::PackedDecimal {
1123 digits,
1124 scale,
1125 signed,
1126 } => {
1127 let scratch = scratch_buffers.get_or_insert_with(crate::memory::ScratchBuffers::new);
1128 let decimal_str = crate::numeric::decode_packed_decimal_to_string_with_scratch(
1129 field_data, *digits, *scale, *signed, scratch,
1130 )?;
1131 Ok(numeric_string_to_value(decimal_str, options))
1132 }
1133 FieldKind::Group => {
1134 Err(Error::new(
1136 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1137 format!(
1138 "Cannot process group field '{name}' as scalar",
1139 name = field.name
1140 ),
1141 ))
1142 }
1143 FieldKind::Condition { values } => {
1144 Ok(condition_value(values, "CONDITION"))
1147 }
1148 FieldKind::Renames { .. } => {
1149 let Some(resolved) = &field.resolved_renames else {
1151 return Err(Error::new(
1152 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1153 format!(
1154 "RENAMES field '{name}' has no resolved metadata",
1155 name = field.name
1156 ),
1157 ));
1158 };
1159 let alias_start = resolved.offset as usize;
1161 let alias_end = alias_start + resolved.length as usize;
1162
1163 if alias_end > field_data.len() {
1164 return Err(Error::new(
1165 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1166 format!(
1167 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
1168 name = field.name,
1169 offset = resolved.offset,
1170 length = resolved.length,
1171 data_len = field_data.len()
1172 ),
1173 ));
1174 }
1175
1176 if resolved.members.len() == 1 {
1179 let alias_data = &field_data[alias_start..alias_end];
1181 let text = crate::charset::ebcdic_to_utf8(
1183 alias_data,
1184 options.codepage,
1185 options.on_decode_unmappable,
1186 )?;
1187 return Ok(Value::String(text));
1188 }
1189 let alias_data = &field_data[alias_start..alias_end];
1191 let text = crate::charset::ebcdic_to_utf8(
1192 alias_data,
1193 options.codepage,
1194 options.on_decode_unmappable,
1195 )?;
1196 Ok(Value::String(text))
1197 }
1198 FieldKind::EditedNumeric {
1199 pic_string, scale, ..
1200 } => {
1201 let raw_str = crate::charset::ebcdic_to_utf8(
1203 field_data,
1204 options.codepage,
1205 options.on_decode_unmappable,
1206 )?;
1207
1208 let pattern = crate::edited_pic::tokenize_edited_pic(pic_string)?;
1210
1211 let numeric_value = crate::edited_pic::decode_edited_numeric(
1213 &raw_str,
1214 &pattern,
1215 *scale,
1216 field.blank_when_zero,
1217 )?;
1218
1219 Ok(numeric_string_to_value(
1221 numeric_value.to_decimal_string(),
1222 options,
1223 ))
1224 }
1225 FieldKind::FloatSingle => {
1226 let value =
1227 crate::numeric::decode_float_single_with_format(field_data, options.float_format)?;
1228 if value.is_nan() || value.is_infinite() {
1229 Ok(Value::Null)
1230 } else {
1231 Ok(Value::Number(
1232 serde_json::Number::from_f64(f64::from(value))
1233 .unwrap_or_else(|| serde_json::Number::from(0)),
1234 ))
1235 }
1236 }
1237 FieldKind::FloatDouble => {
1238 let value =
1239 crate::numeric::decode_float_double_with_format(field_data, options.float_format)?;
1240 if value.is_nan() || value.is_infinite() {
1241 Ok(Value::Null)
1242 } else {
1243 Ok(Value::Number(
1244 serde_json::Number::from_f64(value)
1245 .unwrap_or_else(|| serde_json::Number::from(0)),
1246 ))
1247 }
1248 }
1249 }
1250}
1251
1252fn is_scalar_target_group_redefine(
1257 field: ©book_core::Field,
1258 sibling_fields: &[copybook_core::Field],
1259) -> bool {
1260 let Some(target_path) = field.redefines_of.as_deref() else {
1261 return false;
1262 };
1263
1264 matches!(field.kind, copybook_core::FieldKind::Group)
1265 && find_field_by_path(sibling_fields, target_path)
1266 .is_ok_and(|target| !matches!(target.kind, copybook_core::FieldKind::Group))
1267}
1268
1269#[allow(clippy::too_many_lines)]
1271fn decode_scalar_field_value_with_scratch(
1272 field: ©book_core::Field,
1273 field_data: &[u8],
1274 options: &DecodeOptions,
1275 scratch: &mut crate::memory::ScratchBuffers,
1276) -> Result<Value> {
1277 use copybook_core::FieldKind;
1278
1279 match &field.kind {
1280 FieldKind::Alphanum { .. } => {
1281 let text = crate::charset::ebcdic_to_utf8(
1282 field_data,
1283 options.codepage,
1284 options.on_decode_unmappable,
1285 )?;
1286 Ok(Value::String(text))
1287 }
1288 FieldKind::ZonedDecimal {
1289 digits,
1290 scale,
1291 signed,
1292 sign_separate,
1293 } => {
1294 if let Some(sign_sep) = sign_separate {
1295 let decimal = crate::numeric::decode_zoned_decimal_sign_separate(
1296 field_data,
1297 *digits,
1298 *scale,
1299 sign_sep,
1300 options.codepage,
1301 )?;
1302 Ok(zoned_decimal_to_json_value(
1303 &decimal,
1304 *digits,
1305 *scale,
1306 field.blank_when_zero,
1307 options,
1308 ))
1309 } else {
1310 let decimal_str = crate::numeric::decode_zoned_decimal_to_string_with_scratch(
1311 field_data,
1312 *digits,
1313 *scale,
1314 *signed,
1315 options.codepage,
1316 field.blank_when_zero,
1317 scratch,
1318 )?;
1319 Ok(numeric_string_to_value(decimal_str, options))
1320 }
1321 }
1322 FieldKind::BinaryInt { bits, signed } => {
1323 let int_value = crate::numeric::decode_binary_int(field_data, *bits, *signed)?;
1324 let formatted =
1325 crate::numeric::format_binary_int_to_string_with_scratch(int_value, scratch);
1326 Ok(numeric_string_to_value(formatted, options))
1327 }
1328 FieldKind::PackedDecimal {
1329 digits,
1330 scale,
1331 signed,
1332 } => {
1333 let decimal_str = crate::numeric::decode_packed_decimal_to_string_with_scratch(
1334 field_data, *digits, *scale, *signed, scratch,
1335 )?;
1336 Ok(numeric_string_to_value(decimal_str, options))
1337 }
1338 FieldKind::Group => Err(Error::new(
1339 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1340 format!(
1341 "Cannot process group field '{name}' as scalar",
1342 name = field.name
1343 ),
1344 )),
1345 FieldKind::Condition { values } => Ok(condition_value(values, "CONDITION")),
1346 FieldKind::Renames { .. } => {
1347 let Some(resolved) = &field.resolved_renames else {
1349 return Err(Error::new(
1350 ErrorCode::CBKD101_INVALID_FIELD_TYPE,
1351 format!(
1352 "RENAMES field '{name}' has no resolved metadata",
1353 name = field.name
1354 ),
1355 ));
1356 };
1357 let alias_start = resolved.offset as usize;
1359 let alias_end = alias_start + resolved.length as usize;
1360
1361 if alias_end > field_data.len() {
1362 return Err(Error::new(
1363 ErrorCode::CBKD301_RECORD_TOO_SHORT,
1364 format!(
1365 "RENAMES field '{name}' at offset {offset} with length {length} exceeds data length {data_len}",
1366 name = field.name,
1367 offset = resolved.offset,
1368 length = resolved.length,
1369 data_len = field_data.len()
1370 ),
1371 ));
1372 }
1373
1374 let alias_data = &field_data[alias_start..alias_end];
1376 let text = crate::charset::ebcdic_to_utf8(
1377 alias_data,
1378 options.codepage,
1379 options.on_decode_unmappable,
1380 )?;
1381 Ok(Value::String(text))
1382 }
1383 FieldKind::EditedNumeric {
1384 pic_string, scale, ..
1385 } => {
1386 let raw_str = crate::charset::ebcdic_to_utf8(
1388 field_data,
1389 options.codepage,
1390 options.on_decode_unmappable,
1391 )?;
1392
1393 let pattern = crate::edited_pic::tokenize_edited_pic(pic_string)?;
1395
1396 let numeric_value = crate::edited_pic::decode_edited_numeric(
1398 &raw_str,
1399 &pattern,
1400 *scale,
1401 field.blank_when_zero,
1402 )?;
1403
1404 Ok(numeric_string_to_value(
1406 numeric_value.to_decimal_string(),
1407 options,
1408 ))
1409 }
1410 FieldKind::FloatSingle => {
1411 let value =
1412 crate::numeric::decode_float_single_with_format(field_data, options.float_format)?;
1413 if value.is_nan() || value.is_infinite() {
1414 Ok(Value::Null)
1415 } else {
1416 Ok(Value::Number(
1417 serde_json::Number::from_f64(f64::from(value))
1418 .unwrap_or_else(|| serde_json::Number::from(0)),
1419 ))
1420 }
1421 }
1422 FieldKind::FloatDouble => {
1423 let value =
1424 crate::numeric::decode_float_double_with_format(field_data, options.float_format)?;
1425 if value.is_nan() || value.is_infinite() {
1426 Ok(Value::Null)
1427 } else {
1428 Ok(Value::Number(
1429 serde_json::Number::from_f64(value)
1430 .unwrap_or_else(|| serde_json::Number::from(0)),
1431 ))
1432 }
1433 }
1434 }
1435}
1436
1437#[inline]
1441fn condition_value(values: &[String], prefix: &str) -> Value {
1442 if values.is_empty() {
1443 Value::String(prefix.to_owned())
1444 } else {
1445 Value::String(format!("{prefix}({})", values.join("|")))
1446 }
1447}
1448
1449#[inline]
1477#[must_use = "Handle the Result or propagate the error"]
1478pub fn encode_record(schema: &Schema, json: &Value, options: &EncodeOptions) -> Result<Vec<u8>> {
1479 let root_obj = json.as_object().ok_or_else(|| {
1480 Error::new(
1481 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1482 "Expected JSON object for record envelope",
1483 )
1484 })?;
1485 let fields_value = if let Some(fields_val) = root_obj.get("fields") {
1486 fields_val.as_object().ok_or_else(|| {
1487 Error::new(
1488 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1489 "`fields` must be a JSON object",
1490 )
1491 })?;
1492 fields_val
1493 } else {
1494 json
1495 };
1496
1497 if options.use_raw
1499 && let Some(raw_b64_value) = root_obj
1500 .get("raw_b64")
1501 .or_else(|| root_obj.get("__raw_b64"))
1502 && let Some(raw_str) = raw_b64_value.as_str()
1503 {
1504 let raw_data = base64::engine::general_purpose::STANDARD
1506 .decode(raw_str)
1507 .map_err(|e| {
1508 Error::new(
1509 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1510 format!("Invalid base64 in raw_b64: {e}"),
1511 )
1512 })?;
1513
1514 match options.format {
1515 RecordFormat::RDW => {
1516 if raw_data.len() >= 4 {
1518 let mut rdw_record = raw_data.clone();
1519
1520 let payload = &rdw_record[4..];
1522
1523 let mut should_recompute = false;
1525
1526 let field_payload = encode_fields_to_bytes(schema, fields_value, options)?;
1528 if field_payload != payload {
1529 should_recompute = true;
1530 }
1531
1532 if should_recompute {
1533 let capped_len = field_payload.len().min(u16::MAX as usize);
1535 let new_length = u16::try_from(capped_len).unwrap_or(u16::MAX);
1536 let length_bytes = new_length.to_be_bytes();
1537 rdw_record[0] = length_bytes[0];
1538 rdw_record[1] = length_bytes[1];
1539 rdw_record.splice(4.., field_payload);
1543 }
1544
1545 return Ok(rdw_record);
1546 }
1547 }
1548 RecordFormat::Fixed => {
1549 return Ok(raw_data);
1550 }
1551 }
1552 }
1553
1554 validate_lib_api_redefines_encoding(schema, fields_value, options)?;
1556 validate_lib_api_odo_encoding(schema, fields_value, options)?;
1557
1558 match options.format {
1559 RecordFormat::Fixed => {
1560 let payload = encode_fields_to_bytes(schema, fields_value, options)?;
1561 Ok(payload)
1562 }
1563 RecordFormat::RDW => {
1564 let payload = encode_fields_to_bytes(schema, fields_value, options)?;
1565
1566 let rdw_record = crate::record::RDWRecord::try_new(payload)?;
1568 let mut result = Vec::new();
1569 result.extend_from_slice(&rdw_record.header);
1570 result.extend_from_slice(&rdw_record.payload);
1571 Ok(result)
1572 }
1573 }
1574}
1575
1576fn validate_lib_api_redefines_encoding(
1578 schema: &Schema,
1579 json_value: &Value,
1580 options: &EncodeOptions,
1581) -> Result<()> {
1582 let redefines_context = crate::odo_redefines::build_redefines_context(schema, json_value);
1583
1584 for (cluster_path, non_null_views) in &redefines_context.cluster_views {
1585 let field_path = non_null_views
1586 .first()
1587 .cloned()
1588 .unwrap_or_else(|| cluster_path.clone());
1589
1590 let byte_offset = non_null_views
1591 .iter()
1592 .find_map(|view| schema.find_field(view).map(|field| u64::from(field.offset)))
1593 .or_else(|| {
1594 schema
1595 .find_field(cluster_path)
1596 .map(|field| u64::from(field.offset))
1597 })
1598 .unwrap_or(0);
1599
1600 crate::odo_redefines::validate_redefines_encoding(
1601 &redefines_context,
1602 cluster_path,
1603 &field_path,
1604 json_value,
1605 options.use_raw,
1606 0,
1607 byte_offset,
1608 )?;
1609 }
1610
1611 Ok(())
1612}
1613
1614fn validate_lib_api_odo_encoding(
1616 schema: &Schema,
1617 json_value: &Value,
1618 options: &EncodeOptions,
1619) -> Result<()> {
1620 let Some(tail_odo) = &schema.tail_odo else {
1621 return Ok(());
1622 };
1623
1624 let fields_value = if let Some(fields_value) = json_value.get("fields") {
1625 fields_value
1626 } else {
1627 json_value
1628 };
1629
1630 let has_wrapper = json_value.get("fields").is_some();
1631 if !fields_value.is_object() {
1632 if has_wrapper {
1633 return Err(Error::new(
1634 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
1635 "`fields` must be a JSON object",
1636 ));
1637 }
1638 return Ok(());
1639 }
1640
1641 let array_field =
1642 crate::odo_redefines::find_field_by_path_or_unique_name(schema, &tail_odo.array_path)
1643 .ok_or_else(|| {
1644 Error::new(
1645 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
1646 format!(
1647 "ODO array field '{}' not found in schema",
1648 tail_odo.array_path
1649 ),
1650 )
1651 .with_context(
1652 crate::odo_redefines::create_comprehensive_error_context(
1653 0,
1654 &tail_odo.array_path,
1655 0,
1656 None,
1657 ),
1658 )
1659 })?;
1660
1661 let counter_field =
1662 crate::odo_redefines::find_field_by_path_or_unique_name(schema, &tail_odo.counter_path)
1663 .ok_or_else(|| {
1664 crate::odo_redefines::handle_missing_counter_field(
1665 &tail_odo.counter_path,
1666 &tail_odo.array_path,
1667 schema,
1668 0,
1669 0,
1670 )
1671 })?;
1672
1673 if let Some(array) = json_lookup_array(fields_value, &array_field.path)
1674 .or_else(|| json_lookup_array(fields_value, &tail_odo.array_path))
1675 {
1676 let Some(counter_json_value) = json_lookup_value(fields_value, &counter_field.path)
1677 .or_else(|| json_lookup_value(fields_value, &tail_odo.counter_path))
1678 else {
1679 return Err(crate::odo_redefines::handle_missing_counter_field(
1680 &counter_field.path,
1681 &array_field.path,
1682 schema,
1683 0,
1684 u64::from(counter_field.offset),
1685 ));
1686 };
1687
1688 if let Some(counter_count) = json_counter_value_as_usize(counter_json_value)
1694 && counter_count != array.len()
1695 {
1696 return Err(Error::new(
1697 ErrorCode::CBKE521_ARRAY_LEN_OOB,
1698 format!(
1699 "ODO counter '{}' value ({counter_count}) does not match array '{}' length ({})",
1700 counter_field.path,
1701 array_field.path,
1702 array.len()
1703 ),
1704 )
1705 .with_context(crate::odo_redefines::create_comprehensive_error_context(
1706 0,
1707 &array_field.path,
1708 u64::from(array_field.offset),
1709 Some(format!(
1710 "counter_field={}, counter_value={counter_count}, array_length={}",
1711 counter_field.path,
1712 array.len()
1713 )),
1714 )));
1715 }
1716
1717 let context = crate::odo_redefines::OdoValidationContext {
1718 field_path: array_field.path.clone(),
1719 counter_path: counter_field.path.clone(),
1720 record_index: 0,
1721 byte_offset: u64::from(array_field.offset),
1722 };
1723
1724 crate::odo_redefines::validate_odo_encode(
1725 array.len(),
1726 tail_odo.min_count,
1727 tail_odo.max_count,
1728 &context,
1729 options,
1730 )?;
1731 }
1732
1733 Ok(())
1734}
1735
1736fn json_counter_value_as_usize(value: &Value) -> Option<usize> {
1738 match value {
1739 Value::Number(n) => n.as_u64().and_then(|v| usize::try_from(v).ok()),
1740 Value::String(s) => s.trim().parse::<usize>().ok(),
1741 _ => None,
1742 }
1743}
1744
1745fn json_lookup_value<'a>(value: &'a Value, field_path: &str) -> Option<&'a Value> {
1746 json_lookup_exact_value(value, field_path).or_else(|| {
1747 let (_, path_without_root) = field_path.split_once('.')?;
1748 json_lookup_exact_value(value, path_without_root)
1749 })
1750}
1751
1752fn json_lookup_exact_value<'a>(value: &'a Value, field_path: &str) -> Option<&'a Value> {
1753 let mut current = value;
1754 for segment in field_path.split('.') {
1755 current = current.as_object()?.get(segment)?;
1756 }
1757 Some(current)
1758}
1759
1760fn json_lookup_array<'a>(value: &'a Value, field_path: &str) -> Option<&'a Vec<Value>> {
1761 let leaf = field_path.split('.').next_back().unwrap_or("");
1762 match json_lookup_value(value, field_path) {
1763 Some(Value::Array(array)) => Some(array),
1764 _ => {
1765 if let Value::Object(obj) = value {
1766 obj.get(leaf).and_then(|candidate| candidate.as_array())
1767 } else {
1768 None
1769 }
1770 }
1771 }
1772}
1773
1774fn encode_fields_to_bytes(
1776 schema: &Schema,
1777 json: &Value,
1778 options: &EncodeOptions,
1779) -> Result<Vec<u8>> {
1780 let record_length = schema.lrecl_fixed.unwrap_or_else(|| {
1781 schema.fields.iter().map(|f| f.len).sum::<u32>()
1783 }) as usize;
1784
1785 let mut buffer = vec![0u8; record_length];
1786
1787 if let Some(obj) = json.as_object() {
1788 let encoding_metadata = obj
1789 .get("_encoding_metadata")
1790 .and_then(|value| value.as_object());
1791 encode_fields_recursive(
1792 &schema.fields,
1793 obj,
1794 encoding_metadata,
1795 "",
1796 &mut buffer,
1797 0,
1798 options,
1799 )?;
1800 }
1801
1802 Ok(buffer)
1803}
1804
1805fn encode_fields_recursive(
1807 fields: &[copybook_core::Field],
1808 json_obj: &serde_json::Map<String, Value>,
1809 encoding_metadata: Option<&serde_json::Map<String, Value>>,
1810 path_prefix: &str,
1811 buffer: &mut [u8],
1812 offset: usize,
1813 options: &EncodeOptions,
1814) -> Result<usize> {
1815 let mut current_offset = offset;
1816
1817 for field in fields {
1818 let field_path = if path_prefix.is_empty() {
1819 field.name.clone()
1820 } else {
1821 format!("{path_prefix}.{}", field.name)
1822 };
1823
1824 current_offset = encode_single_field(
1825 field,
1826 &field_path,
1827 json_obj,
1828 encoding_metadata,
1829 buffer,
1830 current_offset,
1831 options,
1832 )?;
1833 }
1834
1835 Ok(current_offset)
1836}
1837
1838#[inline]
1843#[allow(clippy::too_many_lines)]
1844fn encode_single_field(
1845 field: ©book_core::Field,
1846 field_path: &str,
1847 json_obj: &serde_json::Map<String, Value>,
1848 encoding_metadata: Option<&serde_json::Map<String, Value>>,
1849 buffer: &mut [u8],
1850 current_offset: usize,
1851 options: &EncodeOptions,
1852) -> Result<usize> {
1853 use copybook_core::FieldKind;
1854
1855 if let Some(occurs) = &field.occurs {
1856 return encode_occurs_field(
1857 field,
1858 occurs,
1859 field_path,
1860 json_obj,
1861 encoding_metadata,
1862 buffer,
1863 current_offset,
1864 options,
1865 );
1866 }
1867
1868 match &field.kind {
1869 FieldKind::Group => encode_group_field(
1870 field,
1871 field_path,
1872 json_obj,
1873 encoding_metadata,
1874 buffer,
1875 current_offset,
1876 options,
1877 ),
1878 FieldKind::Alphanum { .. } => {
1879 encode_alphanum_field(field, json_obj, buffer, current_offset, options)
1880 }
1881 FieldKind::ZonedDecimal {
1882 digits,
1883 scale,
1884 signed,
1885 sign_separate,
1886 } => {
1887 if let Some(sign_sep) = sign_separate {
1888 let field_len = field.len as usize;
1889 if let Some(text) = json_obj.get(&field.name).and_then(|v| v.as_str()) {
1890 crate::numeric::encode_zoned_decimal_sign_separate(
1891 text,
1892 *digits,
1893 *scale,
1894 sign_sep,
1895 options.codepage,
1896 &mut buffer[current_offset..current_offset + field_len],
1897 )?;
1898 }
1899 Ok(current_offset + field_len)
1900 } else {
1901 encode_zoned_decimal_field(
1902 field,
1903 field_path,
1904 json_obj,
1905 encoding_metadata,
1906 buffer,
1907 current_offset,
1908 options,
1909 DecimalSpec {
1910 digits: *digits,
1911 scale: *scale,
1912 signed: *signed,
1913 },
1914 )
1915 }
1916 }
1917 FieldKind::PackedDecimal {
1918 digits,
1919 scale,
1920 signed,
1921 } => encode_packed_decimal_field(
1922 field,
1923 field_path,
1924 json_obj,
1925 buffer,
1926 current_offset,
1927 options,
1928 DecimalSpec {
1929 digits: *digits,
1930 scale: *scale,
1931 signed: *signed,
1932 },
1933 ),
1934 FieldKind::BinaryInt { bits, signed } => encode_binary_int_field(
1935 field,
1936 field_path,
1937 json_obj,
1938 buffer,
1939 current_offset,
1940 options,
1941 BinarySpec {
1942 bits: *bits,
1943 signed: *signed,
1944 },
1945 ),
1946 FieldKind::Condition { .. } => Ok(current_offset),
1947 FieldKind::Renames { .. } => {
1948 Ok(current_offset)
1952 }
1953 FieldKind::EditedNumeric {
1954 pic_string, scale, ..
1955 } => {
1956 if let Some(text) = json_obj
1958 .get(&field.name)
1959 .and_then(|v| coerce_to_str(v, options.coerce_numbers))
1960 {
1961 let pattern = crate::edited_pic::tokenize_edited_pic(pic_string)?;
1963
1964 let encoded = crate::edited_pic::encode_edited_numeric(
1966 &text,
1967 &pattern,
1968 *scale,
1969 field.blank_when_zero,
1970 )?;
1971
1972 let bytes = crate::charset::utf8_to_ebcdic(&encoded, options.codepage)?;
1974 let field_len = field.len as usize;
1975 let copy_len = bytes.len().min(field_len);
1976
1977 if current_offset + field_len <= buffer.len() {
1978 buffer[current_offset..current_offset + copy_len]
1979 .copy_from_slice(&bytes[..copy_len]);
1980 let space = crate::charset::space_byte(options.codepage);
1982 buffer[current_offset + copy_len..current_offset + field_len].fill(space);
1983 }
1984 }
1985 Ok(current_offset + field.len as usize)
1986 }
1987 FieldKind::FloatSingle => {
1988 let field_len = field.len as usize;
1989 if let Some(val) = json_obj.get(&field.name) {
1990 let f = match val {
1991 Value::Number(n) => {
1992 let f64_val = n.as_f64().unwrap_or(0.0);
1993 if f64_val.is_finite()
1995 && (f64_val > f64::from(f32::MAX) || f64_val < f64::from(f32::MIN))
1996 {
1997 return Err(Error::new(
1998 ErrorCode::CBKE531_FLOAT_ENCODE_OVERFLOW,
1999 format!("Value overflow for COMP-1 field '{}'", field.name),
2000 ));
2001 }
2002 #[allow(clippy::cast_possible_truncation)]
2004 {
2005 f64_val as f32
2006 }
2007 }
2008 Value::String(s) => s.parse::<f32>().map_err(|e| {
2009 Error::new(
2010 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2011 format!(
2012 "Cannot parse '{}' as f32 for field '{}': {}",
2013 s, field.name, e
2014 ),
2015 )
2016 })?,
2017 Value::Null => f32::NAN,
2018 _ => {
2019 return Err(Error::new(
2020 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2021 format!("Expected number for COMP-1 field '{}'", field.name),
2022 ));
2023 }
2024 };
2025 if current_offset + field_len <= buffer.len() {
2026 crate::numeric::encode_float_single_with_format(
2027 f,
2028 &mut buffer[current_offset..current_offset + field_len],
2029 options.float_format,
2030 )?;
2031 }
2032 }
2033 Ok(current_offset + field_len)
2034 }
2035 FieldKind::FloatDouble => {
2036 let field_len = field.len as usize;
2037 if let Some(val) = json_obj.get(&field.name) {
2038 let f = match val {
2039 Value::Number(n) => n.as_f64().unwrap_or(0.0),
2040 Value::String(s) => s.parse::<f64>().map_err(|e| {
2041 Error::new(
2042 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2043 format!(
2044 "Cannot parse '{}' as f64 for field '{}': {}",
2045 s, field.name, e
2046 ),
2047 )
2048 })?,
2049 Value::Null => f64::NAN,
2050 _ => {
2051 return Err(Error::new(
2052 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2053 format!("Expected number for COMP-2 field '{}'", field.name),
2054 ));
2055 }
2056 };
2057 if current_offset + field_len <= buffer.len() {
2058 crate::numeric::encode_float_double_with_format(
2059 f,
2060 &mut buffer[current_offset..current_offset + field_len],
2061 options.float_format,
2062 )?;
2063 }
2064 }
2065 Ok(current_offset + field_len)
2066 }
2067 }
2068}
2069
2070#[allow(clippy::too_many_arguments)]
2071fn encode_occurs_field(
2072 field: ©book_core::Field,
2073 occurs: ©book_core::Occurs,
2074 field_path: &str,
2075 json_obj: &serde_json::Map<String, Value>,
2076 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2077 buffer: &mut [u8],
2078 current_offset: usize,
2079 options: &EncodeOptions,
2080) -> Result<usize> {
2081 let max_count = occurs_max_count(occurs);
2082 let element_len = field.len as usize;
2083 let allocation_len = element_len
2084 .checked_mul(max_count as usize)
2085 .ok_or_else(|| Error::new(ErrorCode::CBKS141_RECORD_TOO_LARGE, "OCCURS size overflow"))?;
2086
2087 let Some(array) = json_obj.get(&field.name).and_then(Value::as_array) else {
2088 return Ok(current_offset + allocation_len);
2089 };
2090
2091 validate_occurs_array_len(array.len(), occurs, field)?;
2092
2093 for (index, element) in array.iter().enumerate() {
2094 let element_offset = current_offset + index * element_len;
2095 encode_occurs_element(
2096 field,
2097 field_path,
2098 element,
2099 encoding_metadata,
2100 buffer,
2101 element_offset,
2102 options,
2103 )?;
2104 }
2105
2106 Ok(current_offset + allocation_len)
2107}
2108
2109fn occurs_max_count(occurs: ©book_core::Occurs) -> u32 {
2110 match occurs {
2111 copybook_core::Occurs::Fixed { count } => *count,
2112 copybook_core::Occurs::ODO { max, .. } => *max,
2113 }
2114}
2115
2116fn validate_occurs_array_len(
2117 actual_len: usize,
2118 occurs: ©book_core::Occurs,
2119 field: ©book_core::Field,
2120) -> Result<()> {
2121 match occurs {
2122 copybook_core::Occurs::Fixed { count } if actual_len != *count as usize => Err(Error::new(
2123 ErrorCode::CBKE521_ARRAY_LEN_OOB,
2124 format!(
2125 "Array length {} doesn't match fixed OCCURS count {} for field '{}'",
2126 actual_len, count, field.path
2127 ),
2128 )
2129 .with_field(field.path.clone())),
2130 copybook_core::Occurs::ODO { max, .. } if actual_len > *max as usize => Err(Error::new(
2131 ErrorCode::CBKE521_ARRAY_LEN_OOB,
2132 format!(
2133 "Array length {} exceeds ODO max {} for field '{}'",
2134 actual_len, max, field.path
2135 ),
2136 )
2137 .with_field(field.path.clone())),
2138 copybook_core::Occurs::Fixed { .. } | copybook_core::Occurs::ODO { .. } => Ok(()),
2139 }
2140}
2141
2142fn encode_occurs_element(
2143 field: ©book_core::Field,
2144 field_path: &str,
2145 element: &Value,
2146 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2147 buffer: &mut [u8],
2148 element_offset: usize,
2149 options: &EncodeOptions,
2150) -> Result<()> {
2151 use copybook_core::FieldKind;
2152
2153 let mut element_field = field.clone();
2154 element_field.occurs = None;
2155
2156 if let FieldKind::Group = &field.kind {
2157 let element_obj = element.as_object().ok_or_else(|| {
2158 Error::new(
2159 ErrorCode::CBKE501_JSON_TYPE_MISMATCH,
2160 format!("Expected object element for OCCURS group '{}'", field.path),
2161 )
2162 .with_field(field.path.clone())
2163 })?;
2164 encode_fields_recursive(
2165 &element_field.children,
2166 element_obj,
2167 encoding_metadata,
2168 field_path,
2169 buffer,
2170 element_offset,
2171 options,
2172 )?;
2173 } else {
2174 let mut element_obj = serde_json::Map::new();
2175 element_obj.insert(field.name.clone(), element.clone());
2176 encode_single_field(
2177 &element_field,
2178 field_path,
2179 &element_obj,
2180 encoding_metadata,
2181 buffer,
2182 element_offset,
2183 options,
2184 )?;
2185 }
2186
2187 Ok(())
2188}
2189
2190#[inline]
2192fn encode_group_field(
2193 field: ©book_core::Field,
2194 field_path: &str,
2195 json_obj: &serde_json::Map<String, Value>,
2196 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2197 buffer: &mut [u8],
2198 current_offset: usize,
2199 options: &EncodeOptions,
2200) -> Result<usize> {
2201 if let Some(sub_obj) = json_obj.get(&field.name).and_then(|v| v.as_object()) {
2202 encode_fields_recursive(
2203 &field.children,
2204 sub_obj,
2205 encoding_metadata,
2206 field_path,
2207 buffer,
2208 current_offset,
2209 options,
2210 )
2211 } else {
2212 encode_fields_recursive(
2213 &field.children,
2214 json_obj,
2215 encoding_metadata,
2216 field_path,
2217 buffer,
2218 current_offset,
2219 options,
2220 )
2221 }
2222}
2223
2224#[inline]
2226fn encode_alphanum_field(
2227 field: ©book_core::Field,
2228 json_obj: &serde_json::Map<String, Value>,
2229 buffer: &mut [u8],
2230 current_offset: usize,
2231 options: &EncodeOptions,
2232) -> Result<usize> {
2233 let field_len = field.len as usize;
2234
2235 if let Some(text) = json_obj.get(&field.name).and_then(|value| value.as_str()) {
2236 let bytes = crate::charset::utf8_to_ebcdic(text, options.codepage)?;
2238 if bytes.len() > field_len {
2239 return Err(Error::new(
2240 ErrorCode::CBKE515_STRING_LENGTH_VIOLATION,
2241 format!(
2242 "Encoded byte length {} exceeds field capacity {} for alphanumeric field {}",
2243 bytes.len(),
2244 field_len,
2245 field.path
2246 ),
2247 )
2248 .with_field(field.path.clone()));
2249 }
2250
2251 let copy_len = bytes.len();
2252
2253 if current_offset + field_len <= buffer.len() {
2254 buffer[current_offset..current_offset + copy_len].copy_from_slice(&bytes);
2255 let space = crate::charset::space_byte(options.codepage);
2257 buffer[current_offset + copy_len..current_offset + field_len].fill(space);
2258 }
2259 }
2260
2261 Ok(current_offset + field_len)
2262}
2263
2264#[derive(Copy, Clone)]
2265struct DecimalSpec {
2266 digits: u16,
2267 scale: i16,
2268 signed: bool,
2269}
2270
2271fn resolve_preserved_zoned_format(
2272 metadata: &serde_json::Map<String, Value>,
2273 field_path: &str,
2274 field_name: &str,
2275) -> Option<ZonedEncodingFormat> {
2276 let candidates = [field_path, field_name];
2277 for key in candidates {
2278 if let Some(format) = metadata
2279 .get(key)
2280 .and_then(parse_zoned_encoding_metadata_value)
2281 {
2282 return Some(format);
2283 }
2284 }
2285 None
2286}
2287
2288fn parse_zoned_encoding_metadata_value(value: &Value) -> Option<ZonedEncodingFormat> {
2289 match value {
2290 Value::String(s) => parse_zoned_encoding_format_str(s),
2291 Value::Object(map) => map
2292 .get("zoned_encoding")
2293 .and_then(Value::as_str)
2294 .and_then(parse_zoned_encoding_format_str),
2295 _ => None,
2296 }
2297}
2298
2299fn parse_zoned_encoding_format_str(value: &str) -> Option<ZonedEncodingFormat> {
2300 match value.trim().to_ascii_lowercase().as_str() {
2301 "ascii" => Some(ZonedEncodingFormat::Ascii),
2302 "ebcdic" => Some(ZonedEncodingFormat::Ebcdic),
2303 "auto" => Some(ZonedEncodingFormat::Auto),
2304 _ => None,
2305 }
2306}
2307
2308fn coerce_to_str(value: &Value, coerce: bool) -> Option<String> {
2314 match value {
2315 Value::String(s) => Some(s.clone()),
2316 Value::Number(n) if coerce => Some(n.to_string()),
2317 _ => None,
2318 }
2319}
2320
2321#[inline]
2322#[allow(clippy::too_many_arguments)]
2323fn encode_zoned_decimal_field(
2324 field: ©book_core::Field,
2325 field_path: &str,
2326 json_obj: &serde_json::Map<String, Value>,
2327 encoding_metadata: Option<&serde_json::Map<String, Value>>,
2328 buffer: &mut [u8],
2329 current_offset: usize,
2330 options: &EncodeOptions,
2331 spec: DecimalSpec,
2332) -> Result<usize> {
2333 let field_len = field.len as usize;
2334
2335 if let Some(text) = json_obj
2336 .get(&field.name)
2337 .and_then(|v| coerce_to_str(v, options.coerce_numbers))
2338 {
2339 if field.blank_when_zero && options.bwz_encode {
2342 let encoded = crate::numeric::encode_zoned_decimal_with_bwz(
2343 &text,
2344 spec.digits,
2345 spec.scale,
2346 spec.signed,
2347 options.codepage,
2348 options.bwz_encode,
2349 )?;
2350 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
2351 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
2352 }
2353 return Ok(current_offset + field_len);
2354 }
2355
2356 let preserved_format = encoding_metadata
2357 .and_then(|meta| resolve_preserved_zoned_format(meta, field_path, &field.name));
2358 let resolved_format = options
2359 .zoned_encoding_override
2360 .or(preserved_format)
2361 .unwrap_or(options.preferred_zoned_encoding);
2362 let (effective_format, zero_policy) = match resolved_format {
2364 ZonedEncodingFormat::Ascii => (ZonedEncodingFormat::Ascii, ZeroSignPolicy::Positive),
2365 ZonedEncodingFormat::Ebcdic => (ZonedEncodingFormat::Ebcdic, ZeroSignPolicy::Preferred),
2366 ZonedEncodingFormat::Auto => {
2367 if options.codepage.is_ascii() {
2368 (ZonedEncodingFormat::Ascii, ZeroSignPolicy::Positive)
2369 } else {
2370 (ZonedEncodingFormat::Ebcdic, ZeroSignPolicy::Preferred)
2371 }
2372 }
2373 };
2374
2375 let encoded = crate::numeric::encode_zoned_decimal_with_format_and_policy(
2376 &text,
2377 spec.digits,
2378 spec.scale,
2379 spec.signed,
2380 options.codepage,
2381 Some(effective_format),
2382 zero_policy,
2383 )?;
2384
2385 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
2386 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
2387 }
2388 }
2389
2390 Ok(current_offset + field_len)
2391}
2392
2393#[inline]
2394fn encode_packed_decimal_field(
2395 field: ©book_core::Field,
2396 _field_path: &str,
2397 json_obj: &serde_json::Map<String, Value>,
2398 buffer: &mut [u8],
2399 current_offset: usize,
2400 options: &EncodeOptions,
2401 spec: DecimalSpec,
2402) -> Result<usize> {
2403 let field_len = field.len as usize;
2404
2405 if let Some(text) = json_obj
2406 .get(&field.name)
2407 .and_then(|v| coerce_to_str(v, options.coerce_numbers))
2408 {
2409 let encoded =
2410 crate::numeric::encode_packed_decimal(&text, spec.digits, spec.scale, spec.signed)?;
2411 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
2412 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
2413 }
2414 }
2415
2416 Ok(current_offset + field_len)
2417}
2418
2419#[derive(Copy, Clone)]
2420struct BinarySpec {
2421 bits: u16,
2422 signed: bool,
2423}
2424
2425#[inline]
2426fn encode_binary_int_field(
2427 field: ©book_core::Field,
2428 _field_path: &str,
2429 json_obj: &serde_json::Map<String, Value>,
2430 buffer: &mut [u8],
2431 current_offset: usize,
2432 options: &EncodeOptions,
2433 spec: BinarySpec,
2434) -> Result<usize> {
2435 let field_len = field.len as usize;
2436
2437 if let Some(num) = json_obj.get(&field.name).and_then(|v| {
2438 if let Some(n) = v.as_i64() {
2439 Some(n)
2441 } else if let Some(s) = coerce_to_str(v, options.coerce_numbers) {
2442 s.parse::<i64>().ok()
2443 } else {
2444 None
2445 }
2446 }) {
2447 let encoded = crate::numeric::encode_binary_int(num, spec.bits, spec.signed)?;
2448 if current_offset + field_len <= buffer.len() && encoded.len() == field_len {
2449 buffer[current_offset..current_offset + field_len].copy_from_slice(&encoded);
2450 }
2451 }
2452
2453 Ok(current_offset + field_len)
2454}
2455
2456#[inline]
2485#[must_use = "Handle the Result or propagate the error"]
2486pub fn decode_file_to_jsonl(
2487 schema: &Schema,
2488 input: impl Read,
2489 mut output: impl Write,
2490 options: &DecodeOptions,
2491) -> Result<RunSummary> {
2492 let start_time = std::time::Instant::now();
2493 let mut summary = RunSummary::with_threads(effective_worker_count(options.threads));
2494 summary.set_schema_fingerprint(schema.fingerprint.clone());
2495
2496 reset_warning_counter();
2497
2498 match options.format {
2499 RecordFormat::Fixed => {
2500 process_fixed_records(schema, input, &mut output, options, &mut summary)?;
2501 }
2502 RecordFormat::RDW => {
2503 process_rdw_records(schema, input, &mut output, options, &mut summary)?;
2504 }
2505 }
2506
2507 let elapsed_ms = start_time.elapsed().as_millis();
2508 summary.processing_time_ms = u64::try_from(elapsed_ms).unwrap_or(u64::MAX);
2509 summary.calculate_throughput();
2510 summary.warnings = warning_count();
2511 telemetry::record_completion(
2512 summary.processing_time_seconds(),
2513 summary.throughput_mbps,
2514 options,
2515 );
2516 info!(
2517 target: "copybook::decode",
2518 records_processed = summary.records_processed,
2519 records_with_errors = summary.records_with_errors,
2520 warnings = summary.warnings,
2521 bytes_processed = summary.bytes_processed,
2522 elapsed_ms = summary.processing_time_ms,
2523 throughput_mibps = summary.throughput_mbps,
2524 schema_fingerprint = %summary.schema_fingerprint,
2525 codepage = %options.codepage,
2526 format = ?options.format,
2527 strict_mode = options.strict_mode,
2528 raw_mode = ?options.emit_raw,
2529 );
2530
2531 Ok(summary)
2532}
2533
2534fn process_fixed_records<R: Read, W: Write>(
2535 schema: &Schema,
2536 reader: R,
2537 output: &mut W,
2538 options: &DecodeOptions,
2539 summary: &mut RunSummary,
2540) -> Result<()> {
2541 if options.threads > 1 {
2542 return process_fixed_records_parallel(schema, reader, output, options, summary);
2543 }
2544
2545 let mut reader = crate::record::FixedRecordReader::new(reader, schema.lrecl_fixed)?;
2546 let mut scratch = crate::memory::ScratchBuffers::new();
2547 let mut record_index = 0u64;
2548
2549 while let Some(record_data) = reader.read_record()? {
2550 record_index += 1;
2551 summary.bytes_processed += record_data.len() as u64;
2552 telemetry::record_read(record_data.len(), options);
2553
2554 let raw_data_for_decode = match options.emit_raw {
2555 crate::options::RawMode::Record => Some(record_data.clone()),
2556 _ => None,
2557 };
2558
2559 match decode_record_with_scratch_and_raw(
2560 schema,
2561 &record_data,
2562 options,
2563 raw_data_for_decode,
2564 record_index,
2565 &mut scratch,
2566 ) {
2567 Ok(json_value) => {
2568 write_json_record(output, &json_value)?;
2569 summary.records_processed += 1;
2570 }
2571 Err(error) => {
2572 summary.records_with_errors += 1;
2573 let family = error.family_prefix();
2574 telemetry::record_error(family);
2575 if options.strict_mode {
2576 return Err(error);
2577 }
2578 }
2579 }
2580 }
2581
2582 Ok(())
2583}
2584
2585struct DecodeWork {
2586 payload: Vec<u8>,
2587 raw_data: Option<Vec<u8>>,
2588 record_index: u64,
2589}
2590
2591struct DecodeOutcome {
2592 result: Result<Value>,
2593 warnings: u64,
2594}
2595
2596fn effective_worker_count(requested: usize) -> usize {
2597 requested.clamp(1, MAX_WORKERS)
2598}
2599
2600fn decode_worker_pool(
2601 schema: &Schema,
2602 options: &DecodeOptions,
2603) -> crate::memory::WorkerPool<DecodeWork, DecodeOutcome> {
2604 let workers = effective_worker_count(options.threads);
2605 let channel_capacity = workers.saturating_mul(4).max(1);
2606 let max_window_size = workers.saturating_mul(2).max(1);
2607 let schema = Arc::new(schema.clone());
2608 let options = Arc::new(options.clone());
2609
2610 crate::memory::WorkerPool::new(
2611 workers,
2612 channel_capacity,
2613 max_window_size,
2614 move |work: DecodeWork, scratch: &mut crate::memory::ScratchBuffers| {
2615 let warning_count_before = warning_count();
2616 let result = decode_record_with_scratch_and_raw(
2617 &schema,
2618 &work.payload,
2619 &options,
2620 work.raw_data,
2621 work.record_index,
2622 scratch,
2623 );
2624 let warnings = warning_count().saturating_sub(warning_count_before);
2625 DecodeOutcome { result, warnings }
2626 },
2627 )
2628}
2629
2630fn process_decode_batch<W: Write>(
2631 pool: &mut crate::memory::WorkerPool<DecodeWork, DecodeOutcome>,
2632 batch_len: usize,
2633 output: &mut W,
2634 options: &DecodeOptions,
2635 summary: &mut RunSummary,
2636) -> Result<()> {
2637 let mut first_error = None;
2638
2639 for _ in 0..batch_len {
2640 let outcome = pool
2641 .recv_ordered()
2642 .map_err(|error| Error::new(ErrorCode::CBKI001_INVALID_STATE, error.to_string()))?
2643 .ok_or_else(|| {
2644 Error::new(
2645 ErrorCode::CBKI001_INVALID_STATE,
2646 "decode worker pool ended before the submitted batch completed",
2647 )
2648 })?;
2649
2650 for _ in 0..outcome.warnings {
2651 increment_warning_counter();
2652 }
2653
2654 if first_error.is_some() {
2655 continue;
2656 }
2657
2658 match outcome.result {
2659 Ok(json_value) => {
2660 write_json_record(output, &json_value)?;
2661 summary.records_processed += 1;
2662 }
2663 Err(error) => {
2664 summary.records_with_errors += 1;
2665 let family = error.family_prefix();
2666 telemetry::record_error(family);
2667 if options.strict_mode {
2668 first_error = Some(error);
2669 }
2670 }
2671 }
2672 }
2673
2674 first_error.map_or(Ok(()), Err)
2675}
2676
2677fn process_fixed_records_parallel<R: Read, W: Write>(
2678 schema: &Schema,
2679 reader: R,
2680 output: &mut W,
2681 options: &DecodeOptions,
2682 summary: &mut RunSummary,
2683) -> Result<()> {
2684 let mut reader = crate::record::FixedRecordReader::new(reader, schema.lrecl_fixed)?;
2685 let workers = effective_worker_count(options.threads);
2686 let batch_capacity = workers.saturating_mul(4).max(1);
2687 let mut pool = decode_worker_pool(schema, options);
2688 let mut record_index = 0_u64;
2689 let mut batch_len = 0_usize;
2690
2691 loop {
2692 let record = match reader.read_record() {
2693 Ok(record) => record,
2694 Err(error) => {
2695 let pending_result = if batch_len > 0 {
2696 process_decode_batch(&mut pool, batch_len, output, options, summary)
2697 } else {
2698 Ok(())
2699 };
2700 let _ = pool.shutdown();
2701 pending_result?;
2702 return Err(error);
2703 }
2704 };
2705 let Some(record_data) = record else { break };
2706
2707 record_index += 1;
2708 summary.bytes_processed += record_data.len() as u64;
2709 telemetry::record_read(record_data.len(), options);
2710 let raw_data = match options.emit_raw {
2711 crate::options::RawMode::Record => Some(record_data.clone()),
2712 _ => None,
2713 };
2714
2715 if let Err(error) = pool.submit(DecodeWork {
2716 payload: record_data,
2717 raw_data,
2718 record_index,
2719 }) {
2720 let pending_result = if batch_len > 0 {
2721 process_decode_batch(&mut pool, batch_len, output, options, summary)
2722 } else {
2723 Ok(())
2724 };
2725 let _ = pool.shutdown();
2726 pending_result?;
2727 return Err(Error::new(
2728 ErrorCode::CBKI001_INVALID_STATE,
2729 error.to_string(),
2730 ));
2731 }
2732 batch_len += 1;
2733
2734 if batch_len == batch_capacity {
2735 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
2736 batch_len = 0;
2737 if let Err(error) = result {
2738 let _ = pool.shutdown();
2739 return Err(error);
2740 }
2741 }
2742 }
2743
2744 if batch_len > 0 {
2745 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
2746 if let Err(error) = result {
2747 let _ = pool.shutdown();
2748 return Err(error);
2749 }
2750 }
2751
2752 pool.shutdown().map_err(|error| {
2753 Error::new(
2754 ErrorCode::CBKI001_INVALID_STATE,
2755 format!("decode worker pool shutdown failed: {error}"),
2756 )
2757 })
2758}
2759
2760fn process_rdw_records<R: Read, W: Write>(
2761 schema: &Schema,
2762 reader: R,
2763 output: &mut W,
2764 options: &DecodeOptions,
2765 summary: &mut RunSummary,
2766) -> Result<()> {
2767 if options.threads > 1 {
2768 return process_rdw_records_parallel(schema, reader, output, options, summary);
2769 }
2770
2771 let mut reader = crate::record::RDWRecordReader::new(reader, options.strict_mode);
2772 let mut scratch = crate::memory::ScratchBuffers::new();
2773 let mut record_index = 0u64;
2774
2775 while let Some(rdw_record) = reader.read_record()? {
2776 record_index += 1;
2777 let record_bytes = rdw_record.header.len() + rdw_record.payload.len();
2778 summary.bytes_processed += record_bytes as u64;
2779 telemetry::record_read(record_bytes, options);
2780 if rdw_record.reserved() != 0 {
2781 increment_warning_counter();
2782 }
2783
2784 if let Some(schema_lrecl) = schema.lrecl_fixed
2792 && schema.tail_odo.is_none()
2793 && rdw_record.payload.len() < schema_lrecl as usize
2794 {
2795 let error = rdw_underflow_error(schema_lrecl, rdw_record.payload.len());
2796
2797 summary.records_with_errors += 1;
2798 let family = error.family_prefix();
2799 telemetry::record_error(family);
2800 if options.strict_mode {
2801 return Err(error);
2802 }
2803 continue;
2804 }
2805
2806 let full_raw_data = rdw_raw_data(&rdw_record, options.emit_raw);
2807
2808 match decode_record_with_scratch_and_raw(
2809 schema,
2810 &rdw_record.payload,
2811 options,
2812 full_raw_data,
2813 record_index,
2814 &mut scratch,
2815 ) {
2816 Ok(json_value) => {
2817 write_json_record(output, &json_value)?;
2818 summary.records_processed += 1;
2819 }
2820 Err(error) => {
2821 summary.records_with_errors += 1;
2822 let family = error.family_prefix();
2823 telemetry::record_error(family);
2824 if options.strict_mode {
2825 return Err(error);
2826 }
2827 }
2828 }
2829 }
2830
2831 Ok(())
2832}
2833
2834fn rdw_underflow_error(schema_lrecl: u32, payload_len: usize) -> Error {
2835 Error::new(
2836 ErrorCode::CBKF221_RDW_UNDERFLOW,
2837 format!("RDW payload too short: {payload_len} bytes, schema requires {schema_lrecl} bytes"),
2838 )
2839}
2840
2841fn rdw_raw_data(
2842 record: &crate::record::RDWRecord,
2843 raw_mode: crate::options::RawMode,
2844) -> Option<Vec<u8>> {
2845 match raw_mode {
2846 crate::options::RawMode::RecordRDW => {
2847 let mut full_data = Vec::with_capacity(record.header.len() + record.payload.len());
2848 full_data.extend_from_slice(&record.header);
2849 full_data.extend_from_slice(&record.payload);
2850 Some(full_data)
2851 }
2852 crate::options::RawMode::Record => Some(record.payload.clone()),
2853 _ => None,
2854 }
2855}
2856
2857fn process_rdw_records_parallel<R: Read, W: Write>(
2858 schema: &Schema,
2859 reader: R,
2860 output: &mut W,
2861 options: &DecodeOptions,
2862 summary: &mut RunSummary,
2863) -> Result<()> {
2864 let mut reader = crate::record::RDWRecordReader::new(reader, options.strict_mode);
2865 let workers = effective_worker_count(options.threads);
2866 let batch_capacity = workers.saturating_mul(4).max(1);
2867 let mut pool = decode_worker_pool(schema, options);
2868 let mut record_index = 0_u64;
2869 let mut batch_len = 0_usize;
2870
2871 loop {
2872 let rdw_record = match reader.read_record() {
2873 Ok(record) => record,
2874 Err(error) => {
2875 let pending_result = if batch_len > 0 {
2876 process_decode_batch(&mut pool, batch_len, output, options, summary)
2877 } else {
2878 Ok(())
2879 };
2880 let _ = pool.shutdown();
2881 pending_result?;
2882 return Err(error);
2883 }
2884 };
2885 let Some(rdw_record) = rdw_record else { break };
2886
2887 record_index += 1;
2888 let record_bytes = rdw_record.header.len() + rdw_record.payload.len();
2889 summary.bytes_processed += record_bytes as u64;
2890 telemetry::record_read(record_bytes, options);
2891 if rdw_record.reserved() != 0 {
2892 increment_warning_counter();
2893 }
2894
2895 if let Some(schema_lrecl) = schema.lrecl_fixed
2896 && rdw_record.payload.len() < schema_lrecl as usize
2897 {
2898 let error = rdw_underflow_error(schema_lrecl, rdw_record.payload.len());
2899
2900 summary.records_with_errors += 1;
2901 let family = error.family_prefix();
2902 telemetry::record_error(family);
2903 if options.strict_mode {
2904 let pending_result = if batch_len > 0 {
2905 process_decode_batch(&mut pool, batch_len, output, options, summary)
2906 } else {
2907 Ok(())
2908 };
2909 let _ = pool.shutdown();
2910 pending_result?;
2911 return Err(error);
2912 }
2913 continue;
2914 }
2915
2916 let full_raw_data = rdw_raw_data(&rdw_record, options.emit_raw);
2917
2918 if let Err(error) = pool.submit(DecodeWork {
2919 payload: rdw_record.payload,
2920 raw_data: full_raw_data,
2921 record_index,
2922 }) {
2923 let pending_result = if batch_len > 0 {
2924 process_decode_batch(&mut pool, batch_len, output, options, summary)
2925 } else {
2926 Ok(())
2927 };
2928 let _ = pool.shutdown();
2929 pending_result?;
2930 return Err(Error::new(
2931 ErrorCode::CBKI001_INVALID_STATE,
2932 error.to_string(),
2933 ));
2934 }
2935 batch_len += 1;
2936
2937 if batch_len == batch_capacity {
2938 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
2939 batch_len = 0;
2940 if let Err(error) = result {
2941 let _ = pool.shutdown();
2942 return Err(error);
2943 }
2944 }
2945 }
2946
2947 if batch_len > 0 {
2948 let result = process_decode_batch(&mut pool, batch_len, output, options, summary);
2949 if let Err(error) = result {
2950 let _ = pool.shutdown();
2951 return Err(error);
2952 }
2953 }
2954
2955 pool.shutdown().map_err(|error| {
2956 Error::new(
2957 ErrorCode::CBKI001_INVALID_STATE,
2958 format!("decode worker pool shutdown failed: {error}"),
2959 )
2960 })
2961}
2962
2963#[inline]
2964fn write_json_record<W: Write>(output: &mut W, value: &Value) -> Result<()> {
2965 if let Err(e) = serde_json::to_writer(&mut *output, value) {
2966 let error = Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string());
2967 telemetry::record_error(error.family_prefix());
2968 return Err(error);
2969 }
2970
2971 if let Err(e) = writeln!(output) {
2972 let error = Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string());
2973 telemetry::record_error(error.family_prefix());
2974 return Err(error);
2975 }
2976
2977 Ok(())
2978}
2979
2980fn encode_worker_pool(
2981 schema: &Schema,
2982 options: &EncodeOptions,
2983) -> crate::memory::WorkerPool<Value, Result<Vec<u8>>> {
2984 let workers = effective_worker_count(options.threads);
2985 let channel_capacity = workers.saturating_mul(4).max(1);
2986 let max_window_size = workers.saturating_mul(2).max(1);
2987 let schema = Arc::new(schema.clone());
2988 let options = Arc::new(options.clone());
2989
2990 crate::memory::WorkerPool::new(
2991 workers,
2992 channel_capacity,
2993 max_window_size,
2994 move |json_value: Value, _scratch: &mut crate::memory::ScratchBuffers| {
2995 encode_record(&schema, &json_value, &options)
2996 },
2997 )
2998}
2999
3000fn process_encode_batch<W: Write>(
3001 pool: &mut crate::memory::WorkerPool<Value, Result<Vec<u8>>>,
3002 batch_len: usize,
3003 records_before_batch: u64,
3004 output: &mut W,
3005 options: &EncodeOptions,
3006 summary: &mut RunSummary,
3007) -> Result<bool> {
3008 let mut first_error = None;
3009
3010 for position in 0..batch_len {
3011 let result = pool
3012 .recv_ordered()
3013 .map_err(|error| Error::new(ErrorCode::CBKI001_INVALID_STATE, error.to_string()))?
3014 .ok_or_else(|| {
3015 Error::new(
3016 ErrorCode::CBKI001_INVALID_STATE,
3017 "encode worker pool ended before the submitted batch completed",
3018 )
3019 })?;
3020
3021 if first_error.is_some() {
3022 continue;
3023 }
3024
3025 match result {
3026 Ok(binary_data) => {
3027 output.write_all(&binary_data).map_err(|error| {
3028 Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, error.to_string())
3029 })?;
3030 summary.bytes_processed += binary_data.len() as u64;
3031 }
3032 Err(error) => {
3033 summary.records_with_errors += 1;
3034 telemetry::record_error(error.family_prefix());
3035 if options.strict_mode {
3036 first_error = Some((position as u64 + 1, error));
3037 }
3038 }
3039 }
3040 }
3041
3042 if let Some((position, _error)) = first_error {
3043 summary.records_processed = records_before_batch + position;
3044 return Ok(true);
3045 }
3046
3047 Ok(false)
3048}
3049
3050fn shutdown_encode_pool(pool: crate::memory::WorkerPool<Value, Result<Vec<u8>>>) -> Result<()> {
3051 pool.shutdown().map_err(|error| {
3052 Error::new(
3053 ErrorCode::CBKI001_INVALID_STATE,
3054 format!("encode worker pool shutdown failed: {error}"),
3055 )
3056 })
3057}
3058
3059fn finish_encode_input_error<W: Write>(
3060 pool: crate::memory::WorkerPool<Value, Result<Vec<u8>>>,
3061 batch_len: usize,
3062 records_before_batch: u64,
3063 output: &mut W,
3064 options: &EncodeOptions,
3065 summary: &mut RunSummary,
3066 error: Error,
3067) -> Result<u64> {
3068 let mut pool = pool;
3069 let pending_result = if batch_len > 0 {
3070 process_encode_batch(
3071 &mut pool,
3072 batch_len,
3073 records_before_batch,
3074 output,
3075 options,
3076 summary,
3077 )
3078 } else {
3079 Ok(false)
3080 };
3081 let shutdown_result = shutdown_encode_pool(pool);
3082 let pending_stop = pending_result?;
3083 shutdown_result?;
3084 if pending_stop {
3085 Ok(summary.records_processed)
3086 } else {
3087 Err(error)
3088 }
3089}
3090
3091fn process_encode_jsonl_parallel<R: BufRead, W: Write>(
3092 schema: &Schema,
3093 reader: R,
3094 output: &mut W,
3095 options: &EncodeOptions,
3096 summary: &mut RunSummary,
3097) -> Result<u64> {
3098 let workers = effective_worker_count(options.threads);
3099 let batch_capacity = workers.saturating_mul(4).max(1);
3100 let mut pool = encode_worker_pool(schema, options);
3101 let mut records_seen = 0_u64;
3102 let mut records_before_batch = 0_u64;
3103 let mut batch_len = 0_usize;
3104
3105 for line in reader.lines() {
3106 let line = match line {
3107 Ok(line) => line,
3108 Err(error) => {
3109 return finish_encode_input_error(
3110 pool,
3111 batch_len,
3112 records_before_batch,
3113 output,
3114 options,
3115 summary,
3116 Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, error.to_string()),
3117 );
3118 }
3119 };
3120
3121 if line.trim().is_empty() {
3122 continue;
3123 }
3124
3125 let json_value: Value = match serde_json::from_str(&line) {
3126 Ok(json_value) => json_value,
3127 Err(error) => {
3128 return finish_encode_input_error(
3129 pool,
3130 batch_len,
3131 records_before_batch,
3132 output,
3133 options,
3134 summary,
3135 Error::new(ErrorCode::CBKE501_JSON_TYPE_MISMATCH, error.to_string()),
3136 );
3137 }
3138 };
3139
3140 records_seen += 1;
3141 if let Err(error) = pool.submit(json_value) {
3142 return finish_encode_input_error(
3143 pool,
3144 batch_len,
3145 records_before_batch,
3146 output,
3147 options,
3148 summary,
3149 Error::new(ErrorCode::CBKI001_INVALID_STATE, error.to_string()),
3150 );
3151 }
3152 batch_len += 1;
3153
3154 if batch_len == batch_capacity {
3155 let batch_result = process_encode_batch(
3156 &mut pool,
3157 batch_len,
3158 records_before_batch,
3159 output,
3160 options,
3161 summary,
3162 );
3163 let stop = match batch_result {
3164 Ok(stop) => stop,
3165 Err(error) => {
3166 let _ = shutdown_encode_pool(pool);
3167 return Err(error);
3168 }
3169 };
3170 batch_len = 0;
3171 records_before_batch = records_seen;
3172 if stop {
3173 shutdown_encode_pool(pool)?;
3174 return Ok(summary.records_processed);
3175 }
3176 }
3177 }
3178
3179 if batch_len > 0 {
3180 let batch_result = process_encode_batch(
3181 &mut pool,
3182 batch_len,
3183 records_before_batch,
3184 output,
3185 options,
3186 summary,
3187 );
3188 let stop = match batch_result {
3189 Ok(stop) => stop,
3190 Err(error) => {
3191 let _ = shutdown_encode_pool(pool);
3192 return Err(error);
3193 }
3194 };
3195 if stop {
3196 shutdown_encode_pool(pool)?;
3197 return Ok(summary.records_processed);
3198 }
3199 }
3200
3201 shutdown_encode_pool(pool)?;
3202 summary.records_processed = records_seen;
3203 Ok(records_seen)
3204}
3205
3206#[inline]
3243#[must_use = "Handle the Result or propagate the error"]
3244pub fn encode_jsonl_to_file(
3245 schema: &Schema,
3246 input: impl Read,
3247 mut output: impl Write,
3248 options: &EncodeOptions,
3249) -> Result<RunSummary> {
3250 let start_time = std::time::Instant::now();
3251 let mut summary = RunSummary::with_threads(effective_worker_count(options.threads));
3252 summary.set_schema_fingerprint(schema.fingerprint.clone());
3253
3254 let reader = BufReader::new(input);
3255 let record_count = if options.threads > 1 {
3256 process_encode_jsonl_parallel(schema, reader, &mut output, options, &mut summary)?
3257 } else {
3258 let mut record_count = 0u64;
3259
3260 for line in reader.lines() {
3261 let line =
3262 line.map_err(|e| Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string()))?;
3263
3264 if line.trim().is_empty() {
3265 continue;
3266 }
3267
3268 record_count += 1;
3269
3270 let json_value: Value = serde_json::from_str(&line)
3272 .map_err(|e| Error::new(ErrorCode::CBKE501_JSON_TYPE_MISMATCH, e.to_string()))?;
3273
3274 if let Ok(binary_data) = encode_record(schema, &json_value, options) {
3276 output
3277 .write_all(&binary_data)
3278 .map_err(|e| Error::new(ErrorCode::CBKC201_JSON_WRITE_ERROR, e.to_string()))?;
3279 summary.bytes_processed += binary_data.len() as u64;
3280 } else {
3281 summary.records_with_errors += 1;
3282 if options.strict_mode {
3283 break;
3284 }
3285 }
3286 }
3287
3288 record_count
3289 };
3290
3291 summary.records_processed = record_count;
3292 let elapsed_ms = start_time.elapsed().as_millis();
3293 summary.processing_time_ms = u64::try_from(elapsed_ms).unwrap_or(u64::MAX);
3294 summary.calculate_throughput();
3295
3296 Ok(summary)
3297}
3298
3299fn format_zoned_decimal_with_digits(
3301 decimal: &crate::numeric::SmallDecimal,
3302 digits: u16,
3303 blank_when_zero: bool,
3304) -> String {
3305 use std::fmt::Write;
3306
3307 if blank_when_zero {
3309 return decimal.to_string();
3310 }
3311
3312 if decimal.value == 0 {
3315 let natural_format = decimal.to_string();
3316 if natural_format == "0" {
3317 return "0".to_string();
3318 }
3319 }
3320
3321 let mut result = String::new();
3323 let value = decimal.value;
3324 let negative = decimal.negative && value != 0;
3325
3326 if negative {
3327 result.push('-');
3328 }
3329
3330 if decimal.scale <= 0 {
3332 let scaled_value = if decimal.scale < 0 {
3333 let exponent = u32::from(decimal.scale.unsigned_abs());
3334 value * 10_i64.pow(exponent)
3335 } else {
3336 value
3337 };
3338 if write!(result, "{:0width$}", scaled_value, width = digits as usize).is_err() {
3339 result.push('0');
3341 }
3342 } else {
3343 result.push_str(&decimal.to_string());
3345 }
3346
3347 result
3348}
3349
3350#[inline]
3351fn small_decimal_to_string(decimal: &crate::numeric::SmallDecimal) -> String {
3352 decimal.to_string()
3353}
3354
3355fn zoned_decimal_to_json_value(
3356 decimal: &crate::numeric::SmallDecimal,
3357 digits: u16,
3358 scale: i16,
3359 blank_when_zero: bool,
3360 options: &DecodeOptions,
3361) -> Value {
3362 let formatted = if scale == 0 {
3363 format_zoned_decimal_with_digits(decimal, digits, blank_when_zero)
3364 } else {
3365 small_decimal_to_string(decimal)
3366 };
3367 numeric_string_to_value(formatted, options)
3368}
3369
3370#[inline]
3371fn decimal_counter_to_u32(
3372 decimal: &crate::numeric::SmallDecimal,
3373 counter_path: &str,
3374) -> Result<u32> {
3375 let text = small_decimal_to_string(decimal);
3376 text.parse::<u32>().map_err(|_| {
3377 Error::new(
3378 ErrorCode::CBKS121_COUNTER_NOT_FOUND,
3379 format!("ODO counter '{counter_path}' has invalid value: {text}"),
3380 )
3381 })
3382}
3383
3384#[cfg(test)]
3385#[allow(clippy::expect_used)]
3386#[allow(clippy::unwrap_used)]
3387mod tests {
3388 use super::*;
3389 use crate::Codepage;
3390 use crate::iterator::RecordIterator;
3391 use copybook_core::{Error, ErrorCode, Result, parse_copybook};
3392 use std::io::Cursor;
3393
3394 #[test]
3395 fn test_decode_record() -> Result<()> {
3396 let copybook_text = r"
3397 01 RECORD.
3398 05 ID PIC 9(3).
3399 05 NAME PIC X(5).
3400 ";
3401
3402 let schema = parse_copybook(copybook_text)?;
3403 let options = DecodeOptions {
3404 codepage: Codepage::ASCII, ..DecodeOptions::default()
3406 };
3407 let data = b"001ALICE";
3408
3409 let result = decode_record(&schema, data, &options)?;
3410 assert!(result.is_object());
3411 let object = result.as_object().ok_or_else(|| {
3412 Error::new(
3413 ErrorCode::CBKP001_SYNTAX,
3414 "decoded record should be an object".to_string(),
3415 )
3416 })?;
3417 assert!(object.len() > 1);
3418 Ok(())
3419 }
3420
3421 #[test]
3422 fn test_encode_record() -> Result<()> {
3423 let copybook_text = r"
3424 01 RECORD.
3425 05 ID PIC 9(3).
3426 05 NAME PIC X(5).
3427 ";
3428
3429 let schema = parse_copybook(copybook_text)?;
3430 let options = EncodeOptions::default();
3431
3432 let mut json_obj = serde_json::Map::new();
3433 json_obj.insert("ID".into(), Value::String("123".into()));
3434 json_obj.insert("NAME".into(), Value::String("HELLO".into()));
3435 let json = Value::Object(json_obj);
3436
3437 let result = encode_record(&schema, &json, &options)?;
3438 assert!(!result.is_empty());
3439 assert_eq!(result.len(), 8); Ok(())
3443 }
3444
3445 #[test]
3446 fn test_record_iterator() -> Result<()> {
3447 let copybook_text = r"
3448 01 RECORD.
3449 05 ID PIC 9(3).
3450 05 NAME PIC X(5).
3451 ";
3452
3453 let schema = parse_copybook(copybook_text)?;
3454 let options = DecodeOptions::default();
3455
3456 let test_data = vec![0u8; 16]; let cursor = Cursor::new(test_data);
3459
3460 let iterator = RecordIterator::new(cursor, &schema, &options)?;
3461 assert_eq!(iterator.current_record_index(), 0);
3462 assert!(!iterator.is_eof());
3463 Ok(())
3464 }
3465
3466 #[test]
3467 fn test_decode_file_to_jsonl() -> Result<()> {
3468 let copybook_text = r"
3469 01 RECORD.
3470 05 ID PIC 9(3).
3471 05 NAME PIC X(5).
3472 ";
3473
3474 let schema = parse_copybook(copybook_text)?;
3475 let options = DecodeOptions {
3476 codepage: Codepage::ASCII, ..DecodeOptions::default()
3478 };
3479
3480 let input_data = b"001ALICE002BOBBY".to_vec(); let input = Cursor::new(input_data);
3483
3484 let mut output = Vec::new();
3486
3487 let summary = decode_file_to_jsonl(&schema, input, &mut output, &options)?;
3488 assert!(summary.records_processed > 0);
3489 assert!(!output.is_empty());
3490 Ok(())
3491 }
3492
3493 #[test]
3494 fn test_encode_jsonl_to_file() -> Result<()> {
3495 let copybook_text = r"
3496 01 RECORD.
3497 05 ID PIC 9(3).
3498 05 NAME PIC X(5).
3499 ";
3500
3501 let schema = parse_copybook(copybook_text)?;
3502 let options = EncodeOptions::default();
3503
3504 let jsonl_data = "{\"__status\":\"test\"}\n{\"__status\":\"test2\"}";
3506 let input = Cursor::new(jsonl_data.as_bytes());
3507
3508 let mut output = Vec::new();
3510
3511 let summary = encode_jsonl_to_file(&schema, input, &mut output, &options)?;
3512 assert_eq!(summary.records_processed, 2);
3513 assert!(!output.is_empty());
3514 Ok(())
3515 }
3516}