File size: 3,302 Bytes
6993919
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
"""Redis-backed Background Task Queue — retry with exponential backoff, visibility."""

import asyncio
import json
import logging
import os
import time
from collections.abc import Callable

logger = logging.getLogger(__name__)
REDIS_URL = os.getenv("REDIS_URL", "redis://localhost:6379/2")

TASKS = {}


def register_task(name: str, fn: Callable):
    """Register a background task handler."""
    TASKS[name] = fn


async def enqueue(name: str, payload: dict, delay: int = 0, max_retries: int = 3):
    """Enqueue a background task."""
    import redis

    try:
        r = redis.from_url(REDIS_URL)
        task = json.dumps(
            {"name": name, "payload": payload, "retries": 0, "max_retries": max_retries, "created_at": time.time()}
        )
        if delay:
            r.zadd("rmi:task_queue:delayed", {task: time.time() + delay})
        else:
            r.lpush("rmi:task_queue", task)
        logger.debug(f"Task enqueued: {name}")
    except Exception as e:
        logger.warning(f"Task enqueue failed: {e}")


async def process_tasks() -> None:
    """Process tasks from the queue. Run as background loop."""
    import redis

    try:
        r = redis.from_url(REDIS_URL, decode_responses=True)
    except Exception:
        return

    while True:
        try:
            # Pop from queue
            task_data = r.brpop("rmi:task_queue", timeout=5)
            if not task_data:
                continue

            task = json.loads(task_data[1])
            name = task["name"]
            payload = task["payload"]
            retries = task["retries"]
            max_retries = task["max_retries"]

            handler = TASKS.get(name)
            if not handler:
                logger.warning(f"No handler for task: {name}")
                continue

            try:
                if asyncio.iscoroutinefunction(handler):
                    await handler(payload)
                else:
                    handler(payload)
                logger.debug(f"Task completed: {name}")
            except Exception as e:
                retries += 1
                if retries <= max_retries:
                    delay = 2**retries  # exponential backoff: 2, 4, 8 seconds
                    task["retries"] = retries
                    r.zadd("rmi:task_queue:delayed", {json.dumps(task): time.time() + delay})
                    logger.warning(f"Task {name} failed (attempt {retries}/{max_retries}), retrying in {delay}s")
                else:
                    # Dead letter queue
                    r.lpush(
                        "rmi:task_queue:dead", json.dumps({"task": task, "error": str(e), "failed_at": time.time()})
                    )
                    logger.error(f"Task {name} permanently failed after {max_retries} retries")
        except Exception as e:
            logger.error(f"Task processor error: {e}")
            await asyncio.sleep(1)


async def get_queue_stats() -> dict:
    """Get task queue statistics."""
    import redis

    try:
        r = redis.from_url(REDIS_URL)
        return {
            "pending": r.llen("rmi:task_queue"),
            "delayed": r.zcard("rmi:task_queue:delayed"),
            "dead": r.llen("rmi:task_queue:dead"),
        }
    except Exception:
        return {"error": "redis unavailable"}