1use std::collections::HashMap;
4use std::sync::Mutex;
5
6use mig_assembly::ConversionService;
7use mig_bo4e::engine::DataBundle;
8use mig_bo4e::MappingEngine;
9
10use crate::data_dir::DataDir;
11use crate::error::MapperError;
12
13pub struct Bo4eResult {
15 pub pid: String,
17 pub message_type: String,
19 pub variant: String,
21 pub bo4e: serde_json::Value,
23}
24
25#[derive(Debug, Clone)]
62pub struct PidListEntry {
63 pub fv: String,
64 pub variant: String,
65 pub pid: String,
66 pub beschreibung: String,
67}
68
69pub struct Mapper {
70 data_dir: DataDir,
71 bundles: Mutex<HashMap<String, DataBundle>>,
72}
73
74fn split_transaktion(tx: &serde_json::Value) -> mig_bo4e::model::MappedTransaktion {
87 let (transaktionsdaten, stammdaten) = if mig_bo4e::model::is_wrapped_transaktion(tx) {
88 (
89 tx.get("transaktionsdaten")
90 .cloned()
91 .unwrap_or(serde_json::Value::Null),
92 tx.get("stammdaten")
93 .cloned()
94 .unwrap_or_else(|| serde_json::Value::Object(Default::default())),
95 )
96 } else {
97 (serde_json::Value::Null, tx.clone())
98 };
99 mig_bo4e::model::MappedTransaktion {
100 transaktionsdaten,
101 stammdaten,
102 nesting_info: Default::default(),
103 }
104}
105
106impl Mapper {
107 pub fn from_data_dir(data_dir: DataDir) -> Result<Self, MapperError> {
112 let mapper = Self {
113 data_dir,
114 bundles: Mutex::new(HashMap::new()),
115 };
116 let eager_fvs: Vec<String> = mapper.data_dir.eager_fvs().to_vec();
117 for fv in &eager_fvs {
118 mapper.ensure_bundle_loaded(fv)?;
119 }
120 Ok(mapper)
121 }
122
123 fn ensure_bundle_loaded(&self, fv: &str) -> Result<(), MapperError> {
125 let mut bundles = self.bundles.lock().unwrap();
126 if bundles.contains_key(fv) {
127 return Ok(());
128 }
129 let path = self.data_dir.bundle_path(fv);
130 if !path.exists() {
131 return Err(MapperError::BundleNotFound { fv: fv.to_string() });
132 }
133 let bundle = DataBundle::load(&path)?;
134 bundles.insert(fv.to_string(), bundle);
135 Ok(())
136 }
137
138 pub fn conversion_service(
142 &self,
143 fv: &str,
144 variant: &str,
145 ) -> Result<ConversionService, MapperError> {
146 self.ensure_bundle_loaded(fv)?;
147 let bundles = self.bundles.lock().unwrap();
148 let bundle = bundles.get(fv).unwrap();
149 let vc = bundle
150 .variant(variant)
151 .ok_or_else(|| MapperError::VariantNotFound {
152 fv: fv.to_string(),
153 variant: variant.to_string(),
154 })?;
155 let mig = vc
156 .mig_schema
157 .as_ref()
158 .ok_or_else(|| MapperError::VariantNotFound {
159 fv: fv.to_string(),
160 variant: format!("{variant} (no MIG schema in bundle)"),
161 })?;
162 Ok(ConversionService::from_mig(mig.clone()))
163 }
164
165 pub fn engine(&self, fv: &str, variant: &str, pid: &str) -> Result<MappingEngine, MapperError> {
169 self.ensure_bundle_loaded(fv)?;
170 let bundles = self.bundles.lock().unwrap();
171 let bundle = bundles.get(fv).unwrap();
172 let vc = bundle
173 .variant(variant)
174 .ok_or_else(|| MapperError::VariantNotFound {
175 fv: fv.to_string(),
176 variant: variant.to_string(),
177 })?;
178 let pid_key = format!("pid_{pid}");
179 let defs = vc
180 .combined_defs
181 .get(&pid_key)
182 .ok_or_else(|| MapperError::PidNotFound {
183 fv: fv.to_string(),
184 variant: variant.to_string(),
185 pid: pid.to_string(),
186 })?;
187 Ok(MappingEngine::from_definitions(defs.clone()))
188 }
189
190 pub fn pid_requirements(
195 &self,
196 fv: &str,
197 variant: &str,
198 pid: &str,
199 ) -> Result<mig_bo4e::pid_requirements::PidRequirements, MapperError> {
200 self.ensure_bundle_loaded(fv)?;
201 let bundles = self.bundles.lock().unwrap();
202 let bundle = bundles.get(fv).unwrap();
203 let vc = bundle
204 .variant(variant)
205 .ok_or_else(|| MapperError::VariantNotFound {
206 fv: fv.to_string(),
207 variant: variant.to_string(),
208 })?;
209 let pid_key = format!("pid_{pid}");
210 vc.pid_requirements
211 .get(&pid_key)
212 .cloned()
213 .ok_or_else(|| MapperError::PidNotFound {
214 fv: fv.to_string(),
215 variant: variant.to_string(),
216 pid: pid.to_string(),
217 })
218 }
219
220 pub fn bo4e_catalog(
226 &self,
227 fv: &str,
228 ) -> Result<mig_bo4e::bo4e_catalog::Bo4eCatalog, MapperError> {
229 self.ensure_bundle_loaded(fv)?;
230 let bundles = self.bundles.lock().unwrap();
231 let bundle = bundles.get(fv).unwrap();
232 Ok(bundle.bo4e_catalog.clone())
233 }
234
235 pub fn list_pids(&self) -> Result<Vec<PidListEntry>, MapperError> {
240 let dir = self.data_dir.data_path();
241 let read_dir = std::fs::read_dir(dir).map_err(|_| MapperError::DataDirNotFound {
242 path: dir.display().to_string(),
243 })?;
244
245 let mut result = Vec::new();
246
247 for entry in read_dir.flatten() {
248 let path = entry.path();
249 if path.extension().is_some_and(|e| e == "bin") {
250 let stem = path
251 .file_stem()
252 .and_then(|s| s.to_str())
253 .unwrap_or("")
254 .to_string();
255 let fv = match stem.strip_prefix("edifact-data-") {
256 Some(v) => v.to_string(),
257 None => continue,
258 };
259 self.ensure_bundle_loaded(&fv)?;
260 let bundles = self.bundles.lock().unwrap();
261 if let Some(bundle) = bundles.get(&fv) {
262 for (variant, vc) in &bundle.variants {
263 for (pid_key, req) in &vc.pid_requirements {
264 let pid = pid_key.strip_prefix("pid_").unwrap_or(pid_key).to_string();
265 result.push(PidListEntry {
266 fv: fv.clone(),
267 variant: variant.clone(),
268 pid,
269 beschreibung: req.beschreibung.clone(),
270 });
271 }
272 }
273 }
274 }
275 }
276
277 result.sort_by(|a, b| a.pid.cmp(&b.pid));
278 Ok(result)
279 }
280
281 pub fn validate_pid(
286 &self,
287 json: &serde_json::Value,
288 fv: &str,
289 variant: &str,
290 pid: &str,
291 ) -> Result<Vec<mig_bo4e::PidValidationError>, MapperError> {
292 self.ensure_bundle_loaded(fv)?;
293 let bundles = self.bundles.lock().unwrap();
294 let bundle = bundles.get(fv).unwrap();
295 let vc = bundle
296 .variant(variant)
297 .ok_or_else(|| MapperError::VariantNotFound {
298 fv: fv.to_string(),
299 variant: variant.to_string(),
300 })?;
301 let pid_key = format!("pid_{pid}");
302 let requirements =
303 vc.pid_requirements
304 .get(&pid_key)
305 .ok_or_else(|| MapperError::PidNotFound {
306 fv: fv.to_string(),
307 variant: variant.to_string(),
308 pid: pid.to_string(),
309 })?;
310
311 Ok(mig_bo4e::pid_validation::validate_pid_json(
312 json,
313 requirements,
314 ))
315 }
316
317 pub fn validate_pid_struct(
329 &self,
330 value: &impl serde::Serialize,
331 fv: &str,
332 variant: &str,
333 pid: &str,
334 ) -> Result<Vec<mig_bo4e::PidValidationError>, MapperError> {
335 let json = serde_json::to_value(value).map_err(|e| {
336 MapperError::Mapping(mig_bo4e::MappingError::TypeConversion(e.to_string()))
337 })?;
338 self.validate_pid(&json, fv, variant, pid)
339 }
340
341 pub fn validate_pid_with_conditions(
349 &self,
350 json: &serde_json::Value,
351 fv: &str,
352 variant: &str,
353 pid: &str,
354 ) -> Result<Vec<mig_bo4e::PidValidationError>, MapperError> {
355 self.ensure_bundle_loaded(fv)?;
356 let bundles = self.bundles.lock().unwrap();
357 let bundle = bundles.get(fv).unwrap();
358 let vc = bundle
359 .variant(variant)
360 .ok_or_else(|| MapperError::VariantNotFound {
361 fv: fv.to_string(),
362 variant: variant.to_string(),
363 })?;
364 let pid_key = format!("pid_{pid}");
365
366 let requirements =
367 vc.pid_requirements
368 .get(&pid_key)
369 .ok_or_else(|| MapperError::PidNotFound {
370 fv: fv.to_string(),
371 variant: variant.to_string(),
372 pid: pid.to_string(),
373 })?;
374
375 let evaluator = crate::evaluator_factory::create_evaluator(variant, fv);
377
378 if let Some(evaluator) = evaluator {
379 let defs = vc
381 .combined_defs
382 .get(&pid_key)
383 .ok_or_else(|| MapperError::PidNotFound {
384 fv: fv.to_string(),
385 variant: variant.to_string(),
386 pid: pid.to_string(),
387 })?;
388 let engine = MappingEngine::from_definitions(defs.clone());
389 let tree = engine.map_all_reverse(json, None);
390
391 let segments = crate::tree_to_segments::tree_to_owned_segments(&tree);
393
394 Ok(crate::evaluator_factory::validate_with_boxed_evaluator(
396 evaluator.as_ref(),
397 json,
398 requirements,
399 pid,
400 &segments,
401 ))
402 } else {
403 Ok(mig_bo4e::pid_validation::validate_pid_json_transaction(
405 json,
406 requirements,
407 ))
408 }
409 }
410
411 pub fn to_edifact(
455 &self,
456 msg_stammdaten: &serde_json::Value,
457 tx_stammdaten: &[serde_json::Value],
458 fv: &str,
459 variant: &str,
460 pid: &str,
461 ) -> Result<String, MapperError> {
462 self.render_message_body(
463 msg_stammdaten,
464 tx_stammdaten,
465 fv,
466 variant,
467 pid,
468 EntrySegmentCheck::Refuse,
469 )
470 }
471
472 pub fn to_edifact_nachricht(
499 &self,
500 nachricht: &mig_bo4e::model::Nachricht<serde_json::Value, serde_json::Value>,
501 fv: &str,
502 variant: &str,
503 pid: &str,
504 ) -> Result<String, MapperError> {
505 let mut msg_stammdaten = nachricht.stammdaten.clone();
506 mig_bo4e::model::restore_message_metadata(&mut msg_stammdaten, &nachricht.nachrichtendaten);
507 self.to_edifact(&msg_stammdaten, &nachricht.transaktionen, fv, variant, pid)
508 }
509
510 fn render_message_body(
518 &self,
519 msg_stammdaten: &serde_json::Value,
520 tx_stammdaten: &[serde_json::Value],
521 fv: &str,
522 variant: &str,
523 pid: &str,
524 check: EntrySegmentCheck,
525 ) -> Result<String, MapperError> {
526 self.ensure_bundle_loaded(fv)?;
527 let bundles = self.bundles.lock().unwrap();
528 let bundle = bundles.get(fv).unwrap();
529 let vc = bundle
530 .variant(variant)
531 .ok_or_else(|| MapperError::VariantNotFound {
532 fv: fv.to_string(),
533 variant: variant.to_string(),
534 })?;
535
536 let tx_group = vc.tx_group(pid).ok_or_else(|| MapperError::PidNotFound {
537 fv: fv.to_string(),
538 variant: variant.to_string(),
539 pid: pid.to_string(),
540 })?;
541
542 let msg_engine = vc.msg_engine(pid);
543 let tx_engine = vc.tx_engine(pid).ok_or_else(|| MapperError::PidNotFound {
544 fv: fv.to_string(),
545 variant: variant.to_string(),
546 pid: pid.to_string(),
547 })?;
548
549 let filtered_mig = vc
550 .filtered_mig(pid)
551 .ok_or_else(|| MapperError::NoMigSchema {
552 fv: fv.to_string(),
553 variant: variant.to_string(),
554 })?;
555
556 let transaktionen: Vec<mig_bo4e::model::MappedTransaktion> =
558 tx_stammdaten.iter().map(split_transaktion).collect();
559 let mapped = mig_bo4e::model::MappedMessage {
560 nachricht_meta: serde_json::Value::Null,
561 stammdaten: msg_stammdaten.clone(),
562 transaktionen,
563 nesting_info: Default::default(),
564 inter_group_segments: Default::default(),
565 };
566
567 let tree = MappingEngine::map_interchange_reverse(
569 &msg_engine,
570 &tx_engine,
571 &mapped,
572 tx_group,
573 Some(&filtered_mig),
574 );
575
576 let disassembler = mig_assembly::disassembler::Disassembler::new(&filtered_mig);
581 let checked = match check {
582 EntrySegmentCheck::Refuse => disassembler.disassemble_checked(&tree),
583 EntrySegmentCheck::Render => Ok(disassembler.disassemble(&tree)),
584 };
585 let segments = checked.map_err(|e| match e {
586 mig_assembly::AssemblyError::MissingGroupEntrySegment {
587 group_path,
588 source_path,
589 entry_segment,
590 present_segments,
591 } => {
592 let (entities, entry_fields) = describe_entry_segment_mappings(
593 [msg_engine.definitions(), tx_engine.definitions()],
594 &source_path,
595 &entry_segment,
596 );
597 MapperError::MissingGroupEntrySegment(Box::new(
598 crate::error::GroupEntrySegmentError {
599 pid: pid.to_string(),
600 group_path,
601 source_path,
602 entry_segment,
603 present_segments,
604 entities,
605 entry_fields,
606 },
607 ))
608 }
609 other => MapperError::Assembly(other),
610 })?;
611
612 let delimiters = edifact_primitives::EdifactDelimiters::default();
614 Ok(mig_assembly::renderer::render_edifact(
615 &segments,
616 &delimiters,
617 ))
618 }
619
620 pub fn to_edifact_struct(
626 &self,
627 nachricht: &impl serde::Serialize,
628 fv: &str,
629 variant: &str,
630 pid: &str,
631 ) -> Result<String, MapperError> {
632 let json = serde_json::to_value(nachricht)
633 .map_err(|e| MapperError::Serialization(e.to_string()))?;
634
635 let msg_stammdaten = json
636 .get("stammdaten")
637 .cloned()
638 .unwrap_or(serde_json::Value::Object(Default::default()));
639
640 let tx_stammdaten: Vec<serde_json::Value> = json
641 .get("transaktionen")
642 .and_then(|v| v.as_array())
643 .cloned()
644 .unwrap_or_default();
645
646 self.to_edifact(&msg_stammdaten, &tx_stammdaten, fv, variant, pid)
647 }
648
649 pub fn from_edifact<M, T>(
677 &self,
678 edifact: &str,
679 fv: &str,
680 variant: &str,
681 pid: &str,
682 ) -> Result<mig_bo4e::model::Interchange<M, T>, MapperError>
683 where
684 M: serde::de::DeserializeOwned,
685 T: serde::de::DeserializeOwned,
686 {
687 let (interchange, diagnostics) =
688 self.from_edifact_with_diagnostics(edifact, fv, variant, pid)?;
689 for d in &diagnostics {
692 tracing::warn!(
693 fv,
694 variant,
695 pid,
696 kind = ?d.kind,
697 segment = %d.segment_id,
698 position = d.position,
699 "from_edifact: {}",
700 d.message
701 );
702 }
703 Ok(interchange)
704 }
705
706 pub fn from_edifact_with_diagnostics<M, T>(
719 &self,
720 edifact: &str,
721 fv: &str,
722 variant: &str,
723 pid: &str,
724 ) -> Result<
725 (
726 mig_bo4e::model::Interchange<M, T>,
727 Vec<mig_assembly::StructureDiagnostic>,
728 ),
729 MapperError,
730 >
731 where
732 M: serde::de::DeserializeOwned,
733 T: serde::de::DeserializeOwned,
734 {
735 self.ensure_bundle_loaded(fv)?;
736 let bundles = self.bundles.lock().unwrap();
737 let bundle = bundles.get(fv).unwrap();
738 let vc = bundle
739 .variant(variant)
740 .ok_or_else(|| MapperError::VariantNotFound {
741 fv: fv.to_string(),
742 variant: variant.to_string(),
743 })?;
744
745 let tx_group = vc.tx_group(pid).ok_or_else(|| MapperError::PidNotFound {
746 fv: fv.to_string(),
747 variant: variant.to_string(),
748 pid: pid.to_string(),
749 })?;
750
751 let msg_engine = vc.msg_engine(pid);
752 let tx_engine = vc.tx_engine(pid).ok_or_else(|| MapperError::PidNotFound {
753 fv: fv.to_string(),
754 variant: variant.to_string(),
755 pid: pid.to_string(),
756 })?;
757
758 let filtered_mig = vc
759 .filtered_mig(pid)
760 .ok_or_else(|| MapperError::NoMigSchema {
761 fv: fv.to_string(),
762 variant: variant.to_string(),
763 })?;
764
765 let svc = ConversionService::from_mig(filtered_mig);
771 let (chunks, trees, assembly_diagnostics) = svc
772 .convert_interchange_to_trees_with_diagnostics(
773 edifact,
774 mig_assembly::assembler::AssemblerConfig {
775 strict_code_matching: true,
776 skip_unknown_segments: true,
777 ..Default::default()
778 },
779 )?;
780
781 let tree = trees.first().ok_or_else(|| {
782 MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
783 "No messages in interchange".to_string(),
784 ))
785 })?;
786
787 let interchangedaten = mig_bo4e::model::extract_interchangedaten(&chunks.envelope);
789 let msg_chunk = chunks.messages.first().ok_or_else(|| {
790 MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
791 "No message chunks".to_string(),
792 ))
793 })?;
794 let (unh_ref, nachrichten_typ) = mig_bo4e::model::extract_unh_fields(&msg_chunk.unh);
795 let nachrichtendaten = mig_bo4e::model::Nachrichtendaten {
796 unh_referenz: unh_ref,
797 nachrichten_typ,
798 nachricht: Default::default(),
799 };
800
801 let interchange = MappingEngine::map_interchange_typed::<M, T>(
803 &msg_engine,
804 &tx_engine,
805 tree,
806 tx_group,
807 true,
808 nachrichtendaten,
809 interchangedaten,
810 )
811 .map_err(|e| MapperError::Serialization(e.to_string()))?;
812
813 Ok((interchange, assembly_diagnostics))
814 }
815
816 pub fn detect_pid(&self, edifact: &str) -> Result<String, MapperError> {
828 let segments = mig_assembly::tokenize::parse_to_segments(edifact.as_bytes())?;
829 let chunks = mig_assembly::split_messages(segments)?;
830 let msg_chunk = chunks.messages.first().ok_or_else(|| {
831 MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
832 "No messages found in EDIFACT content".to_string(),
833 ))
834 })?;
835 let msg_segments = msg_chunk.message_segments();
836 mig_assembly::pid_detect::detect_pid(&msg_segments).map_err(MapperError::Assembly)
837 }
838
839 pub fn validate_edifact(
855 &self,
856 edifact: &str,
857 fv: &str,
858 level: automapper_validation::ValidationLevel,
859 ) -> Result<automapper_validation::ValidationReport, MapperError> {
860 self.validate_edifact_inner(edifact, fv, None, level)
861 }
862
863 pub fn validate_edifact_for_pid(
873 &self,
874 edifact: &str,
875 fv: &str,
876 variant: &str,
877 pid: &str,
878 level: automapper_validation::ValidationLevel,
879 ) -> Result<automapper_validation::ValidationReport, MapperError> {
880 self.validate_edifact_inner(edifact, fv, Some((variant, pid)), level)
881 }
882
883 fn validate_edifact_inner(
884 &self,
885 edifact: &str,
886 fv: &str,
887 known: Option<(&str, &str)>,
888 level: automapper_validation::ValidationLevel,
889 ) -> Result<automapper_validation::ValidationReport, MapperError> {
890 self.ensure_bundle_loaded(fv)?;
891 let bundles = self.bundles.lock().unwrap();
892 let bundle = bundles.get(fv).unwrap();
893
894 let segments = mig_assembly::tokenize::parse_to_segments(edifact.as_bytes())?;
896 let chunks = mig_assembly::split_messages(segments)?;
897 let msg_chunk = chunks.messages.first().ok_or_else(|| {
898 MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
899 "No messages found in EDIFACT content".to_string(),
900 ))
901 })?;
902
903 let (pid, variant, vc) = match known {
909 Some((variant, pid)) => {
910 let vc = bundle
911 .variant(variant)
912 .ok_or_else(|| MapperError::VariantNotFound {
913 fv: fv.to_string(),
914 variant: variant.to_string(),
915 })?;
916 (pid.to_string(), variant.to_string(), vc)
917 }
918 None => {
919 let pid = mig_assembly::pid_detect::detect_pid(&msg_chunk.message_segments())
920 .map_err(MapperError::Assembly)?;
921 let pid_key = format!("pid_{pid}");
922 let (variant, vc) = bundle
923 .variants
924 .iter()
925 .find(|(_, vc)| vc.pid_ahb_workflows.contains_key(&pid_key))
926 .ok_or_else(|| MapperError::PidNotFound {
927 fv: fv.to_string(),
928 variant: "?".to_string(),
929 pid: pid.clone(),
930 })?;
931 (pid, variant.clone(), vc)
932 }
933 };
934 let pid_key = format!("pid_{pid}");
935
936 let workflow =
937 vc.pid_ahb_workflows
938 .get(&pid_key)
939 .ok_or_else(|| MapperError::PidNotFound {
940 fv: fv.to_string(),
941 variant: variant.clone(),
942 pid: pid.clone(),
943 })?;
944 let filtered_mig = vc
945 .filtered_mig(&pid)
946 .ok_or_else(|| MapperError::NoMigSchema {
947 fv: fv.to_string(),
948 variant: variant.clone(),
949 })?;
950
951 let mut all_segments = msg_chunk.segments_for_mig(&filtered_mig);
954 if filtered_mig.segments.iter().any(|s| s.id == "UNZ") {
955 if let Some(unz) = &chunks.unz {
956 all_segments.push(unz.clone());
957 }
958 }
959
960 let evaluator: std::sync::Arc<dyn automapper_validation::ConditionEvaluator> =
964 match crate::evaluator_factory::create_evaluator(&variant, fv) {
965 Some(boxed) => std::sync::Arc::from(boxed),
966 None => std::sync::Arc::new(
967 automapper_validation::UtilmdStromConditionEvaluatorFV2504::default(),
968 ),
969 };
970 let external = automapper_validation::eval::NoOpExternalProvider;
971
972 let mut report = automapper_validation::validate_edifact_message(
973 &all_segments,
974 &filtered_mig,
975 workflow,
976 evaluator,
977 &external,
978 level,
979 );
980
981 if let (Some(mig), Some(defs)) = (vc.mig_schema.as_ref(), vc.combined_defs.get(&pid_key)) {
987 let reverse = mig_bo4e::path_resolver::ReversePathResolver::from_mig(mig);
988 let field_index =
989 mig_bo4e::Bo4eFieldIndex::build_with_resolver(defs, &filtered_mig, &reverse);
990 report.enrich_bo4e_paths(|path, hint| field_index.resolve(path, hint));
991 }
992
993 Ok(report)
994 }
995
996 pub fn validate_bo4e(
1019 &self,
1020 msg_stammdaten: &serde_json::Value,
1021 tx_stammdaten: &[serde_json::Value],
1022 fv: &str,
1023 variant: &str,
1024 pid: &str,
1025 envelope: Option<&InterchangeEnvelope>,
1026 level: automapper_validation::ValidationLevel,
1027 ) -> Result<automapper_validation::ValidationReport, MapperError> {
1028 let placeholder;
1029 let envelope = match envelope {
1030 Some(e) => e,
1031 None => {
1032 placeholder = InterchangeEnvelope {
1033 sender: EdifactParty::bdew("9900000000001"),
1034 receiver: EdifactParty::bdew("9900000000002"),
1035 interchange_ref: "1".to_string(),
1036 };
1037 &placeholder
1038 }
1039 };
1040
1041 let edifact = self.render_interchange(
1046 envelope,
1047 &[InterchangeMessage {
1048 message_ref: "1".to_string(),
1049 msg_stammdaten: msg_stammdaten.clone(),
1050 tx_stammdaten: tx_stammdaten.to_vec(),
1051 fv: fv.to_string(),
1052 variant: variant.to_string(),
1053 pid: pid.to_string(),
1054 }],
1055 EntrySegmentCheck::Render,
1056 &EnvelopeOptions::default(),
1057 )?;
1058
1059 self.validate_edifact_for_pid(&edifact, fv, variant, pid, level)
1062 }
1063
1064 pub fn association_code(&self, fv: &str, variant: &str) -> Result<String, MapperError> {
1075 let meta = self.message_metadata(fv, variant)?;
1076 Ok(meta.association_code)
1077 }
1078
1079 pub fn message_metadata(
1084 &self,
1085 fv: &str,
1086 variant: &str,
1087 ) -> Result<MessageMetadata, MapperError> {
1088 self.ensure_bundle_loaded(fv)?;
1089 let bundles = self.bundles.lock().unwrap();
1090 let bundle = bundles.get(fv).unwrap();
1091 let vc = bundle
1092 .variant(variant)
1093 .ok_or_else(|| MapperError::VariantNotFound {
1094 fv: fv.to_string(),
1095 variant: variant.to_string(),
1096 })?;
1097 let mig = vc
1098 .mig_schema
1099 .as_ref()
1100 .ok_or_else(|| MapperError::NoMigSchema {
1101 fv: fv.to_string(),
1102 variant: variant.to_string(),
1103 })?;
1104 Ok(MessageMetadata {
1105 message_type: mig.message_type.clone(),
1106 release: release_code_for_message_type(&mig.message_type),
1107 association_code: mig.version.clone(),
1108 })
1109 }
1110
1111 pub fn to_edifact_interchange(
1163 &self,
1164 envelope: &InterchangeEnvelope,
1165 messages: &[InterchangeMessage],
1166 ) -> Result<String, MapperError> {
1167 self.render_interchange(
1168 envelope,
1169 messages,
1170 EntrySegmentCheck::Refuse,
1171 &EnvelopeOptions::default(),
1172 )
1173 }
1174
1175 pub fn to_edifact_interchange_with(
1187 &self,
1188 envelope: &InterchangeEnvelope,
1189 messages: &[InterchangeMessage],
1190 options: &EnvelopeOptions,
1191 ) -> Result<String, MapperError> {
1192 self.render_interchange(envelope, messages, EntrySegmentCheck::Refuse, options)
1193 }
1194
1195 fn render_interchange(
1196 &self,
1197 envelope: &InterchangeEnvelope,
1198 messages: &[InterchangeMessage],
1199 check: EntrySegmentCheck,
1200 options: &EnvelopeOptions,
1201 ) -> Result<String, MapperError> {
1202 let delimiters = edifact_primitives::EdifactDelimiters::default();
1203 let sep = delimiters.component as char;
1204 let elem = delimiters.element as char;
1205 let seg_term = delimiters.segment as char;
1206
1207 let mut output = String::new();
1208
1209 if options.emit_una {
1212 output.push_str(&format!(
1213 "UNA{}{}{}{}{}{}",
1214 sep, elem, delimiters.decimal as char, delimiters.release as char, ' ', seg_term, ));
1221 }
1222
1223 let now = chrono::Utc::now();
1226 let date_str = options
1227 .datum
1228 .clone()
1229 .unwrap_or_else(|| now.format("%y%m%d").to_string());
1230 let time_str = options
1231 .zeit
1232 .clone()
1233 .unwrap_or_else(|| now.format("%H%M").to_string());
1234 let sender = &envelope.sender;
1235 let receiver = &envelope.receiver;
1236 let interchange_ref = &envelope.interchange_ref;
1237 output.push_str(&format!(
1238 "UNB{elem}UNOC{sep}3{elem}{sid}{sep}{sq}{elem}{rid}{sep}{rq}{elem}{date_str}{sep}{time_str}{elem}{interchange_ref}{seg_term}",
1239 sid = sender.id,
1240 sq = sender.qualifier,
1241 rid = receiver.id,
1242 rq = receiver.qualifier,
1243 ));
1244
1245 let mut message_count = 0u32;
1246
1247 for msg in messages {
1248 let meta = self.message_metadata(&msg.fv, &msg.variant)?;
1249
1250 let body = self.render_message_body(
1252 &msg.msg_stammdaten,
1253 &msg.tx_stammdaten,
1254 &msg.fv,
1255 &msg.variant,
1256 &msg.pid,
1257 check,
1258 )?;
1259
1260 let body_seg_count = body
1262 .split(seg_term)
1263 .filter(|s: &&str| !s.is_empty())
1264 .count();
1265 let segment_count = body_seg_count + 2;
1267
1268 output.push_str(&format!(
1270 "UNH{elem}{ref}{elem}{msg_type}{sep}D{sep}{release}{sep}UN{sep}{assoc}{seg_term}",
1271 ref = msg.message_ref,
1272 msg_type = meta.message_type,
1273 release = meta.release,
1274 assoc = meta.association_code,
1275 ));
1276
1277 output.push_str(&body);
1279
1280 output.push_str(&format!(
1282 "UNT{elem}{segment_count}{elem}{ref}{seg_term}",
1283 ref = msg.message_ref,
1284 ));
1285
1286 message_count += 1;
1287 }
1288
1289 output.push_str(&format!(
1291 "UNZ{elem}{message_count}{elem}{interchange_ref}{seg_term}",
1292 ));
1293
1294 Ok(output)
1295 }
1296
1297 pub fn loaded_format_versions(&self) -> Vec<String> {
1299 self.bundles.lock().unwrap().keys().cloned().collect()
1300 }
1301
1302 pub fn variants(&self, fv: &str) -> Result<Vec<String>, MapperError> {
1306 self.ensure_bundle_loaded(fv)?;
1307 let bundles = self.bundles.lock().unwrap();
1308 let bundle = bundles.get(fv).unwrap();
1309 Ok(bundle.variants.keys().cloned().collect())
1310 }
1311}
1312
1313#[derive(Debug, Clone)]
1315pub struct MessageMetadata {
1316 pub message_type: String,
1318 pub release: String,
1320 pub association_code: String,
1322}
1323
1324#[derive(Debug, Clone)]
1326pub struct InterchangeEnvelope {
1327 pub sender: EdifactParty,
1329 pub receiver: EdifactParty,
1331 pub interchange_ref: String,
1333}
1334
1335#[derive(Debug, Clone)]
1362pub struct EnvelopeOptions {
1363 emit_una: bool,
1364 datum: Option<String>,
1365 zeit: Option<String>,
1366}
1367
1368impl Default for EnvelopeOptions {
1369 fn default() -> Self {
1370 Self {
1371 emit_una: true,
1372 datum: None,
1373 zeit: None,
1374 }
1375 }
1376}
1377
1378impl EnvelopeOptions {
1379 pub fn emit_una(mut self, emit: bool) -> Self {
1383 self.emit_una = emit;
1384 self
1385 }
1386
1387 pub fn datum_zeit(mut self, datum: impl Into<String>, zeit: impl Into<String>) -> Self {
1390 self.datum = Some(datum.into());
1391 self.zeit = Some(zeit.into());
1392 self
1393 }
1394
1395 pub fn datum_zeit_from(mut self, daten: &mig_bo4e::model::Interchangedaten) -> Self {
1398 self.datum = daten.datum.clone();
1399 self.zeit = daten.zeit.clone();
1400 self
1401 }
1402}
1403
1404#[derive(Debug, Clone)]
1406pub struct EdifactParty {
1407 pub id: String,
1409 pub qualifier: String,
1411}
1412
1413impl EdifactParty {
1414 pub fn bdew(id: &str) -> Self {
1416 Self {
1417 id: id.to_string(),
1418 qualifier: "500".to_string(),
1419 }
1420 }
1421
1422 pub fn gs1(id: &str) -> Self {
1424 Self {
1425 id: id.to_string(),
1426 qualifier: "14".to_string(),
1427 }
1428 }
1429}
1430
1431#[derive(Debug, Clone)]
1434pub struct InterchangeMessage {
1435 pub message_ref: String,
1437 pub msg_stammdaten: serde_json::Value,
1439 pub tx_stammdaten: Vec<serde_json::Value>,
1441 pub fv: String,
1443 pub variant: String,
1445 pub pid: String,
1447}
1448
1449#[derive(Debug, Clone, Copy)]
1451enum EntrySegmentCheck {
1452 Refuse,
1454 Render,
1456}
1457
1458fn describe_entry_segment_mappings<'d>(
1468 definition_sets: impl IntoIterator<Item = &'d [mig_bo4e::definition::MappingDefinition]>,
1469 source_path: &str,
1470 entry_segment: &str,
1471) -> (Vec<String>, Vec<String>) {
1472 fn qualifies(unqualified: &str, qualified: &str) -> bool {
1473 !unqualified.contains('_')
1474 && qualified.len() > unqualified.len()
1475 && qualified.is_char_boundary(unqualified.len())
1476 && qualified[..unqualified.len()].eq_ignore_ascii_case(unqualified)
1477 && qualified.as_bytes()[unqualified.len()] == b'_'
1478 }
1479 fn part_matches(mig_part: &str, def_part: &str) -> bool {
1480 def_part.eq_ignore_ascii_case(mig_part)
1481 || qualifies(mig_part, def_part)
1482 || qualifies(def_part, mig_part)
1483 }
1484 let mig_parts: Vec<&str> = source_path.split('.').collect();
1485
1486 let mut entities: Vec<String> = Vec::new();
1487 let mut entry_fields: Vec<String> = Vec::new();
1488 for def in definition_sets.into_iter().flatten() {
1489 let Some(def_path) = def.meta.source_path.as_deref() else {
1490 continue;
1491 };
1492 let def_parts: Vec<&str> = def_path.split('.').collect();
1493 if def_parts.len() != mig_parts.len()
1494 || !mig_parts
1495 .iter()
1496 .zip(&def_parts)
1497 .all(|(m, d)| part_matches(m, d))
1498 {
1499 continue;
1500 }
1501 if !entities.contains(&def.meta.entity) {
1502 entities.push(def.meta.entity.clone());
1503 }
1504 for (path, mapping) in &def.fields {
1505 let tag = path
1506 .split(['.', '['])
1507 .next()
1508 .unwrap_or_default()
1509 .to_ascii_uppercase();
1510 let target = match mapping {
1511 mig_bo4e::definition::FieldMapping::Simple(t) => t.as_str(),
1512 mig_bo4e::definition::FieldMapping::Structured(f) => f.target.as_str(),
1513 mig_bo4e::definition::FieldMapping::Nested(_) => continue,
1514 };
1515 if tag == entry_segment && !target.is_empty() {
1516 let field = format!("{}.{}", def.meta.entity, target);
1517 if !entry_fields.contains(&field) {
1518 entry_fields.push(field);
1519 }
1520 }
1521 }
1522 }
1523 (entities, entry_fields)
1524}
1525
1526fn release_code_for_message_type(msg_type: &str) -> String {
1530 mig_bo4e::model::release_code_for_message_type(msg_type).to_string()
1531}
1532
1533#[cfg(test)]
1534mod tests {
1535 use super::*;
1536 use std::path::Path;
1537
1538 fn data_dir() -> Option<std::path::PathBuf> {
1539 let dist = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../dist");
1541 if dist.join("edifact-data-FV2504.bin").exists() {
1542 return Some(dist);
1543 }
1544 let cache = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../cache/mappings");
1545 if cache.join("FV2504").exists() {
1546 return Some(cache);
1547 }
1548 eprintln!("Skipping test: no DataBundle files found");
1549 None
1550 }
1551
1552 #[test]
1553 fn test_to_edifact_produces_edifact_output() {
1554 let Some(data_dir) = data_dir() else {
1555 return;
1556 };
1557 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1558
1559 let msg_stammdaten = serde_json::json!({
1560 "marktteilnehmer": [{
1561 "marktrolle": "MS",
1562 "rollencodenummer": "9900123456789",
1563 "codepflegeCode": "293"
1564 }]
1565 });
1566 let tx_stammdaten = serde_json::json!({
1567 "prozessdaten": {
1568 "pruefidentifikator": "55001",
1569 "vorgangId": "ABC123",
1570 "transaktionsgrund": "E01"
1571 }
1572 });
1573
1574 let result = mapper.to_edifact(
1575 &msg_stammdaten,
1576 &[tx_stammdaten],
1577 "FV2504",
1578 "UTILMD_Strom",
1579 "55001",
1580 );
1581 assert!(result.is_ok(), "to_edifact failed: {:?}", result.err());
1582 let edifact = result.unwrap();
1583 assert!(!edifact.is_empty(), "EDIFACT output should not be empty");
1584 assert!(edifact.contains("NAD"), "Should contain NAD segment");
1586 assert!(edifact.contains("IDE"), "Should contain IDE segment");
1588 }
1589
1590 #[test]
1591 fn test_to_edifact_struct_produces_edifact_output() {
1592 let Some(data_dir) = data_dir() else {
1593 return;
1594 };
1595 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1596
1597 let nachricht = serde_json::json!({
1598 "stammdaten": {
1599 "marktteilnehmer": [{
1600 "marktrolle": "MS",
1601 "rollencodenummer": "9900123456789",
1602 "codepflegeCode": "293"
1603 }]
1604 },
1605 "transaktionen": [{
1606 "prozessdaten": {
1607 "pruefidentifikator": "55001",
1608 "vorgangId": "ABC123"
1609 }
1610 }]
1611 });
1612
1613 let result = mapper.to_edifact_struct(&nachricht, "FV2504", "UTILMD_Strom", "55001");
1614 assert!(
1615 result.is_ok(),
1616 "to_edifact_struct failed: {:?}",
1617 result.err()
1618 );
1619 let edifact = result.unwrap();
1620 assert!(!edifact.is_empty(), "EDIFACT output should not be empty");
1621 }
1622
1623 #[test]
1624 fn test_to_edifact_invalid_fv_returns_error() {
1625 let Some(data_dir) = data_dir() else {
1626 return;
1627 };
1628 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1629
1630 let result = mapper.to_edifact(
1631 &serde_json::json!({}),
1632 &[serde_json::json!({})],
1633 "FV9999",
1634 "UTILMD_Strom",
1635 "55001",
1636 );
1637 assert!(result.is_err());
1638 }
1639
1640 #[test]
1641 fn test_to_edifact_invalid_variant_returns_error() {
1642 let Some(data_dir) = data_dir() else {
1643 return;
1644 };
1645 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1646
1647 let result = mapper.to_edifact(
1648 &serde_json::json!({}),
1649 &[serde_json::json!({})],
1650 "FV2504",
1651 "NONEXISTENT",
1652 "55001",
1653 );
1654 assert!(result.is_err());
1655 }
1656
1657 #[test]
1658 fn test_to_edifact_invalid_pid_returns_error() {
1659 let Some(data_dir) = data_dir() else {
1660 return;
1661 };
1662 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1663
1664 let result = mapper.to_edifact(
1665 &serde_json::json!({}),
1666 &[serde_json::json!({})],
1667 "FV2504",
1668 "UTILMD_Strom",
1669 "99999",
1670 );
1671 assert!(result.is_err());
1672 }
1673
1674 #[test]
1675 fn test_association_code() {
1676 let Some(data_dir) = data_dir() else {
1677 return;
1678 };
1679 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1680
1681 let code = mapper.association_code("FV2504", "UTILMD_Strom").unwrap();
1682 assert_eq!(code, "S2.1");
1683
1684 let code = mapper.association_code("FV2504", "MSCONS").unwrap();
1685 assert_eq!(code, "2.4c");
1686 }
1687
1688 #[test]
1689 fn test_message_metadata() {
1690 let Some(data_dir) = data_dir() else {
1691 return;
1692 };
1693 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1694
1695 let meta = mapper.message_metadata("FV2504", "UTILMD_Strom").unwrap();
1696 assert_eq!(meta.message_type, "UTILMD");
1697 assert_eq!(meta.release, "11A");
1698 assert_eq!(meta.association_code, "S2.1");
1699 }
1700
1701 #[test]
1702 fn test_to_edifact_interchange() {
1703 let Some(data_dir) = data_dir() else {
1704 return;
1705 };
1706 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1707
1708 let result = mapper.to_edifact_interchange(
1709 &InterchangeEnvelope {
1710 sender: EdifactParty::bdew("9900000000003"),
1711 receiver: EdifactParty::bdew("9900000000001"),
1712 interchange_ref: "REF001".to_string(),
1713 },
1714 &[InterchangeMessage {
1715 message_ref: "MSG001".to_string(),
1716 msg_stammdaten: serde_json::json!({
1717 "marktteilnehmer": [{
1718 "marktrolle": "MS",
1719 "rollencodenummer": "9900123456789",
1720 "codepflegeCode": "293"
1721 }]
1722 }),
1723 tx_stammdaten: vec![serde_json::json!({
1724 "prozessdaten": {
1725 "pruefidentifikator": "55001",
1726 "vorgangId": "ABC123",
1727 "transaktionsgrund": "E01"
1728 }
1729 })],
1730 fv: "FV2504".to_string(),
1731 variant: "UTILMD_Strom".to_string(),
1732 pid: "55001".to_string(),
1733 }],
1734 );
1735 assert!(
1736 result.is_ok(),
1737 "to_edifact_interchange failed: {:?}",
1738 result.err()
1739 );
1740 let edifact = result.unwrap();
1741
1742 assert!(edifact.starts_with("UNA:+.? '"), "Should start with UNA");
1744 assert!(
1745 edifact.contains("UNB+UNOC:3+9900000000003:500+9900000000001:500+"),
1746 "Should contain UNB with sender/receiver"
1747 );
1748 assert!(
1749 edifact.contains("UNH+MSG001+UTILMD:D:11A:UN:S2.1'"),
1750 "Should contain UNH with correct S009"
1751 );
1752 assert!(edifact.contains("NAD"), "Should contain body NAD segment");
1753 assert!(edifact.contains("UNT+"), "Should contain UNT");
1754 assert!(
1755 edifact.contains("+MSG001'"),
1756 "UNT should reference message ref"
1757 );
1758 assert!(
1759 edifact.contains("UNZ+1+REF001'"),
1760 "Should contain UNZ with count and ref"
1761 );
1762 }
1763
1764 #[test]
1765 fn test_detect_pid_from_rff_z13() {
1766 let Some(data_dir) = data_dir() else {
1767 return;
1768 };
1769 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1770
1771 let edifact = "\
1772 UNB+UNOC:3+9978842000002:500+9900269000000:500+250331:1329+REF001'\
1773 UNH+MSG001+UTILMD:D:11A:UN:S2.1'\
1774 BGM+E01+DOC001'\
1775 DTM+137:202503311329?+00:303'\
1776 NAD+MS+9978842000002::293'\
1777 NAD+MR+9900269000000::293'\
1778 IDE+24+TX001'\
1779 DTM+92:202505312200?+00:303'\
1780 DTM+93:202512312300?+00:303'\
1781 STS+7++E01+ZW4+E03'\
1782 LOC+Z16+12345678900'\
1783 RFF+Z13:55001'\
1784 UNT+12+MSG001'\
1785 UNZ+1+REF001'";
1786
1787 let pid = mapper.detect_pid(edifact).unwrap();
1788 assert_eq!(pid, "55001");
1789 }
1790
1791 #[test]
1792 fn test_detect_pid_no_messages_returns_error() {
1793 let Some(data_dir) = data_dir() else {
1794 return;
1795 };
1796 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1797
1798 let edifact = "UNB+UNOC:3+SENDER:500+RECEIVER:500+250401:1200+REF'\
1799 UNZ+0+REF'";
1800 assert!(mapper.detect_pid(edifact).is_err());
1801 }
1802
1803 #[test]
1804 fn test_list_pids_returns_entries() {
1805 let Some(data_dir) = data_dir() else {
1806 return;
1807 };
1808 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir)).unwrap();
1809 let pids = mapper.list_pids().expect("list_pids should succeed");
1810 assert!(!pids.is_empty(), "should return at least one PID");
1811 assert!(
1812 pids.iter().any(|p| p.pid == "55001"),
1813 "should include PID 55001"
1814 );
1815 assert!(
1816 pids.iter().any(|p| p.fv == "FV2504"),
1817 "should include FV2504"
1818 );
1819 assert!(
1820 pids.iter().any(|p| p.variant == "UTILMD_Strom"),
1821 "should include UTILMD_Strom"
1822 );
1823 }
1824
1825 #[test]
1826 fn test_pid_requirements_returns_requirements() {
1827 let Some(data_dir) = data_dir() else {
1828 return;
1829 };
1830 let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
1831
1832 let req = mapper
1833 .pid_requirements("FV2504", "UTILMD_Strom", "55001")
1834 .expect("pid_requirements should succeed");
1835
1836 assert_eq!(req.pid, "55001");
1837 assert!(
1838 !req.entities.is_empty(),
1839 "55001 should have at least one entity"
1840 );
1841 assert!(
1842 req.entities.iter().any(|e| e.entity == "Prozessdaten"),
1843 "55001 should have a Prozessdaten entity"
1844 );
1845 }
1846}