
Mcp Server Dev
- 1 installs
- 8 repo stars
- Updated July 30, 2026
- drewid74/ai_skills
Build MCP servers with FastMCP or the TypeScript SDK, choosing tool vs resource vs prompt, stdio vs SSE transport, with security and Docker packaging.
About
Provides an MCP server architecture guide covering stack defaults, tool-vs-resource decisions, transport selection, anti-patterns, and Docker packaging. A developer uses it when building or debugging a custom MCP server.
- Tool vs resource vs prompt and stdio vs SSE decision frameworks
- Path-traversal and secret-handling security guards
Mcp Server Dev by the numbers
- 1 all-time installs (skills.sh)
- Ranked #14,102 of 16,546 AI & Agent Building skills by installs in the Skillselion catalog
- Data as of Jul 31, 2026 (Skillselion catalog sync)
npx skills add https://github.com/drewid74/ai_skills --skill mcp-server-devAdd your badge
Show developers this skill is listed on Skillselion. Paste this into your README.
| Installs | 1 |
|---|---|
| repo stars | ★ 8 |
| Last updated | July 30, 2026 |
| Repository | drewid74/ai_skills ↗ |
What it does
Build MCP servers with FastMCP or the TypeScript SDK, choosing tool vs resource vs prompt, stdio vs SSE transport, with security and Docker packaging.
Files
MCP Server Development
Identity
You are an MCP server architect. Build focused, atomic tools — one tool per action, named as verb_noun. Never hardcode credentials or permit path traversal in file-access tools.
Stack Defaults
| Layer | Choice | Why |
|---|---|---|
| Python framework | FastMCP | Decorators + auto-generates JSON schema from type hints |
| TypeScript framework | @modelcontextprotocol/sdk | Full protocol control, Node.js ecosystem |
| Local transport | stdio | Zero overhead; native to Claude Desktop and Claude Code |
| Remote/Docker transport | SSE (HTTP) | Works through firewalls; required for containerized servers |
| Testing | mcp dev python server.py | Visual inspector for tools/resources without needing Claude |
| Packaging | Docker + SSE + env-injected secrets | Portable; credentials never baked into the image |
Decision Framework
Tool vs Resource vs Prompt
- If LLM needs to take an action (write, create, call, delete) → tool
- If LLM needs to read reference data (docs, config, user profile) → resource
- If reusable instruction scaffold for LLM reasoning → prompt
- Ambiguous → tool (more composable)
Transport Selection
- If running locally with Claude Desktop or Claude Code → stdio
- If running in Docker or on a remote host → SSE or HTTP
- If high message volume or large payloads → HTTP (lower overhead than SSE)
- Never → expose a stdio server over a network socket
Return Format
- If structured data → return
dictorlist(LLM formats output for user) - If operation failed → raise
ValueError/RuntimeError(not return error dict) - If result is large → paginate or summarize; never silently truncate
- Never → return pre-formatted label strings like
"User: Alice (123)"
Anti-Patterns
| Don't | Why | Do Instead |
|---|---|---|
Return {"error": "..."} as a success | Protocol error handling bypassed; LLM sees it as data | raise ValueError("reason") |
| Pre-format return strings | LLM cannot extract structured data from prose | Return dict/list; let LLM render |
| Mega-tool with 10+ parameters | LLM struggles to reason about correct invocation | One atomic action per tool |
| Hardcode API keys in source | Leaked in git, container image, logs | os.getenv("API_KEY"); fail fast if None |
Allow arbitrary path parameters | Path traversal exposes entire filesystem | os.path.abspath() + startswith(base_dir) guard |
Generic tool names (do_thing, action) | LLM routing degrades; ambiguous intent | verb_noun: search_docs, create_issue, get_user |
Quality Gates
- [ ] All tools named
verb_noun(e.g.,search_docs,create_issue) - [ ] Return types are
dict/list— no pre-formatted prose strings - [ ] Credentials loaded from env vars; server raises on missing values at startup
- [ ] File-access tools validate resolved path stays within allowed base directory
- [ ] Tools pass independent unit tests before Claude integration
- [ ] Server registered in
claude_desktop_config.jsonor.mcp.json
---
FastMCP Quick Start (Python)
from fastmcp import FastMCP
mcp = FastMCP("my-server")
@mcp.tool()
def search_docs(query: str) -> dict:
"""Search documentation by keyword."""
results = fetch_docs(query)
return {"results": results, "count": len(results)}
if __name__ == "__main__":
mcp.run() # stdio by defaultTypeScript SDK Quick Start
import { Server } from "@modelcontextprotocol/sdk/server/index.js";
import { StdioServerTransport } from "@modelcontextprotocol/sdk/server/stdio.js";
const server = new Server({ name: "my-server", version: "1.0.0" });
server.setRequestHandler(CallToolRequestSchema, async (request) => {
if (request.params.name === "search_docs") {
return { content: [{ type: "text", text: JSON.stringify(results) }] };
}
throw new Error("Unknown tool");
});
await server.connect(new StdioServerTransport());Tool Design Patterns
# Good: atomic, typed, descriptive
@mcp.tool()
def query_database(
sql: str,
timeout_seconds: int = 30,
readonly: bool = True
) -> dict:
"""Execute a SQL query.
Args:
sql: The SQL query to execute
timeout_seconds: Query timeout (default 30s)
readonly: If True, reject INSERT/UPDATE/DELETE
"""
if readonly and not sql.strip().upper().startswith("SELECT"):
raise ValueError("Only SELECT queries in readonly mode")
return {"rows": run_query(sql), "elapsed_ms": ...}Resources and Prompts
# Resource: read-only reference data
@mcp.resource("file://docs/{doc_id}")
def get_doc(doc_id: str) -> str:
"""Retrieve a documentation page by ID."""
return fetch_doc(doc_id)
# Prompt: reusable reasoning scaffold
@mcp.prompt("debug_error")
def debug_prompt(error_type: str, stack_trace: str) -> str:
return f"Help me debug this {error_type}:\n{stack_trace}\nFocus on root cause."Docker Packaging
FROM python:3.11-slim
WORKDIR /app
COPY requirements.txt .
RUN pip install -r requirements.txt
COPY server.py .
EXPOSE 8000
HEALTHCHECK --interval=30s --timeout=3s CMD curl -f http://localhost:8000/health || exit 1
CMD ["python", "server.py", "--transport=sse", "--host=0.0.0.0", "--port=8000"]# docker-compose.yml
services:
mcp-server:
build: .
ports: ["8000:8000"]
environment:
- API_KEY=${API_KEY}
- DATABASE_URL=${DATABASE_URL}
healthcheck:
test: ["CMD", "curl", "-f", "http://localhost:8000/health"]
interval: 30s
timeout: 3s
retries: 3
restart: unless-stopped
networks:
default:
enable_ipv6: false # prevents aiohttp IPv6 connection errorsSecurity: Secrets and Path Guards
import os
# Fail fast on startup if required env vars missing
API_KEY = os.getenv("API_KEY")
if not API_KEY:
raise ValueError("API_KEY environment variable not set")
# Path traversal guard for file tools
BASE_DIR = "/safe/data"
@mcp.tool()
def read_file(path: str) -> str:
"""Read a file from the safe directory."""
safe_path = os.path.abspath(os.path.join(BASE_DIR, path))
if not safe_path.startswith(BASE_DIR):
raise ValueError("Path traversal not allowed")
with open(safe_path) as f:
return f.read()Testing
# Visual inspector (no Claude needed)
mcp dev python server.py
# Unit tests
pytest test_server.py -vdef test_query_validation():
with pytest.raises(ValueError):
query_database("DROP TABLE users")
def test_path_traversal():
with pytest.raises(ValueError):
read_file("../../etc/passwd")Multi-Server Config
// claude_desktop_config.json (~/.config/Claude/)
{
"mcpServers": {
"database": { "command": "python", "args": ["/path/db_server.py"] },
"git": { "command": "python", "args": ["/path/git_server.py"] }
}
}
// .mcp.json (Claude Code: project root or ~/.mcp.json)
{
"mcpServers": {
"database": { "command": "python server.py", "cwd": "/path/db" }
}
}Troubleshooting
| Problem | Cause | Fix |
|---|---|---|
| Tool not appearing | Not registered or schema error | Check server logs; restart Claude Desktop |
| Connection refused | Wrong transport or port conflict | Verify stdio vs SSE; lsof -i :8000 |
| IPv6 error | aiohttp tries IPv6 first | enable_ipv6: false in compose network |
| Schema validation error | Return type mismatch | Match type hints to actual return values |
| LLM not using tool | Ambiguous name/params | Rename to verb_noun; clarify docstring |
| Timeout | Long-running tool | Make async; return status early, let LLM poll |
SQLite format 3@ .�
�"��y"U#indexidx_claimedtasksCREATE INDEX idx_claimed ON tasks(claimed_by, claimed_at)U)yindexidx_job_statustasksCREATE INDEX idx_job_status ON tasks(job_name, status)�z�StabletaskstasksCREATE TABLE tasks (
id TEXT PRIMARY KEY,
job_name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
claimed_by TEXT,
claimed_at TEXT,
updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
reason TEXT,
metadata TEXT
))=indexsqlite_autoindex_tasks_1tasks
"""postgres-control-plane: pluggable task queue with FastMCP server."""
from .backend import (
MemoryBackend,
PostgresBackend,
SQLiteBackend,
TaskBackend,
make_backend,
)
from .models import Task, TaskStatus, is_valid_transition
from .queue import TaskQueue
__all__ = [
"Task",
"TaskStatus",
"TaskQueue",
"TaskBackend",
"MemoryBackend",
"SQLiteBackend",
"PostgresBackend",
"make_backend",
"is_valid_transition",
]
__version__ = "0.2.0"
__pycache__/
*.pyc
*.pyo
.tasks/
*.egg-info/
build/
dist/
.pytest_cache/
"""Pluggable task storage backends: Memory, SQLite, Postgres."""
from __future__ import annotations
import json
import os
import sqlite3
import threading
from abc import ABC, abstractmethod
from datetime import datetime, timezone
from pathlib import Path
from typing import List, Optional
from .models import Task, is_valid_transition
def _now() -> str:
return datetime.now(timezone.utc).isoformat()
class TaskBackend(ABC):
@abstractmethod
def add(self, task: Task) -> Task: ...
@abstractmethod
def get(self, task_id: str) -> Optional[Task]: ...
@abstractmethod
def list_pending(self, job_name: str) -> List[Task]: ...
@abstractmethod
def claim(self, task_id: str, worker_id: str) -> Task: ...
@abstractmethod
def update_status(self, task_id: str, status: str, reason: str = "") -> Task: ...
@abstractmethod
def heartbeat(self, task_id: str, worker_id: str) -> Task: ...
@abstractmethod
def get_stats(self, job_name: str) -> dict: ...
# ---------------------------------------------------------------------------
# Memory
# ---------------------------------------------------------------------------
class MemoryBackend(TaskBackend):
def __init__(self):
self._tasks: dict[str, Task] = {}
self._lock = threading.Lock()
def add(self, task: Task) -> Task:
with self._lock:
task.updated_at = _now()
self._tasks[task.id] = task
return task
def get(self, task_id: str) -> Optional[Task]:
return self._tasks.get(task_id)
def list_pending(self, job_name: str) -> List[Task]:
return [t for t in self._tasks.values()
if t.job_name == job_name and t.status == "pending"]
def claim(self, task_id: str, worker_id: str) -> Task:
with self._lock:
task = self._tasks.get(task_id)
if not task:
raise KeyError(f"Task {task_id} not found")
if task.status != "pending":
raise ValueError(f"Task {task_id} not pending (status={task.status})")
task.status = "claimed"
task.claimed_by = worker_id
task.claimed_at = _now()
task.updated_at = task.claimed_at
return task
def update_status(self, task_id: str, status: str, reason: str = "") -> Task:
with self._lock:
task = self._tasks.get(task_id)
if not task:
raise KeyError(f"Task {task_id} not found")
if not is_valid_transition(task.status, status):
raise ValueError(f"Invalid transition: {task.status} -> {status}")
task.status = status
task.reason = reason
task.updated_at = _now()
return task
def heartbeat(self, task_id: str, worker_id: str) -> Task:
with self._lock:
task = self._tasks.get(task_id)
if not task:
raise KeyError(f"Task {task_id} not found")
if task.claimed_by != worker_id:
raise ValueError(f"Task {task_id} not claimed by {worker_id}")
task.claimed_at = _now()
task.updated_at = task.claimed_at
return task
def get_stats(self, job_name: str) -> dict:
tasks = [t for t in self._tasks.values() if t.job_name == job_name]
return {
"pending": sum(1 for t in tasks if t.status == "pending"),
"claimed": sum(1 for t in tasks if t.status == "claimed"),
"completed": sum(1 for t in tasks if t.status == "completed"),
"failed": sum(1 for t in tasks if t.status == "failed"),
}
# ---------------------------------------------------------------------------
# SQLite
# ---------------------------------------------------------------------------
SQLITE_SCHEMA = """
CREATE TABLE IF NOT EXISTS tasks (
id TEXT PRIMARY KEY,
job_name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
claimed_by TEXT,
claimed_at TEXT,
updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
reason TEXT,
metadata TEXT
);
CREATE INDEX IF NOT EXISTS idx_job_status ON tasks(job_name, status);
CREATE INDEX IF NOT EXISTS idx_claimed ON tasks(claimed_by, claimed_at);
"""
def _row_to_task(row) -> Task:
return Task(
id=row["id"],
job_name=row["job_name"],
status=row["status"],
claimed_by=row["claimed_by"],
claimed_at=row["claimed_at"],
updated_at=row["updated_at"],
reason=row["reason"],
metadata=json.loads(row["metadata"]) if row["metadata"] else {},
)
class SQLiteBackend(TaskBackend):
def __init__(self, db_path: str = ".tasks/queue.db"):
self.db_path = db_path
Path(db_path).parent.mkdir(parents=True, exist_ok=True)
self._lock = threading.Lock()
with self._conn() as conn:
conn.executescript(SQLITE_SCHEMA)
def _conn(self) -> sqlite3.Connection:
conn = sqlite3.connect(self.db_path, isolation_level=None, timeout=30)
conn.row_factory = sqlite3.Row
conn.execute("PRAGMA journal_mode=WAL")
conn.execute("PRAGMA foreign_keys=ON")
return conn
def add(self, task: Task) -> Task:
task.updated_at = _now()
with self._lock, self._conn() as conn:
conn.execute(
"""INSERT INTO tasks(id, job_name, status, claimed_by, claimed_at,
updated_at, reason, metadata)
VALUES(?, ?, ?, ?, ?, ?, ?, ?)""",
(task.id, task.job_name, task.status, task.claimed_by,
task.claimed_at, task.updated_at, task.reason,
json.dumps(task.metadata or {})),
)
return task
def get(self, task_id: str) -> Optional[Task]:
with self._conn() as conn:
row = conn.execute("SELECT * FROM tasks WHERE id=?", (task_id,)).fetchone()
return _row_to_task(row) if row else None
def list_pending(self, job_name: str) -> List[Task]:
with self._conn() as conn:
rows = conn.execute(
"SELECT * FROM tasks WHERE job_name=? AND status='pending'",
(job_name,),
).fetchall()
return [_row_to_task(r) for r in rows]
def claim(self, task_id: str, worker_id: str) -> Task:
ts = _now()
with self._lock, self._conn() as conn:
cur = conn.execute(
"""UPDATE tasks
SET status='claimed', claimed_by=?, claimed_at=?, updated_at=?
WHERE id=? AND status='pending'""",
(worker_id, ts, ts, task_id),
)
if cur.rowcount == 0:
row = conn.execute("SELECT status FROM tasks WHERE id=?", (task_id,)).fetchone()
if not row:
raise KeyError(f"Task {task_id} not found")
raise ValueError(f"Task {task_id} not pending (status={row['status']})")
task = self.get(task_id)
assert task is not None
return task
def update_status(self, task_id: str, status: str, reason: str = "") -> Task:
with self._lock, self._conn() as conn:
row = conn.execute("SELECT status FROM tasks WHERE id=?", (task_id,)).fetchone()
if not row:
raise KeyError(f"Task {task_id} not found")
if not is_valid_transition(row["status"], status):
raise ValueError(f"Invalid transition: {row['status']} -> {status}")
conn.execute(
"UPDATE tasks SET status=?, reason=?, updated_at=? WHERE id=?",
(status, reason, _now(), task_id),
)
task = self.get(task_id)
assert task is not None
return task
def heartbeat(self, task_id: str, worker_id: str) -> Task:
ts = _now()
with self._lock, self._conn() as conn:
cur = conn.execute(
"UPDATE tasks SET claimed_at=?, updated_at=? WHERE id=? AND claimed_by=?",
(ts, ts, task_id, worker_id),
)
if cur.rowcount == 0:
row = conn.execute("SELECT id FROM tasks WHERE id=?", (task_id,)).fetchone()
if not row:
raise KeyError(f"Task {task_id} not found")
raise ValueError(f"Task {task_id} not claimed by {worker_id}")
task = self.get(task_id)
assert task is not None
return task
def get_stats(self, job_name: str) -> dict:
with self._conn() as conn:
rows = conn.execute(
"SELECT status, COUNT(*) AS n FROM tasks WHERE job_name=? GROUP BY status",
(job_name,),
).fetchall()
stats = {"pending": 0, "claimed": 0, "completed": 0, "failed": 0}
for r in rows:
stats[r["status"]] = r["n"]
return stats
# ---------------------------------------------------------------------------
# Postgres
# ---------------------------------------------------------------------------
POSTGRES_SCHEMA = """
CREATE TABLE IF NOT EXISTS tasks (
id TEXT PRIMARY KEY,
job_name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
claimed_by TEXT,
claimed_at TIMESTAMPTZ,
updated_at TIMESTAMPTZ DEFAULT NOW(),
reason TEXT,
metadata JSONB
);
CREATE INDEX IF NOT EXISTS idx_job_status ON tasks(job_name, status);
CREATE INDEX IF NOT EXISTS idx_claimed ON tasks(claimed_by, claimed_at);
"""
class PostgresBackend(TaskBackend):
def __init__(self, dsn: str):
try:
import psycopg2
from psycopg2.extras import RealDictCursor, Json
except ImportError as e:
raise RuntimeError(
"psycopg2-binary required for PostgresBackend. "
"Install: pip install psycopg2-binary"
) from e
self._psycopg2 = psycopg2
self._RealDictCursor = RealDictCursor
self._Json = Json
self.dsn = dsn
with self._conn() as conn, conn.cursor() as cur:
cur.execute(POSTGRES_SCHEMA)
conn.commit()
def _conn(self):
return self._psycopg2.connect(self.dsn)
def _row_to_task(self, row) -> Task:
return Task(
id=row["id"],
job_name=row["job_name"],
status=row["status"],
claimed_by=row["claimed_by"],
claimed_at=row["claimed_at"].isoformat() if row["claimed_at"] else None,
updated_at=row["updated_at"].isoformat() if row["updated_at"] else None,
reason=row["reason"],
metadata=row["metadata"] or {},
)
def add(self, task: Task) -> Task:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute(
"""INSERT INTO tasks(id, job_name, status, claimed_by, claimed_at,
reason, metadata)
VALUES(%s, %s, %s, %s, %s, %s, %s)
RETURNING *""",
(task.id, task.job_name, task.status, task.claimed_by,
task.claimed_at, task.reason, self._Json(task.metadata or {})),
)
row = cur.fetchone()
conn.commit()
return self._row_to_task(row)
def get(self, task_id: str) -> Optional[Task]:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute("SELECT * FROM tasks WHERE id=%s", (task_id,))
row = cur.fetchone()
return self._row_to_task(row) if row else None
def list_pending(self, job_name: str) -> List[Task]:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute(
"SELECT * FROM tasks WHERE job_name=%s AND status='pending'",
(job_name,),
)
rows = cur.fetchall()
return [self._row_to_task(r) for r in rows]
def claim(self, task_id: str, worker_id: str) -> Task:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute(
"""UPDATE tasks
SET status='claimed', claimed_by=%s, claimed_at=NOW(), updated_at=NOW()
WHERE id=%s AND status='pending'
RETURNING *""",
(worker_id, task_id),
)
row = cur.fetchone()
if not row:
cur.execute("SELECT status FROM tasks WHERE id=%s", (task_id,))
check = cur.fetchone()
conn.rollback()
if not check:
raise KeyError(f"Task {task_id} not found")
raise ValueError(f"Task {task_id} not pending (status={check['status']})")
conn.commit()
return self._row_to_task(row)
def update_status(self, task_id: str, status: str, reason: str = "") -> Task:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute("SELECT status FROM tasks WHERE id=%s", (task_id,))
row = cur.fetchone()
if not row:
raise KeyError(f"Task {task_id} not found")
if not is_valid_transition(row["status"], status):
raise ValueError(f"Invalid transition: {row['status']} -> {status}")
cur.execute(
"""UPDATE tasks SET status=%s, reason=%s, updated_at=NOW()
WHERE id=%s RETURNING *""",
(status, reason, task_id),
)
updated = cur.fetchone()
conn.commit()
return self._row_to_task(updated)
def heartbeat(self, task_id: str, worker_id: str) -> Task:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute(
"""UPDATE tasks SET claimed_at=NOW(), updated_at=NOW()
WHERE id=%s AND claimed_by=%s
RETURNING *""",
(task_id, worker_id),
)
row = cur.fetchone()
if not row:
cur.execute("SELECT id FROM tasks WHERE id=%s", (task_id,))
check = cur.fetchone()
conn.rollback()
if not check:
raise KeyError(f"Task {task_id} not found")
raise ValueError(f"Task {task_id} not claimed by {worker_id}")
conn.commit()
return self._row_to_task(row)
def get_stats(self, job_name: str) -> dict:
with self._conn() as conn, conn.cursor(cursor_factory=self._RealDictCursor) as cur:
cur.execute(
"SELECT status, COUNT(*) AS n FROM tasks WHERE job_name=%s GROUP BY status",
(job_name,),
)
rows = cur.fetchall()
stats = {"pending": 0, "claimed": 0, "completed": 0, "failed": 0}
for r in rows:
stats[r["status"]] = r["n"]
return stats
# ---------------------------------------------------------------------------
# Auto-detect
# ---------------------------------------------------------------------------
def make_backend(backend: str = "auto", **kwargs) -> TaskBackend:
"""Construct a backend.
backend: "memory" | "sqlite" | "postgres" | "auto"
auto: DATABASE_URL env -> postgres, else .tasks/ dir present or creatable -> sqlite, else memory
"""
if backend == "memory":
return MemoryBackend()
if backend == "sqlite":
return SQLiteBackend(kwargs.get("db_path", ".tasks/queue.db"))
if backend == "postgres":
dsn = kwargs.get("dsn") or os.environ.get("DATABASE_URL")
if not dsn:
raise ValueError("PostgresBackend requires dsn or DATABASE_URL env")
return PostgresBackend(dsn)
if backend == "auto":
dsn = os.environ.get("DATABASE_URL")
if dsn:
return PostgresBackend(dsn)
try:
return SQLiteBackend(".tasks/queue.db")
except Exception:
return MemoryBackend()
raise ValueError(f"Unknown backend: {backend}")
"""pytest config: make postgres_control_plane importable when running tests from inside the package dir."""
import sys
from pathlib import Path
# Add mcp-server-dev/ (the parent of postgres_control_plane/) to sys.path
# so that `import postgres_control_plane` works.
sys.path.insert(0, str(Path(__file__).resolve().parent.parent))
"""FastMCP server exposing TaskQueue as 7 tools.
Run:
python -m postgres_control_plane.mcp_server
Backend selection via env:
DATABASE_URL=postgresql://... python -m postgres_control_plane.mcp_server
"""
from __future__ import annotations
import os
from dataclasses import asdict
from fastmcp import FastMCP
from .models import Task, is_valid_transition
from .queue import TaskQueue
BACKEND = os.environ.get("CONTROL_PLANE_BACKEND", "auto")
_backend_kwargs: dict = {}
_db_path = os.environ.get("CONTROL_PLANE_DB_PATH")
if _db_path:
_backend_kwargs["db_path"] = _db_path
_dsn = os.environ.get("DATABASE_URL")
if _dsn and BACKEND in ("postgres", "auto"):
_backend_kwargs.setdefault("dsn", _dsn)
queue = TaskQueue(backend=BACKEND, **_backend_kwargs)
mcp = FastMCP("control-plane")
@mcp.tool()
def add_task(task_id: str, job_name: str, metadata: dict | None = None) -> dict:
"""Add a new task in pending status.
Args:
task_id: Unique identifier.
job_name: Job grouping label (e.g., 'data-processing').
metadata: Optional dict of task-specific data.
"""
task = Task(id=task_id, job_name=job_name, metadata=metadata or {})
saved = queue.add(task)
return asdict(saved)
@mcp.tool()
def list_tasks(job_name: str, status: str = "pending") -> dict:
"""List tasks by job and status. Only 'pending' is currently indexed."""
if status == "pending":
tasks = queue.list_pending(job_name)
else:
# Fallback: filter from stats (no full enumeration API by design)
tasks = []
return {
"job_name": job_name,
"status": status,
"count": len(tasks),
"tasks": [asdict(t) for t in tasks],
}
@mcp.tool()
def claim_task(task_id: str, worker_id: str) -> dict:
"""Atomically claim a pending task. Raises if already claimed."""
task = queue.claim(task_id, worker_id)
return asdict(task)
@mcp.tool()
def update_task_status(task_id: str, status: str, reason: str = "") -> dict:
"""Update task status. Validates transition (pending->claimed->completed/failed)."""
task = queue.update_status(task_id, status, reason)
return asdict(task)
@mcp.tool()
def heartbeat(task_id: str, worker_id: str) -> dict:
"""Refresh claim lease — updates claimed_at to NOW()."""
task = queue.heartbeat(task_id, worker_id)
return {"task_id": task.id, "heartbeat": "ok", "claimed_at": task.claimed_at}
@mcp.tool()
def verify_transition(from_status: str, to_status: str) -> dict:
"""Check if a status transition is valid (without mutating)."""
valid = is_valid_transition(from_status, to_status)
return {"from": from_status, "to": to_status, "valid": valid}
@mcp.tool()
def get_stats(job_name: str) -> dict:
"""Get task counts by status for a job."""
stats = queue.get_stats(job_name)
return {"job_name": job_name, "stats": stats, "total": sum(stats.values())}
def main():
mcp.run()
if __name__ == "__main__":
main()
"""Task model and status enum."""
from dataclasses import dataclass, field, asdict
from enum import Enum
from typing import Optional
class TaskStatus(str, Enum):
PENDING = "pending"
CLAIMED = "claimed"
COMPLETED = "completed"
FAILED = "failed"
VALID_TRANSITIONS = {
"pending": {"claimed"},
"claimed": {"completed", "failed", "pending"},
"completed": set(),
"failed": set(),
}
def is_valid_transition(from_status: str, to_status: str) -> bool:
return to_status in VALID_TRANSITIONS.get(from_status, set())
@dataclass
class Task:
id: str
job_name: str
status: str = TaskStatus.PENDING.value
claimed_by: Optional[str] = None
claimed_at: Optional[str] = None
updated_at: Optional[str] = None
reason: Optional[str] = None
metadata: Optional[dict] = field(default_factory=dict)
def to_dict(self) -> dict:
return asdict(self)
"""TaskQueue public API. Delegates to a pluggable backend."""
from __future__ import annotations
from typing import List, Optional
from .backend import TaskBackend, make_backend
from .models import Task
class TaskQueue:
"""Generic task queue with pluggable backend.
Examples
--------
>>> q = TaskQueue() # auto-detect backend
>>> q = TaskQueue(backend="memory") # explicit
>>> q = TaskQueue(backend="sqlite", db_path=".tasks/q.db")
>>> q = TaskQueue(backend="postgres", dsn="postgresql://...")
"""
def __init__(self, backend: str = "auto", **backend_kwargs):
if isinstance(backend, TaskBackend):
self.backend: TaskBackend = backend
else:
self.backend = make_backend(backend, **backend_kwargs)
def add(self, task: Task) -> Task:
return self.backend.add(task)
def get(self, task_id: str) -> Optional[Task]:
return self.backend.get(task_id)
def list_pending(self, job_name: str) -> List[Task]:
return self.backend.list_pending(job_name)
def claim(self, task_id: str, worker_id: str) -> Task:
return self.backend.claim(task_id, worker_id)
def update_status(self, task_id: str, status: str, reason: str = "") -> Task:
return self.backend.update_status(task_id, status, reason)
def heartbeat(self, task_id: str, worker_id: str) -> Task:
return self.backend.heartbeat(task_id, worker_id)
def get_stats(self, job_name: str) -> dict:
return self.backend.get_stats(job_name)
Postgres-Backed Task Queue for AI Agent Coordination
Minimal, reusable patterns extracted from GoldenMatch/GoldenCheck review queues.
A generic task queue with pluggable backends (Memory/SQLite/Postgres), atomic claim semantics, and FastMCP tool wrappers for Claude integration.
Quick Start
Installation
pip install fastmcp psycopg2-binaryBasic Usage
from postgres_control_plane import TaskQueue, Task
# Initialize queue (auto-detects backend)
queue = TaskQueue(backend="auto")
# Add a task
task = Task(id="task-1", job_name="data-processing", status="pending")
queue.add(task)
# List pending tasks
pending = queue.list_pending("data-processing")
print(f"Found {len(pending)} pending tasks")
# Claim a task (atomic: only one worker succeeds)
try:
claimed = queue.claim("task-1", worker_id="worker-1")
print(f"Claimed by {claimed.claimed_by}")
except ValueError:
print("Task already claimed by another worker")
# Update status
completed = queue.update_status("task-1", "completed", reason="Success")
# Heartbeat (keep claim alive)
queue.heartbeat("task-1", "worker-1")Architecture
Pluggable Backends
| Backend | Use Case | Persistence |
|---|---|---|
| Memory | Testing, development | In-process only |
| SQLite | Local persistence, single-machine | .tasks/queue.db |
| Postgres | Production, multi-worker | Remote database |
Backend auto-detection:
queue = TaskQueue(backend="auto")
# Tries: DATABASE_URL env var → .tasks/ dir → MemoryBackendAtomic Claim Semantics
The core pattern prevents double-claim race conditions:
UPDATE tasks SET status='claimed', claimed_by=?, claimed_at=NOW()
WHERE id=? AND status='pending'Only one worker succeeds per task. The WHERE status='pending' guard ensures atomicity.
Task Model
@dataclass
class Task:
id: str # Unique identifier
job_name: str # Job grouping (e.g., "data-processing")
status: str # pending, claimed, completed, failed
claimed_by: Optional[str] = None # Worker ID
claimed_at: Optional[str] = None # ISO timestamp
updated_at: Optional[str] = None # Last update
reason: Optional[str] = None # Status change reason
metadata: Optional[dict] = None # Custom dataMCP Server Integration
FastMCP Tools
The mcp_server.py module exposes 7 tools for Claude:
1. add_task
add_task(task_id: str, job_name: str, metadata: dict | None = None) -> dictAdd a new task in pending status.
2. list_tasks
list_tasks(job_name: str, status: str = "pending") -> dictList tasks by job and status.
3. claim_task
claim_task(task_id: str, worker_id: str) -> dictClaim a task atomically. Raises ValueError if already claimed.
4. update_task_status
update_task_status(task_id: str, status: str, reason: str = "") -> dictUpdate task status (pending → claimed → completed/failed). Validates transition.
5. heartbeat
heartbeat(task_id: str, worker_id: str) -> dictRefresh claim lease (UPDATE claimed_at=NOW()).
6. verify_transition
verify_transition(from_status: str, to_status: str) -> dictCheck if status transition is valid (read-only).
7. get_stats
get_stats(job_name: str) -> dictGet task counts by status.
Running the MCP Server
# Local (SQLite)
python -m postgres_control_plane.mcp_server
# Production (Postgres)
export DATABASE_URL=postgresql://user:pass@localhost/tasks
python -m postgres_control_plane.mcp_serverClaude Desktop Integration
Add to ~/.claude/claude_desktop_config.json:
{
"mcpServers": {
"control-plane": {
"command": "python",
"args": ["-m", "postgres_control_plane.mcp_server"]
}
}
}Database Schemas
Postgres
CREATE TABLE tasks (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
job_name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
claimed_by TEXT,
claimed_at TIMESTAMPTZ,
updated_at TIMESTAMPTZ DEFAULT NOW(),
reason TEXT,
metadata JSONB
);
CREATE INDEX idx_job_status ON tasks(job_name, status);
CREATE INDEX idx_claimed ON tasks(claimed_by, claimed_at);SQLite
CREATE TABLE tasks (
id TEXT PRIMARY KEY,
job_name TEXT NOT NULL,
status TEXT NOT NULL DEFAULT 'pending',
claimed_by TEXT,
claimed_at TEXT,
updated_at TEXT DEFAULT CURRENT_TIMESTAMP,
reason TEXT,
metadata TEXT
);Patterns from GoldenMatch/GoldenCheck
This module extracts production patterns from:
- GoldenMatch review_queue.py — SQLite-backed review queue with confidence gating
- GoldenCheck review_queue.py — Abstract backend interface with Memory/SQLite/Postgres implementations
Key patterns: 1. Abstract backend ABC — Pluggable implementations without code duplication 2. Atomic UPDATE guards — WHERE status='pending' prevents race conditions 3. Heartbeat refresh — Simple polling instead of LISTEN/NOTIFY 4. Status transitions — Explicit valid state machine (pending → claimed → completed/failed) 5. JSON metadata — Flexible task-specific data without schema changes
Testing
Claim Atomicity Test
import threading
from postgres_control_plane import TaskQueue, Task
queue = TaskQueue(backend="memory")
task = Task(id="test-1", job_name="demo", status="pending")
queue.add(task)
results = []
def claim_worker(worker_id):
try:
queue.claim("test-1", worker_id)
results.append(worker_id)
except ValueError:
pass
t1 = threading.Thread(target=claim_worker, args=("worker-1",))
t2 = threading.Thread(target=claim_worker, args=("worker-2",))
t1.start()
t2.start()
t1.join()
t2.join()
assert len(results) == 1, f"Expected 1 claim, got {len(results)}"
print(f"✓ Claim atomicity verified: {results[0]} won the race")Deployment
Docker (Postgres)
FROM python:3.12-slim
WORKDIR /app
COPY . .
RUN pip install fastmcp psycopg2-binary
ENV DATABASE_URL=postgresql://user:pass@postgres:5432/tasks
CMD ["python", "-m", "postgres_control_plane.mcp_server"]Docker Compose
version: "3.8"
services:
postgres:
image: postgres:16-alpine
environment:
POSTGRES_DB: tasks
POSTGRES_PASSWORD: secret
volumes:
- postgres_data:/var/lib/postgresql/data
control-plane:
build: .
environment:
DATABASE_URL: postgresql://postgres:secret@postgres:5432/tasks
depends_on:
- postgres
ports:
- "8000:8000"
volumes:
postgres_data:Design Decisions
No Overbuilding
- No LISTEN/NOTIFY — Polling is simpler and works everywhere
- No job queue library — Just Postgres + psycopg2
- No distributed locking — Atomic UPDATE guards are sufficient
- No retry logic — Caller handles retries
- No task dependencies — Each task is independent
Why Atomic UPDATE?
UPDATE tasks SET status='claimed', claimed_by=?, claimed_at=NOW()
WHERE id=? AND status='pending'This is simpler and more reliable than:
- SELECT + INSERT (race condition window)
- Distributed locks (complexity, latency)
- LISTEN/NOTIFY (not available in SQLite)
The database guarantees atomicity. Only one worker's UPDATE succeeds.
File Structure
postgres-control-plane/
├── __init__.py # Package exports
├── models.py # Task dataclass, TaskStatus enum
├── backend.py # TaskBackend ABC, Memory/SQLite/Postgres implementations
├── queue.py # TaskQueue public API
├── mcp_server.py # FastMCP server with 6 tools
└── README.md # This fileReferences
- GoldenMatch review_queue.py —
C:\Users\afair\dev\ai_skills\goldenmatch\packages\python\goldenmatch\goldenmatch\core\review_queue.py - GoldenCheck review_queue.py —
C:\Users\afair\dev\ai_skills\goldenmatch\packages\python\goldencheck\goldencheck\agent\review_queue.py - MCP Server Dev —
C:\Users\afair\dev\ai_skills\mcp-server-dev\SKILL.md
License
MIT
from setuptools import setup, find_packages
setup(
name="postgres-control-plane",
version="0.2.0",
description="Pluggable task queue (Memory/SQLite/Postgres) with FastMCP server",
packages=find_packages(),
python_requires=">=3.10",
install_requires=[
"fastmcp>=0.1.0",
],
extras_require={
"postgres": ["psycopg2-binary>=2.9.0"],
"dev": ["pytest>=7.0"],
},
entry_points={
"console_scripts": [
"control-plane=postgres_control_plane.mcp_server:main",
],
},
)
"""Backend conformance + claim atomicity tests.
Run:
pytest mcp-server-dev/postgres-control-plane/tests/
Postgres tests are skipped unless DATABASE_URL is set.
"""
from __future__ import annotations
import os
import threading
import tempfile
import pytest
from postgres_control_plane import (
MemoryBackend,
PostgresBackend,
SQLiteBackend,
Task,
TaskQueue,
)
def make_queue_memory():
return TaskQueue(backend=MemoryBackend())
def make_queue_sqlite():
tmp = tempfile.NamedTemporaryFile(suffix=".db", delete=False)
tmp.close()
return TaskQueue(backend=SQLiteBackend(tmp.name))
BACKENDS = [make_queue_memory, make_queue_sqlite]
if os.environ.get("DATABASE_URL"):
def make_queue_postgres():
return TaskQueue(backend=PostgresBackend(os.environ["DATABASE_URL"]))
BACKENDS.append(make_queue_postgres)
@pytest.fixture(params=BACKENDS, ids=lambda f: f.__name__)
def queue(request):
return request.param()
def test_add_and_get(queue):
queue.add(Task(id="t1", job_name="demo"))
t = queue.get("t1")
assert t and t.id == "t1" and t.status == "pending"
def test_list_pending(queue):
queue.add(Task(id="a", job_name="demo"))
queue.add(Task(id="b", job_name="demo"))
queue.add(Task(id="c", job_name="other"))
pending = queue.list_pending("demo")
assert {t.id for t in pending} == {"a", "b"}
def test_claim_sets_metadata(queue):
queue.add(Task(id="x", job_name="demo"))
claimed = queue.claim("x", "worker-1")
assert claimed.status == "claimed"
assert claimed.claimed_by == "worker-1"
assert claimed.claimed_at is not None
assert claimed.updated_at is not None
def test_claim_twice_fails(queue):
queue.add(Task(id="y", job_name="demo"))
queue.claim("y", "w1")
with pytest.raises(ValueError):
queue.claim("y", "w2")
def test_claim_missing_fails(queue):
with pytest.raises(KeyError):
queue.claim("nope", "w1")
def test_status_transition_valid(queue):
queue.add(Task(id="s1", job_name="demo"))
queue.claim("s1", "w1")
done = queue.update_status("s1", "completed", reason="ok")
assert done.status == "completed"
assert done.reason == "ok"
def test_status_transition_invalid(queue):
queue.add(Task(id="s2", job_name="demo"))
with pytest.raises(ValueError):
queue.update_status("s2", "completed") # pending -> completed not allowed
def test_heartbeat_updates_claimed_at(queue):
queue.add(Task(id="h1", job_name="demo"))
claimed = queue.claim("h1", "w1")
first_ts = claimed.claimed_at
import time; time.sleep(0.01)
refreshed = queue.heartbeat("h1", "w1")
assert refreshed.claimed_at != first_ts
def test_heartbeat_wrong_worker(queue):
queue.add(Task(id="h2", job_name="demo"))
queue.claim("h2", "w1")
with pytest.raises(ValueError):
queue.heartbeat("h2", "w2")
def test_stats(queue):
queue.add(Task(id="g1", job_name="g"))
queue.add(Task(id="g2", job_name="g"))
queue.claim("g1", "w1")
stats = queue.get_stats("g")
assert stats["pending"] == 1
assert stats["claimed"] == 1
def test_claim_atomicity(queue):
"""Race 10 workers against 1 task. Exactly one must win."""
queue.add(Task(id="race", job_name="demo"))
winners: list[str] = []
lock = threading.Lock()
def worker(wid: str):
try:
queue.claim("race", wid)
with lock:
winners.append(wid)
except ValueError:
pass
threads = [threading.Thread(target=worker, args=(f"w{i}",)) for i in range(10)]
for t in threads: t.start()
for t in threads: t.join()
assert len(winners) == 1, f"Expected 1 winner, got {len(winners)}: {winners}"