pub mod api;
pub mod auth;
pub mod driver;
pub mod embedded;
#[cfg(feature = "grpc")]
pub mod grpc;
pub mod http;
pub mod registry;
pub mod workflow;
use dtmrs_core::{BranchOp, BranchStatus, GlobalStatus, SagaStep, TransType};
use dtmrs_store::{BranchRow, GlobalRow};
fn global(gid: &str, tt: TransType, status: GlobalStatus, payload: String) -> GlobalRow {
GlobalRow {
gid: gid.to_string(),
trans_type: tt,
status,
payload,
next_cron_time: dtmrs_store::now(),
next_cron_interval: 0,
owner: String::new(),
rollback_reason: String::new(),
query_prepared: String::new(),
create_time: 0,
finish_time: None,
}
}
pub fn saga_rows(gid: &str, steps: &[SagaStep]) -> (GlobalRow, Vec<BranchRow>) {
let payload = serde_json::to_string(steps).unwrap_or_else(|_| "[]".into());
let g = global(gid, TransType::Saga, GlobalStatus::Submitted, payload);
let mut branches = Vec::with_capacity(steps.len() * 2);
for (i, s) in steps.iter().enumerate() {
let bid = driver::branch_id(i);
for (op, url) in [
(BranchOp::Action, &s.action),
(BranchOp::Compensate, &s.compensate),
] {
branches.push(BranchRow {
gid: gid.to_string(),
branch_id: bid.clone(),
op,
url: url.clone(),
payload: s.payload.clone(),
status: BranchStatus::Prepared,
});
}
}
(g, branches)
}
pub fn msg_rows(
gid: &str,
actions: &[String],
query_prepared: &str,
grace_secs: i64,
) -> (GlobalRow, Vec<BranchRow>) {
let steps: Vec<SagaStep> = actions
.iter()
.map(|a| SagaStep {
action: a.clone(),
compensate: String::new(),
payload: String::new(),
})
.collect();
let payload = serde_json::to_string(&steps).unwrap_or_else(|_| "[]".into());
let mut g = global(gid, TransType::Msg, GlobalStatus::Prepared, payload);
g.query_prepared = query_prepared.to_string();
g.next_cron_time = dtmrs_store::now() + grace_secs.max(0);
let branches = actions
.iter()
.enumerate()
.map(|(i, a)| BranchRow {
gid: gid.to_string(),
branch_id: driver::branch_id(i),
op: BranchOp::Action,
url: a.clone(),
payload: String::new(),
status: BranchStatus::Prepared,
})
.collect();
(g, branches)
}
pub fn tcc_rows(gid: &str) -> GlobalRow {
global(gid, TransType::Tcc, GlobalStatus::Prepared, "[]".into())
}