Spaces:
Running
Running
File size: 18,177 Bytes
41b754a 4f5ecd1 41b754a 4f5ecd1 41b754a 4f5ecd1 41b754a 4f5ecd1 41b754a 4d9a0f9 41b754a 4d9a0f9 41b754a 4d9a0f9 41b754a 4d9a0f9 41b754a 4d9a0f9 41b754a 6a5e36e 41b754a 4f5ecd1 41b754a 4f5ecd1 41b754a 6a5e36e 41b754a 6a5e36e 41b754a 4f5ecd1 4d9a0f9 4f5ecd1 41b754a 4d9a0f9 df29e65 41b754a df29e65 41b754a df29e65 6a5e36e 4d9a0f9 6a5e36e | 1 2 3 4 5 6 7 8 9 10 11 12 13 14 15 16 17 18 19 20 21 22 23 24 25 26 27 28 29 30 31 32 33 34 35 36 37 38 39 40 41 42 43 44 45 46 47 48 49 50 51 52 53 54 55 56 57 58 59 60 61 62 63 64 65 66 67 68 69 70 71 72 73 74 75 76 77 78 79 80 81 82 83 84 85 86 87 88 89 90 91 92 93 94 95 96 97 98 99 100 101 102 103 104 105 106 107 108 109 110 111 112 113 114 115 116 117 118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 254 255 256 257 258 259 260 261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 | """Messaging module"""
import time
from db import get_service_client, new_client
from engineHelper import parse_db_timestamp
import engine_user
def execute_with_retry(query, retries=3, delay=0.1):
for i in range(retries):
try:
return query.execute()
except Exception as e:
if "10035" in str(e) and i < retries - 1:
import time
time.sleep(delay)
continue
raise e
# ========================
# CONVERSATIONS
# ========================
def create_conversation(users, is_group, name, admin_username=None):
sb = get_service_client()
avatar = "https://urgedwtmlfsmpkomptmb.supabase.co/storage/v1/object/public/media/23b32b8a24c54bb69068f7bc14db56da.png" if is_group else ""
is_group_int = 1 if is_group else 0
# Use first user as admin if not specified
actual_admin = admin_username if admin_username else (users[0] if users else "")
result = sb.table("conversations").insert({
"name": name,
"is_group": is_group_int,
"avatar": avatar,
"admin": actual_admin
}).execute()
conv_id = result.data[0]["id"]
for u in users:
sb.table("conversation_users").insert({"conversation_id": conv_id, "username": u}).execute()
return conv_id
def add_user_to_conversation(conversation_id, username):
sb = get_service_client()
# Check if already in
existing = sb.table("conversation_users").select("id").eq("conversation_id", conversation_id).eq("username", username).execute()
if existing.data: return False
sb.table("conversation_users").insert({"conversation_id": conversation_id, "username": username}).execute()
return True
def remove_user_from_conversation(conversation_id, username):
sb = get_service_client()
sb.table("conversation_users").delete().eq("conversation_id", conversation_id).eq("username", username).execute()
# Check if conversation is empty
remaining = sb.table("conversation_users").select("id").eq("conversation_id", conversation_id).execute()
if not remaining.data:
delete_conversation(conversation_id)
return True
def delete_conversation(conversation_id):
sb = get_service_client()
# Cascading deletes handle related messages, conversation_users, etc.
sb.table("conversations").delete().eq("id", conversation_id).execute()
def soft_delete_conversation(conversation_id, username):
sb = get_service_client()
res = sb.table("conversations").select("name, is_group").eq("id", conversation_id).execute()
if not res.data: return
conv = res.data[0]
if int(conv.get("is_group", 0)) == 1:
# Group chats are hard-deleted by admin
delete_conversation(conversation_id)
return
curr_name = str(conv.get("name") or "")
if not curr_name.startswith("deleted:"):
# First person deleting
sb.table("conversations").update({"name": f"deleted:{username}"}).eq("id", conversation_id).execute()
else:
# Second person deleting? Check if it's the other user
already_deleted_by = curr_name.replace("deleted:", "").split(",")
if username not in already_deleted_by:
already_deleted_by.append(username)
# If both participants deleted, hard delete
# (Assuming 1-on-1 has exactly 2 members)
users = get_conversation_users(conversation_id)
if all(u in already_deleted_by for u in users):
delete_conversation(conversation_id)
else:
sb.table("conversations").update({"name": f"deleted:{','.join(already_deleted_by)}"}).eq("id", conversation_id).execute()
def get_conversation_avatar(conv_id):
sb = get_service_client()
rows = sb.table("conversations").select("avatar").eq("id", conv_id).execute()
if not rows.data: return ""
return rows.data[0]["avatar"] or ""
def update_conversation_avatar(conv_id, avatar):
sb = get_service_client()
sb.table("conversations").update({"avatar": avatar}).eq("id", conv_id).execute()
def update_conversation_name(conv_id, name):
sb = get_service_client()
sb.table("conversations").update({"name": name}).eq("id", conv_id).execute()
def get_user_conversations(username):
sb = get_service_client()
# Fetch conversation_users and join with conversations in one request
rows = sb.table("conversation_users").select("conversation_id, conversations(id, name, is_group)").eq("username", username).execute()
restricted_users = engine_user.get_blocking_relationship_usernames(username)
results = []
for row in rows.data:
cid = row["conversation_id"]
c = row["conversations"]
if not c: continue
is_group = int(c["is_group"])
if is_group == 1:
final_name = c["name"]
else:
# For 1-on-1, filter if deleted by this user
if str(c.get("name", "")).startswith("deleted:"):
deleted_by = c["name"].replace("deleted:", "").split(",")
if username in deleted_by:
continue
# Find the other person
other_rows = sb.table("conversation_users").select("username").eq("conversation_id", cid).neq("username", username).execute()
other_user = other_rows.data[0]["username"] if other_rows.data else username
# Filter if other user is blocked
if other_user in restricted_users:
continue
final_name = other_user
results.append(f"{cid}|{final_name}|{is_group}")
return results
def get_conversation_users(conversation_id):
sb = get_service_client()
rows = sb.table("conversation_users").select("username").eq("conversation_id", conversation_id).order("id", desc=False).execute()
return [r["username"] for r in rows.data]
def get_conversation_details(conversation_id):
sb = get_service_client()
# Fetch conv info
c_res = sb.table("conversations").select("*").eq("id", conversation_id).execute()
if not c_res.data:
return None
conv = c_res.data[0]
# Fetch users
users = get_conversation_users(conversation_id)
return {
"id": str(conv["id"]),
"name": conv["name"],
"is_group": str(conv["is_group"]),
"avatar": conv["avatar"],
"users": users,
"admin": conv.get("admin") or (users[0] if users else None)
}
def find_conversation(users, name):
if not users: return -1
sb = get_service_client()
sorted_users = sorted(users)
lower_name = name.lower() if name else ""
# Optimize: only check conversations that the first user is part of
res = sb.table("conversation_users").select("conversation_id").eq("username", sorted_users[0]).execute()
cids = [r["conversation_id"] for r in res.data]
if not cids: return -1
conv_rows = sb.table("conversations").select("id, is_group, name").in_("id", cids).execute()
for conv in conv_rows.data:
cid = conv["id"]
is_group = int(conv["is_group"])
conv_name = conv["name"]
# For direct chats (is_group=0), names might have 'deleted:' prefix
# We only care about the users matching
cu = sb.table("conversation_users").select("username").eq("conversation_id", cid).execute()
conv_users = sorted([r["username"] for r in cu.data])
if conv_users == sorted_users:
if is_group == 0 and not name:
return cid
if is_group == 1 and conv_name and conv_name.lower() == lower_name:
return cid
return -1
def restore_conversation(conversation_id, username):
"""Remove a user from the 'deleted:' list in the conversation name."""
sb = get_service_client()
res = sb.table("conversations").select("name, is_group").eq("id", conversation_id).execute()
if not res.data: return
conv = res.data[0]
if int(conv.get("is_group", 0)) == 1: return
curr_name = str(conv.get("name") or "")
if curr_name.startswith("deleted:"):
deleted_by = curr_name.replace("deleted:", "").split(",")
if username in deleted_by:
deleted_by.remove(username)
if not deleted_by:
new_name = ""
else:
new_name = f"deleted:{','.join(deleted_by)}"
sb.table("conversations").update({"name": new_name}).eq("id", conversation_id).execute()
def get_or_create_bitai_conversation(username):
BITAI_USER = 'BitAI'
users = sorted([username, BITAI_USER])
conv_id = find_conversation(users, '')
if conv_id > 0: return conv_id
return create_conversation(users, False, '')
def get_conversations_summary(username):
"""Parallel fetch of all conversation data needed for the chat list."""
sb = get_service_client()
# 1. Get conversation IDs for this user
cu_rows = sb.table("conversation_users").select("conversation_id").eq("username", username).execute()
if not cu_rows.data:
return []
conv_ids = [r["conversation_id"] for r in cu_rows.data]
# 2. Fetch data
convs_data = execute_with_retry(sb.table("conversations").select("id, is_group, name, avatar").in_("id", conv_ids)).data
all_cu_data = execute_with_retry(sb.table("conversation_users").select("conversation_id, username").in_("conversation_id", conv_ids)).data
msgs_data = execute_with_retry(sb.table("messages").select("conversation_id, sender_username, content, type, timestamp") \
.in_("conversation_id", conv_ids).order("timestamp", desc=True)).data
seen_data = execute_with_retry(sb.table("user_last_seen").select("conversation_id, last_seen_timestamp") \
.in_("conversation_id", conv_ids).eq("username", username)).data
restricted_users = engine_user.get_blocking_relationship_usernames(username)
conv_map = {c["id"]: c for c in convs_data}
users_map = {}
for r in all_cu_data:
users_map.setdefault(r["conversation_id"], []).append(r["username"])
# Collect other usernames needing avatars
other_users = set()
for cid in conv_ids:
c = conv_map.get(cid, {})
if not c: continue
if int(c.get("is_group", 0)) == 0:
members = users_map.get(cid, [])
other = next((u for u in members if u != username), None)
if other and other != "BitAI":
other_users.add(other)
other_users.add(username)
av_rows = execute_with_retry(sb.table("users").select("username, avatar").in_("username", list(other_users)))
avatar_map = {r["username"]: r["avatar"] or "" for r in av_rows.data}
# Fetch custom BitAI names for this user
bitai_rows = execute_with_retry(sb.table("bit_ai").select("ai_name").eq("owner_username", username)).data
custom_bitai_name = bitai_rows[0]["ai_name"] if bitai_rows else "BitAI"
last_msg_map = {}
for r in msgs_data:
cid = r["conversation_id"]
if cid not in last_msg_map:
last_msg_map[cid] = r
seen_map = {r["conversation_id"]: parse_db_timestamp(r["last_seen_timestamp"]) for r in seen_data}
results = []
for cid in conv_ids:
c = conv_map.get(cid)
if not c: continue
is_group = int(c.get("is_group", 0))
members = users_map.get(cid, [])
if is_group:
name = c["name"]
avatar = c.get("avatar") or "https://urgedwtmlfsmpkomptmb.supabase.co/storage/v1/object/public/media/23b32b8a24c54bb69068f7bc14db56da.png"
else:
if str(c.get("name", "")).startswith("deleted:"):
deleted_by = c["name"].replace("deleted:", "").split(",")
if username in deleted_by:
continue
other = next((u for u in members if u != username), username)
# Filter 1-on-1 with blocked users
if other in restricted_users:
continue
name = custom_bitai_name if other == "BitAI" else other
avatar = "psychology" if other == "BitAI" else avatar_map.get(other, "")
lm = last_msg_map.get(cid)
last_message = None
if lm:
last_message = {
"sender": lm["sender_username"] or "Deleted User",
"content": lm["content"] or "",
"type": lm["type"] or "text",
"media": "",
"time": parse_db_timestamp(lm["timestamp"]),
}
last_seen = seen_map.get(cid, 0)
has_unread = bool(last_message and last_message["time"] > last_seen and last_message["sender"] != username)
results.append({
"id": str(cid),
"name": name,
"isGroup": str(is_group),
"avatar": avatar,
"lastMessage": last_message,
"unread": has_unread,
"isBitAI": other == "BitAI" if not is_group else False
})
return results
# ========================
# MESSAGES
# ========================
def send_message(sender, conversation_id, content, msg_type, media_path="", reply_to=None):
sb = get_service_client()
# If this is a direct message and was deleted by someone, restore it for the sender
# (Recipient will see it again too if the prefix is cleared, but usually we want
# to ensure the sender can see their own sent message in the list)
restore_conversation(conversation_id, sender)
row = {"conversation_id": conversation_id, "sender_username": sender,
"content": content, "type": msg_type, "media_path": media_path}
if reply_to:
row["reply_to"] = reply_to
# Database handles timestamp automatically via DEFAULT NOW()
res = sb.table("messages").insert(row).execute()
if res.data:
return res.data[0]["id"]
return None
def get_messages(conversation_id):
sb = get_service_client()
rows = sb.table("messages").select("id, sender_username, content, type, media_path, timestamp, reply_to").eq("conversation_id", conversation_id).order("timestamp", desc=False).execute()
results = []
for r in rows.data:
results.append({
"id": r["id"],
"sender": r["sender_username"] or "Deleted User",
"content": r["content"],
"type": r["type"],
"media": r["media_path"] or "",
"time": parse_db_timestamp(r["timestamp"]),
"replyTo": r.get("reply_to")
})
return results
def get_last_message(conversation_id):
sb = get_service_client()
rows = sb.table("messages").select("sender_username, content, type, timestamp").eq("conversation_id", conversation_id).order("timestamp", desc=True).limit(1).execute()
if not rows.data: return ""
r = rows.data[0]
sender = r['sender_username'] or "Deleted User"
return f"{sender}|{r['content']}|{r['type']}||{r['timestamp']}"
def has_unread(conversation_id, username):
last_seen = get_last_interaction(conversation_id, username)
last_msg = get_last_activity(conversation_id)
return last_msg > last_seen
def get_last_activity(conversation_id):
sb = get_service_client()
rows = sb.table("messages").select("timestamp").eq("conversation_id", conversation_id).order("timestamp", desc=True).limit(1).execute()
if not rows.data: return 0
ts = rows.data[0]["timestamp"]
return parse_db_timestamp(ts)
def get_last_interaction(conversation_id, username):
sb = get_service_client()
rows = sb.table("user_last_seen").select("last_seen_timestamp").eq("conversation_id", conversation_id).eq("username", username).execute()
if not rows.data: return 0
ts = rows.data[0]["last_seen_timestamp"]
return parse_db_timestamp(ts)
def mark_seen(username, convo_id):
sb = get_service_client()
# Ensure conversation exists to avoid foreign key violation if it was recently deleted
conv_check = sb.table("conversations").select("id").eq("id", convo_id).execute()
if not conv_check.data:
return
from datetime import datetime, timezone
now_ts = datetime.now(timezone.utc).isoformat()
existing = sb.table("user_last_seen").select("id").eq("conversation_id", convo_id).eq("username", username).execute()
if existing.data:
sb.table("user_last_seen").update({"last_seen_timestamp": now_ts}).eq("conversation_id", convo_id).eq("username", username).execute()
else:
sb.table("user_last_seen").insert({"conversation_id": convo_id, "username": username, "last_seen_timestamp": now_ts}).execute()
# ========================
# BITAI SETTINGS
# ========================
def get_bit_ai_settings(username):
sb = get_service_client()
res = sb.table("bit_ai").select("*").eq("owner_username", username).execute()
if not res.data:
# Create default if doesn't exist
default_settings = {
"owner_username": username,
"ai_name": "BitAI",
"personality": "",
"is_initialized": False
}
try:
sb.table("bit_ai").insert(default_settings).execute()
except:
# Handle race condition
res = sb.table("bit_ai").select("*").eq("owner_username", username).execute()
if res.data: return res.data[0]
return default_settings
return res.data[0]
def update_bit_ai_settings(username, ai_name=None, personality=None, is_initialized=None):
sb = get_service_client()
update_data = {}
if ai_name is not None: update_data["ai_name"] = ai_name
if personality is not None: update_data["personality"] = personality
if is_initialized is not None: update_data["is_initialized"] = is_initialized
if not update_data: return
sb.table("bit_ai").upsert({"owner_username": username, **update_data}).execute()
def clear_bit_ai_messages(username):
sb = get_service_client()
conv_id = get_or_create_bitai_conversation(username)
# Delete all messages in this conversation
sb.table("messages").delete().eq("conversation_id", conv_id).execute()
return True
|