fix(gateway): bound impairment catchup
Verify Data Plane / gateway (push) Failing after 1m39s

This commit is contained in:
sechmachine
2026-07-30 15:27:51 +07:00
parent d3194a85af
commit 719042aa45
2 changed files with 32 additions and 10 deletions
+18
View File
@@ -346,6 +346,24 @@ func TestQualificationLossOnlyDoesNotImplicitlyReorder(t *testing.T) {
}
}
func TestQualificationSourceShaperCarriesAtMostOnePacketOfCatchup(t *testing.T) {
spacing := time.Millisecond
started := time.Unix(0, 0)
now := started.Add(100 * spacing)
first := qualificationBoundedRelease(started, time.Time{}, now, spacing)
if first != now.Add(-spacing) {
t.Fatalf("first catch-up release = %s, want %s", first, now.Add(-spacing))
}
second := qualificationBoundedRelease(started.Add(spacing), first.Add(spacing), now, spacing)
if second != now {
t.Fatalf("second catch-up release = %s, want %s", second, now)
}
third := qualificationBoundedRelease(started.Add(2*spacing), second.Add(spacing), now, spacing)
if third != now.Add(spacing) {
t.Fatalf("catch-up debt was reset: third release = %s, want %s", third, now.Add(spacing))
}
}
func TestQualificationExplicitReorderIsBoundedAndAttributed(t *testing.T) {
observation, err := runQualificationImpairment(
t,
+14 -10
View File
@@ -280,6 +280,17 @@ func qualificationMediaPacerKbps(profile qualificationMediaProfile, reduction in
return (payloadKbps*int64(profile.PacketBytes+frameHeaderSize) + int64(profile.PacketBytes) - 1) / int64(profile.PacketBytes)
}
func qualificationBoundedRelease(target, next, now time.Time, spacing time.Duration) time.Time {
release := target
if next.After(release) {
release = next
}
if now.Sub(release) > spacing {
return now.Add(-spacing)
}
return release
}
type qualificationTracingBackend struct {
native *NativeApolloBackend
setups atomic.Uint64
@@ -1326,20 +1337,13 @@ func runQualificationImpairment(t *testing.T, profile qualificationImpairmentPro
stepAt := make(map[int]time.Time, len(profile.CapacitySteps))
emitted := 0
var previousRelease time.Time
maxCatchup := spacing
var nextRelease time.Time
for _, packet := range jobs {
release := started.Add(packet.target)
if minimum := previousRelease.Add(spacing); !previousRelease.IsZero() && release.Before(minimum) {
release = minimum
}
if lag := time.Since(release); lag > maxCatchup {
release = release.Add(lag - maxCatchup)
}
release := qualificationBoundedRelease(started.Add(packet.target), nextRelease, time.Now(), spacing)
if delay := time.Until(release); delay > 0 {
time.Sleep(delay)
}
previousRelease = release
nextRelease = release.Add(spacing)
if len(profile.CapacitySteps) == 2 {
switch {
case packet.index >= packetCount*2/3 && stepAt[profile.CapacitySteps[1]].IsZero():