1pub mod api;
10pub mod driver;
11pub mod embedded;
12#[cfg(feature = "grpc")]
13pub mod grpc;
14pub mod http;
15pub mod registry;
16pub mod workflow;
17
18use dtmrs_core::{BranchOp, BranchStatus, GlobalStatus, SagaStep, TransType};
19use dtmrs_store::{BranchRow, GlobalRow};
20
21fn global(gid: &str, tt: TransType, status: GlobalStatus, payload: String) -> GlobalRow {
23 GlobalRow {
24 gid: gid.to_string(),
25 trans_type: tt,
26 status,
27 payload,
28 next_cron_time: dtmrs_store::now(),
29 next_cron_interval: 0,
30 owner: String::new(),
31 rollback_reason: String::new(),
32 query_prepared: String::new(),
33 create_time: 0,
34 finish_time: None,
35 }
36}
37
38pub fn saga_rows(gid: &str, steps: &[SagaStep]) -> (GlobalRow, Vec<BranchRow>) {
43 let payload = serde_json::to_string(steps).unwrap_or_else(|_| "[]".into());
44 let g = global(gid, TransType::Saga, GlobalStatus::Submitted, payload);
45 let mut branches = Vec::with_capacity(steps.len() * 2);
46 for (i, s) in steps.iter().enumerate() {
47 let bid = driver::branch_id(i);
48 for (op, url) in [
49 (BranchOp::Action, &s.action),
50 (BranchOp::Compensate, &s.compensate),
51 ] {
52 branches.push(BranchRow {
53 gid: gid.to_string(),
54 branch_id: bid.clone(),
55 op,
56 url: url.clone(),
57 payload: s.payload.clone(),
60 status: BranchStatus::Prepared,
61 });
62 }
63 }
64 (g, branches)
65}
66
67pub fn msg_rows(
72 gid: &str,
73 actions: &[String],
74 query_prepared: &str,
75 grace_secs: i64,
76) -> (GlobalRow, Vec<BranchRow>) {
77 let steps: Vec<SagaStep> = actions
79 .iter()
80 .map(|a| SagaStep {
81 action: a.clone(),
82 compensate: String::new(),
83 payload: String::new(),
84 })
85 .collect();
86 let payload = serde_json::to_string(&steps).unwrap_or_else(|_| "[]".into());
87 let mut g = global(gid, TransType::Msg, GlobalStatus::Prepared, payload);
88 g.query_prepared = query_prepared.to_string();
89 g.next_cron_time = dtmrs_store::now() + grace_secs.max(0);
90 let branches = actions
91 .iter()
92 .enumerate()
93 .map(|(i, a)| BranchRow {
94 gid: gid.to_string(),
95 branch_id: driver::branch_id(i),
96 op: BranchOp::Action,
97 url: a.clone(),
98 payload: String::new(),
99 status: BranchStatus::Prepared,
100 })
101 .collect();
102 (g, branches)
103}
104
105pub fn tcc_rows(gid: &str) -> GlobalRow {
108 global(gid, TransType::Tcc, GlobalStatus::Prepared, "[]".into())
109}