"""RQ worker for AgentForms background tasks.
Connects to Redis and listens on the 'default', 'webhooks', and 'onboarding' queues for jobs:
- Webhook delivery (with retries)
- Email sending
- Onboarding email sequence
- Future: report generation, data exports, etc.
Run via: python3 -m app.worker
"""
import logging
import os
import sys
import time
from dotenv import load_dotenv
# Load .env from data directory
_env_path = os.path.join(os.path.dirname(os.path.abspath(__file__)), "..", "data", ".env")
load_dotenv(_env_path)
# Set up logging
logging.basicConfig(
level=logging.INFO,
format="%(asctime)s [%(name)s] %(levelname)s: %(message)s",
)
logger = logging.getLogger("agentforms.worker")
def get_redis():
"""Get Redis connection from REDIS_URL env var."""
import redis as redis_lib
redis_url = os.environ.get("REDIS_URL", "redis://localhost:6379/0")
r = redis_lib.from_url(redis_url)
return r
def wait_for_redis(timeout=30):
"""Wait for Redis to be available before starting workers."""
r = get_redis()
deadline = time.time() + timeout
while time.time() < deadline:
try:
r.ping()
logger.info("Redis connected")
return r
except Exception:
time.sleep(2)
raise TimeoutError(f"Redis not available after {timeout}s")
def _heartbeat_worker(r):
"""Background thread that writes a periodic heartbeat to Redis.
The relay's /health/worker endpoint checks this key to verify
the worker process is alive. Stale threshold: 3 minutes.
"""
import time as _time
heartbeat_key = "agentforms:worker:heartbeat"
while True:
try:
r.set(heartbeat_key, _time.time())
except Exception as e:
logger.warning("Worker heartbeat failed: %s", e)
time.sleep(30)
def main():
"""Start RQ workers listening on configured queues."""
import threading
import rq
from rq import Queue, Worker
r = wait_for_redis()
# Define queues
queues = [
Queue("default", connection=r),
Queue("webhooks", connection=r),
Queue("onboarding", connection=r),
Queue("campaigns", connection=r),
Queue("chains", connection=r),
Queue("pdf", connection=r),
]
logger.info("Starting RQ workers for queues: %s", [q.name for q in queues])
# Start heartbeat thread
hb = threading.Thread(target=_heartbeat_worker, args=(r,), daemon=True)
hb.start()
logger.info("Worker heartbeat thread started")
worker = Worker(queues, connection=r)
worker.work()
if __name__ == "__main__":
main()