-
Notifications
You must be signed in to change notification settings - Fork 178
/
request_handler_engine.go
95 lines (78 loc) · 2.14 KB
/
request_handler_engine.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
94
95
package synchronization
import (
"fmt"
"github.com/rs/zerolog"
"github.com/onflow/flow-go/engine"
"github.com/onflow/flow-go/model/flow"
"github.com/onflow/flow-go/model/messages"
"github.com/onflow/flow-go/module"
"github.com/onflow/flow-go/network"
"github.com/onflow/flow-go/storage"
)
type ResponseSender interface {
SendResponse(interface{}, flow.Identifier) error
}
type ResponseSenderImpl struct {
con network.Conduit
}
func (r *ResponseSenderImpl) SendResponse(res interface{}, target flow.Identifier) error {
switch res.(type) {
case *messages.BlockResponse:
err := r.con.Unicast(res, target)
if err != nil {
return fmt.Errorf("could not unicast block response to target %x: %w", target, err)
}
case *messages.SyncResponse:
err := r.con.Unicast(res, target)
if err != nil {
return fmt.Errorf("could not unicast sync response to target %x: %w", target, err)
}
default:
return fmt.Errorf("unable to unicast unexpected response %+v", res)
}
return nil
}
func NewResponseSender(con network.Conduit) *ResponseSenderImpl {
return &ResponseSenderImpl{
con: con,
}
}
type RequestHandlerEngine struct {
requestHandler *RequestHandler
}
var _ network.MessageProcessor = (*RequestHandlerEngine)(nil)
func NewRequestHandlerEngine(
logger zerolog.Logger,
metrics module.EngineMetrics,
net network.Network,
me module.Local,
blocks storage.Blocks,
core module.SyncCore,
finalizedHeader *FinalizedHeaderCache,
) (*RequestHandlerEngine, error) {
e := &RequestHandlerEngine{}
con, err := net.Register(engine.PublicSyncCommittee, e)
if err != nil {
return nil, fmt.Errorf("could not register engine: %w", err)
}
e.requestHandler = NewRequestHandler(
logger,
metrics,
NewResponseSender(con),
me,
blocks,
core,
finalizedHeader,
false,
)
return e, nil
}
func (r *RequestHandlerEngine) Process(channel network.Channel, originID flow.Identifier, event interface{}) error {
return r.requestHandler.Process(channel, originID, event)
}
func (r *RequestHandlerEngine) Ready() <-chan struct{} {
return r.requestHandler.Ready()
}
func (r *RequestHandlerEngine) Done() <-chan struct{} {
return r.requestHandler.Done()
}