refactor bandwidth throttling for replication target (#17980)

This refactor is to allow using the bandwidth throttling
for other purposes.
This commit is contained in:
Harshavardhana
2023-09-05 20:21:59 -07:00
committed by GitHub
parent 812f5a02d7
commit 5b114b43f7
8 changed files with 153 additions and 179 deletions

View File

@@ -29,7 +29,7 @@ const (
func TestMonitor_GetReport(t *testing.T) {
type fields struct {
activeBuckets map[string]map[string]*bucketMeasurement
activeBuckets map[BucketOptions]*bucketMeasurement
endTime time.Time
update2 uint64
endTime2 time.Time
@@ -39,6 +39,28 @@ func TestMonitor_GetReport(t *testing.T) {
m0.incrementBytes(0)
m1MiBPS := newBucketMeasurement(start)
m1MiBPS.incrementBytes(oneMiB)
test1Want := make(map[BucketOptions]Details)
test1Want[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = Details{LimitInBytesPerSecond: 1024 * 1024, CurrentBandwidthInBytesPerSecond: 0}
test1Want2 := make(map[BucketOptions]Details)
test1Want2[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = Details{
LimitInBytesPerSecond: 1024 * 1024,
CurrentBandwidthInBytesPerSecond: (1024 * 1024) / start.Add(2*time.Second).Sub(start.Add(1*time.Second)).Seconds(),
}
test2Want := make(map[BucketOptions]Details)
test2Want[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = Details{LimitInBytesPerSecond: 1024 * 1024, CurrentBandwidthInBytesPerSecond: float64(oneMiB)}
test2Want2 := make(map[BucketOptions]Details)
test2Want2[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = Details{
LimitInBytesPerSecond: 1024 * 1024,
CurrentBandwidthInBytesPerSecond: exponentialMovingAverage(betaBucket, float64(oneMiB), 2*float64(oneMiB)),
}
test1ActiveBuckets := make(map[BucketOptions]*bucketMeasurement)
test1ActiveBuckets[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = m0
test1ActiveBuckets2 := make(map[BucketOptions]*bucketMeasurement)
test1ActiveBuckets2[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = m1MiBPS
tests := []struct {
name string
fields fields
@@ -48,46 +70,31 @@ func TestMonitor_GetReport(t *testing.T) {
{
name: "ZeroToOne",
fields: fields{
activeBuckets: map[string]map[string]*bucketMeasurement{
"bucket": {
"arn": m0,
},
},
endTime: start.Add(1 * time.Second),
update2: oneMiB,
endTime2: start.Add(2 * time.Second),
activeBuckets: test1ActiveBuckets,
endTime: start.Add(1 * time.Second),
update2: oneMiB,
endTime2: start.Add(2 * time.Second),
},
want: &BucketBandwidthReport{
BucketStats: map[string]map[string]Details{
"bucket": {
"arn": Details{LimitInBytesPerSecond: 1024 * 1024, CurrentBandwidthInBytesPerSecond: 0},
},
},
BucketStats: test1Want,
},
want2: &BucketBandwidthReport{
BucketStats: map[string]map[string]Details{"bucket": {"arn": Details{LimitInBytesPerSecond: 1024 * 1024, CurrentBandwidthInBytesPerSecond: (1024 * 1024) / start.Add(2*time.Second).Sub(start.Add(1*time.Second)).Seconds()}}},
BucketStats: test1Want2,
},
},
{
name: "OneToTwo",
fields: fields{
activeBuckets: map[string]map[string]*bucketMeasurement{
"bucket": {
"arn": m1MiBPS,
},
},
endTime: start.Add(1 * time.Second),
update2: 2 * oneMiB,
endTime2: start.Add(2 * time.Second),
activeBuckets: test1ActiveBuckets2,
endTime: start.Add(1 * time.Second),
update2: 2 * oneMiB,
endTime2: start.Add(2 * time.Second),
},
want: &BucketBandwidthReport{
BucketStats: map[string]map[string]Details{"bucket": {"arn": Details{LimitInBytesPerSecond: 1024 * 1024, CurrentBandwidthInBytesPerSecond: float64(oneMiB)}}},
BucketStats: test2Want,
},
want2: &BucketBandwidthReport{
BucketStats: map[string]map[string]Details{"bucket": {"arn": Details{
LimitInBytesPerSecond: 1024 * 1024,
CurrentBandwidthInBytesPerSecond: exponentialMovingAverage(betaBucket, float64(oneMiB), 2*float64(oneMiB)),
}}},
BucketStats: test2Want2,
},
},
}
@@ -95,23 +102,23 @@ func TestMonitor_GetReport(t *testing.T) {
tt := tt
t.Run(tt.name, func(t *testing.T) {
t.Parallel()
thr := throttle{
thr := bucketThrottle{
NodeBandwidthPerSec: 1024 * 1024,
}
th := make(map[string]map[string]*throttle)
th["bucket"] = map[string]*throttle{"arn": &thr}
th := make(map[BucketOptions]*bucketThrottle)
th[BucketOptions{Name: "bucket", ReplicationARN: "arn"}] = &thr
m := &Monitor{
activeBuckets: tt.fields.activeBuckets,
bucketThrottle: th,
NodeCount: 1,
bucketsMeasurement: tt.fields.activeBuckets,
bucketsThrottle: th,
NodeCount: 1,
}
m.activeBuckets["bucket"]["arn"].updateExponentialMovingAverage(tt.fields.endTime)
m.bucketsMeasurement[BucketOptions{Name: "bucket", ReplicationARN: "arn"}].updateExponentialMovingAverage(tt.fields.endTime)
got := m.GetReport(SelectBuckets())
if !reflect.DeepEqual(got, tt.want) {
t.Errorf("GetReport() = %v, want %v", got, tt.want)
}
m.activeBuckets["bucket"]["arn"].incrementBytes(tt.fields.update2)
m.activeBuckets["bucket"]["arn"].updateExponentialMovingAverage(tt.fields.endTime2)
m.bucketsMeasurement[BucketOptions{Name: "bucket", ReplicationARN: "arn"}].incrementBytes(tt.fields.update2)
m.bucketsMeasurement[BucketOptions{Name: "bucket", ReplicationARN: "arn"}].updateExponentialMovingAverage(tt.fields.endTime2)
got = m.GetReport(SelectBuckets())
if !reflect.DeepEqual(got.BucketStats, tt.want2.BucketStats) {
t.Errorf("GetReport() = %v, want %v", got.BucketStats, tt.want2.BucketStats)