Concurrency Coordination Reports
Queue Drain Coordination Report
Drain a prefilled queue with nonblocking pops and report whether work was left idle, fully drained, or the drain ran short.
nonblocking pop
`Queue#pop(true)` removes a buffered item without waiting and raises `ThreadError` when the queue is empty, so counting successful pops against requested pops coordinates a consumer with the work that was actually queued.
Queue Drain Coordination Report
queue_drain_coordination_report.rb
Replay: real traced execution (multi-file project)
requests = 2
queue = Queue.new
queue.push("job1")
queue.push("job2")
drained = 0
short = 0
requests.times do
begin
queue.pop(true)
drained += 1
rescue ThreadError
short += 1
end
end
remaining = queue.size
status = if remaining > 0
"idle"
elsif short > 0
"short"
else
"drained"
end
puts "requests=#{requests}"
puts "drained=#{drained}"
puts "remaining=#{remaining}"
puts "status=#{status}"
requests = 1
queue = Queue.new
queue.push("job1")
queue.push("job2")
drained = 0
short = 0
requests.times do
begin
queue.pop(true)
drained += 1
rescue ThreadError
short += 1
end
end
remaining = queue.size
status = if remaining > 0
"idle"
elsif short > 0
"short"
else
"drained"
end
puts "requests=#{requests}"
puts "drained=#{drained}"
puts "remaining=#{remaining}"
puts "status=#{status}"
requests = 4
queue = Queue.new
queue.push("job1")
queue.push("job2")
drained = 0
short = 0
requests.times do
begin
queue.pop(true)
drained += 1
rescue ThreadError
short += 1
end
end
remaining = queue.size
status = if remaining > 0
"idle"
elsif short > 0
"short"
else
"drained"
end
puts "requests=#{requests}"
puts "drained=#{drained}"
puts "remaining=#{remaining}"
puts "status=#{status}"
requests ← 2, queue ← ⟨Thread::Queue A⟩, drained ← 0, short ← 0
1requests→ 2 = 2 #@requests=1, 42queue→ ⟨Thread::Queue A⟩ = Queue.new3queue⟨Thread::Queue A⟩.push("job1")4queue⟨Thread::Queue A⟩.push("job2")5drained→ 0 = 06short→ 0 = 078requests2.times do9 begin10 queue.pop(true)11 drained += 112 rescue ThreadError13 short += 114 end15enddo
pass 1 of 28requests.times do9 begin10 queue.pop(true)do
pass 2 of 28requests2.times do9 begin10 queue.pop(true)11 drained += 112 rescue ThreadError13 short += 114 end15endremaining ← 0, status ← drained
17remaining→ 0 = queue.size01819status→ drained = if remaining0 > 020 "idle"21elsif short0 > 022 "short"23else24 "drained"25end2627puts "requests=#{requests2}"28puts "drained=#{drained2}"29puts "remaining=#{remaining0}"30puts "status=#{statusdrained}"outputrequests=2 drained=2 remaining=0 status=drained
requests ← 1, queue ← ⟨Thread::Queue A⟩, drained ← 0, short ← 0
1requests→ 1 = 12queue→ ⟨Thread::Queue A⟩ = Queue.new3queue⟨Thread::Queue A⟩.push("job1")4queue⟨Thread::Queue A⟩.push("job2")5drained→ 0 = 06short→ 0 = 078requests1.times do9 begin10 queue.pop(true)11 drained += 112 rescue ThreadError13 short += 114 end15enddo
8requests1.times do9 begin10 queue.pop(true)11 drained += 112 rescue ThreadError13 short += 114 end15endremaining ← 1, status ← idle
17remaining→ 1 = queue.size11819status→ idle = if remaining1 > 020 "idle"21elsif short0 > 022 "short"23else24 "drained"25end2627puts "requests=#{requests1}"28puts "drained=#{drained1}"29puts "remaining=#{remaining1}"30puts "status=#{statusidle}"outputrequests=1 drained=1 remaining=1 status=idle
requests ← 4, queue ← ⟨Thread::Queue A⟩, drained ← 0, short ← 0
1requests→ 4 = 42queue→ ⟨Thread::Queue A⟩ = Queue.new3queue⟨Thread::Queue A⟩.push("job1")4queue⟨Thread::Queue A⟩.push("job2")5drained→ 0 = 06short→ 0 = 078requests4.times do9 begin10 queue.pop(true)11 drained += 112 rescue ThreadError13 short += 114 end15enddo
pass 1 of 48requests.times do9 begin10 queue.pop(true)All 4 passes — pass 1 is the card above pass requests1 — 2 — 3 — 4 4 remaining ← 0, status ← short
17remaining→ 0 = queue.size01819status→ short = if remaining0 > 020 "idle"21elsif short2 > 022 "short"23else24 "drained"25end2627puts "requests=#{requests4}"28puts "drained=#{drained2}"29puts "remaining=#{remaining0}"30puts "status=#{statusshort}"outputrequests=4 drained=2 remaining=0 status=short