* Remove specialized handler and archive struct and restructure handlers pkg. * Refactor RPM archive handlers to use a library instead of shelling out * make rpm handling context aware * update test * Refactor AR/deb archive handler to use an existing library instead of shelling out * Update tests * Handle non-archive data within the DefaultHandler * make structs and methods private * Remove non-archive data handling within sources * add max size check * add filename and size to context kvp * move skip file check and is binary check before opening file * fix test * preserve existing funcitonality of not handling non-archive files in HandleFile * Handle non-archive data within the DefaultHandler * rebase * Remove non-archive data handling within sources * Adjust check for rpm/deb archive type * add additional deb mime type * add gzip * move diskbuffered rereader setup into handler pkg * remove DiskBuffereReader creation logic within sources * update comment * move rewind closer * reduce log verbosity * add metrics for file handling * add metrics for errors * make defaultBufferSize a const * add metrics for file handling * add metrics for errors * fix tests * add metrics for max archive depth and skipped files * update error * skip symlinks and dirs * update err * Address incompatible reader to openArchive * remove nil check * fix err assignment * Allow git cat-file blob to complete before trying to handle the file * wrap compReader with DiskbufferReader * Allow git cat-file blob to complete before trying to handle the file * updates * use buffer writer * update * refactor * update context pkg * revert stuff * update test * fix test * remove * use correct reader * add metrics for file handling * add metrics for errors * fix tests * rebase * add metrics for errors * add metrics for max archive depth and skipped files * update error * skip symlinks and dirs * update err * fix err assignment * rebase * remove * Update write method in contentWriter interface * Add bufferReadSeekCloser * update name * update comment * fix lint * Remove specialized handler and archive struct and restructure handlers pkg. * Refactor RPM archive handlers to use a library instead of shelling out * make rpm handling context aware * update test * Refactor AR/deb archive handler to use an existing library instead of shelling out * Update tests * add max size check * add filename and size to context kvp * move skip file check and is binary check before opening file * fix test * preserve existing funcitonality of not handling non-archive files in HandleFile * Handle non-archive data within the DefaultHandler * rebase * Remove non-archive data handling within sources * Handle non-archive data within the DefaultHandler * add gzip * move diskbuffered rereader setup into handler pkg * remove DiskBuffereReader creation logic within sources * update comment * move rewind closer * reduce log verbosity * make defaultBufferSize a const * add metrics for file handling * add metrics for errors * fix tests * add metrics for max archive depth and skipped files * update error * skip symlinks and dirs * update err * Address incompatible reader to openArchive * remove nil check * fix err assignment * wrap compReader with DiskbufferReader * Allow git cat-file blob to complete before trying to handle the file * updates * use buffer writer * update * refactor * update context pkg * revert stuff * update test * remove * rebase * go mod tidy * lint check * update metric to ms * update metric * update comments * dont use ptr * update * fix * Remove specialized handler and archive struct and restructure handlers pkg. * Refactor RPM archive handlers to use a library instead of shelling out * make rpm handling context aware * update test * Refactor AR/deb archive handler to use an existing library instead of shelling out * Update tests * add max size check * add filename and size to context kvp * move skip file check and is binary check before opening file * fix test * preserve existing funcitonality of not handling non-archive files in HandleFile * Adjust check for rpm/deb archive type * add additional deb mime type * update comment * go mod tidy * update go mod * Add a buffered file reader * update comments * use Buffered File Readder * return buffer * update * fix * return * go mod tidy * merge * use a shared pool * use sync.Once * reorganzie * remove unused code * fix double init * fix stuff * nil check * reduce allocations * updates * update metrics * updates * reset buffer instead of putting it back * skip binaries * skip * concurrently process diffs * close chan * concurrently enumerate orgs * increase workers * ignore pbix and vsdx files * add metrics for gitparse's Diffchan * fix metric * update metrics * update * fix checks * fix * inc * update * reduce * Create workers to handle binary files * modify workers * updates * add check * delete code * use custom reader * rename struct * add nonarchive handler * fix break * add comments * add tests * refactor * remove log * do not scan rpm links * simplify * rename var * rename * fix benchmark * add buffer * buffer * buffer * handle panic * merge main * merge main * add recover * revert stuff * revert * revert to using reader * fixes * remove * update * fixes * linter * fix test * move buffers pkg out of writers pkg * rename * [refactor] - move buffer pool logic into own pkg (#2828) * move buffer pool logic into own pkg * fix test * fix test * whoops * [feat] - additional buffer pool (#2829) * move buffer pool logic into own pkg * move * fix test * fix test * fix test * remove * fix test * whoops * revert * fix
108 lines
3.5 KiB
Go
108 lines
3.5 KiB
Go
// Package buffer provides a custom buffer type that includes metrics for tracking buffer usage.
|
|
// It also provides a pool for managing buffer reusability.
|
|
package buffer
|
|
|
|
import (
|
|
"bytes"
|
|
"io"
|
|
"time"
|
|
)
|
|
|
|
// Buffer is a wrapper around bytes.Buffer that includes a timestamp for tracking Buffer checkout duration.
|
|
type Buffer struct {
|
|
*bytes.Buffer
|
|
checkedOutAt time.Time
|
|
}
|
|
|
|
const defaultBufferSize = 1 << 12 // 4KB
|
|
// NewBuffer creates a new instance of Buffer.
|
|
func NewBuffer() *Buffer { return &Buffer{Buffer: bytes.NewBuffer(make([]byte, 0, defaultBufferSize))} }
|
|
|
|
func (b *Buffer) Grow(size int) {
|
|
b.Buffer.Grow(size)
|
|
b.recordGrowth(size)
|
|
}
|
|
|
|
func (b *Buffer) ResetMetric() { b.checkedOutAt = time.Now() }
|
|
|
|
func (b *Buffer) RecordMetric() {
|
|
dur := time.Since(b.checkedOutAt)
|
|
checkoutDuration.Observe(float64(dur.Microseconds()))
|
|
checkoutDurationTotal.Add(float64(dur.Microseconds()))
|
|
totalBufferSize.Add(float64(b.Cap()))
|
|
totalBufferLength.Add(float64(b.Len()))
|
|
}
|
|
|
|
func (b *Buffer) recordGrowth(size int) {
|
|
growCount.Inc()
|
|
growAmount.Add(float64(size))
|
|
}
|
|
|
|
// Write date to the buffer.
|
|
func (b *Buffer) Write(data []byte) (int, error) {
|
|
if b.Buffer == nil {
|
|
// This case should ideally never occur if buffers are properly managed.
|
|
b.Buffer = bytes.NewBuffer(make([]byte, 0, defaultBufferSize))
|
|
b.ResetMetric()
|
|
}
|
|
|
|
size := len(data)
|
|
bufferLength := b.Buffer.Len()
|
|
totalSizeNeeded := bufferLength + size
|
|
|
|
// If the total size is within the threshold, write to the buffer.
|
|
availableSpace := b.Buffer.Cap() - bufferLength
|
|
growSize := totalSizeNeeded - bufferLength
|
|
if growSize > availableSpace {
|
|
// We are manually growing the buffer so we can track the growth via metrics.
|
|
// Knowing the exact data size, we directly resize to fit it, rather than exponential growth
|
|
// which may require multiple allocations and copies if the size required is much larger
|
|
// than double the capacity. Our approach aligns with default behavior when growth sizes
|
|
// happen to match current capacity, retaining asymptotic efficiency benefits.
|
|
b.Grow(growSize)
|
|
}
|
|
|
|
return b.Buffer.Write(data)
|
|
}
|
|
|
|
// Compile time check to make sure readCloser implements io.ReadSeekCloser.
|
|
var _ io.ReadSeekCloser = (*readCloser)(nil)
|
|
|
|
// readCloser is a custom implementation of io.ReadCloser. It wraps a bytes.Reader
|
|
// for reading data from an in-memory buffer and includes an onClose callback.
|
|
// The onClose callback is used to return the buffer to the pool, ensuring buffer re-usability.
|
|
type readCloser struct {
|
|
*bytes.Reader
|
|
onClose func()
|
|
}
|
|
|
|
// ReadCloser creates a new instance of readCloser.
|
|
func ReadCloser(data []byte, onClose func()) *readCloser {
|
|
return &readCloser{Reader: bytes.NewReader(data), onClose: onClose}
|
|
}
|
|
|
|
// Close implements the io.Closer interface. It calls the onClose callback to return the buffer
|
|
// to the pool, enabling buffer reuse. This method should be called by the consumers of ReadCloser
|
|
// once they have finished reading the data to ensure proper resource management.
|
|
func (brc *readCloser) Close() error {
|
|
if brc.onClose == nil {
|
|
return nil
|
|
}
|
|
|
|
brc.onClose() // Return the buffer to the pool
|
|
brc.Reader = nil
|
|
return nil
|
|
}
|
|
|
|
// Read reads up to len(p) bytes into p from the underlying reader.
|
|
// It returns the number of bytes read and any error encountered.
|
|
// On reaching the end of the available data, it returns 0 and io.EOF.
|
|
// Calling Read on a closed reader will also return 0 and io.EOF.
|
|
func (brc *readCloser) Read(p []byte) (int, error) {
|
|
if brc.Reader == nil {
|
|
return 0, io.EOF
|
|
}
|
|
|
|
return brc.Reader.Read(p)
|
|
}
|