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                // 每步自己的业务数据。正向和补偿共用同一份 ——
52                // 补偿需要知道当初做了什么才能撤销
53                payload: s.payload.clone(),
54                status: BranchStatus::Prepared,
55            });
56        }
57    }
58    (g, branches)
59}
60
61/// 二阶段消息:**只有正向分支,没有补偿**。
62///
63/// `grace_secs` 是回查前的宽限期:客户端 prepare 之后正常会很快 submit,
64/// 立刻回查是白问一次。给几秒钟等它自己来。
65pub fn msg_rows(
66    gid: &str,
67    actions: &[String],
68    query_prepared: &str,
69    grace_secs: i64,
70) -> (GlobalRow, Vec<BranchRow>) {
71    // 用 SagaStep 复用 payload 格式,compensate 留空(msg 没有补偿)
72    let steps: Vec<SagaStep> = actions
73        .iter()
74        .map(|a| SagaStep {
75            action: a.clone(),
76            compensate: String::new(),
77            payload: String::new(),
78        })
79        .collect();
80    let payload = serde_json::to_string(&steps).unwrap_or_else(|_| "[]".into());
81    let mut g = global(gid, TransType::Msg, GlobalStatus::Prepared, payload);
82    g.query_prepared = query_prepared.to_string();
83    g.next_cron_time = dtmrs_store::now() + grace_secs.max(0);
84    let branches = actions
85        .iter()
86        .enumerate()
87        .map(|(i, a)| BranchRow {
88            gid: gid.to_string(),
89            branch_id: driver::branch_id(i),
90            op: BranchOp::Action,
91            url: a.clone(),
92            payload: String::new(),
93            status: BranchStatus::Prepared,
94        })
95        .collect();
96    (g, branches)
97}
98
99/// TCC:`prepare` 只建全局事务,**分支是客户端在 try 阶段动态登记的**
100/// (`Store::register_branch`)。所以这里不产生任何分支行。
101pub fn tcc_rows(gid: &str) -> GlobalRow {
102    global(gid, TransType::Tcc, GlobalStatus::Prepared, "[]".into())
103}