chat.py 4.9 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120
  1. from fastapi import APIRouter, WebSocket, WebSocketDisconnect, Depends, Query, HTTPException, File, UploadFile, Form
  2. from services.chat_manager import manager
  3. from services.global_manager import global_manager
  4. import db
  5. import auth_utils
  6. import datetime
  7. import schemas
  8. import locales
  9. from dependencies import get_current_user
  10. import os
  11. import uuid
  12. import config
  13. import json
  14. from typing import Optional
  15. router = APIRouter(tags=["chat"])
  16. # In-memory storage for flood control: {user_id: timestamp}
  17. last_message_times = {}
  18. @router.get("/orders/{order_id}/messages")
  19. async def get_order_messages(order_id: int, user: dict = Depends(get_current_user)):
  20. role = user.get("role")
  21. user_id = user.get("id")
  22. # Fetch user chat status
  23. user_info = db.execute_query("SELECT can_chat FROM users WHERE id = %s", (user_id,))
  24. can_chat = user_info[0]['can_chat'] if user_info else False
  25. order = db.execute_query("SELECT user_id FROM orders WHERE id = %s", (order_id,))
  26. if not order: raise HTTPException(status_code=404, detail="Order not found")
  27. if role != 'admin':
  28. if order[0]['user_id'] != user_id: raise HTTPException(status_code=403, detail="Not authorized")
  29. if not can_chat: raise HTTPException(status_code=403, detail="Chat access disabled for your account")
  30. messages = db.execute_query("SELECT id, is_from_admin, message, created_at FROM order_messages WHERE order_id = %s ORDER BY created_at ASC", (order_id,))
  31. for msg in messages:
  32. if msg.get('created_at'): msg['created_at'] = msg['created_at'].isoformat()
  33. msg['is_from_admin'] = bool(msg['is_from_admin'])
  34. # Mark messages as read
  35. if role == 'admin':
  36. 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,))
  37. await global_manager.notify_admins()
  38. await global_manager.notify_order_read(order_id)
  39. else:
  40. 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,))
  41. await global_manager.notify_user(user_id)
  42. return messages
  43. @router.post("/orders/{order_id}/messages")
  44. async def post_order_message(
  45. order_id: int,
  46. message: Optional[str] = Form(None),
  47. file: Optional[UploadFile] = File(None),
  48. user: dict = Depends(get_current_user),
  49. lang: str = "en"
  50. ):
  51. role = user.get("role")
  52. user_id = user.get("id")
  53. # Flood control for non-admin users
  54. if role != 'admin':
  55. now = datetime.datetime.utcnow().timestamp()
  56. last_time = last_message_times.get(user_id, 0)
  57. if now - last_time < 10:
  58. raise HTTPException(status_code=429, detail=locales.translate_error("flood_control", lang))
  59. last_message_times[user_id] = now
  60. if not message and not file:
  61. raise HTTPException(status_code=400, detail="Empty message")
  62. if message:
  63. message = message.strip()
  64. is_admin = (role == 'admin')
  65. if not is_admin:
  66. user_info = db.execute_query("SELECT can_chat FROM users WHERE id = %s", (user_id,))
  67. if not user_info or not user_info[0]['can_chat']:
  68. raise HTTPException(status_code=403, detail="Chat access disabled")
  69. order = db.execute_query("SELECT user_id FROM orders WHERE id = %s", (order_id,))
  70. if not order: raise HTTPException(status_code=404, detail="Order not found")
  71. if not is_admin and order[0]['user_id'] != user_id: raise HTTPException(status_code=403, detail="Not authorized")
  72. final_message = ""
  73. db_file_path = None
  74. if file and file.filename:
  75. file_ext = os.path.splitext(file.filename)[1]
  76. unique_filename = f"{uuid.uuid4()}{file_ext}"
  77. file_path = os.path.join(config.CHAT_UPLOADS_DIR, unique_filename)
  78. db_file_path = f"uploads/chats/{unique_filename}"
  79. with open(file_path, "wb") as buffer:
  80. while chunk := file.file.read(8192):
  81. buffer.write(chunk)
  82. msg_obj = {}
  83. if db_file_path:
  84. msg_obj["image_url"] = db_file_path
  85. if message:
  86. msg_obj["text"] = message
  87. final_message = json.dumps(msg_obj)
  88. query = "INSERT INTO order_messages (order_id, user_id, is_from_admin, message) VALUES (%s, %s, %s, %s)"
  89. msg_id = db.execute_commit(query, (order_id, user_id, is_admin, final_message))
  90. now = datetime.datetime.utcnow().isoformat()
  91. await manager.broadcast_to_order(order_id, {"id": msg_id, "is_from_admin": is_admin, "message": final_message, "created_at": now})
  92. if is_admin:
  93. await global_manager.notify_user(order[0]['user_id'])
  94. else:
  95. from services import event_hooks
  96. await global_manager.notify_admins()
  97. await global_manager.notify_admins_new_message(order_id, message or "Image uploaded")
  98. event_hooks.on_message_received(order_id, user_id, message or "Image uploaded")
  99. return {"id": msg_id, "status": "sent", "message": final_message}