From ad5c1b6b93f160e3bf34ab954963bbe43d72731b Mon Sep 17 00:00:00 2001 From: yedf2 <120050102@qq.com> Date: Fri, 1 Jul 2022 11:40:47 +0800 Subject: [PATCH] merge workflow prepare and progress --- dtmgrpc/dtmgpb/dtmgimp.pb.go | 19 ++++++++++--------- dtmgrpc/dtmgpb/dtmgimp.proto | 2 +- dtmgrpc/dtmgpb/dtmgimp_grpc.pb.go | 24 ++++++++++++------------ dtmgrpc/workflow/imp.go | 5 +---- dtmgrpc/workflow/rpc.go | 4 ++-- dtmsvr/api.go | 9 +++++++++ dtmsvr/api_grpc.go | 6 +++--- dtmsvr/api_http.go | 11 +++++++---- 8 files changed, 45 insertions(+), 35 deletions(-) diff --git a/dtmgrpc/dtmgpb/dtmgimp.pb.go b/dtmgrpc/dtmgpb/dtmgimp.pb.go index 7f40418..35cd8d4 100644 --- a/dtmgrpc/dtmgpb/dtmgimp.pb.go +++ b/dtmgrpc/dtmgpb/dtmgimp.pb.go @@ -558,7 +558,7 @@ var file_dtmgrpc_dtmgpb_dtmgimp_proto_rawDesc = []byte{ 0x42, 0x69, 0x6e, 0x44, 0x61, 0x74, 0x61, 0x12, 0x1a, 0x0a, 0x08, 0x42, 0x72, 0x61, 0x6e, 0x63, 0x68, 0x49, 0x44, 0x18, 0x03, 0x20, 0x01, 0x28, 0x09, 0x52, 0x08, 0x42, 0x72, 0x61, 0x6e, 0x63, 0x68, 0x49, 0x44, 0x12, 0x0e, 0x0a, 0x02, 0x4f, 0x70, 0x18, 0x04, 0x20, 0x01, 0x28, 0x09, 0x52, - 0x02, 0x4f, 0x70, 0x32, 0xf3, 0x02, 0x0a, 0x03, 0x44, 0x74, 0x6d, 0x12, 0x38, 0x0a, 0x06, 0x4e, + 0x02, 0x4f, 0x70, 0x32, 0xf8, 0x02, 0x0a, 0x03, 0x44, 0x74, 0x6d, 0x12, 0x38, 0x0a, 0x06, 0x4e, 0x65, 0x77, 0x47, 0x69, 0x64, 0x12, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, 0x45, 0x6d, 0x70, 0x74, 0x79, 0x1a, 0x14, 0x2e, 0x64, 0x74, 0x6d, 0x67, 0x69, 0x6d, 0x70, 0x2e, 0x44, 0x74, 0x6d, 0x47, 0x69, 0x64, 0x52, 0x65, @@ -577,12 +577,13 @@ var file_dtmgrpc_dtmgpb_dtmgimp_proto_rawDesc = []byte{ 0x63, 0x68, 0x12, 0x19, 0x2e, 0x64, 0x74, 0x6d, 0x67, 0x69, 0x6d, 0x70, 0x2e, 0x44, 0x74, 0x6d, 0x42, 0x72, 0x61, 0x6e, 0x63, 0x68, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x16, 0x2e, 0x67, 0x6f, 0x6f, 0x67, 0x6c, 0x65, 0x2e, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x62, 0x75, 0x66, 0x2e, - 0x45, 0x6d, 0x70, 0x74, 0x79, 0x22, 0x00, 0x12, 0x40, 0x0a, 0x0a, 0x50, 0x72, 0x6f, 0x67, 0x72, - 0x65, 0x73, 0x73, 0x65, 0x73, 0x12, 0x13, 0x2e, 0x64, 0x74, 0x6d, 0x67, 0x69, 0x6d, 0x70, 0x2e, - 0x44, 0x74, 0x6d, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, 0x1b, 0x2e, 0x64, 0x74, 0x6d, - 0x67, 0x69, 0x6d, 0x70, 0x2e, 0x44, 0x74, 0x6d, 0x50, 0x72, 0x6f, 0x67, 0x72, 0x65, 0x73, 0x73, - 0x65, 0x73, 0x52, 0x65, 0x70, 0x6c, 0x79, 0x22, 0x00, 0x42, 0x0a, 0x5a, 0x08, 0x2e, 0x2f, 0x64, - 0x74, 0x6d, 0x67, 0x70, 0x62, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, 0x6f, 0x33, + 0x45, 0x6d, 0x70, 0x74, 0x79, 0x22, 0x00, 0x12, 0x45, 0x0a, 0x0f, 0x50, 0x72, 0x65, 0x70, 0x61, + 0x72, 0x65, 0x57, 0x6f, 0x72, 0x6b, 0x66, 0x6c, 0x6f, 0x77, 0x12, 0x13, 0x2e, 0x64, 0x74, 0x6d, + 0x67, 0x69, 0x6d, 0x70, 0x2e, 0x44, 0x74, 0x6d, 0x52, 0x65, 0x71, 0x75, 0x65, 0x73, 0x74, 0x1a, + 0x1b, 0x2e, 0x64, 0x74, 0x6d, 0x67, 0x69, 0x6d, 0x70, 0x2e, 0x44, 0x74, 0x6d, 0x50, 0x72, 0x6f, + 0x67, 0x72, 0x65, 0x73, 0x73, 0x65, 0x73, 0x52, 0x65, 0x70, 0x6c, 0x79, 0x22, 0x00, 0x42, 0x0a, + 0x5a, 0x08, 0x2e, 0x2f, 0x64, 0x74, 0x6d, 0x67, 0x70, 0x62, 0x62, 0x06, 0x70, 0x72, 0x6f, 0x74, + 0x6f, 0x33, } var ( @@ -621,13 +622,13 @@ var file_dtmgrpc_dtmgpb_dtmgimp_proto_depIdxs = []int32{ 1, // 7: dtmgimp.Dtm.Prepare:input_type -> dtmgimp.DtmRequest 1, // 8: dtmgimp.Dtm.Abort:input_type -> dtmgimp.DtmRequest 3, // 9: dtmgimp.Dtm.RegisterBranch:input_type -> dtmgimp.DtmBranchRequest - 1, // 10: dtmgimp.Dtm.Progresses:input_type -> dtmgimp.DtmRequest + 1, // 10: dtmgimp.Dtm.PrepareWorkflow:input_type -> dtmgimp.DtmRequest 2, // 11: dtmgimp.Dtm.NewGid:output_type -> dtmgimp.DtmGidReply 9, // 12: dtmgimp.Dtm.Submit:output_type -> google.protobuf.Empty 9, // 13: dtmgimp.Dtm.Prepare:output_type -> google.protobuf.Empty 9, // 14: dtmgimp.Dtm.Abort:output_type -> google.protobuf.Empty 9, // 15: dtmgimp.Dtm.RegisterBranch:output_type -> google.protobuf.Empty - 4, // 16: dtmgimp.Dtm.Progresses:output_type -> dtmgimp.DtmProgressesReply + 4, // 16: dtmgimp.Dtm.PrepareWorkflow:output_type -> dtmgimp.DtmProgressesReply 11, // [11:17] is the sub-list for method output_type 5, // [5:11] is the sub-list for method input_type 5, // [5:5] is the sub-list for extension type_name diff --git a/dtmgrpc/dtmgpb/dtmgimp.proto b/dtmgrpc/dtmgpb/dtmgimp.proto index c777850..6005696 100644 --- a/dtmgrpc/dtmgpb/dtmgimp.proto +++ b/dtmgrpc/dtmgpb/dtmgimp.proto @@ -12,7 +12,7 @@ service Dtm { rpc Prepare(DtmRequest) returns (google.protobuf.Empty) {} rpc Abort(DtmRequest) returns (google.protobuf.Empty) {} rpc RegisterBranch(DtmBranchRequest) returns (google.protobuf.Empty) {} - rpc Progresses(DtmRequest) returns (DtmProgressesReply) {} + rpc PrepareWorkflow(DtmRequest) returns (DtmProgressesReply) {} } message DtmTransOptions { diff --git a/dtmgrpc/dtmgpb/dtmgimp_grpc.pb.go b/dtmgrpc/dtmgpb/dtmgimp_grpc.pb.go index 5045c33..6327fb7 100644 --- a/dtmgrpc/dtmgpb/dtmgimp_grpc.pb.go +++ b/dtmgrpc/dtmgpb/dtmgimp_grpc.pb.go @@ -28,7 +28,7 @@ type DtmClient interface { Prepare(ctx context.Context, in *DtmRequest, opts ...grpc.CallOption) (*emptypb.Empty, error) Abort(ctx context.Context, in *DtmRequest, opts ...grpc.CallOption) (*emptypb.Empty, error) RegisterBranch(ctx context.Context, in *DtmBranchRequest, opts ...grpc.CallOption) (*emptypb.Empty, error) - Progresses(ctx context.Context, in *DtmRequest, opts ...grpc.CallOption) (*DtmProgressesReply, error) + PrepareWorkflow(ctx context.Context, in *DtmRequest, opts ...grpc.CallOption) (*DtmProgressesReply, error) } type dtmClient struct { @@ -84,9 +84,9 @@ func (c *dtmClient) RegisterBranch(ctx context.Context, in *DtmBranchRequest, op return out, nil } -func (c *dtmClient) Progresses(ctx context.Context, in *DtmRequest, opts ...grpc.CallOption) (*DtmProgressesReply, error) { +func (c *dtmClient) PrepareWorkflow(ctx context.Context, in *DtmRequest, opts ...grpc.CallOption) (*DtmProgressesReply, error) { out := new(DtmProgressesReply) - err := c.cc.Invoke(ctx, "/dtmgimp.Dtm/Progresses", in, out, opts...) + err := c.cc.Invoke(ctx, "/dtmgimp.Dtm/PrepareWorkflow", in, out, opts...) if err != nil { return nil, err } @@ -102,7 +102,7 @@ type DtmServer interface { Prepare(context.Context, *DtmRequest) (*emptypb.Empty, error) Abort(context.Context, *DtmRequest) (*emptypb.Empty, error) RegisterBranch(context.Context, *DtmBranchRequest) (*emptypb.Empty, error) - Progresses(context.Context, *DtmRequest) (*DtmProgressesReply, error) + PrepareWorkflow(context.Context, *DtmRequest) (*DtmProgressesReply, error) mustEmbedUnimplementedDtmServer() } @@ -125,8 +125,8 @@ func (UnimplementedDtmServer) Abort(context.Context, *DtmRequest) (*emptypb.Empt func (UnimplementedDtmServer) RegisterBranch(context.Context, *DtmBranchRequest) (*emptypb.Empty, error) { return nil, status.Errorf(codes.Unimplemented, "method RegisterBranch not implemented") } -func (UnimplementedDtmServer) Progresses(context.Context, *DtmRequest) (*DtmProgressesReply, error) { - return nil, status.Errorf(codes.Unimplemented, "method Progresses not implemented") +func (UnimplementedDtmServer) PrepareWorkflow(context.Context, *DtmRequest) (*DtmProgressesReply, error) { + return nil, status.Errorf(codes.Unimplemented, "method PrepareWorkflow not implemented") } func (UnimplementedDtmServer) mustEmbedUnimplementedDtmServer() {} @@ -231,20 +231,20 @@ func _Dtm_RegisterBranch_Handler(srv interface{}, ctx context.Context, dec func( return interceptor(ctx, in, info, handler) } -func _Dtm_Progresses_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { +func _Dtm_PrepareWorkflow_Handler(srv interface{}, ctx context.Context, dec func(interface{}) error, interceptor grpc.UnaryServerInterceptor) (interface{}, error) { in := new(DtmRequest) if err := dec(in); err != nil { return nil, err } if interceptor == nil { - return srv.(DtmServer).Progresses(ctx, in) + return srv.(DtmServer).PrepareWorkflow(ctx, in) } info := &grpc.UnaryServerInfo{ Server: srv, - FullMethod: "/dtmgimp.Dtm/Progresses", + FullMethod: "/dtmgimp.Dtm/PrepareWorkflow", } handler := func(ctx context.Context, req interface{}) (interface{}, error) { - return srv.(DtmServer).Progresses(ctx, req.(*DtmRequest)) + return srv.(DtmServer).PrepareWorkflow(ctx, req.(*DtmRequest)) } return interceptor(ctx, in, info, handler) } @@ -277,8 +277,8 @@ var Dtm_ServiceDesc = grpc.ServiceDesc{ Handler: _Dtm_RegisterBranch_Handler, }, { - MethodName: "Progresses", - Handler: _Dtm_Progresses_Handler, + MethodName: "PrepareWorkflow", + Handler: _Dtm_PrepareWorkflow_Handler, }, }, Streams: []grpc.StreamDesc{}, diff --git a/dtmgrpc/workflow/imp.go b/dtmgrpc/workflow/imp.go index 7f43199..5fd8c7a 100644 --- a/dtmgrpc/workflow/imp.go +++ b/dtmgrpc/workflow/imp.go @@ -95,10 +95,7 @@ func (wf *Workflow) initRestyClient() { } func (wf *Workflow) process(handler WfFunc, data []byte) (err error) { - err = wf.prepare() - if err == nil { - err = wf.loadProgresses() - } + err = wf.loadProgresses() if err == nil { err = handler(wf, data) err = dtmgrpc.GrpcError2DtmError(err) diff --git a/dtmgrpc/workflow/rpc.go b/dtmgrpc/workflow/rpc.go index f76076f..a2204d5 100644 --- a/dtmgrpc/workflow/rpc.go +++ b/dtmgrpc/workflow/rpc.go @@ -14,14 +14,14 @@ import ( func (wf *Workflow) getProgress() ([]*dtmgpb.DtmProgress, error) { if wf.Protocol == dtmimp.ProtocolGRPC { var reply dtmgpb.DtmProgressesReply - err := dtmgimp.MustGetGrpcConn(wf.Dtm, false).Invoke(wf.Context, "/dtmgimp.Dtm/Progresses", + err := dtmgimp.MustGetGrpcConn(wf.Dtm, false).Invoke(wf.Context, "/dtmgimp.Dtm/PrepareWorkflow", dtmgimp.GetDtmRequest(wf.TransBase), &reply) if err == nil { return reply.Progresses, nil } return nil, err } - resp, err := dtmimp.RestyClient.R().SetQueryParam("gid", wf.Gid).Get(wf.Dtm + "/progresses") + resp, err := dtmimp.RestyClient.R().SetBody(wf.TransBase).Post(wf.Dtm + "/prepareWorkflow") var progresses []*dtmgpb.DtmProgress if err == nil { dtmimp.MustUnmarshal(resp.Body(), &progresses) diff --git a/dtmsvr/api.go b/dtmsvr/api.go index 7d14f56..614cdac 100644 --- a/dtmsvr/api.go +++ b/dtmsvr/api.go @@ -50,6 +50,15 @@ func svcPrepare(t *TransGlobal) interface{} { return err } +func svcPrepareWorkflow(t *TransGlobal) ([]TransBranch, error) { + t.Status = dtmcli.StatusPrepared + _, err := t.saveNew() + if err == storage.ErrUniqueConflict { // transaction exists, query the branches + return GetStore().FindBranches(t.Gid), nil + } + return []TransBranch{}, nil +} + func svcAbort(t *TransGlobal) interface{} { dbt := GetTransGlobal(t.Gid) if dbt.TransType == "msg" && dbt.Status == dtmcli.StatusPrepared { diff --git a/dtmsvr/api_grpc.go b/dtmsvr/api_grpc.go index f62ad91..9989577 100644 --- a/dtmsvr/api_grpc.go +++ b/dtmsvr/api_grpc.go @@ -49,8 +49,8 @@ func (s *dtmServer) RegisterBranch(ctx context.Context, in *pb.DtmBranchRequest) return &emptypb.Empty{}, dtmgrpc.DtmError2GrpcError(r) } -func (s *dtmServer) Progresses(ctx context.Context, in *pb.DtmRequest) (*pb.DtmProgressesReply, error) { - branches := GetStore().FindBranches(in.Gid) +func (s *dtmServer) PrepareWorkflow(ctx context.Context, in *pb.DtmRequest) (*pb.DtmProgressesReply, error) { + branches, err := svcPrepareWorkflow(TransFromDtmRequest(ctx, in)) reply := &pb.DtmProgressesReply{ Progresses: []*pb.DtmProgress{}, } @@ -62,5 +62,5 @@ func (s *dtmServer) Progresses(ctx context.Context, in *pb.DtmRequest) (*pb.DtmP BinData: b.BinData, }) } - return reply, nil + return reply, dtmgrpc.DtmError2GrpcError(err) } diff --git a/dtmsvr/api_http.go b/dtmsvr/api_http.go index 52e993e..563006c 100644 --- a/dtmsvr/api_http.go +++ b/dtmsvr/api_http.go @@ -30,8 +30,8 @@ func addRoute(engine *gin.Engine) { engine.POST("/api/dtmsvr/registerBranch", dtmutil.WrapHandler2(registerBranch)) engine.POST("/api/dtmsvr/registerXaBranch", dtmutil.WrapHandler2(registerBranch)) // compatible for old sdk engine.POST("/api/dtmsvr/registerTccBranch", dtmutil.WrapHandler2(registerBranch)) // compatible for old sdk + engine.POST("/api/dtmsvr/prepareWorkflow", dtmutil.WrapHandler2(prepareWorkflow)) engine.GET("/api/dtmsvr/query", dtmutil.WrapHandler2(query)) - engine.GET("/api/dtmsvr/progresses", dtmutil.WrapHandler2(progresses)) engine.GET("/api/dtmsvr/all", dtmutil.WrapHandler2(all)) engine.GET("/api/dtmsvr/resetCronTime", dtmutil.WrapHandler2(resetCronTime)) @@ -86,9 +86,12 @@ func query(c *gin.Context) interface{} { return map[string]interface{}{"transaction": trans, "branches": branches} } -func progresses(c *gin.Context) interface{} { - gid := c.Query("gid") - return GetStore().FindBranches(gid) +func prepareWorkflow(c *gin.Context) interface{} { + branches, err := svcPrepareWorkflow(TransFromContext(c)) + if err != nil { + return err + } + return branches } func all(c *gin.Context) interface{} {