Spaces:
Running
Running
Update app.py
Browse files
app.py
CHANGED
|
@@ -165,6 +165,19 @@ async def init_database():
|
|
| 165 |
)
|
| 166 |
""")
|
| 167 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 168 |
# --- Additive migrations for databases created before v2 ---
|
| 169 |
msg_cols = await _table_columns(db, "messages")
|
| 170 |
if "conversation_id" not in msg_cols:
|
|
@@ -187,6 +200,46 @@ async def init_database():
|
|
| 187 |
if col not in usr_cols:
|
| 188 |
await _ensure_column(db, "users", col, ddl)
|
| 189 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 190 |
rc_cols = await _table_columns(db, "read_receipts")
|
| 191 |
if "read_at" not in rc_cols:
|
| 192 |
await _ensure_column(db, "read_receipts", "read_at",
|
|
@@ -221,8 +274,14 @@ async def init_database():
|
|
| 221 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_receipts_user_msg ON read_receipts(user_id, message_id)")
|
| 222 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_receipts_message ON read_receipts(message_id)")
|
| 223 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_members_user ON conversation_members(user_id)")
|
|
|
|
|
|
|
| 224 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_conv_dm_pair "
|
| 225 |
"ON conversations(type, user_low_id, user_high_id) WHERE type = 'dm'")
|
|
|
|
|
|
|
|
|
|
|
|
|
| 226 |
await db.commit()
|
| 227 |
|
| 228 |
upload_database(DATABASE_URL)
|
|
@@ -364,29 +423,123 @@ async def authenticate_user(token: str) -> Optional[dict]:
|
|
| 364 |
# ------------------------------------------------------------------------
|
| 365 |
# Conversation helpers
|
| 366 |
# ------------------------------------------------------------------------
|
| 367 |
-
async def
|
| 368 |
-
"""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 369 |
if cid == GLOBAL_CONVERSATION_ID:
|
| 370 |
-
return
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 371 |
cursor = await db.execute(
|
| 372 |
-
"SELECT
|
| 373 |
(cid, uid)
|
| 374 |
)
|
| 375 |
-
|
|
|
|
| 376 |
|
| 377 |
|
| 378 |
-
async def
|
| 379 |
-
"""
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 380 |
if cid == GLOBAL_CONVERSATION_ID:
|
| 381 |
-
# Global lobby membership is implicit for every account.
|
| 382 |
cursor = await db.execute("SELECT id FROM users")
|
| 383 |
return [r["id"] for r in await cursor.fetchall()]
|
| 384 |
cursor = await db.execute(
|
| 385 |
-
"SELECT user_id FROM conversation_members
|
|
|
|
|
|
|
| 386 |
)
|
| 387 |
return [r["user_id"] for r in await cursor.fetchall()]
|
| 388 |
|
| 389 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 390 |
async def _ensure_global_member(db: aiosqlite.Connection, uid: int):
|
| 391 |
"""First sighting of a user: mark the whole existing lobby as already-read."""
|
| 392 |
cursor = await db.execute(
|
|
@@ -397,7 +550,8 @@ async def _ensure_global_member(db: aiosqlite.Connection, uid: int):
|
|
| 397 |
max_id = row["max_id"] or 0
|
| 398 |
await db.execute(
|
| 399 |
"INSERT OR IGNORE INTO conversation_members "
|
| 400 |
-
"(conversation_id, user_id, last_read_message_id
|
|
|
|
| 401 |
(GLOBAL_CONVERSATION_ID, uid, max_id)
|
| 402 |
)
|
| 403 |
|
|
@@ -446,7 +600,7 @@ async def _peer_for(db: aiosqlite.Connection, cid: int, uid: int) -> Optional[di
|
|
| 446 |
async def _dm_number(db: aiosqlite.Connection, cid: int) -> int:
|
| 447 |
cursor = await db.execute(
|
| 448 |
"""SELECT id FROM conversations
|
| 449 |
-
WHERE type = 'dm' AND user_low_id = (
|
| 450 |
SELECT user_low_id FROM conversations WHERE id = ?
|
| 451 |
) AND user_high_id = (
|
| 452 |
SELECT user_high_id FROM conversations WHERE id = ?
|
|
@@ -472,14 +626,34 @@ async def _build_conversation_summary(
|
|
| 472 |
raise HTTPException(404, "Conversation not found")
|
| 473 |
|
| 474 |
online_ids = online_ids if online_ids is not None else set()
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 475 |
summary = {
|
| 476 |
"id": conv["id"],
|
| 477 |
"type": conv["type"],
|
| 478 |
-
"
|
|
|
|
|
|
|
| 479 |
"created_at": conv["created_at"],
|
| 480 |
"peer": None,
|
| 481 |
"peer_online": False,
|
| 482 |
-
"dm_number":
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 483 |
"last_message_id": None,
|
| 484 |
"last_message_preview": "",
|
| 485 |
"last_message_ts": None,
|
|
@@ -489,7 +663,24 @@ async def _build_conversation_summary(
|
|
| 489 |
"unread_count": 0,
|
| 490 |
}
|
| 491 |
|
| 492 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 493 |
peer = await _peer_for(db, cid, uid)
|
| 494 |
if peer:
|
| 495 |
summary["peer"] = {
|
|
@@ -502,6 +693,10 @@ async def _build_conversation_summary(
|
|
| 502 |
}
|
| 503 |
summary["peer_online"] = peer["id"] in online_ids
|
| 504 |
summary["dm_number"] = await _dm_number(db, cid)
|
|
|
|
|
|
|
|
|
|
|
|
|
| 505 |
|
| 506 |
# Newest visible message (encrypted in DB - decrypt only what we preview)
|
| 507 |
cursor = await db.execute(
|
|
@@ -522,28 +717,36 @@ async def _build_conversation_summary(
|
|
| 522 |
summary["last_sender_name"] = last["display_name"] or last["username"]
|
| 523 |
summary["last_message_type"] = last["file_type"] or "text"
|
| 524 |
|
| 525 |
-
# Unread = messages from other people past the point the user last read to
|
| 526 |
-
|
| 527 |
-
|
| 528 |
-
|
| 529 |
-
|
| 530 |
-
|
| 531 |
-
|
| 532 |
-
|
| 533 |
-
|
| 534 |
-
|
| 535 |
-
|
|
|
|
|
|
|
|
|
|
|
|
|
| 536 |
return summary
|
| 537 |
|
| 538 |
|
| 539 |
async def _conversations_for_user(db: aiosqlite.Connection, uid: int, online_ids=None) -> List[dict]:
|
| 540 |
-
"""Global lobby first, then the user's private chats newest-activity first.
|
|
|
|
|
|
|
|
|
|
| 541 |
cursor = await db.execute(
|
| 542 |
"""SELECT c.id FROM conversations c
|
| 543 |
WHERE c.id = ?
|
| 544 |
-
OR
|
| 545 |
SELECT 1 FROM conversation_members cm
|
| 546 |
-
WHERE cm.conversation_id = c.id AND cm.user_id = ?
|
|
|
|
| 547 |
ORDER BY c.id""",
|
| 548 |
(GLOBAL_CONVERSATION_ID, uid)
|
| 549 |
)
|
|
@@ -551,6 +754,10 @@ async def _conversations_for_user(db: aiosqlite.Connection, uid: int, online_ids
|
|
| 551 |
summaries = []
|
| 552 |
for cid in ids:
|
| 553 |
try:
|
|
|
|
|
|
|
|
|
|
|
|
|
| 554 |
summaries.append(await _build_conversation_summary(db, cid, uid, online_ids))
|
| 555 |
except HTTPException:
|
| 556 |
continue
|
|
@@ -558,6 +765,9 @@ async def _conversations_for_user(db: aiosqlite.Connection, uid: int, online_ids
|
|
| 558 |
def sort_key(s):
|
| 559 |
if s["id"] == GLOBAL_CONVERSATION_ID:
|
| 560 |
return (0, 0)
|
|
|
|
|
|
|
|
|
|
| 561 |
return (1, -(s["last_message_ts"] or 0))
|
| 562 |
|
| 563 |
return sorted(summaries, key=sort_key)
|
|
@@ -935,22 +1145,26 @@ async def _create_dm(db: aiosqlite.Connection, uid: int, target_id: int,
|
|
| 935 |
"""Create a private chat. Returns summary or None if invalid/at limit.
|
| 936 |
|
| 937 |
Up to MAX_PRIVATE_CHATS_PER_PAIR private chats are allowed per pair of users.
|
|
|
|
|
|
|
| 938 |
"""
|
| 939 |
if target_id == uid:
|
| 940 |
raise HTTPException(400, "You can't start a private chat with yourself")
|
| 941 |
|
| 942 |
async with _dm_create_lock:
|
| 943 |
cursor = await db.execute(
|
| 944 |
-
"SELECT id, username FROM users WHERE id = ?", (target_id,)
|
| 945 |
)
|
| 946 |
target = await cursor.fetchone()
|
| 947 |
if not target:
|
| 948 |
raise HTTPException(404, "User not found")
|
|
|
|
|
|
|
| 949 |
|
| 950 |
low, high = sorted([uid, target_id])
|
| 951 |
cursor = await db.execute(
|
| 952 |
"""SELECT COUNT(*) AS cnt FROM conversations
|
| 953 |
-
WHERE type = 'dm' AND user_low_id = ? AND user_high_id = ?""",
|
| 954 |
(low, high)
|
| 955 |
)
|
| 956 |
existing = (await cursor.fetchone())["cnt"]
|
|
@@ -963,15 +1177,15 @@ async def _create_dm(db: aiosqlite.Connection, uid: int, target_id: int,
|
|
| 963 |
)
|
| 964 |
|
| 965 |
cursor = await db.execute(
|
| 966 |
-
"""INSERT INTO conversations (type, title, created_by, user_low_id, user_high_id)
|
| 967 |
-
VALUES ('dm', '', ?, ?, ?)""",
|
| 968 |
(uid, low, high)
|
| 969 |
)
|
| 970 |
cid = cursor.lastrowid
|
| 971 |
now = int(time.time())
|
| 972 |
await db.executemany(
|
| 973 |
"INSERT OR IGNORE INTO conversation_members "
|
| 974 |
-
"(conversation_id, user_id, joined_at) VALUES (?, ?, ?)",
|
| 975 |
[(cid, uid, now), (cid, target_id, now)]
|
| 976 |
)
|
| 977 |
await db.commit()
|
|
@@ -1006,6 +1220,418 @@ async def create_dm_rest(
|
|
| 1006 |
await db.close()
|
| 1007 |
|
| 1008 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1009 |
# ------------------------------------------------------------------------
|
| 1010 |
# Messages (REST fallbacks for reliable edit/delete/receipts)
|
| 1011 |
# ------------------------------------------------------------------------
|
|
@@ -1036,6 +1662,8 @@ async def edit_message_rest(
|
|
| 1036 |
row = await _get_message_row(db, message_id)
|
| 1037 |
if not row or row["sender_id"] != user['id'] or row["is_deleted"]:
|
| 1038 |
raise HTTPException(404, "Message not found")
|
|
|
|
|
|
|
| 1039 |
new_enc = encrypt_message(content)
|
| 1040 |
await db.execute(
|
| 1041 |
"UPDATE messages SET encrypted_content = ?, is_edited = 1 WHERE id = ?",
|
|
@@ -1069,8 +1697,15 @@ async def delete_message_rest(
|
|
| 1069 |
db = await get_db()
|
| 1070 |
try:
|
| 1071 |
row = await _get_message_row(db, message_id)
|
| 1072 |
-
if not row
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1073 |
raise HTTPException(404, "Message not found")
|
|
|
|
|
|
|
| 1074 |
|
| 1075 |
if row["file_path"]:
|
| 1076 |
try:
|
|
@@ -1295,6 +1930,23 @@ async def download_file(
|
|
| 1295 |
if '..' in file_path or '\\' in file_path:
|
| 1296 |
raise HTTPException(400, "Invalid path")
|
| 1297 |
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
| 1298 |
try:
|
| 1299 |
# AES-GCM needs the whole ciphertext to verify its auth tag, so one full
|
| 1300 |
# read is required for security - the response is streamed in chunks.
|
|
@@ -1940,4 +2592,4 @@ async def ws_endpoint(ws: WebSocket, token: str = Query(...)):
|
|
| 1940 |
|
| 1941 |
@app.get("/")
|
| 1942 |
async def root():
|
| 1943 |
-
return FileResponse("static/index.html")
|
|
|
|
| 165 |
)
|
| 166 |
""")
|
| 167 |
|
| 168 |
+
# Blocking: one-way rows (A blocks B). Chat creation/access is restricted
|
| 169 |
+
# when either user in a pair has blocked the other.
|
| 170 |
+
await db.execute("""
|
| 171 |
+
CREATE TABLE IF NOT EXISTS user_blocks (
|
| 172 |
+
blocker_id INTEGER NOT NULL,
|
| 173 |
+
blocked_id INTEGER NOT NULL,
|
| 174 |
+
created_at INTEGER DEFAULT (strftime('%s','now')),
|
| 175 |
+
PRIMARY KEY (blocker_id, blocked_id),
|
| 176 |
+
FOREIGN KEY(blocker_id) REFERENCES users(id) ON DELETE CASCADE,
|
| 177 |
+
FOREIGN KEY(blocked_id) REFERENCES users(id) ON DELETE CASCADE
|
| 178 |
+
)
|
| 179 |
+
""")
|
| 180 |
+
|
| 181 |
# --- Additive migrations for databases created before v2 ---
|
| 182 |
msg_cols = await _table_columns(db, "messages")
|
| 183 |
if "conversation_id" not in msg_cols:
|
|
|
|
| 200 |
if col not in usr_cols:
|
| 201 |
await _ensure_column(db, "users", col, ddl)
|
| 202 |
|
| 203 |
+
conv_cols = await _table_columns(db, "conversations")
|
| 204 |
+
for col, ddl in {
|
| 205 |
+
"is_group": "INTEGER NOT NULL DEFAULT 0",
|
| 206 |
+
"custom_name": "TEXT",
|
| 207 |
+
}.items():
|
| 208 |
+
if col not in conv_cols:
|
| 209 |
+
await _ensure_column(db, "conversations", col, ddl)
|
| 210 |
+
|
| 211 |
+
cm_cols = await _table_columns(db, "conversation_members")
|
| 212 |
+
for col, ddl in {
|
| 213 |
+
"status": "TEXT NOT NULL DEFAULT 'accepted'",
|
| 214 |
+
"role": "TEXT NOT NULL DEFAULT 'member'",
|
| 215 |
+
}.items():
|
| 216 |
+
if col not in cm_cols:
|
| 217 |
+
await _ensure_column(db, "conversation_members", col, ddl)
|
| 218 |
+
|
| 219 |
+
# Existing memberships predate invites/roles: make them accepted.
|
| 220 |
+
await db.execute(
|
| 221 |
+
"UPDATE conversation_members SET status = 'accepted' "
|
| 222 |
+
"WHERE status IS NULL OR status = '' OR status = 'pending'"
|
| 223 |
+
)
|
| 224 |
+
await db.execute(
|
| 225 |
+
"UPDATE conversation_members SET role = 'member' "
|
| 226 |
+
"WHERE role IS NULL OR role = ''"
|
| 227 |
+
)
|
| 228 |
+
# Group owner is the first accepted member of each group with no owner.
|
| 229 |
+
await db.execute("""
|
| 230 |
+
UPDATE conversation_members
|
| 231 |
+
SET role = 'owner'
|
| 232 |
+
WHERE conversation_id IN (
|
| 233 |
+
SELECT id FROM conversations WHERE is_group = 1
|
| 234 |
+
)
|
| 235 |
+
AND user_id = (
|
| 236 |
+
SELECT user_id FROM conversation_members cm
|
| 237 |
+
WHERE cm.conversation_id = conversation_members.conversation_id
|
| 238 |
+
AND cm.status = 'accepted'
|
| 239 |
+
ORDER BY cm.joined_at ASC, cm.user_id ASC LIMIT 1
|
| 240 |
+
)
|
| 241 |
+
""")
|
| 242 |
+
|
| 243 |
rc_cols = await _table_columns(db, "read_receipts")
|
| 244 |
if "read_at" not in rc_cols:
|
| 245 |
await _ensure_column(db, "read_receipts", "read_at",
|
|
|
|
| 274 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_receipts_user_msg ON read_receipts(user_id, message_id)")
|
| 275 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_receipts_message ON read_receipts(message_id)")
|
| 276 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_members_user ON conversation_members(user_id)")
|
| 277 |
+
await db.execute("CREATE INDEX IF NOT EXISTS idx_members_status "
|
| 278 |
+
"ON conversation_members(user_id, status)")
|
| 279 |
await db.execute("CREATE INDEX IF NOT EXISTS idx_conv_dm_pair "
|
| 280 |
"ON conversations(type, user_low_id, user_high_id) WHERE type = 'dm'")
|
| 281 |
+
await db.execute("CREATE INDEX IF NOT EXISTS idx_blocks_blocked "
|
| 282 |
+
"ON user_blocks(blocked_id)")
|
| 283 |
+
await db.execute("CREATE INDEX IF NOT EXISTS idx_blocks_blocker "
|
| 284 |
+
"ON user_blocks(blocker_id)")
|
| 285 |
await db.commit()
|
| 286 |
|
| 287 |
upload_database(DATABASE_URL)
|
|
|
|
| 423 |
# ------------------------------------------------------------------------
|
| 424 |
# Conversation helpers
|
| 425 |
# ------------------------------------------------------------------------
|
| 426 |
+
async def _block_exists(db: aiosqlite.Connection, a: int, b: int) -> bool:
|
| 427 |
+
"""True when blocking exists in either direction (mutual restriction)."""
|
| 428 |
+
if a == b:
|
| 429 |
+
return False
|
| 430 |
+
cursor = await db.execute(
|
| 431 |
+
"SELECT 1 FROM user_blocks WHERE (blocker_id = ? AND blocked_id = ?) "
|
| 432 |
+
"OR (blocker_id = ? AND blocked_id = ?) LIMIT 1",
|
| 433 |
+
(a, b, b, a)
|
| 434 |
+
)
|
| 435 |
+
return await cursor.fetchone() is not None
|
| 436 |
+
|
| 437 |
+
|
| 438 |
+
async def _conversation_blocked(db: aiosqlite.Connection, cid: int, uid: int) -> bool:
|
| 439 |
+
"""A non-global conversation is hidden when the user blocks (or is blocked by)
|
| 440 |
+
any other accepted member."""
|
| 441 |
if cid == GLOBAL_CONVERSATION_ID:
|
| 442 |
+
return False
|
| 443 |
+
member_ids = await _raw_conversation_member_ids(db, cid)
|
| 444 |
+
for mid in member_ids:
|
| 445 |
+
if mid != uid and await _block_exists(db, uid, mid):
|
| 446 |
+
return True
|
| 447 |
+
return False
|
| 448 |
+
|
| 449 |
+
|
| 450 |
+
async def _membership(db: aiosqlite.Connection, cid: int, uid: int) -> Optional[dict]:
|
| 451 |
cursor = await db.execute(
|
| 452 |
+
"SELECT * FROM conversation_members WHERE conversation_id = ? AND user_id = ?",
|
| 453 |
(cid, uid)
|
| 454 |
)
|
| 455 |
+
row = await cursor.fetchone()
|
| 456 |
+
return dict(row) if row else None
|
| 457 |
|
| 458 |
|
| 459 |
+
async def user_in_conversation(db: aiosqlite.Connection, uid: int, cid: int) -> bool:
|
| 460 |
+
"""Everyone is in the global lobby; other chats require an accepted membership
|
| 461 |
+
and no active block with another member."""
|
| 462 |
+
if cid == GLOBAL_CONVERSATION_ID:
|
| 463 |
+
return True
|
| 464 |
+
member = await _membership(db, cid, uid)
|
| 465 |
+
if not member or member.get("status") != "accepted":
|
| 466 |
+
return False
|
| 467 |
+
if await _conversation_blocked(db, cid, uid):
|
| 468 |
+
return False
|
| 469 |
+
return True
|
| 470 |
+
|
| 471 |
+
|
| 472 |
+
async def _raw_conversation_member_ids(db: aiosqlite.Connection, cid: int) -> List[int]:
|
| 473 |
if cid == GLOBAL_CONVERSATION_ID:
|
|
|
|
| 474 |
cursor = await db.execute("SELECT id FROM users")
|
| 475 |
return [r["id"] for r in await cursor.fetchall()]
|
| 476 |
cursor = await db.execute(
|
| 477 |
+
"SELECT user_id FROM conversation_members "
|
| 478 |
+
"WHERE conversation_id = ? AND status = 'accepted'",
|
| 479 |
+
(cid,)
|
| 480 |
)
|
| 481 |
return [r["user_id"] for r in await cursor.fetchall()]
|
| 482 |
|
| 483 |
|
| 484 |
+
async def _blocked_member_ids(db: aiosqlite.Connection, cid: int) -> set:
|
| 485 |
+
"""Members of a non-global conversation who are in a mutual block pair with
|
| 486 |
+
another accepted member. Those members are excluded from live broadcasts so
|
| 487 |
+
hidden/blocked chats never leak messages."""
|
| 488 |
+
if cid == GLOBAL_CONVERSATION_ID:
|
| 489 |
+
return set()
|
| 490 |
+
member_ids = await _raw_conversation_member_ids(db, cid)
|
| 491 |
+
blocked = set()
|
| 492 |
+
for mid in member_ids:
|
| 493 |
+
for other in member_ids:
|
| 494 |
+
if mid != other and await _block_exists(db, mid, other):
|
| 495 |
+
blocked.add(mid)
|
| 496 |
+
blocked.add(other)
|
| 497 |
+
return blocked
|
| 498 |
+
|
| 499 |
+
|
| 500 |
+
async def conversation_member_ids(db: aiosqlite.Connection, cid: int) -> List[int]:
|
| 501 |
+
"""Accepted member ids of a conversation (for scoped broadcasts).
|
| 502 |
+
|
| 503 |
+
Global broadcasts include everyone. Non-global broadcasts exclude users who
|
| 504 |
+
are in a block pair so hidden chats never leak messages."""
|
| 505 |
+
member_ids = await _raw_conversation_member_ids(db, cid)
|
| 506 |
+
if cid == GLOBAL_CONVERSATION_ID:
|
| 507 |
+
return member_ids
|
| 508 |
+
blocked = await _blocked_member_ids(db, cid)
|
| 509 |
+
return [mid for mid in member_ids if mid not in blocked]
|
| 510 |
+
|
| 511 |
+
|
| 512 |
+
async def _group_update_recipient_ids(db: aiosqlite.Connection, cid: int) -> List[int]:
|
| 513 |
+
"""Recipients for group summary/update/deleted events.
|
| 514 |
+
|
| 515 |
+
Includes accepted members and pending invitees, but still excludes anyone
|
| 516 |
+
who has an active block with another member of the group so hidden chats
|
| 517 |
+
never leak through live events."""
|
| 518 |
+
if cid == GLOBAL_CONVERSATION_ID:
|
| 519 |
+
return []
|
| 520 |
+
ids = await _raw_conversation_member_ids(db, cid)
|
| 521 |
+
cursor = await db.execute(
|
| 522 |
+
"SELECT user_id FROM conversation_members WHERE conversation_id = ? AND status = 'pending'",
|
| 523 |
+
(cid,)
|
| 524 |
+
)
|
| 525 |
+
ids += [r["user_id"] for r in await cursor.fetchall()]
|
| 526 |
+
return [mid for mid in ids if not await _conversation_blocked(db, cid, mid)]
|
| 527 |
+
|
| 528 |
+
|
| 529 |
+
async def _conversation_member_rows(db: aiosqlite.Connection, cid: int, include_pending: bool = False) -> List[dict]:
|
| 530 |
+
cond = "" if include_pending else "AND cm.status = 'accepted'"
|
| 531 |
+
cursor = await db.execute(
|
| 532 |
+
f"""SELECT u.id, u.username, u.display_name, u.avatar_path,
|
| 533 |
+
cm.status, cm.role, cm.joined_at, u.status AS user_status,
|
| 534 |
+
u.last_seen
|
| 535 |
+
FROM conversation_members cm JOIN users u ON u.id = cm.user_id
|
| 536 |
+
WHERE cm.conversation_id = ? {cond}
|
| 537 |
+
ORDER BY cm.joined_at ASC, u.id ASC""",
|
| 538 |
+
(cid,)
|
| 539 |
+
)
|
| 540 |
+
return [dict(r) for r in await cursor.fetchall()]
|
| 541 |
+
|
| 542 |
+
|
| 543 |
async def _ensure_global_member(db: aiosqlite.Connection, uid: int):
|
| 544 |
"""First sighting of a user: mark the whole existing lobby as already-read."""
|
| 545 |
cursor = await db.execute(
|
|
|
|
| 550 |
max_id = row["max_id"] or 0
|
| 551 |
await db.execute(
|
| 552 |
"INSERT OR IGNORE INTO conversation_members "
|
| 553 |
+
"(conversation_id, user_id, last_read_message_id, status, role) "
|
| 554 |
+
"VALUES (?, ?, ?, 'accepted', 'member')",
|
| 555 |
(GLOBAL_CONVERSATION_ID, uid, max_id)
|
| 556 |
)
|
| 557 |
|
|
|
|
| 600 |
async def _dm_number(db: aiosqlite.Connection, cid: int) -> int:
|
| 601 |
cursor = await db.execute(
|
| 602 |
"""SELECT id FROM conversations
|
| 603 |
+
WHERE type = 'dm' AND is_group = 0 AND user_low_id = (
|
| 604 |
SELECT user_low_id FROM conversations WHERE id = ?
|
| 605 |
) AND user_high_id = (
|
| 606 |
SELECT user_high_id FROM conversations WHERE id = ?
|
|
|
|
| 626 |
raise HTTPException(404, "Conversation not found")
|
| 627 |
|
| 628 |
online_ids = online_ids if online_ids is not None else set()
|
| 629 |
+
membership = await _membership(db, cid, uid) or {}
|
| 630 |
+
is_group = bool(conv["is_group"] or 0) if "is_group" in conv.keys() else False
|
| 631 |
+
|
| 632 |
+
# Display title: custom rename wins for DMs/groups, otherwise group title
|
| 633 |
+
# or the name of the other person in a DM.
|
| 634 |
+
title = conv["title"] or "Chat"
|
| 635 |
+
custom_name = conv["custom_name"] if "custom_name" in conv.keys() else None
|
| 636 |
+
dm_number = 1
|
| 637 |
+
if custom_name:
|
| 638 |
+
title = custom_name
|
| 639 |
+
|
| 640 |
summary = {
|
| 641 |
"id": conv["id"],
|
| 642 |
"type": conv["type"],
|
| 643 |
+
"is_group": is_group,
|
| 644 |
+
"title": title,
|
| 645 |
+
"custom_name": custom_name,
|
| 646 |
"created_at": conv["created_at"],
|
| 647 |
"peer": None,
|
| 648 |
"peer_online": False,
|
| 649 |
+
"dm_number": dm_number,
|
| 650 |
+
"status": membership.get("status", "accepted"),
|
| 651 |
+
"role": membership.get("role", "member"),
|
| 652 |
+
"is_owner": membership.get("role") == "owner",
|
| 653 |
+
"is_pending": membership.get("status") == "pending",
|
| 654 |
+
"member_count": 0,
|
| 655 |
+
"pending_member_count": 0,
|
| 656 |
+
"members": [],
|
| 657 |
"last_message_id": None,
|
| 658 |
"last_message_preview": "",
|
| 659 |
"last_message_ts": None,
|
|
|
|
| 663 |
"unread_count": 0,
|
| 664 |
}
|
| 665 |
|
| 666 |
+
members = await _conversation_member_rows(db, cid, include_pending=is_group)
|
| 667 |
+
summary["member_count"] = sum(1 for m in members if m["status"] == "accepted")
|
| 668 |
+
summary["pending_member_count"] = sum(1 for m in members if m["status"] == "pending")
|
| 669 |
+
public_members = []
|
| 670 |
+
for m in members:
|
| 671 |
+
public_members.append({
|
| 672 |
+
"id": m["id"],
|
| 673 |
+
"username": m["username"],
|
| 674 |
+
"display_name": m["display_name"],
|
| 675 |
+
"avatar_path": m["avatar_path"],
|
| 676 |
+
"online": m["id"] in online_ids,
|
| 677 |
+
"last_seen": m["last_seen"],
|
| 678 |
+
"status": m["status"],
|
| 679 |
+
"role": m["role"],
|
| 680 |
+
})
|
| 681 |
+
summary["members"] = public_members
|
| 682 |
+
|
| 683 |
+
if not is_group and conv["type"] == "dm":
|
| 684 |
peer = await _peer_for(db, cid, uid)
|
| 685 |
if peer:
|
| 686 |
summary["peer"] = {
|
|
|
|
| 693 |
}
|
| 694 |
summary["peer_online"] = peer["id"] in online_ids
|
| 695 |
summary["dm_number"] = await _dm_number(db, cid)
|
| 696 |
+
if not custom_name:
|
| 697 |
+
title = peer["display_name"] or peer["username"]
|
| 698 |
+
summary["title"] = title
|
| 699 |
+
# If no custom name and no peer (unlikely), keep generic title.
|
| 700 |
|
| 701 |
# Newest visible message (encrypted in DB - decrypt only what we preview)
|
| 702 |
cursor = await db.execute(
|
|
|
|
| 717 |
summary["last_sender_name"] = last["display_name"] or last["username"]
|
| 718 |
summary["last_message_type"] = last["file_type"] or "text"
|
| 719 |
|
| 720 |
+
# Unread = messages from other people past the point the user last read to.
|
| 721 |
+
# Pending invites do not count as unread messages.
|
| 722 |
+
if summary["is_pending"]:
|
| 723 |
+
summary["unread_count"] = 0
|
| 724 |
+
else:
|
| 725 |
+
cursor = await db.execute(
|
| 726 |
+
"""SELECT COUNT(*) AS cnt FROM messages m
|
| 727 |
+
WHERE m.conversation_id = ? AND m.sender_id != ? AND m.is_deleted = 0
|
| 728 |
+
AND m.id > (SELECT COALESCE(last_read_message_id, 0)
|
| 729 |
+
FROM conversation_members
|
| 730 |
+
WHERE conversation_id = ? AND user_id = ?)""",
|
| 731 |
+
(cid, uid, cid, uid)
|
| 732 |
+
)
|
| 733 |
+
unread = await cursor.fetchone()
|
| 734 |
+
summary["unread_count"] = unread["cnt"] or 0
|
| 735 |
return summary
|
| 736 |
|
| 737 |
|
| 738 |
async def _conversations_for_user(db: aiosqlite.Connection, uid: int, online_ids=None) -> List[dict]:
|
| 739 |
+
"""Global lobby first, then the user's private chats + groups newest-activity first.
|
| 740 |
+
|
| 741 |
+
Pending group invites are included so clients can accept/reject them. Active
|
| 742 |
+
blocks hide the conversation entirely (but leave the underlying rows intact)."""
|
| 743 |
cursor = await db.execute(
|
| 744 |
"""SELECT c.id FROM conversations c
|
| 745 |
WHERE c.id = ?
|
| 746 |
+
OR EXISTS (
|
| 747 |
SELECT 1 FROM conversation_members cm
|
| 748 |
+
WHERE cm.conversation_id = c.id AND cm.user_id = ?
|
| 749 |
+
AND cm.status IN ('accepted', 'pending'))
|
| 750 |
ORDER BY c.id""",
|
| 751 |
(GLOBAL_CONVERSATION_ID, uid)
|
| 752 |
)
|
|
|
|
| 754 |
summaries = []
|
| 755 |
for cid in ids:
|
| 756 |
try:
|
| 757 |
+
# Hide chats when either user has blocked the other; never delete them.
|
| 758 |
+
if _conversation_blocked is not None and cid != GLOBAL_CONVERSATION_ID \
|
| 759 |
+
and await _conversation_blocked(db, cid, uid):
|
| 760 |
+
continue
|
| 761 |
summaries.append(await _build_conversation_summary(db, cid, uid, online_ids))
|
| 762 |
except HTTPException:
|
| 763 |
continue
|
|
|
|
| 765 |
def sort_key(s):
|
| 766 |
if s["id"] == GLOBAL_CONVERSATION_ID:
|
| 767 |
return (0, 0)
|
| 768 |
+
# Pending invites float to the top so they are easy to accept.
|
| 769 |
+
if s.get("is_pending"):
|
| 770 |
+
return (0, 1)
|
| 771 |
return (1, -(s["last_message_ts"] or 0))
|
| 772 |
|
| 773 |
return sorted(summaries, key=sort_key)
|
|
|
|
| 1145 |
"""Create a private chat. Returns summary or None if invalid/at limit.
|
| 1146 |
|
| 1147 |
Up to MAX_PRIVATE_CHATS_PER_PAIR private chats are allowed per pair of users.
|
| 1148 |
+
Blocked users cannot create or be added to a private chat; the global room
|
| 1149 |
+
remains the only shared conversation.
|
| 1150 |
"""
|
| 1151 |
if target_id == uid:
|
| 1152 |
raise HTTPException(400, "You can't start a private chat with yourself")
|
| 1153 |
|
| 1154 |
async with _dm_create_lock:
|
| 1155 |
cursor = await db.execute(
|
| 1156 |
+
"SELECT id, username, display_name FROM users WHERE id = ?", (target_id,)
|
| 1157 |
)
|
| 1158 |
target = await cursor.fetchone()
|
| 1159 |
if not target:
|
| 1160 |
raise HTTPException(404, "User not found")
|
| 1161 |
+
if await _block_exists(db, uid, target_id):
|
| 1162 |
+
raise HTTPException(403, "You can't start a private chat with this user")
|
| 1163 |
|
| 1164 |
low, high = sorted([uid, target_id])
|
| 1165 |
cursor = await db.execute(
|
| 1166 |
"""SELECT COUNT(*) AS cnt FROM conversations
|
| 1167 |
+
WHERE type = 'dm' AND is_group = 0 AND user_low_id = ? AND user_high_id = ?""",
|
| 1168 |
(low, high)
|
| 1169 |
)
|
| 1170 |
existing = (await cursor.fetchone())["cnt"]
|
|
|
|
| 1177 |
)
|
| 1178 |
|
| 1179 |
cursor = await db.execute(
|
| 1180 |
+
"""INSERT INTO conversations (type, title, created_by, user_low_id, user_high_id, is_group)
|
| 1181 |
+
VALUES ('dm', '', ?, ?, ?, 0)""",
|
| 1182 |
(uid, low, high)
|
| 1183 |
)
|
| 1184 |
cid = cursor.lastrowid
|
| 1185 |
now = int(time.time())
|
| 1186 |
await db.executemany(
|
| 1187 |
"INSERT OR IGNORE INTO conversation_members "
|
| 1188 |
+
"(conversation_id, user_id, joined_at, status, role) VALUES (?, ?, ?, 'accepted', 'member')",
|
| 1189 |
[(cid, uid, now), (cid, target_id, now)]
|
| 1190 |
)
|
| 1191 |
await db.commit()
|
|
|
|
| 1220 |
await db.close()
|
| 1221 |
|
| 1222 |
|
| 1223 |
+
_group_create_lock = asyncio.Lock()
|
| 1224 |
+
_MAX_GROUP_NAME = 60
|
| 1225 |
+
|
| 1226 |
+
|
| 1227 |
+
async def _create_group(db: aiosqlite.Connection, creator_id: int, name: str,
|
| 1228 |
+
member_ids: List[int]) -> dict:
|
| 1229 |
+
if not name or not name.strip():
|
| 1230 |
+
raise HTTPException(400, "Group name is required")
|
| 1231 |
+
name = name.strip()[:_MAX_GROUP_NAME]
|
| 1232 |
+
member_ids = sorted({int(m) for m in member_ids if int(m) != creator_id})
|
| 1233 |
+
if not member_ids:
|
| 1234 |
+
raise HTTPException(400, "Select at least one member")
|
| 1235 |
+
if len(member_ids) > 50:
|
| 1236 |
+
raise HTTPException(400, "A group can have at most 50 members")
|
| 1237 |
+
|
| 1238 |
+
placeholders = ",".join("?" * len(member_ids))
|
| 1239 |
+
cursor = await db.execute(
|
| 1240 |
+
f"SELECT id, username, display_name FROM users WHERE id IN ({placeholders})",
|
| 1241 |
+
member_ids
|
| 1242 |
+
)
|
| 1243 |
+
users = {r["id"]: r for r in await cursor.fetchall()}
|
| 1244 |
+
for mid in member_ids:
|
| 1245 |
+
if mid not in users:
|
| 1246 |
+
raise HTTPException(404, "Could not find a member")
|
| 1247 |
+
if await _block_exists(db, creator_id, mid):
|
| 1248 |
+
raise HTTPException(403, f"You can't include a blocked user: {users[mid]['username']}")
|
| 1249 |
+
for i in range(len(member_ids)):
|
| 1250 |
+
for j in range(i + 1, len(member_ids)):
|
| 1251 |
+
if await _block_exists(db, member_ids[i], member_ids[j]):
|
| 1252 |
+
raise HTTPException(
|
| 1253 |
+
403,
|
| 1254 |
+
f"{users[member_ids[i]]['username']} and {users[member_ids[j]]['username']} "
|
| 1255 |
+
"have a block; they cannot be in the same group"
|
| 1256 |
+
)
|
| 1257 |
+
|
| 1258 |
+
async with _group_create_lock:
|
| 1259 |
+
cursor = await db.execute(
|
| 1260 |
+
"""INSERT INTO conversations (type, title, created_by, is_group)
|
| 1261 |
+
VALUES ('dm', ?, ?, 1)""",
|
| 1262 |
+
(name, creator_id)
|
| 1263 |
+
)
|
| 1264 |
+
cid = cursor.lastrowid
|
| 1265 |
+
now = int(time.time())
|
| 1266 |
+
await db.execute(
|
| 1267 |
+
"INSERT INTO conversation_members "
|
| 1268 |
+
"(conversation_id, user_id, joined_at, status, role) VALUES (?, ?, ?, 'accepted', 'owner')",
|
| 1269 |
+
(cid, creator_id, now)
|
| 1270 |
+
)
|
| 1271 |
+
await db.executemany(
|
| 1272 |
+
"INSERT OR IGNORE INTO conversation_members "
|
| 1273 |
+
"(conversation_id, user_id, joined_at, status, role) VALUES (?, ?, ?, 'pending', 'member')",
|
| 1274 |
+
[(cid, mid, now) for mid in member_ids]
|
| 1275 |
+
)
|
| 1276 |
+
await db.commit()
|
| 1277 |
+
schedule_db_sync()
|
| 1278 |
+
|
| 1279 |
+
online = {u["id"] for u in manager.get_online_users()}
|
| 1280 |
+
summary = await _build_conversation_summary(db, cid, creator_id, online)
|
| 1281 |
+
|
| 1282 |
+
# Invitees see a pending conversation so they can accept or decline.
|
| 1283 |
+
for mid in member_ids:
|
| 1284 |
+
invite = await _build_conversation_summary(db, cid, mid, online)
|
| 1285 |
+
await manager.send_to_user(mid, {"type": "conversation_updated", "conversation": invite})
|
| 1286 |
+
return summary
|
| 1287 |
+
|
| 1288 |
+
|
| 1289 |
+
@app.post("/api/conversations/group")
|
| 1290 |
+
async def create_group_rest(
|
| 1291 |
+
name: str = Query(..., max_length=_MAX_GROUP_NAME),
|
| 1292 |
+
member_ids: str = Query(...),
|
| 1293 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1294 |
+
):
|
| 1295 |
+
user = await authenticate_user(token)
|
| 1296 |
+
if not user:
|
| 1297 |
+
raise HTTPException(401)
|
| 1298 |
+
ids = [int(x) for x in member_ids.split(",") if x.strip()]
|
| 1299 |
+
db = await get_db()
|
| 1300 |
+
try:
|
| 1301 |
+
summary = await _create_group(db, user['id'], name, ids)
|
| 1302 |
+
await manager.send_to_user(user['id'], {"type": "dm_created", "conversation": summary})
|
| 1303 |
+
return {"conversation": summary}
|
| 1304 |
+
finally:
|
| 1305 |
+
await db.close()
|
| 1306 |
+
|
| 1307 |
+
|
| 1308 |
+
async def _accept_conversation_invite(db: aiosqlite.Connection, uid: int, cid: int) -> dict:
|
| 1309 |
+
member = await _membership(db, cid, uid)
|
| 1310 |
+
if not member or member.get("status") != "pending":
|
| 1311 |
+
raise HTTPException(404, "You don't have a pending invite to this group")
|
| 1312 |
+
if await _conversation_blocked(db, cid, uid):
|
| 1313 |
+
raise HTTPException(403, "This conversation is hidden because of a block")
|
| 1314 |
+
await db.execute(
|
| 1315 |
+
"UPDATE conversation_members SET status = 'accepted' WHERE conversation_id = ? AND user_id = ?",
|
| 1316 |
+
(cid, uid)
|
| 1317 |
+
)
|
| 1318 |
+
await db.commit()
|
| 1319 |
+
schedule_db_sync()
|
| 1320 |
+
return await _build_conversation_summary(db, cid, uid, {u["id"] for u in manager.get_online_users()})
|
| 1321 |
+
|
| 1322 |
+
|
| 1323 |
+
async def _reject_conversation_invite(db: aiosqlite.Connection, uid: int, cid: int) -> bool:
|
| 1324 |
+
member = await _membership(db, cid, uid)
|
| 1325 |
+
if not member or member.get("status") != "pending":
|
| 1326 |
+
raise HTTPException(404, "You don't have a pending invite to this group")
|
| 1327 |
+
await db.execute("DELETE FROM conversation_members WHERE conversation_id = ? AND user_id = ?",
|
| 1328 |
+
(cid, uid))
|
| 1329 |
+
await db.commit()
|
| 1330 |
+
schedule_db_sync()
|
| 1331 |
+
return True
|
| 1332 |
+
|
| 1333 |
+
|
| 1334 |
+
@app.post("/api/conversations/{conversation_id}/accept")
|
| 1335 |
+
async def accept_conversation_invite_rest(
|
| 1336 |
+
conversation_id: int,
|
| 1337 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1338 |
+
):
|
| 1339 |
+
user = await authenticate_user(token)
|
| 1340 |
+
if not user:
|
| 1341 |
+
raise HTTPException(401)
|
| 1342 |
+
db = await get_db()
|
| 1343 |
+
try:
|
| 1344 |
+
summary = await _accept_conversation_invite(db, user['id'], conversation_id)
|
| 1345 |
+
online = {u["id"] for u in manager.get_online_users()}
|
| 1346 |
+
members = await _group_update_recipient_ids(db, conversation_id)
|
| 1347 |
+
for mid in members:
|
| 1348 |
+
await manager.send_to_user(mid, {"type": "conversation_updated",
|
| 1349 |
+
"conversation": await _build_conversation_summary(
|
| 1350 |
+
db, conversation_id, mid, online)})
|
| 1351 |
+
return {"conversation": summary}
|
| 1352 |
+
finally:
|
| 1353 |
+
await db.close()
|
| 1354 |
+
|
| 1355 |
+
|
| 1356 |
+
@app.post("/api/conversations/{conversation_id}/reject")
|
| 1357 |
+
async def reject_conversation_invite_rest(
|
| 1358 |
+
conversation_id: int,
|
| 1359 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1360 |
+
):
|
| 1361 |
+
user = await authenticate_user(token)
|
| 1362 |
+
if not user:
|
| 1363 |
+
raise HTTPException(401)
|
| 1364 |
+
db = await get_db()
|
| 1365 |
+
try:
|
| 1366 |
+
await _reject_conversation_invite(db, user['id'], conversation_id)
|
| 1367 |
+
return {"status": "rejected"}
|
| 1368 |
+
finally:
|
| 1369 |
+
await db.close()
|
| 1370 |
+
|
| 1371 |
+
|
| 1372 |
+
@app.post("/api/conversations/{conversation_id}/members")
|
| 1373 |
+
async def add_group_members_rest(
|
| 1374 |
+
conversation_id: int,
|
| 1375 |
+
member_ids: str = Query(...),
|
| 1376 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1377 |
+
):
|
| 1378 |
+
user = await authenticate_user(token)
|
| 1379 |
+
if not user:
|
| 1380 |
+
raise HTTPException(401)
|
| 1381 |
+
db = await get_db()
|
| 1382 |
+
try:
|
| 1383 |
+
member = await _membership(db, conversation_id, user['id'])
|
| 1384 |
+
conv = await db.execute("SELECT * FROM conversations WHERE id = ?", (conversation_id,))
|
| 1385 |
+
row = await conv.fetchone()
|
| 1386 |
+
if not row or not row["is_group"]:
|
| 1387 |
+
raise HTTPException(400, "Only groups can have members added")
|
| 1388 |
+
if not member or member.get("status") != "accepted":
|
| 1389 |
+
raise HTTPException(403, "Accept the group invite first")
|
| 1390 |
+
if not await user_in_conversation(db, user['id'], conversation_id):
|
| 1391 |
+
raise HTTPException(403, "You can no longer access this conversation")
|
| 1392 |
+
|
| 1393 |
+
ids = [int(x) for x in member_ids.split(",") if x.strip()]
|
| 1394 |
+
existing = await _raw_conversation_member_ids(db, conversation_id)
|
| 1395 |
+
new_ids = [mid for mid in ids if mid != user['id']]
|
| 1396 |
+
if len(existing) + len(set(new_ids)) > 50:
|
| 1397 |
+
raise HTTPException(400, "A group can have at most 50 members")
|
| 1398 |
+
now = int(time.time())
|
| 1399 |
+
for i in range(len(ids)):
|
| 1400 |
+
for j in range(i + 1, len(ids)):
|
| 1401 |
+
if await _block_exists(db, ids[i], ids[j]):
|
| 1402 |
+
raise HTTPException(403, "Selected users have a block between them")
|
| 1403 |
+
for mid in ids:
|
| 1404 |
+
if not mid or mid == user['id']:
|
| 1405 |
+
continue
|
| 1406 |
+
if await _block_exists(db, user['id'], mid):
|
| 1407 |
+
raise HTTPException(403, "You can't add a blocked user")
|
| 1408 |
+
for existing_id in existing:
|
| 1409 |
+
if await _block_exists(db, mid, existing_id):
|
| 1410 |
+
raise HTTPException(
|
| 1411 |
+
403, "That user has a block with someone already in this group")
|
| 1412 |
+
await db.execute(
|
| 1413 |
+
"INSERT OR IGNORE INTO conversation_members "
|
| 1414 |
+
"(conversation_id, user_id, joined_at, status, role) VALUES (?, ?, ?, 'pending', 'member')",
|
| 1415 |
+
(conversation_id, mid, now)
|
| 1416 |
+
)
|
| 1417 |
+
await db.commit()
|
| 1418 |
+
schedule_db_sync()
|
| 1419 |
+
online = {u["id"] for u in manager.get_online_users()}
|
| 1420 |
+
summary = await _build_conversation_summary(db, conversation_id, user['id'], online)
|
| 1421 |
+
members = await _group_update_recipient_ids(db, conversation_id)
|
| 1422 |
+
for mid in ids:
|
| 1423 |
+
if mid == user['id'] or await _conversation_blocked(db, conversation_id, mid):
|
| 1424 |
+
continue
|
| 1425 |
+
invite = await _build_conversation_summary(db, conversation_id, mid, online)
|
| 1426 |
+
await manager.send_to_user(mid, {"type": "conversation_updated", "conversation": invite})
|
| 1427 |
+
for mid in members:
|
| 1428 |
+
await manager.send_to_user(mid, {"type": "conversation_updated",
|
| 1429 |
+
"conversation": await _build_conversation_summary(
|
| 1430 |
+
db, conversation_id, mid, online)})
|
| 1431 |
+
return {"conversation": summary}
|
| 1432 |
+
finally:
|
| 1433 |
+
await db.close()
|
| 1434 |
+
|
| 1435 |
+
|
| 1436 |
+
@app.patch("/api/conversations/{conversation_id}")
|
| 1437 |
+
async def rename_conversation_rest(
|
| 1438 |
+
conversation_id: int,
|
| 1439 |
+
name: str = Query(..., max_length=_MAX_GROUP_NAME),
|
| 1440 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1441 |
+
):
|
| 1442 |
+
user = await authenticate_user(token)
|
| 1443 |
+
if not user:
|
| 1444 |
+
raise HTTPException(401)
|
| 1445 |
+
name = name.strip()
|
| 1446 |
+
if not name:
|
| 1447 |
+
raise HTTPException(400, "Name cannot be empty")
|
| 1448 |
+
db = await get_db()
|
| 1449 |
+
try:
|
| 1450 |
+
if not await user_in_conversation(db, user['id'], conversation_id):
|
| 1451 |
+
raise HTTPException(403, "You can't rename this conversation")
|
| 1452 |
+
await db.execute(
|
| 1453 |
+
"UPDATE conversations SET custom_name = ?, title = ? WHERE id = ?",
|
| 1454 |
+
(name, name, conversation_id)
|
| 1455 |
+
)
|
| 1456 |
+
await db.commit()
|
| 1457 |
+
schedule_db_sync()
|
| 1458 |
+
online = {u["id"] for u in manager.get_online_users()}
|
| 1459 |
+
summary = await _build_conversation_summary(db, conversation_id, user['id'], online)
|
| 1460 |
+
members = await _group_update_recipient_ids(db, conversation_id)
|
| 1461 |
+
for mid in members:
|
| 1462 |
+
await manager.send_to_user(mid, {"type": "conversation_updated",
|
| 1463 |
+
"conversation": await _build_conversation_summary(
|
| 1464 |
+
db, conversation_id, mid, online)})
|
| 1465 |
+
return {"conversation": summary}
|
| 1466 |
+
finally:
|
| 1467 |
+
await db.close()
|
| 1468 |
+
|
| 1469 |
+
|
| 1470 |
+
@app.delete("/api/conversations/{conversation_id}")
|
| 1471 |
+
async def delete_or_leave_conversation_rest(
|
| 1472 |
+
conversation_id: int,
|
| 1473 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1474 |
+
):
|
| 1475 |
+
user = await authenticate_user(token)
|
| 1476 |
+
if not user:
|
| 1477 |
+
raise HTTPException(401)
|
| 1478 |
+
db = await get_db()
|
| 1479 |
+
try:
|
| 1480 |
+
member = await _membership(db, conversation_id, user['id'])
|
| 1481 |
+
cursor = await db.execute("SELECT * FROM conversations WHERE id = ?", (conversation_id,))
|
| 1482 |
+
conv = await cursor.fetchone()
|
| 1483 |
+
if not conv:
|
| 1484 |
+
raise HTTPException(404, "Conversation not found")
|
| 1485 |
+
if conv["is_group"] != 1:
|
| 1486 |
+
raise HTTPException(400, "Only group chats can be left or deleted")
|
| 1487 |
+
if not member or member.get("status") != "accepted":
|
| 1488 |
+
raise HTTPException(403, "You are not a member of this group")
|
| 1489 |
+
|
| 1490 |
+
notify_ids = await _group_update_recipient_ids(db, conversation_id)
|
| 1491 |
+
# Owner deletes the group permanently; members leave the group.
|
| 1492 |
+
if member and member.get("role") == "owner" and member.get("status") == "accepted":
|
| 1493 |
+
await db.execute("DELETE FROM read_receipts WHERE message_id IN "
|
| 1494 |
+
"(SELECT id FROM messages WHERE conversation_id = ?)", (conversation_id,))
|
| 1495 |
+
await db.execute("DELETE FROM messages WHERE conversation_id = ?", (conversation_id,))
|
| 1496 |
+
await db.execute("DELETE FROM conversation_members WHERE conversation_id = ?", (conversation_id,))
|
| 1497 |
+
await db.execute("DELETE FROM conversations WHERE id = ?", (conversation_id,))
|
| 1498 |
+
await db.commit()
|
| 1499 |
+
schedule_db_sync()
|
| 1500 |
+
for mid in notify_ids:
|
| 1501 |
+
await manager.send_to_user(mid, {"type": "conversation_deleted",
|
| 1502 |
+
"conversation_id": conversation_id})
|
| 1503 |
+
return {"status": "deleted", "conversation_id": conversation_id}
|
| 1504 |
+
|
| 1505 |
+
# Non-owner: leave (removes membership, keeps the group + history).
|
| 1506 |
+
await db.execute("DELETE FROM conversation_members WHERE conversation_id = ? AND user_id = ?",
|
| 1507 |
+
(conversation_id, user['id']))
|
| 1508 |
+
await db.commit()
|
| 1509 |
+
schedule_db_sync()
|
| 1510 |
+
# Only notify the members who remain in the group.
|
| 1511 |
+
remaining = await conversation_member_ids(db, conversation_id)
|
| 1512 |
+
for mid in remaining:
|
| 1513 |
+
await manager.send_to_user(mid, {"type": "conversation_left",
|
| 1514 |
+
"conversation_id": conversation_id,
|
| 1515 |
+
"user_id": user['id'],
|
| 1516 |
+
"conversation": await _build_conversation_summary(
|
| 1517 |
+
db, conversation_id, mid,
|
| 1518 |
+
{u["id"] for u in manager.get_online_users()})})
|
| 1519 |
+
return {"status": "left", "conversation_id": conversation_id}
|
| 1520 |
+
finally:
|
| 1521 |
+
await db.close()
|
| 1522 |
+
|
| 1523 |
+
|
| 1524 |
+
# ------------------------------------------------------------------------
|
| 1525 |
+
# Users: search + blocking
|
| 1526 |
+
# ------------------------------------------------------------------------
|
| 1527 |
+
@app.get("/api/users")
|
| 1528 |
+
async def list_users(
|
| 1529 |
+
query: str = Query("", max_length=100),
|
| 1530 |
+
token: str = Header(..., alias="X-Auth-Token")
|
| 1531 |
+
):
|
| 1532 |
+
user = await authenticate_user(token)
|
| 1533 |
+
if not user:
|
| 1534 |
+
raise HTTPException(401)
|
| 1535 |
+
q = query.strip().lower()
|
| 1536 |
+
db = await get_db()
|
| 1537 |
+
try:
|
| 1538 |
+
if q:
|
| 1539 |
+
cursor = await db.execute(
|
| 1540 |
+
"""SELECT id, username, display_name, avatar_path, status
|
| 1541 |
+
FROM users WHERE id != ? AND (lower(username) LIKE ? OR lower(display_name) LIKE ?)
|
| 1542 |
+
ORDER BY CASE WHEN status = 'online' THEN 0 ELSE 1 END, username LIMIT 100""",
|
| 1543 |
+
(user['id'], f"%{q}%", f"%{q}%")
|
| 1544 |
+
)
|
| 1545 |
+
else:
|
| 1546 |
+
cursor = await db.execute(
|
| 1547 |
+
"""SELECT id, username, display_name, avatar_path, status FROM users
|
| 1548 |
+
WHERE id != ? ORDER BY CASE WHEN status = 'online' THEN 0 ELSE 1 END, username LIMIT 100""",
|
| 1549 |
+
(user['id'],)
|
| 1550 |
+
)
|
| 1551 |
+
out = []
|
| 1552 |
+
for r in await cursor.fetchall():
|
| 1553 |
+
out.append({
|
| 1554 |
+
"id": r["id"],
|
| 1555 |
+
"username": r["username"],
|
| 1556 |
+
"display_name": r["display_name"],
|
| 1557 |
+
"avatar_path": r["avatar_path"],
|
| 1558 |
+
"online": r["status"] == "online" and manager.is_online(r["id"]),
|
| 1559 |
+
"blocked": await _block_exists(db, user['id'], r["id"]),
|
| 1560 |
+
})
|
| 1561 |
+
return {"users": out}
|
| 1562 |
+
finally:
|
| 1563 |
+
await db.close()
|
| 1564 |
+
|
| 1565 |
+
|
| 1566 |
+
@app.get("/api/users/blocked")
|
| 1567 |
+
async def blocked_users(token: str = Header(..., alias="X-Auth-Token")):
|
| 1568 |
+
user = await authenticate_user(token)
|
| 1569 |
+
if not user:
|
| 1570 |
+
raise HTTPException(401)
|
| 1571 |
+
db = await get_db()
|
| 1572 |
+
try:
|
| 1573 |
+
cursor = await db.execute(
|
| 1574 |
+
"""SELECT u.id, u.username, u.display_name, u.avatar_path, b.created_at
|
| 1575 |
+
FROM user_blocks b JOIN users u ON u.id = b.blocked_id
|
| 1576 |
+
WHERE b.blocker_id = ? ORDER BY b.created_at DESC""",
|
| 1577 |
+
(user['id'],)
|
| 1578 |
+
)
|
| 1579 |
+
return {"users": [dict(r) for r in await cursor.fetchall()]}
|
| 1580 |
+
finally:
|
| 1581 |
+
await db.close()
|
| 1582 |
+
|
| 1583 |
+
|
| 1584 |
+
@app.post("/api/users/{user_id}/block")
|
| 1585 |
+
async def block_user_rest(user_id: int, token: str = Header(..., alias="X-Auth-Token")):
|
| 1586 |
+
user = await authenticate_user(token)
|
| 1587 |
+
if not user:
|
| 1588 |
+
raise HTTPException(401)
|
| 1589 |
+
if user_id == user['id']:
|
| 1590 |
+
raise HTTPException(400, "You can't block yourself")
|
| 1591 |
+
db = await get_db()
|
| 1592 |
+
try:
|
| 1593 |
+
cursor = await db.execute("SELECT id FROM users WHERE id = ?", (user_id,))
|
| 1594 |
+
if not await cursor.fetchone():
|
| 1595 |
+
raise HTTPException(404, "User not found")
|
| 1596 |
+
await db.execute(
|
| 1597 |
+
"INSERT OR IGNORE INTO user_blocks (blocker_id, blocked_id) VALUES (?, ?)",
|
| 1598 |
+
(user['id'], user_id)
|
| 1599 |
+
)
|
| 1600 |
+
await db.commit()
|
| 1601 |
+
schedule_db_sync()
|
| 1602 |
+
await manager.send_to_user(user_id, {
|
| 1603 |
+
"type": "block_changed",
|
| 1604 |
+
"blocker_id": user['id'],
|
| 1605 |
+
"blocked": True,
|
| 1606 |
+
})
|
| 1607 |
+
return {"status": "blocked", "user_id": user_id}
|
| 1608 |
+
finally:
|
| 1609 |
+
await db.close()
|
| 1610 |
+
|
| 1611 |
+
|
| 1612 |
+
@app.delete("/api/users/{user_id}/block")
|
| 1613 |
+
async def unblock_user_rest(user_id: int, token: str = Header(..., alias="X-Auth-Token")):
|
| 1614 |
+
user = await authenticate_user(token)
|
| 1615 |
+
if not user:
|
| 1616 |
+
raise HTTPException(401)
|
| 1617 |
+
db = await get_db()
|
| 1618 |
+
try:
|
| 1619 |
+
await db.execute(
|
| 1620 |
+
"DELETE FROM user_blocks WHERE blocker_id = ? AND blocked_id = ?",
|
| 1621 |
+
(user['id'], user_id)
|
| 1622 |
+
)
|
| 1623 |
+
await db.commit()
|
| 1624 |
+
schedule_db_sync()
|
| 1625 |
+
await manager.send_to_user(user_id, {
|
| 1626 |
+
"type": "block_changed",
|
| 1627 |
+
"blocker_id": user['id'],
|
| 1628 |
+
"blocked": False,
|
| 1629 |
+
})
|
| 1630 |
+
return {"status": "unblocked", "user_id": user_id}
|
| 1631 |
+
finally:
|
| 1632 |
+
await db.close()
|
| 1633 |
+
|
| 1634 |
+
|
| 1635 |
# ------------------------------------------------------------------------
|
| 1636 |
# Messages (REST fallbacks for reliable edit/delete/receipts)
|
| 1637 |
# ------------------------------------------------------------------------
|
|
|
|
| 1662 |
row = await _get_message_row(db, message_id)
|
| 1663 |
if not row or row["sender_id"] != user['id'] or row["is_deleted"]:
|
| 1664 |
raise HTTPException(404, "Message not found")
|
| 1665 |
+
if not await user_in_conversation(db, user['id'], row["conversation_id"]):
|
| 1666 |
+
raise HTTPException(403, "You can no longer access this conversation")
|
| 1667 |
new_enc = encrypt_message(content)
|
| 1668 |
await db.execute(
|
| 1669 |
"UPDATE messages SET encrypted_content = ?, is_edited = 1 WHERE id = ?",
|
|
|
|
| 1697 |
db = await get_db()
|
| 1698 |
try:
|
| 1699 |
row = await _get_message_row(db, message_id)
|
| 1700 |
+
if not row:
|
| 1701 |
+
# Idempotent: an already-removed message should not make the
|
| 1702 |
+
# client keep showing "delete failed" when a resend/reconnect
|
| 1703 |
+
# replays the request.
|
| 1704 |
+
return {"status": "already_deleted", "message_id": message_id}
|
| 1705 |
+
if row["sender_id"] != user['id'] or row["is_deleted"]:
|
| 1706 |
raise HTTPException(404, "Message not found")
|
| 1707 |
+
if not await user_in_conversation(db, user['id'], row["conversation_id"]):
|
| 1708 |
+
raise HTTPException(403, "You can no longer access this conversation")
|
| 1709 |
|
| 1710 |
if row["file_path"]:
|
| 1711 |
try:
|
|
|
|
| 1930 |
if '..' in file_path or '\\' in file_path:
|
| 1931 |
raise HTTPException(400, "Invalid path")
|
| 1932 |
|
| 1933 |
+
# Avatars are public profile data. Chat files belong to a conversation, so
|
| 1934 |
+
# a user blocked out of it must not be able to keep downloading the files.
|
| 1935 |
+
if not file_path.startswith("avatars/"):
|
| 1936 |
+
db = await get_db()
|
| 1937 |
+
try:
|
| 1938 |
+
cursor = await db.execute(
|
| 1939 |
+
"SELECT conversation_id FROM messages WHERE file_path = ? ORDER BY id DESC LIMIT 1",
|
| 1940 |
+
(file_path,)
|
| 1941 |
+
)
|
| 1942 |
+
row = await cursor.fetchone()
|
| 1943 |
+
if not row:
|
| 1944 |
+
raise HTTPException(404, "File not found")
|
| 1945 |
+
if not await user_in_conversation(db, user['id'], row["conversation_id"]):
|
| 1946 |
+
raise HTTPException(403, "You can no longer access this file")
|
| 1947 |
+
finally:
|
| 1948 |
+
await db.close()
|
| 1949 |
+
|
| 1950 |
try:
|
| 1951 |
# AES-GCM needs the whole ciphertext to verify its auth tag, so one full
|
| 1952 |
# read is required for security - the response is streamed in chunks.
|
|
|
|
| 2592 |
|
| 2593 |
@app.get("/")
|
| 2594 |
async def root():
|
| 2595 |
+
return FileResponse("static/index.html", headers={"Cache-Control": "no-cache"})
|