import asyncio import json import logging import firebase_admin import httpx from fastapi import Depends, FastAPI, File, HTTPException, Request, UploadFile, Header from fastapi.security import HTTPAuthorizationCredentials, HTTPBearer from firebase_admin import auth, credentials from string import Template from openai import OpenAI from typing import Optional import base64 # Importazioni dai moduli locali import config import prompts # --- INIT --- logging.basicConfig(level=logging.INFO) logger = logging.getLogger("monito_gateway") # Inizializzazione Firebase try: cred = credentials.Certificate(json.loads(config.FIREBASE_SERVICE_ACCOUNT_JSON)) firebase_admin.initialize_app(cred) except Exception as e: logger.error(f"Errore inizializzazione Firebase: {e}") app = FastAPI() security = HTTPBearer() client = OpenAI( api_key=config.HF_TOKEN, base_url="https://s5jfvadty4tkdlkr.eu-west-1.aws.endpoints.huggingface.cloud/v1" ) # --- UTILITY --- async def send_to_telegram(message: str): url = f"https://api.telegram.org/bot{config.TELEGRAM_BOT_TOKEN}/sendMessage" payload = {"chat_id": config.TELEGRAM_CHAT_ID, "text": f"🖥️ [Gateway]: {message}"} async with httpx.AsyncClient() as client: try: await client.post(url, json=payload, timeout=5.0) except Exception as e: logger.error(f"Telegram Log Fallito: {e}") async def verify_firebase_token(creds: HTTPAuthorizationCredentials = Depends(security)): try: return auth.verify_id_token(creds.credentials)["uid"] except Exception: raise HTTPException(status_code=401, detail="Token non valido") async def call_deepseek_with_retry(messages: list, max_retries=3): """ Versione Forense: Logga ogni passaggio, ogni byte ricevuto e ogni errore di parsing. """ logger.info(f"--- INIZIO CALL_DEEPSEEK: {len(messages)} messaggi ---") logger.debug(f"FULL_INPUT_MESSAGES: {json.dumps(messages, ensure_ascii=False)}") # Timeout aumentato a 300s per gestire risposte lunghe (es. generazione quiz) async with httpx.AsyncClient(timeout=300.0) as client: for attempt in range(1, max_retries + 1): try: logger.info(f"Tentativo {attempt}/{max_retries} in corso verso {config.DEEPSEEK_BASE_URL}...") response = await client.post( f"{config.DEEPSEEK_BASE_URL}/chat/completions", headers={"Authorization": f"Bearer {config.DEEPSEEK_API_KEY}"}, json={ "model": "deepseek-chat", "messages": messages, "temperature": 0.1, "response_format": {"type": "json_object"} } ) # Registriamo lo status e gli header per debuggare limitazioni (Rate Limit) logger.info(f"Ricevuta risposta. Status: {response.status_code}") if response.status_code != 200: error_body = response.text logger.error(f"Errore HTTP {response.status_code}: {error_body}") response.raise_for_status() data = response.json() content = data['choices'][0]['message']['content'] # LOGGING MASSIVO: vediamo cosa c'è davvero dentro logger.info(f"RAW_CONTENT_LENGTH: {len(content)}") logger.debug(f"RAW_CONTENT_PREVIEW: {content[:500]}") # Pulizia aggressiva ma tracciata cleaned = content.replace("```json", "").replace("```", "").strip() try: result = json.loads(cleaned) logger.info("Parsing JSON completato con successo.") return result except json.JSONDecodeError as json_e: logger.error(f"FALLIMENTO PARSING JSON: {json_e}") logger.error(f"Testo che ha causato il fallimento: {cleaned}") raise json_e except httpx.HTTPStatusError as http_e: logger.error(f"Errore di rete HTTP al tentativo {attempt}: {http_e}") if attempt == max_retries: raise http_e except Exception as e: logger.error(f"Eccezione generica al tentativo {attempt}: {type(e).__name__} - {str(e)}", exc_info=True) if attempt == max_retries: raise e # Backoff esponenziale sleep_time = 2 ** attempt logger.info(f"Attesa {sleep_time} secondi prima di riprovare...") await asyncio.sleep(sleep_time) # Se arriviamo qui, è perché abbiamo esaurito i tentativi logger.critical("call_deepseek_with_retry: ESAURITI TUTTI I TENTATIVI.") raise HTTPException(status_code=500, detail="DeepSeek non ha risposto dopo diversi tentativi.") # --- ENDPOINTS --- @app.get("/") async def root(): return {"status": "ok", "message": "Gateway AI Attivo"} @app.post("/extract-image") async def extract_image(uid: str = Depends(verify_firebase_token), file: UploadFile = File(...)): logger.info(f"Ricevuta richiesta /extract-image da UID: {uid}") await send_to_telegram(f"Nuova richiesta VLM da {uid}") try: image_bytes = await file.read() encoded_image = base64.b64encode(image_bytes).decode('utf-8') logger.debug(f"Immagine caricata, dimensione bytes: {len(image_bytes)}") response = client.chat.completions.create( model="ggml-org/Qwen2.5-VL-7B-Instruct-GGUF", messages=[{ "role": "user", "content": [ {"type": "image_url", "image_url": {"url": f"data:image/jpeg;base64,{encoded_image}"}}, {"type": "text", "text": "Describe this image in one sentence."} ] }], max_tokens=100 ) result = {"markdown": response.choices[0].message.content} logger.info(f"Risposta VLM ottenuta con successo per UID: {uid}") return result except Exception as e: logger.error(f"Errore critico in /extract-image: {str(e)}", exc_info=True) await send_to_telegram(f"🚨 Errore VLM: {str(e)}") raise HTTPException(status_code=502, detail=str(e)) @app.post("/analyze-document") async def analyze_document(request: Request, uid: str = Depends(verify_firebase_token)): try: body = await request.json() logger.info(f"INPUT_/analyze-document UID={uid} BODY_LEN={len(str(body))}") messages = [ {"role": "system", "content": prompts.PROMPT_ANALISI}, {"role": "user", "content": f"Analizza:\n{body.get('text', '')}"} ] result = await call_deepseek_with_retry(messages) logger.info(f"OUTPUT_/analyze-document UID={uid} RESULT={json.dumps(result, ensure_ascii=False)}") return result except Exception as e: logger.error(f"Errore fatale /analyze-document UID={uid}: {str(e)}", exc_info=True) raise HTTPException(status_code=500, detail=str(e)) @app.post("/validate-document") async def validate_document(request: Request, uid: str = Depends(verify_firebase_token)): try: body = await request.json() logger.info(f"INPUT_/validate-document UID={uid} BODY_LEN={len(str(body))}") messages = [ {"role": "system", "content": prompts.PROMPT_VALIDAZIONE}, {"role": "user", "content": f"Valida:\n{body.get('text', '')}"} ] result = await call_deepseek_with_retry(messages) logger.info(f"OUTPUT_/validate-document UID={uid} RESULT={json.dumps(result, ensure_ascii=False)}") return result except Exception as e: logger.error(f"Errore fatale /validate-document UID={uid}: {str(e)}", exc_info=True) raise HTTPException(status_code=500, detail=str(e)) @app.post("/generate-questions") async def generate_questions(request: Request, uid: str = Depends(verify_firebase_token)): body = await request.json() logger.info(f"INPUT_/generate-questions UID={uid} BODY={json.dumps(body)}") try: t = Template(prompts.PROMPT_QUIZ) prompt_finale = t.substitute( targetLanguage=body.get("language", "Italiano"), count=body.get("count", 3), macro=body.get("macro", "Generico"), argomento=body.get("argomento", "Generico") ) messages = [{"role": "system", "content": prompt_finale}, {"role": "user", "content": f"Genera:\n{body.get('content', '')}"}] result = await call_deepseek_with_retry(messages) logger.info(f"OUTPUT_/generate-questions UID={uid} RESULT={json.dumps(result, ensure_ascii=False)}") return result except Exception as e: logger.error(f"Errore fatale /generate-questions UID={uid}: {str(e)}", exc_info=True) raise HTTPException(status_code=500, detail=str(e)) @app.post("/generate-flashcards") async def generate_flashcards(request: Request, uid: str = Depends(verify_firebase_token)): body = await request.json() logger.info(f"INPUT_/generate-flashcards UID={uid} BODY={json.dumps(body)}") try: t = Template(prompts.PROMPT_FLASHCARD) prompt_finale = t.substitute( targetLanguage=body.get("language", "Italiano"), tag=body.get("tag", "Studio") ) messages = [{"role": "system", "content": prompt_finale}, {"role": "user", "content": f"Genera:\n{body.get('content', '')}"}] result = await call_deepseek_with_retry(messages) logger.info(f"OUTPUT_/generate-flashcards UID={uid} RESULT={json.dumps(result, ensure_ascii=False)}") return result except Exception as e: logger.error(f"Errore fatale /generate-flashcards UID={uid}: {str(e)}", exc_info=True) raise HTTPException(status_code=500, detail=str(e)) # --- ENDPOINTS KEEP-ALIVE --- @app.get("/keep-alive") async def keep_alive(authorization: Optional[str] = Header(None)): """ Endpoint unico /keep-alive che funge da: 1. Ping pubblico (se non viene inviato l'header Authorization) 2. Ping autenticato (se viene inviato il token Firebase) """ # CASO 1: RICHIESTA AUTENTICATA (Android pingServerAuth) if authorization and authorization.startswith("Bearer "): try: token = authorization.split(" ")[1] # Valida il token Firebase decoded_token = auth.verify_id_token(token) uid = decoded_token["uid"] logger.info(f"Keep-alive autenticato ricevuto da UID: {uid}") return {"status": "authenticated", "uid": uid} except Exception as e: logger.warning(f"Token di keep-alive non valido: {str(e)}") # Se il token è presente ma fallisce, blocchiamo la richiesta raise HTTPException(status_code=401, detail="Token non valido") # CASO 2: RICHIESTA PUBBLICA (Android debugConnection) logger.info("Keep-alive pubblico (ping) ricevuto.") return {"status": "alive", "type": "public"}