Skip to main content

dtmrs_server/
lib.rs

1//! TC 的可复用部分。做成 lib 是为了让集成测试能直接驱动推进器 ——
2//! bin crate 是不能被 tests/ import 的。
3//!
4//! **两个协议层都放在这儿**(`http` 和 `grpc`),理由同上:它们早先一个在
5//! 这里、一个在 bin crate 的 main.rs 里,结果 gRPC 层有 86% 覆盖率而 HTTP 层
6//! 是 0% —— 而「两边不许漂移」恰恰是这个项目最要紧的结构约束之一。
7//! 现在两边都能被 tests/ 拿到,可以用同一组用例做等价性测试。
8
9pub 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
21/// 各模式共用的全局事务骨架
22fn 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
38/// 把一组 SAGA 步骤展开成"1 个全局事务 + 2N 个分支"。
39///
40/// HTTP `submit` 和集成测试都走这里,避免两处构造逻辑漂移 ——
41/// 分支号规则错位是那种测试全绿、线上补偿补错对象的 bug。
42pub 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                // 每步自己的业务数据。正向和补偿共用同一份 ——
58                // 补偿需要知道当初做了什么才能撤销
59                payload: s.payload.clone(),
60                status: BranchStatus::Prepared,
61            });
62        }
63    }
64    (g, branches)
65}
66
67/// 二阶段消息:**只有正向分支,没有补偿**。
68///
69/// `grace_secs` 是回查前的宽限期:客户端 prepare 之后正常会很快 submit,
70/// 立刻回查是白问一次。给几秒钟等它自己来。
71pub fn msg_rows(
72    gid: &str,
73    actions: &[String],
74    query_prepared: &str,
75    grace_secs: i64,
76) -> (GlobalRow, Vec<BranchRow>) {
77    // 用 SagaStep 复用 payload 格式,compensate 留空(msg 没有补偿)
78    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
105/// TCC:`prepare` 只建全局事务,**分支是客户端在 try 阶段动态登记的**
106/// (`Store::register_branch`)。所以这里不产生任何分支行。
107pub fn tcc_rows(gid: &str) -> GlobalRow {
108    global(gid, TransType::Tcc, GlobalStatus::Prepared, "[]".into())
109}