Browse Source

merge workflow prepare and progress

pull/328/head
yedf2 4 years ago
parent
commit
ad5c1b6b93
  1. 19
      dtmgrpc/dtmgpb/dtmgimp.pb.go
  2. 2
      dtmgrpc/dtmgpb/dtmgimp.proto
  3. 24
      dtmgrpc/dtmgpb/dtmgimp_grpc.pb.go
  4. 5
      dtmgrpc/workflow/imp.go
  5. 4
      dtmgrpc/workflow/rpc.go
  6. 9
      dtmsvr/api.go
  7. 6
      dtmsvr/api_grpc.go
  8. 11
      dtmsvr/api_http.go

19
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

2
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 {

24
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{},

5
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)

4
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)

9
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 {

6
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)
}

11
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{} {

Loading…
Cancel
Save