mirror of
https://github.com/libp2p/go-libp2p-resource-manager.git
synced 2026-08-22 11:33:28 +08:00
update for reservation interface changes
This commit is contained in:
6
rcmgr.go
6
rcmgr.go
@@ -339,7 +339,7 @@ func (s *connectionScope) SetPeer(p peer.ID) error {
|
||||
|
||||
// juggle resources from transient scope to peer scope
|
||||
stat := s.resourceScope.rc.stat()
|
||||
if _, err := s.peer.ReserveForChild(stat); err != nil {
|
||||
if err := s.peer.ReserveForChild(stat); err != nil {
|
||||
s.peer.DecRef()
|
||||
s.peer = nil
|
||||
return err
|
||||
@@ -376,7 +376,7 @@ func (s *streamScope) SetProtocol(proto protocol.ID) error {
|
||||
|
||||
// juggle resources from transient scope to protocol scope
|
||||
stat := s.resourceScope.rc.stat()
|
||||
if _, err := s.proto.ReserveForChild(stat); err != nil {
|
||||
if err := s.proto.ReserveForChild(stat); err != nil {
|
||||
s.proto.DecRef()
|
||||
s.proto = nil
|
||||
return err
|
||||
@@ -416,7 +416,7 @@ func (s *streamScope) SetService(svc string) error {
|
||||
s.svc = s.rcmgr.getServiceScope(svc)
|
||||
|
||||
// reserve resources in service
|
||||
if _, err := s.svc.ReserveForChild(s.resourceScope.rc.stat()); err != nil {
|
||||
if err := s.svc.ReserveForChild(s.resourceScope.rc.stat()); err != nil {
|
||||
s.svc.DecRef()
|
||||
s.svc = nil
|
||||
return err
|
||||
|
||||
83
scope.go
83
scope.go
@@ -58,36 +58,26 @@ func newTxnResourceScope(owner *resourceScope) *resourceScope {
|
||||
}
|
||||
|
||||
// Resources implementation
|
||||
func (rc *resources) checkMemory(rsvp int64) (network.MemoryStatus, error) {
|
||||
func (rc *resources) checkMemory(rsvp int64, prio uint8) error {
|
||||
// overflow check; this also has the side effect that we cannot reserve negative memory.
|
||||
newmem := rc.memory + rsvp
|
||||
if newmem < rc.memory {
|
||||
return network.MemoryStatusOK, fmt.Errorf("memory reservation overflow: %w", network.ErrResourceLimitExceeded)
|
||||
}
|
||||
|
||||
// limit check
|
||||
limit := rc.limit.GetMemoryLimit()
|
||||
switch {
|
||||
case newmem > limit:
|
||||
return network.MemoryStatusOK, fmt.Errorf("cannot reserve memory: %w", network.ErrResourceLimitExceeded)
|
||||
case newmem > int64(0.8*float64(limit)):
|
||||
return network.MemoryStatusCritical, nil
|
||||
threshold := (1 + int64(prio)) * limit / 256
|
||||
|
||||
case newmem > int64(0.5*float64(limit)):
|
||||
return network.MemoryStatusCaution, nil
|
||||
|
||||
default:
|
||||
return network.MemoryStatusOK, nil
|
||||
if newmem > threshold {
|
||||
return network.ErrResourceLimitExceeded
|
||||
}
|
||||
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rc *resources) reserveMemory(size int64) (status network.MemoryStatus, err error) {
|
||||
if status, err = rc.checkMemory(size); err != nil {
|
||||
return status, err
|
||||
func (rc *resources) reserveMemory(size int64, prio uint8) error {
|
||||
if err := rc.checkMemory(size, prio); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
rc.memory += size
|
||||
return status, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
func (rc *resources) releaseMemory(size int64) {
|
||||
@@ -204,47 +194,38 @@ func (rc *resources) stat() network.ScopeStat {
|
||||
}
|
||||
|
||||
// resourceScope implementation
|
||||
func (s *resourceScope) ReserveMemory(size int) (status network.MemoryStatus, err error) {
|
||||
func (s *resourceScope) ReserveMemory(size int, prio uint8) error {
|
||||
s.Lock()
|
||||
defer s.Unlock()
|
||||
|
||||
if s.done {
|
||||
return network.MemoryStatusOK, network.ErrResourceScopeClosed
|
||||
return network.ErrResourceScopeClosed
|
||||
}
|
||||
|
||||
if status, err = s.rc.reserveMemory(int64(size)); err != nil {
|
||||
return status, err
|
||||
if err := s.rc.reserveMemory(int64(size), prio); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
var statusCst network.MemoryStatus
|
||||
if statusCst, err = s.reserveMemoryForConstraints(size); err != nil {
|
||||
if err := s.reserveMemoryForConstraints(size, prio); err != nil {
|
||||
s.rc.releaseMemory(int64(size))
|
||||
return status, err
|
||||
return err
|
||||
}
|
||||
|
||||
if statusCst > status {
|
||||
status = statusCst
|
||||
}
|
||||
|
||||
return status, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *resourceScope) reserveMemoryForConstraints(size int) (status network.MemoryStatus, err error) {
|
||||
func (s *resourceScope) reserveMemoryForConstraints(size int, prio uint8) error {
|
||||
if s.owner != nil {
|
||||
return s.owner.ReserveMemory(size)
|
||||
return s.owner.ReserveMemory(size, prio)
|
||||
}
|
||||
|
||||
var reserved int
|
||||
var err error
|
||||
for _, cst := range s.constraints {
|
||||
var statusCst network.MemoryStatus
|
||||
if statusCst, err = cst.ReserveMemoryForChild(int64(size)); err != nil {
|
||||
if err = cst.ReserveMemoryForChild(int64(size), prio); err != nil {
|
||||
break
|
||||
}
|
||||
|
||||
if statusCst > status {
|
||||
status = statusCst
|
||||
}
|
||||
|
||||
reserved++
|
||||
}
|
||||
|
||||
@@ -255,7 +236,7 @@ func (s *resourceScope) reserveMemoryForConstraints(size int) (status network.Me
|
||||
}
|
||||
}
|
||||
|
||||
return status, err
|
||||
return err
|
||||
}
|
||||
|
||||
func (s *resourceScope) releaseMemoryForConstraints(size int) {
|
||||
@@ -269,15 +250,15 @@ func (s *resourceScope) releaseMemoryForConstraints(size int) {
|
||||
}
|
||||
}
|
||||
|
||||
func (s *resourceScope) ReserveMemoryForChild(size int64) (network.MemoryStatus, error) {
|
||||
func (s *resourceScope) ReserveMemoryForChild(size int64, prio uint8) error {
|
||||
s.Lock()
|
||||
defer s.Unlock()
|
||||
|
||||
if s.done {
|
||||
return network.MemoryStatusOK, network.ErrResourceScopeClosed
|
||||
return network.ErrResourceScopeClosed
|
||||
}
|
||||
|
||||
return s.rc.reserveMemory(size)
|
||||
return s.rc.reserveMemory(size, prio)
|
||||
}
|
||||
|
||||
func (s *resourceScope) ReleaseMemory(size int) {
|
||||
@@ -478,30 +459,30 @@ func (s *resourceScope) RemoveConnForChild(dir network.Direction, usefd bool) {
|
||||
s.rc.removeConn(dir, usefd)
|
||||
}
|
||||
|
||||
func (s *resourceScope) ReserveForChild(st network.ScopeStat) (status network.MemoryStatus, err error) {
|
||||
func (s *resourceScope) ReserveForChild(st network.ScopeStat) error {
|
||||
s.Lock()
|
||||
defer s.Unlock()
|
||||
|
||||
if s.done {
|
||||
return network.MemoryStatusOK, network.ErrResourceScopeClosed
|
||||
return network.ErrResourceScopeClosed
|
||||
}
|
||||
|
||||
if status, err = s.rc.reserveMemory(st.Memory); err != nil {
|
||||
return status, err
|
||||
if err := s.rc.reserveMemory(st.Memory, network.ReservationPriorityAlways); err != nil {
|
||||
return err
|
||||
}
|
||||
|
||||
if err := s.rc.addStreams(st.NumStreamsInbound, st.NumStreamsOutbound); err != nil {
|
||||
s.rc.releaseMemory(st.Memory)
|
||||
return status, err
|
||||
return err
|
||||
}
|
||||
|
||||
if err := s.rc.addConns(st.NumConnsInbound, st.NumConnsOutbound, st.NumFD); err != nil {
|
||||
s.rc.releaseMemory(st.Memory)
|
||||
s.rc.removeStreams(st.NumStreamsInbound, st.NumStreamsOutbound)
|
||||
return status, err
|
||||
return err
|
||||
}
|
||||
|
||||
return status, nil
|
||||
return nil
|
||||
}
|
||||
|
||||
func (s *resourceScope) ReleaseForChild(st network.ScopeStat) {
|
||||
|
||||
Reference in New Issue
Block a user