Skip to main content

dtmrs_server/
lib.rs

1//! TC 的可复用部分。做成 lib 是为了让集成测试能直接驱动推进器 ——
2//! bin crate 是不能被 tests/ import 的。
3
4pub mod api;
5pub mod driver;
6pub mod embedded;
7#[cfg(feature = "grpc")]
8pub mod grpc;
9pub mod registry;
10pub mod workflow;
11
12use dtmrs_core::{BranchOp, BranchStatus, GlobalStatus, SagaStep, TransType};
13use dtmrs_store::{BranchRow, GlobalRow};
14
15/// 各模式共用的全局事务骨架
16fn global(gid: &str, tt: TransType, status: GlobalStatus, payload: String) -> GlobalRow {
17    GlobalRow {
18        gid: gid.to_string(),
19        trans_type: tt,
20        status,
21        payload,
22        next_cron_time: dtmrs_store::now(),
23        next_cron_interval: 0,
24        owner: String::new(),
25        rollback_reason: String::new(),
26        query_prepared: String::new(),
27        create_time: 0,
28        finish_time: None,
29    }
30}
31
32/// 把一组 SAGA 步骤展开成"1 个全局事务 + 2N 个分支"。
33///
34/// HTTP `submit` 和集成测试都走这里,避免两处构造逻辑漂移 ——
35/// 分支号规则错位是那种测试全绿、线上补偿补错对象的 bug。
36pub fn saga_rows(gid: &str, steps: &[SagaStep]) -> (GlobalRow, Vec<BranchRow>) {
37    let payload = serde_json::to_string(steps).unwrap_or_else(|_| "[]".into());
38    let g = global(gid, TransType::Saga, GlobalStatus::Submitted, payload);
39    let mut branches = Vec::with_capacity(steps.len() * 2);
40    for (i, s) in steps.iter().enumerate() {
41        let bid = driver::branch_id(i);
42        for (op, url) in [
43            (BranchOp::Action, &s.action),
44            (BranchOp::Compensate, &s.compensate),
45        ] {
46            branches.push(BranchRow {
47                gid: gid.to_string(),
48                branch_id: bid.clone(),
49                op,
50                url: url.clone(),
51                payload: String::new(),
52                status: BranchStatus::Prepared,
53            });
54        }
55    }
56    (g, branches)
57}
58
59/// 二阶段消息:**只有正向分支,没有补偿**。
60///
61/// `grace_secs` 是回查前的宽限期:客户端 prepare 之后正常会很快 submit,
62/// 立刻回查是白问一次。给几秒钟等它自己来。
63pub fn msg_rows(
64    gid: &str,
65    actions: &[String],
66    query_prepared: &str,
67    grace_secs: i64,
68) -> (GlobalRow, Vec<BranchRow>) {
69    // 用 SagaStep 复用 payload 格式,compensate 留空(msg 没有补偿)
70    let steps: Vec<SagaStep> = actions
71        .iter()
72        .map(|a| SagaStep {
73            action: a.clone(),
74            compensate: String::new(),
75        })
76        .collect();
77    let payload = serde_json::to_string(&steps).unwrap_or_else(|_| "[]".into());
78    let mut g = global(gid, TransType::Msg, GlobalStatus::Prepared, payload);
79    g.query_prepared = query_prepared.to_string();
80    g.next_cron_time = dtmrs_store::now() + grace_secs.max(0);
81    let branches = actions
82        .iter()
83        .enumerate()
84        .map(|(i, a)| BranchRow {
85            gid: gid.to_string(),
86            branch_id: driver::branch_id(i),
87            op: BranchOp::Action,
88            url: a.clone(),
89            payload: String::new(),
90            status: BranchStatus::Prepared,
91        })
92        .collect();
93    (g, branches)
94}
95
96/// TCC:`prepare` 只建全局事务,**分支是客户端在 try 阶段动态登记的**
97/// (`Store::register_branch`)。所以这里不产生任何分支行。
98pub fn tcc_rows(gid: &str) -> GlobalRow {
99    global(gid, TransType::Tcc, GlobalStatus::Prepared, "[]".into())
100}