-
Notifications
You must be signed in to change notification settings - Fork 0
/
service.go
93 lines (81 loc) · 2.59 KB
/
service.go
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
/*
Copyright IBM Corp. 2017 All Rights Reserved.
SPDX-License-Identifier: Apache-2.0
*/
package cluster
import (
"context"
"io"
"github.com/hyperledger/fabric/common/flogging"
"github.com/hyperledger/fabric/common/util"
"github.com/hyperledger/fabric/protos/orderer"
"google.golang.org/grpc"
)
//go:generate mockery -dir . -name Dispatcher -case underscore -output ./mocks/
// Dispatcher dispatches requests
type Dispatcher interface {
DispatchSubmit(ctx context.Context, request *orderer.SubmitRequest) (*orderer.SubmitResponse, error)
DispatchStep(ctx context.Context, request *orderer.StepRequest) (*orderer.StepResponse, error)
}
//go:generate mockery -dir . -name SubmitStream -case underscore -output ./mocks/
// SubmitStream defines the gRPC stream for sending
// transactions, and receiving corresponding responses
type SubmitStream interface {
Send(response *orderer.SubmitResponse) error
Recv() (*orderer.SubmitRequest, error)
grpc.ServerStream
}
// Service defines the raft Service
type Service struct {
Dispatcher Dispatcher
Logger *flogging.FabricLogger
StepLogger *flogging.FabricLogger
}
// Step forwards a message to a raft FSM located in this server
func (s *Service) Step(ctx context.Context, request *orderer.StepRequest) (*orderer.StepResponse, error) {
addr := util.ExtractRemoteAddress(ctx)
s.StepLogger.Debugf("Connection from %s", addr)
defer s.StepLogger.Debugf("Closing connection from %s", addr)
response, err := s.Dispatcher.DispatchStep(ctx, request)
if err != nil {
s.Logger.Warningf("Handling of Step() from %s failed: %+v", addr, err)
}
return response, err
}
// Submit accepts transactions
func (s *Service) Submit(stream orderer.Cluster_SubmitServer) error {
addr := util.ExtractRemoteAddress(stream.Context())
s.Logger.Debugf("Connection from %s", addr)
defer s.Logger.Debugf("Closing connection from %s", addr)
for {
err := s.handleSubmit(stream, addr)
if err == io.EOF {
s.Logger.Debugf("%s disconnected", addr)
return nil
}
if err != nil {
return err
}
// Else, no error occurred, so we continue to the next iteration
}
}
func (s *Service) handleSubmit(stream SubmitStream, addr string) error {
request, err := stream.Recv()
if err == io.EOF {
return err
}
if err != nil {
s.Logger.Warningf("Stream read from %s failed: %v", addr, err)
return err
}
response, err := s.Dispatcher.DispatchSubmit(stream.Context(), request)
if err != nil {
s.Logger.Warningf("Handling of Propose() from %s failed: %+v", addr, err)
return err
}
err = stream.Send(response)
if err != nil {
s.Logger.Warningf("Send() failed: %v", err)
}
return err
}