Files

652 lines
15 KiB
Go
Raw Permalink Normal View History

2019-08-05 21:55:50 +02:00
// Package node contains node logic.
2019-06-10 12:13:37 -07:00
package node
2019-04-20 20:09:20 +02:00
2019-04-30 18:49:39 +02:00
// @todo this is a very rough implementation that needs cleanup
2019-04-28 15:56:06 +02:00
import (
2019-06-14 09:49:56 -04:00
"context"
2019-08-20 20:06:35 +02:00
"database/sql"
2019-08-29 08:10:45 +02:00
"encoding/hex"
2019-05-30 10:44:49 -04:00
"fmt"
2024-06-18 14:48:40 +08:00
"math/rand"
2025-12-19 20:10:02 +00:00
"sync"
2019-05-28 10:26:20 -04:00
"sync/atomic"
"time"
2019-06-07 14:11:22 -07:00
2019-08-29 08:10:45 +02:00
"go.uber.org/zap"
2024-01-08 15:02:48 +00:00
"github.com/status-im/mvds/peers"
"github.com/status-im/mvds/protobuf"
"github.com/status-im/mvds/state"
"github.com/status-im/mvds/store"
"github.com/status-im/mvds/transport"
2019-04-28 15:56:06 +02:00
)
2019-04-28 05:22:46 +02:00
2019-06-13 13:39:03 -04:00
// Mode represents the synchronization mode.
2019-06-13 10:20:48 -04:00
type Mode int
2019-06-04 17:50:39 -05:00
const (
2019-06-13 10:20:48 -04:00
INTERACTIVE Mode = iota
BATCH
2019-06-04 17:50:39 -05:00
)
2019-07-16 01:07:10 +02:00
// CalculateNextEpoch is a function used to calculate the next `SendEpoch` for a given message.
type CalculateNextEpoch func(count uint64, epoch int64) int64
2019-04-30 01:27:07 +02:00
2024-06-18 14:48:40 +08:00
type EventStatus int
const (
OnlineStatus EventStatus = iota
OfflineStatus
2024-06-18 14:48:40 +08:00
)
const FreshEventPeriod = 10 // seconds
const MaxSendCount = 14 // stop resend the message after 14 times (~10 days)
2024-06-18 14:48:40 +08:00
type PeerStatusChangeEvent struct {
PeerID state.PeerID
Status EventStatus
EventTime uint64
}
2019-06-13 13:35:59 -04:00
// Node represents an MVDS node, it runs all the logic like sending and receiving protocol messages.
2019-04-20 20:09:20 +02:00
type Node struct {
2019-08-29 08:10:45 +02:00
// This needs to be declared first: https://github.com/golang/go/issues/9959
epoch int64
2019-06-14 09:49:56 -04:00
ctx context.Context
cancel context.CancelFunc
2025-12-19 20:10:02 +00:00
wg sync.WaitGroup
2019-06-14 09:49:56 -04:00
2019-06-10 12:13:37 -07:00
store store.MessageStore
transport transport.Transport
2019-04-28 15:56:06 +02:00
2019-06-10 12:13:37 -07:00
syncState state.SyncState
2019-06-10 11:52:49 -07:00
peers peers.Persistence
2019-05-28 10:26:20 -04:00
2019-05-28 20:59:31 -04:00
payloads payloads
2019-04-30 02:35:05 +02:00
2019-07-16 01:07:10 +02:00
nextEpoch CalculateNextEpoch
2019-04-30 01:27:07 +02:00
2019-06-10 12:13:37 -07:00
ID state.PeerID
2019-05-06 19:38:42 +02:00
2025-10-20 13:34:18 +02:00
epochPersistence EpochPersistence
2019-08-20 20:06:35 +02:00
mode Mode
2019-08-06 11:18:46 +02:00
subscription chan protobuf.Message
2019-08-29 08:10:45 +02:00
2024-06-18 14:48:40 +08:00
peerStatusChangeEvent chan PeerStatusChangeEvent
2019-08-29 08:10:45 +02:00
logger *zap.Logger
2019-04-21 01:05:57 +02:00
}
2019-04-27 15:12:16 +02:00
2025-10-20 13:34:18 +02:00
type Persistence interface {
MessageStore() store.MessageStore
PeersStore() peers.Persistence
StateStore() state.SyncState
EpochStore() EpochPersistence
}
type sqlitePersistence struct {
db *sql.DB
}
var _ Persistence = (*sqlitePersistence)(nil)
func NewSQLitePersistence(db *sql.DB) Persistence {
return &sqlitePersistence{db: db}
}
func (p *sqlitePersistence) MessageStore() store.MessageStore {
return store.NewPersistentMessageStore(p.db)
}
func (p *sqlitePersistence) PeersStore() peers.Persistence {
return peers.NewSQLitePersistence(p.db)
}
func (p *sqlitePersistence) StateStore() state.SyncState {
return state.NewPersistentSyncState(p.db)
}
func (p *sqlitePersistence) EpochStore() EpochPersistence {
return NewEpochSQLitePersistence(p.db)
}
2019-08-20 20:06:35 +02:00
func NewPersistentNode(
2025-10-20 13:34:18 +02:00
persistence Persistence,
2019-08-20 20:06:35 +02:00
st transport.Transport,
id state.PeerID,
mode Mode,
nextEpoch CalculateNextEpoch,
2024-06-18 14:48:40 +08:00
peerStatusChangeEvent chan PeerStatusChangeEvent,
2019-08-29 08:10:45 +02:00
logger *zap.Logger,
2019-08-20 20:06:35 +02:00
) (*Node, error) {
ctx, cancel := context.WithCancel(context.Background())
2019-08-29 08:10:45 +02:00
if logger == nil {
logger = zap.NewNop()
}
2019-08-20 20:06:35 +02:00
node := Node{
2024-06-18 14:48:40 +08:00
ID: id,
ctx: ctx,
cancel: cancel,
2025-10-20 13:34:18 +02:00
store: persistence.MessageStore(),
2024-06-18 14:48:40 +08:00
transport: st,
2025-10-20 13:34:18 +02:00
peers: persistence.PeersStore(),
syncState: persistence.StateStore(),
2024-06-18 14:48:40 +08:00
payloads: newPayloads(),
2025-10-20 13:34:18 +02:00
epochPersistence: persistence.EpochStore(),
2024-06-18 14:48:40 +08:00
nextEpoch: nextEpoch,
peerStatusChangeEvent: peerStatusChangeEvent,
logger: logger.Named("mvds"),
2024-06-18 14:48:40 +08:00
mode: mode,
2019-08-20 20:06:35 +02:00
}
if currentEpoch, err := node.epochPersistence.Get(id); err != nil {
return nil, err
} else {
node.epoch = currentEpoch
}
return &node, nil
}
func NewEphemeralNode(
id state.PeerID,
t transport.Transport,
nextEpoch CalculateNextEpoch,
currentEpoch int64,
mode Mode,
2019-08-29 08:10:45 +02:00
logger *zap.Logger,
2019-08-20 20:06:35 +02:00
) *Node {
ctx, cancel := context.WithCancel(context.Background())
2019-08-29 08:10:45 +02:00
if logger == nil {
logger = zap.NewNop()
}
2019-08-20 20:06:35 +02:00
return &Node{
ID: id,
ctx: ctx,
cancel: cancel,
store: store.NewDummyStore(),
transport: t,
syncState: state.NewSyncState(),
peers: peers.NewMemoryPersistence(),
payloads: newPayloads(),
nextEpoch: nextEpoch,
epoch: currentEpoch,
2025-12-19 20:10:02 +00:00
logger: logger.Named("mvds"),
2019-08-20 20:06:35 +02:00
mode: mode,
}
}
2019-06-13 13:35:59 -04:00
// NewNode returns a new node.
2019-06-11 20:50:14 -07:00
func NewNode(
ms store.MessageStore,
st transport.Transport,
ss state.SyncState,
2019-07-16 01:07:10 +02:00
nextEpoch CalculateNextEpoch,
2019-06-11 20:50:14 -07:00
currentEpoch int64,
id state.PeerID,
mode Mode,
pp peers.Persistence,
2019-08-29 08:10:45 +02:00
logger *zap.Logger,
2019-06-11 20:50:14 -07:00
) *Node {
2019-06-14 09:49:56 -04:00
ctx, cancel := context.WithCancel(context.Background())
2019-08-29 08:10:45 +02:00
if logger == nil {
logger = zap.NewNop()
}
2019-05-14 02:16:55 +02:00
return &Node{
ctx: ctx,
cancel: cancel,
store: ms,
transport: st,
syncState: ss,
peers: pp,
payloads: newPayloads(),
nextEpoch: nextEpoch,
ID: id,
epoch: currentEpoch,
2025-12-19 20:10:02 +00:00
logger: logger.Named("mvds"),
mode: mode,
2019-04-30 18:37:39 +02:00
}
}
2019-08-20 20:06:35 +02:00
func (n *Node) CurrentEpoch() int64 {
return atomic.LoadInt64(&n.epoch)
}
2019-06-14 09:49:56 -04:00
// Start listens for new messages received by the node and sends out those required every epoch.
func (n *Node) Start(duration time.Duration) {
2025-12-19 20:10:02 +00:00
n.wg.Add(1)
2019-05-28 10:26:20 -04:00
go func() {
2025-12-19 20:10:02 +00:00
defer n.wg.Done()
2019-05-28 10:26:20 -04:00
for {
2025-12-19 20:10:02 +00:00
p, ok := n.transport.Watch(n.ctx)
if !ok {
n.logger.Debug("watching transport stopped")
2019-06-14 09:49:56 -04:00
return
}
2025-12-19 20:10:02 +00:00
n.wg.Add(1)
go func() {
defer n.wg.Done()
n.onPayload(p.Sender, p.Payload)
}()
2019-05-28 10:26:20 -04:00
}
}()
2019-05-06 19:38:42 +02:00
2025-12-19 20:10:02 +00:00
n.wg.Add(1)
2024-06-18 14:48:40 +08:00
go func() {
2025-12-19 20:10:02 +00:00
defer n.wg.Done()
2024-06-18 14:48:40 +08:00
for {
select {
case <-n.ctx.Done():
2025-12-19 20:10:02 +00:00
n.logger.Debug("reset data sync for peer stopped")
2024-06-18 14:48:40 +08:00
return
case event := <-n.peerStatusChangeEvent:
if event.Status == OnlineStatus && event.EventTime > uint64(time.Now().Unix())-FreshEventPeriod {
2024-06-18 14:48:40 +08:00
n.logger.Debug("resetting peer epoch", zap.String("peerID", hex.EncodeToString(event.PeerID[:4])))
n.resetPeerEpoch(event.PeerID)
}
}
}
}()
2025-12-19 20:10:02 +00:00
n.wg.Add(1)
2019-05-28 10:26:20 -04:00
go func() {
2025-12-19 20:10:02 +00:00
defer n.wg.Done()
2019-05-28 10:26:20 -04:00
for {
2019-06-14 09:49:56 -04:00
select {
case <-n.ctx.Done():
2025-12-19 20:10:02 +00:00
n.logger.Debug("epoch processing stopped")
2019-06-14 09:49:56 -04:00
return
2025-12-19 20:10:02 +00:00
case <-time.After(duration):
err := n.sendMessages()
if err != nil {
2025-12-19 20:10:02 +00:00
n.logger.Error("error sending messages.", zap.Error(err))
}
2024-07-25 13:33:57 +08:00
err = n.syncState.Clear(MaxSendCount)
if err != nil {
2025-12-19 20:10:02 +00:00
n.logger.Error("error clearing sync state.", zap.Error(err))
2024-07-25 13:33:57 +08:00
}
2019-06-14 09:49:56 -04:00
atomic.AddInt64(&n.epoch, 1)
2019-08-20 20:06:35 +02:00
// When a persistent node is used, the epoch needs to be saved.
if n.epochPersistence != nil {
if err := n.epochPersistence.Set(n.ID, n.epoch); err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("Failed to persisten epoch", zap.Error(err))
2019-08-20 20:06:35 +02:00
}
}
2019-06-14 09:49:56 -04:00
}
2019-05-28 10:26:20 -04:00
}
}()
2019-05-06 19:38:42 +02:00
}
2019-06-14 09:49:56 -04:00
// Stop message reading and epoch processing
func (n *Node) Stop() {
2025-12-19 20:10:02 +00:00
n.logger.Info("stopping node")
2019-08-20 20:06:35 +02:00
n.Unsubscribe()
2019-06-14 09:49:56 -04:00
n.cancel()
2025-12-19 20:10:02 +00:00
n.wg.Wait()
2019-06-14 09:49:56 -04:00
}
2019-08-06 11:18:46 +02:00
// Subscribe subscribes to incoming messages.
func (n *Node) Subscribe() chan protobuf.Message {
n.subscription = make(chan protobuf.Message)
return n.subscription
}
// Unsubscribe closes the listening channels
func (n *Node) Unsubscribe() {
2019-08-20 20:06:35 +02:00
if n.subscription != nil {
close(n.subscription)
}
n.subscription = nil
2019-08-06 11:18:46 +02:00
}
2019-05-11 16:01:40 +02:00
// AppendMessage sends a message to a given group.
func (n *Node) AppendMessage(groupID state.GroupID, data []byte) (state.MessageID, error) {
2019-06-07 14:11:22 -07:00
m := protobuf.Message{
GroupId: groupID[:],
2019-05-06 19:38:42 +02:00
Timestamp: time.Now().Unix(),
Body: data,
}
2019-07-12 12:20:30 -04:00
id := m.ID()
2019-05-30 10:44:49 -04:00
peers, err := n.peers.GetByGroupID(groupID)
if err != nil {
return state.MessageID{}, fmt.Errorf("trying to send to unknown group %x", groupID[:4])
2019-05-30 10:44:49 -04:00
}
err = n.store.Add(&m)
2019-05-06 19:38:42 +02:00
if err != nil {
2019-06-10 12:13:37 -07:00
return state.MessageID{}, err
2019-05-06 19:38:42 +02:00
}
for _, p := range peers {
t := state.OFFER
if n.mode == BATCH {
t = state.MESSAGE
2019-05-06 19:38:42 +02:00
}
n.insertSyncState(&groupID, id, p, t)
}
2019-08-29 08:10:45 +02:00
n.logger.Debug("Sending message",
zap.String("node", hex.EncodeToString(n.ID[:4])),
zap.String("groupID", hex.EncodeToString(groupID[:4])),
zap.String("id", hex.EncodeToString(id[:4])))
2019-05-06 19:38:42 +02:00
// @todo think about a way to insta trigger send messages when send was selected, we don't wanna wait for ticks here
2019-05-11 16:01:40 +02:00
return id, nil
2019-04-27 15:12:16 +02:00
}
2019-04-28 05:22:46 +02:00
2019-07-13 23:18:00 -04:00
// RequestMessage adds a REQUEST record to the next payload for a given message ID.
func (n *Node) RequestMessage(group state.GroupID, id state.MessageID) error {
peers, err := n.peers.GetByGroupID(group)
if err != nil {
2019-07-13 23:18:00 -04:00
return fmt.Errorf("trying to request from an unknown group %x", group[:4])
}
for _, p := range peers {
exist, err := n.IsPeerInGroup(group, p)
if err != nil {
return err
2019-07-13 23:18:00 -04:00
}
if exist {
continue
}
n.insertSyncState(&group, id, p, state.REQUEST)
}
2019-07-13 23:18:00 -04:00
return nil
}
2019-06-13 13:35:59 -04:00
// AddPeer adds a peer to a specific group making it a recipient of messages.
func (n *Node) AddPeer(group state.GroupID, id state.PeerID) error {
return n.peers.Add(group, id)
2019-05-06 19:38:42 +02:00
}
2019-04-30 23:53:40 +02:00
2019-06-13 13:35:59 -04:00
// IsPeerInGroup checks whether a peer is in the specified group.
func (n *Node) IsPeerInGroup(g state.GroupID, p state.PeerID) (bool, error) {
return n.peers.Exists(g, p)
2019-05-28 20:58:38 -04:00
}
func (n *Node) sendMessages() error {
err := n.syncState.Map(n.epoch, func(s state.State) state.State {
m := s.MessageID
p := s.PeerID
2019-06-15 16:42:54 -04:00
switch s.Type {
case state.OFFER:
n.payloads.AddOffers(p, m[:])
2019-06-15 16:42:54 -04:00
case state.REQUEST:
n.payloads.AddRequests(p, m[:])
2019-08-29 08:10:45 +02:00
n.logger.Debug("sending REQUEST",
zap.String("from", hex.EncodeToString(n.ID[:4])),
zap.String("to", hex.EncodeToString(p[:4])),
zap.String("messageID", hex.EncodeToString(m[:4])),
)
2019-06-18 15:51:38 -04:00
case state.MESSAGE:
g := *s.GroupID
// TODO: Handle errors
exist, err := n.IsPeerInGroup(g, p)
if err != nil {
return s
}
if !exist {
return s
}
2019-06-18 15:51:38 -04:00
msg, err := n.store.Get(m)
if err != nil {
2025-12-19 20:10:02 +00:00
n.logger.Error("failed to retreive message",
2019-08-29 08:10:45 +02:00
zap.String("messageID", hex.EncodeToString(m[:4])),
zap.Error(err),
)
2019-06-18 15:51:38 -04:00
return s
}
n.payloads.AddMessages(p, msg)
2019-08-29 08:10:45 +02:00
n.logger.Debug("sending MESSAGE",
zap.String("groupID", hex.EncodeToString(g[:4])),
zap.String("from", hex.EncodeToString(n.ID[:4])),
zap.String("to", hex.EncodeToString(p[:4])),
zap.String("messageID", hex.EncodeToString(m[:4])),
)
2019-06-15 16:42:54 -04:00
}
2019-05-28 10:26:20 -04:00
return n.updateSendEpoch(s)
})
2019-06-10 09:45:00 -07:00
if err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("error while mapping sync state", zap.Error(err))
return err
2019-06-10 09:45:00 -07:00
}
2024-01-11 11:03:39 +00:00
return n.payloads.MapAndClear(func(peer state.PeerID, payload *protobuf.Payload) error {
err := n.transport.Send(n.ID, peer, payload)
2019-05-28 10:26:20 -04:00
if err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("error sending message", zap.Error(err))
return err
2019-05-28 10:26:20 -04:00
}
return nil
2019-05-28 10:26:20 -04:00
})
2019-04-30 18:15:05 +02:00
}
2024-01-11 11:03:39 +00:00
func (n *Node) onPayload(sender state.PeerID, payload *protobuf.Payload) {
2019-06-21 15:37:34 +02:00
// Acks, Requests and Offers are all arrays of bytes as protobuf doesn't allow type aliases otherwise arrays of messageIDs would be nicer.
if err := n.onAck(sender, payload.Acks); err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("error processing acks", zap.Error(err))
}
if err := n.onRequest(sender, payload.Requests); err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("error processing requests", zap.Error(err))
}
if err := n.onOffer(sender, payload.Offers); err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("error processing offers", zap.Error(err))
}
messageIds := n.onMessages(sender, payload.Messages)
n.payloads.AddAcks(sender, messageIds)
2019-04-30 00:11:12 +02:00
}
func (n *Node) onOffer(sender state.PeerID, offers [][]byte) error {
2019-06-21 15:37:34 +02:00
for _, raw := range offers {
2019-05-02 04:53:56 +02:00
id := toMessageID(raw)
2019-08-29 08:10:45 +02:00
n.logger.Debug("OFFER received",
zap.String("from", hex.EncodeToString(sender[:4])),
zap.String("to", hex.EncodeToString(n.ID[:4])),
zap.String("messageID", hex.EncodeToString(id[:4])),
)
2019-05-28 10:26:20 -04:00
exist, err := n.store.Has(id)
2019-05-28 10:26:20 -04:00
// @todo maybe ack?
if err != nil {
return err
}
if exist {
2019-05-28 10:26:20 -04:00
continue
}
n.insertSyncState(nil, id, sender, state.REQUEST)
2019-04-30 13:22:04 +02:00
}
return nil
2019-04-30 00:11:12 +02:00
}
func (n *Node) onRequest(sender state.PeerID, requests [][]byte) error {
2019-06-21 15:37:34 +02:00
for _, raw := range requests {
2019-05-26 22:24:50 -04:00
id := toMessageID(raw)
2019-08-29 08:10:45 +02:00
n.logger.Debug("REQUEST received",
zap.String("from", hex.EncodeToString(sender[:4])),
zap.String("to", hex.EncodeToString(n.ID[:4])),
zap.String("messageID", hex.EncodeToString(id[:4])),
)
2019-05-28 10:26:20 -04:00
message, err := n.store.Get(id)
if err != nil {
return err
2019-05-28 10:26:20 -04:00
}
if message == nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("message does not exist", zap.String("messageID", hex.EncodeToString(id[:4])))
2019-05-28 10:26:20 -04:00
continue
}
groupID := toGroupID(message.GroupId)
2019-04-28 05:22:46 +02:00
exist, err := n.IsPeerInGroup(groupID, sender)
2019-06-10 09:46:44 -07:00
if err != nil {
return err
}
if !exist {
2019-08-29 08:10:45 +02:00
n.logger.Error("peer is not in group",
zap.String("groupID", hex.EncodeToString(groupID[:4])),
zap.String("peer", hex.EncodeToString(sender[:4])),
)
2019-06-10 09:46:44 -07:00
continue
}
2019-05-26 22:24:50 -04:00
n.insertSyncState(&groupID, id, sender, state.MESSAGE)
2019-04-28 14:00:23 +02:00
}
return nil
2019-04-28 05:22:46 +02:00
}
func (n *Node) onAck(sender state.PeerID, acks [][]byte) error {
for _, ack := range acks {
id := toMessageID(ack)
err := n.syncState.Remove(id, sender)
if err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("Error while removing sync state.", zap.Error(err))
return err
}
2019-08-29 08:10:45 +02:00
n.logger.Debug("ACK received",
zap.String("from", hex.EncodeToString(sender[:4])),
zap.String("to", hex.EncodeToString(n.ID[:4])),
zap.String("messageID", hex.EncodeToString(id[:4])),
)
}
return nil
}
func (n *Node) onMessages(sender state.PeerID, messages []*protobuf.Message) [][]byte {
2019-05-28 10:26:20 -04:00
a := make([][]byte, 0)
2019-05-26 22:24:50 -04:00
2019-05-28 10:26:20 -04:00
for _, m := range messages {
groupID := toGroupID(m.GroupId)
err := n.onMessage(sender, *m)
2019-05-28 10:26:20 -04:00
if err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("Error processing message", zap.Error(err))
2019-05-28 10:26:20 -04:00
continue
}
2019-07-12 12:20:30 -04:00
id := m.ID()
2019-08-29 08:10:45 +02:00
n.logger.Debug("sending ACK",
zap.String("groupID", hex.EncodeToString(groupID[:4])),
zap.String("from", hex.EncodeToString(n.ID[:4])),
zap.String("", hex.EncodeToString(sender[:4])),
zap.String("messageID", hex.EncodeToString(id[:4])),
)
2019-05-28 10:26:20 -04:00
a = append(a, id[:])
}
return a
}
func (n *Node) onMessage(sender state.PeerID, msg protobuf.Message) error {
2019-07-12 12:20:30 -04:00
id := msg.ID()
groupID := toGroupID(msg.GroupId)
2019-08-29 08:10:45 +02:00
n.logger.Debug("MESSAGE received",
zap.String("from", hex.EncodeToString(sender[:4])),
zap.String("to", hex.EncodeToString(n.ID[:4])),
zap.String("messageID", hex.EncodeToString(id[:4])),
)
2019-05-28 10:26:20 -04:00
err := n.syncState.Remove(id, sender)
if err != nil && err != state.ErrStateNotFound {
2019-06-15 16:42:54 -04:00
return err
}
err = n.store.Add(&msg)
2019-04-30 02:56:58 +02:00
if err != nil {
2019-05-28 10:26:20 -04:00
return err
2019-04-30 02:56:58 +02:00
// @todo process, should this function ever even have an error?
}
2019-04-30 23:53:40 +02:00
peers, err := n.peers.GetByGroupID(groupID)
if err != nil {
return err
}
2019-08-06 11:18:46 +02:00
for _, peer := range peers {
if peer == sender {
continue
}
n.insertSyncState(&groupID, id, peer, state.OFFER)
}
2019-08-06 11:18:46 +02:00
if n.subscription != nil {
n.subscription <- msg
}
2019-05-28 10:26:20 -04:00
return nil
2019-04-28 05:22:46 +02:00
}
2019-04-28 15:58:58 +02:00
func (n *Node) insertSyncState(groupID *state.GroupID, messageID state.MessageID, peerID state.PeerID, t state.RecordType) {
2019-06-18 15:51:38 -04:00
s := state.State{
GroupID: groupID,
MessageID: messageID,
PeerID: peerID,
2019-07-15 05:32:51 +02:00
Type: t,
2019-06-18 15:51:38 -04:00
SendEpoch: n.epoch + 1,
}
err := n.syncState.Add(s)
2019-06-18 15:51:38 -04:00
if err != nil {
2019-08-29 08:10:45 +02:00
n.logger.Error("error setting sync states",
zap.Error(err),
zap.String("groupID", hex.EncodeToString(groupID[:4])),
zap.String("messageID", hex.EncodeToString(messageID[:4])),
zap.String("peerID", hex.EncodeToString(peerID[:4])),
)
2019-06-18 15:51:38 -04:00
}
}
func (n *Node) updateSendEpoch(s state.State) state.State {
2019-05-03 17:47:21 +02:00
s.SendCount += 1
2020-11-23 22:42:35 +01:00
s.SendEpoch = n.nextEpoch(s.SendCount, n.epoch)
2019-05-28 10:26:20 -04:00
return s
2019-04-30 17:20:23 +02:00
}
2024-06-18 14:48:40 +08:00
func (n *Node) resetPeerEpoch(peerID state.PeerID) {
n.syncState.MapWithPeerId(peerID, func(s state.State) state.State {
s.SendEpoch = n.epoch + int64(rand.Intn(60))
return s
})
}
2019-06-10 12:13:37 -07:00
func toMessageID(b []byte) state.MessageID {
var id state.MessageID
2019-05-02 04:53:56 +02:00
copy(id[:], b)
return id
}
func toGroupID(b []byte) state.GroupID {
var id state.GroupID
copy(id[:], b)
return id
}