| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144 |
- from fastapi import APIRouter, WebSocket, WebSocketDisconnect, Depends, Query, HTTPException, File, UploadFile, Form
- from services.chat_manager import manager
- from services.global_manager import global_manager
- import db
- import auth_utils
- import datetime
- import schemas
- import locales
- from dependencies import get_current_user
- import os
- import uuid
- import config
- import json
- from typing import Optional
- router = APIRouter(tags=["chat"])
- # In-memory storage for flood control: {user_id: timestamp}
- last_message_times = {}
- @router.get("/orders/{order_id}/messages")
- async def get_order_messages(order_id: int, user: dict = Depends(get_current_user)):
- role = user.get("role")
- user_id = user.get("id")
- # Fetch user chat status
- user_info = db.execute_query("SELECT can_chat FROM users WHERE id = %s", (user_id,))
- can_chat = user_info[0]['can_chat'] if user_info else False
- order = db.execute_query("SELECT user_id FROM orders WHERE id = %s", (order_id,))
- if not order: raise HTTPException(status_code=404, detail="Order not found")
- if role != 'admin':
- if order[0]['user_id'] != user_id: raise HTTPException(status_code=403, detail="Not authorized")
- if not can_chat: raise HTTPException(status_code=403, detail="Chat access disabled for your account")
- query = """
- SELECT m.id, m.is_from_admin, m.message, m.created_at, u.first_name, u.email
- FROM order_messages m
- LEFT JOIN users u ON m.user_id = u.id
- WHERE m.order_id = %s ORDER BY m.created_at ASC
- """
- messages = db.execute_query(query, (order_id,))
- for msg in messages:
- if msg.get('created_at'): msg['created_at'] = msg['created_at'].isoformat()
- msg['is_from_admin'] = bool(msg['is_from_admin'])
-
- # Mark messages as read
- if role == 'admin':
- db.execute_commit("UPDATE order_messages SET is_read = TRUE WHERE order_id = %s AND is_from_admin = FALSE AND is_read = FALSE", (order_id,))
- await global_manager.notify_admins()
- await global_manager.notify_order_read(order_id)
- else:
- db.execute_commit("UPDATE order_messages SET is_read = TRUE WHERE order_id = %s AND is_from_admin = TRUE AND is_read = FALSE", (order_id,))
- await global_manager.notify_user(user_id)
-
- return messages
- @router.post("/orders/{order_id}/messages")
- async def post_order_message(
- order_id: int,
- message: Optional[str] = Form(None),
- file: Optional[UploadFile] = File(None),
- user: dict = Depends(get_current_user),
- lang: str = "en"
- ):
- role = user.get("role")
- user_id = user.get("id")
-
- # Flood control for non-admin users
- if role != 'admin':
- now = datetime.datetime.utcnow().timestamp()
- last_time = last_message_times.get(user_id, 0)
- if now - last_time < 10:
- raise HTTPException(status_code=429, detail=locales.translate_error("flood_control", lang))
- last_message_times[user_id] = now
- if not message and not file:
- raise HTTPException(status_code=400, detail="Empty message")
- if message:
- message = message.strip()
- is_admin = (role == 'admin')
-
- if not is_admin:
- user_info = db.execute_query("SELECT can_chat FROM users WHERE id = %s", (user_id,))
- if not user_info or not user_info[0]['can_chat']:
- raise HTTPException(status_code=403, detail="Chat access disabled")
- order = db.execute_query("SELECT user_id FROM orders WHERE id = %s", (order_id,))
- if not order: raise HTTPException(status_code=404, detail="Order not found")
- if not is_admin and order[0]['user_id'] != user_id: raise HTTPException(status_code=403, detail="Not authorized")
- final_message = ""
- db_file_path = None
-
- if file and file.filename:
- file_ext = os.path.splitext(file.filename)[1]
- unique_filename = f"{uuid.uuid4()}{file_ext}"
- file_path = os.path.join(config.CHAT_UPLOADS_DIR, unique_filename)
- db_file_path = f"uploads/chats/{unique_filename}"
-
- with open(file_path, "wb") as buffer:
- while chunk := file.file.read(8192):
- buffer.write(chunk)
-
- msg_obj = {}
- if db_file_path:
- msg_obj["image_url"] = db_file_path
- if message:
- msg_obj["text"] = message
-
- final_message = json.dumps(msg_obj)
- query = "INSERT INTO order_messages (order_id, user_id, is_from_admin, message) VALUES (%s, %s, %s, %s)"
- msg_id = db.execute_commit(query, (order_id, user_id, is_admin, final_message))
- now = datetime.datetime.utcnow().isoformat()
-
- user_info = db.execute_query("SELECT first_name, email FROM users WHERE id = %s", (user_id,))
- first_name = user_info[0]['first_name'] if user_info else None
- email = user_info[0]['email'] if user_info else None
-
- await manager.broadcast_to_order(order_id, {
- "id": msg_id,
- "is_from_admin": is_admin,
- "message": final_message,
- "created_at": now,
- "first_name": first_name,
- "email": email
- })
-
- if is_admin:
- await global_manager.notify_user(order[0]['user_id'])
- else:
- from services import event_hooks
- await global_manager.notify_admins()
- await global_manager.notify_admins_new_message(order_id, message or "Image uploaded")
- event_hooks.on_message_received(order_id, user_id, message or "Image uploaded", db_file_path)
-
- return {
- "id": msg_id,
- "status": "sent",
- "message": final_message,
- "first_name": first_name,
- "email": email
- }
|