dtmrs_server/grpc/
server.rs1use tonic::{Request, Response, Status};
9
10use super::pb;
11use crate::api::{Api, ApiError, RegisterBranch};
12use dtmrs_core::SagaStep;
13
14impl From<ApiError> for Status {
15 fn from(e: ApiError) -> Self {
16 match &e {
17 ApiError::BadRequest(m) => Status::invalid_argument(m.clone()),
18 ApiError::NotFound(m) => Status::not_found(m.clone()),
19 ApiError::Conflict(m) => Status::failed_precondition(m.clone()),
20 ApiError::Internal(m) => Status::internal(m.clone()),
21 }
22 }
23}
24
25pub struct TcService {
26 api: Api,
27}
28
29impl TcService {
30 pub fn new(api: Api) -> Self {
31 Self { api }
32 }
33
34 pub fn into_server(self) -> pb::tc_server::TcServer<Self> {
36 pb::tc_server::TcServer::new(self)
37 }
38}
39
40#[tonic::async_trait]
41impl pb::tc_server::Tc for TcService {
42 async fn new_gid(
43 &self,
44 _req: Request<pb::NewGidRequest>,
45 ) -> Result<Response<pb::NewGidReply>, Status> {
46 Ok(Response::new(pb::NewGidReply {
47 gid: self.api.new_gid(),
48 }))
49 }
50
51 async fn prepare(
52 &self,
53 req: Request<pb::PrepareRequest>,
54 ) -> Result<Response<pb::Empty>, Status> {
55 let r = req.into_inner();
56 let grace = if r.grace_secs > 0 {
59 Some(r.grace_secs)
60 } else {
61 None
62 };
63 self.api
64 .prepare(&r.gid, &r.trans_type, &r.actions, &r.query_prepared, grace)
65 .await?;
66 Ok(Response::new(pb::Empty {}))
67 }
68
69 async fn register_branch(
70 &self,
71 req: Request<pb::RegisterBranchRequest>,
72 ) -> Result<Response<pb::Empty>, Status> {
73 let r = req.into_inner();
74 self.api
75 .register_branch(&RegisterBranch {
76 gid: r.gid,
77 branch_id: r.branch_id,
78 confirm: r.confirm,
79 cancel: r.cancel,
80 r#try: r.r#try,
81 commit: r.commit,
82 rollback: r.rollback,
83 })
84 .await?;
85 Ok(Response::new(pb::Empty {}))
86 }
87
88 async fn submit(&self, req: Request<pb::SubmitRequest>) -> Result<Response<pb::Empty>, Status> {
89 let r = req.into_inner();
90 let tt = if r.trans_type.is_empty() {
91 "saga"
92 } else {
93 &r.trans_type
94 };
95 let steps: Vec<SagaStep> = r
96 .steps
97 .into_iter()
98 .map(|s| SagaStep {
99 action: s.action,
100 compensate: s.compensate,
101 payload: s.payload,
102 })
103 .collect();
104 self.api.submit(&r.gid, tt, &steps).await?;
105 Ok(Response::new(pb::Empty {}))
106 }
107
108 async fn abort(&self, req: Request<pb::AbortRequest>) -> Result<Response<pb::Empty>, Status> {
109 self.api.abort(&req.into_inner().gid).await?;
110 Ok(Response::new(pb::Empty {}))
111 }
112
113 async fn retry(&self, req: Request<pb::RetryRequest>) -> Result<Response<pb::Empty>, Status> {
114 self.api.retry(&req.into_inner().gid).await?;
115 Ok(Response::new(pb::Empty {}))
116 }
117
118 async fn query(
119 &self,
120 req: Request<pb::QueryRequest>,
121 ) -> Result<Response<pb::TransView>, Status> {
122 let v = self.api.query(&req.into_inner().gid).await?;
123 Ok(Response::new(pb::TransView {
124 gid: v.gid,
125 trans_type: v.trans_type,
126 status: v.status,
127 rollback_reason: v.rollback_reason,
128 create_time: v.create_time,
129 finish_time: v.finish_time,
130 branches: v
131 .branches
132 .into_iter()
133 .map(|b| pb::BranchView {
134 branch_id: b.branch_id,
135 op: b.op,
136 url: b.url,
137 status: b.status,
138 })
139 .collect(),
140 }))
141 }
142}