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:
objectRepresent 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:
objectStore 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
PENDINGstate.Parameters
Name
Type
Description
name
str
Application-defined task type.
payload
Dict[str, Any] | None
JSON-serializable dictionary;
Noneis 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
Nonefor immediate eligibility.task_id
str | None
Explicit identifier, or
Noneto generate UUID text.Returns
Type
Description
str
Persisted task identifier.
Raises
Exception
Description
TypeError
payloadcannot 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.
Noneor 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
Nonewhen no eligible task arrives in time.Note
Normal block exit marks the task
DONE. AnExceptionraised 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, andfailedkeys, 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
Nonewhen 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;
Noneincludes 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_atvalues 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.Nonedeletes 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
Truewhen one pending row was updated; otherwiseFalse.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()