ddp_utils.task_queue

import ddp_utils.task_queue

Provide a small SQLite-backed task queue with retry state transitions.

File-backed queues retain tasks across process restarts. A queue instance serializes claims and mutations with an in-process lock and gives each thread its own SQLite connection. Tasks progress from PENDING to RUNNING and then to DONE, PENDING for retry, or FAILED.

Examples

Enqueue and process one task:

from ddp_utils.task_queue import TaskQueue

with TaskQueue("tasks.db") as queue:
    queue.enqueue("send_email", {"to": "user@example.com"})
    with queue.dequeue(timeout=5) as task:
        if task is not None:
            send_email(**task.payload)
class ddp_utils.task_queue.Task(id: str, name: str, payload: Dict[str, Any], priority: int, status: str, attempts: int, max_attempts: int, created_at: float, updated_at: float, error: str | None = None, scheduled_at: float | None = None)

Bases: object

Represent one persisted queue record.

Variables

Name

Type

Description

id

str

Unique task identifier.

name

str

Application-defined task type.

payload

Dict[str, Any]

JSON-decoded business data.

priority

int

Higher values are claimed first.

status

str

Current queue state.

attempts

int

Number of recorded processing failures.

max_attempts

int

Failure count that makes the task terminal.

created_at

float

Creation time as a Unix timestamp.

updated_at

float

Last state-change time as a Unix timestamp.

error

str | None

Most recent failure message, if any.

scheduled_at

float | None

Earliest claim time as a Unix timestamp, if constrained.

Examples

Inspect a dequeued task:

print(task.id, task.name, task.attempts)
class ddp_utils.task_queue.TaskQueue(db_path: str | Path = ':memory:', max_attempts: int = 3, retry_delay: float = 0.0)

Bases: object

Store and process application tasks in SQLite.

Each thread receives a separate connection. Mutating operations are serialized by this instance’s lock; the class does not provide a cross-process claim lock beyond SQLite’s normal statement behavior.

Examples

Persist work and process it through a context manager:

queue = TaskQueue("myapp-tasks.db", max_attempts=3)
queue.enqueue("process", {"file": "data.csv"})
with queue.dequeue() as task:
    if task is not None:
        process(task.payload)

Initialize queue defaults, thread-local storage, and the schema.

Parameters

Name

Type

Description

db_path

Union[str, Path]

SQLite file path or ":memory:" for the creating thread’s in-memory database.

max_attempts

int

Default terminal failure threshold for new tasks.

retry_delay

float

Seconds before a nonterminal failed task is eligible for another claim.

Examples

Create a file-backed queue with delayed retries:

queue = TaskQueue("tasks.db", max_attempts=5, retry_delay=30)
enqueue(name: str, payload: Dict[str, Any] | None = None, priority: int = 0, max_attempts: int | None = None, scheduled_at: float | None = None, task_id: str | None = None) → str

Insert a new task in PENDING state.

Parameters

Name

Type

Description

name

str

Application-defined task type.

payload

Dict[str, Any] | None

JSON-serializable dictionary; None is stored as {}.

priority

int

Claim order weight; higher values run first.

max_attempts

int | None

Per-task terminal failure threshold, or the queue default when omitted.

scheduled_at

float | None

Earliest eligible Unix timestamp, or None for immediate eligibility.

task_id

str | None

Explicit identifier, or None to generate UUID text.

Returns

Type

Description

str

Persisted task identifier.

Raises

Exception

Description

TypeError

payload cannot be JSON serialized.

sqlite3.Error

SQLite rejects or cannot persist the record.

Examples

Prioritize an outgoing email task:

task_id = queue.enqueue(
    "email",
    {"to": "user@example.com"},
    priority=5,
)
dequeue(timeout: float | None = None, names: List[str] | None = None) → Iterator[Task | None]

Claim one eligible task and finalize it when the block exits.

Parameters

Name

Type

Description

timeout

float | None

Maximum seconds to poll an empty queue. None or zero performs only the initial claim.

names

List[str] | None

Optional task-name allowlist. Values remain bound SQL parameters rather than interpolated query text.

Yields

Claimed task, or None when no eligible task arrives in time.

Note

Normal block exit marks the task DONE. An Exception raised by the block is recorded, schedules retry or terminal failure, and is suppressed by this context manager.

Examples

Process only email tasks:

with queue.dequeue(timeout=5, names=["email"]) as task:
    if task is not None:
        send_email(**task.payload)
stats() → Dict[str, int]

Count tasks in each supported state.

Returns

Type

Description

Dict[str, int]

Dictionary containing lowercase pending, running, done, and failed keys, including zeros for absent states.

Examples

Display pending work:

print(queue.stats()["pending"])
list_tasks(status: str | None = None, limit: int = 50) → List[Task]

List tasks in priority and creation order.

Parameters

Name

Type

Description

status

str | None

Optional state filter, normalized to uppercase.

limit

int

Maximum number of rows returned.

Returns

Type

Description

List[Task]

Task snapshots matching the filter.

Examples

Inspect the first ten failed tasks:

failed = queue.list_tasks("FAILED", limit=10)
get_task(task_id: str) → Task | None

Look up one task by its exact identifier.

Parameters

Name

Type

Description

task_id

str

Persisted task identifier.

Returns

Type

Description

Task | None

Task snapshot, or None when no row matches.

Examples

Inspect a task after enqueueing it:

task = queue.get_task(task_id)
retry_failed(names: List[str] | None = None) → int

Return selected failed tasks to pending state.

Parameters

Name

Type

Description

names

List[str] | None

Optional task-name allowlist; None includes all failed tasks.

Returns

Type

Description

int

Number of rows moved to PENDING.

Note

Attempts reset to zero and error text is cleared. Existing scheduled_at values are preserved.

Examples

Retry only failed email tasks:

count = queue.retry_failed(names=["email"])
purge_done(older_than: float | None = None) → int

Delete completed tasks, optionally by completion age.

Parameters

Name

Type

Description

older_than

float | None

Minimum age in seconds based on updated_at. None deletes every completed task.

Returns

Type

Description

int

Number of deleted rows.

Examples

Keep the most recent day of completed tasks:

removed = queue.purge_done(older_than=86_400)
cancel(task_id: str) → bool

Cancel a pending task by moving it to failed state.

Parameters

Name

Type

Description

task_id

str

Identifier of the pending task.

Returns

Type

Description

bool

True when one pending row was updated; otherwise False.

Examples

Cancel work that has not started:

cancelled = queue.cancel(task_id)
close() → None

Close and clear the current thread’s SQLite connection.

Examples

Release the current thread’s database handle explicitly:

queue.close()