// Package ratelimit begrenzt den Datendurchsatz eines Sicherungslaufs. // // Der Zweck ist betrieblich, nicht technisch: Eine Sicherung, die eine // Standleitung vollständig belegt, legt die Anwendungen lahm, für deren Schutz // sie läuft. Ein Backup, das man abschalten muss, um arbeiten zu können, wird // abgeschaltet. package ratelimit import ( "context" "errors" "fmt" "io" "sync" "time" ) // Unlimited steht für einen abgeschalteten Begrenzer. const Unlimited int64 = 0 // minimumBurstBytes ist die kleinste sinnvolle Eimergröße. // // Kleiner als ein typischer Netzwerkblock zu begrenzen führt zu ständigem // Warten in winzigen Schritten: Der Durchsatz bräche stärker ein als vorgesehen, // weil jeder Aufruf mehr Zeit mit Warten als mit Übertragen verbrächte. const minimumBurstBytes int64 = 32 * 1024 // defaultBurstDuration bestimmt die Eimergröße aus der Rate. // // Eine Sekunde Vorrat ist der übliche Ausgleich: Er glättet kurze Spitzen, // ohne dass sich über längere Zeit mehr als die vereinbarte Rate ansammelt. const defaultBurstDuration = time.Second // Limiter begrenzt den Durchsatz nach dem Token-Bucket-Verfahren. // // Gewählt wurde das Token-Bucket-Verfahren und kein festes Zeitfenster: Ein // Fenster erlaubt an seiner Grenze die doppelte Rate — 10 MB am Ende des einen // und 10 MB am Anfang des nächsten Fensters ergeben 20 MB in kurzer Folge. Der // Eimer verhindert das, weil er kontinuierlich nachfüllt. type Limiter struct { // mutex schützt den Zustand. // // Ein Begrenzer wird von allen Arbeitern der Pipeline gemeinsam genutzt; // je Arbeiter einen eigenen zu führen ergäbe ein Vielfaches der Rate. mutex sync.Mutex // bytesPerSecond ist die zugelassene Rate. bytesPerSecond int64 // burstBytes ist die Größe des Eimers. burstBytes int64 // availableTokens ist der aktuelle Füllstand in Byte. availableTokens float64 // lastRefillTime ist der Zeitpunkt der letzten Nachfüllung. lastRefillTime time.Time // nowFunction liefert die aktuelle Zeit. // // Austauschbar, damit Tests das Zeitverhalten prüfen können, ohne // tatsächlich zu warten. Ein Test, der echte Sekunden verstreichen lässt, // wird entweder langsam oder unzuverlässig. nowFunction func() time.Time // sleepFunction wartet eine Dauer ab. sleepFunction func(waitContext context.Context, waitDuration time.Duration) error } // ErrInvalidRate meldet eine unbrauchbare Rate. var ErrInvalidRate = errors.New("die bandbreitengrenze ist unbrauchbar") // NewLimiter erzeugt einen Begrenzer. // // Eine Rate von Unlimited ergibt einen Begrenzer, der nichts begrenzt. Das ist // bequemer als überall zu prüfen, ob überhaupt einer vorhanden ist — und // verhindert den Fehler, den Begrenzer versehentlich zu übergehen. func NewLimiter(bytesPerSecond int64) (*Limiter, error) { if bytesPerSecond < 0 { return nil, fmt.Errorf("%w: %d byte je sekunde", ErrInvalidRate, bytesPerSecond) } if bytesPerSecond > 0 && bytesPerSecond < minimumBurstBytes { return nil, fmt.Errorf("%w: weniger als %d byte je sekunde sind nicht sinnvoll begrenzbar", ErrInvalidRate, minimumBurstBytes) } burstBytes := int64(float64(bytesPerSecond) * defaultBurstDuration.Seconds()) if burstBytes < minimumBurstBytes { burstBytes = minimumBurstBytes } return &Limiter{ bytesPerSecond: bytesPerSecond, burstBytes: burstBytes, availableTokens: float64(burstBytes), lastRefillTime: time.Now(), nowFunction: time.Now, sleepFunction: sleepWithContext, }, nil } // IsUnlimited meldet einen abgeschalteten Begrenzer. func (limiter *Limiter) IsUnlimited() bool { return limiter == nil || limiter.bytesPerSecond == Unlimited } // BytesPerSecond liefert die eingestellte Rate. func (limiter *Limiter) BytesPerSecond() int64 { if limiter == nil { return Unlimited } return limiter.bytesPerSecond } // Wait wartet, bis die angegebene Menge übertragen werden darf. // // Der Kontext wird beachtet: Ein Abbruch während einer langen Wartezeit muss // wirken. Andernfalls hinge ein abgebrochener Sicherungslauf noch minutenlang // in einer Wartezeit, die niemand mehr braucht. func (limiter *Limiter) Wait(waitContext context.Context, byteCount int64) error { if limiter.IsUnlimited() || byteCount <= 0 { return waitContext.Err() } // Eine Anforderung größer als der Eimer liesse sich nie erfüllen und // blockierte für immer. Sie wird in Teilstücke zerlegt. remainingBytes := byteCount for remainingBytes > 0 { if contextError := waitContext.Err(); contextError != nil { return contextError } requestedBytes := remainingBytes if requestedBytes > limiter.burstBytes { requestedBytes = limiter.burstBytes } waitDuration := limiter.reserve(requestedBytes) if waitDuration > 0 { if sleepError := limiter.sleepFunction(waitContext, waitDuration); sleepError != nil { return sleepError } } remainingBytes -= requestedBytes } return nil } // reserve entnimmt Token und liefert die nötige Wartezeit. func (limiter *Limiter) reserve(byteCount int64) time.Duration { limiter.mutex.Lock() defer limiter.mutex.Unlock() currentTime := limiter.nowFunction() elapsedSeconds := currentTime.Sub(limiter.lastRefillTime).Seconds() if elapsedSeconds > 0 { limiter.availableTokens += elapsedSeconds * float64(limiter.bytesPerSecond) if limiter.availableTokens > float64(limiter.burstBytes) { limiter.availableTokens = float64(limiter.burstBytes) } } limiter.lastRefillTime = currentTime // Die Token werden auch dann entnommen, wenn der Eimer leer ist: Der // Füllstand wird negativ und die Wartezeit ergibt sich daraus. So teilen // mehrere Arbeiter die Rate gerecht, statt dass der erste alles bekommt und // die übrigen wiederholt leer ausgehen. limiter.availableTokens -= float64(byteCount) if limiter.availableTokens >= 0 { return 0 } secondsToWait := -limiter.availableTokens / float64(limiter.bytesPerSecond) return time.Duration(secondsToWait * float64(time.Second)) } // sleepWithContext wartet und beachtet dabei einen Abbruch. func sleepWithContext(waitContext context.Context, waitDuration time.Duration) error { waitTimer := time.NewTimer(waitDuration) defer waitTimer.Stop() select { case <-waitContext.Done(): return waitContext.Err() case <-waitTimer.C: return nil } } // LimitedReader begrenzt den Durchsatz eines Datenstroms. type LimitedReader struct { // innerReader ist die eigentliche Quelle. innerReader io.Reader // limiter ist der gemeinsame Begrenzer. limiter *Limiter // readContext erlaubt den Abbruch während einer Wartezeit. // // io.Reader kennt keinen Kontext; ohne dieses Feld liesse sich ein // wartender Lesevorgang nicht beenden. readContext context.Context } // NewLimitedReader umhüllt einen Datenstrom mit einer Bandbreitengrenze. func NewLimitedReader(readContext context.Context, innerReader io.Reader, limiter *Limiter) io.Reader { // Ohne Begrenzung wird der Datenstrom unverändert durchgereicht: eine Hülle, // die nichts tut, kostet bei jedem Block einen Aufruf mehr. if limiter.IsUnlimited() { return innerReader } return &LimitedReader{ innerReader: innerReader, limiter: limiter, readContext: readContext, } } // Read liest und wartet entsprechend der Bandbreitengrenze. // // Gewartet wird **nach** dem Lesen, nicht davor. Vorher zu warten hiesse, für // eine Menge zu zahlen, die möglicherweise gar nicht kommt — am Ende einer // Datei wartete der Lauf für Bytes, die es nicht mehr gibt. func (limitedReader *LimitedReader) Read(targetBuffer []byte) (int, error) { readCount, readError := limitedReader.innerReader.Read(targetBuffer) if readCount > 0 { if waitError := limitedReader.limiter.Wait(limitedReader.readContext, int64(readCount)); waitError != nil { // Der Abbruch wiegt schwerer als ein gleichzeitiges Leseende: Wer // abbricht, will kein Ergebnis mehr. return readCount, waitError } } return readCount, readError } // LimitedWriter begrenzt den Durchsatz eines Schreibziels. type LimitedWriter struct { // innerWriter ist das eigentliche Ziel. innerWriter io.Writer // limiter ist der gemeinsame Begrenzer. limiter *Limiter // writeContext erlaubt den Abbruch während einer Wartezeit. writeContext context.Context } // NewLimitedWriter umhüllt ein Schreibziel mit einer Bandbreitengrenze. func NewLimitedWriter(writeContext context.Context, innerWriter io.Writer, limiter *Limiter) io.Writer { if limiter.IsUnlimited() { return innerWriter } return &LimitedWriter{ innerWriter: innerWriter, limiter: limiter, writeContext: writeContext, } } // Write schreibt und wartet entsprechend der Bandbreitengrenze. // // Hier wird **vor** dem Schreiben gewartet: Die Menge steht bereits fest, und // ein Schreibvorgang, der erst hinterher wartet, hätte die Leitung schon // belegt. func (limitedWriter *LimitedWriter) Write(sourceBuffer []byte) (int, error) { if waitError := limitedWriter.limiter.Wait(limitedWriter.writeContext, int64(len(sourceBuffer))); waitError != nil { return 0, waitError } return limitedWriter.innerWriter.Write(sourceBuffer) } // ParseBandwidthLimit liest eine Bandbreitenangabe wie "50MB" oder "1Gbit". // // Die doppelte Schreibweise ist Absicht: Netzwerkleute rechnen in Bit je // Sekunde, Speicherleute in Byte. Wer „100Mbit" meint und „100MB" einträgt, // vergibt das Achtfache — deshalb werden beide Einheiten angenommen und // auseinandergehalten. func ParseBandwidthLimit(limitText string) (int64, error) { trimmedText := trimSpaceAndLower(limitText) if trimmedText == "" || trimmedText == "0" || trimmedText == "unbegrenzt" { return Unlimited, nil } unitSuffixes := []struct { // suffix ist die erkannte Einheit. suffix string // multiplier ist der Umrechnungsfaktor in Byte je Sekunde. multiplier float64 }{ {"gbit", 1000 * 1000 * 1000 / 8}, {"mbit", 1000 * 1000 / 8}, {"kbit", 1000 / 8}, {"gb", 1 << 30}, {"mb", 1 << 20}, {"kb", 1 << 10}, {"b", 1}, } for _, unitSuffix := range unitSuffixes { numericPart, hasSuffix := cutSuffix(trimmedText, unitSuffix.suffix) if !hasSuffix { continue } numericValue, parseError := parseFloatValue(numericPart) if parseError != nil { return 0, fmt.Errorf("%w: %q enthält keine gültige zahl", ErrInvalidRate, limitText) } return int64(numericValue * unitSuffix.multiplier), nil } // Ohne Einheit wird Byte je Sekunde angenommen. numericValue, parseError := parseFloatValue(trimmedText) if parseError != nil { return 0, fmt.Errorf("%w: %q ist keine bandbreitenangabe (erwartet etwa 50MB oder 400Mbit)", ErrInvalidRate, limitText) } return int64(numericValue), nil }