2019-08-16 17:13:04 +00:00
|
|
|
package sync
|
|
|
|
|
|
|
|
import (
|
|
|
|
"context"
|
|
|
|
"fmt"
|
2019-08-24 01:15:02 +00:00
|
|
|
"runtime/debug"
|
2019-08-22 05:34:25 +00:00
|
|
|
"time"
|
2019-08-16 17:13:04 +00:00
|
|
|
|
|
|
|
"github.com/gogo/protobuf/proto"
|
|
|
|
"github.com/prysmaticlabs/prysm/beacon-chain/p2p"
|
2019-10-03 16:33:16 +00:00
|
|
|
"github.com/prysmaticlabs/prysm/shared/hashutil"
|
2019-09-17 21:14:51 +00:00
|
|
|
"github.com/prysmaticlabs/prysm/shared/roughtime"
|
2019-09-29 18:48:55 +00:00
|
|
|
"github.com/prysmaticlabs/prysm/shared/traceutil"
|
2019-09-03 18:06:35 +00:00
|
|
|
"go.opencensus.io/trace"
|
2019-08-16 17:13:04 +00:00
|
|
|
)
|
|
|
|
|
2019-08-28 15:14:22 +00:00
|
|
|
const oneYear = 365 * 24 * time.Hour
|
2019-09-20 17:54:32 +00:00
|
|
|
const pubsubMessageTimeout = 10 * time.Second
|
2019-08-22 05:34:25 +00:00
|
|
|
|
|
|
|
// prefix to add to keys, so that we can represent invalid objects
|
2019-08-28 15:14:22 +00:00
|
|
|
const invalid = "invalidObject"
|
2019-08-22 05:34:25 +00:00
|
|
|
|
2019-08-16 17:13:04 +00:00
|
|
|
// subHandler represents handler for a given subscription.
|
|
|
|
type subHandler func(context.Context, proto.Message) error
|
|
|
|
|
|
|
|
// validator should verify the contents of the message, propagate the message
|
|
|
|
// as expected, and return true or false to continue the message processing
|
2019-09-04 00:22:15 +00:00
|
|
|
// pipeline. FromSelf indicates whether or not this is a message received from our
|
|
|
|
// node in pubsub.
|
2019-10-01 15:13:04 +00:00
|
|
|
type validator func(ctx context.Context, msg proto.Message, broadcaster p2p.Broadcaster, fromSelf bool) (bool, error)
|
2019-08-16 17:13:04 +00:00
|
|
|
|
|
|
|
// noopValidator is a no-op that always returns true and does not propagate any
|
|
|
|
// message.
|
2019-10-01 15:13:04 +00:00
|
|
|
func noopValidator(_ context.Context, _ proto.Message, _ p2p.Broadcaster, _ bool) (bool, error) {
|
|
|
|
return true, nil
|
2019-08-16 17:13:04 +00:00
|
|
|
}
|
|
|
|
|
|
|
|
// Register PubSub subscribers
|
|
|
|
func (r *RegularSync) registerSubscribers() {
|
2019-09-17 21:14:51 +00:00
|
|
|
go func() {
|
|
|
|
ch := make(chan time.Time)
|
|
|
|
sub := r.chain.StateInitializedFeed().Subscribe(ch)
|
|
|
|
defer sub.Unsubscribe()
|
2019-09-20 17:08:32 +00:00
|
|
|
|
2019-09-17 21:14:51 +00:00
|
|
|
// Wait until chain start.
|
|
|
|
genesis := <-ch
|
|
|
|
if genesis.After(roughtime.Now()) {
|
|
|
|
time.Sleep(roughtime.Until(genesis))
|
|
|
|
}
|
|
|
|
r.chainStarted = true
|
|
|
|
}()
|
2019-08-16 17:13:04 +00:00
|
|
|
r.subscribe(
|
|
|
|
"/eth2/beacon_block",
|
2019-08-22 18:11:52 +00:00
|
|
|
r.validateBeaconBlockPubSub,
|
|
|
|
r.beaconBlockSubscriber,
|
2019-08-16 17:13:04 +00:00
|
|
|
)
|
|
|
|
r.subscribe(
|
|
|
|
"/eth2/beacon_attestation",
|
2019-08-23 19:46:04 +00:00
|
|
|
r.validateBeaconAttestation,
|
2019-08-23 21:34:03 +00:00
|
|
|
r.beaconAttestationSubscriber,
|
2019-08-16 17:13:04 +00:00
|
|
|
)
|
|
|
|
r.subscribe(
|
|
|
|
"/eth2/voluntary_exit",
|
2019-08-18 15:33:58 +00:00
|
|
|
r.validateVoluntaryExit,
|
|
|
|
r.voluntaryExitSubscriber,
|
2019-08-16 17:13:04 +00:00
|
|
|
)
|
|
|
|
r.subscribe(
|
|
|
|
"/eth2/proposer_slashing",
|
2019-08-22 14:31:56 +00:00
|
|
|
r.validateProposerSlashing,
|
|
|
|
r.proposerSlashingSubscriber,
|
2019-08-16 17:13:04 +00:00
|
|
|
)
|
|
|
|
r.subscribe(
|
|
|
|
"/eth2/attester_slashing",
|
2019-08-22 05:34:25 +00:00
|
|
|
r.validateAttesterSlashing,
|
|
|
|
r.attesterSlashingSubscriber,
|
2019-08-16 17:13:04 +00:00
|
|
|
)
|
|
|
|
}
|
|
|
|
|
|
|
|
// subscribe to a given topic with a given validator and subscription handler.
|
|
|
|
// The base protobuf message is used to initialize new messages for decoding.
|
|
|
|
func (r *RegularSync) subscribe(topic string, validate validator, handle subHandler) {
|
|
|
|
base := p2p.GossipTopicMappings[topic]
|
|
|
|
if base == nil {
|
|
|
|
panic(fmt.Sprintf("%s is not mapped to any message in GossipTopicMappings", topic))
|
|
|
|
}
|
|
|
|
|
|
|
|
topic += r.p2p.Encoding().ProtocolSuffix()
|
|
|
|
log := log.WithField("topic", topic)
|
|
|
|
|
|
|
|
sub, err := r.p2p.PubSub().Subscribe(topic)
|
|
|
|
if err != nil {
|
|
|
|
// Any error subscribing to a PubSub topic would be the result of a misconfiguration of
|
|
|
|
// libp2p PubSub library. This should not happen at normal runtime, unless the config
|
|
|
|
// changes to a fatal configuration.
|
|
|
|
panic(err)
|
|
|
|
}
|
|
|
|
|
|
|
|
// Pipeline decodes the incoming subscription data, runs the validation, and handles the
|
|
|
|
// message.
|
2019-09-04 00:22:15 +00:00
|
|
|
pipeline := func(data []byte, fromSelf bool) {
|
2019-09-29 18:48:55 +00:00
|
|
|
ctx, _ := context.WithTimeout(context.Background(), pubsubMessageTimeout)
|
|
|
|
ctx, span := trace.StartSpan(ctx, "sync.pubsub")
|
|
|
|
defer span.End()
|
|
|
|
|
2019-08-24 01:15:02 +00:00
|
|
|
defer func() {
|
|
|
|
if r := recover(); r != nil {
|
2019-09-29 18:48:55 +00:00
|
|
|
traceutil.AnnotateError(span, fmt.Errorf("panic occurred: %v", r))
|
2019-08-24 01:15:02 +00:00
|
|
|
log.WithField("error", r).Error("Panic occurred")
|
|
|
|
debug.PrintStack()
|
|
|
|
}
|
|
|
|
}()
|
2019-09-29 18:48:55 +00:00
|
|
|
|
2019-09-03 18:06:35 +00:00
|
|
|
span.AddAttributes(trace.StringAttribute("topic", topic))
|
2019-09-29 18:48:55 +00:00
|
|
|
span.AddAttributes(trace.BoolAttribute("fromSelf", fromSelf))
|
2019-08-24 12:50:43 +00:00
|
|
|
|
2019-08-16 17:13:04 +00:00
|
|
|
if data == nil {
|
|
|
|
log.Warn("Received nil message on pubsub")
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2019-10-03 16:33:16 +00:00
|
|
|
if span.IsRecordingEvents() {
|
|
|
|
id := hashutil.FastSum64(data)
|
|
|
|
messageLen := int64(len(data))
|
|
|
|
span.AddMessageReceiveEvent(int64(id), messageLen /*uncompressed*/, messageLen /*compressed*/)
|
|
|
|
}
|
|
|
|
|
2019-08-16 17:13:04 +00:00
|
|
|
msg := proto.Clone(base)
|
2019-09-09 02:34:52 +00:00
|
|
|
if err := r.p2p.Encoding().Decode(data, msg); err != nil {
|
2019-09-29 18:48:55 +00:00
|
|
|
traceutil.AnnotateError(span, err)
|
2019-08-16 17:13:04 +00:00
|
|
|
log.WithError(err).Warn("Failed to decode pubsub message")
|
|
|
|
return
|
|
|
|
}
|
|
|
|
|
2019-10-01 15:13:04 +00:00
|
|
|
valid, err := validate(ctx, msg, r.p2p, fromSelf)
|
|
|
|
if err != nil {
|
2019-09-29 18:48:55 +00:00
|
|
|
if !fromSelf {
|
2019-10-01 15:13:04 +00:00
|
|
|
log.WithError(err).Error("Message failed to verify")
|
2019-09-30 05:23:19 +00:00
|
|
|
messageFailedValidationCounter.WithLabelValues(topic).Inc()
|
2019-09-29 18:48:55 +00:00
|
|
|
}
|
2019-08-16 17:13:04 +00:00
|
|
|
return
|
|
|
|
}
|
2019-10-01 15:13:04 +00:00
|
|
|
if !valid {
|
|
|
|
return
|
|
|
|
}
|
2019-08-16 17:13:04 +00:00
|
|
|
|
2019-09-03 18:06:35 +00:00
|
|
|
if err := handle(ctx, msg); err != nil {
|
2019-09-29 18:48:55 +00:00
|
|
|
traceutil.AnnotateError(span, err)
|
2019-08-16 17:13:04 +00:00
|
|
|
log.WithError(err).Error("Failed to handle p2p pubsub")
|
2019-09-30 05:23:19 +00:00
|
|
|
messageFailedProcessingCounter.WithLabelValues(topic).Inc()
|
2019-08-16 17:13:04 +00:00
|
|
|
return
|
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
// The main message loop for receiving incoming messages from this subscription.
|
|
|
|
messageLoop := func() {
|
|
|
|
for {
|
|
|
|
msg, err := sub.Next(r.ctx)
|
|
|
|
if err != nil {
|
|
|
|
log.WithError(err).Error("Subscription next failed")
|
|
|
|
// TODO(3147): Mark status unhealthy.
|
|
|
|
return
|
|
|
|
}
|
2019-09-17 21:14:51 +00:00
|
|
|
if !r.chainStarted {
|
|
|
|
messageReceivedBeforeChainStartCounter.WithLabelValues(topic + r.p2p.Encoding().ProtocolSuffix()).Inc()
|
|
|
|
continue
|
|
|
|
}
|
2019-09-04 00:22:15 +00:00
|
|
|
// Special validation occurs on messages received from ourselves.
|
|
|
|
fromSelf := msg.GetFrom() == r.p2p.PeerID()
|
2019-08-24 18:41:24 +00:00
|
|
|
|
2019-08-28 15:14:22 +00:00
|
|
|
messageReceivedCounter.WithLabelValues(topic + r.p2p.Encoding().ProtocolSuffix()).Inc()
|
|
|
|
|
2019-09-04 00:22:15 +00:00
|
|
|
go pipeline(msg.Data, fromSelf)
|
2019-08-16 17:13:04 +00:00
|
|
|
}
|
|
|
|
}
|
|
|
|
|
|
|
|
go messageLoop()
|
|
|
|
}
|