@cryptotaxi247 / kubo / commits / d5331e7dc

feat(gcr/s) add eventlogs

Brian Tiger Chow committed Jan 28, 2015 at 08:06 UTC d5331e7dc761d52e17a2cb66306b37607fb9066a
3 files changed +32 -8
routing/grandcentral/client.go
+6
@@ -36,6 +36,7 @@ func NewClient(px proxy.Proxy, h host.Host, ps peer.Peerstore, local peer.ID) (*
36 }
37
38 func (c *Client) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-chan peer.PeerInfo {
39 + defer log.EventBegin(ctx, "findProviders", &k).Done()
40 ch := make(chan peer.PeerInfo)
41 go func() {
42 defer close(ch)
@@ -58,6 +59,7 @@ func (c *Client) FindProvidersAsync(ctx context.Context, k u.Key, max int) <-cha
59 }
60
61 func (c *Client) PutValue(ctx context.Context, k u.Key, v []byte) error {
62 + defer log.EventBegin(ctx, "putValue", &k).Done()
63 r, err := makeRecord(c.peerstore, c.local, k, v)
64 if err != nil {
65 return err
@@ -68,6 +70,7 @@ func (c *Client) PutValue(ctx context.Context, k u.Key, v []byte) error {
70 }
71
72 func (c *Client) GetValue(ctx context.Context, k u.Key) ([]byte, error) {
73 + defer log.EventBegin(ctx, "getValue", &k).Done()
74 msg := pb.NewMessage(pb.Message_GET_VALUE, string(k), 0)
75 response, err := c.proxy.SendRequest(ctx, msg) // TODO wrap to hide the remote
76 if err != nil {
@@ -77,6 +80,7 @@ func (c *Client) GetValue(ctx context.Context, k u.Key) ([]byte, error) {
80 }
81
82 func (c *Client) Provide(ctx context.Context, k u.Key) error {
83 + defer log.EventBegin(ctx, "provide", &k).Done()
84 msg := pb.NewMessage(pb.Message_ADD_PROVIDER, string(k), 0)
85 // FIXME how is connectedness defined for the local node
86 pri := []pb.PeerRoutingInfo{
@@ -92,6 +96,7 @@ func (c *Client) Provide(ctx context.Context, k u.Key) error {
96 }
97
98 func (c *Client) FindPeer(ctx context.Context, id peer.ID) (peer.PeerInfo, error) {
99 + defer log.EventBegin(ctx, "findPeer", id).Done()
100 request := pb.NewMessage(pb.Message_FIND_NODE, string(id), 0)
101 response, err := c.proxy.SendRequest(ctx, request) // hide remote
102 if err != nil {
@@ -121,6 +126,7 @@ func makeRecord(ps peer.Peerstore, p peer.ID, k u.Key, v []byte) (*pb.Record, er
126 }
127
128 func (c *Client) Ping(ctx context.Context, id peer.ID) (time.Duration, error) {
129 + defer log.EventBegin(ctx, "ping", id).Done()
130 return time.Nanosecond, errors.New("grandcentral routing does not support the ping method")
131 }
132
routing/grandcentral/proxy/standard.go
+22 -8
@@ -13,10 +13,10 @@ import (
13 errors "github.com/jbenet/go-ipfs/util/debugerror"
14 )
15
16 -var log = eventlog.Logger("proxy")
17 -
16 const ProtocolGCR = "/ipfs/grandcentral"
17
18 +var log = eventlog.Logger("grandcentral/proxy")
19 +
20 type Proxy interface {
21 HandleStream(inet.Stream)
22 SendMessage(ctx context.Context, m *dhtpb.Message) error
@@ -48,8 +48,15 @@ func (px *standard) SendMessage(ctx context.Context, m *dhtpb.Message) error {
48 return err // NB: returns the last error
49 }
50
51 -func (px *standard) sendMessage(ctx context.Context, m *dhtpb.Message, remote peer.ID) error {
52 - if err := px.Host.Connect(ctx, peer.PeerInfo{ID: remote}); err != nil {
51 +func (px *standard) sendMessage(ctx context.Context, m *dhtpb.Message, remote peer.ID) (err error) {
52 + e := log.EventBegin(ctx, "sendRoutingMessage", px.Host.ID(), remote, m)
53 + defer func() {
54 + if err != nil {
55 + e.SetError(err)
56 + }
57 + e.Done()
58 + }()
59 + if err = px.Host.Connect(ctx, peer.PeerInfo{ID: remote}); err != nil {
60 return err
61 }
62 s, err := px.Host.NewStream(ProtocolGCR, remote)
@@ -78,8 +85,15 @@ func (px *standard) SendRequest(ctx context.Context, m *dhtpb.Message) (*dhtpb.M
85 return nil, err // NB: returns the last error
86 }
87
81 -func (px *standard) sendRequest(ctx context.Context, m *dhtpb.Message, remote peer.ID) (*dhtpb.Message, error) {
82 - if err := px.Host.Connect(ctx, peer.PeerInfo{ID: remote}); err != nil {
88 +func (px *standard) sendRequest(ctx context.Context, m *dhtpb.Message, remote peer.ID) (_ *dhtpb.Message, err error) {
89 + e := log.EventBegin(ctx, "sendRoutingRequest", px.Host.ID(), remote, m)
90 + defer func() {
91 + if err != nil {
92 + e.SetError(err)
93 + }
94 + e.Done()
95 + }()
96 + if err = px.Host.Connect(ctx, peer.PeerInfo{ID: remote}); err != nil {
97 return nil, err
98 }
99 s, err := px.Host.NewStream(ProtocolGCR, remote)
@@ -89,12 +103,12 @@ func (px *standard) sendRequest(ctx context.Context, m *dhtpb.Message, remote pe
103 defer s.Close()
104 r := ggio.NewDelimitedReader(s, inet.MessageSizeMax)
105 w := ggio.NewDelimitedWriter(s)
92 - if err := w.WriteMsg(m); err != nil {
106 + if err = w.WriteMsg(m); err != nil {
107 return nil, err
108 }
109
110 var reply dhtpb.Message
97 - if err := r.ReadMsg(&reply); err != nil {
111 + if err = r.ReadMsg(&reply); err != nil {
112 return nil, err
113 }
114 // need ctx expiration?
routing/grandcentral/server.go
+4
@@ -42,6 +42,8 @@ func (s *Server) HandleRequest(ctx context.Context, p peer.ID, req *dhtpb.Messag
42 func (s *Server) handleMessage(
43 ctx context.Context, p peer.ID, req *dhtpb.Message) (peer.ID, *dhtpb.Message) {
44
45 + log.EventBegin(ctx, "routingMessageReceived", req, p, s.local).Done() // TODO may need to differentiate between local and remote
46 +
47 // FIXME threw everything into this switch statement to get things going.
48 // Once each operation is well-defined, extract pluggable backend so any
49 // database may be used.
@@ -131,6 +133,7 @@ func putRoutingRecord(ds datastore.Datastore, k util.Key, value *dhtpb.Record) e
133 }
134
135 func putRoutingProviders(ds datastore.Datastore, k util.Key, providers []*dhtpb.Message_Peer) error {
136 + log.Event(context.Background(), "putRoutingProviders", &k)
137 pkey := datastore.KeyWithNamespaces([]string{"routing", "providers", k.String()})
138 if v, err := ds.Get(pkey); err == nil {
139 if msg, ok := v.([]byte); ok {
@@ -166,6 +169,7 @@ func storeProvidersToPeerstore(ps peer.Peerstore, p peer.ID, providers []*dhtpb.
169 }
170
171 func getRoutingProviders(local peer.ID, ds datastore.Datastore, k util.Key) ([]*dhtpb.Message_Peer, error) {
172 + log.Event(context.Background(), "getProviders", local, &k)
173 var providers []*dhtpb.Message_Peer
174 exists, err := ds.Has(k.DsKey()) // TODO store values in a local datastore?
175 if err == nil && exists {