Repository navigation
Expand file tree
/
Copy pathlimit_test.go
More file actions
124 lines (96 loc) · 2.84 KB
/
Copy pathlimit_test.go
File metadata and controls
124 lines (96 loc) · 2.84 KB
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
package httpstream
import (
"context"
"errors"
"net/http"
"sync"
"testing"
"time"
"github.com/stretchr/testify/assert"
)
type mockRoundTripper struct {
roundTripFunc func(req *http.Request) (*http.Response, error)
}
func (m *mockRoundTripper) RoundTrip(req *http.Request) (*http.Response, error) {
return m.roundTripFunc(req)
}
func TestConcurrencyLimit(t *testing.T) {
// Semaphore limit is 2
middleware := ConcurrencyMiddleware(2)
var wg sync.WaitGroup
var activeCount int
var maxActive int
var mu sync.Mutex
// Mock roundtripper that simulates work by sleeping
mockRT := &mockRoundTripper{
roundTripFunc: func(req *http.Request) (*http.Response, error) {
mu.Lock()
activeCount++
if activeCount > maxActive {
maxActive = activeCount
}
mu.Unlock()
// Simulate request duration
time.Sleep(50 * time.Millisecond)
mu.Lock()
activeCount--
mu.Unlock()
return &http.Response{StatusCode: 200}, nil
},
}
limiter := middleware(mockRT)
// Launch 5 concurrent requests
for i := 0; i < 5; i++ {
wg.Add(1)
go func() {
defer wg.Done()
req, _ := http.NewRequestWithContext(context.Background(), "GET", "http://example.com", nil)
resp, err := limiter.RoundTrip(req)
assert.NoError(t, err)
assert.Equal(t, 200, resp.StatusCode)
}()
}
wg.Wait()
// Max concurrent active requests must not exceed 2
assert.True(t, maxActive <= 2, "max active requests was %d, expected <= 2", maxActive)
}
func TestConcurrencyLimitContextCancellation(t *testing.T) {
// Semaphore limit is 1
middleware := ConcurrencyMiddleware(1)
// Block channel to keep the first request active
blockChan := make(chan struct{})
mockRT := &mockRoundTripper{
roundTripFunc: func(req *http.Request) (*http.Response, error) {
<-blockChan
return &http.Response{StatusCode: 200}, nil
},
}
limiter := middleware(mockRT)
// Start 1st request (occupies the only slot)
go func() {
req, _ := http.NewRequestWithContext(context.Background(), "GET", "http://example.com", nil)
_, _ = limiter.RoundTrip(req)
}()
// Give the 1st request time to start and acquire the slot
time.Sleep(10 * time.Millisecond)
// Start 2nd request with a cancellable context
ctx, cancel := context.WithCancel(context.Background())
req2, _ := http.NewRequestWithContext(ctx, "GET", "http://example.com", nil)
var err2 error
var wg sync.WaitGroup
wg.Add(1)
go func() {
defer wg.Done()
_, err2 = limiter.RoundTrip(req2)
}()
// Give the 2nd request time to block on acquiring slot
time.Sleep(10 * time.Millisecond)
// Cancel the context of the 2nd request
cancel()
wg.Wait()
// Verify that the 2nd request failed with context.Canceled immediately
assert.Error(t, err2)
assert.True(t, errors.Is(err2, context.Canceled), "expected context.Canceled, got %v", err2)
// Cleanup block channel to release the 1st request
close(blockChan)
}