Skip to main content

FastPluggy Task Runner

Task Runner Release Pipeline Status Coverage

A powerful and extensible task execution framework for Python, built on top of FastPluggy.
Register, run, monitor, and schedule background tasks over a pluggable broker (memory, local, postgres, rabbitmq), with per-task logs, progress, locks, and an admin UI.


✨ Features

  • 🔧 Task registration with metadata (description, tags, schedule, topic, allow_concurrent)
  • 🧠 Dynamic form generation from the function signature
  • 📅 CRON / interval scheduler (beat), with standby takeover and health gauges
  • 🔁 Manual retry from the UI, linked to the original task via parent_task_id
  • 🔒 Non-concurrent execution and per-entity locks (lock_key_resolver), with a lock-conflict policy
  • 🧵 Per-topic concurrency limits, drain, purge, and dead-letter queues
  • 📝 Captured task logs and progress (TaskWorker.set_task_progression), stored with the report
  • 📊 Admin UI for tasks, running tasks, schedules, locks, broker debug, and duration analytics
  • 💾 Persistent task context and report (fp_task_contexts, fp_task_reports)
  • 📈 Prometheus gauges via the FastPluggy metrics capability (see docs/observability.md)

Not implemented today: automatic per-task retries (max_retries / retry_delay are stored but not applied, #42), notifications (#28), and live log streaming over WebSocket (#29).


🛠️ How It Works

from fastpluggy_plugin.tasks_worker import TaskWorker

@TaskWorker.register(
    description="Sync data every 5 mins",
    schedule="*/5 * * * *",   # cron string, or an int interval in seconds
    allow_concurrent=False,
    topic="sync",
)
def sync_data_task():
    print("Sync running...")

# Submit from code; returns the task_id
task_id = TaskWorker.submit(sync_data_task)
result = TaskWorker.wait_for_task(task_id, timeout=30)

Tasks with a schedule are registered as scheduled tasks automatically when task discovery runs (enable_auto_task_discovery, on by default).

Further reading:

  • docs/README.md: index of all technical docs.
  • Broker matrix: which broker gives which guarantees.
  • Orphan recovery: how a dead worker's in-flight messages are recovered, and why an unreaped one silently wedges a topic's concurrency.
  • Jinja template globals: task_submit_url and the JS globals for triggering tasks from a page.

📋 Roadmap

✅ Completed / In Progress

  • Task registration with metadata (description, tags, schedule, topic, allow_concurrent)
  • Dynamic task form rendering via metadata
  • Internal event bus (TaskEventBus) with DB, telemetry and broker-broadcast subscribers
  • Context/report tracking in DB
  • Task trees via parent_task_id (retry linking + pipeline chaining / fan-out / completion joins — see docs/orchestration.md)
  • CRON-based scheduler loop
  • Web UI for:
    • Task logs
    • Task reports
    • Scheduled tasks
    • Locks
    • Running task status
  • Broker-level task locks (per task or per entity via lock_key_resolver), with a force-release UI
  • Cancel button (only stops a task this worker has claimed but not yet started; see #31)

📌 Upcoming Features

🔁 Task Queue Enhancements

  • Automatic per-task retries honouring max_retries / retry_delay (#42)
  • Notifications on task events (#28)
  • Live log streaming over WebSocket (#29)
  • Priority & rate-limit execution
  • Per-user concurrency limits
  • Task dependencies / DAG runner

🧠 Task Registry & Detection

  • Auto-discovery of task definitions from modules
  • Celery-style shared task detection

💾 Persistence & Rehydration

  • Save function reference + args for replay/retry
  • Task dependency tree and retry visualization

🌐 Remote Workers

  • Register and manage remote workers
  • Assign tasks based on tags/strategies
  • Remote heartbeat & health monitoring

📈 Observability

  • Prometheus gauges for broker, scheduler and task telemetry
  • Per-task resource metrics via psutil (CPU, memory, threads)
  • UI views for thread/process diagnostics

Standalone Worker

Run task workers as a standalone long-running process, independent of the FastAPI dev server:

# Worker only (consumes and executes tasks)
fastpluggy tasks-worker start

# Worker + scheduler (beat) — simple setups, dev
fastpluggy tasks-worker start --beat

# Standalone scheduler (recommended for production)
fastpluggy tasks-worker beat

# Use RabbitMQ in production
fastpluggy tasks-worker start --broker-type rabbitmq --broker-dsn amqp://user:pass@rabbit:5672/

# Consume only specific topics with 4 threads per worker
fastpluggy tasks-worker start --topics email,reports --max-workers 4

# Verbose logging for debugging
fastpluggy tasks-worker start --log-level DEBUG

# Multiple workers with PostgreSQL broker
fastpluggy tasks-worker start -n 3 --broker-type postgres --broker-dsn postgresql+psycopg2://localhost/tasks

# Inspect: registered workers, and the effective broker configuration
fastpluggy tasks-worker workers
fastpluggy tasks-worker info

The process blocks until interrupted with Ctrl+C or SIGTERM, then performs a graceful shutdown.

RabbitMQ vhost auto-creation

When using the RabbitMQ broker, the worker automatically creates the vhost specified in the DSN if it does not already exist. This uses the RabbitMQ Management HTTP API (port 15672) and grants full permissions to the connecting user. If the management API is unreachable (not exposed, firewalled, or the user lacks admin rights), the check is silently skipped — the worker will connect normally if the vhost already exists, or fail with a clear error if it doesn't.

Production deployment

For production, run the scheduler (beat) and workers as separate processes:

# One beat process — reads scheduled tasks from DB, submits when due
fastpluggy tasks-worker beat --broker-type rabbitmq --broker-dsn amqp://...

# N worker processes — consume and execute tasks
fastpluggy tasks-worker start -n 4 --broker-type rabbitmq --broker-dsn amqp://...

Beat resilience (standby takeover, ≥0.3.288)

Only one beat runs at a time (a dedup guard skips beat startup when a live beat registration exists). Since 0.3.288 that guard is no longer one-shot: a worker that skipped beat keeps standing by — it re-runs the guard every beat_standby_poll_seconds (default 15s) and takes over the moment no live beat remains (claim, then smallest-live-beat-worker_id-wins verify, so concurrent standbys can't double-start). This fixes the container-recreate race (#14) where the successor booted seconds after the dying beat's last heartbeat, saw it as "live", skipped beat permanently — and every scheduled task silently stopped until the next restart. Opt out with beat_standby_enabled=false.

Monitoring the scheduler

Two gauges ship via the FastPluggy metrics capability (aggregator route):

  • fastpluggy_broker_beats — live beat-role workers. 0 with enabled schedules = scheduler dead; alert on it.
  • fastpluggy_schedule_overdue_seconds — worst now − expected_next_run across enabled schedules (0 = on time). Catches both a dead beat and a single wedged schedule, cron or interval.

Options

Option Commands Description
-n, --workers start Number of workers to start (default: $WORKER_NUMBER or 1)
--beat start Also start the scheduler alongside workers
--broker-type start, beat, workers Broker backend: local, memory, rabbitmq, postgres (overrides $BROKER_TYPE)
--broker-dsn start, beat, workers Broker connection string (overrides $BROKER_DSN)
--topics start Comma-separated list of topics to consume (default: all)
--max-workers start Thread pool size per worker (default: 8)
--log-level start, beat Logging level: DEBUG, INFO, WARNING, ERROR (default: INFO)

Topic Routing

Topics determine which queue tasks are published to and consumed from. Resolution order:

  1. FORCE_TASK_TOPIC — if set, overrides everything. Both publish and consume use this value.
  2. Explicit topic= argument — passed to TaskWorker.submit(my_task, topic="email").
  3. Function metadata — set via @TaskWorker.register(topic="reports").
  4. DEFAULT_TOPIC — fallback (default: "default").
Setting Env var Description
default_topic DEFAULT_TOPIC Fallback topic for publish and consume (default: "default")
force_task_topic FORCE_TASK_TOPIC Hard override — locks both publish and consume to this value

Examples:

# Pin a worker to a specific queue (useful for dedicated workers or debugging)
FORCE_TASK_TOPIC=gpu-worker fastpluggy tasks-worker start

# Two workers on the same machine, each consuming a different queue
FORCE_TASK_TOPIC=worker-a fastpluggy tasks-worker start &
FORCE_TASK_TOPIC=worker-b fastpluggy tasks-worker start &

# Normal mode — worker consumes all topics, tasks route via metadata or default
fastpluggy tasks-worker start

🧪 Testing

This plugin includes comprehensive test coverage with pytest.

Running Tests Locally

# Install development dependencies (CI installs ".[tests,postgres]")
pip install -e ".[dev,postgres,rabbitmq]"

# Run all tests (Postgres/RabbitMQ broker tests start testcontainers, so Docker is needed)
pytest tests/

# Run tests with coverage report
pytest tests/ --cov=src --cov-report=term-missing

# Run one file
pytest tests/brokers/test_local.py -v

Every test has a 120 s deadline (pytest.ini, #25); see CLAUDE.md before changing it.

CI/CD Integration

run_tests runs on merge requests and on main, with a 55% coverage floor. A Playwright e2e job (e2e/) drives the admin UI against the memory broker and refreshes docs/screenshots/e2e/. Merging a pyproject.toml version bump to main tags and publishes the release.


📦 Tech Stack

  • FastAPI + FastPluggy
  • SQLAlchemy + SQLite/PostgreSQL
  • Jinja2 + Bootstrap (Tabler)
  • Brokers: in-process memory, multiprocessing BaseManager (local), PostgreSQL, RabbitMQ (pika)

🧠 Philosophy

This runner is built to be:

  • Introspective: auto-generate UIs from functions
  • Composable: integrate with your FastPluggy app
  • Scalable: support single-machine and multi-worker environments
  • Extensible: event-bus subscribers, pluggable brokers, message converters, CRON

📎 License

MIT – Use freely and contribute 💙


🚀 Contributions Welcome!

Open issues, send PRs, share ideas —
Let’s build the most pluggable Python task runner together.

Database support

The models use the generic SQLAlchemy JSON type, so the plugin runs on both SQLite and PostgreSQL. CI runs without DATABASE_URL, i.e. on FastPluggy's SQLite fallback, apart from the broker tests that start a PostgreSQL testcontainer. The postgres broker needs PostgreSQL.

Metadata

Release files for fastpluggy-tasks-worker 0.3.311

For a detailed explanation of source distributions (sdists) and built distributions (wheels), please see the package formats documentation.

Source distribution (sdist)

Source distribution for fastpluggy-tasks-worker 0.3.311
File Size Uploaded
fastpluggy_tasks_worker-0.3.311.tar.gz 263.0 kB Details

Built distribution (wheel)

Table of built distributions (wheels) for fastpluggy-tasks-worker 0.3.311
File Interpreter ABI Platform
fastpluggy_tasks_worker-0.3.311-py3-none-any.whl Python 3 none any Details

Total release size: 539.5 kB

Release files / fastpluggy_tasks_worker-0.3.311.tar.gz

Download URL fastpluggy_tasks_worker-0.3.311.tar.gz
Size 263.0 kB
Tags Source
SHA-256 checksum
How to use checksums
41e4957481f63fb22548dae522799f5a02db7fc84b089ad10ef35057ea3df09f
BLAKE2b-256 checksum
How to use checksums
f89a0216074007aec69a68c027fca1bdebb57bcb6ce4fd755dbb32c343f80971
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.13.13

Release files / fastpluggy_tasks_worker-0.3.311-py3-none-any.whl

Download URL fastpluggy_tasks_worker-0.3.311-py3-none-any.whl
Size 276.5 kB
Tags Python 3
SHA-256 checksum
How to use checksums
d4fe0267b88410b0a25e8fcae12275458d90b051b14a5975639d5c99be526b85
BLAKE2b-256 checksum
How to use checksums
e05af3770732156f17abe4001583386343bbf12102362e7cf396a651ee0b8653
Upload date
Uploaded using Trusted Publishing?
What is trusted publishing?
No
Uploaded via twine/6.2.0 CPython/3.13.13

Release history Release notifications | RSS feed

This release

0.3.311 This release

2 release files

Anthropic, PBC Visionary sponsor Bloomberg Visionary sponsor Hudson River Trading Visionary sponsor Meta Visionary sponsor NVIDIA Visionary sponsor Microsoft Sustainability sponsor Depot Continuous Integration AWS Cloud computing and Security Sponsor Datadog Monitoring Fastly CDN Google Download Analytics Sentry Error logging StatusPage Status page