///|
async test "mutex serializes access" {
let mutex = @semaphore.Semaphore(1)
let log = []
@async.with_task_group group => {
for i in 0..<3 {
group.spawn_bg() () => {
mutex.acquire()
defer mutex.release()
log.push("task \{i}: enter")
@async.sleep(50)
log.push("task \{i}: leave")
}
}
}
// Every "enter" is followed by the same task's "leave" before
// any other task enters.
json_inspect(log, content=[
"task 0: enter", "task 0: leave", "task 1: enter", "task 1: leave", "task 2: enter",
"task 2: leave",
])
}///|
async test "bounded concurrency caps in-flight tasks" {
let limit = @semaphore.Semaphore(3)
let mut active = 0
let mut peak = 0
@async.with_task_group group => {
for _ in 0..<8 {
group.spawn_bg() () => {
limit.acquire()
defer limit.release()
active = active + 1
if active > peak {
peak = active
}
@async.sleep(30)
active = active - 1
}
}
}
// 8 tasks competed for 3 permits; peak in-flight is exactly 3.
inspect(peak, content="3")
// All tasks finished — the counter is back to zero.
inspect(active, content="0")
}///|
async test "one-shot signal" {
// size 1, but starts empty — acquire blocks until release().
let ready = @semaphore.Semaphore(1, initial_value=0)
let log = []
@async.with_task_group group => {
group.spawn_bg() () => {
log.push("waiter: parking")
ready.acquire()
log.push("waiter: unblocked")
}
group.spawn_bg() () => {
log.push("signaller: setup")
@async.sleep(50)
log.push("signaller: release")
ready.release()
}
}
json_inspect(log, content=[
"waiter: parking", "signaller: setup", "signaller: release", "waiter: unblocked",
])
}Caveat: release() wakes one waiter. If multiple tasks need to react to the same event, use @async.Cond (broadcast notification) instead.
///|
async test "count-down latch waits for N events" {
let n = 3
let done = @semaphore.Semaphore(n, initial_value=0)
let log = []
@async.with_task_group group => {
for i in 0..<n {
group.spawn_bg() () => {
@async.sleep(20 * (i + 1))
log.push("worker \{i} done")
done.release()
}
}
for _ in 0..<n {
done.acquire()
}
log.push("all workers reported")
}
json_inspect(log, content=[
"worker 0 done", "worker 1 done", "worker 2 done", "all workers reported",
])
}In most real code, prefer with_task_group's implicit join over a count-down latch — the group already waits for every child to finish. Reach for a latch only when you want the coordinator to proceed mid-group, before the children's resources are released.
///|
test "try_acquire never blocks" {
let sem = @semaphore.Semaphore(2)
assert_true(sem.try_acquire()) // 2 -> 1
assert_true(sem.try_acquire()) // 1 -> 0
assert_false(sem.try_acquire()) // 0 -> 0 (no permits, returns false)
sem.release() // 0 -> 1
assert_true(sem.try_acquire()) // 1 -> 0
}///|
async test "waiters are woken FIFO" {
let sem = @semaphore.Semaphore(1)
sem.acquire() // main holds the permit
let log = []
@async.with_task_group group => {
// Task A starts waiting at ~0 ms.
group.spawn_bg() () => {
sem.acquire()
log.push("A acquired")
sem.release()
}
// Task B starts waiting at ~50 ms.
group.spawn_bg() () => {
@async.sleep(50)
sem.acquire()
log.push("B acquired")
sem.release()
}
@async.sleep(100)
sem.release() // main hands off the permit
}
// A queued first, so A wakes first.
json_inspect(log, content=["A acquired", "B acquired"])
}pub struct Semaphore {
size : Int
// private fields
}Asynchronous programming library for MoonBit