// ------------------------------------------------------------------------ // Project atila // Active Thing (activething.com) git.activething.com/go // // File name buffer_pool.go // Created by DEV // Modified 16/01/2024 // // Copyright 2024 activething.com // ------------------------------------------------------------------------ // Licensed under the Apache License, Version 2.0 (the "License"); // you may not use this file except in compliance with the License. // You may obtain a copy of the License at // // http://www.apache.org/licenses/LICENSE-2.0 // // Unless required by applicable law or agreed to in writing, software // distributed under the License is distributed on an "AS IS" BASIS, // WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. // See the License for the specific language governing permissions and // limitations under the License. // ------------------------------------------------------------------------ package mems import ( "sort" "sync" "sync/atomic" "syncs" ) type ( Pool struct { syncs.PoolStatsInfo inner sync.Pool calls [poolSteps]uint64 calibra uint64 defSize uint64 // buffer default size maxSize uint64 threshold uint64 percentile float32 } PoolFnc func (*Pool) ) func NewPool(options ...PoolFnc) *Pool { return applyBufferPoolOptions(&Pool{percentile: poolPercentile}, options...) } func (p *Pool) MaxSize() uint64 { return atomic.LoadUint64(&p.maxSize) } func (p *Pool) DefSize() uint64 { return atomic.LoadUint64(&p.defSize) } func (p *Pool) MaxPercentile() float32 { return p.percentile } func (p *Pool) CallsThreshold() uint64 { return atomic.LoadUint64(&p.threshold) } func (p *Pool) Get() *ABuffer { b := p.inner.Get() p.AddAcquired() if b == nil { p.AddCreated() return NewABuffer() } return b.(*ABuffer) } func (p *Pool) Put(buffer *ABuffer) { if buffer == nil { return } i := index(len(buffer.data)) if atomic.AddUint64(&p.calls[i], 1) > atomic.LoadUint64(&p.threshold) { p.calibrate() } if m := int(atomic.LoadUint64(&p.maxSize)); m == 0 || buffer.Cap() <= m { buffer.Reset() p.inner.Put(buffer) } p.AddReleased() } func (p *Pool) calibrate() { if atomic.CompareAndSwapUint64(&p.calibra, 0, 1) { return } ar := make(syncs.PoolCallSizes, 0, poolSteps) cs := uint64(0) for i := uint64(0); i < poolSteps; i++ { cl := atomic.SwapUint64(&p.calls[i], 0) cs += cl ar = append(ar, syncs.PoolCallSize{ Calls: cl, Size: minSize << i}) } sort.Sort(ar) ds := ar[0].Size ms := ds mx := uint64(float32(cs) * p.percentile) cs = 0 for i := 0; i < poolSteps; i++ { if cs > mx { break } cs += ar[i].Calls sz := ar[i].Size if sz > ms { ms = sz } } atomic.StoreUint64(&p.defSize, ds) atomic.StoreUint64(&p.maxSize, ms) atomic.StoreUint64(&p.calibra, 0) } func index(n int) int { n-- n >>= minBitSize i := 0 for n > 0 { n >>= 1 i++ } if i >= poolSteps { return poolSteps - 1 } return i } func applyBufferPoolOptions (b *Pool, options ...PoolFnc) *Pool { for _,o := range options { o(b) } return b } func WithMaxPercentile(percentile float32) PoolFnc { return func(b *Pool) { b.percentile = percentile } } func CallsThreshold(threshold uint64) PoolFnc { return func(b *Pool) { atomic.StoreUint64(&b.threshold, threshold) } }