110 lines
3.4 KiB
Diff
110 lines
3.4 KiB
Diff
|
diff --git a/whisper/whisperv6/doc.go b/whisper/whisperv6/doc.go
|
||
|
index 5ad660616..9659e6c46 100644
|
||
|
--- a/whisper/whisperv6/doc.go
|
||
|
+++ b/whisper/whisperv6/doc.go
|
||
|
@@ -109,3 +109,40 @@ type NotificationServer interface {
|
||
|
// Stop stops notification sending loop, releasing related resources
|
||
|
Stop() error
|
||
|
}
|
||
|
+
|
||
|
+type envelopeSource int
|
||
|
+
|
||
|
+const (
|
||
|
+ _ = iota
|
||
|
+ // peerSource indicates a source as a regular peer.
|
||
|
+ peerSource envelopeSource = iota
|
||
|
+ // p2pSource indicates that envelop was received from a trusted peer.
|
||
|
+ p2pSource
|
||
|
+)
|
||
|
+
|
||
|
+// EnvelopeMeta keeps metadata of received envelopes.
|
||
|
+type EnvelopeMeta struct {
|
||
|
+ Hash string
|
||
|
+ Topic TopicType
|
||
|
+ Size uint32
|
||
|
+ Source envelopeSource
|
||
|
+ IsNew bool
|
||
|
+ Peer string
|
||
|
+}
|
||
|
+
|
||
|
+// SourceString converts source to string.
|
||
|
+func (m *EnvelopeMeta) SourceString() string {
|
||
|
+ switch m.Source {
|
||
|
+ case peerSource:
|
||
|
+ return "peer"
|
||
|
+ case p2pSource:
|
||
|
+ return "p2p"
|
||
|
+ default:
|
||
|
+ return "unknown"
|
||
|
+ }
|
||
|
+}
|
||
|
+
|
||
|
+// EnvelopeTracer tracks received envelopes.
|
||
|
+type EnvelopeTracer interface {
|
||
|
+ Trace(*EnvelopeMeta)
|
||
|
+}
|
||
|
diff --git a/whisper/whisperv6/whisper.go b/whisper/whisperv6/whisper.go
|
||
|
index 54d7d0f24..ce9405dff 100644
|
||
|
--- a/whisper/whisperv6/whisper.go
|
||
|
+++ b/whisper/whisperv6/whisper.go
|
||
|
@@ -87,6 +87,7 @@ type Whisper struct {
|
||
|
|
||
|
mailServer MailServer // MailServer interface
|
||
|
notificationServer NotificationServer
|
||
|
+ envelopeTracer EnvelopeTracer // Service collecting envelopes metadata
|
||
|
}
|
||
|
|
||
|
// New creates a Whisper client ready to communicate through the Ethereum P2P network.
|
||
|
@@ -215,6 +216,12 @@ func (whisper *Whisper) RegisterNotificationServer(server NotificationServer) {
|
||
|
whisper.notificationServer = server
|
||
|
}
|
||
|
|
||
|
+// RegisterEnvelopeTracer registers an EnveloperTracer to collect information
|
||
|
+// about received envelopes.
|
||
|
+func (whisper *Whisper) RegisterEnvelopeTracer(tracer EnvelopeTracer) {
|
||
|
+ whisper.envelopeTracer = tracer
|
||
|
+}
|
||
|
+
|
||
|
// Protocols returns the whisper sub-protocols ran by this particular client.
|
||
|
func (whisper *Whisper) Protocols() []p2p.Protocol {
|
||
|
return []p2p.Protocol{whisper.protocol}
|
||
|
@@ -756,6 +763,7 @@ func (whisper *Whisper) runMessageLoop(p *Peer, rw p2p.MsgReadWriter) error {
|
||
|
|
||
|
trouble := false
|
||
|
for _, env := range envelopes {
|
||
|
+ whisper.traceEnvelope(env, !whisper.isEnvelopeCached(env.Hash()), peerSource, p)
|
||
|
cached, err := whisper.add(env)
|
||
|
if err != nil {
|
||
|
trouble = true
|
||
|
@@ -810,6 +818,7 @@ func (whisper *Whisper) runMessageLoop(p *Peer, rw p2p.MsgReadWriter) error {
|
||
|
return errors.New("invalid direct message")
|
||
|
}
|
||
|
whisper.postEvent(&envelope, true)
|
||
|
+ whisper.traceEnvelope(&envelope, false, p2pSource, p)
|
||
|
}
|
||
|
case p2pRequestCode:
|
||
|
// Must be processed if mail server is implemented. Otherwise ignore.
|
||
|
@@ -906,6 +915,22 @@ func (whisper *Whisper) add(envelope *Envelope) (bool, error) {
|
||
|
return true, nil
|
||
|
}
|
||
|
|
||
|
+// traceEnvelope collects basic metadata about an envelope and sender peer.
|
||
|
+func (whisper *Whisper) traceEnvelope(envelope *Envelope, isNew bool, source envelopeSource, peer *Peer) {
|
||
|
+ if whisper.envelopeTracer == nil {
|
||
|
+ return
|
||
|
+ }
|
||
|
+
|
||
|
+ whisper.envelopeTracer.Trace(&EnvelopeMeta{
|
||
|
+ Hash: envelope.Hash().String(),
|
||
|
+ Topic: BytesToTopic(envelope.Topic[:]),
|
||
|
+ Size: uint32(envelope.size()),
|
||
|
+ Source: source,
|
||
|
+ IsNew: isNew,
|
||
|
+ Peer: peer.peer.Info().ID,
|
||
|
+ })
|
||
|
+}
|
||
|
+
|
||
|
// postEvent queues the message for further processing.
|
||
|
func (whisper *Whisper) postEvent(envelope *Envelope, isP2P bool) {
|
||
|
if isP2P {
|