Bound accepted work and make executor shutdown predictable
Application and assignment
A reporting API runs database jobs using four available connections. A burst of submissions can arrive while all four are occupied. Starting a new thread for every request only moves the waiting into more threads. You need a bounded set of active workers and an explicit limit on accepted waiting work.
Implement the executor contract below, using the supplied implementation as a later reference. Decide what a blocked submitter, a queued task, and a running task each observe during shutdown. This is an in-process exercise. Accepted work does not survive a process crash unless you add durable storage.
Contract and starting evidence
“A report service accepts bursts faster than its four database connections can work. Today every request starts a thread. Build an executor with four workers and room for eight waiting tasks. A blocked caller must be released when shutdown starts. Which tasks have we promised to finish?”
This constructed backend/infra extension assumes runtime fundamentals and basic maps/queues. Read the infra candidate brief before opening the implementation (download file, source below) or assessor notes.
Read the supplied code · executor.py
"""A fixed worker pool with bounded admission; stdlib only, Python 3.10+."""
from collections import deque
from threading import Condition, Event, Thread, TIMEOUT_MAX, local
import math
import time
class Closed(RuntimeError):
pass
class NestedSubmission(RuntimeError):
pass
class Cancelled(RuntimeError):
pass
class DeadlineExceeded(TimeoutError):
pass
def _time_value(value, name, *, nonnegative=False):
if value is None:
return None
try:
valid = type(value) in (int, float) and math.isfinite(value)
except OverflowError:
valid = False
if not valid or (nonnegative and value < 0):
raise ValueError(f"{name} must be finite" + (" and nonnegative" if nonnegative else ""))
return value
class Context:
def __init__(self, clock, deadline):
self.clock, self.deadline, self.cancelled = clock, deadline, Event()
def check(self):
if self.cancelled.is_set():
raise Cancelled("cooperative cancellation")
if self.deadline is not None and self.clock() >= self.deadline:
raise DeadlineExceeded("total task deadline")
class Task:
def __init__(self, fn, context):
self.fn, self.context = fn, context
self.done = Event()
self.state, self.value, self.error, self.terminal_count = "queued", None, None, 0
def result(self, timeout=None):
timeout = _time_value(timeout, "timeout", nonnegative=True)
end = None if timeout is None else time.monotonic() + timeout
while not self.done.is_set():
remaining = None if end is None else end - time.monotonic()
if remaining is not None and remaining <= 0:
raise TimeoutError("waiting did not cancel the task")
self.done.wait(None if remaining is None else min(remaining, TIMEOUT_MAX))
if self.error is not None:
raise self.error
return self.value
class Executor:
def __init__(self, workers=4, capacity=8, clock=time.monotonic):
if type(workers) is not int or type(capacity) is not int or workers <= 0 or capacity <= 0:
raise ValueError("positive workers and queue capacity required")
self.capacity, self.clock = capacity, clock
self.cv, self.queue, self.worker_local = Condition(), deque(), local()
self.closed, self.active = False, 0
self.peak_active, self.peak_queued, self.accepted, self.completed = 0, 0, 0, 0
self.waiting_producers = 0
self.threads = [Thread(target=self._worker, name=f"executor-{i}", daemon=True) for i in range(workers)]
for thread in self.threads:
thread.start()
def wake(self):
"""Notify predicate rechecks after an injected logical clock advances."""
with self.cv:
self.cv.notify_all()
def submit(self, fn, *, deadline=None):
deadline = _time_value(deadline, "deadline")
if getattr(self.worker_local, "inside", False):
raise NestedSubmission("workers cannot submit into their own pool")
with self.cv:
# Predicate: admission is possible OR shutdown/deadline requires refusal.
while len(self.queue) >= self.capacity and not self.closed:
remaining = None if deadline is None else deadline - self.clock()
if remaining is not None and remaining <= 0:
raise DeadlineExceeded("admission deadline")
self.waiting_producers += 1
self.cv.notify_all()
try:
self.cv.wait(None if remaining is None else min(remaining, TIMEOUT_MAX))
finally:
self.waiting_producers -= 1
if self.closed:
raise Closed("executor closed")
if deadline is not None and self.clock() >= deadline:
raise DeadlineExceeded("admission deadline")
task = Task(fn, Context(self.clock, deadline))
self.queue.append(task)
self.accepted += 1
self.peak_queued = max(self.peak_queued, len(self.queue))
self.cv.notify_all()
return task
def _finish(self, task, state, value=None, error=None):
# Caller owns cv. Only internal state and Event are changed under this lock.
assert not task.done.is_set()
# cancel() and this publication share cv: an accepted cancellation wins.
if task.context.cancelled.is_set():
state, value, error = "cancelled", None, Cancelled("cancelled before publication")
elif state == "succeeded" and task.context.deadline is not None and self.clock() >= task.context.deadline:
state, value, error = "expired", None, DeadlineExceeded("deadline before publication")
task.state, task.value, task.error = state, value, error
task.terminal_count += 1
self.completed += 1
task.done.set()
def cancel(self, task):
with self.cv:
if task.done.is_set():
return False
if task in self.queue:
self.queue.remove(task)
self._finish(task, "cancelled", error=Cancelled("cancelled before start"))
else:
task.context.cancelled.set()
self.cv.notify_all()
return True
def _worker(self):
self.worker_local.inside = True
while True:
with self.cv:
# A notification is not a promise: always recheck this predicate.
while not self.queue and not self.closed:
self.cv.wait()
if not self.queue:
return
task = self.queue.popleft()
self.active += 1
self.peak_active = max(self.peak_active, self.active)
task.state = "running"
self.cv.notify_all() # free queue slot wakes a blocked producer
state, value, error = "succeeded", None, None
try:
task.context.check()
value = task.fn(task.context) # never run arbitrary user code under cv
task.context.check()
except Cancelled as exc:
state, error = "cancelled", exc
except DeadlineExceeded as exc:
state, error = "expired", exc
except BaseException as exc:
state, error = "failed", exc
finally:
with self.cv:
self.active -= 1 # release capacity even after failure
self._finish(task, state, value, error)
self.cv.notify_all()
def shutdown(self, *, cancel_queued=False, wait=True, timeout=None):
timeout = _time_value(timeout, "timeout", nonnegative=True)
if wait and getattr(self.worker_local, "inside", False):
raise NestedSubmission("worker cannot join its own executor")
with self.cv:
self.closed = True
if cancel_queued:
while self.queue:
task = self.queue.popleft()
self._finish(task, "cancelled", error=Cancelled("shutdown before start"))
self.cv.notify_all() # includes blocked producers and idle workers
if wait:
end = None if timeout is None else time.monotonic() + timeout
for thread in self.threads:
while thread.is_alive():
remaining = None if end is None else end - time.monotonic()
if remaining is not None and remaining <= 0:
break
thread.join(None if remaining is None else min(remaining, TIMEOUT_MAX))
if any(thread.is_alive() for thread in self.threads):
raise TimeoutError("running work still owns resources")
def snapshot(self):
with self.cv:
return dict(active=self.active, queued=len(self.queue), accepted=self.accepted,
completed=self.completed, peak_active=self.peak_active, peak_queued=self.peak_queued,
waiting_producers=self.waiting_producers)
| Contract | Observable result |
|---|---|
| Capacity | At most 4 active tasks and 8 accepted waiting tasks; a 13th submit blocks while all slots are occupied |
| Task interface | Callable receives a context with check(), cancellation event and optional absolute monotonic deadline |
| Accepted task | Exactly one terminal outcome: succeeded, failed, cancelled or expired |
| Shutdown | Stop admission; either drain queued work or cancel it; let running work finish cooperatively |
| Deadline | One finite absolute monotonic deadline covers admission through terminal publication; a wait timeout alone does not stop a task |
| Cancellation | cancel(task) == True wins over an unpublished outcome; False means the task was already terminal |
| Nested work | A worker submitting into this same pool fails immediately; worker joining this pool is also rejected |
| Excluded | Forced thread termination, FIFO fairness among blocked producers, process-crash durability and fleet-wide limits |
Worked trace: hold four tasks at a barrier, submit eight more, then start a producer.
The producer waits and accepted == 12. Shutdown with cancel_queued=True releases
that producer with Closed, gives eight queued tasks Cancelled, and allows four
running tasks to finish. A refused submission is not an accepted task and receives no
task handle. A task raising ValueError still frees its worker slot.
Run the supplied checks
From the repository root, Python 3.10+ / standard library:
python -m unittest discover -s curriculum/02-applications/01-backend/labs/bounded-executor -v
The tests use barriers/events and a logical clock to control ordering. Real timeouts are watchdogs preventing a broken implementation from hanging the test runner; there are no sleep-only scheduling assertions. One subprocess intentionally deadlocks and is killed only after all four parents report they are waiting for queued children. Do not run that deliberately broken demo (download file, source below) directly.
Read the supplied code · deadlock_demo.py
"""Deliberately wrong: run only in the subprocess watchdog test."""
from concurrent.futures import ThreadPoolExecutor
from threading import Barrier
pool = ThreadPoolExecutor(max_workers=4)
barrier = Barrier(4)
def parent():
barrier.wait()
child = pool.submit(lambda: "child")
print("parent waiting for queued child", flush=True)
return child.result() # all workers wait; no worker can run any child
parents = [pool.submit(parent) for _ in range(4)]
for future in parents:
future.result()
Derive the synchronization before adding features
Start with a single queue and one condition variable: a condition coordinates waiting for a predicate under a mutex. A notification means “recheck”; it does not transfer ownership or reserve a slot. A task itself is arbitrary user code, so it must execute outside the lock.
Write safety properties before implementation: 0 ≤ active ≤ 4, 0 ≤ queued ≤ 8,
and accepted = queued + active + completed. Each terminal transition occurs once
under the condition lock. The exception path and cancellation path must preserve the
same accounting. Here there is no separate semaphore permit: active slots and queue
slots are the counted permits, released in finally or queued cancellation.
Liveness needs assumptions: workers are scheduled and running callables eventually return or check cancellation. It does not follow from bounded memory alone. The producer loop waits while the queue is full and admission remains open; after waking it rechecks shutdown and deadline before admission. The worker loop waits while the queue is empty and the executor is open; closed plus empty terminates.
Follow-ups that change the schedule
Question 1: Cancel a running task blocked in a library that ignores cancellation.
Can its slot be handed to another task immediately? Expected answer: no. Signal
cancellation, retain active accounting, wait for actual return, then record a terminal
outcome. A shutdown timeout reports that work still owns resources. It must not
pretend to release a socket or thread. Cancellation and terminal publication use
the same lock. A successful cancel() before publication makes the eventual
outcome cancelled, including the gap after the callable's final check. A success
is also checked for deadline expiry under that lock. Neither can undo an external
effect already performed. Deadlines and wait timeouts reject NaN/infinity/bool;
wait timeouts must be nonnegative. Individual OS waits are capped and rechecked.
Question 2: Four running parents each submit a child and wait for it. More queue space is proposed. Expected answer: the children can queue but no worker is free to run them. More queue capacity does not break the wait cycle. This implementation rejects same-pool submission from workers; alternatives require explicit semantics (inline execution/reentrancy, separate dependency pools, or nonblocking composition). Predict the deadlock demo (shown in this lesson) output before reading its watchdog test.
Question 3: A deadline expires while the producer is blocked on a full queue.
Expected answer: wake/recheck the remaining budget, refuse admission, and do not
create a terminal task for a promise never accepted. The test advances an injected
clock and calls wake(); it does not sleep until a wall-clock deadline.
Queue operations are O(1) except cancelling an arbitrary queued task, O(q), and retained memory is O(workers + capacity) plus payloads/future handles held by callers. Output history is not retained by the executor. Local limits multiply across replicas: 20 instances can run 80 tasks and queue 160. Lead follow-up: allocate dependency and tenant budgets, choose overload responses, and state whether waiting producers can starve. This implementation promises bounds and eventual progress under its assumptions, not strict fairness or a fleet-wide concurrency cap.
| Saturated one-process snapshot | Count |
|---|---|
| Active | 4 |
| Accepted in queue | 8 |
| Blocked producer outside accepted work | 1 |
| Accepted total before release | 12 |
snapshot() exposes waiting_producers separately. A service must also bound
incoming callers, apply a total caller deadline and return an overload response.
For 20 replicas, budget the actual 80 active + 160 queued tasks against downstream
capacity; local bounds do not allocate a fleet quota. FIFO dequeue order does not
promise completion order or fairness among blocked callers.