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.
Easily register, run, monitor, and schedule background tasks with full support for retries, logging, live WebSocket updates, and notifications.


✨ Features

  • 🔧 Task registration with metadata, retries, scheduling, and custom parameters
  • 🧠 Dynamic form generation from metadata
  • 📡 Live logs and WebSocket updates
  • 📅 CRON-based scheduler with optional notification rules
  • 🔁 Retry logic with auto-link to parent task
  • 🔒 Non-concurrent task execution with lock tracking
  • 🧩 Extensible subscribers system (Console, Slack, Webhook...)
  • 📊 Admin UI to manage tasks, schedules, locks, and reports
  • 💾 Persistent task context and rehydration
  • 📈 Task metrics from process/thread info

🛠️ How It Works

@TaskWorker.register(
    description="Sync data every 5 mins",
    schedule="*/5 * * * *",
    max_retries=3,
    allow_concurrent=False
)
def sync_data_task():
    print("Sync running...")

For detailed instructions on creating tasks and triggering them from JavaScript, see the Task Creation and JS Triggering Guide.

For information about Jinja template global variables available for task triggering, see the Jinja Template Globals documentation.

For how a dead worker's in-flight messages are recovered — and why an unreaped one silently wedges a topic's concurrency — see Orphan recovery.


📋 Roadmap

✅ Completed / In Progress

  • Task registration with metadata (description, tags, max_retries, schedule, allow_concurrent)
  • Dynamic task form rendering via metadata
  • Notification/subscribers system with:
    • Console / webhook / Slack (optional)
    • Selectable events: task_started, task_failed, logs, etc.
  • 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
  • Lock manager (TaskLockManager) with DB tracking
  • Cancel button for live-running tasks

📌 Upcoming Features

🔁 Task Queue Enhancements

  • 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

  • Task 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://localhost/tasks

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 Description
-n, --workers Number of workers to start (default: $WORKER_NUMBER or 1)
--beat Also start the scheduler alongside workers (for start command)
--broker-type Broker backend: local, memory, rabbitmq, postgres (overrides $BROKER_TYPE)
--broker-dsn Broker connection string (overrides $BROKER_DSN)
--topics Comma-separated list of topics to consume (default: all)
--max-workers Thread pool size per worker (default: 8)
--log-level 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
pip install -e ".[dev]"

# Run all tests
pytest tests/

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

# Run specific test file
pytest tests/test_runner_topics.py -v

# Run tests with specific markers
pytest tests/ -m unit  # Only unit tests
pytest tests/ -m "not slow"  # Skip slow tests

CI/CD Integration

Tests are automatically run in the GitLab CI/CD pipeline on:

  • Merge requests
  • Main branch commits

Coverage reports are generated and stored as artifacts for 30 days.


📦 Tech Stack

  • FastAPI + FastPluggy
  • SQLAlchemy + SQLite/PostgreSQL
  • WTForms + Jinja2 + Bootstrap (Tabler)
  • WebSockets for real-time feedback
  • Plugin-ready & modular architecture

🧠 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: notifiers, hooks, CRON, logs

📎 License

MIT – Use freely and contribute 💙


🚀 Contributions Welcome!

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

Warning:

Does not work with SQLite due to JSONB field requirements.

Metadata

Release files for fastpluggy-tasks-worker 0.3.306

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.306
File Size Uploaded
fastpluggy_tasks_worker-0.3.306.tar.gz 261.5 kB Details

Built distribution (wheel)

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

Total release size: 541.7 kB

Release files / fastpluggy_tasks_worker-0.3.306.tar.gz

Download URL fastpluggy_tasks_worker-0.3.306.tar.gz
Size 261.5 kB
Tags Source
SHA-256 checksum
How to use checksums
2230a0a4f5db4af2c4cebfc786e6ab40791e5435f55bc1f01a6ce67eb619149c
BLAKE2b-256 checksum
How to use checksums
f92dcac997ed93f79638e118ef1f1f8b3178c38833902130f2bdb96ba3e82582
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.306-py3-none-any.whl

Download URL fastpluggy_tasks_worker-0.3.306-py3-none-any.whl
Size 280.2 kB
Tags Python 3
SHA-256 checksum
How to use checksums
3f46fb629e1338256a4269dbd53083bd5e3bc53c17c78f06db650263778d0848
BLAKE2b-256 checksum
How to use checksums
a6f9e70cfc213be3f833a54b4916bd0af94cc5abffc711da6c73934d53712755
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.306 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