Atlas › Test

TestClockedRateLimiter_Wait_DeadlineWouldExceed

Exact test identity: go.temporal.io/server/common/quotas/TestClockedRateLimiter_Wait_DeadlineWouldExceed

Package
go.temporal.io/server/common/quotas
Suite / test hierarchy
TestClockedRateLimiter_Wait_DeadlineWouldExceed
Test
TestClockedRateLimiter_Wait_DeadlineWouldExceed
Introduced at
clocked_rate_limiter.go ×1 Frontier kind: Joint frontier
Covered ranges
16
Covered lines
54
Covered files
2

Covered source

Expand a file to inspect source; the > gutter marks covered lines.

go.temporal.io/server/common/quotas/clocked_rate_limiter.go 43 covered LOC · 14 ranges

Open complete file

25 )
26
27 > func NewClockedRateLimiter(rateLimiter *rate.Limiter, timeSource clock.TimeSource) ClockedRateLimiter { clocked_rate_limiter.go
28 > return ClockedRateLimiter{
29 > rateLimiter: rateLimiter,
30 > timeSource: timeSource,
31 > recycleCh: make(chan struct{}),
32 > }
33 > }
34
35 > func (l ClockedRateLimiter) Allow() bool { clocked_rate_limiter.go
36 > return l.AllowN(l.timeSource.Now(), 1)
37 > }
38
39 > func (l ClockedRateLimiter) AllowN(now time.Time, token int) bool { clocked_rate_limiter.go
40 > return l.rateLimiter.AllowN(now, token)
41 > }
42
43 // ClockedReservation wraps a rate.Reservation with a clockwork.Clock. It is used to ensure that the reservation
48 }
49
50 > func (r ClockedReservation) OK() bool { clocked_rate_limiter.go
51 > return r.reservation.OK()
52 > }
53
54 > func (r ClockedReservation) Delay() time.Duration { clocked_rate_limiter.go
55 > return r.DelayFrom(r.timeSource.Now())
56 > }
57
58 > func (r ClockedReservation) DelayFrom(t time.Time) time.Duration { clocked_rate_limiter.go
59 > return r.reservation.DelayFrom(t)
60 > }
61
62 > func (r ClockedReservation) Cancel() { clocked_rate_limiter.go
63 > r.CancelAt(r.timeSource.Now())
64 > }
65
66 > func (r ClockedReservation) CancelAt(t time.Time) { clocked_rate_limiter.go
67 > r.reservation.CancelAt(t)
68 > }
69
70 func (l ClockedRateLimiter) Reserve() ClockedReservation {
77 }
78
79 > func (l ClockedRateLimiter) Wait(ctx context.Context) error { clocked_rate_limiter.go
80 > return l.WaitN(ctx, 1)
81 > }
82
83 // WaitN is the only method that is different from rate.Limiter. We need to fully reimplement this method because
84 // the original method uses time.Now(), and does not allow us to pass in a time.Time. Fortunately, it can be built on
85 // top of ReserveN. However, there are some optimizations that we can make.
86 > func (l ClockedRateLimiter) WaitN(ctx context.Context, token int) error { clocked_rate_limiter.go
87 > reservation := ClockedReservation{l.rateLimiter.ReserveN(l.timeSource.Now(), token), l.timeSource}
88 > if !reservation.OK() {
89 return fmt.Errorf("%w: WaitN(n=%d)", ErrRateLimiterReservationCannotBeMade, token)
90 }
91
92 > waitDuration := reservation.Delay() clocked_rate_limiter.go
93 >
94 > // Optimization: if the waitDuration is 0, we don't need to start a timer.
95 > if waitDuration <= 0 {
96 return nil
97 }
98
99 // Optimization: if the waitDuration is longer than the context deadline, we can immediately return an error.
100 > if deadline, ok := ctx.Deadline(); ok { clocked_rate_limiter.go
101 > if l.timeSource.Now().Add(waitDuration).After(deadline) { clocked_rate_limiter.go
102 > reservation.Cancel() clocked_rate_limiter.go
103 > return fmt.Errorf("%w: WaitN(n=%d)", ErrRateLimiterReservationWouldExceedContextDeadline, token)
104 > }
105 }
106
go.temporal.io/server/common/clock/event_time_source.go 11 covered LOC · 2 ranges

Open complete file

39
40 // NewEventTimeSource returns a EventTimeSource with the current time set to Unix zero: 1970-01-01 00:00:00 +0000 UTC.
41 > func NewEventTimeSource() *EventTimeSource { event_time_source.go
42 > return &EventTimeSource{
43 > now: time.Unix(0, 0),
44 > }
45 > }
46
47 // Some clients depend on the fact that the runtime's timers do _not_ run synchronously.
55
56 // Now return the current time.
57 > func (ts *EventTimeSource) Now() time.Time { event_time_source.go
58 > ts.mu.RLock()
59 > defer ts.mu.RUnlock()
60 >
61 > return ts.now
62 > }
63
64 func (ts *EventTimeSource) Since(t time.Time) time.Duration {