Atlas › Test

TestTryAcquire_HighAfterWaiting

Exact test identity: go.temporal.io/server/common/locks/TestPrioritySemaphoreSuite/TestTryAcquire_HighAfterWaiting

Package
go.temporal.io/server/common/locks
Suite / test hierarchy
TestPrioritySemaphoreSuite/TestTryAcquire_HighAfterWaiting
Test
TestTryAcquire_HighAfterWaiting
Introduced at
TestTryAcquire_LowAfterWaiting, TestTryAcquire_HighAfterWaiting Frontier kind: Test frontier
Covered ranges
20
Covered lines
62
Covered files
1

Co-introduced tests

1 other test enter at the same concept.

Covered source

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

go.temporal.io/server/common/locks/priority_semaphore_impl.go 62 covered LOC · 20 ranges

Open complete file

67 // maximum combined weight for concurrent access, capable of handling multiple priority levels.
68 // Most of the logic is taken directly from golang's semaphore.Weighted.
69 > func NewPrioritySemaphore(n int) *PrioritySemaphoreImpl { priority_semaphore_impl.go
70 > waitLists := make([]*list.List, NumPriorities)
71 > for i := range waitLists {
72 > waitLists[i] = list.New()
73 > }
74 > return &PrioritySemaphoreImpl{
75 > size: n,
76 > waitLists: waitLists,
77 > }
78 }
79
81 // are available or ctx is done. On success, returns nil. On failure, returns
82 // ctx.Err() and leaves the semaphore unchanged.
83 > func (s *PrioritySemaphoreImpl) Acquire(ctx context.Context, priority Priority, n int) error { priority_semaphore_impl.go
84 > if priority >= NumPriorities {
85 // nolint:forbidigo
86 panic(fmt.Sprintf("semaphore: invalid priority %v, priority must be less than %v", priority, NumPriorities))
87 }
88
89 > done := ctx.Done() priority_semaphore_impl.go
90 >
91 > s.mu.Lock()
92 > select {
93 case <-done:
94 // ctx becoming done has "happened before" acquiring the semaphore,
98 s.mu.Unlock()
99 return ctx.Err()
101 }
102 // Check if acquisition can proceed without waiting
103 > if s.size-s.cur >= n && s.noWaiters(priority) { priority_semaphore_impl.go
104 // Since we hold s.mu and haven't synchronized since checking done, if
105 // ctx becomes done before we return here, it becoming done must have
111 }
112
113 > if n > s.size { priority_semaphore_impl.go
114 s.mu.Unlock()
115 return ErrRequestTooLarge
116 }
117
118 > ready := make(chan struct{}) priority_semaphore_impl.go
119 > w := waiter{n: n, ready: ready}
120 > elem := s.waitLists[priority].PushBack(w)
121 > s.mu.Unlock()
122 >
123 > select {
124 case <-done:
125 s.mu.Lock()
141 return ctx.Err()
142
143 > case <-ready: priority_semaphore_impl.go
144 > // Acquired the semaphore. Check that ctx isn't already done.
145 > // We check the done channel instead of calling ctx.Err because we
146 > // already have the channel, and ctx.Err is O(n) with the nesting
147 > // depth of ctx.
148 > select {
149 case <-done:
150 s.Release(n)
151 return ctx.Err()
153 }
154 > return nil priority_semaphore_impl.go
155 }
156 }
158 // TryAcquire acquires the semaphore with a weight of n without blocking.
159 // On success, returns true. On failure, returns false and leaves the semaphore unchanged.
160 > func (s *PrioritySemaphoreImpl) TryAcquire(priority Priority, n int) bool { priority_semaphore_impl.go
161 > if priority >= NumPriorities {
162 // nolint:forbidigo
163 panic(fmt.Sprintf("semaphore: invalid priority %v, priority must be less than %v", priority, NumPriorities))
164 }
165
166 > s.mu.Lock() priority_semaphore_impl.go
167 > defer s.mu.Unlock()
168 > if s.size-s.cur >= n && s.noWaiters(priority) {
169 > s.cur += n
170 > return true
171 > }
172 return false
173 }
174
175 > func (s *PrioritySemaphoreImpl) Release(n int) { priority_semaphore_impl.go
176 > s.mu.Lock()
177 > defer s.mu.Unlock()
178 > s.cur -= n
179 > if s.cur < 0 {
180 s.mu.Unlock()
181 panic("semaphore: released more than held")
182 }
183 > s.notifyWaiters() priority_semaphore_impl.go
184 }
185
186 > func (s *PrioritySemaphoreImpl) notifyWaiters() { priority_semaphore_impl.go
187 > for _, l := range s.waitLists {
188 > for {
189 > next := l.Front()
190 > if next == nil {
191 > break // No more waiters blocked.
192 }
193
194 > w, ok := next.Value.(waiter) priority_semaphore_impl.go
195 > if !ok {
196 panic("semaphore: failed to cast waiter")
197 }
198 > if s.size-s.cur < w.n { priority_semaphore_impl.go
199 // Not enough tokens for the next waiter. We could keep going (to try to
200 // find a waiter with a smaller request), but under load that could cause
219
220 // noWaiters returns if there is no waiter that has priority higher or equal to lowestPriority.
221 > func (s *PrioritySemaphoreImpl) noWaiters(lowestPriority Priority) bool { priority_semaphore_impl.go
222 > for _, l := range s.waitLists[:lowestPriority+1] {
223 > if l.Len() > 0 {
224 return false
225 }
226 }
227 > return true priority_semaphore_impl.go
228 }