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