Atlas › Test

TestDelayedRequestRateLimiter_Allow

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

Package
go.temporal.io/server/common/quotas
Suite / test hierarchy
TestDelayedRequestRateLimiter_Allow
Test
TestDelayedRequestRateLimiter_Allow
Introduced at
TestDelayedRequestRateLimiter_Allow Frontier kind: Test frontier
Covered ranges
23
Covered lines
68
Covered files
4

Covered source

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

go.temporal.io/server/common/clock/event_time_source.go 44 covered LOC · 15 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 {
71 // wrap all such calls in a goroutine. If the duration is non-positive, the callback will fire immediately before
72 // AfterFunc returns.
73 > func (ts *EventTimeSource) AfterFunc(d time.Duration, f func()) Timer { event_time_source.go
74 > if d < 0 {
75 d = 0
76 }
77 > timer := &fakeTimer{timeSource: ts, deadline: ts.Now().Add(d), callback: f} event_time_source.go
78 > ts.addTimer(timer)
79 > return timer
80 }
81
96 }
97
98 > func (ts *EventTimeSource) addTimer(t *fakeTimer) { event_time_source.go
99 > ts.mu.Lock()
100 > defer ts.mu.Unlock()
101 > t.index = len(ts.timers)
102 > ts.timers = append(ts.timers, t)
103 > ts.fireTimers()
104 > }
105
106 // Update the fake current time. It returns the timeSource so that you can chain calls like this:
116
117 // Advance the timer by the specified duration.
118 > func (ts *EventTimeSource) Advance(d time.Duration) { event_time_source.go
119 > ts.mu.Lock()
120 > defer ts.mu.Unlock()
121 >
122 > ts.now = ts.now.Add(d)
123 > ts.fireTimers()
124 > }
125
126 // AdvanceNext advances to the next timer.
156
157 // fireTimers fires all timers that are ready.
158 > func (ts *EventTimeSource) fireTimers() { event_time_source.go
159 > n := 0
160 > for _, t := range ts.timers {
161 > if t.deadline.After(ts.now) { event_time_source.go
162 > ts.timers[n] = t event_time_source.go
163 > t.index = n
164 > n++
165 > } else { event_time_source.go
166 > if ts.async { event_time_source.go
167 go t.callback()
168 > } else { event_time_source.go
169 > t.callback() event_time_source.go
170 > }
171 > t.done = true event_time_source.go
172 }
173 }
174 > ts.timers = ts.timers[:n] event_time_source.go
175 }
176
go.temporal.io/server/common/quotas/delayed_request_rate_limiter.go 12 covered LOC · 4 ranges

Open complete file

30 delay time.Duration,
31 timeSource clock.TimeSource,
32 > ) (*DelayedRequestRateLimiter, error) { delayed_request_rate_limiter.go
33 > if delay < 0 {
34 return nil, fmt.Errorf("%w: %v", ErrNegativeDelay, delay)
35 }
36
37 > delegator := RequestRateLimiterDelegator{} delayed_request_rate_limiter.go
38 > delegator.SetRateLimiter(NoopRequestRateLimiter)
39 >
40 > timer := timeSource.AfterFunc(delay, func() {
41 > delegator.SetRateLimiter(rl) delayed_request_rate_limiter.go
42 > })
43
44 > return &DelayedRequestRateLimiter{ delayed_request_rate_limiter.go
45 > RequestRateLimiter: &delegator,
46 > timer: timer,
47 > }, nil
48 }
49
go.temporal.io/server/common/quotas/request_rate_limiter_delegator.go 9 covered LOC · 3 ranges

Open complete file

24
25 // SetRateLimiter sets the rate limiter to delegate to.
26 > func (d *RequestRateLimiterDelegator) SetRateLimiter(rl RequestRateLimiter) { request_rate_limiter_delegator.go
27 > d.delegate.Store(monomorphicRequestRateLimiter{rl})
28 > }
29
30 // loadDelegate returns the rate limiter that this rate limiter delegates to.
31 > func (d *RequestRateLimiterDelegator) loadDelegate() RequestRateLimiter { request_rate_limiter_delegator.go
32 > return d.delegate.Load().(RequestRateLimiter)
33 > }
34
35 // The following methods just delegate to the underlying rate limiter.
36
37 > func (d *RequestRateLimiterDelegator) Allow(now time.Time, request Request) bool { request_rate_limiter_delegator.go
38 > return d.loadDelegate().Allow(now, request)
39 > }
40
41 func (d *RequestRateLimiterDelegator) Reserve(now time.Time, request Request) Reservation {
go.temporal.io/server/common/quotas/noop_request_rate_limiter_impl.go 3 covered LOC · 1 range

Open complete file