mirror of
https://github.com/buger/goreplay.git
synced 2024-04-21 12:32:02 +00:00
Added `—input-raw-buffer-size` - Controls size of the OS buffer (in bytes) which holds packets until they dispatched. Default value depends by system: in Linux around 2MB. If you see big package drop, increase this value. Additionally snaplen (max number of bytes being read for each packet) now dynamically set based on interface MTU + max header size. In most situations it should reduce package drop, because each packet will consume less space in buffer.
101 lines
2.0 KiB
Go
101 lines
2.0 KiB
Go
package main
|
|
|
|
import (
|
|
"fmt"
|
|
"io"
|
|
"math/rand"
|
|
"strconv"
|
|
"strings"
|
|
"time"
|
|
)
|
|
|
|
// Limiter is a wrapper for input or output plugin which adds rate limiting
|
|
type Limiter struct {
|
|
plugin interface{}
|
|
limit int
|
|
isPercent bool
|
|
|
|
currentRPS int
|
|
currentTime int64
|
|
}
|
|
|
|
func parseLimitOptions(options string) (limit int, isPercent bool) {
|
|
if strings.Contains(options, "%") {
|
|
limit, _ = strconv.Atoi(strings.Split(options, "%")[0])
|
|
isPercent = true
|
|
} else {
|
|
limit, _ = strconv.Atoi(options)
|
|
isPercent = false
|
|
}
|
|
|
|
return
|
|
}
|
|
|
|
// NewLimiter constructor for Limiter, accepts plugin and options
|
|
// `options` allow to sprcify relatve or absolute limiting
|
|
func NewLimiter(plugin interface{}, options string) io.ReadWriter {
|
|
l := new(Limiter)
|
|
l.limit, l.isPercent = parseLimitOptions(options)
|
|
l.plugin = plugin
|
|
l.currentTime = time.Now().UnixNano()
|
|
|
|
// FileInput have its own rate limiting. Unlike other inputs we not just dropping requests, we can slow down or speed up request emittion.
|
|
if fi, ok := l.plugin.(*FileInput); ok && l.isPercent {
|
|
fi.speedFactor = float64(l.limit) / float64(100)
|
|
}
|
|
|
|
return l
|
|
}
|
|
|
|
func (l *Limiter) isLimited() bool {
|
|
// File input have its own limiting algorithm
|
|
if _, ok := l.plugin.(*FileInput); ok && l.isPercent {
|
|
return false
|
|
}
|
|
|
|
if l.isPercent {
|
|
return l.limit <= rand.Intn(100)
|
|
}
|
|
|
|
if (time.Now().UnixNano() - l.currentTime) > time.Second.Nanoseconds() {
|
|
l.currentTime = time.Now().UnixNano()
|
|
l.currentRPS = 0
|
|
}
|
|
|
|
if l.currentRPS >= l.limit {
|
|
return true
|
|
}
|
|
|
|
l.currentRPS++
|
|
|
|
return false
|
|
}
|
|
|
|
func (l *Limiter) Write(data []byte) (n int, err error) {
|
|
if l.isLimited() {
|
|
return 0, nil
|
|
}
|
|
|
|
n, err = l.plugin.(io.Writer).Write(data)
|
|
|
|
return
|
|
}
|
|
|
|
func (l *Limiter) Read(data []byte) (n int, err error) {
|
|
if r, ok := l.plugin.(io.Reader); ok {
|
|
n, err = r.Read(data)
|
|
} else {
|
|
return 0, nil
|
|
}
|
|
|
|
if l.isLimited() {
|
|
return 0, nil
|
|
}
|
|
|
|
return
|
|
}
|
|
|
|
func (l *Limiter) String() string {
|
|
return fmt.Sprintf("Limiting %s to: %d (isPercent: %v)", l.plugin, l.limit, l.isPercent)
|
|
}
|