from fastapi import FastAPI, Request, HTTPException, Depends from fastapi.middleware.cors import CORSMiddleware from loguru import logger import json from typing import Dict, Any, List import asyncio from concurrent.futures import ThreadPoolExecutor from .config import Settings, get_settings from .facebook import FacebookClient from .sheets import SheetsClient from .supabase_db import SupabaseClient from .embedding import EmbeddingClient from .utils import setup_logging, extract_command, extract_keywords, timing_decorator, ensure_log_dir, validate_config from .constants import VEHICLE_KEYWORDS, SHEET_RANGE from .health import router as health_router app = FastAPI(title="WeBot Facebook Messenger API") # Add CORS middleware app.add_middleware( CORSMiddleware, allow_origins=["*"], allow_credentials=True, allow_methods=["*"], allow_headers=["*"], ) # Initialize clients settings = get_settings() setup_logging(settings.log_level) facebook_client = FacebookClient(settings.facebook_app_secret) sheets_client = SheetsClient( settings.google_sheets_credentials_file, settings.google_sheets_token_file, settings.conversation_sheet_id ) supabase_client = SupabaseClient(settings.supabase_url, settings.supabase_key) embedding_client = EmbeddingClient() # Keywords to look for in messages VEHICLE_KEYWORDS = ["xe máy", "ô tô", "xe đạp", "xe hơi"] app.include_router(health_router) ensure_log_dir() validate_config(settings) executor = ThreadPoolExecutor(max_workers=4) @app.get("/webhook") async def verify_webhook(request: Request): """ Xác thực webhook Facebook Messenger. Input: request (Request) - request từ Facebook với các query params. Output: Trả về challenge nếu verify thành công, lỗi nếu thất bại. """ params = dict(request.query_params) mode = params.get("hub.mode") token = str(params.get("hub.verify_token", "")) challenge = str(params.get("hub.challenge", "")) if not all([mode, token, challenge]): raise HTTPException(status_code=400, detail="Missing parameters") return await facebook_client.verify_webhook( token, challenge, settings.facebook_verify_token ) @app.post("/webhook") @timing_decorator async def webhook(request: Request): """ Nhận và xử lý message từ Facebook Messenger webhook. Input: request (Request) - request chứa payload JSON từ Facebook. Output: JSON status. """ body_bytes = await request.body() # Verify request is from Facebook if not facebook_client.verify_signature(request, body_bytes): raise HTTPException(status_code=403, detail="Invalid signature") try: body = json.loads(body_bytes) message_data = facebook_client.parse_message(body) if not message_data: return {"status": "ok"} # Process the message await process_message(message_data) return {"status": "ok"} except Exception as e: logger.error(f"Error processing webhook: {e}") raise HTTPException(status_code=500, detail="Internal server error") async def process_message(message_data: Dict[str, Any]): """ Xử lý message từ người dùng Facebook, phân tích, truy vấn, gửi phản hồi và log lại. Input: message_data (dict) - thông tin message đã parse từ Facebook. Output: None (gửi message và log hội thoại). """ sender_id = message_data["sender_id"] page_id = message_data["page_id"] message_text = message_data["text"] logger.bind(user_id=sender_id, page_id=page_id, message=message_text).info("Processing message") if not message_text: return # Get page access token page_token = await supabase_client.get_page_token(page_id) if not page_token: logger.error(f"No access token found for page {page_id}") return # Extract command and keywords command, remaining_text = extract_command(message_text) keywords = extract_keywords(message_text, VEHICLE_KEYWORDS) # Get conversation history (run in thread pool) history = await asyncio.get_event_loop().run_in_executor( executor, lambda: sheets_client.get_conversation_history(sender_id, page_id).result() ) response = "" if command == "xong": if not keywords: response = "Vui lòng cho biết loại phương tiện bạn cần tìm (xe máy, ô tô...)" else: # Create embedding from message embedding = await embedding_client.create_embedding(message_text) # Search for similar documents matches = await supabase_client.match_documents(embedding) if matches: response = format_search_results(matches) else: response = "Xin lỗi, tôi không tìm thấy thông tin phù hợp." else: response = "Vui lòng cung cấp thêm thông tin và gõ lệnh \\xong khi hoàn tất." # Send response await facebook_client.send_message(page_token, sender_id, response) # Log conversation (run in thread pool) await asyncio.get_event_loop().run_in_executor( executor, lambda: sheets_client.log_conversation(sender_id, page_id, message_text, keywords, response).result() ) def format_search_results(matches: List[Dict[str, Any]]) -> str: """ Format kết quả truy vấn vector search thành chuỗi gửi về user. Input: matches (list[dict]) - danh sách kết quả từ Supabase. Output: Chuỗi kết quả đã format. """ if not matches: return "Không tìm thấy kết quả phù hợp." result = "Đây là một số kết quả phù hợp:\n\n" for i, match in enumerate(matches, 1): result += f"{i}. {match['content']}\n" if match.get('metadata', {}).get('url'): result += f" Link: {match['metadata']['url']}\n" result += "\n" return result.strip() if __name__ == "__main__": import uvicorn uvicorn.run( "app.main:app", host=settings.host, port=settings.port )