A bounded FIFO queue can block producers when full and consumers when empty. Closing it needs an explicit state transition: reject any later offer, wake producers waiting for space, and let consumers drain work already accepted before returning a closed signal. This implementation guards a deque and closed flag with one condition variable. Waits occur inside while loops because waking does not guarantee the desired condition still holds. None is reserved as the drained-and-closed result. The example is a thread-safe in-process buffer with a single queue state; it is not a lock-free queue, a durable message broker, or a task-completion tracker.
Bounded queues: close, wake waiters, and drain accepted work
Operational case
A two-slot intake accepts J-47 and J-52. A nonblocking offer of J-61 reports false while the queue is full. Closing intake leaves the two accepted jobs in FIFO order: take returns J-47, then J-52, then None. Offering J-83 after closure raises QueueClosed. In a threaded run, a producer blocked on a full queue wakes when close broadcasts the state change and raises instead of waiting forever. A consumer blocked on an empty queue also wakes and sees the drained signal. Draining is a policy choice; discarding accepted work would need a different documented shutdown path.
Working Python program
from collections import deque
from threading import Condition
class QueueClosed(Exception):
pass
class ClosingWorkQueue:
def __init__(self, capacity):
if capacity < 1:
raise ValueError("capacity must be positive")
self.capacity = capacity
self.jobs = deque()
self.closed = False
self.changed = Condition()
def offer(self, job_id, wait=True):
if job_id is None:
raise ValueError("None is reserved for drained closure")
with self.changed:
while len(self.jobs) == self.capacity and not self.closed:
if not wait:
return False
self.changed.wait()
if self.closed:
raise QueueClosed("queue is closed")
self.jobs.append(job_id)
self.changed.notify()
return True
def take(self):
with self.changed:
while not self.jobs and not self.closed:
self.changed.wait()
if not self.jobs:
return None
job_id = self.jobs.popleft()
self.changed.notify()
return job_id
def close(self):
with self.changed:
self.closed = True
self.changed.notify_all()
intake = ClosingWorkQueue(2)
intake.offer("J-47")
intake.offer("J-52")
print(intake.offer("J-61", wait=False))
intake.close()
print(intake.take(), intake.take(), intake.take())
try:
intake.offer("J-83")
except QueueClosed:
print("closed")Output
False
J-47 J-52 None
closedTime, space, and tradeoff
Appending and removing from the deque are O(1) while the lock is held. A full producer or empty consumer can wait for an unbounded wall-clock time until another thread changes state or close wakes it. The buffer uses O(C) space for capacity C, plus condition-waiter bookkeeping owned by the runtime. Close itself changes state in O(1) and broadcasts to waiting threads; waking many threads has scheduling cost outside a constant-time queue-operation claim. A taken job is merely removed from memory. If processing fails afterward, this structure has no retry record, acknowledgment, or crash recovery.
Common Mistakes
- Do not use an if check before wait; recheck the queue state after every wake.
- Do not close without waking blocked producers and consumers.
- Do not treat queue removal as proof that a worker completed its job.
- Do not call a condition-variable queue lock-free or persistent.
Connected lessons
- Stacks and Queues
- Data Structures
- Bounded thread queues: separate FIFO removal from task completion
- Ring buffers: make capacity and overwrite rules explicit
- Queues: preserve arrival order without front shifts
- Projects
- Quizzes
Apply this contract in the reversible depot release project, then check the operations quiz.
CAS queue protocol: link, help, and reclaim safely examines the next boundary.
