AXe Skills HubSearch /

← All skills

queue-background-jobs

AXe First-party 

Reference: full SKILL.md

Below is the complete skill definition this hub loads when the skill is triggered — what the agent sees as its instructions, verbatim and unabridged.

Queue & Background Jobs

Role

You are an elite background processing engineer. You design reliable job queues that handle

failures gracefully, scale horizontally, and process millions of tasks with exactly-once

semantics and zero data loss.

Part 1: Redis Queue (RQ) Setup

from redis import Redis
from rq import Queue, Worker, Retry
from rq.job import Job
from datetime import timedelta

redis_conn = Redis(host="localhost", port=6379, db=0)

# Create queues with different priorities
high_q = Queue("high", connection=redis_conn)
default_q = Queue("default", connection=redis_conn)
low_q = Queue("low", connection=redis_conn)

# Enqueue jobs
def send_welcome_email(user_id: str, email: str):
    """This function runs in the worker process."""
    # ... send email logic
    pass

# Basic enqueue
job = default_q.enqueue(send_welcome_email, "user123", "[email protected]")

# With options
job = high_q.enqueue(
    send_welcome_email,
    "user123", "[email protected]",
    job_timeout=300,           # 5 min timeout
    result_ttl=86400,          # keep result 24h
    retry=Retry(max=3, interval=[10, 30, 60]),  # retry with backoff
    meta={"source": "signup"},
)

# Scheduled jobs
job = default_q.enqueue_in(
    timedelta(minutes=30),
    send_welcome_email,
    "user123", "[email protected]",
)

# Check job status
print(f"Status: {job.get_status()}")  # queued, started, finished, failed
print(f"Result: {job.result}")

Worker startup:

rq worker high default low --with-scheduler

Part 2: Celery Configuration

from celery import Celery, chain, group, chord
from celery.schedules import crontab

app = Celery("myapp")
app.config_from_object({
    "broker_url": "redis://localhost:6379/0",
    "result_backend": "redis://localhost:6379/1",
    "task_serializer": "json",
    "result_serializer": "json",
    "accept_content": ["json"],
    "task_acks_late": True,              # ack after completion (safer)
    "worker_prefetch_multiplier": 1,     # one task at a time per worker
    "task_reject_on_worker_lost": True,  # re-queue if worker dies
    "task_default_queue": "default",
    "task_routes": {
        "tasks.email.*": {"queue": "email"},
        "tasks.heavy.*": {"queue": "heavy"},
    },
    "beat_schedule": {
        "cleanup-every-hour": {
            "task": "tasks.maintenance.cleanup",
            "schedule": crontab(minute=0),
        },
        "daily-report": {
            "task": "tasks.reports.daily_summary",
            "schedule": crontab(hour=8, minute=0),
        },
    },
})

@app.task(bind=True, max_retries=3, default_retry_delay=60)
def process_upload(self, file_id: str):
    try:
        # ... process file
        pass
    except TemporaryError as e:
        self.retry(exc=e, countdown=2 ** self.request.retries * 10)

@app.task
def send_notification(user_id: str, message: str):
    pass

@app.task
def aggregate_results(results: list):
    return {"total": len(results), "results": results}

# Task composition
workflow = chain(
    process_upload.s("file1"),
    send_notification.s("Processing complete"),
)

# Fan-out / fan-in
batch = chord(
    group(process_upload.s(fid) for fid in file_ids),
    aggregate_results.s(),
)
# Start workers
celery -A tasks worker -Q default,email --concurrency=4 --loglevel=info
celery -A tasks worker -Q heavy --concurrency=2 --loglevel=info
celery -A tasks beat --loglevel=info  # scheduler

Part 3: Custom Async Worker Pool

import asyncio, uuid, time, json
from dataclasses import dataclass, field
from typing import Any, Callable
from enum import Enum

class JobStatus(str, Enum):
    PENDING = "pending"
    RUNNING = "running"
    COMPLETED = "completed"
    FAILED = "failed"
    RETRYING = "retrying"

@dataclass
class Job:
    id: str = field(default_factory=lambda: str(uuid.uuid4()))
    name: str = ""
    payload: dict = field(default_factory=dict)
    status: JobStatus = JobStatus.PENDING
    priority: int = 0
    max_retries: int = 3
    attempts: int = 0
    result: Any = None
    error: str | None = None
    created_at: float = field(default_factory=time.time)
    started_at: float = 0
    completed_at: float = 0

class AsyncWorkerPool:
    def __init__(self, num_workers: int = 4):
        self.num_workers = num_workers
        self.queue = asyncio.PriorityQueue()
        self.handlers: dict[str, Callable] = {}
        self.jobs: dict[str, Job] = {}
        self.dlq: list[Job] = []
        self._workers: list[asyncio.Task] = []
        self._shutdown = asyncio.Event()

    def register(self, name: str, handler: Callable):
        self.handlers[name] = handler

    async def submit(self, name: str, payload: dict, priority: int = 5,
                     max_retries: int = 3) -> str:
        job = Job(name=name, payload=payload, priority=priority, max_retries=max_retries)
        self.jobs[job.id] = job
        await self.queue.put((priority, job.created_at, job.id))
        return job.id

    async def _worker(self, worker_id: int):
        while not self._shutdown.is_set():
            try:
                _, _, job_id = await asyncio.wait_for(self.queue.get(), timeout=1.0)
            except asyncio.TimeoutError:
                continue

            job = self.jobs.get(job_id)
            if not job:
                continue

            handler = self.handlers.get(job.name)
            if not handler:
                job.status = JobStatus.FAILED
                job.error = f"No handler for '{job.name}'"
                self.dlq.append(job)
                continue

            job.status = JobStatus.RUNNING
            job.started_at = time.time()
            job.attempts += 1

            try:
                job.result = await handler(job.payload)
                job.status = JobStatus.COMPLETED
                job.completed_at = time.time()
            except Exception as e:
                if job.attempts < job.max_retries:
                    job.status = JobStatus.RETRYING
                    delay = 2 ** job.attempts
                    await asyncio.sleep(delay)
                    await self.queue.put((job.priority + 1, time.time(), job.id))
                else:
                    job.status = JobStatus.FAILED
                    job.error = str(e)
                    self.dlq.append(job)

    async def start(self):
        self._workers = [
            asyncio.create_task(self._worker(i))
            for i in range(self.num_workers)
        ]

    async def shutdown(self, timeout: float = 30):
        self._shutdown.set()
        await asyncio.wait(self._workers, timeout=timeout)

    def get_status(self, job_id: str) -> dict | None:
        job = self.jobs.get(job_id)
        if not job:
            return None
        return {
            "id": job.id, "status": job.status, "attempts": job.attempts,
            "result": job.result, "error": job.error,
        }

Part 4: Idempotent Workers

import hashlib

class IdempotentWorker:
    """Ensure jobs are processed exactly once."""

    def __init__(self, redis_conn):
        self.redis = redis_conn
        self.ttl = 86400 * 7  # 7 days

    def _job_key(self, job_name: str, payload: dict) -> str:
        content = json.dumps({"name": job_name, "payload": payload}, sort_keys=True)
        return f"idempotent:{hashlib.sha256(content.encode()).hexdigest()}"

    async def process_if_new(self, job_name: str, payload: dict, handler: Callable) -> dict:
        key = self._job_key(job_name, payload)

        # Check if already processed
        existing = self.redis.get(key)
        if existing:
            return {"status": "duplicate", "result": json.loads(existing)}

        # Acquire lock (prevents concurrent duplicates)
        lock_key = f"{key}:lock"
        if not self.redis.set(lock_key, "1", nx=True, ex=300):
            return {"status": "in_progress"}

        try:
            result = await handler(payload)
            self.redis.set(key, json.dumps(result), ex=self.ttl)
            return {"status": "processed", "result": result}
        finally:
            self.redis.delete(lock_key)

Part 5: Job Monitoring & Dashboard

from fastapi import FastAPI
from datetime import datetime

app = FastAPI()

class JobMonitor:
    def __init__(self, pool: AsyncWorkerPool):
        self.pool = pool

    def stats(self) -> dict:
        jobs = list(self.pool.jobs.values())
        return {
            "total": len(jobs),
            "pending": sum(1 for j in jobs if j.status == JobStatus.PENDING),
            "running": sum(1 for j in jobs if j.status == JobStatus.RUNNING),
            "completed": sum(1 for j in jobs if j.status == JobStatus.COMPLETED),
            "failed": sum(1 for j in jobs if j.status == JobStatus.FAILED),
            "dlq_size": len(self.pool.dlq),
            "queue_size": self.pool.queue.qsize(),
            "avg_duration": self._avg_duration(jobs),
        }

    def _avg_duration(self, jobs: list[Job]) -> float:
        completed = [j for j in jobs if j.completed_at > 0]
        if not completed:
            return 0
        return sum(j.completed_at - j.started_at for j in completed) / len(completed)

    def failed_jobs(self) -> list[dict]:
        return [
            {"id": j.id, "name": j.name, "error": j.error, "attempts": j.attempts,
             "payload": j.payload}
            for j in self.pool.dlq
        ]

monitor = JobMonitor(pool)

@app.get("/jobs/stats")
async def job_stats():
    return monitor.stats()

@app.get("/jobs/{job_id}")
async def job_status(job_id: str):
    return pool.get_status(job_id)

@app.get("/jobs/failed")
async def failed_jobs():
    return monitor.failed_jobs()

@app.post("/jobs/retry/{job_id}")
async def retry_job(job_id: str):
    job = pool.jobs.get(job_id)
    if not job:
        return {"error": "not found"}
    job.attempts = 0
    job.status = JobStatus.PENDING
    await pool.queue.put((job.priority, time.time(), job.id))
    return {"status": "requeued"}

Part 6: Graceful Shutdown

import signal

class GracefulWorkerPool(AsyncWorkerPool):
    def __init__(self, *args, **kwargs):
        super().__init__(*args, **kwargs)
        self._draining = False

    async def start(self):
        loop = asyncio.get_event_loop()
        for sig in (signal.SIGTERM, signal.SIGINT):
            loop.add_signal_handler(sig, lambda: asyncio.create_task(self.graceful_shutdown()))
        await super().start()

    async def graceful_shutdown(self, timeout: float = 30):
        if self._draining:
            return
        self._draining = True
        print(f"Graceful shutdown initiated. Waiting up to {timeout}s for running jobs...")

        # Stop accepting new jobs
        self._shutdown.set()

        # Wait for running jobs to finish
        start = time.time()
        while time.time() - start < timeout:
            running = [j for j in self.jobs.values() if j.status == JobStatus.RUNNING]
            if not running:
                break
            print(f"  {len(running)} jobs still running...")
            await asyncio.sleep(1)

        # Force-cancel remaining
        running = [j for j in self.jobs.values() if j.status == JobStatus.RUNNING]
        for job in running:
            job.status = JobStatus.FAILED
            job.error = "Shutdown timeout"
            self.dlq.append(job)

        for worker in self._workers:
            worker.cancel()

        print(f"Shutdown complete. {len(self.dlq)} jobs in DLQ.")

Part 7: Redis-Backed Persistent Queue

import redis

class RedisPersistentQueue:
    """Production queue using Redis lists for persistence."""

    def __init__(self, redis_url: str = "redis://localhost:6379", queue_name: str = "jobs"):
        self.redis = redis.from_url(redis_url)
        self.queue = f"queue:{queue_name}"
        self.processing = f"processing:{queue_name}"
        self.dlq_key = f"dlq:{queue_name}"

    def enqueue(self, job_name: str, payload: dict, priority: int = 0):
        job = json.dumps({"id": str(uuid.uuid4()), "name": job_name,
                          "payload": payload, "priority": priority, "attempts": 0})
        if priority > 5:
            self.redis.lpush(self.queue, job)   # high priority -> front
        else:
            self.redis.rpush(self.queue, job)   # normal -> back

    def dequeue(self, timeout: int = 5) -> dict | None:
        result = self.redis.brpoplpush(self.queue, self.processing, timeout=timeout)
        if result:
            return json.loads(result)
        return None

    def complete(self, job: dict):
        self.redis.lrem(self.processing, 1, json.dumps(job))

    def fail(self, job: dict, error: str):
        job["error"] = error
        job["attempts"] = job.get("attempts", 0) + 1
        self.redis.lrem(self.processing, 1, json.dumps(job))
        if job["attempts"] < 3:
            self.redis.rpush(self.queue, json.dumps(job))
        else:
            self.redis.rpush(self.dlq_key, json.dumps(job))

    def recover_stuck(self, max_age: int = 300):
        """Move stuck processing jobs back to queue."""
        stuck = self.redis.lrange(self.processing, 0, -1)
        for raw in stuck:
            self.redis.lrem(self.processing, 1, raw)
            self.redis.rpush(self.queue, raw)

    def stats(self) -> dict:
        return {
            "queued": self.redis.llen(self.queue),
            "processing": self.redis.llen(self.processing),
            "dead": self.redis.llen(self.dlq_key),
        }

AXE MCP Server Integration

Every skill in the AXE Skills Hub runs with access to the AXE MCP Server — giving it the full fleet intelligence toolkit automatically. No setup required; tools are available in any AXE-powered session.

Core Tools Available

CategoryToolsUse Case
Memoryread_memory, write_memory, list_memoryPersist context across sessions
Webweb_search, web_fetchLive data, docs, research
File Opsread_file, write_fileRead/write any local file
Fleetfleet_ssh, axe_pushRun commands on JL2/JL3/JL4, send notifications
AI Modelsquery_team_channel, get_partner_stateCross-agent coordination
Dataqdrant_search, qdrant_storeSemantic memory & vector search
Pipelinehydra_addAdd high-quality outputs to Edge training
Skillshub_list_skills, hub_get_skill, hub_search_skills, hub_get_registry, hub_skill_metadataChain skills together
Secretsget_secretRetrieve API keys securely

Quick Start

# In any AXE session, tools are pre-loaded. Example chaining:

# 1. Search for context
results = qdrant_search("user query here", collection="axe_persistent_memory")

# 2. Fetch live data if needed
content = web_fetch("https://docs.example.com/api")

# 3. Write result to memory for next session
write_memory("shared/last_result.md", output)

# 4. Log quality output to Edge training pipeline
hydra_add(prompt=user_query, response=output, score=0.9, source="skill-name")

Edge Training Integration

High-quality skill outputs are automatically eligible for Edge model training via hydra_add. When a response scores ≥0.85 in evals, pipe it to the Hydra pipeline to compound Edge's knowledge. This is how skills make Edge smarter over time.

# After generating a high-quality response:
hydra_add(
    prompt=user_input,
    response=final_output,
    score=0.9,          # eval score
    source="skill-name" # tracks provenance
)

Metadata

Category
Infrastructure
Tier
community
Version
1.0.0
License
MIT
Path
skills/queue-background-jobs/SKILL.md

Use with an agent

Fetch this skill’s definition over the open API — no key required.

curl -s /v1/skills/queue-background-jobs

View source ↗