use std::collections::HashMap;
use std::sync::Mutex;
use mig_assembly::ConversionService;
use mig_bo4e::engine::DataBundle;
use mig_bo4e::MappingEngine;
use crate::data_dir::DataDir;
use crate::error::MapperError;
pub struct Bo4eResult {
pub pid: String,
pub message_type: String,
pub variant: String,
pub bo4e: serde_json::Value,
}
#[derive(Debug, Clone)]
pub struct PidListEntry {
pub fv: String,
pub variant: String,
pub pid: String,
pub beschreibung: String,
}
pub struct Mapper {
data_dir: DataDir,
bundles: Mutex<HashMap<String, DataBundle>>,
}
fn split_transaktion(tx: &serde_json::Value) -> mig_bo4e::model::MappedTransaktion {
let (transaktionsdaten, stammdaten) = if mig_bo4e::model::is_wrapped_transaktion(tx) {
(
tx.get("transaktionsdaten")
.cloned()
.unwrap_or(serde_json::Value::Null),
tx.get("stammdaten")
.cloned()
.unwrap_or_else(|| serde_json::Value::Object(Default::default())),
)
} else {
(serde_json::Value::Null, tx.clone())
};
mig_bo4e::model::MappedTransaktion {
transaktionsdaten,
stammdaten,
nesting_info: Default::default(),
}
}
impl Mapper {
pub fn from_data_dir(data_dir: DataDir) -> Result<Self, MapperError> {
let mapper = Self {
data_dir,
bundles: Mutex::new(HashMap::new()),
};
let eager_fvs: Vec<String> = mapper.data_dir.eager_fvs().to_vec();
for fv in &eager_fvs {
mapper.ensure_bundle_loaded(fv)?;
}
Ok(mapper)
}
fn ensure_bundle_loaded(&self, fv: &str) -> Result<(), MapperError> {
let mut bundles = self.bundles.lock().unwrap();
if bundles.contains_key(fv) {
return Ok(());
}
let path = self.data_dir.bundle_path(fv);
if !path.exists() {
return Err(MapperError::BundleNotFound { fv: fv.to_string() });
}
let bundle = DataBundle::load(&path)?;
let expected = DataBundle::PRODUCING_VERSION;
if !self.data_dir.allows_bundle_from_other_release()
&& bundle.built_by.as_deref() != Some(expected)
{
return Err(MapperError::BundleFromOtherRelease {
fv: fv.to_string(),
built_by: bundle.built_by.clone(),
expected: expected.to_string(),
path: path.display().to_string(),
});
}
bundles.insert(fv.to_string(), bundle);
Ok(())
}
pub fn conversion_service(
&self,
fv: &str,
variant: &str,
) -> Result<ConversionService, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let mig = vc
.mig_schema
.as_ref()
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: format!("{variant} (no MIG schema in bundle)"),
})?;
Ok(ConversionService::from_mig(mig.clone()))
}
pub fn engine(&self, fv: &str, variant: &str, pid: &str) -> Result<MappingEngine, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let pid_key = format!("pid_{pid}");
let defs = vc
.combined_defs
.get(&pid_key)
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
Ok(MappingEngine::from_definitions_with_code_lists(
std::sync::Arc::clone(&vc.code_lists),
defs.clone(),
))
}
pub fn pid_requirements(
&self,
fv: &str,
variant: &str,
pid: &str,
) -> Result<mig_bo4e::pid_requirements::PidRequirements, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let pid_key = format!("pid_{pid}");
vc.pid_requirements
.get(&pid_key)
.cloned()
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})
}
pub fn bo4e_catalog(
&self,
fv: &str,
) -> Result<mig_bo4e::bo4e_catalog::Bo4eCatalog, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
Ok(bundle.bo4e_catalog.clone())
}
pub fn list_pids(&self) -> Result<Vec<PidListEntry>, MapperError> {
let dir = self.data_dir.data_path();
let read_dir = std::fs::read_dir(dir).map_err(|_| MapperError::DataDirNotFound {
path: dir.display().to_string(),
})?;
let mut result = Vec::new();
for entry in read_dir.flatten() {
let path = entry.path();
if path.extension().is_some_and(|e| e == "bin") {
let stem = path
.file_stem()
.and_then(|s| s.to_str())
.unwrap_or("")
.to_string();
let fv = match stem.strip_prefix("edifact-data-") {
Some(v) => v.to_string(),
None => continue,
};
self.ensure_bundle_loaded(&fv)?;
let bundles = self.bundles.lock().unwrap();
if let Some(bundle) = bundles.get(&fv) {
for (variant, vc) in &bundle.variants {
for (pid_key, req) in &vc.pid_requirements {
let pid = pid_key.strip_prefix("pid_").unwrap_or(pid_key).to_string();
result.push(PidListEntry {
fv: fv.clone(),
variant: variant.clone(),
pid,
beschreibung: req.beschreibung.clone(),
});
}
}
}
}
}
result.sort_by(|a, b| a.pid.cmp(&b.pid));
Ok(result)
}
pub fn validate_pid(
&self,
json: &serde_json::Value,
fv: &str,
variant: &str,
pid: &str,
) -> Result<Vec<mig_bo4e::PidValidationError>, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let pid_key = format!("pid_{pid}");
let requirements =
vc.pid_requirements
.get(&pid_key)
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
Ok(mig_bo4e::pid_validation::validate_pid_json(
json,
requirements,
))
}
pub fn validate_pid_struct(
&self,
value: &impl serde::Serialize,
fv: &str,
variant: &str,
pid: &str,
) -> Result<Vec<mig_bo4e::PidValidationError>, MapperError> {
let json = serde_json::to_value(value).map_err(|e| {
MapperError::Mapping(mig_bo4e::MappingError::TypeConversion(e.to_string()))
})?;
self.validate_pid(&json, fv, variant, pid)
}
pub fn validate_pid_with_conditions(
&self,
json: &serde_json::Value,
fv: &str,
variant: &str,
pid: &str,
) -> Result<Vec<mig_bo4e::PidValidationError>, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let pid_key = format!("pid_{pid}");
let requirements =
vc.pid_requirements
.get(&pid_key)
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
let evaluator = crate::evaluator_factory::create_evaluator(variant, fv);
if let Some(evaluator) = evaluator {
let defs = vc
.combined_defs
.get(&pid_key)
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
let engine = MappingEngine::from_definitions_with_code_lists(
std::sync::Arc::clone(&vc.code_lists),
defs.clone(),
);
let tree = engine.map_all_reverse(json, None);
let segments = crate::tree_to_segments::tree_to_owned_segments(&tree);
Ok(crate::evaluator_factory::validate_with_boxed_evaluator(
evaluator.as_ref(),
json,
requirements,
pid,
&segments,
))
} else {
Ok(mig_bo4e::pid_validation::validate_pid_json_transaction(
json,
requirements,
))
}
}
pub fn to_edifact(
&self,
msg_stammdaten: &serde_json::Value,
tx_stammdaten: &[serde_json::Value],
fv: &str,
variant: &str,
pid: &str,
) -> Result<String, MapperError> {
self.render_message_body(
msg_stammdaten,
tx_stammdaten,
fv,
variant,
pid,
EntrySegmentCheck::Refuse,
)
}
pub fn to_edifact_nachricht(
&self,
nachricht: &mig_bo4e::model::Nachricht<serde_json::Value, serde_json::Value>,
fv: &str,
variant: &str,
pid: &str,
) -> Result<String, MapperError> {
let mut msg_stammdaten = nachricht.stammdaten.clone();
mig_bo4e::model::restore_message_metadata(&mut msg_stammdaten, &nachricht.nachrichtendaten);
self.to_edifact(&msg_stammdaten, &nachricht.transaktionen, fv, variant, pid)
}
fn render_message_body(
&self,
msg_stammdaten: &serde_json::Value,
tx_stammdaten: &[serde_json::Value],
fv: &str,
variant: &str,
pid: &str,
check: EntrySegmentCheck,
) -> Result<String, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let tx_group = vc.tx_group(pid).ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
let msg_engine = vc.msg_engine(pid);
let tx_engine = vc.tx_engine(pid).ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
let filtered_mig = vc
.filtered_mig(pid)
.ok_or_else(|| MapperError::NoMigSchema {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let transaktionen: Vec<mig_bo4e::model::MappedTransaktion> =
tx_stammdaten.iter().map(split_transaktion).collect();
let mapped = mig_bo4e::model::MappedMessage {
nachricht_meta: serde_json::Value::Null,
stammdaten: msg_stammdaten.clone(),
transaktionen,
nesting_info: Default::default(),
inter_group_segments: Default::default(),
};
let tree = MappingEngine::map_interchange_reverse(
&msg_engine,
&tx_engine,
&mapped,
tx_group,
Some(&filtered_mig),
);
let disassembler = mig_assembly::disassembler::Disassembler::new(&filtered_mig);
let checked = match check {
EntrySegmentCheck::Refuse => disassembler.disassemble_checked(&tree),
EntrySegmentCheck::Render => Ok(disassembler.disassemble(&tree)),
};
let segments = checked.map_err(|e| match e {
mig_assembly::AssemblyError::MissingGroupEntrySegment {
group_path,
source_path,
entry_segment,
present_segments,
} => {
let (entities, entry_fields) = describe_entry_segment_mappings(
[msg_engine.definitions(), tx_engine.definitions()],
&source_path,
&entry_segment,
);
MapperError::MissingGroupEntrySegment(Box::new(
crate::error::GroupEntrySegmentError {
pid: pid.to_string(),
group_path,
source_path,
entry_segment,
present_segments,
entities,
entry_fields,
},
))
}
other => MapperError::Assembly(other),
})?;
let delimiters = edifact_primitives::EdifactDelimiters::default();
Ok(mig_assembly::renderer::render_edifact(
&segments,
&delimiters,
))
}
pub fn to_edifact_struct(
&self,
nachricht: &impl serde::Serialize,
fv: &str,
variant: &str,
pid: &str,
) -> Result<String, MapperError> {
let json = serde_json::to_value(nachricht)
.map_err(|e| MapperError::Serialization(e.to_string()))?;
let msg_stammdaten = json
.get("stammdaten")
.cloned()
.unwrap_or(serde_json::Value::Object(Default::default()));
let tx_stammdaten: Vec<serde_json::Value> = json
.get("transaktionen")
.and_then(|v| v.as_array())
.cloned()
.unwrap_or_default();
self.to_edifact(&msg_stammdaten, &tx_stammdaten, fv, variant, pid)
}
pub fn from_edifact<M, T>(
&self,
edifact: &str,
fv: &str,
variant: &str,
pid: &str,
) -> Result<mig_bo4e::model::Interchange<M, T>, MapperError>
where
M: serde::de::DeserializeOwned,
T: serde::de::DeserializeOwned,
{
let (interchange, diagnostics) =
self.from_edifact_with_diagnostics(edifact, fv, variant, pid)?;
for d in &diagnostics {
tracing::warn!(
fv,
variant,
pid,
kind = ?d.kind,
segment = %d.segment_id,
position = d.position,
"from_edifact: {}",
d.message
);
}
Ok(interchange)
}
pub fn from_edifact_with_diagnostics<M, T>(
&self,
edifact: &str,
fv: &str,
variant: &str,
pid: &str,
) -> Result<
(
mig_bo4e::model::Interchange<M, T>,
Vec<mig_assembly::StructureDiagnostic>,
),
MapperError,
>
where
M: serde::de::DeserializeOwned,
T: serde::de::DeserializeOwned,
{
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let tx_group = vc.tx_group(pid).ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
let msg_engine = vc.msg_engine(pid);
let tx_engine = vc.tx_engine(pid).ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
})?;
let filtered_mig = vc
.filtered_mig(pid)
.ok_or_else(|| MapperError::NoMigSchema {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let svc = ConversionService::from_mig(filtered_mig);
let (chunks, trees, assembly_diagnostics) = svc
.convert_interchange_to_trees_with_diagnostics(
edifact,
mig_assembly::assembler::AssemblerConfig {
strict_code_matching: true,
skip_unknown_segments: true,
..Default::default()
},
)?;
let tree = trees.first().ok_or_else(|| {
MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
"No messages in interchange".to_string(),
))
})?;
let interchangedaten = mig_bo4e::model::extract_interchangedaten(&chunks.envelope);
let msg_chunk = chunks.messages.first().ok_or_else(|| {
MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
"No message chunks".to_string(),
))
})?;
let (unh_ref, nachrichten_typ) = mig_bo4e::model::extract_unh_fields(&msg_chunk.unh);
let nachrichtendaten = mig_bo4e::model::Nachrichtendaten {
unh_referenz: unh_ref,
nachrichten_typ,
nachricht: Default::default(),
};
let interchange = MappingEngine::map_interchange_typed::<M, T>(
&msg_engine,
&tx_engine,
tree,
tx_group,
true,
nachrichtendaten,
interchangedaten,
)
.map_err(|e| MapperError::Serialization(e.to_string()))?;
Ok((interchange, assembly_diagnostics))
}
pub fn detect_pid(&self, edifact: &str) -> Result<String, MapperError> {
let segments = mig_assembly::tokenize::parse_to_segments(edifact.as_bytes())?;
let chunks = mig_assembly::split_messages(segments)?;
let msg_chunk = chunks.messages.first().ok_or_else(|| {
MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
"No messages found in EDIFACT content".to_string(),
))
})?;
let msg_segments = msg_chunk.message_segments();
mig_assembly::pid_detect::detect_pid(&msg_segments).map_err(MapperError::Assembly)
}
pub fn validate_edifact(
&self,
edifact: &str,
fv: &str,
level: automapper_validation::ValidationLevel,
) -> Result<automapper_validation::ValidationReport, MapperError> {
self.validate_edifact_inner(edifact, fv, None, level)
}
pub fn validate_edifact_for_pid(
&self,
edifact: &str,
fv: &str,
variant: &str,
pid: &str,
level: automapper_validation::ValidationLevel,
) -> Result<automapper_validation::ValidationReport, MapperError> {
self.validate_edifact_inner(edifact, fv, Some((variant, pid)), level)
}
fn validate_edifact_inner(
&self,
edifact: &str,
fv: &str,
known: Option<(&str, &str)>,
level: automapper_validation::ValidationLevel,
) -> Result<automapper_validation::ValidationReport, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let segments = mig_assembly::tokenize::parse_to_segments(edifact.as_bytes())?;
let chunks = mig_assembly::split_messages(segments)?;
let msg_chunk = chunks.messages.first().ok_or_else(|| {
MapperError::Assembly(mig_assembly::AssemblyError::ParseError(
"No messages found in EDIFACT content".to_string(),
))
})?;
let (pid, variant, vc) = match known {
Some((variant, pid)) => {
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
(pid.to_string(), variant.to_string(), vc)
}
None => {
let pid = mig_assembly::pid_detect::detect_pid(&msg_chunk.message_segments())
.map_err(MapperError::Assembly)?;
let pid_key = format!("pid_{pid}");
let (variant, vc) = bundle
.variants
.iter()
.find(|(_, vc)| vc.pid_ahb_workflows.contains_key(&pid_key))
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: "?".to_string(),
pid: pid.clone(),
})?;
(pid, variant.clone(), vc)
}
};
let pid_key = format!("pid_{pid}");
let workflow =
vc.pid_ahb_workflows
.get(&pid_key)
.ok_or_else(|| MapperError::PidNotFound {
fv: fv.to_string(),
variant: variant.clone(),
pid: pid.clone(),
})?;
let filtered_mig = vc
.filtered_mig(&pid)
.ok_or_else(|| MapperError::NoMigSchema {
fv: fv.to_string(),
variant: variant.clone(),
})?;
let mut all_segments = msg_chunk.segments_for_mig(&filtered_mig);
if filtered_mig.segments.iter().any(|s| s.id == "UNZ") {
if let Some(unz) = &chunks.unz {
all_segments.push(unz.clone());
}
}
let evaluator: std::sync::Arc<dyn automapper_validation::ConditionEvaluator> =
match crate::evaluator_factory::create_evaluator(&variant, fv) {
Some(boxed) => std::sync::Arc::from(boxed),
None => std::sync::Arc::new(
automapper_validation::UtilmdStromConditionEvaluatorFV2504::default(),
),
};
let external = automapper_validation::eval::NoOpExternalProvider;
let mut report = automapper_validation::validate_edifact_message(
&all_segments,
&filtered_mig,
workflow,
evaluator,
&external,
level,
);
if let (Some(mig), Some(defs)) = (vc.mig_schema.as_ref(), vc.combined_defs.get(&pid_key)) {
let reverse = mig_bo4e::path_resolver::ReversePathResolver::from_mig(mig);
let field_index =
mig_bo4e::Bo4eFieldIndex::build_with_resolver(defs, &filtered_mig, &reverse);
report.enrich_bo4e_paths(|path, hint| field_index.resolve(path, hint));
}
Ok(report)
}
pub fn validate_bo4e(
&self,
msg_stammdaten: &serde_json::Value,
tx_stammdaten: &[serde_json::Value],
fv: &str,
variant: &str,
pid: &str,
envelope: Option<&InterchangeEnvelope>,
level: automapper_validation::ValidationLevel,
) -> Result<automapper_validation::ValidationReport, MapperError> {
let placeholder;
let envelope = match envelope {
Some(e) => e,
None => {
placeholder = InterchangeEnvelope {
sender: EdifactParty::bdew("9900000000001"),
receiver: EdifactParty::bdew("9900000000002"),
interchange_ref: "1".to_string(),
};
&placeholder
}
};
let edifact = self.render_interchange(
envelope,
&[InterchangeMessage {
message_ref: "1".to_string(),
msg_stammdaten: msg_stammdaten.clone(),
tx_stammdaten: tx_stammdaten.to_vec(),
fv: fv.to_string(),
variant: variant.to_string(),
pid: pid.to_string(),
}],
EntrySegmentCheck::Render,
&EnvelopeOptions::default(),
)?;
self.validate_edifact_for_pid(&edifact, fv, variant, pid, level)
}
pub fn association_code(&self, fv: &str, variant: &str) -> Result<String, MapperError> {
let meta = self.message_metadata(fv, variant)?;
Ok(meta.association_code)
}
pub fn message_metadata(
&self,
fv: &str,
variant: &str,
) -> Result<MessageMetadata, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
let vc = bundle
.variant(variant)
.ok_or_else(|| MapperError::VariantNotFound {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
let mig = vc
.mig_schema
.as_ref()
.ok_or_else(|| MapperError::NoMigSchema {
fv: fv.to_string(),
variant: variant.to_string(),
})?;
Ok(MessageMetadata {
message_type: mig.message_type.clone(),
release: release_code_for_message_type(&mig.message_type),
association_code: mig.version.clone(),
})
}
pub fn to_edifact_interchange(
&self,
envelope: &InterchangeEnvelope,
messages: &[InterchangeMessage],
) -> Result<String, MapperError> {
self.render_interchange(
envelope,
messages,
EntrySegmentCheck::Refuse,
&EnvelopeOptions::default(),
)
}
pub fn to_edifact_interchange_with(
&self,
envelope: &InterchangeEnvelope,
messages: &[InterchangeMessage],
options: &EnvelopeOptions,
) -> Result<String, MapperError> {
self.render_interchange(envelope, messages, EntrySegmentCheck::Refuse, options)
}
fn render_interchange(
&self,
envelope: &InterchangeEnvelope,
messages: &[InterchangeMessage],
check: EntrySegmentCheck,
options: &EnvelopeOptions,
) -> Result<String, MapperError> {
let delimiters = edifact_primitives::EdifactDelimiters::default();
let sep = delimiters.component as char;
let elem = delimiters.element as char;
let seg_term = delimiters.segment as char;
let mut output = String::new();
if options.emit_una {
output.push_str(&format!(
"UNA{}{}{}{}{}{}",
sep, elem, delimiters.decimal as char, delimiters.release as char, ' ', seg_term, ));
}
check_unb_field("datum", "yymmdd", 6, options.datum.as_deref())?;
check_unb_field("zeit", "hhmm", 4, options.zeit.as_deref())?;
let now = chrono::Utc::now();
let date_str = options
.datum
.clone()
.unwrap_or_else(|| now.format("%y%m%d").to_string());
let time_str = options
.zeit
.clone()
.unwrap_or_else(|| now.format("%H%M").to_string());
let sender = &envelope.sender;
let receiver = &envelope.receiver;
let interchange_ref = &envelope.interchange_ref;
output.push_str(&format!(
"UNB{elem}UNOC{sep}3{elem}{sid}{sep}{sq}{elem}{rid}{sep}{rq}{elem}{date_str}{sep}{time_str}{elem}{interchange_ref}{seg_term}",
sid = sender.id,
sq = sender.qualifier,
rid = receiver.id,
rq = receiver.qualifier,
));
let mut message_count = 0u32;
for msg in messages {
let meta = self.message_metadata(&msg.fv, &msg.variant)?;
let body = self.render_message_body(
&msg.msg_stammdaten,
&msg.tx_stammdaten,
&msg.fv,
&msg.variant,
&msg.pid,
check,
)?;
let body_seg_count = body
.split(seg_term)
.filter(|s: &&str| !s.is_empty())
.count();
let segment_count = body_seg_count + 2;
output.push_str(&format!(
"UNH{elem}{ref}{elem}{msg_type}{sep}D{sep}{release}{sep}UN{sep}{assoc}{seg_term}",
ref = msg.message_ref,
msg_type = meta.message_type,
release = meta.release,
assoc = meta.association_code,
));
output.push_str(&body);
output.push_str(&format!(
"UNT{elem}{segment_count}{elem}{ref}{seg_term}",
ref = msg.message_ref,
));
message_count += 1;
}
output.push_str(&format!(
"UNZ{elem}{message_count}{elem}{interchange_ref}{seg_term}",
));
Ok(output)
}
pub fn loaded_format_versions(&self) -> Vec<String> {
self.bundles.lock().unwrap().keys().cloned().collect()
}
pub fn variants(&self, fv: &str) -> Result<Vec<String>, MapperError> {
self.ensure_bundle_loaded(fv)?;
let bundles = self.bundles.lock().unwrap();
let bundle = bundles.get(fv).unwrap();
Ok(bundle.variants.keys().cloned().collect())
}
}
#[derive(Debug, Clone)]
pub struct MessageMetadata {
pub message_type: String,
pub release: String,
pub association_code: String,
}
#[derive(Debug, Clone)]
pub struct InterchangeEnvelope {
pub sender: EdifactParty,
pub receiver: EdifactParty,
pub interchange_ref: String,
}
fn check_unb_field(
field: &'static str,
expected: &'static str,
digits: usize,
value: Option<&str>,
) -> Result<(), MapperError> {
let Some(value) = value else {
return Ok(());
};
if value.len() == digits && value.bytes().all(|b| b.is_ascii_digit()) {
return Ok(());
}
Err(MapperError::MalformedEnvelopeDateTime {
field,
expected,
digits,
value: value.to_string(),
})
}
#[derive(Debug, Clone)]
pub struct EnvelopeOptions {
emit_una: bool,
datum: Option<String>,
zeit: Option<String>,
}
impl Default for EnvelopeOptions {
fn default() -> Self {
Self {
emit_una: true,
datum: None,
zeit: None,
}
}
}
impl EnvelopeOptions {
pub fn emit_una(mut self, emit: bool) -> Self {
self.emit_una = emit;
self
}
pub fn datum_zeit(mut self, datum: impl Into<String>, zeit: impl Into<String>) -> Self {
self.datum = Some(datum.into());
self.zeit = Some(zeit.into());
self
}
pub fn datum_zeit_from(mut self, daten: &mig_bo4e::model::Interchangedaten) -> Self {
self.datum = daten.datum.clone();
self.zeit = daten.zeit.clone();
self
}
}
#[derive(Debug, Clone)]
pub struct EdifactParty {
pub id: String,
pub qualifier: String,
}
impl EdifactParty {
pub fn bdew(id: &str) -> Self {
Self {
id: id.to_string(),
qualifier: "500".to_string(),
}
}
pub fn gs1(id: &str) -> Self {
Self {
id: id.to_string(),
qualifier: "14".to_string(),
}
}
}
#[derive(Debug, Clone)]
pub struct InterchangeMessage {
pub message_ref: String,
pub msg_stammdaten: serde_json::Value,
pub tx_stammdaten: Vec<serde_json::Value>,
pub fv: String,
pub variant: String,
pub pid: String,
}
#[derive(Debug, Clone, Copy)]
enum EntrySegmentCheck {
Refuse,
Render,
}
fn describe_entry_segment_mappings<'d>(
definition_sets: impl IntoIterator<Item = &'d [mig_bo4e::definition::MappingDefinition]>,
source_path: &str,
entry_segment: &str,
) -> (Vec<String>, Vec<String>) {
fn qualifies(unqualified: &str, qualified: &str) -> bool {
!unqualified.contains('_')
&& qualified.len() > unqualified.len()
&& qualified.is_char_boundary(unqualified.len())
&& qualified[..unqualified.len()].eq_ignore_ascii_case(unqualified)
&& qualified.as_bytes()[unqualified.len()] == b'_'
}
fn part_matches(mig_part: &str, def_part: &str) -> bool {
def_part.eq_ignore_ascii_case(mig_part)
|| qualifies(mig_part, def_part)
|| qualifies(def_part, mig_part)
}
let mig_parts: Vec<&str> = source_path.split('.').collect();
let mut entities: Vec<String> = Vec::new();
let mut entry_fields: Vec<String> = Vec::new();
for def in definition_sets.into_iter().flatten() {
let Some(def_path) = def.meta.source_path.as_deref() else {
continue;
};
let def_parts: Vec<&str> = def_path.split('.').collect();
if def_parts.len() != mig_parts.len()
|| !mig_parts
.iter()
.zip(&def_parts)
.all(|(m, d)| part_matches(m, d))
{
continue;
}
if !entities.contains(&def.meta.entity) {
entities.push(def.meta.entity.clone());
}
for (path, mapping) in &def.fields {
let tag = path
.split(['.', '['])
.next()
.unwrap_or_default()
.to_ascii_uppercase();
let target = match mapping {
mig_bo4e::definition::FieldMapping::Simple(t) => t.as_str(),
mig_bo4e::definition::FieldMapping::Structured(f) => f.target.as_str(),
mig_bo4e::definition::FieldMapping::Nested(_) => continue,
};
if tag == entry_segment && !target.is_empty() {
let field = format!("{}.{}", def.meta.entity, target);
if !entry_fields.contains(&field) {
entry_fields.push(field);
}
}
}
}
(entities, entry_fields)
}
fn release_code_for_message_type(msg_type: &str) -> String {
mig_bo4e::model::release_code_for_message_type(msg_type).to_string()
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::Path;
fn data_dir() -> Option<std::path::PathBuf> {
let dist = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../dist");
if dist.join("edifact-data-FV2504.bin").exists() {
return Some(dist);
}
let cache = Path::new(env!("CARGO_MANIFEST_DIR")).join("../../cache/mappings");
if cache.join("FV2504").exists() {
return Some(cache);
}
eprintln!("Skipping test: no DataBundle files found");
None
}
#[test]
fn test_to_edifact_produces_edifact_output() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let msg_stammdaten = serde_json::json!({
"marktteilnehmer": [{
"marktrolle": "MS",
"rollencodenummer": "9900123456789",
"codepflegeCode": "293"
}]
});
let tx_stammdaten = serde_json::json!({
"prozessdaten": {
"pruefidentifikator": "55001",
"vorgangId": "ABC123",
"transaktionsgrund": "E01"
}
});
let result = mapper.to_edifact(
&msg_stammdaten,
&[tx_stammdaten],
"FV2504",
"UTILMD_Strom",
"55001",
);
assert!(result.is_ok(), "to_edifact failed: {:?}", result.err());
let edifact = result.unwrap();
assert!(!edifact.is_empty(), "EDIFACT output should not be empty");
assert!(edifact.contains("NAD"), "Should contain NAD segment");
assert!(edifact.contains("IDE"), "Should contain IDE segment");
}
#[test]
fn test_to_edifact_struct_produces_edifact_output() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let nachricht = serde_json::json!({
"stammdaten": {
"marktteilnehmer": [{
"marktrolle": "MS",
"rollencodenummer": "9900123456789",
"codepflegeCode": "293"
}]
},
"transaktionen": [{
"prozessdaten": {
"pruefidentifikator": "55001",
"vorgangId": "ABC123"
}
}]
});
let result = mapper.to_edifact_struct(&nachricht, "FV2504", "UTILMD_Strom", "55001");
assert!(
result.is_ok(),
"to_edifact_struct failed: {:?}",
result.err()
);
let edifact = result.unwrap();
assert!(!edifact.is_empty(), "EDIFACT output should not be empty");
}
#[test]
fn test_to_edifact_invalid_fv_returns_error() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let result = mapper.to_edifact(
&serde_json::json!({}),
&[serde_json::json!({})],
"FV9999",
"UTILMD_Strom",
"55001",
);
assert!(result.is_err());
}
#[test]
fn test_to_edifact_invalid_variant_returns_error() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let result = mapper.to_edifact(
&serde_json::json!({}),
&[serde_json::json!({})],
"FV2504",
"NONEXISTENT",
"55001",
);
assert!(result.is_err());
}
#[test]
fn test_to_edifact_invalid_pid_returns_error() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let result = mapper.to_edifact(
&serde_json::json!({}),
&[serde_json::json!({})],
"FV2504",
"UTILMD_Strom",
"99999",
);
assert!(result.is_err());
}
#[test]
fn test_association_code() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let code = mapper.association_code("FV2504", "UTILMD_Strom").unwrap();
assert_eq!(code, "S2.1");
let code = mapper.association_code("FV2504", "MSCONS").unwrap();
assert_eq!(code, "2.4c");
}
#[test]
fn test_message_metadata() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let meta = mapper.message_metadata("FV2504", "UTILMD_Strom").unwrap();
assert_eq!(meta.message_type, "UTILMD");
assert_eq!(meta.release, "11A");
assert_eq!(meta.association_code, "S2.1");
}
#[test]
fn test_to_edifact_interchange() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let result = mapper.to_edifact_interchange(
&InterchangeEnvelope {
sender: EdifactParty::bdew("9900000000003"),
receiver: EdifactParty::bdew("9900000000001"),
interchange_ref: "REF001".to_string(),
},
&[InterchangeMessage {
message_ref: "MSG001".to_string(),
msg_stammdaten: serde_json::json!({
"marktteilnehmer": [{
"marktrolle": "MS",
"rollencodenummer": "9900123456789",
"codepflegeCode": "293"
}]
}),
tx_stammdaten: vec![serde_json::json!({
"prozessdaten": {
"pruefidentifikator": "55001",
"vorgangId": "ABC123",
"transaktionsgrund": "E01"
}
})],
fv: "FV2504".to_string(),
variant: "UTILMD_Strom".to_string(),
pid: "55001".to_string(),
}],
);
assert!(
result.is_ok(),
"to_edifact_interchange failed: {:?}",
result.err()
);
let edifact = result.unwrap();
assert!(edifact.starts_with("UNA:+.? '"), "Should start with UNA");
assert!(
edifact.contains("UNB+UNOC:3+9900000000003:500+9900000000001:500+"),
"Should contain UNB with sender/receiver"
);
assert!(
edifact.contains("UNH+MSG001+UTILMD:D:11A:UN:S2.1'"),
"Should contain UNH with correct S009"
);
assert!(edifact.contains("NAD"), "Should contain body NAD segment");
assert!(edifact.contains("UNT+"), "Should contain UNT");
assert!(
edifact.contains("+MSG001'"),
"UNT should reference message ref"
);
assert!(
edifact.contains("UNZ+1+REF001'"),
"Should contain UNZ with count and ref"
);
}
#[test]
fn test_detect_pid_from_rff_z13() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let edifact = "\
UNB+UNOC:3+9978842000002:500+9900269000000:500+250331:1329+REF001'\
UNH+MSG001+UTILMD:D:11A:UN:S2.1'\
BGM+E01+DOC001'\
DTM+137:202503311329?+00:303'\
NAD+MS+9978842000002::293'\
NAD+MR+9900269000000::293'\
IDE+24+TX001'\
DTM+92:202505312200?+00:303'\
DTM+93:202512312300?+00:303'\
STS+7++E01+ZW4+E03'\
LOC+Z16+12345678900'\
RFF+Z13:55001'\
UNT+12+MSG001'\
UNZ+1+REF001'";
let pid = mapper.detect_pid(edifact).unwrap();
assert_eq!(pid, "55001");
}
#[test]
fn test_detect_pid_no_messages_returns_error() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let edifact = "UNB+UNOC:3+SENDER:500+RECEIVER:500+250401:1200+REF'\
UNZ+0+REF'";
assert!(mapper.detect_pid(edifact).is_err());
}
#[test]
fn test_list_pids_returns_entries() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir)).unwrap();
let pids = mapper.list_pids().expect("list_pids should succeed");
assert!(!pids.is_empty(), "should return at least one PID");
assert!(
pids.iter().any(|p| p.pid == "55001"),
"should include PID 55001"
);
assert!(
pids.iter().any(|p| p.fv == "FV2504"),
"should include FV2504"
);
assert!(
pids.iter().any(|p| p.variant == "UTILMD_Strom"),
"should include UTILMD_Strom"
);
}
#[test]
fn test_pid_requirements_returns_requirements() {
let Some(data_dir) = data_dir() else {
return;
};
let mapper = Mapper::from_data_dir(DataDir::path(&data_dir).eager(&["FV2504"])).unwrap();
let req = mapper
.pid_requirements("FV2504", "UTILMD_Strom", "55001")
.expect("pid_requirements should succeed");
assert_eq!(req.pid, "55001");
assert!(
!req.entities.is_empty(),
"55001 should have at least one entity"
);
assert!(
req.entities.iter().any(|e| e.entity == "Prozessdaten"),
"55001 should have a Prozessdaten entity"
);
}
}