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 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
22/// 各模式共用的全局事务骨架
23fn 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
39/// 把一组 SAGA 步骤展开成"1 个全局事务 + 2N 个分支"。
40///
41/// HTTP `submit` 和集成测试都走这里,避免两处构造逻辑漂移 ——
42/// 分支号规则错位是那种测试全绿、线上补偿补错对象的 bug。
43pub 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                // 每步自己的业务数据。正向和补偿共用同一份 ——
59                // 补偿需要知道当初做了什么才能撤销
60                payload: s.payload.clone(),
61                status: BranchStatus::Prepared,
62            });
63        }
64    }
65    (g, branches)
66}
67
68/// 二阶段消息:**只有正向分支,没有补偿**。
69///
70/// `grace_secs` 是回查前的宽限期:客户端 prepare 之后正常会很快 submit,
71/// 立刻回查是白问一次。给几秒钟等它自己来。
72pub fn msg_rows(
73    gid: &str,
74    actions: &[String],
75    query_prepared: &str,
76    grace_secs: i64,
77) -> (GlobalRow, Vec<BranchRow>) {
78    // 用 SagaStep 复用 payload 格式,compensate 留空(msg 没有补偿)
79    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
106/// TCC:`prepare` 只建全局事务,**分支是客户端在 try 阶段动态登记的**
107/// (`Store::register_branch`)。所以这里不产生任何分支行。
108pub fn tcc_rows(gid: &str) -> GlobalRow {
109    global(gid, TransType::Tcc, GlobalStatus::Prepared, "[]".into())
110}