mirror of
https://github.com/libp2p/go-libp2p-resource-manager.git
synced 2026-08-22 11:33:28 +08:00
update for interface changes; RIP buffers.
This commit is contained in:
8
rcmgr.go
8
rcmgr.go
@@ -80,7 +80,7 @@ type ConnectionScope struct {
|
||||
}
|
||||
|
||||
var _ network.ConnectionScope = (*ConnectionScope)(nil)
|
||||
var _ network.UserConnectionScope = (*ConnectionScope)(nil)
|
||||
var _ network.ConnectionManagementScope = (*ConnectionScope)(nil)
|
||||
|
||||
type StreamScope struct {
|
||||
*ResourceScope
|
||||
@@ -93,7 +93,7 @@ type StreamScope struct {
|
||||
}
|
||||
|
||||
var _ network.StreamScope = (*StreamScope)(nil)
|
||||
var _ network.UserStreamScope = (*StreamScope)(nil)
|
||||
var _ network.StreamManagementScope = (*StreamScope)(nil)
|
||||
|
||||
func NewResourceManager(limits Limiter) *ResourceManager {
|
||||
r := &ResourceManager{
|
||||
@@ -187,7 +187,7 @@ func (r *ResourceManager) getPeerScope(p peer.ID) *PeerScope {
|
||||
return s
|
||||
}
|
||||
|
||||
func (r *ResourceManager) OpenConnection(dir network.Direction, usefd bool) (network.ConnectionScope, error) {
|
||||
func (r *ResourceManager) OpenConnection(dir network.Direction, usefd bool) (network.ConnectionManagementScope, error) {
|
||||
conn := NewConnectionScope(dir, usefd, r.limits.GetConnLimits(), r)
|
||||
|
||||
if err := conn.AddConn(dir); err != nil {
|
||||
@@ -205,7 +205,7 @@ func (r *ResourceManager) OpenConnection(dir network.Direction, usefd bool) (net
|
||||
return conn, nil
|
||||
}
|
||||
|
||||
func (r *ResourceManager) OpenStream(p peer.ID, dir network.Direction) (network.StreamScope, error) {
|
||||
func (r *ResourceManager) OpenStream(p peer.ID, dir network.Direction) (network.StreamManagementScope, error) {
|
||||
peer := r.getPeerScope(p)
|
||||
stream := NewStreamScope(dir, r.limits.GetStreamLimits(p), peer)
|
||||
peer.DecRef() // we have the reference in constraints
|
||||
|
||||
152
scope.go
152
scope.go
@@ -4,7 +4,6 @@ import (
|
||||
"fmt"
|
||||
"sync"
|
||||
|
||||
"github.com/libp2p/go-buffer-pool"
|
||||
"github.com/libp2p/go-libp2p-core/network"
|
||||
)
|
||||
|
||||
@@ -16,9 +15,7 @@ type resources struct {
|
||||
nstreamsIn, nstreamsOut int
|
||||
nfd int
|
||||
|
||||
memory int64
|
||||
buffers map[int][]byte
|
||||
nextBuf int
|
||||
memory int64
|
||||
}
|
||||
|
||||
// ResourceScopes.
|
||||
@@ -27,7 +24,7 @@ type resources struct {
|
||||
// using a linearized parent set.
|
||||
// A ResourceScope can be a txn scope, where it has a specific owner; txn scopes create a tree rooted
|
||||
// at the owner (which can be a DAG scope) and can outlive their parents -- this is important because
|
||||
// txn scopes are the main *user* interface for buffer/memory management, and the user may call
|
||||
// txn scopes are the main *user* interface for memory management, and the user may call
|
||||
// Done in a txn scope after the system has closed the root of the txn tree in some background
|
||||
// goroutine.
|
||||
// If we didn't make this distinction we would have a double release problem in that case.
|
||||
@@ -44,14 +41,6 @@ type ResourceScope struct {
|
||||
var _ network.ResourceScope = (*ResourceScope)(nil)
|
||||
var _ network.TransactionalScope = (*ResourceScope)(nil)
|
||||
|
||||
type Buffer struct {
|
||||
s *ResourceScope
|
||||
data []byte
|
||||
key int
|
||||
}
|
||||
|
||||
var _ network.Buffer = (*Buffer)(nil)
|
||||
|
||||
func newResources(limit Limit) *resources {
|
||||
return &resources{
|
||||
limit: limit,
|
||||
@@ -91,13 +80,6 @@ func (rc *resources) checkMemory(rsvp int64) error {
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rc *resources) releaseBuffers() {
|
||||
for _, buf := range rc.buffers {
|
||||
pool.Put(buf)
|
||||
}
|
||||
rc.buffers = nil
|
||||
}
|
||||
|
||||
func (rc *resources) reserveMemory(size int64) error {
|
||||
if err := rc.checkMemory(size); err != nil {
|
||||
return err
|
||||
@@ -116,61 +98,6 @@ func (rc *resources) releaseMemory(size int64) {
|
||||
}
|
||||
}
|
||||
|
||||
func (rc *resources) getBuffer(size int) ([]byte, int, error) {
|
||||
if err := rc.checkMemory(int64(size)); err != nil {
|
||||
return nil, -1, err
|
||||
}
|
||||
|
||||
buf := pool.Get(size)
|
||||
key := rc.nextBuf
|
||||
|
||||
rc.memory += int64(size)
|
||||
if rc.buffers == nil {
|
||||
rc.buffers = make(map[int][]byte)
|
||||
}
|
||||
rc.buffers[key] = buf
|
||||
rc.nextBuf++
|
||||
|
||||
return buf, key, nil
|
||||
}
|
||||
|
||||
func (rc *resources) growBuffer(key int, newsize int) ([]byte, error) {
|
||||
oldbuf, ok := rc.buffers[key]
|
||||
if !ok {
|
||||
return nil, fmt.Errorf("invalid buffer; cannot grow buffer not allocated through this scope")
|
||||
}
|
||||
|
||||
grow := newsize - len(oldbuf)
|
||||
if err := rc.checkMemory(int64(grow)); err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
newbuf := pool.Get(newsize)
|
||||
copy(newbuf, oldbuf)
|
||||
|
||||
rc.memory += int64(grow)
|
||||
rc.buffers[key] = newbuf
|
||||
|
||||
return newbuf, nil
|
||||
}
|
||||
|
||||
func (rc *resources) releaseBuffer(key int) {
|
||||
buf, ok := rc.buffers[key]
|
||||
if !ok {
|
||||
panic("BUG: release unknown buffer")
|
||||
}
|
||||
|
||||
rc.memory -= int64(len(buf))
|
||||
|
||||
// sanity check for bugs upstream
|
||||
if rc.memory < 0 {
|
||||
panic("BUG: too much memory released")
|
||||
}
|
||||
|
||||
delete(rc.buffers, key)
|
||||
pool.Put(buf)
|
||||
}
|
||||
|
||||
func (rc *resources) addStream(dir network.Direction) error {
|
||||
if dir == network.DirInbound {
|
||||
return rc.addStreams(1, 0)
|
||||
@@ -363,76 +290,6 @@ func (s *ResourceScope) ReleaseMemoryForChild(size int64) {
|
||||
s.rc.releaseMemory(size)
|
||||
}
|
||||
|
||||
func (s *ResourceScope) GetBuffer(size int) (network.Buffer, error) {
|
||||
s.Lock()
|
||||
defer s.Unlock()
|
||||
|
||||
if s.done {
|
||||
return nil, ErrResourceScopeClosed
|
||||
}
|
||||
|
||||
buf, key, err := s.rc.getBuffer(size)
|
||||
if err != nil {
|
||||
return nil, err
|
||||
}
|
||||
|
||||
if err := s.reserveMemoryForConstraints(size); err != nil {
|
||||
s.rc.releaseBuffer(key)
|
||||
return nil, err
|
||||
}
|
||||
|
||||
return &Buffer{s: s, data: buf, key: key}, nil
|
||||
}
|
||||
|
||||
func (b *Buffer) Data() []byte { return b.data }
|
||||
|
||||
func (b *Buffer) Grow(newsize int) error {
|
||||
b.s.Lock()
|
||||
defer b.s.Unlock()
|
||||
|
||||
if b.s.done {
|
||||
return ErrResourceScopeClosed
|
||||
}
|
||||
|
||||
grow := newsize - len(b.data)
|
||||
if err := b.s.reserveMemoryForConstraints(grow); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
newbuf, err := b.s.rc.growBuffer(b.key, newsize)
|
||||
if err != nil {
|
||||
b.s.releaseMemoryForConstraints(grow)
|
||||
return err
|
||||
}
|
||||
|
||||
b.data = newbuf
|
||||
return nil
|
||||
}
|
||||
|
||||
func (b *Buffer) Release() {
|
||||
b.s.Lock()
|
||||
defer b.s.Unlock()
|
||||
|
||||
if b.data == nil {
|
||||
return
|
||||
}
|
||||
|
||||
if b.s.done {
|
||||
b.data = nil
|
||||
return
|
||||
}
|
||||
|
||||
if b.s.owner != nil {
|
||||
b.s.owner.ReleaseMemory(len(b.data))
|
||||
} else {
|
||||
for _, cst := range b.s.constraints {
|
||||
cst.ReleaseMemoryForChild(int64(len(b.data)))
|
||||
}
|
||||
}
|
||||
b.s.rc.releaseBuffer(b.key)
|
||||
b.data = nil
|
||||
}
|
||||
|
||||
func (s *ResourceScope) AddStream(dir network.Direction) error {
|
||||
s.Lock()
|
||||
defer s.Unlock()
|
||||
@@ -765,7 +622,7 @@ func (s *ResourceScope) ReleaseResources(st network.ScopeStat) {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *ResourceScope) BeginTxn() (network.TransactionalScope, error) {
|
||||
func (s *ResourceScope) BeginTransaction() (network.TransactionalScope, error) {
|
||||
s.Lock()
|
||||
defer s.Unlock()
|
||||
|
||||
@@ -796,15 +653,12 @@ func (s *ResourceScope) Done() {
|
||||
}
|
||||
}
|
||||
|
||||
s.rc.releaseBuffers()
|
||||
|
||||
s.rc.nstreamsIn = 0
|
||||
s.rc.nstreamsOut = 0
|
||||
s.rc.nconnsIn = 0
|
||||
s.rc.nconnsOut = 0
|
||||
s.rc.nfd = 0
|
||||
s.rc.memory = 0
|
||||
s.rc.buffers = nil
|
||||
|
||||
s.done = true
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user