use mako_engine::{
error::WorkflowError,
ids::DeadlineId,
types::{MarktpartnerCode, MessageRef},
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const COMDIS_PIDS: &[u32] = &[29001, 29002];
pub const WORKFLOW_NAME: &str = "gpke-comdis";
pub const COMDIS_APERAK_WINDOW_LABEL: &str = "gpke-comdis-aperak-next-workday";
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum GpkeComdisEvent {
Comdis29001Received {
message_ref: MessageRef,
rechnungsnummer: String,
lf_code: MarktpartnerCode,
nb_code: MarktpartnerCode,
disputed_amount_ct: i64,
received_at: time::OffsetDateTime,
},
ValidationPassed {
message_ref: MessageRef,
},
ValidationFailed {
errors: Vec<String>,
},
Comdis29002Dispatched {
counter_ref: MessageRef,
accepted: bool,
reason_code: Option<String>,
},
ComdisResolved {
outcome: ComdisOutcome,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for GpkeComdisEvent {
fn event_type(&self) -> &'static str {
match self {
Self::Comdis29001Received { .. } => "GpkeComdis29001Received",
Self::ValidationPassed { .. } => "GpkeComdisValidationPassed",
Self::ValidationFailed { .. } => "GpkeComdisValidationFailed",
Self::Comdis29002Dispatched { .. } => "GpkeComdis29002Dispatched",
Self::ComdisResolved { .. } => "GpkeComdisResolved",
Self::DeadlineExpired { .. } => "GpkeComdisDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(rename_all = "snake_case")]
pub enum ComdisOutcome {
Settled,
Withdrawn,
EscalatedBnetza,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct GpkeComdisData {
pub message_ref: MessageRef,
pub rechnungsnummer: String,
pub lf_code: MarktpartnerCode,
pub nb_code: MarktpartnerCode,
pub disputed_amount_ct: i64,
pub received_at: time::OffsetDateTime,
}
#[derive(Debug, Clone, Default, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum GpkeComdisState {
#[default]
New,
Initiated(GpkeComdisData),
ValidationPassed(GpkeComdisData),
Answered(GpkeComdisData),
Resolved(GpkeComdisData),
Rejected {
reason: String,
},
}
#[derive(Debug, Clone)]
pub enum GpkeComdisCommand {
Receive29001 {
message_ref: MessageRef,
rechnungsnummer: String,
lf_code: MarktpartnerCode,
nb_code: MarktpartnerCode,
disputed_amount_ct: i64,
received_at: time::OffsetDateTime,
validation_passed: bool,
validation_errors: Vec<String>,
},
Dispatch29002 {
counter_ref: MessageRef,
accepted: bool,
reason_code: Option<String>,
},
Resolve {
outcome: ComdisOutcome,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for GpkeComdisCommand {}
pub struct GpkeComdisWorkflow;
impl Workflow for GpkeComdisWorkflow {
type State = GpkeComdisState;
type Event = GpkeComdisEvent;
type Command = GpkeComdisCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(
COMDIS_APERAK_WINDOW_LABEL,
GpkeComdisState::Initiated(_) | GpkeComdisState::ValidationPassed(_),
) => Some(GpkeComdisCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
}),
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
GpkeComdisEvent::Comdis29001Received {
message_ref,
rechnungsnummer,
lf_code,
nb_code,
disputed_amount_ct,
received_at,
} => GpkeComdisState::Initiated(GpkeComdisData {
message_ref: message_ref.clone(),
rechnungsnummer: rechnungsnummer.clone(),
lf_code: lf_code.clone(),
nb_code: nb_code.clone(),
disputed_amount_ct: *disputed_amount_ct,
received_at: *received_at,
}),
GpkeComdisEvent::ValidationPassed { .. } => match state {
GpkeComdisState::Initiated(data) => GpkeComdisState::ValidationPassed(data),
other => other,
},
GpkeComdisEvent::ValidationFailed { errors } => GpkeComdisState::Rejected {
reason: errors.join("; "),
},
GpkeComdisEvent::Comdis29002Dispatched { .. } => match state {
GpkeComdisState::ValidationPassed(data) => GpkeComdisState::Answered(data),
other => other,
},
GpkeComdisEvent::ComdisResolved { .. } => match state {
GpkeComdisState::Answered(data) => GpkeComdisState::Resolved(data),
other => other,
},
GpkeComdisEvent::DeadlineExpired { label, .. } => GpkeComdisState::Rejected {
reason: format!("APERAK deadline expired: {label}"),
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
GpkeComdisCommand::Receive29001 {
message_ref,
rechnungsnummer,
lf_code,
nb_code,
disputed_amount_ct,
received_at,
validation_passed,
validation_errors,
} => {
if !matches!(state, GpkeComdisState::New) {
return Err(WorkflowError::CommandRejected {
reason: "Receive29001 requires New state".into(),
});
}
let received_event = GpkeComdisEvent::Comdis29001Received {
message_ref: message_ref.clone(),
rechnungsnummer,
lf_code,
nb_code,
disputed_amount_ct,
received_at,
};
let validation_event = if validation_passed {
GpkeComdisEvent::ValidationPassed { message_ref }
} else {
GpkeComdisEvent::ValidationFailed {
errors: validation_errors,
}
};
Ok(WorkflowOutput::from(vec![received_event, validation_event]))
}
GpkeComdisCommand::Dispatch29002 {
counter_ref,
accepted,
reason_code,
} => {
if !matches!(state, GpkeComdisState::ValidationPassed(_)) {
return Err(WorkflowError::CommandRejected {
reason: "Dispatch29002 requires ValidationPassed state".into(),
});
}
Ok(WorkflowOutput::from(vec![
GpkeComdisEvent::Comdis29002Dispatched {
counter_ref,
accepted,
reason_code,
},
]))
}
GpkeComdisCommand::Resolve { outcome } => {
if !matches!(state, GpkeComdisState::Answered(_)) {
return Err(WorkflowError::CommandRejected {
reason: "Resolve requires Answered state".into(),
});
}
Ok(WorkflowOutput::from(vec![
GpkeComdisEvent::ComdisResolved { outcome },
]))
}
GpkeComdisCommand::TimeoutExpired { deadline_id, label } => {
Ok(WorkflowOutput::from(vec![
GpkeComdisEvent::DeadlineExpired { deadline_id, label },
]))
}
}
}
}
pub const COMDIS_RECORDS_DDL: &str = r"
CREATE TABLE IF NOT EXISTS comdis_records (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
tenant TEXT NOT NULL,
rechnungsnummer TEXT NOT NULL,
comdis_ref TEXT NOT NULL,
lf_mp_id TEXT NOT NULL,
nb_mp_id TEXT NOT NULL,
disputed_amount_ct BIGINT NOT NULL,
status TEXT NOT NULL DEFAULT 'open'
CHECK (status IN ('open','answered','resolved','escalated')),
counter_ref TEXT,
outcome TEXT
CHECK (outcome IN ('settled','withdrawn','escalated_bnetza') OR outcome IS NULL),
received_at TIMESTAMPTZ NOT NULL DEFAULT now(),
answered_at TIMESTAMPTZ,
resolved_at TIMESTAMPTZ,
UNIQUE (tenant, rechnungsnummer, comdis_ref)
);
CREATE INDEX IF NOT EXISTS comdis_records_status ON comdis_records (tenant, status);
CREATE INDEX IF NOT EXISTS comdis_records_lf_mp_id ON comdis_records (tenant, lf_mp_id);
";