Skip to main content

dtmrs_server/
driver.rs

1//! 事务推进器:把状态机的决策落成真实的 HTTP 调用和状态更新。
2//!
3//! 这里是**唯一**会修改全局事务状态的地方,决策全部来自 `dtmrs_core::saga_advance`,
4//! 本文件只负责 I/O。这样状态迁移的正确性可以在 core 里纯单测覆盖。
5
6use crate::registry::{parse_target, BranchCtx, Registry, Target};
7use dtmrs_core::{
8    msg_advance, saga_advance, tcc_advance, xa_advance, Advance, BranchOp, BranchResult,
9    BranchStatus, GlobalStatus, SagaStep, TransType,
10};
11use dtmrs_store::{GlobalRow, Store};
12use std::sync::Arc;
13use std::time::Duration;
14use tracing::{error, info, warn};
15
16#[derive(Clone)]
17pub struct Driver {
18    pub store: Store,
19    pub http: reqwest::Client,
20    pub owner: String,
21    /// 租约时长(秒)。持租约的实例崩了,这么久之后别的实例接手
22    pub lease: i64,
23    /// 重试退避策略。默认 10s 起、300s 封顶,可用环境变量改
24    pub retry: dtmrs_core::RetryPolicy,
25    /// 分支调用超时(秒),只为可观测性保留一份
26    branch_timeout_secs: u64,
27    /// 并行推进的 worker 数
28    pub workers: usize,
29    /// 进程内分支注册表。嵌入式模式用,纯 HTTP 部署时是空表
30    pub registry: Arc<Registry>,
31    /// gRPC 分支调用器(带 channel 缓存)
32    #[cfg(feature = "grpc")]
33    pub grpc: crate::grpc::client::GrpcCaller,
34    /// workflow 函数注册表。跟 registry 一样,纯 HTTP 部署时是空表
35    pub workflows: Arc<crate::workflow::WorkflowRegistry>,
36}
37
38impl Driver {
39    /// 默认配置的推进器。分支超时 10s、租约 30s、退避 10s→300s。
40    /// 想按环境变量配就用 [`Driver::from_env`]
41    pub fn new(store: Store, owner: String) -> Self {
42        Self::with_config(store, owner, DriverConfig::default())
43    }
44
45    /// 按环境变量配置:
46    ///
47    /// | 变量 | 默认 | 说明 |
48    /// |---|---|---|
49    /// | `DTMRS_BRANCH_TIMEOUT` | 10 | 调一个分支最多等几秒 |
50    /// | `DTMRS_LEASE` | 30 | 租约时长(秒) |
51    /// | `DTMRS_RETRY_INTERVAL` | 10 | 首次重试间隔(秒) |
52    /// | `DTMRS_RETRY_MAX_INTERVAL` | 300 | 退避上限(秒) |
53    /// | `DTMRS_WORKERS` | 16 | 并行推进的 worker 数 |
54    ///
55    /// 存储连接池另有 `DTMRS_DB_POOL`(默认 32),跟 worker 数**要一起调** ——
56    /// 池子小于 worker 数时,多出来的 worker 只会排队等连接
57    ///
58    /// **非法值一律退回默认**,绝不因为配置写错就让推进器起不来
59    pub fn from_env(store: Store, owner: String) -> Self {
60        Self::with_config(store, owner, DriverConfig::from_env())
61    }
62
63    pub fn with_config(store: Store, owner: String, cfg: DriverConfig) -> Self {
64        Self {
65            store,
66            http: reqwest::Client::builder()
67                .timeout(Duration::from_secs(cfg.branch_timeout_secs.max(1) as u64))
68                .build()
69                .expect("build http client"),
70            owner,
71            lease: cfg.lease_secs,
72            retry: cfg.retry,
73            branch_timeout_secs: cfg.branch_timeout_secs.max(1) as u64,
74            workers: cfg.workers.max(1),
75            registry: Arc::new(Registry::new()),
76            #[cfg(feature = "grpc")]
77            grpc: crate::grpc::client::GrpcCaller::new(Duration::from_secs(
78                cfg.branch_timeout_secs.max(1) as u64,
79            )),
80            workflows: Arc::new(crate::workflow::WorkflowRegistry::new()),
81        }
82    }
83
84    /// 当前的分支调用超时(秒),启动日志里打出来方便确认配置生效
85    pub fn http_timeout_secs(&self) -> u64 {
86        self.branch_timeout_secs
87    }
88
89    /// 挂上进程内分支注册表 —— 嵌入式模式的入口
90    pub fn with_registry(mut self, r: Arc<Registry>) -> Self {
91        self.registry = r;
92        self
93    }
94
95    /// 挂上 workflow 函数注册表
96    pub fn with_workflows(mut self, w: Arc<crate::workflow::WorkflowRegistry>) -> Self {
97        self.workflows = w;
98        self
99    }
100
101    /// 给 `grpcs://` 分支加一个额外信任的 CA(PEM 内容,不是路径)。
102    ///
103    /// 独立部署时用环境变量 `DTMRS_GRPC_CA` 就够了;这个方法是给
104    /// **嵌入式宿主**和测试用的 —— 宿主往往已经从别处拿到了证书,
105    /// 不想为了传一份 PEM 再去设进程级环境变量(并行测试里还会互相打架)。
106    ///
107    /// PEM 里没有证书块时会被忽略并打警告,见 `grpc::client` 里的 `check_ca_pem`。
108    #[cfg(feature = "grpc")]
109    pub fn with_grpc_ca_pem(mut self, pem: impl Into<Vec<u8>>) -> Self {
110        self.grpc = self.grpc.with_ca_pem(pem);
111        self
112    }
113
114    /// 常驻推进器。起 `workers` 个并行的抢占循环。
115    ///
116    /// # 为什么可以直接并行,不需要新的并发控制
117    ///
118    /// 每个 worker 都走 `lock_one_due` —— 那是一次**原子抢占**(SQL 靠带条件的
119    /// UPDATE,Redis 靠 Lua 脚本),抢到才推。所以进程内 N 个 worker
120    /// 跟部署 N 个实例是**完全相同的情形**,而后者的正确性已经有测试钉死了
121    /// (`两个实例并发不会重复推进` / `redis_多实例并发不重复推进`)。
122    ///
123    /// 换句话说:这里没有引入新的竞态,只是把「多实例才能用上的并行」
124    /// 在单进程内也用上。
125    ///
126    /// # 为什么不是把 process() 内部并行
127    ///
128    /// 一笔事务内部的分支**必须按序**(SAGA 就是顺序语义),并行只能跨事务。
129    ///
130    /// # ⚠ 为什么必须用 JoinSet 而不是 Vec<JoinHandle>
131    ///
132    /// 调用方是 `tokio::spawn(driver.run_forever(..))`,靠 **abort 这个外层
133    /// 任务**来停推进器(`Embedded` 的 Drop 就是这么干的)。
134    /// `tokio::spawn` 出来的子任务是**游离的**:外层被 abort 掉,它们照跑不误。
135    ///
136    /// 这个坑实测撞出来过:`跨进程重启_事务不丢且已完成的步骤不重做` 里
137    /// 第一个「进程」析构后,它那些僵尸 worker 还在抢同一笔事务,
138    /// 而它们的 handler 永远返回 Unknown —— 于是第二个「进程」怎么等都推不完。
139    ///
140    /// `JoinSet` 被 drop 时会把里面所有任务一并 abort,正好是我们要的语义。
141    pub async fn run_forever(self, tick: Duration) {
142        let mut set = tokio::task::JoinSet::new();
143        for _ in 0..self.workers.max(1) {
144            let d = self.clone();
145            set.spawn(async move { d.worker_loop(tick).await });
146        }
147        // 任一 worker 意外退出就整体结束 —— 静默少几个 worker 比直接挂更难查
148        set.join_next().await;
149    }
150
151    /// 单个 worker:抢一个到期事务推一下,没活就睡
152    async fn worker_loop(&self, tick: Duration) {
153        loop {
154            match self.store.lock_one_due(&self.owner, self.lease).await {
155                Ok(Some(g)) => {
156                    if let Err(e) = self.process(&g).await {
157                        warn!(gid = %g.gid, error = %e, "推进出错,等下轮重试");
158                    }
159                }
160                Ok(None) => tokio::time::sleep(tick).await,
161                Err(e) => {
162                    warn!(error = %e, "取待办失败");
163                    tokio::time::sleep(tick).await;
164                }
165            }
166        }
167    }
168
169    /// 推进一个全局事务,直到它落终态或需要等待。
170    ///
171    /// 可以被重复调用(崩溃恢复就靠这个)—— 分支的幂等由客户端屏障保证。
172    pub async fn process(&self, g: &GlobalRow) -> anyhow::Result<()> {
173        match g.trans_type {
174            TransType::Saga => self.process_saga(g).await,
175            TransType::Tcc => self.process_tcc(g).await,
176            TransType::Msg => self.process_msg(g).await,
177            TransType::Xa => self.process_xa(g).await,
178            TransType::Workflow => self.process_workflow(g).await,
179        }
180    }
181
182    // ---------------- workflow ----------------
183
184    /// 跑用户的 workflow 函数,失败则逆序补偿它**已经登记过**的分支。
185    ///
186    /// 跟另外四种模式的差别见 [`dtmrs_core::workflow_advance`]:正向走向由
187    /// 用户函数决定,状态机只管「什么时候跑、什么时候补、按什么顺序补」。
188    async fn process_workflow(&self, g: &GlobalRow) -> anyhow::Result<()> {
189        let (name, input) = crate::workflow::decode_payload(&g.payload);
190        let mut status = g.status;
191
192        loop {
193            let rows = self.store.list_branches(&g.gid).await?;
194            let compensates = compensate_states(&rows);
195
196            match dtmrs_core::workflow_advance(status, &compensates) {
197                Advance::Finish(s) => {
198                    info!(gid = %g.gid, status = s.as_str(), "workflow 事务终结");
199                    self.store
200                        .set_global_status(&g.gid, s, g.trans_type, "")
201                        .await?;
202                    return Ok(());
203                }
204                Advance::Wait => return Ok(()),
205
206                Advance::RunWorkflow => {
207                    let Some(f) = self.workflows.get(&name) else {
208                        // 漏注册(新版本删了 workflow / 换了名字)。
209                        // **按结果未知处理** —— 这是部署问题,改回来重试就好,
210                        // 判失败会白白触发回滚
211                        warn!(gid = %g.gid, workflow = %name,
212                              "workflow 未注册,按结果未知处理(会重试,不回滚)");
213                        self.retry_later(g).await?;
214                        return Ok(());
215                    };
216                    let ctx =
217                        crate::workflow::WorkflowCtx::new(&g.gid, &input, self.store.clone(), rows);
218                    match f(ctx).await {
219                        Ok(()) => {
220                            info!(gid = %g.gid, workflow = %name, "workflow 跑完");
221                            self.store
222                                .set_global_status(&g.gid, GlobalStatus::Succeed, g.trans_type, "")
223                                .await?;
224                            return Ok(());
225                        }
226                        Err(crate::workflow::WorkflowError::Rollback(reason)) => {
227                            info!(gid = %g.gid, workflow = %name, %reason, "workflow 要求回滚");
228                            self.store
229                                .set_global_status(
230                                    &g.gid,
231                                    GlobalStatus::Aborting,
232                                    g.trans_type,
233                                    &reason,
234                                )
235                                .await?;
236                            status = GlobalStatus::Aborting;
237                            continue;
238                        }
239                        Err(crate::workflow::WorkflowError::Diverged {
240                            branch_id: bid,
241                            recorded,
242                            got,
243                        }) => {
244                            // **绝不能继续**:按位置记忆化会张冠李戴,回滚也会补错对象。
245                            // 也不回滚 —— 我们已经不知道真实进度了,硬回滚更危险。
246                            // 停在这里等人:改回确定性的代码,重启就能接着推。
247                            warn!(gid = %g.gid, workflow = %name, branch = %bid,
248                                  %recorded, %got,
249                                  "workflow 重放走岔了,已停止推进,需要人工介入");
250                            self.retry_later(g).await?;
251                            return Ok(());
252                        }
253                        Err(e) => {
254                            // Retry / Internal:只重试,绝不回滚
255                            warn!(gid = %g.gid, workflow = %name, error = %e, "workflow 需要重试");
256                            self.retry_later(g).await?;
257                            return Ok(());
258                        }
259                    }
260                }
261
262                Advance::Call { index, op } => {
263                    let bid = branch_id(index);
264                    let Some(url) = url_of(&rows, &bid, op) else {
265                        // 补偿行不见了,不该发生
266                        warn!(gid = %g.gid, branch = %bid, "workflow 补偿地址缺失");
267                        self.retry_later(g).await?;
268                        return Ok(());
269                    };
270                    let bp = payload_of(&rows, &bid, op);
271                    match self.call_branch(g, &bid, op, &url, &bp).await {
272                        BranchResult::Success => {
273                            self.store
274                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Succeed)
275                                .await?;
276                        }
277                        // 补偿失败只能不停重试 —— 漏掉就是真的漏了副作用
278                        _ => {
279                            warn!(gid = %g.gid, branch = %bid, "workflow 补偿未成功,会重试");
280                            self.retry_later(g).await?;
281                            return Ok(());
282                        }
283                    }
284                }
285            }
286        }
287    }
288
289    // ---------------- SAGA ----------------
290
291    async fn process_saga(&self, g: &GlobalRow) -> anyhow::Result<()> {
292        let steps: Vec<SagaStep> = serde_json::from_str(&g.payload).unwrap_or_default();
293        if steps.is_empty() {
294            self.store
295                .set_global_status(&g.gid, GlobalStatus::Succeed, g.trans_type, "")
296                .await?;
297            return Ok(());
298        }
299        let mut status = g.status;
300
301        loop {
302            let (actions, compensates) = self.branch_states(&g.gid, steps.len()).await?;
303            match saga_advance(status, &actions, &compensates) {
304                Advance::Finish(s) => {
305                    if s == GlobalStatus::Aborting {
306                        // 防御性分支:状态机发现有 failed 分支但全局还没转 aborting
307                        status = s;
308                        self.store
309                            .set_global_status(&g.gid, s, g.trans_type, "分支已判失败")
310                            .await?;
311                        continue;
312                    }
313                    info!(gid = %g.gid, status = s.as_str(), "事务终结");
314                    self.store
315                        .set_global_status(&g.gid, s, g.trans_type, "")
316                        .await?;
317                    return Ok(());
318                }
319                Advance::Wait => return Ok(()),
320                // 只有 workflow 模式会出现,别的模式走到这里说明状态机接错了。
321                // **不 panic** —— 推进器是常驻的,崩了整个 TC 就停了
322                Advance::RunWorkflow => {
323                    warn!(gid = %g.gid, "非 workflow 事务收到 RunWorkflow 决策,跳过");
324                    return Ok(());
325                }
326                Advance::Call { index, op } => {
327                    let branch_id = branch_id(index);
328                    let url = match op {
329                        BranchOp::Action => &steps[index].action,
330                        _ => &steps[index].compensate,
331                    };
332                    match self
333                        .call_branch(g, &branch_id, op, url, &steps[index].payload)
334                        .await
335                    {
336                        BranchResult::Success => {
337                            self.store
338                                .set_branch_status(&g.gid, &branch_id, op, BranchStatus::Succeed)
339                                .await?;
340                        }
341                        BranchResult::Failure => {
342                            self.store
343                                .set_branch_status(&g.gid, &branch_id, op, BranchStatus::Failed)
344                                .await?;
345                            if op == BranchOp::Action {
346                                // 只有业务**明确**说失败才回滚
347                                info!(gid = %g.gid, branch = %branch_id, "分支要求回滚");
348                                status = GlobalStatus::Aborting;
349                                self.store
350                                    .set_global_status(
351                                        &g.gid,
352                                        GlobalStatus::Aborting,
353                                        g.trans_type,
354                                        &format!("分支 {branch_id} 返回 FAILURE"),
355                                    )
356                                    .await?;
357                            } else {
358                                // 补偿都失败了,只能不停重试 —— 这时候需要人介入
359                                warn!(gid = %g.gid, branch = %branch_id, "补偿失败,需要人工介入");
360                                self.retry_later(g).await?;
361                                return Ok(());
362                            }
363                        }
364                        BranchResult::Ongoing | BranchResult::Unknown => {
365                            // **绝不能当成失败**:对方可能已经成功了。退避重试。
366                            self.retry_later(g).await?;
367                            return Ok(());
368                        }
369                    }
370                }
371            }
372        }
373    }
374
375    // ---------------- TCC ----------------
376
377    /// TCC 的 try 阶段是**客户端驱动**的(客户端先 registerBranch 再调 try),
378    /// TC 只负责 confirm / cancel。所以分支的 URL 来自 `trans_branch_op` 表,
379    /// 不是全局 payload。
380    async fn process_tcc(&self, g: &GlobalRow) -> anyhow::Result<()> {
381        self.drive_two_phase(g, BranchOp::Confirm, BranchOp::Cancel, "TCC")
382            .await
383    }
384
385    // ---------------- XA ----------------
386
387    /// XA 的一阶段(业务 SQL + `PREPARE TRANSACTION`)由**客户端**做,
388    /// TC 只负责统一决定 `COMMIT PREPARED` 还是 `ROLLBACK PREPARED`。
389    ///
390    /// 形状跟 TCC 一样,只是 op 换成 commit/rollback。语义上那条铁律也一样:
391    /// **commit 失败绝不能转 rollback** —— 别的分支可能已经提交了。
392    ///
393    /// XA 独有的严重性:没解决的 prepared 事务会**永久持锁**,在 Postgres 里
394    /// 还阻塞 VACUUM。所以这里的重试比 SAGA 的补偿重试要紧得多。
395    async fn process_xa(&self, g: &GlobalRow) -> anyhow::Result<()> {
396        self.drive_two_phase(g, BranchOp::Commit, BranchOp::Rollback, "XA")
397            .await
398    }
399
400    /// TCC 和 XA 共用的二阶段推进:正向 op 全做完就成功,反向 op 全做完就失败,
401    /// **任一方向的失败都只重试,绝不改变方向**。
402    async fn drive_two_phase(
403        &self,
404        g: &GlobalRow,
405        fwd: BranchOp,
406        bwd: BranchOp,
407        label: &str,
408    ) -> anyhow::Result<()> {
409        let rows = self.store.list_branches(&g.gid).await?;
410        let n = rows
411            .iter()
412            .filter_map(|r| index_of(&r.branch_id))
413            .max()
414            .map(|m| m + 1)
415            .unwrap_or(0);
416        if n == 0 {
417            // ⚠ `n == 0` 有两种成因,**结论完全相反,绝不能混为一谈**:
418            //
419            //   a) rows 是空的 —— 真的一个分支都没登记就 submit/abort 了,空事务,落终态;
420            //   b) rows 不空,但分支号一个都解析不出下标 —— 这时候落终态是**灾难**:
421            //      TCC 会被判成 succeed 而 confirm 一次都没调,那份已经 try 冻结的资源
422            //      永远没人收尾。实测过:登记 branch_id="inventory" 再 submit,
423            //      事务直接变 succeed,分支停在 prepared。
424            //
425            // 现在 register_branch 会挡住 (b),但库里可能存着更早写进去的坏行,
426            // 所以这里也得挡。按项目的一贯原则处理:**既不判成功也不回滚,停下等人。**
427            if !rows.is_empty() {
428                error!(
429                    gid = %g.gid,
430                    branches = rows.len(),
431                    ids = ?rows.iter().map(|r| r.branch_id.as_str()).take(5).collect::<Vec<_>>(),
432                    "分支号全都无法解析成下标,无法推进。这笔事务需要人工介入 —— \
433                     合法的分支号形如 01 / 02 / 100(见 is_canonical_branch_id)"
434                );
435                return Ok(());
436            }
437            let s = if g.status == GlobalStatus::Aborting {
438                GlobalStatus::Failed
439            } else {
440                GlobalStatus::Succeed
441            };
442            self.store
443                .set_global_status(&g.gid, s, g.trans_type, "")
444                .await?;
445            return Ok(());
446        }
447
448        let status = g.status;
449        loop {
450            let rows = self.store.list_branches(&g.gid).await?;
451            let (f, b) = split_by_op(&rows, n, fwd, bwd);
452            let adv = if fwd == BranchOp::Commit {
453                xa_advance(status, &f, &b)
454            } else {
455                tcc_advance(status, &f, &b)
456            };
457            match adv {
458                Advance::Finish(s) => {
459                    info!(gid = %g.gid, status = s.as_str(), mode = label, "事务终结");
460                    self.store
461                        .set_global_status(&g.gid, s, g.trans_type, "")
462                        .await?;
463                    return Ok(());
464                }
465                Advance::Wait => return Ok(()),
466                // 只有 workflow 模式会出现,别的模式走到这里说明状态机接错了。
467                // **不 panic** —— 推进器是常驻的,崩了整个 TC 就停了
468                Advance::RunWorkflow => {
469                    warn!(gid = %g.gid, "非 workflow 事务收到 RunWorkflow 决策,跳过");
470                    return Ok(());
471                }
472                Advance::Call { index, op } => {
473                    let bid = branch_id(index);
474                    let Some(url) = url_of(&rows, &bid, op) else {
475                        warn!(gid = %g.gid, branch = %bid, op = op.as_str(), mode = label,
476                              "分支没登记这个操作的 URL,无法调用");
477                        self.retry_later(g).await?;
478                        return Ok(());
479                    };
480                    let bp = payload_of(&rows, &bid, op);
481                    match self.call_branch(g, &bid, op, &url, &bp).await {
482                        BranchResult::Success => {
483                            self.store
484                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Succeed)
485                                .await?;
486                        }
487                        // 二阶段失败**绝不改变全局方向**:一阶段已经成功、
488                        // 方向已经定了,反向操作会造成一半提交一半回滚。
489                        // 唯一正确处理是无限重试 + 报警。
490                        BranchResult::Failure => {
491                            self.store
492                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Failed)
493                                .await?;
494                            warn!(gid = %g.gid, branch = %bid, op = op.as_str(), mode = label,
495                                  "二阶段失败,会持续重试,需要人工介入");
496                            self.retry_later(g).await?;
497                            return Ok(());
498                        }
499                        BranchResult::Ongoing | BranchResult::Unknown => {
500                            self.retry_later(g).await?;
501                            return Ok(());
502                        }
503                    }
504                }
505            }
506        }
507    }
508
509    // ---------------- 二阶段消息 ----------------
510
511    /// 流程:`prepare` 落库 → 业务提交本地事务 → `submit`。
512    ///
513    /// 如果进程在这两步之间崩了,事务会一直停在 prepared。这时 TC 靠回查
514    /// `query_prepared` 问业务方"你那个本地事务到底提交了没有",据此决定
515    /// 是往前推还是整单作废。**这是取代 MQ 事务消息的关键一环。**
516    async fn process_msg(&self, g: &GlobalRow) -> anyhow::Result<()> {
517        let mut status = g.status;
518
519        if status == GlobalStatus::Prepared {
520            if g.query_prepared.is_empty() {
521                // 没给回查地址就没法自动决断。不能瞎猜 —— 猜错要么丢单要么重复扣款
522                warn!(gid = %g.gid, "msg 事务没提供 query_prepared,无法回查,等人处理");
523                self.retry_later(g).await?;
524                return Ok(());
525            }
526            // 借用分支调用的通道做回查,branch_id 用 "00" 跟真实分支区分开
527            match self
528                .call_branch(g, "00", BranchOp::Action, &g.query_prepared, "")
529                .await
530            {
531                BranchResult::Success => {
532                    info!(gid = %g.gid, "回查:本地事务已提交 → 继续推进");
533                    self.store
534                        .set_global_status(&g.gid, GlobalStatus::Submitted, g.trans_type, "")
535                        .await?;
536                    status = GlobalStatus::Submitted;
537                }
538                BranchResult::Failure => {
539                    // 业务方明确说"这单没提交" → 整单作废。msg 没有补偿分支,
540                    // 但也不需要:正向分支压根还没跑过
541                    info!(gid = %g.gid, "回查:本地事务未提交 → 整单作废");
542                    self.store
543                        .set_global_status(
544                            &g.gid,
545                            GlobalStatus::Failed,
546                            g.trans_type,
547                            "回查得到 FAILURE:本地事务未提交",
548                        )
549                        .await?;
550                    return Ok(());
551                }
552                BranchResult::Ongoing | BranchResult::Unknown => {
553                    // 回查本身失败了,不能当作"没提交"。退避重试。
554                    self.retry_later(g).await?;
555                    return Ok(());
556                }
557            }
558        }
559
560        let steps: Vec<SagaStep> = serde_json::from_str(&g.payload).unwrap_or_default();
561        if steps.is_empty() {
562            self.store
563                .set_global_status(&g.gid, GlobalStatus::Succeed, g.trans_type, "")
564                .await?;
565            return Ok(());
566        }
567        loop {
568            let (actions, _) = self.branch_states(&g.gid, steps.len()).await?;
569            match msg_advance(status, &actions) {
570                Advance::Finish(s) => {
571                    info!(gid = %g.gid, status = s.as_str(), "消息事务终结");
572                    self.store
573                        .set_global_status(&g.gid, s, g.trans_type, "")
574                        .await?;
575                    return Ok(());
576                }
577                Advance::Wait => return Ok(()),
578                // 只有 workflow 模式会出现,别的模式走到这里说明状态机接错了。
579                // **不 panic** —— 推进器是常驻的,崩了整个 TC 就停了
580                Advance::RunWorkflow => {
581                    warn!(gid = %g.gid, "非 workflow 事务收到 RunWorkflow 决策,跳过");
582                    return Ok(());
583                }
584                Advance::Call { index, op } => {
585                    let bid = branch_id(index);
586                    match self
587                        .call_branch(g, &bid, op, &steps[index].action, &steps[index].payload)
588                        .await
589                    {
590                        BranchResult::Success => {
591                            self.store
592                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Succeed)
593                                .await?;
594                        }
595                        // msg 保证"最终一定送达",没有补偿一说。失败只能重试。
596                        BranchResult::Failure | BranchResult::Ongoing | BranchResult::Unknown => {
597                            self.retry_later(g).await?;
598                            return Ok(());
599                        }
600                    }
601                }
602            }
603        }
604    }
605
606    async fn retry_later(&self, g: &GlobalRow) -> anyhow::Result<()> {
607        let iv = dtmrs_core::next_interval_with(g.next_cron_interval, self.retry);
608        self.store.schedule_retry(&g.gid, iv).await?;
609        Ok(())
610    }
611
612    /// 取每一步的 action / compensate 分支当前状态,按步序对齐
613    async fn branch_states(
614        &self,
615        gid: &str,
616        n: usize,
617    ) -> anyhow::Result<(Vec<BranchStatus>, Vec<BranchStatus>)> {
618        let rows = self.store.list_branches(gid).await?;
619        let mut actions = vec![BranchStatus::Prepared; n];
620        let mut compensates = vec![BranchStatus::Prepared; n];
621        for r in rows {
622            let Some(i) = index_of(&r.branch_id) else {
623                continue;
624            };
625            if i >= n {
626                continue;
627            }
628            match r.op {
629                BranchOp::Action | BranchOp::Try => actions[i] = r.status,
630                BranchOp::Compensate | BranchOp::Cancel | BranchOp::Rollback => {
631                    compensates[i] = r.status
632                }
633                _ => {}
634            }
635        }
636        Ok((actions, compensates))
637    }
638
639    /// 调一个分支。`local://名字` 走进程内函数,`grpc://` 走 gRPC,其它走 HTTP。
640    async fn call_branch(
641        &self,
642        g: &GlobalRow,
643        branch_id: &str,
644        op: BranchOp,
645        url: &str,
646        payload: &str,
647    ) -> BranchResult {
648        match parse_target(url) {
649            Target::Local(name) => self.call_local(g, branch_id, op, &name).await,
650            Target::Http(u) => self.call_http(g, branch_id, op, &u, payload).await,
651            #[cfg(feature = "grpc")]
652            Target::Grpc(t) => {
653                self.grpc
654                    .call(
655                        &t,
656                        &g.gid,
657                        &g.trans_type.to_string(),
658                        branch_id,
659                        op.as_str(),
660                    )
661                    .await
662            }
663            // 编译时关掉了 grpc feature,却遇到 grpc:// 分支。
664            // 按「结果未知」处理(重试,不回滚)—— 这是构建配置问题,不是业务失败
665            #[cfg(not(feature = "grpc"))]
666            Target::Grpc(t) => {
667                warn!(gid = %g.gid, branch = %branch_id, endpoint = %t.endpoint,
668                      "遇到 grpc:// 分支但本次构建关掉了 grpc feature,按结果未知处理");
669                BranchResult::Unknown
670            }
671        }
672    }
673
674    /// 进程内调用:没有网络、没有序列化,一次函数调用。
675    async fn call_local(
676        &self,
677        g: &GlobalRow,
678        branch_id: &str,
679        op: BranchOp,
680        name: &str,
681    ) -> BranchResult {
682        let Some(h) = self.registry.get(name) else {
683            // 漏注册(比如新版本删了 handler)。**必须当 Unknown 而不是 Failure**:
684            // 判失败会触发回滚,而这其实是部署问题,改回来重试才对。
685            warn!(gid = %g.gid, branch = %branch_id, handler = name,
686                  "本地分支未注册,按结果未知处理(会重试,不回滚)");
687            return BranchResult::Unknown;
688        };
689        let ctx = BranchCtx {
690            gid: g.gid.clone(),
691            branch_id: branch_id.to_string(),
692            op,
693            trans_type: g.trans_type.to_string(),
694        };
695        let r = h(ctx).await;
696        info!(gid = %g.gid, branch = %branch_id, op = op.as_str(),
697              handler = name, result = ?r, "本地分支返回");
698        r
699    }
700
701    /// 远端调用。查询参数跟 DTM 保持一致,客户端屏障库可以直接复用。
702    async fn call_http(
703        &self,
704        g: &GlobalRow,
705        branch_id: &str,
706        op: BranchOp,
707        url: &str,
708        payload: &str,
709    ) -> BranchResult {
710        let req = self
711            .http
712            .post(url)
713            .query(&[
714                ("gid", g.gid.as_str()),
715                ("trans_type", &g.trans_type.to_string()),
716                ("branch_id", branch_id),
717                ("op", op.as_str()),
718            ])
719            .header("content-type", "application/json")
720            .body(branch_payload(payload));
721        match req.send().await {
722            Ok(resp) => {
723                let code = resp.status().as_u16();
724                let body = resp.text().await.unwrap_or_default();
725                let r = BranchResult::from_http(code, &body);
726                info!(gid = %g.gid, branch = %branch_id, op = op.as_str(), code, result = ?r, "分支返回");
727                r
728            }
729            Err(e) => {
730                // 超时/连不上 —— 结果未知,必须重试而不是回滚
731                warn!(gid = %g.gid, branch = %branch_id, error = %e, "分支不可达,结果未知");
732                BranchResult::Unknown
733            }
734        }
735    }
736}
737
738/// 分支号:第 0 步是 "01",跟 DTM 一致。
739///
740/// ⚠ **补零只有两位,第 100 个分支起会变成三位** —— 别据此以为顺序会乱。
741/// 执行顺序不是按 `branch_id` 的字符串序决定的:`split_by_op` /
742/// `compensate_states` 都是用 `index_of()` 解析出的**整数下标**写进数组,
743/// 再由 core 的状态机按下标推进。`ORDER BY branch_id` 只影响 `/query`
744/// 的输出排版(105 个分支时列表里 "100" 会显示在 "99" 前面,仅此而已)。
745///
746/// 宽度**不能随便改宽**:TCC / XA 的分支是客户端自己登记 id 的,而 driver
747/// 在推进时用 `branch_id(index)` 反查那一行。改成 `{:04}` 会让所有登记
748/// `"01"` 的现有客户端对不上行 —— 试过,TCC 的三个用例当场挂掉。
749pub fn branch_id(index: usize) -> String {
750    format!("{:02}", index + 1)
751}
752
753/// 分支下标的上限。
754///
755/// ⚠ 这不是「够用就行」的产品限制,是**内存安全的护栏**。
756/// `n` 是拿所有分支行里的最大下标算出来的,紧接着就 `vec![...; n]`。
757/// 分支号是 TCC / XA 客户端自己给的字符串 —— 登记一个 `"2000000000"`,
758/// TC 就会去分配二十亿个元素。实测过:一次 submit 把 RSS 从 38 MB
759/// 顶到 **3.4 GB**,而且那行留在库里,推进器每次轮询都再来一遍。
760///
761/// SAGA 那边有 payload 的 8192 字符上限兜着(最多装几百步),
762/// TCC / XA 的分支是一条条登记的,没有任何天然上限,只能在这里挡。
763pub const MAX_BRANCH_INDEX: usize = 9999;
764
765/// 把分支号解析回下标。**认不出来就返回 None,调用方必须当心 —— 见 `n == 0` 那段。**
766///
767/// 超过 `MAX_BRANCH_INDEX` 一律不认:库里可能已经存着历史坏数据
768/// (早先 `register_branch` 只校验非空),不在这里挡住的话,
769/// 光是把它捞起来推进就能把 TC 撑爆。
770fn index_of(branch_id: &str) -> Option<usize> {
771    branch_id
772        .parse::<usize>()
773        .ok()
774        .and_then(|v| v.checked_sub(1))
775        .filter(|&i| i <= MAX_BRANCH_INDEX)
776}
777
778/// 分支号是否是 driver 能**原样还原**的形式。
779///
780/// 判据就一句:`branch_id(index_of(s)) == s`。这不是洁癖 ——
781/// driver 推进时并不用库里存的那个字符串,而是拿下标**重新生成**
782/// (`Advance::Call { index }` → `branch_id(index)` → 反查那一行)。
783/// 所以只要还原不出原样,那一行就永远对不上:
784///
785/// | 客户端登记 | 会发生什么 |
786/// |---|---|
787/// | `"01"` / `"100"` | ✅ 正常 |
788/// | `"1"` / `"001"` | 存进去是 `"1"`,driver 找的是 `"01"` —— 状态更新全部落空,无限重试 |
789/// | `"inventory"` | 解析不出下标 → 事务被当成**空事务直接判成功**,confirm 一次都不会调,资源永久泄漏 |
790/// | `"2000000000"` | 推进时按下标开数组,一次 submit 吃掉 3.4 GB |
791///
792/// 后两行都是实测出来的,不是推演。
793pub fn is_canonical_branch_id(s: &str) -> bool {
794    index_of(s).map(branch_id).as_deref() == Some(s)
795}
796
797/// 按 op 把分支行拆成两列(正向 / 反向),按步序对齐
798fn split_by_op(
799    rows: &[dtmrs_store::BranchRow],
800    n: usize,
801    fwd: BranchOp,
802    bwd: BranchOp,
803) -> (Vec<BranchStatus>, Vec<BranchStatus>) {
804    let mut a = vec![BranchStatus::Prepared; n];
805    let mut b = vec![BranchStatus::Prepared; n];
806    for r in rows {
807        let Some(i) = index_of(&r.branch_id) else {
808            continue;
809        };
810        if i >= n {
811            continue;
812        }
813        if r.op == fwd {
814            a[i] = r.status;
815        } else if r.op == bwd {
816            b[i] = r.status;
817        }
818    }
819    (a, b)
820}
821
822/// 按分支序取出各 compensate 行的状态,用于 workflow 的逆序补偿。
823///
824/// 跟 `split_by_op` 的区别:workflow 的分支数是**运行时长出来的**,
825/// 提交时并不知道有几个,所以长度得从已登记的行里推出来。
826fn compensate_states(rows: &[dtmrs_store::BranchRow]) -> Vec<BranchStatus> {
827    let n = rows
828        .iter()
829        .filter(|r| r.op == BranchOp::Compensate)
830        .filter_map(|r| index_of(&r.branch_id))
831        .max()
832        .map(|m| m + 1)
833        .unwrap_or(0);
834    // 没登记补偿的分支(比如纯查询步骤)留成 Succeed,逆序扫描时会跳过它 ——
835    // 本来就没有副作用要收拾
836    let mut v = vec![BranchStatus::Succeed; n];
837    for r in rows.iter().filter(|r| r.op == BranchOp::Compensate) {
838        if let Some(i) = index_of(&r.branch_id) {
839            if i < n {
840                v[i] = r.status;
841            }
842        }
843    }
844    v
845}
846
847/// 取某个分支行上存的 payload(TCC / XA / workflow 的分支是动态登记的,
848/// 业务数据跟着行走,不在全局 payload 里)
849fn payload_of(rows: &[dtmrs_store::BranchRow], branch_id: &str, op: BranchOp) -> String {
850    rows.iter()
851        .find(|r| r.branch_id == branch_id && r.op == op)
852        .map(|r| r.payload.clone())
853        .unwrap_or_default()
854}
855
856fn url_of(rows: &[dtmrs_store::BranchRow], branch_id: &str, op: BranchOp) -> Option<String> {
857    rows.iter()
858        .find(|r| r.branch_id == branch_id && r.op == op)
859        .map(|r| r.url.clone())
860}
861
862/// 这一步要发给分支的请求体。
863///
864/// 每步各自独立 —— 扣款那步要金额、发货那步要地址,本来就不该收到同一份数据。
865/// 步骤没写 payload 就发 `{}`(很多分支只靠 gid/branch_id/op 做幂等,不需要请求体)。
866fn branch_payload(step_payload: &str) -> String {
867    if step_payload.trim().is_empty() {
868        "{}".to_string()
869    } else {
870        step_payload.to_string()
871    }
872}
873
874/// 推进器的可配置项
875#[derive(Debug, Clone, Copy)]
876pub struct DriverConfig {
877    /// 调一个分支最多等几秒
878    pub branch_timeout_secs: i64,
879    /// 租约时长(秒)
880    pub lease_secs: i64,
881    pub retry: dtmrs_core::RetryPolicy,
882    /// 并行推进的 worker 数。**一笔事务内部仍然按序**,并行只发生在事务之间
883    pub workers: usize,
884}
885
886impl Default for DriverConfig {
887    fn default() -> Self {
888        // 跟 0.2 的写死值一致,不配置的人行为不变
889        Self {
890            branch_timeout_secs: 10,
891            lease_secs: 30,
892            retry: dtmrs_core::RetryPolicy::default(),
893            // 推一笔事务的时间几乎全花在等 I/O(存储往返 + 分支调用)上,
894            // 所以 worker 数可以明显高于核数。
895            //
896            // 16 基本就是这台机器上的天花板了:实测(20 核,空库,
897            //   三次取中位数,bench/)
898            //   Postgres  1→267 笔/秒  16→3424  64→3184(一样,没收益)
899            //   Redis     1→965        16→4695  32→4974
900            //   sqlite    1→435        16→682 —— 写是全库串行的,并行收益有限
901            //   MySQL     1→19         16→129   32→202
902            // 再往上加 worker 不涨,连接池也不是瓶颈。要继续提升得**减少
903            // 每笔事务的存储往返次数**,不是加并发。
904            // 另外 worker 开多少就要占多少条数据库连接,而 TC 常常和业务
905            // 共用一个库,所以默认值往保守取
906            workers: 16,
907        }
908    }
909}
910
911impl DriverConfig {
912    pub fn from_env() -> Self {
913        let d = Self::default();
914        let get = |k: &str, fallback: i64| {
915            std::env::var(k)
916                .ok()
917                .and_then(|v| v.parse::<i64>().ok())
918                .filter(|v| *v > 0)
919                .unwrap_or(fallback)
920        };
921        Self {
922            branch_timeout_secs: get("DTMRS_BRANCH_TIMEOUT", d.branch_timeout_secs),
923            lease_secs: get("DTMRS_LEASE", d.lease_secs),
924            retry: dtmrs_core::RetryPolicy::from_env(),
925            workers: get("DTMRS_WORKERS", d.workers as i64) as usize,
926        }
927    }
928}
929
930#[cfg(test)]
931mod tests {
932    use super::*;
933
934    #[test]
935    fn 分支号与下标互转() {
936        assert_eq!(branch_id(0), "01");
937        assert_eq!(branch_id(9), "10");
938        assert_eq!(index_of("01"), Some(0));
939        assert_eq!(index_of("10"), Some(9));
940        assert_eq!(index_of("00"), None);
941        assert_eq!(index_of("xx"), None);
942    }
943
944    /// 超过 99 个分支时分支号会变成三位(`"100"`),**字符串序就跟数值序分家了**
945    /// (`"100" < "99"`)。这条测试钉的是:即便如此,**下标解析仍然正确** ——
946    /// 因为执行顺序走的是 `index_of()` 解析出的整数,不是字符串排序。
947    ///
948    /// 曾经因为看到 `/query` 的输出乱序就误判成执行顺序 bug,实际只是列表排版。
949    /// 真要改这里的补零宽度前先读 `branch_id()` 的注释(会破坏 TCC 客户端)。
950    #[test]
951    fn 分支号超过99后下标解析仍然正确() {
952        let ids: Vec<String> = (0..500).map(branch_id).collect();
953        for (i, id) in ids.iter().enumerate() {
954            assert_eq!(index_of(id), Some(i), "下标解析错了,执行顺序会乱");
955        }
956        // 字符串序确实不等于数值序 —— 这是已知且无害的,别改这条断言去"修"它
957        let mut sorted = ids.clone();
958        sorted.sort();
959        assert_ne!(ids, sorted);
960    }
961}