Simplify Athena runtime and document current architecture

This commit is contained in:
Mikei386 committed 2026-08-30 08:45:52 +02:00
1 parent c721db47d0
commit c6517ee137
56 files changed
+1275 -4713

No files matched your search

-12
View File
@@ -1,12 +0,0 @@
FROM python:3.13-slim AS builder
COPY --from=ghcr.io/astral-sh/uv:0.11.7 /uv /uvx /bin/
RUN uv pip install --system --break-system-packages "arr-mcp[mcp]==1.0.1"
FROM python:3.13-slim
COPY --from=builder /usr/local /usr/local
RUN groupadd --system --gid 10001 mcp \
&& useradd --system --uid 10001 --gid 10001 --no-create-home mcp
USER 10001:10001
EXPOSE 8000
ENTRYPOINT ["arr-mcp"]
CMD ["--transport", "streamable-http", "--host", "0.0.0.0", "--port", "8000", "--auth-type", "none"]
-25
View File
@@ -1,25 +0,0 @@
FROM ghcr.io/github/github-mcp-server@sha256:1817b57d43916532dc002bdc5f344d639bd9fb54a9148d42168458f7c3280567 AS github
FROM python:3.13-slim@sha256:ffb752e139c0a19692a43af8d8523b274222dd68eebad5d583b45c2201c6e30a
ARG MCP_PROXY_VERSION=0.12.0
RUN pip install --no-cache-dir "mcp-proxy==${MCP_PROXY_VERSION}" "mcp==1.29.0"
# GitHub publishes a minimal image containing only the official Go binary.
# mcp-proxy contributes transport conversion only; GitHub API behavior and
# every exposed tool remain implemented by GitHub's official MCP server.
COPY --from=github /server/github-mcp-server /usr/local/bin/github-mcp-server
RUN useradd --system --uid 10001 --create-home --home-dir /app mcp
USER 10001:10001
WORKDIR /app
EXPOSE 8000
# This is the same OpenWebUI-compatible stateless transport used by Athena's
# other Python/stdio MCP adapters. The official GitHub binary remains the only
# component implementing GitHub operations.
# mcp-proxy intentionally starts stdio children with a minimal environment.
# Explicit pass-through is required so the GitHub subprocess receives the PAT
# already injected into this container by Docker. The value is never placed on
# the command line, image, logs or Open WebUI connection record.
ENTRYPOINT ["mcp-proxy", "--host", "0.0.0.0", "--port", "8000", "--stateless", "--pass-environment", "--"]
CMD ["/usr/local/bin/github-mcp-server", "stdio", "--read-only", "--tools", "search_repositories,get_file_contents,search_code"]
@@ -1,8 +0,0 @@
FROM nginx:1.29-alpine
COPY platform/mcp/homeassistant.conf.template /etc/nginx/templates/homeassistant.conf.template
COPY platform/mcp/ha-relay-entrypoint.sh /usr/local/bin/ha-relay-entrypoint
RUN chmod 0755 /usr/local/bin/ha-relay-entrypoint \
&& mkdir -p /tmp/client_temp /tmp/proxy_temp \
&& chown -R nginx:nginx /tmp/client_temp /tmp/proxy_temp
EXPOSE 8000
ENTRYPOINT ["/usr/local/bin/ha-relay-entrypoint"]
-8
View File
@@ -1,8 +0,0 @@
FROM ghcr.io/blakeem/navidrome-mcp:2.2.0@sha256:047f911a5a8f7cc8f185bb4d6e7ca6c435542edefff4694a00c2f718ab0ee7f5
# llama.cpp's tool-schema converter requires every JSON-Schema regex to be
# fully anchored. Upstream 2.2.0 omits the trailing `$` on exactly two radio
# URL fields. Fail the build if upstream changes instead of patching blindly.
USER root
RUN node -e 'const fs=require("node:fs"); const p="/app/dist/tools/handlers/radio-handlers.js"; let s=fs.readFileSync(p,"utf8"); const a="pattern: '\''^https?://.+'\''"; const b="pattern: '\''^https?://.+$'\''"; const n=s.split(a).length-1; if(n!==2) throw new Error(`expected 2 schema patterns, found ${n}`); fs.writeFileSync(p,s.split(a).join(b));'
USER node
-15
View File
@@ -1,15 +0,0 @@
FROM python:3.13-slim
ARG MCP_PROXY_VERSION=0.12.0
RUN apt-get update \
&& apt-get install -y --no-install-recommends openssh-client \
&& rm -rf /var/lib/apt/lists/* \
&& pip install --no-cache-dir "mcp-proxy==${MCP_PROXY_VERSION}" "mcp>=1.17,<2" \
&& useradd --system --uid 10001 --create-home --home-dir /app mcp
RUN touch /app/unraid_mcp.py && chown 10001:10001 /app/unraid_mcp.py
USER 10001:10001
WORKDIR /app
EXPOSE 8000
ENTRYPOINT ["mcp-proxy", "--host", "0.0.0.0", "--port", "8000", "--stateless", "--"]
CMD ["python", "/app/unraid_mcp.py"]
-18
View File
@@ -1,18 +0,0 @@
FROM python:3.13-slim
ARG MCP_PROXY_VERSION=0.12.0
ARG YT_DLP_VERSION=2026.7.4
RUN pip install --no-cache-dir \
"mcp-proxy==${MCP_PROXY_VERSION}" \
"mcp>=1.17,<2" \
"yt-dlp==${YT_DLP_VERSION}"
RUN useradd --system --uid 10001 --create-home --home-dir /app mcp
COPY web-search/web_search_mcp.py /app/web_search_mcp.py
RUN chown -R 10001:10001 /app
USER 10001:10001
WORKDIR /app
EXPOSE 8000
ENTRYPOINT ["mcp-proxy", "--host", "0.0.0.0", "--port", "8000", "--stateless", "--"]
CMD ["python", "/app/web_search_mcp.py"]
@@ -1,12 +0,0 @@
services:
mcp-github:
# Temporary, operator-controlled maintenance mode. It deliberately omits
# delete, merge, repository creation, issue mutation and workflow tools.
command:
- /usr/local/bin/github-mcp-server
- stdio
- --tools
- search_repositories,get_repository_tree,get_file_contents,search_code,list_branches,create_branch,create_or_update_file,push_files,create_pull_request
environment:
GITHUB_READ_ONLY: "0"
GITHUB_TOOLS: search_repositories,get_repository_tree,get_file_contents,search_code,list_branches,create_branch,create_or_update_file,push_files,create_pull_request
+3 -193
View File
@@ -12,147 +12,8 @@ x-tool-common: &tool-common
max-file: "3"
services:
mcp-web:
<<: *tool-common
# Historical site-specific facade. Kept only for rollback while the
# default portable endpoint points directly at TinySearch's broad MCP.
profiles: [legacy-web]
build:
context: ..
dockerfile: mcp/Dockerfile.web
image: mike-ai/mcp-web:local
container_name: mike-ai-mcp-web
# The relay fetches and validates public result pages itself. It therefore
# needs both the private tool network and the explicitly separated egress
# network; keeping it on `tools` only makes search discovery work while
# every page fetch fails.
networks: [tools, tools-egress]
dns: ["${AI_DNS:-1.1.1.1}"]
environment:
TINYSEARCH_MCP_URL: http://tinysearch:8000/mcp
SEARXNG_URL: http://searxng:8080
WEB_SEARCH_BUDGET_MAX_RELATED: "3"
YOUTUBE_TIMEOUT: "45"
depends_on:
tinysearch:
condition: service_started
healthcheck:
test: ["CMD", "python", "-c", "import socket; s=socket.create_connection(('127.0.0.1', 8000), 2); s.close()"]
interval: 30s
timeout: 5s
retries: 5
start_period: 15s
searxng:
<<: *tool-common
image: searxng/searxng@sha256:e45d5894bfaa0bf8773b9f283795ae57f1c15ddb29c8cecb70b3665b0ce9ec60
container_name: mike-ai-tools-searxng
dns: ["${AI_DNS:-1.1.1.1}"]
volumes:
- ${SEARXNG_SETTINGS_FILE:-../web-search/searxng-settings.example.yml}:/etc/searxng/settings.yml:ro
networks: [tools, tools-egress]
tinysearch:
<<: *tool-common
# TinySearch v0.6.1. This release fixes the crawler behavior observed with
# v0.5.1 and adds the current search -> scrape_urls workflow.
image: marcellm01/tinysearch@sha256:7a7d0585f5000f462e699e42b97409715826a9e2edcd09a166afa93a4b7cba31
container_name: mike-ai-tools-tinysearch
dns: ["${AI_DNS:-1.1.1.1}"]
# Crawl4AI keeps transient browser/session state here. The container stays
# read-only; only this disposable runtime directory (and /tmp from the
# common hardening block) is writable.
tmpfs:
- /tmp:rw,noexec,nosuid,nodev,size=64m
- /home/tinysearch/.crawl4ai:rw,nosuid,nodev,size=256m,mode=1777
shm_size: 1gb
volumes:
- tinysearch-models:/data/models
- ../web-search/tinysearch_config.json:/config/tinysearch_config.json:ro
environment:
MCP_TRANSPORT: streamable-http
MCP_HOST: 0.0.0.0
MCP_PORT: "8000"
TINYSEARCH_CONFIG_PATH: /config/tinysearch_config.json
TINYSEARCH_SEARCH_BACKEND: searxng
SEARXNG_URL: http://searxng:8080/search
depends_on: [searxng]
cap_add: [SETUID, SETGID, CHOWN]
networks: [tools, tools-egress]
# The image's built-in `tinysearch doctor` also requires a writable
# configuration directory, although normal server operation does not.
# Check the service socket instead so read-only hardening remains intact.
healthcheck:
test: ["CMD", "python", "-c", "import socket; s=socket.create_connection(('127.0.0.1', 8000), 2); s.close()"]
interval: 30s
timeout: 5s
retries: 5
start_period: 20s
mcp-homeassistant:
<<: *tool-common
build:
context: ../..
dockerfile: platform/mcp/Dockerfile.homeassistant-relay
image: mike-ai/mcp-homeassistant-relay:local
container_name: mike-ai-mcp-homeassistant
profiles: [homeassistant]
# Keep the public TLS hostname for SNI/certificate validation, but route it
# to the private reverse proxy through WireGuard. Public DNS may otherwise
# resolve to the Fritzbox WAN address, which is unreachable/hairpinned from
# Athena's remote-site containers.
extra_hosts:
- "ha.casaderoll.de:${HOME_LAN_PROXY_IP:-192.168.1.2}"
volumes:
- ${HA_ENV_FILE:-/etc/mike-ai/homeassistant-admin-mcp.env}:/run/secrets/homeassistant.env:ro
cap_add: [CHOWN, SETUID, SETGID]
networks: [tools, tools-egress]
mcp-arr:
<<: *tool-common
build:
context: .
dockerfile: Dockerfile.arr
image: mike-ai/mcp-arr:1.0.1-patched
container_name: mike-ai-mcp-arr
profiles: [arr]
env_file:
- ${ARR_ENV_FILE:-/etc/mike-ai/arr-mcp.env}
volumes:
# The local fork adds bounded read-only Sonarr pseudo-actions. Keep the
# patch explicit until upstream publishes a self-contained 2.x image.
- ${ARR_SONARR_PATCH:-./patches/mcp_sonarr.py}:/usr/local/lib/python3.13/site-packages/arr_mcp/mcp/mcp_sonarr.py:ro
# Upstream's generic "Execute any Radarr API action" text gives small
# models no routing boundary. This overlay changes guidance only.
- ${ARR_RADARR_PATCH:-./patches/mcp_radarr.py}:/usr/local/lib/python3.13/site-packages/arr_mcp/mcp/mcp_radarr.py:ro
networks: [tools, tools-egress]
mcp-navidrome:
<<: *tool-common
# Version and amd64 manifest are pinned. The image contains no mpv, so it
# cannot play audio on the headless AI host and does not expose playback
# controls. It talks to Navidrome only through its authenticated API.
build:
context: .
dockerfile: Dockerfile.navidrome
image: mike-ai/mcp-navidrome:2.2.0-schemafix1
container_name: mike-ai-mcp-navidrome
profiles: [navidrome]
env_file:
- ${NAVIDROME_MCP_ENV_FILE:-/etc/mike-ai/navidrome-mcp.env}
environment:
MCP_TRANSPORT: http
MCP_HTTP_EXPOSE: "true"
MCP_HTTP_PORT: "3000"
# OpenWebUI uses the Docker name; Pi/Hermes may reach the same endpoint
# directly through Athena's WireGuard address and VPN port 8207.
MCP_HTTP_ALLOWED_HOSTS: "mike-ai-mcp-navidrome:3000,mike-ai-mcp-navidrome,${VPN_SERVICE_IP:-192.168.1.212}:8207,${VPN_SERVICE_IP:-192.168.1.212}"
WEBUI_ENABLED: "false"
tmpfs:
- /tmp:rw,noexec,nosuid,nodev,size=64m
- /config:rw,noexec,nosuid,nodev,size=4m,mode=0700
networks: [tools, tools-egress]
# Einziger MCP auf Athena: die schmale Fassade zum root-eigenen Operator.
# Alle portablen Fach-MCPs laufen als eigene Container auf Unraid.
mcp-athena-operator:
<<: *tool-common
build:
@@ -163,61 +24,10 @@ services:
environment:
ATHENA_OPERATOR_SOCKET: /operator/operator.sock
volumes:
# The unprivileged MCP facade sees only the root-owned executor socket.
# Docker, source, models, Git credentials and host paths remain on the
# executor side and are reachable only through structured operations.
- /run/mike-ai-operator:/operator:ro
networks: [tools]
healthcheck:
test: ["CMD", "python", "-c", "import socket; s=socket.create_connection(('127.0.0.1',8000),2); s.close()"]
test: [CMD, python, -c, "import socket; s=socket.create_connection(('127.0.0.1',8000),2); s.close()"]
interval: 30s
timeout: 5s
retries: 5
start_period: 10s
mcp-github:
<<: *tool-common
build:
context: .
dockerfile: Dockerfile.github
image: mike-ai/mcp-github:github-v1.10.1-mcp-proxy-v0.12.0
container_name: mike-ai-mcp-github
profiles: [github]
env_file:
- ${GITHUB_MCP_ENV_FILE:-/etc/mike-ai/github-mcp.env}
environment:
# These server-side limits remain authoritative even if a client asks
# for broader toolsets. The token itself must also remain read-only.
GITHUB_TOOLS: search_repositories,get_file_contents,search_code
GITHUB_READ_ONLY: "1"
networks: [tools, tools-egress]
healthcheck:
test: ["CMD", "python", "-c", "import socket; s=socket.create_connection(('127.0.0.1',8000),2); s.close()"]
interval: 30s
timeout: 5s
retries: 5
start_period: 15s
mcp-unraid-ssh:
<<: *tool-common
profiles: [extended]
build:
context: .
dockerfile: Dockerfile.unraid-ssh
image: mike-ai/mcp-unraid-ssh:local
container_name: mike-ai-mcp-unraid-ssh
environment:
UNRAID_MCP_CONFIG: /run/config/unraid-mcp.json
volumes:
- ${UNRAID_MCP_SOURCE:-/opt/mike-ai/unraid-agent/unraid_mcp.py}:/app/unraid_mcp.py:ro
- ${UNRAID_MCP_CONFIG:-/etc/mike-ai/unraid-mcp.json}:/run/config/unraid-mcp.json:ro
- ${UNRAID_SSH_KEY:-/etc/mike-ai/keys/unraid_root}:/etc/mike-ai/keys/unraid_root:ro
- ${UNRAID_KNOWN_HOSTS:-/etc/mike-ai/ssh/known_hosts_unraid_ai}:/etc/mike-ai/ssh/known_hosts_unraid_ai:ro
- unraid-audit:/var/log/mike-ai
networks: [tools, tools-egress]
volumes:
tinysearch-models:
name: mike-ai-tools_tinysearch-models
external: true
unraid-audit:
-49
View File
@@ -1,49 +0,0 @@
#!/usr/bin/env bash
set -Eeuo pipefail
MODE=${1:-}
CONFIRM=${2:-}
ROOT_DIR=$(cd "$(dirname "${BASH_SOURCE[0]}")" && pwd)
PROJECT=mike-ai-tools
BASE=(docker compose -p "$PROJECT" --profile github -f "$ROOT_DIR/compose.yaml")
MAINTENANCE=(-f "$ROOT_DIR/compose.github-maintenance.yaml")
die() { printf 'FEHLER: %s\n' "$*" >&2; exit 1; }
[[ $EUID -eq 0 ]] || die "Bitte als root ausführen."
case "$MODE" in
read)
"${BASE[@]}" up -d --no-deps --force-recreate mcp-github
;;
maintenance)
[[ $CONFIRM == --confirm ]] || die \
"Wartungsmodus nur mit: $0 maintenance --confirm"
"${BASE[@]}" "${MAINTENANCE[@]}" up -d --no-deps --force-recreate mcp-github
;;
status)
command_line=$(docker inspect -f '{{json .Config.Cmd}}' mike-ai-mcp-github 2>/dev/null || true)
if [[ $command_line == *create_branch* ]]; then
printf 'GitHub MCP: WARTUNGSMODUS (begrenzter Schreibzugriff)\n'
elif [[ $command_line == *--read-only* ]]; then
printf 'GitHub MCP: NUR LESEN\n'
else
die "GitHub-MCP-Modus ist nicht eindeutig; Konfiguration prüfen."
fi
exit 0
;;
*)
die "Aufruf: $0 {status|read|maintenance --confirm}"
;;
esac
# Open WebUI caches MCP capabilities. A short backend restart makes the new
# allowlist deterministic for all profiles without touching model services.
docker restart mike-ai-open-webui >/dev/null
for _ in $(seq 1 30); do
[[ $(docker inspect -f '{{.State.Health.Status}}' mike-ai-mcp-github 2>/dev/null || true) == healthy ]] && break
sleep 1
done
[[ $(docker inspect -f '{{.State.Health.Status}}' mike-ai-mcp-github 2>/dev/null || true) == healthy ]] || \
die "GitHub MCP wurde nicht gesund. Zurücksetzen mit: $0 read"
"$0" status
-23
View File
@@ -1,23 +0,0 @@
#!/bin/sh
set -eu
config=/run/secrets/homeassistant.env
if [ ! -r "$config" ]; then
echo "Home Assistant secret file is missing" >&2
exit 1
fi
set -a
. "$config"
set +a
: "${HASS_URL:?HASS_URL is required}"
: "${HASS_TOKEN:?HASS_TOKEN is required}"
upstream=${HASS_URL%/}
escaped_token=$(printf '%s' "$HASS_TOKEN" | sed 's/[&/]/\\&/g')
escaped_upstream=$(printf '%s' "$upstream" | sed 's/[&/]/\\&/g')
sed -e "s/__HASS_TOKEN__/$escaped_token/g" \
-e "s/__HASS_UPSTREAM__/$escaped_upstream/g" \
/etc/nginx/templates/homeassistant.conf.template \
> /tmp/nginx.conf
unset HASS_TOKEN
exec nginx -c /tmp/nginx.conf -g 'daemon off;'
-30
View File
@@ -1,30 +0,0 @@
worker_processes 1;
pid /tmp/nginx.pid;
error_log /dev/stderr warn;
events { worker_connections 128; }
http {
access_log /dev/stdout;
client_body_temp_path /tmp/client_temp;
proxy_temp_path /tmp/proxy_temp;
fastcgi_temp_path /tmp/fastcgi_temp;
uwsgi_temp_path /tmp/uwsgi_temp;
scgi_temp_path /tmp/scgi_temp;
proxy_buffering off;
proxy_read_timeout 600s;
proxy_send_timeout 600s;
server {
listen 8000;
location /mcp {
proxy_pass __HASS_UPSTREAM__/api/hass_mcp;
proxy_http_version 1.1;
proxy_ssl_server_name on;
proxy_ssl_name $proxy_host;
proxy_set_header Authorization "Bearer __HASS_TOKEN__";
proxy_set_header Host $proxy_host;
proxy_set_header Connection "";
}
}
}
@@ -1,27 +0,0 @@
#!/usr/bin/env bash
set -Eeuo pipefail
umask 077
repo_dir=$(cd "$(dirname "${BASH_SOURCE[0]}")/../.." && pwd)
source_file=$repo_dir/platform/mcp/patches/hass_mcp/yaml_config.py
ha_config_dir=${1:-}
die() { printf 'FEHLER: %s\n' "$*" >&2; exit 1; }
[[ $EUID -eq 0 ]] || die "Bitte als root auf dem Home-Assistant-Host ausführen."
[[ -n $ha_config_dir ]] || die "Aufruf: $0 /pfad/zum/home-assistant-config"
[[ -s $source_file ]] || die "Patchdatei fehlt: $source_file"
target=$ha_config_dir/custom_components/hass_mcp/tools/yaml_config.py
[[ -s $target ]] || die "Native hass_mcp-Installation nicht gefunden: $target"
backup_dir=$ha_config_dir/.hass_mcp_patch_backups
mkdir -p "$backup_dir"
chmod 700 "$backup_dir"
stamp=$(date +%Y%m%d-%H%M%S)
backup=$backup_dir/yaml_config.py.before-guard-$stamp
cp -p "$target" "$backup"
chmod 600 "$backup"
install -m 0644 "$source_file" "$target"
printf 'YAML Guard installiert. Sicherung: %s\n' "$backup"
printf 'Home Assistant muss jetzt kontrolliert neu gestartet werden.\n'
+3 -7
View File
@@ -16,14 +16,10 @@ docker network inspect mike-ai-tools >/dev/null 2>&1 || \
docker network inspect mike-ai-tools-egress >/dev/null 2>&1 || \
docker network create --subnet 172.30.50.0/24 mike-ai-tools-egress >/dev/null
# One administrative MCP exposes ATHENA.md plus the bounded host operator.
# Einziger MCP auf Athena: ATHENA.md plus der begrenzte Host-Operator.
"$MCP_DIR/../operator/install-operator.sh"
services=(mcp-athena-operator)
# Portable Fach-MCPs laufen zentral im MCPHub auf Unraid. Ihre Compose-Blöcke
# bleiben vorläufig als explizite Rollback-Profile erhalten, werden bei einer
# normalen Athena-Installation aber weder gebaut noch gestartet.
echo "ARR, GitHub, Home Assistant und Navidrome werden über MCPHub auf Unraid bereitgestellt."
"${COMPOSE[@]}" up -d --build "${services[@]}"
echo "Portable Fach-MCPs laufen als eigene Container auf Unraid."
"${COMPOSE[@]}" up -d --build mcp-athena-operator
"${COMPOSE[@]}" ps
@@ -1,707 +0,0 @@
"""Guarded YAML access for Home Assistant configuration files.
The language model is untrusted. Reads are bounded and redact likely inline
credentials. Every mutation is previewed first, bound to the exact current file
hash, backed up, written atomically, checked by Home Assistant and rolled back
when validation or reload fails. ``secrets.yaml`` is never addressable.
"""
from __future__ import annotations
import copy
import difflib
import hashlib
import json
import os
import re
import secrets
import time
from pathlib import Path
from typing import Any
from homeassistant.core import HomeAssistant
from homeassistant.util import slugify
from ..identity import user_context
from ..protocol import ToolError, internal_error
from ..registry import LIMIT_FIELD, OFFSET_FIELD, paginate, schema, tool
# kind -> (filename, parsed structure, reload service domain)
_KINDS: dict[str, tuple[str, str, str | None]] = {
"automation": ("automations.yaml", "list", "automation"),
"script": ("scripts.yaml", "dict", "script"),
"scene": ("scenes.yaml", "list", "scene"),
# configuration.yaml is intentionally raw-read/replace only. Treating its
# top-level keys as CRUD records would be dangerously misleading.
"configuration": ("configuration.yaml", "dict", None),
}
_STRUCTURED_KINDS = frozenset({"automation", "script", "scene"})
_OPS = (
"list",
"get",
"read_source",
"find_source",
"find_commented_blocks",
"list_backups",
"create",
"update",
"delete",
"replace_source_text",
"restore_backup",
"reload",
)
_MUTATING_OPS = frozenset({"create", "update", "delete", "replace_source_text", "restore_backup"})
_TICKET_TTL_SECONDS = 600
_MAX_SOURCE_BYTES = 2 * 1024 * 1024
_MAX_REPLACEMENT_CHARS = 50_000
_PENDING: dict[str, dict[str, Any]] = {}
_SENSITIVE_LINE = re.compile(
r"(?i)^(?P<prefix>\s*[^#\n]*(?:password|passwd|token|secret|api[_-]?key|authorization)[^:]*:\s*).*$"
)
_SECRET_REFERENCE = re.compile(r"!secret\s+[^\s#]+", re.IGNORECASE)
@tool(
name="ha_yaml_config",
description=(
"Safely inspect and edit Home Assistant YAML. Structured CRUD is limited to "
"automations.yaml, scripts.yaml and scenes.yaml. Raw source operations also "
"allow configuration.yaml so commented-out blocks can be found and reviewed. "
"secrets.yaml and arbitrary paths are impossible. Use find_commented_blocks once "
"to inventory fully commented YAML entries; use read_source/find_source only for "
"other comments or exact YAML text. Every mutation first returns a diff/preview "
"and one-time approval_ticket; only repeat the exact unchanged call with "
"confirm=true after explicit user approval. Writes create a backup, are atomic, "
"run Home Assistant config validation, roll back on failure, and reload the "
"affected domain when supported. Never claim a preview changed Home Assistant."
),
input_schema=schema(
properties={
"kind": {"type": "string", "enum": list(_KINDS)},
"op": {"type": "string", "enum": list(_OPS)},
"id": {
"type": "string",
"description": "Entry id for structured get/create/update/delete.",
},
"config": {
"type": "object",
"additionalProperties": True,
"description": "Complete entry config for structured create/update.",
},
"query": {
"type": "string",
"maxLength": 500,
"description": "Case-insensitive literal text for find_source.",
},
"start_line": {
"type": "integer",
"minimum": 1,
"default": 1,
"description": "First 1-based line returned by read_source.",
},
"max_lines": {
"type": "integer",
"minimum": 1,
"maximum": 400,
"default": 120,
"description": "Bounded source lines returned by read_source/find_source.",
},
"old_text": {
"type": "string",
"minLength": 1,
"maxLength": _MAX_REPLACEMENT_CHARS,
"description": "Exact unique YAML source text to replace.",
},
"new_text": {
"type": "string",
"maxLength": _MAX_REPLACEMENT_CHARS,
"description": "Replacement YAML source text; may be empty to remove a block.",
},
"backup_id": {
"type": "string",
"pattern": r"^[a-z0-9_.-]+$",
"description": "Opaque filename returned by list_backups.",
},
"confirm": {
"type": "boolean",
"default": False,
"description": "True only after the user approved the exact preview.",
},
"approval_ticket": {
"type": "string",
"description": "One-time ticket from the unchanged mutation preview.",
},
"limit": LIMIT_FIELD,
"offset": OFFSET_FIELD,
},
required=["kind", "op"],
),
read_only=False,
idempotent=False,
requires_admin=True,
write_ops=["create", "update", "replace_source_text", "reload"],
destructive_ops=["delete", "restore_backup"],
admin_ops=["list", "get", "read_source", "find_source", "find_commented_blocks", "list_backups"],
)
async def ha_yaml_config(
hass: HomeAssistant,
kind: str,
op: str,
id: str | None = None,
config: dict[str, Any] | None = None,
query: str | None = None,
start_line: int = 1,
max_lines: int = 120,
old_text: str | None = None,
new_text: str | None = None,
backup_id: str | None = None,
confirm: bool = False,
approval_ticket: str | None = None,
limit: int = 100,
offset: int = 0,
) -> dict[str, Any]:
if kind not in _KINDS:
raise ToolError(f"unknown kind '{kind}'")
if op not in _OPS:
raise ToolError(f"unknown op '{op}'")
filename, structure, reload_domain = _KINDS[kind]
path = Path(hass.config.path(filename))
if op in {"list", "get", "create", "update", "delete", "reload"} and kind not in _STRUCTURED_KINDS:
raise ToolError(
"configuration.yaml supports only read_source, find_source, list_backups, "
"replace_source_text and restore_backup"
)
if op == "read_source":
return await _read_source(hass, path, start_line, max_lines)
if op == "find_source":
if not query:
raise ToolError("op=find_source requires query")
return await _find_source(hass, path, query, max_lines)
if op == "find_commented_blocks":
if kind not in {"automation", "scene"}:
raise ToolError("find_commented_blocks supports automation and scene list files")
return await _find_commented_blocks(hass, path, max_lines)
if op == "list_backups":
return await _list_backups(hass, filename, limit, offset)
if op == "replace_source_text":
if old_text is None or new_text is None:
raise ToolError("op=replace_source_text requires old_text and new_text")
_reject_sensitive_replacement(old_text, new_text)
current = await _read_text(hass, path)
if current.count(old_text) != 1:
raise ToolError(
f"old_text must occur exactly once in {filename}; found {current.count(old_text)} occurrences"
)
proposed = current.replace(old_text, new_text, 1)
await _validate_yaml_text(hass, proposed, structure, filename)
change = _change_record(kind, op, current, {"old_text": old_text, "new_text": new_text})
if not confirm:
return _preview(change, _source_diff(filename, current, proposed))
_consume_ticket(change, approval_ticket)
return await _commit_text(hass, path, filename, structure, reload_domain, current, proposed)
if op == "restore_backup":
if not backup_id:
raise ToolError("op=restore_backup requires backup_id from list_backups")
current = await _read_text(hass, path)
restored = await _read_backup(hass, filename, backup_id)
await _validate_yaml_text(hass, restored, structure, filename)
change = _change_record(kind, op, current, {"backup_id": backup_id})
if not confirm:
return _preview(change, _source_diff(filename, current, restored), extra={"backup_id": backup_id})
_consume_ticket(change, approval_ticket)
return await _commit_text(hass, path, filename, structure, reload_domain, current, restored)
data = await _load(hass, path, structure)
if op == "list":
return paginate(_to_list(data, structure), limit, offset)
if op == "get":
if not id:
raise ToolError("op=get requires id")
item = _find(data, structure, id)
if item is None:
raise ToolError(f"{kind} '{id}' not found in {filename}")
return item
if op == "reload":
await _reload(hass, reload_domain)
return {"reloaded": reload_domain, "changed_file": False}
current_text = await _read_text(hass, path)
proposed_data = _copy_data(data)
result: dict[str, Any]
if op == "create":
if not config:
raise ToolError("op=create requires config")
# The generated id must be deterministic so the exact preview can be
# confirmed in a second call without silently proposing another entry.
generated_id = f"mcp_{hashlib.sha256(json.dumps(config, sort_keys=True).encode()).hexdigest()[:16]}"
new_id = id or config.get("id") or generated_id
if _find(proposed_data, structure, new_id) is not None:
raise ToolError(f"{kind} '{new_id}' already exists")
if structure == "list":
proposed_data.append({"id": new_id, **{k: v for k, v in config.items() if k != "id"}})
else:
proposed_data[new_id] = config
result = {
"operation": "create",
"id": new_id,
"proposed_entry": _find(proposed_data, structure, new_id),
}
elif op == "update":
if not id or not config:
raise ToolError("op=update requires id and config")
before = _find(proposed_data, structure, id)
if before is None or not _replace(proposed_data, structure, id, config):
raise ToolError(f"{kind} '{id}' not found")
result = {"operation": "update", "id": id, "current_entry": before, "proposed_entry": _find(proposed_data, structure, id)}
elif op == "delete":
if not id:
raise ToolError("op=delete requires id")
before = _find(proposed_data, structure, id)
if before is None or not _remove(proposed_data, structure, id):
raise ToolError(f"{kind} '{id}' not found")
result = {"operation": "delete", "id": id, "current_entry": before}
else:
raise ToolError(f"unsupported op '{op}'")
change = _change_record(kind, op, current_text, {"id": id, "config": config, "result": result})
if not confirm:
return _preview(change, extra=result)
_consume_ticket(change, approval_ticket)
committed = await _commit_data(hass, path, filename, structure, reload_domain, current_text, proposed_data)
return {**result, **committed}
def _copy_data(data: Any) -> Any:
return copy.deepcopy(data)
def _fingerprint(text: str) -> str:
return hashlib.sha256(text.encode("utf-8")).hexdigest()
def _change_record(kind: str, op: str, current: str, arguments: dict[str, Any]) -> dict[str, Any]:
return {
"kind": kind,
"op": op,
"current_sha256": _fingerprint(current),
"arguments": arguments,
}
def _new_ticket(change: dict[str, Any]) -> str:
now = time.time()
for key, value in list(_PENDING.items()):
if value["expires_at"] <= now:
_PENDING.pop(key, None)
ticket = secrets.token_urlsafe(18)
_PENDING[ticket] = {
"fingerprint": _fingerprint(json.dumps(change, sort_keys=True, separators=(",", ":"))),
"expires_at": now + _TICKET_TTL_SECONDS,
}
return ticket
def _consume_ticket(change: dict[str, Any], ticket: str | None) -> None:
record = _PENDING.pop(ticket, None) if ticket else None
expected = _fingerprint(json.dumps(change, sort_keys=True, separators=(",", ":")))
if not record or record["expires_at"] <= time.time() or record["fingerprint"] != expected:
raise ToolError(
"approval_ticket is missing, expired, already used, or does not match the exact "
"change/current file. Run the same operation without confirm, show the preview, "
"then repeat unchanged with confirm=true only after explicit user approval."
)
def _preview(change: dict[str, Any], diff: list[str] | None = None, extra: dict[str, Any] | None = None) -> dict[str, Any]:
return {
"changed": False,
"confirmation_required": True,
"approval_ticket": _new_ticket(change),
"ticket_expires_in_seconds": _TICKET_TTL_SECONDS,
"current_sha256": change["current_sha256"],
**(extra or {}),
**({"diff": diff, "diff_truncated": len(diff) >= 120} if diff is not None else {}),
"model_instruction": (
"This is a preview only. Show it to the user and stop. Do not claim anything was "
"changed. After explicit approval repeat the exact call with confirm=true and approval_ticket."
),
}
def _redact_line(line: str) -> str:
match = _SENSITIVE_LINE.match(line)
if match:
return f"{match.group('prefix')}<redacted>"
return _SECRET_REFERENCE.sub("!secret <redacted-reference>", line)
def _reject_sensitive_replacement(*values: str) -> None:
for value in values:
if any(_SENSITIVE_LINE.match(line) for line in value.splitlines()) or _SECRET_REFERENCE.search(value):
raise ToolError(
"Raw replacement containing credential-like keys or !secret references is refused. "
"Edit that material locally outside the LLM context."
)
async def _read_text(hass: HomeAssistant, path: Path) -> str:
def _read() -> str:
if not path.exists():
return ""
if path.stat().st_size > _MAX_SOURCE_BYTES:
raise ToolError(f"{path.name} exceeds the {_MAX_SOURCE_BYTES} byte safety limit")
return path.read_text(encoding="utf-8")
return await hass.async_add_executor_job(_read)
async def _read_source(hass: HomeAssistant, path: Path, start_line: int, max_lines: int) -> dict[str, Any]:
text = await _read_text(hass, path)
lines = text.splitlines()
start = max(1, start_line)
count = max(1, min(max_lines, 400))
selected = lines[start - 1 : start - 1 + count]
return {
"file": path.name,
"sha256": _fingerprint(text),
"total_lines": len(lines),
"start_line": start,
"returned_lines": len(selected),
"has_more": start - 1 + len(selected) < len(lines),
"lines": [{"line": start + index, "text": _redact_line(line)} for index, line in enumerate(selected)],
"redaction_note": "Credential-like values and !secret reference names are redacted.",
}
async def _find_source(hass: HomeAssistant, path: Path, query: str, max_lines: int) -> dict[str, Any]:
text = await _read_text(hass, path)
lines = text.splitlines()
hits = [index for index, line in enumerate(lines) if query.casefold() in line.casefold()]
cap = max(1, min(max_lines, 400))
selected = hits[:cap]
return {
"file": path.name,
"sha256": _fingerprint(text),
"authoritative_match_count": len(hits),
"returned_count": len(selected),
"has_more": len(hits) > len(selected),
"matches": [{"line": index + 1, "text": _redact_line(lines[index])} for index in selected],
"redaction_note": "Credential-like values and !secret reference names are redacted.",
}
def _extract_commented_blocks(text: str, max_lines: int) -> tuple[list[dict[str, Any]], bool]:
"""Return top-level YAML list entries whose every source line is commented.
This intentionally recognizes only the conservative ``# - id:`` form used
by Home Assistant's automations/scenes editor. Ordinary prose comments,
partially disabled entries and nested comments are not treated as entries.
"""
lines = text.splitlines()
start_pattern = re.compile(r"^\s*#\s*-\s+id\s*:\s*(.*?)\s*$", re.IGNORECASE)
alias_pattern = re.compile(r"^\s*#\s+alias\s*:\s*(.*?)\s*$", re.IGNORECASE)
blocks: list[dict[str, Any]] = []
consumed = 0
index = 0
truncated = False
def clean_scalar(value: str) -> str:
value = value.strip()
if len(value) >= 2 and value[0] == value[-1] and value[0] in {"'", '"'}:
return value[1:-1]
return value
while index < len(lines):
match = start_pattern.match(lines[index])
if not match:
index += 1
continue
start = index
block_lines = [lines[index]]
index += 1
while index < len(lines):
if start_pattern.match(lines[index]):
break
if not lines[index].strip() or not re.match(r"^\s*#", lines[index]):
break
block_lines.append(lines[index])
index += 1
if consumed + len(block_lines) > max_lines:
truncated = True
break
alias = None
for line in block_lines:
alias_match = alias_pattern.match(line)
if alias_match:
alias = clean_scalar(alias_match.group(1))
break
blocks.append(
{
"start_line": start + 1,
"end_line": start + len(block_lines),
"id": clean_scalar(match.group(1)),
"alias": alias,
"source": [
{"line": start + offset + 1, "text": _redact_line(line)}
for offset, line in enumerate(block_lines)
],
}
)
consumed += len(block_lines)
return blocks, truncated
async def _find_commented_blocks(
hass: HomeAssistant, path: Path, max_lines: int
) -> dict[str, Any]:
text = await _read_text(hass, path)
cap = max(1, min(max_lines, 400))
blocks, truncated = _extract_commented_blocks(text, cap)
return {
"file": path.name,
"sha256": _fingerprint(text),
"authoritative_block_count": len(blocks) if not truncated else None,
"returned_block_count": len(blocks),
"returned_source_lines": sum(len(block["source"]) for block in blocks),
"has_more": truncated,
"blocks": blocks,
"recognition_rule": "Only fully commented top-level '# - id:' YAML list entries are returned.",
"redaction_note": "Credential-like values and !secret reference names are redacted.",
}
def _source_diff(filename: str, before: str, after: str) -> list[str]:
return list(
difflib.unified_diff(
before.splitlines(),
after.splitlines(),
fromfile=f"{filename}:before",
tofile=f"{filename}:after",
lineterm="",
n=3,
)
)[:120]
def _backup_dir(hass: HomeAssistant) -> Path:
return Path(hass.config.path(".hass_mcp_backups", "yaml"))
async def _create_backup(hass: HomeAssistant, filename: str, content: str) -> str:
backup_id = f"{filename}.{time.strftime('%Y%m%d-%H%M%S')}.{_fingerprint(content)[:10]}.bak"
directory = _backup_dir(hass)
def _write() -> None:
directory.mkdir(mode=0o700, parents=True, exist_ok=True)
target = directory / backup_id
target.write_text(content, encoding="utf-8")
target.chmod(0o600)
await hass.async_add_executor_job(_write)
return backup_id
async def _list_backups(hass: HomeAssistant, filename: str, limit: int, offset: int) -> dict[str, Any]:
directory = _backup_dir(hass)
def _list() -> list[dict[str, Any]]:
if not directory.exists():
return []
rows = []
for path in directory.glob(f"{filename}.*.bak"):
stat = path.stat()
rows.append({"backup_id": path.name, "size": stat.st_size, "created_unix": int(stat.st_mtime)})
return sorted(rows, key=lambda row: row["created_unix"], reverse=True)
return paginate(await hass.async_add_executor_job(_list), limit, offset)
async def _read_backup(hass: HomeAssistant, filename: str, backup_id: str) -> str:
if Path(backup_id).name != backup_id or not backup_id.startswith(f"{filename}.") or not backup_id.endswith(".bak"):
raise ToolError("backup_id is not valid for this YAML kind")
path = _backup_dir(hass) / backup_id
def _read() -> str:
if not path.is_file():
raise ToolError("backup_id not found")
if path.stat().st_size > _MAX_SOURCE_BYTES:
raise ToolError("backup exceeds safety limit")
return path.read_text(encoding="utf-8")
return await hass.async_add_executor_job(_read)
async def _validate_yaml_text(hass: HomeAssistant, content: str, structure: str, filename: str) -> Any:
from homeassistant.util.yaml import parse_yaml
def _parse() -> Any:
parsed = parse_yaml(content) if content.strip() else ([] if structure == "list" else {})
if structure == "list" and not isinstance(parsed, list):
raise ToolError(f"{filename} must be a YAML list, got {type(parsed).__name__}")
if structure == "dict" and not isinstance(parsed, dict):
raise ToolError(f"{filename} must be a YAML mapping, got {type(parsed).__name__}")
return parsed
return await hass.async_add_executor_job(_parse)
async def _check_full_config(hass: HomeAssistant) -> dict[str, Any]:
try:
from homeassistant.components.config.core import async_check_ha_config_file
except ImportError:
from homeassistant.config import async_check_ha_config_file
result = await async_check_ha_config_file(hass)
if result is None:
return {"valid": True}
if isinstance(result, str):
return {"valid": not bool(result), "error": result or None}
errors = getattr(result, "errors", None)
if errors:
return {"valid": False, "error": str(errors)}
return {"valid": True, "result": str(result)}
async def _atomic_write(hass: HomeAssistant, path: Path, content: str) -> None:
def _write() -> None:
temporary = path.with_name(f".{path.name}.hass-mcp-{secrets.token_hex(6)}.tmp")
try:
with temporary.open("w", encoding="utf-8") as handle:
handle.write(content)
handle.flush()
os.fsync(handle.fileno())
os.replace(temporary, path)
finally:
if temporary.exists():
temporary.unlink()
await hass.async_add_executor_job(_write)
async def _commit_text(
hass: HomeAssistant,
path: Path,
filename: str,
structure: str,
reload_domain: str | None,
before: str,
proposed: str,
) -> dict[str, Any]:
current = await _read_text(hass, path)
if _fingerprint(current) != _fingerprint(before):
raise ToolError("YAML file changed after preview; refusing stale write and requiring a new preview")
await _validate_yaml_text(hass, proposed, structure, filename)
backup_id = await _create_backup(hass, filename, before)
await _atomic_write(hass, path, proposed)
validation = await _check_full_config(hass)
if not validation["valid"]:
await _atomic_write(hass, path, before)
raise ToolError(f"Home Assistant config validation failed; original restored from {backup_id}: {validation.get('error')}")
try:
if reload_domain:
await _reload(hass, reload_domain)
except Exception:
await _atomic_write(hass, path, before)
if reload_domain:
try:
await _reload(hass, reload_domain)
except Exception:
pass
raise
readback = await _read_text(hass, path)
return {
"changed": readback == proposed,
"file": filename,
"backup_id": backup_id,
"full_config_valid": True,
"reloaded": reload_domain,
"restart_required": reload_domain is None,
"new_sha256": _fingerprint(readback),
"exact_readback_match": readback == proposed,
}
async def _commit_data(
hass: HomeAssistant,
path: Path,
filename: str,
structure: str,
reload_domain: str | None,
before: str,
data: Any,
) -> dict[str, Any]:
from homeassistant.util.yaml import save_yaml
def _render() -> str:
temporary = path.with_name(f".{path.name}.hass-mcp-render-{secrets.token_hex(6)}.tmp")
try:
save_yaml(str(temporary), data)
return temporary.read_text(encoding="utf-8")
finally:
if temporary.exists():
temporary.unlink()
proposed = await hass.async_add_executor_job(_render)
return await _commit_text(hass, path, filename, structure, reload_domain, before, proposed)
async def _load(hass: HomeAssistant, path: Path, structure: str) -> Any:
return await _validate_yaml_text(hass, await _read_text(hass, path), structure, path.name)
async def _reload(hass: HomeAssistant, domain: str | None) -> None:
if not domain:
return
try:
await hass.services.async_call(domain, "reload", {}, blocking=True, context=user_context())
except Exception as error:
raise internal_error(f"{domain}.reload failed", error) from error
def _derive_entity_id(domain: str, structure: str, new_id: str, config: dict[str, Any]) -> str:
slug = slugify(new_id) if structure == "dict" else slugify(config.get("alias") or new_id)
return f"{domain}.{slug}"
def _to_list(data: Any, structure: str) -> list[dict[str, Any]]:
if structure == "list":
return list(data)
return [{"id": key, **value} for key, value in data.items()]
def _find(data: Any, structure: str, id: str) -> dict[str, Any] | None:
if structure == "list":
for entry in data:
if entry.get("id") == id or entry.get("alias") == id:
return entry
return None
return {"id": id, **data[id]} if id in data else None
def _replace(data: Any, structure: str, id: str, new: dict[str, Any]) -> bool:
if structure == "list":
for index, entry in enumerate(data):
if entry.get("id") == id:
data[index] = {"id": id, **{key: value for key, value in new.items() if key != "id"}}
return True
return False
if id in data:
data[id] = new
return True
return False
def _remove(data: Any, structure: str, id: str) -> bool:
if structure == "list":
for index, entry in enumerate(data):
if entry.get("id") == id:
del data[index]
return True
return False
if id in data:
del data[id]
return True
return False
-187
View File
@@ -1,187 +0,0 @@
"""Small, explicit, read-only Radarr tools built on the upstream API client."""
import asyncio
from typing import Any
from fastmcp import FastMCP
from pydantic import Field
from arr_mcp.auth import get_radarr_client
async def _call(client: Any, action: str, kwargs: dict[str, Any] | None = None) -> Any:
"""Call one known upstream API method without exposing dynamic dispatch."""
return await asyncio.to_thread(getattr(client, action), **(kwargs or {}))
def _plain(value: Any) -> Any:
if hasattr(value, "model_dump") and callable(value.model_dump):
return value.model_dump()
if hasattr(value, "dict") and callable(value.dict):
return value.dict()
return value
def _movies(value: Any) -> list[dict[str, Any]]:
value = _plain(value)
if isinstance(value, dict) and "result" in value:
value = value["result"]
return [item for item in value if isinstance(item, dict)] if isinstance(value, list) else []
def _compact_movie(movie: dict[str, Any]) -> dict[str, Any]:
return {
key: movie[key]
for key in ("id", "title", "originalTitle", "year", "status", "monitored", "hasFile", "path", "tmdbId")
if movie.get(key) is not None
}
def _compact_release(item: dict[str, Any]) -> dict[str, Any]:
quality = item.get("quality") or {}
quality_name = (quality.get("quality") or {}).get("name") if isinstance(quality, dict) else None
result = {
key: item[key]
for key in ("guid", "title", "indexer", "indexerId", "size", "age", "seeders", "leechers", "protocol", "downloadAllowed", "releaseGroup")
if item.get(key) is not None
}
if quality_name:
result["quality"] = quality_name
if isinstance(item.get("rejections"), list) and item["rejections"]:
result["rejections"] = [str(reason)[:180] for reason in item["rejections"][:5]]
return result
def _codec_aliases(value: str) -> set[str]:
aliases = {
"h264": {"h264", "x264", "avc"},
"x264": {"h264", "x264", "avc"},
"avc": {"h264", "x264", "avc"},
"h265": {"h265", "x265", "hevc"},
"x265": {"h265", "x265", "hevc"},
"hevc": {"h265", "x265", "hevc"},
}
requested: set[str] = set()
for item in value.split(","):
key = item.strip().casefold()
if key:
requested.update(aliases.get(key, {key}))
return requested
def _compact_inventory(
movies: list[dict[str, Any]], *, codecs: str = "", query: str = "",
offset: int = 0, limit: int = 200,
) -> dict[str, Any]:
wanted = _codec_aliases(codecs)
needle = query.strip().casefold()
rows: list[dict[str, Any]] = []
for movie in movies:
movie_file = movie.get("movieFile") or {}
if not movie.get("hasFile") or not movie_file:
continue
media = movie_file.get("mediaInfo") or {}
codec = str(media.get("videoCodec") or "unknown")
if wanted and codec.casefold() not in wanted:
continue
title = str(movie.get("title") or "")
if needle and needle not in title.casefold():
continue
quality = movie_file.get("quality") or {}
quality_name = (quality.get("quality") or {}).get("name") if isinstance(quality, dict) else None
size = int(movie_file.get("size") or 0)
rows.append({
"radarrId": movie.get("id"),
"title": title,
"year": movie.get("year"),
"movieFileId": movie_file.get("id"),
"relativePath": movie_file.get("relativePath"),
"sizeBytes": size,
"sizeGiB": round(size / 1073741824, 2),
"quality": quality_name,
"resolution": media.get("resolution"),
"videoCodec": codec,
"videoBitDepth": media.get("videoBitDepth"),
"audioCodec": media.get("audioCodec"),
"audioLanguages": media.get("audioLanguages"),
"subtitles": media.get("subtitles"),
})
rows.sort(key=lambda row: (str(row["title"]).casefold(), row.get("year") or 0))
total = len(rows)
start = max(0, int(offset))
count = min(500, max(1, int(limit)))
selected = rows[start:start + count]
return {
"totalMatched": total,
"offset": start,
"returned": len(selected),
"hasMore": start + len(selected) < total,
"movies": selected,
}
def register_radarr_tools(mcp: FastMCP) -> None:
@mcp.tool(tags={"radarr"})
async def radarr_find_movie(
query: str = Field(description="Movie title or title fragment."),
limit: int = Field(default=10, ge=1, le=25),
) -> Any:
"""Find a movie already managed by Radarr. READ ONLY. Never starts a search or download."""
needle = query.strip().casefold()
if len(needle) < 2:
raise ValueError("query must contain at least two characters")
movies = _movies(await _call(get_radarr_client(), "get_movie"))
matches = [
_compact_movie(movie)
for movie in movies
if needle in " ".join(str(movie.get(key, "")) for key in ("title", "originalTitle", "sortTitle")).casefold()
]
return {"query": query, "match_count": len(matches), "matches": matches[:limit], "truncated": len(matches) > limit}
@mcp.tool(tags={"radarr"})
async def radarr_movie_codec_inventory(
video_codecs: str = Field(
default="",
description="Optional comma-separated filter, e.g. h264, x264, h265, x265 or hevc.",
),
query: str = Field(default="", description="Optional case-insensitive title fragment."),
offset: int = Field(default=0, ge=0),
limit: int = Field(default=200, ge=1, le=500),
) -> Any:
"""Compact authoritative Radarr movie-file inventory. Use for codec, resolution, language and size questions instead of get_movie, raw API requests or filesystem scans. Results are valid bounded JSON without alternate titles, images, overviews or ratings."""
client = get_radarr_client()
response = _plain(await _call(client, "get_movie"))
movies = _movies(response)
if not movies and response not in ([], {"result": []}):
raise RuntimeError("Radarr get_movie returned an unexpected response")
return _compact_inventory(
movies, codecs=video_codecs, query=query, offset=offset, limit=limit,
)
@mcp.tool(tags={"radarr"})
async def radarr_search_releases(
movie_id: int = Field(ge=1, description="Exact Radarr movie id returned by radarr_find_movie."),
release_group: str = Field(default="", description="Optional release-group filter."),
limit: int = Field(default=50, ge=1, le=100),
) -> Any:
"""Search Radarr's configured indexers for one movie. READ ONLY: never grabs or downloads a release."""
raw = _plain(await _call(get_radarr_client(), "get_release", {"movieId": movie_id}))
if isinstance(raw, dict) and "result" in raw:
raw = raw["result"]
releases = [item for item in raw if isinstance(item, dict)] if isinstance(raw, list) else []
needle = release_group.strip().casefold()
if needle:
releases = [
item for item in releases
if needle in (str(item.get("releaseGroup", "")) + " " + str(item.get("title", ""))).casefold()
]
compact = [_compact_release(item) for item in releases[:limit]]
return {
"movie_id": movie_id,
"release_group_filter": release_group or None,
"total": len(releases),
"returned": len(compact),
"truncated": len(releases) > limit,
"results": compact,
"download_started": False,
}
-676
View File
@@ -1,676 +0,0 @@
"""Sonarr condensed action-routed MCP tool.
CONCEPT:ECO-4.82 — gitlab-style organized per-service tool surface.
"""
import asyncio
import os
import json
import re
import secrets
import time
from typing import Any
from fastmcp import FastMCP
from pydantic import Field
from arr_mcp.auth import get_sonarr_client
MAX_COLLECTION_ITEMS = 50
APPROVAL_TTL_SECONDS = 600
_APPROVALS: dict[str, tuple[float, str]] = {}
async def _call(client: Any, action: str, kwargs: dict[str, Any] | None = None) -> Any:
"""Call one known upstream API method without exposing dynamic dispatch."""
method = getattr(client, action)
return await asyncio.to_thread(method, **(kwargs or {}))
def _plain(value: Any) -> Any:
if hasattr(value, "model_dump") and callable(value.model_dump):
return value.model_dump()
if hasattr(value, "dict") and callable(value.dict):
return value.dict()
if isinstance(value, list):
return [_plain(item) for item in value]
if isinstance(value, dict):
return {str(key): _plain(item) for key, item in value.items()}
return value
def _unwrap(value: Any) -> Any:
value = _plain(value)
if isinstance(value, dict) and set(value) == {"result"}:
return value["result"]
return value
def _pick(item: dict[str, Any], fields: tuple[str, ...]) -> dict[str, Any]:
return {field: item[field] for field in fields if item.get(field) is not None}
def _compact_series(item: dict[str, Any], include_seasons: bool = False) -> dict[str, Any]:
result = _pick(
item,
("id", "title", "sortTitle", "year", "status", "monitored", "path", "tvdbId"),
)
statistics = item.get("statistics") or {}
if isinstance(statistics, dict):
result["statistics"] = _pick(
statistics,
("seasonCount", "episodeFileCount", "episodeCount", "totalEpisodeCount", "sizeOnDisk", "percentOfEpisodes"),
)
if include_seasons:
result["seasons"] = [
{
**_pick(season, ("seasonNumber", "monitored")),
"statistics": _pick(
season.get("statistics") or {},
("episodeFileCount", "episodeCount", "totalEpisodeCount", "sizeOnDisk", "percentOfEpisodes"),
),
}
for season in item.get("seasons", [])
if isinstance(season, dict)
]
return result
def _compact_episode(item: dict[str, Any]) -> dict[str, Any]:
return _pick(
item,
("id", "seriesId", "seasonNumber", "episodeNumber", "title", "airDate", "airDateUtc", "monitored", "hasFile", "episodeFileId"),
)
def _compact_file(item: dict[str, Any]) -> dict[str, Any]:
quality = item.get("quality") or {}
quality_name = (quality.get("quality") or {}).get("name") if isinstance(quality, dict) else None
result = _pick(
item,
("id", "seriesId", "seasonNumber", "relativePath", "path", "size", "dateAdded", "releaseGroup"),
)
if quality_name:
result["quality"] = quality_name
return result
def _compact_release(item: dict[str, Any]) -> dict[str, Any]:
quality = item.get("quality") or {}
quality_name = (quality.get("quality") or {}).get("name") if isinstance(quality, dict) else None
result = _pick(
item,
(
"guid", "title", "indexer", "indexerId", "size", "age", "ageHours",
"seeders", "leechers", "protocol", "downloadAllowed", "releaseWeight",
"releaseGroup", "seasonNumber", "fullSeason",
),
)
if quality_name:
result["quality"] = quality_name
rejections = item.get("rejections")
if isinstance(rejections, list) and rejections:
result["rejections"] = [str(reason)[:180] for reason in rejections[:5]]
return result
def _bounded(items: list[Any], compact) -> dict[str, Any]:
total = len(items)
return {
"total": total,
"returned": min(total, MAX_COLLECTION_ITEMS),
"truncated": total > MAX_COLLECTION_ITEMS,
"items": [compact(item) for item in items[:MAX_COLLECTION_ITEMS] if isinstance(item, dict)],
"next_step": (
"Use find_series or narrower Sonarr parameters; do not repeat the same broad request."
if total > MAX_COLLECTION_ITEMS else None
),
}
def _compact_result(action: str, value: Any) -> Any:
value = _unwrap(value)
if isinstance(value, list):
if action in {"get_series", "get_series_lookup", "lookup_series"}:
return _bounded(value, _compact_series)
if action in {"get_episode", "get_calendar", "get_wanted_missing", "get_wanted_cutoff"}:
return _bounded(value, _compact_episode)
if action == "get_episodefile":
return _bounded(value, _compact_file)
if action == "get_release":
return _bounded(value, _compact_release)
return _bounded(value, lambda item: item)
if isinstance(value, dict) and action in {"get_series_id"}:
return _compact_series(value, include_seasons=True)
if isinstance(value, dict) and action in {"get_episode_id"}:
return _compact_episode(value)
if isinstance(value, dict) and action in {"get_episodefile_id"}:
return _compact_file(value)
return value
async def _find_series(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
query = str(kwargs.get("query", "")).strip()
if len(query) < 2:
raise ValueError("query must contain at least two characters")
limit = max(1, min(int(kwargs.get("limit", 8)), 15))
raw = _unwrap(await _call(client, "get_series"))
words = [word for word in re.findall(r"[a-z0-9]+", query.casefold()) if len(word) > 1]
matches = []
for item in raw if isinstance(raw, list) else []:
haystack = " ".join(
str(item.get(field, "")) for field in ("title", "sortTitle", "originalTitle", "alternateTitles")
).casefold()
if all(word in haystack for word in words):
matches.append(_compact_series(item, include_seasons=False))
return {
"query": query,
"matches": matches[:limit],
"match_count": len(matches),
"truncated": len(matches) > limit,
"task_complete": True,
"instruction": "Use the returned series id for details. Do not call get_series for discovery.",
}
async def _season_summary(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
series_id = int(kwargs["series_id"])
season_number = int(kwargs["season_number"])
series = _unwrap(await _call(client, "get_series_id", {"id": series_id}))
episodes = _unwrap(
await _call(client, "get_episode", {"seriesId": series_id, "seasonNumber": season_number})
)
files = _unwrap(
await _call(client, "get_episodefile", {"seriesId": series_id})
)
selected_episodes = [
_compact_episode(item) for item in episodes
if isinstance(item, dict) and item.get("seasonNumber") == season_number
] if isinstance(episodes, list) else []
selected_files = [
_compact_file(item) for item in files
if isinstance(item, dict) and item.get("seasonNumber") == season_number
] if isinstance(files, list) else []
groups = sorted({str(item.get("releaseGroup")) for item in selected_files if item.get("releaseGroup")})
return {
"series": _compact_series(series) if isinstance(series, dict) else {"id": series_id},
"season_number": season_number,
"episode_count": len(selected_episodes),
"file_count": len(selected_files),
"release_groups": groups,
"episodes": selected_episodes[:30],
"files": selected_files[:30],
"task_complete": True,
"instruction": "This is the complete compact season answer. Do not repeat broad series or episode queries.",
}
async def _search_releases(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
series_id = kwargs.get("series_id")
episode_id = kwargs.get("episode_id")
season_number = kwargs.get("season_number")
release_group = str(kwargs.get("release_group", "")).strip()
season_pack_only = kwargs.get("season_pack_only") is True
if series_id is None and episode_id is None:
raise ValueError("search_releases requires series_id or episode_id")
query: dict[str, Any] = {}
if series_id is not None:
query["seriesId"] = int(series_id)
if episode_id is not None:
query["episodeId"] = int(episode_id)
if season_number is not None:
query["seasonNumber"] = int(season_number)
raw = await _call(client, "get_release", query)
raw = _unwrap(raw)
if release_group and isinstance(raw, list):
needle = release_group.casefold()
raw = [
item for item in raw
if isinstance(item, dict)
and needle in (
str(item.get("releaseGroup", "")) + " " + str(item.get("title", ""))
).casefold()
]
if season_pack_only:
if season_number is None:
raise ValueError("season_pack_only=true requires season_number")
season_token = rf"(?:^|[. _-])S0*{int(season_number)}(?:[. _-]|$)"
episode_token = rf"S0*{int(season_number)}E\d+"
raw = [
item for item in raw
if isinstance(item, dict)
and (
item.get("fullSeason") is True
or (
re.search(season_token, str(item.get("title", "")), re.IGNORECASE)
and not re.search(episode_token, str(item.get("title", "")), re.IGNORECASE)
)
)
] if isinstance(raw, list) else raw
compact = _compact_result("get_release", raw)
return {
"task_complete": True,
"search_scope": {
"series_id": series_id,
"episode_id": episode_id,
"season_number": season_number,
"release_group_filter": release_group or None,
"season_pack_only": season_pack_only,
},
"monitoring_changed": False,
"download_started": False,
"results": compact,
"instruction": (
"These are Sonarr indexer results. Do not use web search to replace them. "
"If the user requests one specific release, release group, or complete season pack, "
"NEVER substitute an automatic episode search: preview that exact result with "
"preview_release_grab, then wait for explicit approval before grab_release."
),
}
def _episode_numbers(value: Any) -> list[int]:
if value is None:
return []
if not isinstance(value, list) or len(value) > 100:
raise ValueError("episode_numbers must be a JSON list with at most 100 entries")
numbers = sorted({int(item) for item in value})
if any(item < 0 or item > 9999 for item in numbers):
raise ValueError("episode_numbers contains an invalid episode number")
return numbers
async def _resolve_episode_search(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
series_id = int(kwargs["series_id"])
season_number = int(kwargs["season_number"])
requested_numbers = _episode_numbers(kwargs.get("episode_numbers"))
if series_id < 1 or season_number < 0:
raise ValueError("series_id and season_number must be non-negative identifiers")
series = _unwrap(await _call(client, "get_series_id", {"id": series_id}))
episodes = _unwrap(
await _call(client, "get_episode", {"seriesId": series_id, "seasonNumber": season_number})
)
candidates = [
item for item in episodes
if isinstance(item, dict)
and int(item.get("seasonNumber", -1)) == season_number
and (not requested_numbers or int(item.get("episodeNumber", -1)) in requested_numbers)
] if isinstance(episodes, list) else []
if requested_numbers:
found_numbers = {int(item.get("episodeNumber", -1)) for item in candidates}
missing_metadata = sorted(set(requested_numbers) - found_numbers)
if missing_metadata:
raise ValueError(f"Sonarr has no episode metadata for episode numbers: {missing_metadata}")
missing = [item for item in candidates if not bool(item.get("hasFile"))]
if not candidates:
raise ValueError("No Sonarr episodes match the requested scope")
if len(missing) > 100:
raise ValueError("Refusing to search more than 100 missing episodes at once")
compact = [_compact_episode(item) for item in missing]
scope = {
"series_id": series_id,
"series_title": str(series.get("title", "")) if isinstance(series, dict) else "",
"season_number": season_number,
"requested_episode_numbers": requested_numbers,
"missing_episode_ids": [int(item["id"]) for item in missing],
"missing_episode_numbers": [int(item["episodeNumber"]) for item in missing],
}
fingerprint = json.dumps(scope, ensure_ascii=False, sort_keys=True, separators=(",", ":"))
return {
"scope": scope,
"fingerprint": fingerprint,
"selected_episode_count": len(candidates),
"already_present_count": len(candidates) - len(missing),
"unmonitored_missing_count": sum(not bool(item.get("monitored")) for item in missing),
"missing_episodes": compact,
}
async def _preview_episode_search(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
resolved = await _resolve_episode_search(client, kwargs)
ticket = secrets.token_urlsafe(24)
now = time.monotonic()
for old_ticket, (expires, _) in list(_APPROVALS.items()):
if expires <= now:
_APPROVALS.pop(old_ticket, None)
_APPROVALS[ticket] = (now + APPROVAL_TTL_SECONDS, resolved["fingerprint"])
return {
"action": "preview-only",
**{key: value for key, value in resolved.items() if key != "fingerprint"},
"monitoring_changed": False,
"download_started": False,
"approval_ticket": ticket,
"approval_expires_in_seconds": APPROVAL_TTL_SECONDS,
"next_step": (
"Review series, season and episode list. Only after explicit approval call "
"start_episode_search with exactly the same scope, confirm=true and this ticket."
),
}
async def _start_episode_search(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
if kwargs.get("confirm") is not True:
raise PermissionError("confirm=true is required after reviewing preview_episode_search")
ticket = str(kwargs.get("approval_ticket", ""))
if not ticket:
raise PermissionError("approval_ticket is required")
resolved = await _resolve_episode_search(client, kwargs)
approval = _APPROVALS.pop(ticket, None)
if approval is None or approval[0] <= time.monotonic():
raise PermissionError("Approval ticket is missing, expired or already used")
if not secrets.compare_digest(approval[1], resolved["fingerprint"]):
raise PermissionError("Approval ticket does not match this exact episode search")
episode_ids = resolved["scope"]["missing_episode_ids"]
if not episode_ids:
return {
"ok": True,
"command_started": False,
"reason": "All selected episodes already have files",
"scope": resolved["scope"],
}
command = _unwrap(
await _call(client, "post_command", {"data": {"name": "EpisodeSearch", "episodeIds": episode_ids}})
)
return {
"ok": True,
"command_started": True,
"sonarr_command": _pick(command, ("id", "name", "status", "queued", "startedOn")) if isinstance(command, dict) else command,
"scope": resolved["scope"],
"monitoring_changed": False,
"download_may_start_immediately": True,
"instruction": (
"Sonarr is now searching its configured indexers and may immediately grab/download "
"the best acceptable release for every approved episode. EpisodeSearch is NOT a "
"read-only manual-search preview and is NOT limited by the monitored flag. Check the "
"queue before claiming that nothing was downloaded."
),
}
def _release_query(kwargs: dict[str, Any]) -> dict[str, Any]:
series_id = int(kwargs["series_id"])
if series_id < 1:
raise ValueError("series_id must be a positive Sonarr identifier")
query: dict[str, Any] = {"seriesId": series_id}
if kwargs.get("season_number") is not None:
season_number = int(kwargs["season_number"])
if season_number < 0:
raise ValueError("season_number must be non-negative")
query["seasonNumber"] = season_number
if kwargs.get("episode_id") is not None:
episode_id = int(kwargs["episode_id"])
if episode_id < 1:
raise ValueError("episode_id must be a positive Sonarr identifier")
query["episodeId"] = episode_id
return query
async def _resolve_release_grab(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
guid = str(kwargs.get("guid", "")).strip()
if not guid:
raise ValueError(
"guid is required; copy it from the exact search_releases result the user selected"
)
query = _release_query(kwargs)
raw = _unwrap(
await _call(client, "get_release", query)
)
matches = [
item for item in raw if isinstance(item, dict) and str(item.get("guid", "")) == guid
] if isinstance(raw, list) else []
if len(matches) != 1:
raise ValueError(
"The exact release GUID is no longer present in Sonarr's current indexer results; "
"run search_releases again and do not guess or substitute another release"
)
release = matches[0]
compact = _compact_release(release)
rejections = compact.get("rejections") or []
blocked_without_force = release.get("downloadAllowed") is False or bool(rejections)
force = kwargs.get("force") is True
existing_file_count = None
season_number = query.get("seasonNumber")
if season_number is not None:
episodes = _unwrap(
await _call(client, "get_episode", {"seriesId": query["seriesId"], "seasonNumber": season_number})
)
if isinstance(episodes, list):
existing_file_count = sum(
bool(item.get("hasFile")) for item in episodes if isinstance(item, dict)
)
stable_release = _pick(
release,
(
"guid", "title", "indexer", "indexerId", "size", "protocol",
"downloadAllowed", "releaseGroup", "seasonNumber", "fullSeason",
),
)
raw_rejections = release.get("rejections")
stable_release["rejections"] = (
[str(reason) for reason in raw_rejections]
if isinstance(raw_rejections, list)
else []
)
scope = {
"series_id": query["seriesId"],
"season_number": query.get("seasonNumber"),
"episode_id": query.get("episodeId"),
"guid": guid,
"force": force,
}
fingerprint = json.dumps(
{"scope": scope, "release": stable_release},
ensure_ascii=False,
sort_keys=True,
separators=(",", ":"),
)
return {
"scope": scope,
"fingerprint": fingerprint,
"release": compact,
"raw_release": release,
"existing_episode_files_in_season": existing_file_count,
"blocked_without_force": blocked_without_force,
}
async def _preview_release_grab(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
resolved = await _resolve_release_grab(client, kwargs)
ticket = secrets.token_urlsafe(24)
now = time.monotonic()
for old_ticket, (expires, _) in list(_APPROVALS.items()):
if expires <= now:
_APPROVALS.pop(old_ticket, None)
force_required = resolved["blocked_without_force"] and not resolved["scope"]["force"]
if not force_required:
_APPROVALS[ticket] = (now + APPROVAL_TTL_SECONDS, resolved["fingerprint"])
return {
"action": "preview-only",
"scope": resolved["scope"],
"release": resolved["release"],
"existing_episode_files_in_season": resolved["existing_episode_files_in_season"],
"monitoring_changed": False,
"download_started": False,
"existing_files_deleted": False,
"replacement_guaranteed": False,
"warning": (
"Grabbing a season pack does not itself delete or guarantee replacement of existing "
"episode files. Sonarr applies its import, quality-profile and upgrade rules after download."
),
"force_required": force_required,
"approval_ticket": None if force_required else ticket,
"approval_expires_in_seconds": None if force_required else APPROVAL_TTL_SECONDS,
"next_step": (
"This result has Sonarr rejections or downloadAllowed=false. Explain the rejections and "
"only after the user explicitly accepts them call preview_release_grab again with force=true."
if force_required else
"Show the exact title, size, indexer, rejections and overwrite warning. Only after explicit "
"approval call grab_release with exactly the same scope, confirm=true and this ticket."
),
}
async def _grab_release(client: Any, kwargs: dict[str, Any]) -> dict[str, Any]:
if kwargs.get("confirm") is not True:
raise PermissionError("confirm=true is required after reviewing preview_release_grab")
ticket = str(kwargs.get("approval_ticket", ""))
if not ticket:
raise PermissionError("approval_ticket is required")
resolved = await _resolve_release_grab(client, kwargs)
if resolved["blocked_without_force"] and not resolved["scope"]["force"]:
raise PermissionError(
"This release has Sonarr rejections or downloadAllowed=false; an explicitly approved "
"force=true preview is required"
)
approval = _APPROVALS.pop(ticket, None)
if approval is None or approval[0] <= time.monotonic():
raise PermissionError("Approval ticket is missing, expired or already used")
if not secrets.compare_digest(approval[1], resolved["fingerprint"]):
raise PermissionError("Approval ticket does not match this exact release grab")
result = _unwrap(
await _call(client, "post_release", {"data": resolved["raw_release"]})
)
return {
"ok": True,
"download_started": True,
"selected_release": resolved["release"],
"sonarr_result": _compact_release(result) if isinstance(result, dict) else result,
"monitoring_changed": False,
"existing_files_deleted": False,
"replacement_guaranteed": False,
"instruction": (
"The exact approved release was sent to Sonarr's configured download client. "
"Do not claim that an existing episode was overwritten; verify queue/import history later."
),
}
def register_sonarr_tools(mcp: FastMCP) -> None:
@mcp.tool(tags={"sonarr"})
async def sonarr_find_series(
query: str = Field(description="Series title or title fragment, for example Mord ist ihr Hobby."),
limit: int = Field(default=8, ge=1, le=15),
) -> Any:
"""Find a Sonarr series by name. READ ONLY. Never changes monitoring and never starts a search or download."""
client = get_sonarr_client()
return await _find_series(client, {"query": query, "limit": limit})
@mcp.tool(tags={"sonarr"})
async def sonarr_get_season_summary(
series_id: int = Field(description="Exact Sonarr series id returned by sonarr_find_series."),
season_number: int = Field(ge=0, description="Season number."),
) -> Any:
"""Return episodes and existing files for one Sonarr season. READ ONLY. File names do not prove audio language."""
return await _season_summary(
get_sonarr_client(),
{"series_id": series_id, "season_number": season_number},
)
@mcp.tool(tags={"sonarr"})
async def sonarr_search_releases(
series_id: int = Field(description="Exact Sonarr series id."),
season_number: int | None = Field(default=None, ge=0),
episode_id: int | None = Field(default=None, ge=1),
release_group: str = Field(default="", description="Optional release-group filter, for example FuN."),
season_pack_only: bool = Field(default=False, description="Only complete season packs. Requires season_number."),
) -> Any:
"""Search Sonarr's configured indexers and return compact matching releases. READ ONLY: does not alter monitoring, start automatic search, grab, or download anything."""
return await _search_releases(
get_sonarr_client(),
{
"series_id": series_id,
"season_number": season_number,
"episode_id": episode_id,
"release_group": release_group,
"season_pack_only": season_pack_only,
},
)
@mcp.tool(tags={"sonarr"})
async def sonarr_system_status() -> Any:
"""Return compact Sonarr version and runtime status. READ ONLY."""
return _compact_result("get_system_status", await _call(get_sonarr_client(), "get_system_status"))
if os.environ.get("ARR_MCP_WRITE", "").strip().lower() not in ("1", "true", "yes", "on"):
return
@mcp.tool(tags={"sonarr", "write"})
async def sonarr_preview_release_grab(
series_id: int = Field(description="Exact Sonarr series id."),
guid: str = Field(description="Exact GUID returned by sonarr_search_releases."),
season_number: int | None = Field(default=None, ge=0),
episode_id: int | None = Field(default=None, ge=1),
force: bool = Field(default=False),
) -> Any:
"""Preview one exact release grab and issue a short-lived approval ticket. Does not download anything."""
return await _preview_release_grab(
get_sonarr_client(),
{
"series_id": series_id,
"guid": guid,
"season_number": season_number,
"episode_id": episode_id,
"force": force,
},
)
@mcp.tool(tags={"sonarr", "write"})
async def sonarr_grab_release(
series_id: int = Field(description="Same series id used for the preview."),
guid: str = Field(description="Same exact release GUID used for the preview."),
approval_ticket: str = Field(description="Ticket returned by sonarr_preview_release_grab."),
confirm: bool = Field(description="Must be true after explicit user approval."),
season_number: int | None = Field(default=None, ge=0),
episode_id: int | None = Field(default=None, ge=1),
force: bool = Field(default=False),
) -> Any:
"""Grab exactly one previously previewed release. WRITE: can immediately start a download."""
return await _grab_release(
get_sonarr_client(),
{
"series_id": series_id,
"guid": guid,
"approval_ticket": approval_ticket,
"confirm": confirm,
"season_number": season_number,
"episode_id": episode_id,
"force": force,
},
)
@mcp.tool(tags={"sonarr", "write"})
async def sonarr_preview_episode_search(
series_id: int = Field(description="Exact Sonarr series id."),
season_number: int = Field(ge=0),
episode_numbers: list[int] | None = Field(default=None, description="Optional episode numbers; omit for all missing episodes in the season."),
) -> Any:
"""Preview an automatic Sonarr episode search. Does not change monitoring or download anything."""
return await _preview_episode_search(
get_sonarr_client(),
{"series_id": series_id, "season_number": season_number, "episode_numbers": episode_numbers},
)
@mcp.tool(tags={"sonarr", "write"})
async def sonarr_start_episode_search(
series_id: int = Field(description="Same series id used for the preview."),
season_number: int = Field(ge=0),
approval_ticket: str = Field(description="Ticket returned by sonarr_preview_episode_search."),
confirm: bool = Field(description="Must be true after explicit user approval."),
episode_numbers: list[int] | None = Field(default=None),
) -> Any:
"""Start a previously previewed automatic episode search. WRITE: may immediately download releases."""
return await _start_episode_search(
get_sonarr_client(),
{
"series_id": series_id,
"season_number": season_number,
"episode_numbers": episode_numbers,
"approval_ticket": approval_ticket,
"confirm": confirm,
},
)
-177
View File
@@ -1,177 +0,0 @@
#!/usr/bin/env python3
"""Generate Hermes and OpenWebUI MCP registrations from one JSON registry."""
from __future__ import annotations
import argparse
import json
import os
import pathlib
import sqlite3
import time
BEGIN = "# BEGIN MANAGED MCP SERVERS"
END = "# END MANAGED MCP SERVERS"
CLIENT_TOKEN = ""
def env_file(path: str) -> dict[str, str]:
values: dict[str, str] = {}
source = pathlib.Path(path)
if not source.is_file():
return values
for raw in source.read_text(encoding="utf-8", errors="replace").splitlines():
line = raw.strip()
if not line or line.startswith("#") or "=" not in line:
continue
key, value = line.split("=", 1)
values[key.strip()] = value.strip().strip('"').strip("'")
return values
def enabled(item: dict) -> bool:
required = item.get("required_file")
if required and not pathlib.Path(required).is_file():
return False
source = item.get("env_file")
if source:
values = env_file(source)
url_ready = bool(item.get("url")) or bool(values.get(item.get("url_env", "")))
key_ready = (not item.get("key_env")
or bool(values.get(item["key_env"]))
or (item.get("key_env") == "MCPHUB_BEARER_TOKEN"
and bool(CLIENT_TOKEN)))
return url_ready and key_ready
return True
def resolved(item: dict) -> tuple[str, str]:
if item.get("env_file"):
values = env_file(item["env_file"])
url = item.get("url") or values[item["url_env"]]
key = values.get(item.get("key_env", ""), "")
if not key and item.get("key_env") == "MCPHUB_BEARER_TOKEN":
key = CLIENT_TOKEN
return url, key
return item["url"], ""
def active(registry: pathlib.Path, client: str) -> list[dict]:
document = json.loads(registry.read_text(encoding="utf-8"))
if document.get("version") != 1 or not isinstance(document.get("servers"), list):
raise SystemExit("Unsupported MCP registry schema")
return [item for item in document["servers"] if client in item.get("clients", []) and enabled(item)]
def yaml_quote(value: str) -> str:
return json.dumps(value, ensure_ascii=False)
def hermes_block(items: list[dict]) -> str:
lines = [BEGIN, "mcp_servers:"]
for item in items:
url, key = resolved(item)
lines.extend([
f" {item.get('hermes_id', item['id'])}:",
f" url: {yaml_quote(url)}",
])
if key:
lines.extend([" headers:", f" Authorization: {yaml_quote('Bearer ' + key)}"])
if "tool_include" in item:
lines.append(" tools:")
lines.append(" include:")
for tool in item["tool_include"]:
lines.append(f" - {yaml_quote(str(tool))}")
elif "tool_exclude" in item:
lines.append(" tools:")
lines.append(" exclude:")
for tool in item["tool_exclude"]:
lines.append(f" - {yaml_quote(str(tool))}")
lines.extend([
f" timeout: {int(item.get('timeout', 300))}",
" connect_timeout: 30",
" supports_parallel_tool_calls: false",
])
lines.append(END)
return "\n".join(lines) + "\n"
def update_hermes(path: pathlib.Path, block: str) -> None:
if not path.is_file():
return
text = path.read_text(encoding="utf-8")
if BEGIN in text and END in text:
prefix, rest = text.split(BEGIN, 1)
_, suffix = rest.split(END, 1)
text = prefix.rstrip() + "\n\n" + block + suffix.lstrip("\n")
else:
marker = "\nmcp_servers:"
if marker in text:
text = text.split(marker, 1)[0].rstrip() + "\n\n" + block
else:
text = text.rstrip() + "\n\n" + block
path.write_text(text, encoding="utf-8")
def openwebui_connection(item: dict) -> dict:
url, key = resolved(item)
config = {"enable": True, "access_grants": []}
if item.get("functions"):
config["function_name_filter_list"] = item["functions"]
elif item.get("tool_include"):
config["function_name_filter_list"] = ",".join(item["tool_include"])
return {
"url": url, "path": "", "type": "mcp",
"auth_type": item.get("auth_type", "none"), "headers": None,
"key": key, "config": config,
"info": {"id": item["id"], "name": item["name"], "description": item["description"]},
}
def update_openwebui(db: pathlib.Path, items: list[dict]) -> None:
con = sqlite3.connect(db)
now = int(time.time())
row = con.execute("select value from config where key=?", ("tool_server.connections",)).fetchone()
old = json.loads(row[0]) if row else []
if not isinstance(old, list):
raise SystemExit("Unexpected OpenWebUI tool_server.connections format")
managed_ids = {
"athena-platform", "athena-operator-local", "web-general-local", "github-local",
"homeassistant-local", "arr-local", "navidrome-local", "mua",
"mua-readonly-local", "athena-terminal-local", "unraid-readonly-local", "web-local",
}
keep = [entry for entry in old if str((entry.get("info") or {}).get("id", "")) not in managed_ids]
keep.extend(openwebui_connection(item) for item in items)
with con:
con.execute(
"""insert into config (key,value,updated_at) values (?,?,?)
on conflict(key) do update set value=excluded.value,updated_at=excluded.updated_at""",
("tool_server.connections", json.dumps(keep, ensure_ascii=False), now),
)
con.close()
def main() -> None:
global CLIENT_TOKEN
parser = argparse.ArgumentParser()
parser.add_argument("--registry", type=pathlib.Path, required=True)
parser.add_argument("--hermes", type=pathlib.Path, action="append", default=[])
parser.add_argument("--openwebui-db", type=pathlib.Path)
parser.add_argument("--mcphub-token-file", type=pathlib.Path)
args = parser.parse_args()
if args.mcphub_token_file:
CLIENT_TOKEN = args.mcphub_token_file.read_text(encoding="utf-8").strip()
if not CLIENT_TOKEN:
raise SystemExit("MCPHub token file is empty")
if args.hermes:
block = hermes_block(active(args.registry, "hermes"))
for path in args.hermes:
update_hermes(path, block)
if args.openwebui_db:
update_openwebui(args.openwebui_db, active(args.registry, "openwebui"))
print("MCP_CLIENT_SYNC_OK")
if __name__ == "__main__":
main()
-96
View File
@@ -1,96 +0,0 @@
#!/usr/bin/env bash
set -Eeuo pipefail
CONTAINER=${OPENWEBUI_CONTAINER:-mike-ai-open-webui}
ENV_FILE=${NAVIDROME_MCP_ENV_FILE:-/etc/mike-ai/navidrome-mcp.env}
die() { printf 'FEHLER: %s\n' "$*" >&2; exit 1; }
[[ $EUID -eq 0 ]] || die "Bitte als root ausführen."
[[ -s $ENV_FILE ]] || die "Navidrome-Secret-Datei fehlt."
[[ $(docker inspect -f '{{.State.Health.Status}}' mike-ai-mcp-navidrome 2>/dev/null || true) == healthy ]] || \
die "Navidrome-MCP ist nicht gesund."
[[ $(docker inspect -f '{{.State.Health.Status}}' "$CONTAINER" 2>/dev/null || true) == healthy ]] || \
die "OpenWebUI ist nicht gesund."
expect_lastfm=false
grep -q '^LASTFM_API_KEY=..' "$ENV_FILE" && expect_lastfm=true
result=$(docker exec -i -e EXPECT_LASTFM="$expect_lastfm" "$CONTAINER" python - <<'PY'
import asyncio
import json
import os
import sqlite3
from mcp import ClientSession
from mcp.client.streamable_http import streamablehttp_client
EXPECTED_LASTFM = {
"get_similar_artists", "get_similar_tracks", "get_artist_info",
"get_top_tracks_by_artist", "get_trending_music", "get_artist_albums",
"get_album_info",
}
PLAYBACK = {"play_songs", "pause", "set_volume"}
def find_unanchored(value, path=""):
bad = []
if isinstance(value, dict):
for key, child in value.items():
here = f"{path}.{key}" if path else key
if key == "pattern" and (
not isinstance(child, str)
or not child.startswith("^")
or not child.endswith("$")
):
bad.append(here)
bad.extend(find_unanchored(child, here))
elif isinstance(value, list):
for index, child in enumerate(value):
bad.extend(find_unanchored(child, f"{path}[{index}]"))
return bad
async def verify():
async with streamablehttp_client(
"http://mike-ai-mcp-navidrome:3000/mcp"
) as (read, write, _):
async with ClientSession(read, write) as session:
await session.initialize()
result = await session.list_tools()
names = {tool.name for tool in result.tools}
bad = []
for tool in result.tools:
bad.extend(find_unanchored(tool.inputSchema, tool.name))
if bad:
raise SystemExit("Unverankerte JSON-Schema-Patterns: " + ", ".join(bad))
if PLAYBACK & names:
raise SystemExit("Playback-Werkzeuge sind auf dem Headless-Host aktiv.")
expect_lastfm = os.environ.get("EXPECT_LASTFM") == "true"
if expect_lastfm and not EXPECTED_LASTFM <= names:
raise SystemExit("Last.fm-Werkzeugkatalog ist unvollständig.")
if not expect_lastfm and EXPECTED_LASTFM & names:
raise SystemExit("Last.fm-Werkzeuge sind ohne konfigurierten Schlüssel aktiv.")
if expect_lastfm:
# Public metadata only. Do not print the returned chart data.
response = await session.call_tool(
"get_trending_music", {"type": "artists", "limit": 1}
)
if response.isError:
raise SystemExit("Öffentliche Last.fm-Testabfrage ist fehlgeschlagen.")
con = sqlite3.connect("/app/backend/data/webui.db")
row = con.execute(
"select value from config where key=?", ("tool_server.connections",)
).fetchone()
connections = json.loads(row[0]) if row else []
ids = {
str((connection.get("info") or {}).get("id", ""))
for connection in connections if isinstance(connection, dict)
}
if "navidrome-local" not in ids:
raise SystemExit("OpenWebUI-Verbindung navidrome-local fehlt.")
print(f"NAVIDROME_ACCEPTANCE_OK tools={len(names)} lastfm={str(expect_lastfm).lower()}")
asyncio.run(verify())
PY
)
[[ $result == NAVIDROME_ACCEPTANCE_OK\ * ]] || \
die "Navidrome-Abnahme lieferte keinen gültigen Erfolgsmarker."
printf '%s\n' "$result"