Skip to content
AITroveRead. Build. Understand.
Make this comfortable

Bounded queues: close, wake waiters, and drain accepted work

Last updated: 3 Oct 20269 min read
tutorial
IntermediateBy AITrove Editorial

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.

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

python
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

Output
False
J-47 J-52 None
closed

Time, 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

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.

data structures
range-query-structures
Storage details