stash
Juan Batiz-Benet committed
Dec 14, 2014 at 09:28 UTC
9d304768fc1292ffff9a334e0957aaa99039d8e1
1 file changed
+77
-35
net/service/service.go
+77
-35
@@ -7,9 +7,10 @@ import (
7
8
msg "github.com/jbenet/go-ipfs/net/message"
9
u "github.com/jbenet/go-ipfs/util"
10
- ctxc "github.com/jbenet/go-ipfs/util/ctxcloser"
10
11
+ ctxgroup "github.com/jbenet/go-ctxgroup"
12
context "github.com/jbenet/go-ipfs/Godeps/_workspace/src/code.google.com/p/go.net/context"
13
+ router "github.com/jbenet/go-router"
14
)
15
16
var log = u.Logger("service")
@@ -40,11 +41,12 @@ type Sender interface {
41
// Service is an interface for a net resource with both outgoing (sender) and
42
// incomig (SetHandler) requests.
43
type Service interface {
43
- Sender
44
- ctxc.ContextCloser
44
+ Sender // can use it to send out msgs
45
+ router.Node // it is a Node in the net topology.
46
46
- // GetPipe
47
- GetPipe() *msg.Pipe
47
+ // SetUplink assigns the Node to send packets out
48
+ SetUplink(router.Node)
49
+ Uplink() router.Node
50
51
// SetHandler assigns the request Handler for this service.
52
SetHandler(Handler)
@@ -62,30 +64,23 @@ type service struct {
64
Requests RequestMap
65
RequestsLock sync.RWMutex
66
65
- // Message Pipe (connected to the outside world)
66
- *msg.Pipe
67
- ctxc.ContextCloser
67
+ // the connection to the outside world
68
+ uplink router.Node
69
+ uplinkLock sync.RWMutex
70
+ addr router.Address
71
}
72
73
// NewService creates a service object with given type ID and Handler
71
-func NewService(ctx context.Context, h Handler) Service {
74
+func NewService(addr router.Address, uplink router.Node, h Handler) Service {
75
s := &service{
73
- Handler: h,
74
- Requests: RequestMap{},
75
- Pipe: msg.NewPipe(10),
76
- ContextCloser: ctxc.NewContextCloser(ctx, nil),
76
+ Handler: h,
77
+ Requests: RequestMap{},
78
+ uplink: uplink,
79
+ addr: addr,
80
}
78
-
79
- s.Children().Add(1)
80
- go s.handleIncomingMessages()
81
return s
82
}
83
84
-// GetPipe implements the mux.Protocol interface
85
-func (s *service) GetPipe() *msg.Pipe {
86
- return s.Pipe
87
-}
88
-
84
// sendMessage sends a message out (actual leg work. SendMessage is to export w/o rid)
85
func (s *service) sendMessage(ctx context.Context, m msg.NetMessage, rid RequestID) error {
86
@@ -99,12 +94,21 @@ func (s *service) sendMessage(ctx context.Context, m msg.NetMessage, rid Request
94
95
// send message
96
m2 := msg.New(m.Peer(), data)
97
+
98
+ pkt := msg.Packet{
99
+ Src:
100
+ }
101
+
102
select {
103
case s.Outgoing <- m2:
104
case <-ctx.Done():
105
return ctx.Err()
106
}
107
108
+ pkt := msg.Packet{
109
+ Src: m.
110
+ }
111
+
112
return nil
113
}
114
@@ -193,34 +197,57 @@ func (s *service) handleIncomingMessages() {
197
}
198
}
199
196
-func (s *service) handleIncomingMessage(m msg.NetMessage) {
197
- defer s.Children().Done()
200
+func (s *service) handleIncomingMessage(pkt *msg.Packet) error {
201
+
202
+ // check the packet has a valid Context
203
+ ctx := pkt.Context
204
+ if ctx == nil {
205
+ return fmt.Errorf("service got pkt without valid Context")
206
+ }
207
+
208
+ // check the source is a peer
209
+ srcPeer, ok := pkt.Src.(peer.Peer)
210
+ if !ok {
211
+ return fmt.Errorf("service got pkt from non-Peer src: %v", pkt.Src)
212
+ }
213
214
// unwrap the incoming message
200
- data, rid, err := unwrapData(m.Data())
215
+ data, rid, err := unwrapData(pkt.Data)
216
if err != nil {
202
- log.Errorf("service de-serializing error: %v", err)
203
- return
217
+ return fmt.Errorf("service de-serializing error: %v", err)
218
}
219
206
- m2 := msg.New(m.Peer(), data)
220
+ // convert to msg.NetMessage, which the rest of the system expects.
221
+ m2 := msg.New(srcPeer, data)
222
223
// if it's a request (or has no RequestID), handle it
224
if rid == nil || rid.IsRequest() {
225
handler := s.GetHandler()
226
if handler == nil {
227
log.Errorf("service dropped msg: %v", m)
213
- return // no handler, drop it.
228
+ log.Event()
229
+ return nil
230
+ // no handler, drop it.
231
}
232
216
- // should this be "go HandleMessage ... ?"
217
- r1 := handler.HandleMessage(s.Context(), m2)
218
-
219
- // if handler gave us a response, send it back out!
233
+ // this go routine is developer friendliness to keep their stacks
234
+ // separate (and more readable) from the network goroutine. If
235
+ // problems arise and you'd like to see _the full_ stack of where
236
+ // this message is coming from, just remove the goroutine part.
237
+ response := make(chan msg.NetMessage)
238
+ go func() msg.NetMessage {
239
+ return handler.HandleMessage(ctx, m2)
240
+ }()
241
+ r1 := <-response
242
+ // Note: HandleMessage *must* respect context. We could co-opt it
243
+ // and do a select {} here on the context, BUT that would just drop
244
+ // a packet and free up the goroutine to return to the network. the
245
+ // problem is still there: the Service handler hasn't returned yet.
246
+
247
+ // if handler gave us a response, send it out!
248
if r1 != nil {
221
- err := s.sendMessage(s.Context(), r1, rid.Response())
222
- if err != nil {
223
- log.Errorf("error sending response message: %v", err)
249
+ if err := s.sendMessage(ctx, r1, rid.Response()); err != nil {
250
+ return fmt.Errorf("error sending response message: %v", err)
251
}
252
}
253
return
@@ -247,6 +274,21 @@ func (s *service) handleIncomingMessage(m msg.NetMessage) {
274
}
275
}
276
277
+
278
+// Address is the router.Node address
279
+func (s *service) Address() router.Address {
280
+ return s.addr
281
+}
282
+
283
+// HandlePacket implements router.Node
284
+// service only receives packets in HandlePacket
285
+func (s *service) HandlePacket(p router.Packet, from router.Node) error {
286
+ pkt, ok := p.(*msg.Packet)
287
+ if !ok {
288
+ return msg.ErrInvalidPayload
289
+ }
290
+}
291
+
292
// SetHandler assigns the request Handler for this service.
293
func (s *service) SetHandler(h Handler) {
294
s.HandlerLock.Lock()