dtmrs-server 0.2.0

Distributed transaction coordinator: SAGA / TCC / two-phase messaging / XA / workflow, over HTTP and gRPC, embeddable as a library
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
478
479
480
481
482
483
484
485
486
487
488
489
490
491
492
493
494
495
496
497
498
499
500
501
502
503
504
505
506
507
508
509
510
511
512
513
514
515
516
517
518
519
520
521
522
523
524
525
526
527
528
529
530
531
532
533
534
535
536
537
538
539
540
541
542
543
544
545
546
547
548
549
550
551
552
553
554
555
556
557
558
559
560
561
562
563
564
565
566
567
568
569
570
571
572
573
574
575
576
577
578
579
580
581
582
583
584
585
586
587
588
589
590
591
592
593
594
595
596
597
598
599
600
601
602
603
604
605
606
607
608
609
610
611
612
613
614
615
616
617
618
619
620
621
622
623
624
625
626
627
628
629
630
631
632
633
634
635
636
637
638
639
640
641
642
643
644
645
646
647
648
649
650
651
652
653
654
655
656
657
658
659
660
661
662
663
664
665
666
667
668
669
670
671
672
673
674
675
676
677
678
679
680
681
682
683
684
685
686
687
688
689
//! 事务推进器:把状态机的决策落成真实的 HTTP 调用和状态更新。
//!
//! 这里是**唯一**会修改全局事务状态的地方,决策全部来自 `dtmrs_core::saga_advance`,
//! 本文件只负责 I/O。这样状态迁移的正确性可以在 core 里纯单测覆盖。

use crate::registry::{parse_target, BranchCtx, Registry, Target};
use dtmrs_core::{
    msg_advance, saga_advance, tcc_advance, xa_advance, Advance, BranchOp, BranchResult,
    BranchStatus, GlobalStatus, SagaStep, TransType,
};
use dtmrs_store::{GlobalRow, Store};
use std::sync::Arc;
use std::time::Duration;
use tracing::{info, warn};

#[derive(Clone)]
pub struct Driver {
    pub store: Store,
    pub http: reqwest::Client,
    pub owner: String,
    /// 租约时长(秒)。持租约的实例崩了,这么久之后别的实例接手
    pub lease: i64,
    /// 进程内分支注册表。嵌入式模式用,纯 HTTP 部署时是空表
    pub registry: Arc<Registry>,
    /// gRPC 分支调用器(带 channel 缓存)
    #[cfg(feature = "grpc")]
    pub grpc: crate::grpc::client::GrpcCaller,
    /// workflow 函数注册表。跟 registry 一样,纯 HTTP 部署时是空表
    pub workflows: Arc<crate::workflow::WorkflowRegistry>,
}

impl Driver {
    pub fn new(store: Store, owner: String) -> Self {
        Self {
            store,
            http: reqwest::Client::builder()
                .timeout(Duration::from_secs(10))
                .build()
                .expect("build http client"),
            owner,
            lease: 30,
            registry: Arc::new(Registry::new()),
            #[cfg(feature = "grpc")]
            grpc: crate::grpc::client::GrpcCaller::new(Duration::from_secs(10)),
            workflows: Arc::new(crate::workflow::WorkflowRegistry::new()),
        }
    }

    /// 挂上进程内分支注册表 —— 嵌入式模式的入口
    pub fn with_registry(mut self, r: Arc<Registry>) -> Self {
        self.registry = r;
        self
    }

    /// 挂上 workflow 函数注册表
    pub fn with_workflows(mut self, w: Arc<crate::workflow::WorkflowRegistry>) -> Self {
        self.workflows = w;
        self
    }

    /// 常驻循环:抢一个到期事务推一下,没活就睡
    pub async fn run_forever(self, tick: Duration) {
        loop {
            match self.store.lock_one_due(&self.owner, self.lease).await {
                Ok(Some(g)) => {
                    if let Err(e) = self.process(&g).await {
                        warn!(gid = %g.gid, error = %e, "推进出错,等下轮重试");
                    }
                }
                Ok(None) => tokio::time::sleep(tick).await,
                Err(e) => {
                    warn!(error = %e, "取待办失败");
                    tokio::time::sleep(tick).await;
                }
            }
        }
    }

    /// 推进一个全局事务,直到它落终态或需要等待。
    ///
    /// 可以被重复调用(崩溃恢复就靠这个)—— 分支的幂等由客户端屏障保证。
    pub async fn process(&self, g: &GlobalRow) -> anyhow::Result<()> {
        match g.trans_type {
            TransType::Saga => self.process_saga(g).await,
            TransType::Tcc => self.process_tcc(g).await,
            TransType::Msg => self.process_msg(g).await,
            TransType::Xa => self.process_xa(g).await,
            TransType::Workflow => self.process_workflow(g).await,
        }
    }

    // ---------------- workflow ----------------

    /// 跑用户的 workflow 函数,失败则逆序补偿它**已经登记过**的分支。
    ///
    /// 跟另外四种模式的差别见 [`dtmrs_core::workflow_advance`]:正向走向由
    /// 用户函数决定,状态机只管「什么时候跑、什么时候补、按什么顺序补」。
    async fn process_workflow(&self, g: &GlobalRow) -> anyhow::Result<()> {
        let (name, input) = crate::workflow::decode_payload(&g.payload);
        let mut status = g.status;

        loop {
            let rows = self.store.list_branches(&g.gid).await?;
            let compensates = compensate_states(&rows);

            match dtmrs_core::workflow_advance(status, &compensates) {
                Advance::Finish(s) => {
                    info!(gid = %g.gid, status = s.as_str(), "workflow 事务终结");
                    self.store.set_global_status(&g.gid, s, "").await?;
                    return Ok(());
                }
                Advance::Wait => return Ok(()),

                Advance::RunWorkflow => {
                    let Some(f) = self.workflows.get(&name) else {
                        // 漏注册(新版本删了 workflow / 换了名字)。
                        // **按结果未知处理** —— 这是部署问题,改回来重试就好,
                        // 判失败会白白触发回滚
                        warn!(gid = %g.gid, workflow = %name,
                              "workflow 未注册,按结果未知处理(会重试,不回滚)");
                        self.retry_later(g).await?;
                        return Ok(());
                    };
                    let ctx =
                        crate::workflow::WorkflowCtx::new(&g.gid, &input, self.store.clone(), rows);
                    match f(ctx).await {
                        Ok(()) => {
                            info!(gid = %g.gid, workflow = %name, "workflow 跑完");
                            self.store
                                .set_global_status(&g.gid, GlobalStatus::Succeed, "")
                                .await?;
                            return Ok(());
                        }
                        Err(crate::workflow::WorkflowError::Rollback(reason)) => {
                            info!(gid = %g.gid, workflow = %name, %reason, "workflow 要求回滚");
                            self.store
                                .set_global_status(&g.gid, GlobalStatus::Aborting, &reason)
                                .await?;
                            status = GlobalStatus::Aborting;
                            continue;
                        }
                        Err(crate::workflow::WorkflowError::Diverged {
                            branch_id: bid,
                            recorded,
                            got,
                        }) => {
                            // **绝不能继续**:按位置记忆化会张冠李戴,回滚也会补错对象。
                            // 也不回滚 —— 我们已经不知道真实进度了,硬回滚更危险。
                            // 停在这里等人:改回确定性的代码,重启就能接着推。
                            warn!(gid = %g.gid, workflow = %name, branch = %bid,
                                  %recorded, %got,
                                  "workflow 重放走岔了,已停止推进,需要人工介入");
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                        Err(e) => {
                            // Retry / Internal:只重试,绝不回滚
                            warn!(gid = %g.gid, workflow = %name, error = %e, "workflow 需要重试");
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                    }
                }

                Advance::Call { index, op } => {
                    let bid = branch_id(index);
                    let Some(url) = url_of(&rows, &bid, op) else {
                        // 补偿行不见了,不该发生
                        warn!(gid = %g.gid, branch = %bid, "workflow 补偿地址缺失");
                        self.retry_later(g).await?;
                        return Ok(());
                    };
                    match self.call_branch(g, &bid, op, &url).await {
                        BranchResult::Success => {
                            self.store
                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Succeed)
                                .await?;
                        }
                        // 补偿失败只能不停重试 —— 漏掉就是真的漏了副作用
                        _ => {
                            warn!(gid = %g.gid, branch = %bid, "workflow 补偿未成功,会重试");
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                    }
                }
            }
        }
    }

    // ---------------- SAGA ----------------

    async fn process_saga(&self, g: &GlobalRow) -> anyhow::Result<()> {
        let steps: Vec<SagaStep> = serde_json::from_str(&g.payload).unwrap_or_default();
        if steps.is_empty() {
            self.store
                .set_global_status(&g.gid, GlobalStatus::Succeed, "")
                .await?;
            return Ok(());
        }
        let mut status = g.status;

        loop {
            let (actions, compensates) = self.branch_states(&g.gid, steps.len()).await?;
            match saga_advance(status, &actions, &compensates) {
                Advance::Finish(s) => {
                    if s == GlobalStatus::Aborting {
                        // 防御性分支:状态机发现有 failed 分支但全局还没转 aborting
                        status = s;
                        self.store
                            .set_global_status(&g.gid, s, "分支已判失败")
                            .await?;
                        continue;
                    }
                    info!(gid = %g.gid, status = s.as_str(), "事务终结");
                    self.store.set_global_status(&g.gid, s, "").await?;
                    return Ok(());
                }
                Advance::Wait => return Ok(()),
                // 只有 workflow 模式会出现,别的模式走到这里说明状态机接错了。
                // **不 panic** —— 推进器是常驻的,崩了整个 TC 就停了
                Advance::RunWorkflow => {
                    warn!(gid = %g.gid, "非 workflow 事务收到 RunWorkflow 决策,跳过");
                    return Ok(());
                }
                Advance::Call { index, op } => {
                    let branch_id = branch_id(index);
                    let url = match op {
                        BranchOp::Action => &steps[index].action,
                        _ => &steps[index].compensate,
                    };
                    match self.call_branch(g, &branch_id, op, url).await {
                        BranchResult::Success => {
                            self.store
                                .set_branch_status(&g.gid, &branch_id, op, BranchStatus::Succeed)
                                .await?;
                        }
                        BranchResult::Failure => {
                            self.store
                                .set_branch_status(&g.gid, &branch_id, op, BranchStatus::Failed)
                                .await?;
                            if op == BranchOp::Action {
                                // 只有业务**明确**说失败才回滚
                                info!(gid = %g.gid, branch = %branch_id, "分支要求回滚");
                                status = GlobalStatus::Aborting;
                                self.store
                                    .set_global_status(
                                        &g.gid,
                                        GlobalStatus::Aborting,
                                        &format!("分支 {branch_id} 返回 FAILURE"),
                                    )
                                    .await?;
                            } else {
                                // 补偿都失败了,只能不停重试 —— 这时候需要人介入
                                warn!(gid = %g.gid, branch = %branch_id, "补偿失败,需要人工介入");
                                self.retry_later(g).await?;
                                return Ok(());
                            }
                        }
                        BranchResult::Ongoing | BranchResult::Unknown => {
                            // **绝不能当成失败**:对方可能已经成功了。退避重试。
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                    }
                }
            }
        }
    }

    // ---------------- TCC ----------------

    /// TCC 的 try 阶段是**客户端驱动**的(客户端先 registerBranch 再调 try),
    /// TC 只负责 confirm / cancel。所以分支的 URL 来自 `trans_branch_op` 表,
    /// 不是全局 payload。
    async fn process_tcc(&self, g: &GlobalRow) -> anyhow::Result<()> {
        self.drive_two_phase(g, BranchOp::Confirm, BranchOp::Cancel, "TCC")
            .await
    }

    // ---------------- XA ----------------

    /// XA 的一阶段(业务 SQL + `PREPARE TRANSACTION`)由**客户端**做,
    /// TC 只负责统一决定 `COMMIT PREPARED` 还是 `ROLLBACK PREPARED`。
    ///
    /// 形状跟 TCC 一样,只是 op 换成 commit/rollback。语义上那条铁律也一样:
    /// **commit 失败绝不能转 rollback** —— 别的分支可能已经提交了。
    ///
    /// XA 独有的严重性:没解决的 prepared 事务会**永久持锁**,在 Postgres 里
    /// 还阻塞 VACUUM。所以这里的重试比 SAGA 的补偿重试要紧得多。
    async fn process_xa(&self, g: &GlobalRow) -> anyhow::Result<()> {
        self.drive_two_phase(g, BranchOp::Commit, BranchOp::Rollback, "XA")
            .await
    }

    /// TCC 和 XA 共用的二阶段推进:正向 op 全做完就成功,反向 op 全做完就失败,
    /// **任一方向的失败都只重试,绝不改变方向**。
    async fn drive_two_phase(
        &self,
        g: &GlobalRow,
        fwd: BranchOp,
        bwd: BranchOp,
        label: &str,
    ) -> anyhow::Result<()> {
        let rows = self.store.list_branches(&g.gid).await?;
        let n = rows
            .iter()
            .filter_map(|r| index_of(&r.branch_id))
            .max()
            .map(|m| m + 1)
            .unwrap_or(0);
        if n == 0 {
            // 一个分支都没登记就 submit/abort 了 —— 空事务,直接落终态
            let s = if g.status == GlobalStatus::Aborting {
                GlobalStatus::Failed
            } else {
                GlobalStatus::Succeed
            };
            self.store.set_global_status(&g.gid, s, "").await?;
            return Ok(());
        }

        let status = g.status;
        loop {
            let rows = self.store.list_branches(&g.gid).await?;
            let (f, b) = split_by_op(&rows, n, fwd, bwd);
            let adv = if fwd == BranchOp::Commit {
                xa_advance(status, &f, &b)
            } else {
                tcc_advance(status, &f, &b)
            };
            match adv {
                Advance::Finish(s) => {
                    info!(gid = %g.gid, status = s.as_str(), mode = label, "事务终结");
                    self.store.set_global_status(&g.gid, s, "").await?;
                    return Ok(());
                }
                Advance::Wait => return Ok(()),
                // 只有 workflow 模式会出现,别的模式走到这里说明状态机接错了。
                // **不 panic** —— 推进器是常驻的,崩了整个 TC 就停了
                Advance::RunWorkflow => {
                    warn!(gid = %g.gid, "非 workflow 事务收到 RunWorkflow 决策,跳过");
                    return Ok(());
                }
                Advance::Call { index, op } => {
                    let bid = branch_id(index);
                    let Some(url) = url_of(&rows, &bid, op) else {
                        warn!(gid = %g.gid, branch = %bid, op = op.as_str(), mode = label,
                              "分支没登记这个操作的 URL,无法调用");
                        self.retry_later(g).await?;
                        return Ok(());
                    };
                    match self.call_branch(g, &bid, op, &url).await {
                        BranchResult::Success => {
                            self.store
                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Succeed)
                                .await?;
                        }
                        // 二阶段失败**绝不改变全局方向**:一阶段已经成功、
                        // 方向已经定了,反向操作会造成一半提交一半回滚。
                        // 唯一正确处理是无限重试 + 报警。
                        BranchResult::Failure => {
                            self.store
                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Failed)
                                .await?;
                            warn!(gid = %g.gid, branch = %bid, op = op.as_str(), mode = label,
                                  "二阶段失败,会持续重试,需要人工介入");
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                        BranchResult::Ongoing | BranchResult::Unknown => {
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                    }
                }
            }
        }
    }

    // ---------------- 二阶段消息 ----------------

    /// 流程:`prepare` 落库 → 业务提交本地事务 → `submit`。
    ///
    /// 如果进程在这两步之间崩了,事务会一直停在 prepared。这时 TC 靠回查
    /// `query_prepared` 问业务方"你那个本地事务到底提交了没有",据此决定
    /// 是往前推还是整单作废。**这是取代 MQ 事务消息的关键一环。**
    async fn process_msg(&self, g: &GlobalRow) -> anyhow::Result<()> {
        let mut status = g.status;

        if status == GlobalStatus::Prepared {
            if g.query_prepared.is_empty() {
                // 没给回查地址就没法自动决断。不能瞎猜 —— 猜错要么丢单要么重复扣款
                warn!(gid = %g.gid, "msg 事务没提供 query_prepared,无法回查,等人处理");
                self.retry_later(g).await?;
                return Ok(());
            }
            // 借用分支调用的通道做回查,branch_id 用 "00" 跟真实分支区分开
            match self
                .call_branch(g, "00", BranchOp::Action, &g.query_prepared)
                .await
            {
                BranchResult::Success => {
                    info!(gid = %g.gid, "回查:本地事务已提交 → 继续推进");
                    self.store
                        .set_global_status(&g.gid, GlobalStatus::Submitted, "")
                        .await?;
                    status = GlobalStatus::Submitted;
                }
                BranchResult::Failure => {
                    // 业务方明确说"这单没提交" → 整单作废。msg 没有补偿分支,
                    // 但也不需要:正向分支压根还没跑过
                    info!(gid = %g.gid, "回查:本地事务未提交 → 整单作废");
                    self.store
                        .set_global_status(
                            &g.gid,
                            GlobalStatus::Failed,
                            "回查得到 FAILURE:本地事务未提交",
                        )
                        .await?;
                    return Ok(());
                }
                BranchResult::Ongoing | BranchResult::Unknown => {
                    // 回查本身失败了,不能当作"没提交"。退避重试。
                    self.retry_later(g).await?;
                    return Ok(());
                }
            }
        }

        let steps: Vec<SagaStep> = serde_json::from_str(&g.payload).unwrap_or_default();
        if steps.is_empty() {
            self.store
                .set_global_status(&g.gid, GlobalStatus::Succeed, "")
                .await?;
            return Ok(());
        }
        loop {
            let (actions, _) = self.branch_states(&g.gid, steps.len()).await?;
            match msg_advance(status, &actions) {
                Advance::Finish(s) => {
                    info!(gid = %g.gid, status = s.as_str(), "消息事务终结");
                    self.store.set_global_status(&g.gid, s, "").await?;
                    return Ok(());
                }
                Advance::Wait => return Ok(()),
                // 只有 workflow 模式会出现,别的模式走到这里说明状态机接错了。
                // **不 panic** —— 推进器是常驻的,崩了整个 TC 就停了
                Advance::RunWorkflow => {
                    warn!(gid = %g.gid, "非 workflow 事务收到 RunWorkflow 决策,跳过");
                    return Ok(());
                }
                Advance::Call { index, op } => {
                    let bid = branch_id(index);
                    match self.call_branch(g, &bid, op, &steps[index].action).await {
                        BranchResult::Success => {
                            self.store
                                .set_branch_status(&g.gid, &bid, op, BranchStatus::Succeed)
                                .await?;
                        }
                        // msg 保证"最终一定送达",没有补偿一说。失败只能重试。
                        BranchResult::Failure | BranchResult::Ongoing | BranchResult::Unknown => {
                            self.retry_later(g).await?;
                            return Ok(());
                        }
                    }
                }
            }
        }
    }

    async fn retry_later(&self, g: &GlobalRow) -> anyhow::Result<()> {
        let iv = dtmrs_core::next_interval(g.next_cron_interval);
        self.store.schedule_retry(&g.gid, iv).await?;
        Ok(())
    }

    /// 取每一步的 action / compensate 分支当前状态,按步序对齐
    async fn branch_states(
        &self,
        gid: &str,
        n: usize,
    ) -> anyhow::Result<(Vec<BranchStatus>, Vec<BranchStatus>)> {
        let rows = self.store.list_branches(gid).await?;
        let mut actions = vec![BranchStatus::Prepared; n];
        let mut compensates = vec![BranchStatus::Prepared; n];
        for r in rows {
            let Some(i) = index_of(&r.branch_id) else {
                continue;
            };
            if i >= n {
                continue;
            }
            match r.op {
                BranchOp::Action | BranchOp::Try => actions[i] = r.status,
                BranchOp::Compensate | BranchOp::Cancel | BranchOp::Rollback => {
                    compensates[i] = r.status
                }
                _ => {}
            }
        }
        Ok((actions, compensates))
    }

    /// 调一个分支。`local://名字` 走进程内函数,`grpc://` 走 gRPC,其它走 HTTP。
    async fn call_branch(
        &self,
        g: &GlobalRow,
        branch_id: &str,
        op: BranchOp,
        url: &str,
    ) -> BranchResult {
        match parse_target(url) {
            Target::Local(name) => self.call_local(g, branch_id, op, &name).await,
            Target::Http(u) => self.call_http(g, branch_id, op, &u).await,
            #[cfg(feature = "grpc")]
            Target::Grpc(t) => {
                self.grpc
                    .call(
                        &t,
                        &g.gid,
                        &g.trans_type.to_string(),
                        branch_id,
                        op.as_str(),
                    )
                    .await
            }
            // 编译时关掉了 grpc feature,却遇到 grpc:// 分支。
            // 按「结果未知」处理(重试,不回滚)—— 这是构建配置问题,不是业务失败
            #[cfg(not(feature = "grpc"))]
            Target::Grpc(t) => {
                warn!(gid = %g.gid, branch = %branch_id, endpoint = %t.endpoint,
                      "遇到 grpc:// 分支但本次构建关掉了 grpc feature,按结果未知处理");
                BranchResult::Unknown
            }
        }
    }

    /// 进程内调用:没有网络、没有序列化,一次函数调用。
    async fn call_local(
        &self,
        g: &GlobalRow,
        branch_id: &str,
        op: BranchOp,
        name: &str,
    ) -> BranchResult {
        let Some(h) = self.registry.get(name) else {
            // 漏注册(比如新版本删了 handler)。**必须当 Unknown 而不是 Failure**:
            // 判失败会触发回滚,而这其实是部署问题,改回来重试才对。
            warn!(gid = %g.gid, branch = %branch_id, handler = name,
                  "本地分支未注册,按结果未知处理(会重试,不回滚)");
            return BranchResult::Unknown;
        };
        let ctx = BranchCtx {
            gid: g.gid.clone(),
            branch_id: branch_id.to_string(),
            op,
            trans_type: g.trans_type.to_string(),
        };
        let r = h(ctx).await;
        info!(gid = %g.gid, branch = %branch_id, op = op.as_str(),
              handler = name, result = ?r, "本地分支返回");
        r
    }

    /// 远端调用。查询参数跟 DTM 保持一致,客户端屏障库可以直接复用。
    async fn call_http(
        &self,
        g: &GlobalRow,
        branch_id: &str,
        op: BranchOp,
        url: &str,
    ) -> BranchResult {
        let req = self
            .http
            .post(url)
            .query(&[
                ("gid", g.gid.as_str()),
                ("trans_type", &g.trans_type.to_string()),
                ("branch_id", branch_id),
                ("op", op.as_str()),
            ])
            .header("content-type", "application/json")
            .body(branch_payload(&g.payload));
        match req.send().await {
            Ok(resp) => {
                let code = resp.status().as_u16();
                let body = resp.text().await.unwrap_or_default();
                let r = BranchResult::from_http(code, &body);
                info!(gid = %g.gid, branch = %branch_id, op = op.as_str(), code, result = ?r, "分支返回");
                r
            }
            Err(e) => {
                // 超时/连不上 —— 结果未知,必须重试而不是回滚
                warn!(gid = %g.gid, branch = %branch_id, error = %e, "分支不可达,结果未知");
                BranchResult::Unknown
            }
        }
    }
}

/// 分支号:第 0 步是 "01",跟 DTM 一致
pub fn branch_id(index: usize) -> String {
    format!("{:02}", index + 1)
}

fn index_of(branch_id: &str) -> Option<usize> {
    branch_id
        .parse::<usize>()
        .ok()
        .and_then(|v| v.checked_sub(1))
}

/// 按 op 把分支行拆成两列(正向 / 反向),按步序对齐
fn split_by_op(
    rows: &[dtmrs_store::BranchRow],
    n: usize,
    fwd: BranchOp,
    bwd: BranchOp,
) -> (Vec<BranchStatus>, Vec<BranchStatus>) {
    let mut a = vec![BranchStatus::Prepared; n];
    let mut b = vec![BranchStatus::Prepared; n];
    for r in rows {
        let Some(i) = index_of(&r.branch_id) else {
            continue;
        };
        if i >= n {
            continue;
        }
        if r.op == fwd {
            a[i] = r.status;
        } else if r.op == bwd {
            b[i] = r.status;
        }
    }
    (a, b)
}

/// 按分支序取出各 compensate 行的状态,用于 workflow 的逆序补偿。
///
/// 跟 `split_by_op` 的区别:workflow 的分支数是**运行时长出来的**,
/// 提交时并不知道有几个,所以长度得从已登记的行里推出来。
fn compensate_states(rows: &[dtmrs_store::BranchRow]) -> Vec<BranchStatus> {
    let n = rows
        .iter()
        .filter(|r| r.op == BranchOp::Compensate)
        .filter_map(|r| index_of(&r.branch_id))
        .max()
        .map(|m| m + 1)
        .unwrap_or(0);
    // 没登记补偿的分支(比如纯查询步骤)留成 Succeed,逆序扫描时会跳过它 ——
    // 本来就没有副作用要收拾
    let mut v = vec![BranchStatus::Succeed; n];
    for r in rows.iter().filter(|r| r.op == BranchOp::Compensate) {
        if let Some(i) = index_of(&r.branch_id) {
            if i < n {
                v[i] = r.status;
            }
        }
    }
    v
}

fn url_of(rows: &[dtmrs_store::BranchRow], branch_id: &str, op: BranchOp) -> Option<String> {
    rows.iter()
        .find(|r| r.branch_id == branch_id && r.op == op)
        .map(|r| r.url.clone())
}

/// MVP 阶段各分支共用同一份请求体。真实场景应该每步独立 payload,
/// 那是第二版的事(见 DESIGN.md 范围表)。
fn branch_payload(_global_payload: &str) -> String {
    "{}".to_string()
}

#[cfg(test)]
mod tests {
    use super::*;

    #[test]
    fn 分支号与下标互转() {
        assert_eq!(branch_id(0), "01");
        assert_eq!(branch_id(9), "10");
        assert_eq!(index_of("01"), Some(0));
        assert_eq!(index_of("10"), Some(9));
        assert_eq!(index_of("00"), None);
        assert_eq!(index_of("xx"), None);
    }
}