use mako_engine::{
error::WorkflowError,
ids::DeadlineId,
types::MarktpartnerCode,
workflow::{CommandPayload, EventPayload, Workflow, WorkflowOutput},
};
pub const WORKFLOW_NAME: &str = "wim-steuerungsauftrag";
pub const STEUERUNGSAUFTRAG_DEADLINE_LABEL: &str = "wim-steuerungsauftrag-deadline";
#[derive(Debug, Clone, PartialEq, Eq, serde::Serialize, serde::Deserialize)]
pub enum SteuerungsCommandType {
Konfiguration,
InitialZustand,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(deny_unknown_fields)]
pub struct SteuerungsauftragData {
pub tx_id: String,
pub sender_gln: MarktpartnerCode,
pub location_id: String,
pub command_type: SteuerungsCommandType,
pub execution_time_from: String,
pub max_power_kw: Option<String>,
pub execution_time_until: Option<String>,
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "type", content = "data")]
pub enum SteuerungsauftragEvent {
KonfigurationReceived {
tx_id: String,
sender_gln: MarktpartnerCode,
location_id: String,
execution_time_from: String,
max_power_kw: String,
execution_time_until: Option<String>,
},
InitialZustandReceived {
tx_id: String,
sender_gln: MarktpartnerCode,
location_id: String,
execution_time_from: String,
},
EndantwortPositiv {
reference_id: String,
},
EndantwortNegativ {
reason: Option<String>,
},
DeadlineExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl EventPayload for SteuerungsauftragEvent {
fn event_type(&self) -> &'static str {
match self {
Self::KonfigurationReceived { .. } => "WimSteuerungsauftragKonfigurationReceived",
Self::InitialZustandReceived { .. } => "WimSteuerungsauftragInitialZustandReceived",
Self::EndantwortPositiv { .. } => "WimSteuerungsauftragEndantwortPositiv",
Self::EndantwortNegativ { .. } => "WimSteuerungsauftragEndantwortNegativ",
Self::DeadlineExpired { .. } => "WimSteuerungsauftragDeadlineExpired",
}
}
}
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize)]
#[serde(tag = "status", content = "data")]
pub enum SteuerungsauftragState {
New,
Received(SteuerungsauftragData),
Completed(SteuerungsauftragData),
Rejected {
tx_id: Option<String>,
reason: String,
},
}
impl Default for SteuerungsauftragState {
fn default() -> Self {
Self::New
}
}
impl SteuerungsauftragState {
#[must_use]
pub fn status_str(&self) -> &'static str {
match self {
Self::New => "New",
Self::Received(_) => "Received",
Self::Completed(_) => "Completed",
Self::Rejected { .. } => "Rejected",
}
}
}
#[derive(Clone)]
pub enum SteuerungsauftragCommand {
ReceiveKonfiguration {
tx_id: String,
sender_gln: MarktpartnerCode,
location_id: String,
execution_time_from: String,
max_power_kw: String,
execution_time_until: Option<String>,
},
ReceiveInitialZustand {
tx_id: String,
sender_gln: MarktpartnerCode,
location_id: String,
execution_time_from: String,
},
SendEndantwortPositiv {
reference_id: String,
},
SendEndantwortNegativ {
reason: Option<String>,
},
TimeoutExpired {
deadline_id: DeadlineId,
label: Box<str>,
},
}
impl CommandPayload for SteuerungsauftragCommand {}
pub struct WimSteuerungsauftragWorkflow;
impl Workflow for WimSteuerungsauftragWorkflow {
type State = SteuerungsauftragState;
type Event = SteuerungsauftragEvent;
type Command = SteuerungsauftragCommand;
fn on_deadline(
deadline: &mako_engine::deadline::Deadline,
state: &Self::State,
) -> Option<Self::Command> {
match (deadline.label(), state) {
(STEUERUNGSAUFTRAG_DEADLINE_LABEL, SteuerungsauftragState::Received(_)) => {
Some(SteuerungsauftragCommand::TimeoutExpired {
deadline_id: deadline.deadline_id(),
label: deadline.label().into(),
})
}
_ => None,
}
}
fn apply(state: Self::State, event: &Self::Event) -> Self::State {
match event {
SteuerungsauftragEvent::KonfigurationReceived {
tx_id,
sender_gln,
location_id,
execution_time_from,
max_power_kw,
execution_time_until,
} => SteuerungsauftragState::Received(SteuerungsauftragData {
tx_id: tx_id.clone(),
sender_gln: sender_gln.clone(),
location_id: location_id.clone(),
command_type: SteuerungsCommandType::Konfiguration,
execution_time_from: execution_time_from.clone(),
max_power_kw: Some(max_power_kw.clone()),
execution_time_until: execution_time_until.clone(),
}),
SteuerungsauftragEvent::InitialZustandReceived {
tx_id,
sender_gln,
location_id,
execution_time_from,
} => SteuerungsauftragState::Received(SteuerungsauftragData {
tx_id: tx_id.clone(),
sender_gln: sender_gln.clone(),
location_id: location_id.clone(),
command_type: SteuerungsCommandType::InitialZustand,
execution_time_from: execution_time_from.clone(),
max_power_kw: None,
execution_time_until: None,
}),
SteuerungsauftragEvent::EndantwortPositiv { .. } => {
if let SteuerungsauftragState::Received(data) = state {
SteuerungsauftragState::Completed(data)
} else {
state
}
}
SteuerungsauftragEvent::EndantwortNegativ { reason } => {
let tx_id = match &state {
SteuerungsauftragState::Received(d) => Some(d.tx_id.clone()),
_ => None,
};
SteuerungsauftragState::Rejected {
tx_id,
reason: reason
.clone()
.unwrap_or_else(|| "negative response".to_owned()),
}
}
SteuerungsauftragEvent::DeadlineExpired { label, .. } => match state {
SteuerungsauftragState::Completed(_) | SteuerungsauftragState::Rejected { .. } => {
state
}
_ => {
let tx_id = if let SteuerungsauftragState::Received(ref d) = state {
Some(d.tx_id.clone())
} else {
None
};
SteuerungsauftragState::Rejected {
tx_id,
reason: format!("5-Werktage deadline expired: {label}"),
}
}
},
}
}
fn handle(
state: &Self::State,
command: Self::Command,
) -> Result<WorkflowOutput<Self::Event>, WorkflowError> {
match command {
SteuerungsauftragCommand::ReceiveKonfiguration {
tx_id,
sender_gln,
location_id,
execution_time_from,
max_power_kw,
execution_time_until,
} => {
if !matches!(state, SteuerungsauftragState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
Ok(vec![SteuerungsauftragEvent::KonfigurationReceived {
tx_id,
sender_gln,
location_id,
execution_time_from,
max_power_kw,
execution_time_until,
}]
.into())
}
SteuerungsauftragCommand::ReceiveInitialZustand {
tx_id,
sender_gln,
location_id,
execution_time_from,
} => {
if !matches!(state, SteuerungsauftragState::New) {
return Err(WorkflowError::invalid_state("New", state.status_str()));
}
Ok(vec![SteuerungsauftragEvent::InitialZustandReceived {
tx_id,
sender_gln,
location_id,
execution_time_from,
}]
.into())
}
SteuerungsauftragCommand::SendEndantwortPositiv { reference_id } => {
if !matches!(state, SteuerungsauftragState::Received(_)) {
return Err(WorkflowError::invalid_state("Received", state.status_str()));
}
Ok(vec![SteuerungsauftragEvent::EndantwortPositiv { reference_id }].into())
}
SteuerungsauftragCommand::SendEndantwortNegativ { reason } => {
if !matches!(state, SteuerungsauftragState::Received(_)) {
return Err(WorkflowError::invalid_state("Received", state.status_str()));
}
Ok(vec![SteuerungsauftragEvent::EndantwortNegativ { reason }].into())
}
SteuerungsauftragCommand::TimeoutExpired { deadline_id, label } => {
if matches!(
state,
SteuerungsauftragState::Completed(_) | SteuerungsauftragState::Rejected { .. }
) {
return Ok(WorkflowOutput::events(vec![]));
}
Ok(vec![SteuerungsauftragEvent::DeadlineExpired { deadline_id, label }].into())
}
}
}
}