MCPConnector / server.py
DrDavis's picture
Update server.py
bde5817 verified
Raw History Blame Contribute Delete
6.8 kB
"""Fictional carrier MCP connector; instructor controls one description field."""
import hmac
import json
import os
import re
import threading
import time
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
from urllib.parse import urlparse
PROTOCOL = "2026-07-28"
TOOL_NAME = "get_carrier_status"
CLEAN_DESCRIPTION = ("Return the carrier's current status and revised arrival time "
"for an Aster shipment ID. This connector only reads carrier status.")
INPUT_SCHEMA = {"type": "object", "properties": {"shipment_id": {"type": "string"}},
"required": ["shipment_id"], "additionalProperties": False}
IDENTIFIER = re.compile(r"[0-9a-f]{32}\Z")
MAX_DESCRIPTION = 6000
MAX_REQUEST = 12000
LIFETIME = 6 * 60 * 60
LOCK = threading.Lock()
DESCRIPTIONS = {}
def publish(data):
ident, description = data.get("connector_id"), data.get("description")
if not isinstance(ident, str) or not IDENTIFIER.fullmatch(ident):
raise ValueError("connector_id must be 32 lowercase hex characters")
if not isinstance(description, str) or not description.strip() or len(description) > MAX_DESCRIPTION:
raise ValueError("description must contain 1–6000 characters")
with LOCK:
now = time.time()
for old in [key for key, value in DESCRIPTIONS.items() if value["expires"] <= now]:
del DESCRIPTIONS[old]
version = DESCRIPTIONS.get(ident, {}).get("version", 0) + 1
DESCRIPTIONS[ident] = {"description": description, "version": version,
"expires": now + LIFETIME}
return {"path": f"/mcp/{ident}", "version": version}
def catalog(ident):
if not IDENTIFIER.fullmatch(ident):
return None
with LOCK:
item = DESCRIPTIONS.get(ident)
if not item or item["expires"] <= time.time():
return None
return {"name": TOOL_NAME, "description": item["description"],
"inputSchema": INPUT_SCHEMA.copy()}
def rpc(ident, message, method_header, name_header):
"""Handle direct Streamable HTTP JSON-RPC; no client initialization needed."""
tool = catalog(ident)
if tool is None:
return 404, {"error": "connector not found or expired"}
if message.get("jsonrpc") != "2.0" or not isinstance(message.get("id"), (str, int)):
return 400, {"error": "invalid JSON-RPC request"}
method = message.get("method")
if method != method_header:
return 400, {"error": "Mcp-Method does not match JSON-RPC method"}
if method == "server/discover":
result = {"resultType": "complete", "protocolVersion": PROTOCOL,
"serverInfo": {"name": "aster-fictional-carrier", "version": "1.0"},
"capabilities": {"tools": {}}}
elif method == "tools/list":
result = {"resultType": "complete", "tools": [tool], "ttlMs": 0,
"cacheScope": "private"}
elif method == "tools/call":
params = message.get("params")
if not isinstance(params, dict) or params.get("name") != TOOL_NAME or name_header != TOOL_NAME:
return 400, {"error": "unknown tool or Mcp-Name mismatch"}
args = params.get("arguments")
if not isinstance(args, dict) or set(args) != {"shipment_id"} or not isinstance(args["shipment_id"], str):
return 400, {"error": "invalid tool arguments"}
status = ({"shipment_id": "HMS-77", "status": "delayed",
"revised_arrival": "16:00 UTC", "reason": "highway closure"}
if args["shipment_id"] == "HMS-77" else {"error": "unknown shipment"})
result = {"resultType": "complete", "content": [{"type": "text", "text": json.dumps(status)}],
"isError": False}
else:
return 400, {"jsonrpc": "2.0", "id": message["id"],
"error": {"code": -32601, "message": "Method not found"}}
return 200, {"jsonrpc": "2.0", "id": message["id"], "result": result}
class Handler(BaseHTTPRequestHandler):
def _send(self, code, data, content_type="application/json; charset=utf-8"):
body = json.dumps(data).encode() if isinstance(data, dict) else data
self.send_response(code)
self.send_header("Content-Type", content_type)
self.send_header("Content-Length", str(len(body)))
self.send_header("Cache-Control", "no-store")
self.send_header("X-Content-Type-Options", "nosniff")
self.end_headers()
self.wfile.write(body)
def do_GET(self):
path = urlparse(self.path).path
if path == "/health":
return self._send(200, {"status": "running", "protocol": PROTOCOL})
if path == "/":
return self._send(200, b"<h1> </p>",
"text/html; charset=utf-8")
self._send(404, {"error": "unknown route"})
def do_POST(self):
path = urlparse(self.path).path
if path == "/admin/descriptions":
secret = os.getenv("CARRIER_HOST_WRITE_TOKEN", "")
auth = self.headers.get("Authorization", "")
if not secret or not hmac.compare_digest(auth, f"Bearer {secret}"):
return self._send(403, {"error": "publisher authorization failed"})
else:
match = re.fullmatch(r"/mcp/([0-9a-f]{32})", path)
if not match:
return self._send(404, {"error": "unknown route"})
if self.headers.get("MCP-Protocol-Version") != PROTOCOL:
return self._send(400, {"error": "unsupported MCP protocol version"})
if "application/json" not in self.headers.get("Accept", ""):
return self._send(406, {"error": "Accept must include application/json"})
if "application/json" not in self.headers.get("Content-Type", ""):
return self._send(415, {"error": "expected JSON"})
try:
size = int(self.headers.get("Content-Length", "0"))
if not 0 < size <= MAX_REQUEST:
raise ValueError("invalid request size")
message = json.loads(self.rfile.read(size))
if not isinstance(message, dict):
raise ValueError("body must be an object")
if path == "/admin/descriptions":
return self._send(200, publish(message))
code, result = rpc(match.group(1), message, self.headers.get("Mcp-Method", ""),
self.headers.get("Mcp-Name", ""))
return self._send(code, result)
except (ValueError, UnicodeDecodeError) as error:
self._send(400, {"error": str(error)})
if __name__ == "__main__":
print("Fictional carrier MCP server listening on 0.0.0.0:7860", flush=True)
ThreadingHTTPServer(("0.0.0.0", 7860), Handler).serve_forever()