Work Queue
The current background-work API is fluvius.workq, implemented with Taskiq. fluvius.worker is a compatibility shim for applications migrating from the retired ARQ worker package.
from fluvius.workq import WorkqService
class BillingTasks(WorkqService):
@WorkqService.task(name="process-payment")
async def process_payment(context, payment_id: str):
return {"payment_id": payment_id, "status": "processed"}
@WorkqService.cron(name="daily-report", hour=0, minute=0)
async def daily_report(context):
return {"status": "generated"}
Create the service with a durable broker URL:
service = BillingTasks(
url="postgresql://user:password@localhost/app",
queue_name="billing",
)
await service.startup()
The workq extra installs Redis, NATS, and PostgreSQL Taskiq adapters. Set FLUVIUS_WORKQ_URL or FLUVIUS_BUS_URL; the framework intentionally does not fall back to an in-memory broker.
Use DomainWorkqService and DomainWorkqClient to execute domain commands through the queue. WorkqResultBackend, WorkqHandle, and WorkqTaskScope provide result storage, handles, and scoped context injection.