Bound accepted work and make executor shutdown predictableLESSON 4.06 · 6 OF 13 IN CHAPTER
PART B / APIs and background work
Step 92 of 252
LESSON 4.06 · 6 OF 13 IN CHAPTERHands-on

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
the implementation · 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
that deliberately broken demo · 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.

Diagram: Derive the synchronization before adding features

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.

Diagram: Derive the synchronization before adding features

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.

Diagram: Follow-ups that change the schedule

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.