"""server.py — MCP per query sui dati di business ICR (Emesso/Anagrafica/Diritti...). Nato da una discussione con Mauro (09/08/2026): Claude.ai dovrebbe poter interrogare i dati di business direttamente, con query generate da lui stesso — questo server resta puramente meccanico ("agnostico", stesso principio di gate-ufficio): esegue la query, non interpreta il dominio. La conoscenza di schema/business-logic non vive qui — Claude.ai la legge da archivio/Adrian/agenti/data-expert.md, già nel mirror Drive "SecondBrain", stessa fonte usata dal subagente data-expert. Autenticazione OAuth: recuperata da git (vault-secondbrain/server.py, commit 5b461f6~1, cancellato il 29/07/2026 con lo smantellamento del vecchio vault) — stesso pattern collaudato con Elon/Claude.ai, adattato al nuovo SDK (mcp 2.0.0, MCPServer al posto di FastMCP — API di auth invariata, verificato compatibile 09/08/2026). Nessuna logica di dominio recuperata insieme, solo l'infrastruttura OAuth (generica, non legata al contenuto del vault). Sicurezza query: stesse regole del Frank-relay (readonly_actions.py in PYTHON/gate-ufficio/) — solo SELECT/WITH, nessuna keyword di scrittura, nessun punto e virgola interno (niente query multiple). """ from __future__ import annotations import json import os import re import secrets import time from pathlib import Path import duckdb from dotenv import load_dotenv from mcp.server.auth.provider import ( AccessToken, AuthorizationCode, AuthorizationParams, OAuthAuthorizationServerProvider, RefreshToken, construct_redirect_uri, ) from mcp.server.auth.settings import AuthSettings, ClientRegistrationOptions, RevocationOptions from mcp.server.mcpserver import MCPServer from mcp.server.transport_security import TransportSecuritySettings from mcp.shared.auth import OAuthClientInformationFull, OAuthToken from pydantic import AnyHttpUrl, AnyUrl HERE = Path(__file__).resolve().parent load_dotenv(HERE / ".env") PARQUET_DIR = ( HERE.parent.parent / "PYTHON" / "MyICR_Suite" / "local_db" / "parquet" ) _WRITE_KEYWORDS = { "INSERT", "UPDATE", "DELETE", "DROP", "CREATE", "ALTER", "TRUNCATE", "REPLACE", "MERGE", "COPY", "EXPORT", "CALL", "PRAGMA", "ATTACH", "DETACH", } DEFAULT_LIMIT = 1000 # ── OAuth (recuperato da vault-secondbrain/server.py, generico, adattato) ────── TRUSTED_REDIRECT_HOSTS = {"claude.ai", "claude.com", "localhost", "127.0.0.1"} class TrustedClient(OAuthClientInformationFull): """Client il cui redirect_uri non è fissato a una stringa esatta: viene accettato qualsiasi URI con host in TRUSTED_REDIRECT_HOSTS. Il client_secret resta obbligatorio allo scambio del token, quindi l'autenticità non dipende da questo.""" def validate_redirect_uri(self, redirect_uri: AnyUrl | None) -> AnyUrl: if redirect_uri is not None and redirect_uri.host in TRUSTED_REDIRECT_HOSTS: return redirect_uri return super().validate_redirect_uri(redirect_uri) ISSUER_URL = os.environ.get("ISSUER_URL", "https://query.privcloud.dev") CLIENT_ID = os.environ["MCP_QUERY_OAUTH_CLIENT_ID"] CLIENT_SECRET = os.environ["MCP_QUERY_OAUTH_CLIENT_SECRET"] ACCESS_TOKEN_TTL = 3600 * 24 * 30 # 30 giorni — stesso valore del vault, non serve rotazione frequente TOKENS_PATH = HERE / "oauth_tokens.json" # locale al progetto, gitignored (era in adrian-ops/, non più esistente) class McpQueryOAuthProvider(OAuthAuthorizationServerProvider): """Provider OAuth con persistenza token su disco. Client pre-registrato per Claude.ai + registrazione dinamica per Claude Code (RFC 7591).""" def __init__(self) -> None: self._auth_codes: dict[str, AuthorizationCode] = {} self._access_tokens: dict[str, AccessToken] = {} self._refresh_tokens: dict[str, RefreshToken] = {} self._load_tokens() self._clients: dict[str, OAuthClientInformationFull] = { CLIENT_ID: TrustedClient( client_id=CLIENT_ID, client_secret=CLIENT_SECRET, redirect_uris=[AnyUrl("https://claude.ai/api/mcp/auth_callback")], grant_types=["authorization_code", "refresh_token"], response_types=["code"], token_endpoint_auth_method="client_secret_post", ) } def _load_tokens(self) -> None: if not TOKENS_PATH.exists(): return try: data = json.loads(TOKENS_PATH.read_text()) now = time.time() for tok, d in data.get("access_tokens", {}).items(): if d.get("expires_at", 0) > now: self._access_tokens[tok] = AccessToken(**d) for tok, d in data.get("refresh_tokens", {}).items(): self._refresh_tokens[tok] = RefreshToken(**d) except Exception: pass # token corrotti o incompatibili — ripartenza pulita def _save_tokens(self) -> None: try: data = { "access_tokens": {t: v.model_dump() for t, v in self._access_tokens.items()}, "refresh_tokens": {t: v.model_dump() for t, v in self._refresh_tokens.items()}, } TOKENS_PATH.write_text(json.dumps(data)) except Exception: pass async def get_client(self, client_id: str) -> OAuthClientInformationFull | None: return self._clients.get(client_id) async def register_client(self, client_info: OAuthClientInformationFull) -> None: # Registrazione dinamica disattivata (30/08/2026, trovato in audit di sicurezza): # con ClientRegistrationOptions(enabled=True) questo endpoint accettava QUALUNQUE # client si registrasse, senza alcuna verifica di identita' - unico gate reale era # poi authorize() che pero' non richiede consenso e rilascia un codice a chiunque # presenti un client, anche appena auto-registrato. Nessun uso reale mai osservato # nei log (un solo client_id, quello pre-registrato da .env, sempre dalla stessa # sorgente WireGuard) - la registrazione dinamica era per "Claude Code" ma non # risultava effettivamente necessaria. Sollevare invece di accettare in silenzio. raise NotImplementedError("Registrazione dinamica client disattivata su mcp-query") async def authorize(self, client: OAuthClientInformationFull, params: AuthorizationParams) -> str: # Difesa in profondita' (30/08/2026): con registrazione dinamica disattivata questo # client dovrebbe gia' essere solo CLIENT_ID, ma un controllo esplicito qui non # costa nulla e non dipende dal fatto che get_client() venga chiamato prima da SDK. if client.client_id != CLIENT_ID: raise ValueError(f"Client non autorizzato: {client.client_id}") code = secrets.token_urlsafe(32) self._auth_codes[code] = AuthorizationCode( code=code, scopes=params.scopes or [], expires_at=time.time() + 600, client_id=client.client_id, code_challenge=params.code_challenge, redirect_uri=params.redirect_uri, redirect_uri_provided_explicitly=params.redirect_uri_provided_explicitly, resource=params.resource, ) return construct_redirect_uri(str(params.redirect_uri), code=code, state=params.state) async def load_authorization_code( self, client: OAuthClientInformationFull, authorization_code: str ) -> AuthorizationCode | None: return self._auth_codes.get(authorization_code) async def exchange_authorization_code( self, client: OAuthClientInformationFull, authorization_code: AuthorizationCode ) -> OAuthToken: del self._auth_codes[authorization_code.code] access_token = secrets.token_urlsafe(32) refresh_token = secrets.token_urlsafe(32) self._access_tokens[access_token] = AccessToken( token=access_token, client_id=client.client_id, scopes=authorization_code.scopes, expires_at=int(time.time() + ACCESS_TOKEN_TTL), resource=authorization_code.resource, ) self._refresh_tokens[refresh_token] = RefreshToken( token=refresh_token, client_id=client.client_id, scopes=authorization_code.scopes, ) self._save_tokens() return OAuthToken( access_token=access_token, token_type="bearer", expires_in=ACCESS_TOKEN_TTL, refresh_token=refresh_token, scope=" ".join(authorization_code.scopes), ) async def load_refresh_token(self, client: OAuthClientInformationFull, refresh_token: str) -> RefreshToken | None: return self._refresh_tokens.get(refresh_token) async def exchange_refresh_token( self, client: OAuthClientInformationFull, refresh_token: RefreshToken, scopes: list[str], ) -> OAuthToken: del self._refresh_tokens[refresh_token.token] access_token = secrets.token_urlsafe(32) new_refresh = secrets.token_urlsafe(32) self._access_tokens[access_token] = AccessToken( token=access_token, client_id=client.client_id, scopes=scopes or refresh_token.scopes, expires_at=int(time.time() + ACCESS_TOKEN_TTL), ) self._refresh_tokens[new_refresh] = RefreshToken( token=new_refresh, client_id=client.client_id, scopes=scopes or refresh_token.scopes, ) self._save_tokens() return OAuthToken( access_token=access_token, token_type="bearer", expires_in=ACCESS_TOKEN_TTL, refresh_token=new_refresh, scope=" ".join(scopes or refresh_token.scopes), ) async def load_access_token(self, token: str) -> AccessToken | None: at = self._access_tokens.get(token) if at and at.expires_at and at.expires_at < time.time(): del self._access_tokens[token] return None return at async def revoke_token(self, token) -> None: self._access_tokens.pop(getattr(token, "token", token), None) self._refresh_tokens.pop(getattr(token, "token", token), None) self._save_tokens() def _build_server(with_auth: bool) -> MCPServer: """with_auth=False per i test locali via stdio (test_client.py) — l'auth OAuth ha senso solo sul transport HTTP esposto pubblicamente, e costruire il provider a vuoto ogni test rallenta inutilmente senza aggiungere nulla.""" kwargs = {} if with_auth: kwargs["auth_server_provider"] = McpQueryOAuthProvider() kwargs["auth"] = AuthSettings( issuer_url=AnyHttpUrl(ISSUER_URL), resource_server_url=AnyHttpUrl(f"{ISSUER_URL}/mcp"), # Disattivata 30/08/2026 (audit sicurezza): era aperta a chiunque, nessun # controllo di identita' - vedi commento su register_client() sopra. client_registration_options=ClientRegistrationOptions(enabled=False), revocation_options=RevocationOptions(enabled=True), ) return MCPServer( "mcp-query", description=( "Query di sola lettura sui dati di business ICR (diritti, anagrafica " "prodotti, emesso/palinsesto, boxoffice, cast, IMDb) — mirror parquet " "sincronizzati sulla Nave. Prima di interrogare, leggi lo schema/le " "regole di business in archivio/Adrian/agenti/data-expert.md (mirror " "Drive 'SecondBrain') — questo server esegue la query così com'è, " "senza conoscenza propria del dominio." ), **kwargs, ) # Server usato per i test locali (stdio, no auth) — vedi __main__ sotto per # la variante con auth usata in produzione (streamable HTTP). server = _build_server(with_auth=False) def _connect() -> duckdb.DuckDBPyConnection: """Apre una connessione DuckDB in-memory con una view per ogni parquet disponibile — così le query usano nomi tabella semplici (es. "emesso") invece del path completo al file.""" con = duckdb.connect(":memory:") for path in sorted(PARQUET_DIR.glob("*.parquet")): view_name = path.stem # CREATE VIEW non supporta parametri preparati in DuckDB — path # comunque non controllato da input esterno (solo glob locale). escaped_path = str(path).replace("'", "''") con.execute(f"CREATE VIEW \"{view_name}\" AS SELECT * FROM read_parquet('{escaped_path}')") return con def _validate_readonly(sql: str) -> str: """Stessa validazione di readonly_actions.py::execute_readonly (azione duckdb_query) — solo SELECT/WITH, niente keyword di scrittura, niente query multiple. Ritorna la query "pulita" (senza whitespace/`;` finali) o solleva ValueError.""" sql_stripped = sql.strip() first_word = sql_stripped.upper().split()[0] if sql_stripped else "" if first_word not in ("SELECT", "WITH"): raise ValueError(f"Solo SELECT/WITH ammessi come inizio query (trovato: '{first_word}')") sql_stripped = sql_stripped.rstrip().rstrip(";") if ";" in sql_stripped: raise ValueError("Query multiple non ammesse (';' interno non consentito)") sql_upper = sql_stripped.upper() for kw in _WRITE_KEYWORDS: if re.search(r"\b" + kw + r"\b", sql_upper): raise ValueError(f"Keyword non ammessa: '{kw}'") return sql_stripped @server.tool() def list_tables() -> list[dict]: """Elenca le tabelle (view) disponibili, con conteggio righe. Usalo per sapere cosa esiste — per capire COSA contiene ciascuna tabella e come interrogarla bene, leggi archivio/Adrian/agenti/data-expert.md su Drive prima di scrivere query complesse.""" con = _connect() try: tables = [] for path in sorted(PARQUET_DIR.glob("*.parquet")): view_name = path.stem count = con.execute(f'SELECT COUNT(*) FROM "{view_name}"').fetchone()[0] tables.append({"table": view_name, "rows": count}) return tables finally: con.close() @server.tool() def query_business_data(sql: str, limit: int = DEFAULT_LIMIT) -> dict: """Esegue una query SQL di sola lettura (SELECT/WITH) sui dati di business ICR. Le tabelle disponibili si vedono con list_tables(); nomi colonne/join/regole di business vanno letti da archivio/Adrian/agenti/data-expert.md (mirror Drive "SecondBrain") PRIMA di scrivere la query — questo tool la esegue così com'è, senza correggerla né interpretarla.""" sql_clean = _validate_readonly(sql) con = _connect() try: rel = con.execute(sql_clean) cols = [d[0] for d in rel.description] rows = rel.fetchmany(limit) return { "columns": cols, "rows": [list(r) for r in rows], "truncated": len(rows) == limit, } finally: con.close() if __name__ == "__main__": import sys if "--real" in sys.argv: # Modalita' produzione: HTTP + OAuth, bind sull'IP WireGuard — MAI # lanciata da test_client.py (quello usa stdio, senza --real). import uvicorn prod_server = _build_server(with_auth=True) prod_server.tool()(list_tables) prod_server.tool()(query_business_data) # transport_security non è più un parametro del costruttore MCPServer # (mcp 2.0.0) — si passa qui, a streamable_http_app(). Bug reale # trovato il 09/08/2026 al primo collegamento vero da Claude.ai: # senza questo, "Invalid Host header" -> 421 su ogni richiesta POST # /mcp (l'header Host reale, query.privcloud.dev, non è né # 127.0.0.1 né localhost, i soli default impliciti). app = prod_server.streamable_http_app( transport_security=TransportSecuritySettings( allowed_hosts=["query.privcloud.dev", "10.0.0.2:8765", "localhost:8765"], allowed_origins=["https://query.privcloud.dev"], ) ) uvicorn.run(app, host="10.0.0.2", port=8765) else: # Modalita' test locale: stdio, nessuna auth — quella usata da # test_client.py per verificare il protocollo end-to-end. server.run()