~/systems/work-queue

Fault-tolerant work queue

Hand tasks to workers with leases that run out, retry what failed, and give up on what keeps failing. A state machine plus a heap of deadlines.

what

Give every task an explicit state (queued, leased, done, dead). A claim issues a lease with a deadline kept in a min-heap; every call first expires overdue leases. Heartbeats issue a new lease id, so old heap entries go stale and are skipped.

use when

Job queues, visibility timeouts, worker heartbeats, retries with a limit, dead-letter queues, and fair scheduling between customers or tasks.

time

O(log n) per call, amortized

space

O(tasks + leases issued)

You’ll recognise it when

  • Workers claim tasks and must finish or renew them before a deadline, or the task goes back to the queue.
  • The words lease, heartbeat, visibility timeout, retry, attempts or dead letter appear.
  • Each task moves through states, and some calls are only legal in some states (you can’t complete a task you no longer hold).
  • Or the queue must be fair: customers or tasks take turns instead of first come, first served.

It’s often confused with a plain FIFO queue. The queue part is easy; the work is in what happens when a worker goes quiet, and in making sure a late message from it can’t undo someone else’s work.

The idea

A library lends a study room for an hour. If you want longer, you ask at the desk and get a new slip with a new end time; the old slip is no good any more. If your hour runs out and you haven’t renewed, the room goes to the next person in line, even if you come back two minutes later waving the old slip.

That’s a lease. Claiming a task gives the worker a lease_id and a deadline. Every task has an explicit state, and every call starts by expiring leases whose deadline has passed, so the rest of the call sees the truth. Deadlines live in a min-heap. A renewal (a heartbeat) issues a new lease id and pushes a new heap entry instead of finding and editing the old one. The old entry becomes stale: when it reaches the top of the heap, its id is no longer live, so it’s skipped. That’s lazy deletion, and it keeps every operation O(log n).

How it works

WorkQueue(lease_time, max_attempts) keeps ready (a queue of tasks waiting), leases (live lease id to task), and deadlines (a heap of (deadline, lease_id)).

  1. Expire first. Every method begins with _expire(now): pop heap entries with deadline <= now. If the lease_id is still in leases, the worker ran out of time: release the task. If not, the entry is stale; drop it.
  2. Claim. Take the front of ready, mark it leased, count an attempt, and grant a lease: a fresh lease_id, stored in leases and pushed onto deadlines with now + lease_time.
  3. Heartbeat. If the lease_id is live, remove it and grant a new one. The task never leaves the leased state, and the old id now fails every check.
  4. Complete or fail. Only a live lease_id counts. Complete marks the task done. Fail releases it like an expiry.
  5. Release. A released task goes to the back of ready, unless it has already used max_attempts attempts; then it’s dead and never runs again.

With lease_time = 10 and max_attempts = 2:

time call result
0 add a, add b both queued
1, 2 claim, claim a with lease 1 (due 11), b with lease 2 (due 12)
8 heartbeat lease 1 lease 3, due 18; lease 1 is now stale
12 status of b queued: lease 2 ran out (lease 1’s entry was skipped)
12 claim b with lease 4, its 2nd attempt
13 complete with lease 2 false: that lease is over
14 complete with lease 3 true: a is done
21 fail with lease 4 true, and b is dead: 2 attempts used

Why it’s correct: a task is leased exactly when one of its lease ids is in leases, and only calls that hold that id can complete, fail or renew it. Expiry runs before anything else in every call, so no call can act on a lease that has already run out, and a stale heap entry can’t end a lease that was renewed.

import heapq
from collections import deque
class WorkQueue:
"""Tasks are QUEUED -> LEASED -> DONE. A lease that runs out, or a failure,
sends the task back to QUEUED, or to DEAD once its attempts are used up."""
def __init__(self, lease_time, max_attempts):
self.lease_time = lease_time
self.max_attempts = max_attempts
self.state = {} # task -> "queued" | "leased" | "done" | "dead"
self.attempts = {} # task -> how many times it was claimed
self.ready = deque() # queued tasks, oldest first
self.leases = {} # live lease_id -> task
self.deadlines = [] # heap of (deadline, lease_id); may hold stale entries
self.next_id = 1
def _expire(self, now):
# Runs first in every call: leases that ran out end before anything else.
while self.deadlines and self.deadlines[0][0] <= now:
deadline, lease_id = heapq.heappop(self.deadlines)
if lease_id in self.leases: # skip stale entries (lazy deletion)
self._release(lease_id)
def _release(self, lease_id):
task = self.leases.pop(lease_id)
if self.attempts[task] >= self.max_attempts:
self.state[task] = "dead"
else:
self.state[task] = "queued"
self.ready.append(task) # back of the line
def _grant(self, now, task):
lease_id, self.next_id = self.next_id, self.next_id + 1
self.leases[lease_id] = task
heapq.heappush(self.deadlines, (now + self.lease_time, lease_id))
return lease_id
def add(self, now, task):
self._expire(now)
if task in self.state:
return False
self.state[task], self.attempts[task] = "queued", 0
self.ready.append(task)
return True
def claim(self, now):
self._expire(now)
if not self.ready:
return None
task = self.ready.popleft()
self.state[task] = "leased"
self.attempts[task] += 1
return task, self._grant(now, task)
def heartbeat(self, now, lease_id):
self._expire(now)
task = self.leases.pop(lease_id, None)
if task is None:
return None # unknown, expired or finished
return self._grant(now, task) # a new id: the old heap entry goes stale
def complete(self, now, lease_id):
self._expire(now)
task = self.leases.pop(lease_id, None)
if task is None:
return False
self.state[task] = "done"
return True
def fail(self, now, lease_id):
self._expire(now)
if lease_id not in self.leases:
return False
self._release(lease_id)
return True
def status(self, now, task):
self._expire(now)
return self.state.get(task, "unknown")
#include <deque>
#include <functional>
#include <optional>
#include <queue>
#include <string>
#include <unordered_map>
#include <utility>
#include <vector>
using namespace std;
// Tasks are queued -> leased -> done. A lease that runs out, or a failure,
// sends the task back to queued, or to dead once its attempts are used up.
class WorkQueue {
long long lease_time;
int max_attempts;
unordered_map<string, string> state; // task -> "queued" | "leased" | "done" | "dead"
unordered_map<string, int> attempts; // task -> how many times it was claimed
deque<string> ready; // queued tasks, oldest first
unordered_map<long long, string> leases; // live lease_id -> task
using Entry = pair<long long, long long>; // (deadline, lease_id); may be stale
priority_queue<Entry, vector<Entry>, greater<>> deadlines;
long long next_id = 1;
// Runs first in every call: leases that ran out end before anything else.
void expire(long long now) {
while (!deadlines.empty() && deadlines.top().first <= now) {
long long lease_id = deadlines.top().second;
deadlines.pop();
if (leases.count(lease_id)) release(lease_id); // skip stale entries (lazy deletion)
}
}
void release(long long lease_id) {
string task = leases[lease_id];
leases.erase(lease_id);
if (attempts[task] >= max_attempts) {
state[task] = "dead";
} else {
state[task] = "queued";
ready.push_back(task); // back of the line
}
}
long long grant(long long now, const string& task) {
long long lease_id = next_id++;
leases[lease_id] = task;
deadlines.push({now + lease_time, lease_id});
return lease_id;
}
public:
WorkQueue(long long lease_time, int max_attempts) : lease_time(lease_time), max_attempts(max_attempts) {}
bool add(long long now, const string& task) {
expire(now);
if (state.count(task)) return false;
state[task] = "queued";
attempts[task] = 0;
ready.push_back(task);
return true;
}
optional<pair<string, long long>> claim(long long now) {
expire(now);
if (ready.empty()) return nullopt;
string task = ready.front();
ready.pop_front();
state[task] = "leased";
attempts[task]++;
return make_pair(task, grant(now, task));
}
optional<long long> heartbeat(long long now, long long lease_id) {
expire(now);
auto it = leases.find(lease_id);
if (it == leases.end()) return nullopt; // unknown, expired or finished
string task = it->second;
leases.erase(it);
return grant(now, task); // a new id: the old heap entry goes stale
}
bool complete(long long now, long long lease_id) {
expire(now);
auto it = leases.find(lease_id);
if (it == leases.end()) return false;
state[it->second] = "done";
leases.erase(it);
return true;
}
bool fail(long long now, long long lease_id) {
expire(now);
if (!leases.count(lease_id)) return false;
release(lease_id);
return true;
}
string status(long long now, const string& task) {
expire(now);
auto it = state.find(task);
return it == state.end() ? "unknown" : it->second;
}
};
import java.util.*;
// Tasks are queued -> leased -> done. A lease that runs out, or a failure,
// sends the task back to queued, or to dead once its attempts are used up.
class WorkQueue {
private final long leaseTime;
private final int maxAttempts;
private final Map<String, String> state = new HashMap<>(); // task -> "queued" | "leased" | "done" | "dead"
private final Map<String, Integer> attempts = new HashMap<>(); // task -> how many times it was claimed
private final ArrayDeque<String> ready = new ArrayDeque<>(); // queued tasks, oldest first
private final Map<Long, String> leases = new HashMap<>(); // live lease_id -> task
// (deadline, lease_id); may hold stale entries
private final PriorityQueue<long[]> deadlines = new PriorityQueue<>((a, b) ->
a[0] != b[0] ? Long.compare(a[0], b[0]) : Long.compare(a[1], b[1]));
private long nextId = 1;
WorkQueue(long leaseTime, int maxAttempts) {
this.leaseTime = leaseTime;
this.maxAttempts = maxAttempts;
}
// Runs first in every call: leases that ran out end before anything else.
private void expire(long now) {
while (!deadlines.isEmpty() && deadlines.peek()[0] <= now) {
long leaseId = deadlines.poll()[1];
if (leases.containsKey(leaseId)) release(leaseId); // skip stale entries (lazy deletion)
}
}
private void release(long leaseId) {
String task = leases.remove(leaseId);
if (attempts.get(task) >= maxAttempts) {
state.put(task, "dead");
} else {
state.put(task, "queued");
ready.addLast(task); // back of the line
}
}
private long grant(long now, String task) {
long leaseId = nextId++;
leases.put(leaseId, task);
deadlines.add(new long[]{now + leaseTime, leaseId});
return leaseId;
}
boolean add(long now, String task) {
expire(now);
if (state.containsKey(task)) return false;
state.put(task, "queued");
attempts.put(task, 0);
ready.addLast(task);
return true;
}
// Returns {task, lease id} or null.
String[] claim(long now) {
expire(now);
if (ready.isEmpty()) return null;
String task = ready.pollFirst();
state.put(task, "leased");
attempts.merge(task, 1, Integer::sum);
return new String[]{task, String.valueOf(grant(now, task))};
}
Long heartbeat(long now, long leaseId) {
expire(now);
String task = leases.remove(leaseId);
if (task == null) return null; // unknown, expired or finished
return grant(now, task); // a new id: the old heap entry goes stale
}
boolean complete(long now, long leaseId) {
expire(now);
String task = leases.remove(leaseId);
if (task == null) return false;
state.put(task, "done");
return true;
}
boolean fail(long now, long leaseId) {
expire(now);
if (!leases.containsKey(leaseId)) return false;
release(leaseId);
return true;
}
String status(long now, String task) {
expire(now);
return state.getOrDefault(task, "unknown");
}
}

Fair turns between customers

When one customer submits 500 jobs and another submits 1, first come, first served makes the second wait for all 500. A round-robin queue keeps one queue per customer and a rotation of customers who have work: serve the front customer’s oldest job, then move that customer to the back if they still have jobs. A CPU scheduler does the same with time slices: run the front task for min(quantum, remaining) and requeue it if it isn’t finished.

Why it’s O(log n)

Each claim and heartbeat pushes one heap entry: O(log n). Each entry is popped at most once, either to expire a live lease or to drop a stale one, so the expiry loop costs O(log n) per entry over the queue’s whole life, even though one call may pop several. Everything else is dict and deque work, O(1).

call cost
add, complete, fail, status O(1), plus any expiries due
claim, heartbeat O(log n) for the heap push
expiring one entry O(log n), at most once per entry

Stale entries cost memory until they surface: one per heartbeat. If workers heartbeat very often with long leases, the heap can hold many dead entries. Rebuilding it from the live leases once stale entries outnumber live ones keeps it bounded.

Common mistakes

Not expiring before acting

If complete checks the lease before expiring overdue ones, a worker whose lease ran out a minute ago can still complete a task that was already handed to someone else.

if lease_id in self.leases: ... # ✗ the lease may be long over
self._expire(now) # ✓ first, in every method

Reusing the lease id on renewal

If a heartbeat keeps the same id and pushes a later deadline, the old heap entry still carries that id and expires the renewed lease early.

heapq.heappush(self.deadlines, (now + T, lease_id)) # ✗ the old entry still matches
new_id = self._grant(now, task) # ✓ new id; the old entry goes stale

Searching the heap to remove an entry

Finding a lease in the heap to delete or update it is O(n), and editing a heap entry in place breaks the heap order. Leave it and skip it when it surfaces.

self.deadlines.remove(entry); heapq.heapify(self.deadlines) # ✗ O(n) per call
self.leases.pop(lease_id) # ✓ the heap entry is now stale, skipped later

Counting attempts in the wrong place

An attempt is a claim. Counting failures instead lets a task whose leases keep expiring retry forever, because expiries never call fail.

if self.failures[task] >= self.max_attempts: ... # ✗ expiries never count
self.attempts[task] += 1 # ✓ in claim, checked on every release

Variations

  • Priorities. Claim the highest priority first, oldest first among equals: a heap of (-priority, seq, task), where seq keeps the original position when a task is requeued.
  • Backoff. A failed task becomes ready again only after a delay that doubles each attempt. Keep a second heap of (ready_at, task) and move due tasks into ready during the expiry step.
  • Dependencies. A task is claimable only when everything it depends on is done; failure of a dependency fails its dependents. Count unfinished dependencies per task, as in topological sort.
  • Redrive. Dead tasks go to a dead-letter list for a human to inspect. “Redrive” puts one back in the queue with its attempts reset.
  • Idempotency. Leases give at-least-once delivery: a slow worker may finish a task after it was re-leased. Make the work safe to repeat, or record completed task ids and ignore duplicates.

Climb the ladder

Our work-queue problems in ladder order.

  1. Expire leases with a min-heap: the heap and lazy deletion on their own.
  2. Round-robin print queue across customers: fair turns between customers.
  3. Fault-tolerant work queue: reserve, complete and fail, then timeouts, retries and dead letters.
  4. Round-robin task scheduler: time slices, pause and resume, then weights.

Check yourself

5 quick questions. Pick an answer to see why it's right or wrong.

  1. 1

    Leases are a heap of (deadline, lease_id) with lazy deletion. What does this print?

    import heapq
    heap, live = [], set()
    def grant(lease_id, deadline):
    live.add(lease_id)
    heapq.heappush(heap, (deadline, lease_id))
    grant(1, 10)
    grant(2, 12)
    live.discard(1) # heartbeat: lease 1 is replaced...
    grant(3, 20) # ...by lease 3
    expired = []
    while heap and heap[0][0] <= 15:
    _, lease_id = heapq.heappop(heap)
    if lease_id in live:
    live.discard(lease_id)
    expired.append(lease_id)
    print(expired, len(heap))
  2. 2

    A heartbeat keeps the same lease id and pushes a second heap entry (now + T, id). The first entry, with deadline 10, is still in the heap. What happens at time 10?

  3. 3

    With max_attempts = 3, a task’s first two leases run out without a heartbeat. A worker claims it a third time and calls fail. What state is the task in now?

  4. 4

    Worker A’s lease on task t ran out, and worker B now holds a new lease on t. Then A’s complete(old_lease_id) finally arrives. What should the queue do?

  5. 5

    When a lease is renewed, why skip its old heap entry later instead of removing it from the heap right away?

Practice problems

Solve these right here, in Python, C++ or Java. Tests run as you go.

Further reading

esc