Atlas › Test
TestAcquire_Low_Success
Exact test identity: go.temporal.io/server/common/locks/TestPrioritySemaphoreSuite/TestAcquire_Low_Success
- Package
go.temporal.io/server/common/locks
- Suite / test hierarchy
TestPrioritySemaphoreSuite/TestAcquire_Low_Success
- Test
TestAcquire_Low_Success
- Introduced at
- TestAcquire_Low_Success, TestAcquire_High_Success Frontier kind: Test frontier
- Covered ranges
- 14
- Covered lines
- 50
- 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 50 covered LOC · 14 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.
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.
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
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
105
>
// ctx becomes done before we return here, it becoming done must have
106
>
// "happened concurrently" with this call - it cannot "happen before"
107
>
// we return in this branch. So, we're ok to always acquire here.
108
>
s.cur += n
109
>
s.mu.Unlock()
110
>
return nil
111
>
}
112
113
if n > s.size {
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.
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
167
>
defer s.mu.Unlock()
168
>
if s.size-s.cur >= n && s.noWaiters(priority) {
169
>
s.cur += n
170
>
return true
171
>
}
173
}
174
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
}
184
}
185
187
>
for _, l := range s.waitLists {
188
>
for {
189
>
next := l.Front()
190
>
if next == nil {
191
>
break // No more waiters blocked.
192
}
193
219
220
// noWaiters returns if there is no waiter that has priority higher or equal to lowestPriority.
222
>
for _, l := range s.waitLists[:lowestPriority+1] {
223
>
if l.Len() > 0 {
224
return false
225
}
226
}
228
}