@cryptotaxi247 / kubo / commits / 338b03723

clean up and add inet.Network to bitswap

new Service interface

Juan Batiz-Benet committed Oct 10, 2014 at 05:15 UTC 338b037238ab893da94223222714b0aa457ff07a
6 files changed +59 -40
core/core.go
+4 -4
@@ -126,7 +126,7 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
126 return nil, err
127 }
128
129 - net, err = inet.NewIpfsNetwork(context.TODO(), local, peerstore, &mux.ProtocolMap{
129 + net, err = inet.NewIpfsNetwork(ctx, local, peerstore, &mux.ProtocolMap{
130 mux.ProtocolID_Routing: dhtService,
131 mux.ProtocolID_Exchange: exchangeService,
132 mux.ProtocolID_Diagnostic: diagService,
@@ -137,14 +137,14 @@ func NewIpfsNode(cfg *config.Config, online bool) (*IpfsNode, error) {
137 }
138
139 diagnostics = diag.NewDiagnostics(local, net, diagService)
140 - diagService.Handler = diagnostics
140 + diagService.SetHandler(diagnostics)
141
142 route = dht.NewDHT(local, peerstore, net, dhtService, d)
143 // TODO(brian): perform this inside NewDHT factory method
144 - dhtService.Handler = route // wire the handler to the service.
144 + dhtService.SetHandler(route) // wire the handler to the service.
145
146 const alwaysSendToPeer = true // use YesManStrategy
147 - exchangeSession = bitswap.NetMessageSession(ctx, local, exchangeService, route, d, alwaysSendToPeer)
147 + exchangeSession = bitswap.NetMessageSession(ctx, local, net, exchangeService, route, d, alwaysSendToPeer)
148
149 // TODO(brian): pass a context to initConnections
150 go initConnections(ctx, cfg, peerstore, route)
exchange/bitswap/bitswap.go
+10 -3
@@ -11,6 +11,7 @@ import (
11 bsnet "github.com/jbenet/go-ipfs/exchange/bitswap/network"
12 notifications "github.com/jbenet/go-ipfs/exchange/bitswap/notifications"
13 strategy "github.com/jbenet/go-ipfs/exchange/bitswap/strategy"
14 + inet "github.com/jbenet/go-ipfs/net"
15 peer "github.com/jbenet/go-ipfs/peer"
16 u "github.com/jbenet/go-ipfs/util"
17 )
@@ -19,14 +20,17 @@ var log = u.Logger("bitswap")
20
21 // NetMessageSession initializes a BitSwap session that communicates over the
22 // provided NetMessage service
22 -func NetMessageSession(parent context.Context, p *peer.Peer, s bsnet.NetMessageService, directory bsnet.Routing, d ds.Datastore, nice bool) exchange.Interface {
23 +func NetMessageSession(parent context.Context, p *peer.Peer,
24 + net inet.Network, srv inet.Service, directory bsnet.Routing,
25 + d ds.Datastore, nice bool) exchange.Interface {
26
24 - networkAdapter := bsnet.NetMessageAdapter(s, nil)
27 + networkAdapter := bsnet.NetMessageAdapter(srv, nil)
28 bs := &bitswap{
29 blockstore: blockstore.NewBlockstore(d),
30 notifications: notifications.New(),
31 strategy: strategy.New(nice),
32 routing: directory,
33 + network: net,
34 sender: networkAdapter,
35 wantlist: u.NewKeySet(),
36 }
@@ -38,6 +42,9 @@ func NetMessageSession(parent context.Context, p *peer.Peer, s bsnet.NetMessageS
42 // bitswap instances implement the bitswap protocol.
43 type bitswap struct {
44
45 + // network maintains connections to the outside world.
46 + network inet.Network
47 +
48 // sender delivers messages on behalf of the session
49 sender bsnet.Adapter
50
@@ -79,7 +86,7 @@ func (bs *bitswap) Block(parent context.Context, k u.Key) (*blocks.Block, error)
86 }
87 message.AppendWanted(k)
88 for iiiii := range peersToQuery {
82 - // log.Debug("bitswap got peersToQuery: %s", iiiii)
89 + log.Debug("bitswap got peersToQuery: %s", iiiii)
90 go func(p *peer.Peer) {
91 response, err := bs.sender.SendRequest(ctx, p, message)
92 if err != nil {
exchange/bitswap/network/interface.go
-9
@@ -2,10 +2,8 @@ package network
2
3 import (
4 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5 - netservice "github.com/jbenet/go-ipfs/net/service"
5
6 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
8 - netmsg "github.com/jbenet/go-ipfs/net/message"
7 peer "github.com/jbenet/go-ipfs/peer"
8 u "github.com/jbenet/go-ipfs/util"
9 )
@@ -38,13 +36,6 @@ type Receiver interface {
36 ReceiveError(error)
37 }
38
41 -// TODO(brian): move this to go-ipfs/net package
42 -type NetMessageService interface {
43 - SendRequest(ctx context.Context, m netmsg.NetMessage) (netmsg.NetMessage, error)
44 - SendMessage(ctx context.Context, m netmsg.NetMessage) error
45 - SetHandler(netservice.Handler)
46 -}
47 -
39 // TODO rename -> Router?
40 type Routing interface {
41 // FindProvidersAsync returns a channel of providers for the given key
exchange/bitswap/network/net_message_adapter.go
+3 -2
@@ -4,12 +4,13 @@ import (
4 context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
5
6 bsmsg "github.com/jbenet/go-ipfs/exchange/bitswap/message"
7 + inet "github.com/jbenet/go-ipfs/net"
8 netmsg "github.com/jbenet/go-ipfs/net/message"
9 peer "github.com/jbenet/go-ipfs/peer"
10 )
11
12 // NetMessageAdapter wraps a NetMessage network service
12 -func NetMessageAdapter(s NetMessageService, r Receiver) Adapter {
13 +func NetMessageAdapter(s inet.Service, r Receiver) Adapter {
14 adapter := impl{
15 nms: s,
16 receiver: r,
@@ -20,7 +21,7 @@ func NetMessageAdapter(s NetMessageService, r Receiver) Adapter {
21
22 // implements an Adapter that integrates with a NetMessage network service
23 type impl struct {
23 - nms NetMessageService
24 + nms inet.Service
25
26 // inbound messages from the network are forwarded to the receiver
27 receiver Receiver
net/interface.go
+4 -10
@@ -5,8 +5,6 @@ import (
5 mux "github.com/jbenet/go-ipfs/net/mux"
6 srv "github.com/jbenet/go-ipfs/net/service"
7 peer "github.com/jbenet/go-ipfs/peer"
8 -
9 - context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
8 )
9
10 // Network is the interface IPFS uses for connecting to the world.
@@ -39,14 +37,10 @@ type Network interface {
37 }
38
39 // Sender interface for network services.
42 -type Sender interface {
43 - // SendMessage sends out a given message, without expecting a response.
44 - SendMessage(ctx context.Context, m msg.NetMessage) error
45 -
46 - // SendRequest sends out a given message, and awaits a response.
47 - // Set Deadlines or cancellations in the context.Context you pass in.
48 - SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMessage, error)
49 -}
40 +type Sender srv.Sender
41
42 // Handler interface for network services.
43 type Handler srv.Handler
44 +
45 +// Service interface for network resources.
46 +type Service srv.Service
net/service/service.go
+38 -12
@@ -25,9 +25,35 @@ type Handler interface {
25 HandleMessage(context.Context, msg.NetMessage) msg.NetMessage
26 }
27
28 +// Sender interface for network services.
29 +type Sender interface {
30 + // SendMessage sends out a given message, without expecting a response.
31 + SendMessage(ctx context.Context, m msg.NetMessage) error
32 +
33 + // SendRequest sends out a given message, and awaits a response.
34 + // Set Deadlines or cancellations in the context.Context you pass in.
35 + SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMessage, error)
36 +}
37 +
38 +// Service is an interface for a net resource with both outgoing (sender) and
39 +// incomig (SetHandler) requests.
40 +type Service interface {
41 + Sender
42 +
43 + // Start + Stop Service
44 + Start(ctx context.Context) error
45 + Stop()
46 +
47 + // GetPipe
48 + GetPipe() *msg.Pipe
49 +
50 + // SetHandler assigns the request Handler for this service.
51 + SetHandler(Handler)
52 +}
53 +
54 // Service is a networking component that protocols can use to multiplex
55 // messages over the same channel, and to issue + handle requests.
30 -type Service struct {
56 +type service struct {
57 // Handler is the object registered to handle incoming requests.
58 Handler Handler
59
@@ -43,8 +69,8 @@ type Service struct {
69 }
70
71 // NewService creates a service object with given type ID and Handler
46 -func NewService(h Handler) *Service {
47 - return &Service{
72 +func NewService(h Handler) Service {
73 + return &service{
74 Handler: h,
75 Requests: RequestMap{},
76 Pipe: msg.NewPipe(10),
@@ -52,7 +78,7 @@ func NewService(h Handler) *Service {
78 }
79
80 // Start kicks off the Service goroutines.
55 -func (s *Service) Start(ctx context.Context) error {
81 +func (s *service) Start(ctx context.Context) error {
82 if s.cancel != nil {
83 return errors.New("Service already started.")
84 }
@@ -65,18 +91,18 @@ func (s *Service) Start(ctx context.Context) error {
91 }
92
93 // Stop stops Service activity.
68 -func (s *Service) Stop() {
94 +func (s *service) Stop() {
95 s.cancel()
96 s.cancel = context.CancelFunc(nil)
97 }
98
99 // GetPipe implements the mux.Protocol interface
74 -func (s *Service) GetPipe() *msg.Pipe {
100 +func (s *service) GetPipe() *msg.Pipe {
101 return s.Pipe
102 }
103
104 // sendMessage sends a message out (actual leg work. SendMessage is to export w/o rid)
79 -func (s *Service) sendMessage(ctx context.Context, m msg.NetMessage, rid RequestID) error {
105 +func (s *service) sendMessage(ctx context.Context, m msg.NetMessage, rid RequestID) error {
106
107 // serialize ServiceMessage wrapper
108 data, err := wrapData(m.Data(), rid)
@@ -98,12 +124,12 @@ func (s *Service) sendMessage(ctx context.Context, m msg.NetMessage, rid Request
124 }
125
126 // SendMessage sends a message out
101 -func (s *Service) SendMessage(ctx context.Context, m msg.NetMessage) error {
127 +func (s *service) SendMessage(ctx context.Context, m msg.NetMessage) error {
128 return s.sendMessage(ctx, m, nil)
129 }
130
131 // SendRequest sends a request message out and awaits a response.
106 -func (s *Service) SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMessage, error) {
132 +func (s *service) SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMessage, error) {
133
134 // create a request
135 r, err := NewRequest(m.Peer().ID)
@@ -151,7 +177,7 @@ func (s *Service) SendRequest(ctx context.Context, m msg.NetMessage) (msg.NetMes
177
178 // handleIncoming consumes the messages on the s.Incoming channel and
179 // routes them appropriately (to requests, or handler).
154 -func (s *Service) handleIncomingMessages(ctx context.Context) {
180 +func (s *service) handleIncomingMessages(ctx context.Context) {
181 for {
182 select {
183 case m, more := <-s.Incoming:
@@ -166,7 +192,7 @@ func (s *Service) handleIncomingMessages(ctx context.Context) {
192 }
193 }
194
169 -func (s *Service) handleIncomingMessage(ctx context.Context, m msg.NetMessage) {
195 +func (s *service) handleIncomingMessage(ctx context.Context, m msg.NetMessage) {
196
197 // unwrap the incoming message
198 data, rid, err := unwrapData(m.Data())
@@ -217,6 +243,6 @@ func (s *Service) handleIncomingMessage(ctx context.Context, m msg.NetMessage) {
243 }
244
245 // SetHandler assigns the request Handler for this service.
220 -func (s *Service) SetHandler(h Handler) {
246 +func (s *service) SetHandler(h Handler) {
247 s.Handler = h
248 }