codeBOKER/customer_service
1
1import os2from datetime import datetime3from typing import List, Dict, Optional4from supabase import create_async_client, AsyncClient5import logging6from config import SUPABASE_URL, SUPABASE_KEY7 8class DatabaseManager:9 def __init__(self, supabase_url: str = SUPABASE_URL, supabase_key: str = SUPABASE_KEY):10 if not supabase_url or not supabase_key:11 raise ValueError("SUPABASE_URL and SUPABASE_KEY must be set")12 self.supabase_url = supabase_url13 self.supabase_key = supabase_key14 self.supabase: Optional[AsyncClient] = None15 self.logger = logging.getLogger(__name__)16 17 async def _ensure_connection(self):18 if self.supabase is None:19 await self.connect()20 21 async def connect(self):22 """Initialize the async client"""23 if not self.supabase:24 self.supabase = await create_async_client(self.supabase_url, self.supabase_key)25 26 async def create_or_update_user(self, telegram_id: int, username: str = None, 27 first_name: str = None, last_name: str = None):28 await self._ensure_connection()29 try:30 existing_user = await self.supabase.table("users").select("id").eq("telegram_id", telegram_id).execute()31 32 user_data = {33 "telegram_id": telegram_id,34 "username": username,35 "first_name": first_name,36 "last_name": last_name,37 "updated_at": datetime.utcnow().isoformat()38 }39 40 if existing_user.data:41 result = await self.supabase.table("users").update(user_data).eq("telegram_id", telegram_id).execute()42 else:43 user_data["created_at"] = datetime.utcnow().isoformat()44 result = await self.supabase.table("users").insert(user_data).execute()45 46 return result.data[0] if result.data else None47 except Exception as e:48 self.logger.error(f"Error creating/updating user: {e}")49 return None50 51 async def save_message(self, telegram_id: int, message_text: str, message_type: str):52 await self._ensure_connection()53 try:54 await self.create_or_update_user(telegram_id)55 56 message_data = {57 "telegram_id": telegram_id,58 "message_text": message_text,59 "message_type": message_type,60 "created_at": datetime.utcnow().isoformat()61 }62 63 result = await self.supabase.table("messages").insert(message_data).execute()64 await self._ensure_active_session(telegram_id)65 66 return result.data[0] if result.data else None67 except Exception as e:68 self.logger.error(f"Error saving message: {e}")69 return None70 71 async def get_conversation_history(self, telegram_id: int, limit: int = 10) -> List[Dict]:72 await self._ensure_connection()73 try:74 result = await (self.supabase.table("messages")75 .select("message_text, message_type, created_at")76 .eq("telegram_id", telegram_id)77 .order("created_at", desc=True)78 .limit(limit)79 .execute())80 return result.data if result.data else []81 except Exception as e:82 self.logger.error(f"Error getting history: {e}")83 return []84 85 async def _ensure_active_session(self, telegram_id: int):86 await self._ensure_connection()87 try:88 active = await (self.supabase.table("conversation_sessions")89 .select("id")90 .eq("telegram_id", telegram_id)91 .is_("session_end", "null")92 .execute())93 94 if not active.data:95 session_data = {96 "telegram_id": telegram_id,97 "session_start": datetime.utcnow().isoformat(),98 "created_at": datetime.utcnow().isoformat()99 }100 await self.supabase.table("conversation_sessions").insert(session_data).execute()101 except Exception as e:102 self.logger.error(f"Error ensuring session: {e}")103 104db_manager = DatabaseManager()