Metadata-Version: 2.4
Name: fastpluggy-tasks-worker
Version: 0.3.316
Summary: Task Runner plugin for Fastpluggy
Author: FastPluggy Team
License-Expression: MIT
Requires-Python: >=3.10
Description-Content-Type: text/markdown
Requires-Dist: FastPluggy>=0.4.60
Requires-Dist: fastpluggy-crud-tools>=0.2.12
Requires-Dist: croniter
Requires-Dist: psutil
Requires-Dist: python-dateutil>=2.9.0
Provides-Extra: rabbitmq
Requires-Dist: pika>=1.3.0; extra == "rabbitmq"
Provides-Extra: postgres
Requires-Dist: SQLAlchemy>=2.0; extra == "postgres"
Requires-Dist: psycopg[binary]>=3.1; extra == "postgres"
Provides-Extra: dev
Requires-Dist: pytest>=7.0; extra == "dev"
Requires-Dist: pytest-cov>=4.0; extra == "dev"
Requires-Dist: pytest-asyncio>=0.21; extra == "dev"
Requires-Dist: pytest-timeout>=2.1; extra == "dev"
Requires-Dist: testcontainers[postgres,rabbitmq]>=4.0.0; extra == "dev"
Provides-Extra: tests
Requires-Dist: pytest>=7.0; extra == "tests"
Requires-Dist: pytest-cov>=4.0; extra == "tests"
Requires-Dist: pytest-asyncio>=0.21; extra == "tests"
Requires-Dist: pytest-timeout>=2.1; extra == "tests"
Requires-Dist: testcontainers[postgres,rabbitmq]>=4.0.0; extra == "tests"
Provides-Extra: e2e
Requires-Dist: playwright>=1.40.0; extra == "e2e"
Requires-Dist: fastpluggy-example-plugin; extra == "e2e"

# FastPluggy Task Runner

![Task Runner](https://img.shields.io/badge/FastPluggy-Task%20Runner-blue)
[![Release](https://gitlab.ggcorp.fr/open/fastpluggy/plugins/tasks_worker/-/badges/release.svg)](https://gitlab.ggcorp.fr/open/fastpluggy/plugins/tasks_worker/-/releases)
[![Pipeline Status](https://gitlab.ggcorp.fr/open/fastpluggy/plugins/tasks_worker/badges/main/pipeline.svg?key_text=CI)](https://gitlab.ggcorp.fr/open/fastpluggy/plugins/tasks_worker/-/pipelines?ignore_skipped=true)
[![Coverage](https://gitlab.ggcorp.fr/open/fastpluggy/plugins/tasks_worker/badges/main/coverage.svg)](https://gitlab.ggcorp.fr/open/fastpluggy/plugins/tasks_worker/-/pipelines)

A powerful and extensible **task execution framework** for Python, built on top of [FastPluggy](https://fastpluggy.xyz).  
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](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

```python
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](docs/README.md): index of all technical docs.
- [Broker matrix](docs/broker/broker-matrix.md): which broker gives which guarantees.
- [Orphan recovery](docs/broker/orphan-recovery.md): how a dead worker's in-flight messages are
  recovered, and why an unreaped one silently wedges a topic's concurrency.
- [Jinja template globals](docs/jinja_template_globals.md): `task_submit_url` and the JS globals for
  triggering tasks from a page.

---

## 📋 Roadmap

### ✅ Completed / In Progress

- [x] Task registration with metadata (`description`, `tags`, `schedule`, `topic`, `allow_concurrent`)
- [x] Dynamic task form rendering via metadata
- [x] Internal event bus (`TaskEventBus`) with DB, telemetry and broker-broadcast subscribers
- [x] Context/report tracking in DB
- [x] Task trees via `parent_task_id` (retry linking + pipeline chaining / fan-out / completion joins — see [docs/orchestration.md](docs/orchestration.md))
- [x] CRON-based scheduler loop
- [x] Web UI for:
  - Task logs
  - Task reports
  - Scheduled tasks
  - Locks
  - Running task status
- [x] Broker-level task locks (per task or per entity via `lock_key_resolver`), with a force-release UI
- [x] Cancel button (only stops a task this worker has claimed but not yet started; otherwise 409, 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
- [x] Auto-discovery of task definitions from modules
- [x] Celery-style shared task detection


#### 💾 Persistence & Rehydration
- [x] Save function reference + args for replay/retry
- [x] Task dependency tree and retry visualization

#### 🌐 Remote Workers
- [ ] Register and manage remote workers
- [ ] Assign tasks based on tags/strategies
- [x] Remote heartbeat & health monitoring

#### 📈 Observability
- [x] 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:

```bash
# 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:

```bash
# 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.

The **Scheduled Monitoring** page (sidebar entry, also linked from the Dashboard; only mounted
when `store_task_db` is on) shows every schedule with its last runs:

| Status | Meaning |
|---|---|
| Operational | Last finished run succeeded and the schedule is on time |
| Issues | Last finished run failed, or the beat fired it but no run was ever recorded |
| Late | More than 1 min past its expected next run (the beat polls every `scheduler_frequency` s) |
| Waiting for first run | Never fired yet, or just fired and still queued |
| Disabled | `enabled = false`; listed last, no next run |

Uptime is the share of *finished* runs (last N, "Reports per task") that succeeded; queued and
running runs are not counted. Times are shown in the browser's local time zone.

### 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:**

```bash
# 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
```

---

## 🔐 Authentication

Every route of the plugin (API, admin pages, debug actions, scheduler, monitoring) requires an
authenticated user when the FastPluggy app has an `auth_manager`. `get_router()` mounts all routers
under one parent carrying `require_authentication`, so a new router is covered automatically (#27).
Apps without an `auth_manager` stay open, which is FastPluggy's normal behaviour.

---

## 🧪 Testing

This plugin includes comprehensive test coverage with pytest.

### Running Tests Locally

```bash
# 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.
