First-party
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
```python
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:
```bash
rq worker high default low --with-scheduler
```
---
## Part 2: Celery Configuration
```python
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(),
)
```
```bash
# 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
```python
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
```python
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
```python
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
```python
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
```python
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
| Category | Tools | Use Case |
|----------|-------|----------|
| **Memory** | `read_memory`, `write_memory`, `list_memory` | Persist context across sessions |
| **Web** | `web_search`, `web_fetch` | Live data, docs, research |
| **File Ops** | `read_file`, `write_file` | Read/write any local file |
| **Fleet** | `fleet_ssh`, `axe_push` | Run commands on JL2/JL3/JL4, send notifications |
| **AI Models** | `query_team_channel`, `get_partner_state` | Cross-agent coordination |
| **Data** | `qdrant_search`, `qdrant_store` | Semantic memory & vector search |
| **Pipeline** | `hydra_add` | Add high-quality outputs to Edge training |
| **Skills** | `hub_list_skills`, `hub_get_skill`, `hub_search_skills`, `hub_get_registry`, `hub_skill_metadata` | Chain skills together |
| **Secrets** | `get_secret` | Retrieve API keys securely |
### Quick Start
```python
# 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.
```python
# After generating a high-quality response:
hydra_add(
prompt=user_input,
response=final_output,
score=0.9, # eval score
source="skill-name" # tracks provenance
)
```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.
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
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
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,
}
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)
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"}
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.")
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),
}
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.
| Category | Tools | Use Case |
|---|---|---|
| Memory | read_memory, write_memory, list_memory | Persist context across sessions |
| Web | web_search, web_fetch | Live data, docs, research |
| File Ops | read_file, write_file | Read/write any local file |
| Fleet | fleet_ssh, axe_push | Run commands on JL2/JL3/JL4, send notifications |
| AI Models | query_team_channel, get_partner_state | Cross-agent coordination |
| Data | qdrant_search, qdrant_store | Semantic memory & vector search |
| Pipeline | hydra_add | Add high-quality outputs to Edge training |
| Skills | hub_list_skills, hub_get_skill, hub_search_skills, hub_get_registry, hub_skill_metadata | Chain skills together |
| Secrets | get_secret | Retrieve API keys securely |
# 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")
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
)
Fetch this skill’s definition over the open API — no key required.
curl -s /v1/skills/queue-background-jobs