Bounded blocking queue with shutdown
“Several producer threads submit items to worker threads. Memory must stay bounded: producers wait when the queue is full, consumers wait when empty. Shutdown must wake blocked callers and either drain accepted items or cancel queued items. Which state change lets a waiting operation proceed, and what if a wake-up does not mean that state is still true?”
Constructed practice question; backend/infrastructure extension. Prerequisites: explicit mutable state and time boundaries. Concurrency means operations overlap; it does not imply parallel CPU execution. A condition variable lets a thread release a mutex while waiting, then reacquire it before checking state.
| Contract | Required behavior |
|---|---|
| Input | Positive capacity; arbitrary items; put/get(timeout=None) |
| Ordering | FIFO removal; successful puts are accepted once; processing order is outside scope |
| Blocking | Full put/empty get wait; zero timeout attempts immediately; expiry raises TimeoutError |
| Shutdown | First call selects drain or cancel; rejects future puts and wakes all waiters |
| Drain/cancel | Drain permits existing gets; cancel returns removed items; closed-empty get raises QueueClosed |
| Scope | One process; no fairness, task execution, acknowledgement, or cross-replica limit |
Optional refresher · the underlying tool
A condition variable lets a thread sleep while releasing the lock, then reacquire it to inspect state. Waking does not reserve a slot or item: another thread may have taken it.
with condition:
while len(items) >= capacity and not closed:
condition.wait() # releases and reacquires the mutex
if closed:
raise QueueClosed()
For capacity 1 containing A, a put(B) blocks until A is removed or shutdown rejects it. Use a loop, a single deadline across wakes, and explicit drain/cancel behavior.
A design choice worth saying aloud
Separate closed from items: shutdown may reject new puts while still permitting consumers to drain queued work. Recheck capacity and closure in a while after every wake; a signal is not a reservation. Use one monotonic deadline for timed operations, or repeated wakes could extend the advertised timeout.
What the interviewer expects
The interviewer gives you the scenario and the contract above. Explain what a successful call returns, walk through one example below, and name what your state means before choosing a data structure.
Done means: Preserve FIFO and bounded capacity while giving every put, get, timeout, drain, and cancellation race the documented outcome.
Now predict each output before looking at the reference; invalid input should leave any existing state unchanged unless the contract says otherwise.
Test-case scenarios to settle before coding
01 · FIFO
- Input / starting state
- put A, put B, then two gets
- Expected result
- A then B
What it is testing: Acceptance order is preserved.
02 · Full producer
- Input / starting state
- capacity 1 contains A; producer puts B
- Expected result
- producer waits until A is removed, then succeeds once
What it is testing: Use a condition loop, not one wakeup assumption.
03 · Empty consumer
- Input / starting state
- consumer gets from empty queue
- Expected result
- waits until item, close, cancel, or deadline
What it is testing: Every wake path rechecks state.
04 · Timeout zero
- Input / starting state
- full put or empty get with timeout 0
- Expected result
- immediate
TimeoutError
What it is testing: Zero is a nonblocking attempt.
05 · Drain/cancel
- Input / starting state
- close with items present
- Expected result
- drain serves them; cancel returns/removes them
What it is testing: Shutdown policy is explicit and first call wins.
06 · Race safety
- Input / starting state
- competing producers/consumers plus spurious wakeups
- Expected result
- no loss/duplication; deadline budget not reset
What it is testing: Concurrency tests target schedules, not only values.
For each case, show which branch or state change produces that result.
At capacity 1, put A succeeds; put B waits. Getting A frees B's slot. If shutdown
begins while B waits, B instead raises QueueClosed; drain still permits reading A.
Cancel mode returns [A] and leaves no item to read. Invalid capacity/timeout raises
ValueError. Ask which shutdown policy callers expect before choosing one.
Implement independently. Write the exact put/get predicates and identify the instant each successful operation becomes visible to other threads.
Solution, schedules, and follow-ups
Checking capacity outside a lock is the baseline race: two producers both see one free slot and both append. Holding a lock while sleeping blocks the consumer that could free a slot. A condition's atomic release-and-wait operation solves this coordination problem, but a notification itself does not reserve an item or slot.
Use one mutex/condition around deque and closed flag. Put waits while full and open, then rejects closed or appends. Get waits while empty and open, then removes an item or rejects closed-empty. Every mutation notifies waiting callers. Shutdown sets closed and optionally clears queued items while holding that same mutex, then notifies all. No arbitrary user callback executes under the lock.
Safety invariant: 0 <= size <= capacity; each successful removal takes one
previously accepted item, and no append linearizes after closure. Liveness
condition: a waiter can progress when its predicate changes, provided the runtime
eventually schedules it; strict fairness is not promised. Predicate loops tolerate
spurious notifications and another thread consuming the opportunity first.
Compute one monotonic deadline per operation; each repeated wait uses the remaining
budget. A timeout limits waiting while the predicate is false, not OS scheduling,
mutex-acquisition latency, or task execution. The implementation caps individual
waits at the runtime's supported duration. Deque mutation is O(1); notify_all
may wake O(w) waiters, so total notification/scheduling work is not constant. Queue
storage is O(capacity), plus O(w) waiting-thread resources; cancellation copies O(q)
queued items. Blocking duration is unbounded without timeout or progress.
Follow-up 1 — build a worker pool. Remove items under the queue lock, execute tasks outside it, and publish one success/failure outcome even when a task raises. Cancellation can return queued items but cannot forcibly stop arbitrary running Python code. Track accepted, running, completed, failed, and cancelled work separately.
Follow-up 2 — a worker submits a child and waits. Predict deadlock when all workers do this: parents hold every worker slot while their children wait in the queue. Queue correctness cannot create an available worker. Prohibit nested waits, execute children inline under a documented rule, or use a separate scheduling model.
Senior depth is explaining safety, conditional liveness, deadlines, and shutdown races. Lead depth adds fairness, tenant budgets, and cross-process coordination; this queue alone does not establish executor or distributed delivery guarantees.
Reference: solution.py (download file, source below). Tests use barriers and condition-entry events rather than sleeps, explicitly inject a useless wake-up, verify competing threads lose/duplicate no items, and simulate repeated deadline wake-ups with a fake clock. Every thread join is bounded and workers are daemonized as a hang guard.
"""Bounded FIFO with condition predicate loops, deadlines, and explicit shutdown."""
from collections import deque
import math
import threading
import time
class QueueClosed(Exception):
pass
class BoundedBlockingQueue:
def __init__(self, capacity, *, condition=None):
if not isinstance(capacity, int) or capacity <= 0:
raise ValueError("positive integer capacity required")
self.capacity = capacity
self._items = deque()
self._closed = False
self._condition = condition if condition is not None else threading.Condition()
@staticmethod
def _deadline(timeout):
if timeout is None:
return None
if not isinstance(timeout, (int, float)) or not math.isfinite(timeout) or timeout < 0:
raise ValueError("finite nonnegative timeout required")
return time.monotonic() + timeout
def _wait(self, deadline):
if deadline is None:
self._condition.wait()
else:
remaining = deadline - time.monotonic()
if remaining <= 0:
raise TimeoutError("queue operation timed out")
self._condition.wait(min(remaining, threading.TIMEOUT_MAX))
def put(self, item, timeout=None):
deadline = self._deadline(timeout)
with self._condition:
while len(self._items) == self.capacity and not self._closed:
self._wait(deadline)
if self._closed:
raise QueueClosed("queue is closed")
self._items.append(item)
self._condition.notify_all()
def get(self, timeout=None):
deadline = self._deadline(timeout)
with self._condition:
while not self._items and not self._closed:
self._wait(deadline)
if not self._items:
raise QueueClosed("closed queue is drained")
item = self._items.popleft()
self._condition.notify_all()
return item
def shutdown(self, *, cancel_pending=False):
"""First call chooses drain or cancel; return removed items on cancellation."""
with self._condition:
if self._closed:
return []
self._closed = True
cancelled = list(self._items) if cancel_pending else []
if cancel_pending:
self._items.clear()
self._condition.notify_all()
return cancelled
python -m unittest discover -s curriculum/02-applications/01-backend/problems/42-bounded-blocking-queue -p 'test_*.py'