-
Notifications
You must be signed in to change notification settings - Fork 110
New issue
Have a question about this project? Sign up for a free GitHub account to open an issue and contact its maintainers and the community.
By clicking “Sign up for GitHub”, you agree to our terms of service and privacy statement. We’ll occasionally send you account related emails.
Already on GitHub? Sign in to your account
chore(dot/network): use sync.Pool for network message buffers #1600
Changes from 12 commits
5a38145
117ebb8
2fca3d8
09e8ca9
328dcd6
040cdc4
3f5a7b6
4eceebf
0680d7c
1ca19c6
4d747f2
e648cf1
e7abd0f
19034b8
0c71acd
823b5f8
66ebd2b
5c33f4a
2955598
5634551
ecfa3c7
e666056
a4199e9
File filter
Filter by extension
Conversations
Jump to
Diff view
Diff view
There are no files selected for viewing
Original file line number | Diff line number | Diff line change |
---|---|---|
|
@@ -46,6 +46,8 @@ const ( | |
lightID = "/light/2" | ||
blockAnnounceID = "/block-announces/1" | ||
transactionsID = "/transactions/1" | ||
|
||
maxMessageSize = maxBlockResponseSize | ||
) | ||
|
||
var ( | ||
|
@@ -72,6 +74,9 @@ type Service struct { | |
mdns *mdns | ||
gossip *gossip | ||
syncQueue *syncQueue | ||
bufPool *sync.Pool | ||
bufMap *sync.Map //map[*[maxMessageSize]byte]struct{} | ||
bufCount int | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. Use atomic counter instead. This might lead to race condition. There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. You can also use channels instead of There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. removed in favour of |
||
|
||
notificationsProtocols map[byte]*notificationsProtocol // map of sub-protocol msg ID to protocol info | ||
notificationsMu sync.RWMutex | ||
|
@@ -132,6 +137,30 @@ func NewService(cfg *Config) (*Service, error) { | |
return nil, err | ||
} | ||
|
||
// pre-allocate pool of buffers used to read from streams. | ||
// initially allocate as many buffers as minimally necessary which is the number inbound streams we will have, | ||
// which should equal min peers times the number of notifications protocols, which is currently 3. | ||
// | ||
// the purpose of the bufMap is to keep the reference to the buffers inside the *Service | ||
// at all times to prevent the buffers from being garbage collected. | ||
// with this addition, they are only garbage collected once the *Service is no longer | ||
// used, which is when the node shuts down. | ||
bufPool := new(sync.Pool) | ||
bufPool.New = func() interface{} { | ||
return [maxMessageSize]byte{} | ||
} | ||
bufMap := new(sync.Map) | ||
|
||
var bufCount int | ||
if !cfg.noPreAllocate { | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. If you switch
|
||
for i := 0; i < cfg.MinPeers*3; i++ { | ||
buf := [maxBlockResponseSize]byte{} | ||
bufPool.Put(buf) //nolint | ||
bufMap.Store(&buf, struct{}{}) | ||
} | ||
bufCount = cfg.MinPeers * 3 | ||
} | ||
|
||
network := &Service{ | ||
ctx: ctx, | ||
cancel: cancel, | ||
|
@@ -148,6 +177,9 @@ func NewService(cfg *Config) (*Service, error) { | |
lightRequest: make(map[peer.ID]struct{}), | ||
telemetryInterval: cfg.telemetryInterval, | ||
closeCh: make(chan interface{}), | ||
bufPool: bufPool, | ||
bufMap: bufMap, | ||
bufCount: bufCount, | ||
} | ||
|
||
network.syncQueue = newSyncQueue(network) | ||
|
@@ -550,15 +582,39 @@ func isInbound(stream libp2pnetwork.Stream) bool { | |
return stream.Stat().Direction == libp2pnetwork.DirInbound | ||
} | ||
|
||
func (s *Service) getMessageBuffer() [maxMessageSize]byte { | ||
buf := s.bufPool.Get() | ||
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. You can directly do this.
There was a problem hiding this comment. Choose a reason for hiding this commentThe reason will be displayed to describe this comment to others. Learn more. removed in favour of |
||
if buf == nil { | ||
return [maxMessageSize]byte{} | ||
} | ||
|
||
msgBytes, ok := buf.([maxMessageSize]byte) | ||
if !ok { | ||
return [maxMessageSize]byte{} | ||
} | ||
|
||
return msgBytes | ||
} | ||
|
||
func (s *Service) cleanupMessageBuffer(buf [maxMessageSize]byte) { | ||
s.bufPool.Put(buf) //nolint | ||
|
||
logger.Trace("message buffers allocated", "count", s.bufCount) | ||
if s.bufCount >= s.cfg.MaxPeers*3 { | ||
return | ||
} | ||
|
||
s.bufMap.Store(&buf, struct{}{}) | ||
s.bufCount++ | ||
} | ||
|
||
func (s *Service) readStream(stream libp2pnetwork.Stream, decoder messageDecoder, handler messageHandler) { | ||
var ( | ||
maxMessageSize uint64 = maxBlockResponseSize // TODO: determine actual max message size | ||
msgBytes = make([]byte, maxMessageSize) | ||
peer = stream.Conn().RemotePeer() | ||
) | ||
peer := stream.Conn().RemotePeer() | ||
msgBytes := s.getMessageBuffer() | ||
defer s.cleanupMessageBuffer(msgBytes) | ||
|
||
for { | ||
tot, err := readStream(stream, msgBytes) | ||
tot, err := readStream(stream, msgBytes[:]) | ||
if err == io.EOF { | ||
continue | ||
} else if err != nil { | ||
|
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
Rename it to
preAllocate
Because this is difficult to read
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
we want to pre-allocate by default so if this is changed to
preAllocate
then it will be false by default, but we won't know if it was set to false or if it defaulted to false, hence why it'snoPreAllocate