smodusermc commited on
Commit
a40ecea
·
verified ·
1 Parent(s): 4ad9c44

Update app.py

Browse files
Files changed (1) hide show
  1. app.py +992 -32
app.py CHANGED
@@ -8,12 +8,14 @@ import logging
8
  import tempfile
9
  import shutil
10
  import asyncio
 
 
11
  from typing import Optional, Dict, List, Any, Tuple
12
  from contextlib import asynccontextmanager
13
 
14
  import aiosqlite
15
 
16
- from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect, Header, Query, UploadFile, File
17
  from fastapi.staticfiles import StaticFiles
18
  from fastapi.responses import FileResponse, StreamingResponse
19
  from fastapi.middleware.cors import CORSMiddleware
@@ -25,17 +27,27 @@ from cryptography.hazmat.backends import default_backend
25
 
26
  from storage_handler import (
27
  store_file, retrieve_file, delete_file,
28
- download_database, upload_database, start_db_sync
 
 
29
  )
30
 
31
  # ------------------------------------------------------------------------
32
  # Configuration
33
  # ------------------------------------------------------------------------
34
- APP_VERSION = "2.0.0"
35
  DATABASE_URL = os.environ.get("DATABASE_URL", "/data/infinitychat.db")
36
  MESSAGE_KEY_B64 = os.environ.get("SECRET_KEY", None)
37
  FILE_ENCRYPTION_KEY_B64 = os.environ.get("FILE_ENCRYPTION_KEY", None)
38
 
 
 
 
 
 
 
 
 
39
  logging.basicConfig(level=logging.INFO)
40
  logger = logging.getLogger("InfinityChat")
41
 
@@ -83,9 +95,12 @@ async def _ensure_column(db: aiosqlite.Connection, table: str, column: str, ddl:
83
 
84
  async def init_database():
85
  os.makedirs(os.path.dirname(DATABASE_URL), exist_ok=True)
 
86
  db_existed = download_database(DATABASE_URL)
87
  if db_existed:
88
- logger.info("✅ Restored database from bucket")
 
 
89
  else:
90
  logger.info("🆕 Starting with fresh database")
91
 
@@ -178,6 +193,78 @@ async def init_database():
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:
@@ -196,6 +283,9 @@ async def init_database():
196
  "avatar_path": "TEXT",
197
  "last_seen": "INTEGER DEFAULT (strftime('%s','now'))",
198
  "status": "TEXT DEFAULT 'offline'",
 
 
 
199
  }.items():
200
  if col not in usr_cols:
201
  await _ensure_column(db, "users", col, ddl)
@@ -286,11 +376,30 @@ async def init_database():
286
  "ON user_blocks(blocked_id)")
287
  await db.execute("CREATE INDEX IF NOT EXISTS idx_blocks_blocker "
288
  "ON user_blocks(blocker_id)")
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
289
  await db.commit()
290
 
291
  upload_database(DATABASE_URL)
292
  start_db_sync(DATABASE_URL)
293
- logger.info("✅ Database initialized successfully (schema v2)")
 
294
 
295
  # Clean up any leftover temp upload dirs from previous runs
296
  _cleanup_temp_dir()
@@ -408,7 +517,7 @@ async def authenticate_user(token: str) -> Optional[dict]:
408
  db = await get_db()
409
  try:
410
  cursor = await db.execute(
411
- "SELECT id, username, display_name, avatar_path, status FROM users WHERE token = ?",
412
  (token,)
413
  )
414
  user = await cursor.fetchone()
@@ -993,7 +1102,10 @@ async def get_profile(token: str = Header(..., alias="X-Auth-Token")):
993
 
994
  @app.patch("/api/profile")
995
  async def update_profile(
996
- display_name: str = Query(..., max_length=50),
 
 
 
997
  token: str = Header(..., alias="X-Auth-Token")
998
  ):
999
  user = await authenticate_user(token)
@@ -1001,21 +1113,38 @@ async def update_profile(
1001
  raise HTTPException(401)
1002
  db = await get_db()
1003
  try:
1004
- updated_name = display_name.strip() or user['username']
1005
- await db.execute(
1006
- "UPDATE users SET display_name = ? WHERE id = ?",
1007
- (updated_name, user['id'])
1008
- )
1009
- await db.commit()
1010
- schedule_db_sync()
1011
- user['display_name'] = updated_name
1012
- # Live presence uses in-memory state; keep the online-users list fresh
1013
- manager.update_profile(user['id'], display_name=updated_name, avatar_path=user.get('avatar_path'))
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1014
  await manager.broadcast({
1015
  "type": "profile_updated",
1016
  "user_id": user['id'],
1017
  "username": user['username'],
1018
- "display_name": updated_name,
1019
  "avatar_path": user.get('avatar_path'),
1020
  })
1021
  return {"user": user}
@@ -1118,7 +1247,10 @@ async def upload_avatar(
1118
  elif detected == "image/webp":
1119
  ext = ".webp"
1120
  remote_path = f"avatars/{user['username']}_{uuid.uuid4().hex}{ext}"
1121
- await asyncio.to_thread(store_file, remote_path, data)
 
 
 
1122
  finally:
1123
  os.unlink(tmp_path)
1124
 
@@ -1797,6 +1929,815 @@ async def unblock_user_rest(user_id: int, token: str = Header(..., alias="X-Auth
1797
  await db.close()
1798
 
1799
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
 
1800
  # ------------------------------------------------------------------------
1801
  # Messages (REST fallbacks for reliable edit/delete/receipts)
1802
  # ------------------------------------------------------------------------
@@ -2057,8 +2998,12 @@ async def upload_chunk(
2057
  with open(assembled_path, 'rb') as f:
2058
  full_data = f.read()
2059
  # Blocking bucket write happens off the event loop
2060
- await asyncio.to_thread(store_file, remote_path, full_data)
2061
- del full_data
 
 
 
 
2062
 
2063
  finally:
2064
  shutil.rmtree(session_dir, ignore_errors=True)
@@ -2095,20 +3040,30 @@ async def download_file(
2095
  if '..' in file_path or '\\' in file_path:
2096
  raise HTTPException(400, "Invalid path")
2097
 
2098
- # Avatars are public profile data. Chat files belong to a conversation, so
2099
- # a user blocked out of it must not be able to keep downloading the files.
 
 
2100
  if not file_path.startswith("avatars/"):
2101
  db = await get_db()
2102
  try:
2103
- cursor = await db.execute(
2104
- "SELECT conversation_id FROM messages WHERE file_path = ? ORDER BY id DESC LIMIT 1",
2105
- (file_path,)
 
2106
  )
2107
- row = await cursor.fetchone()
2108
- if not row:
2109
- raise HTTPException(404, "File not found")
2110
- if not await user_in_conversation(db, user['id'], row["conversation_id"]):
2111
- raise HTTPException(403, "You can no longer access this file")
 
 
 
 
 
 
 
2112
  finally:
2113
  await db.close()
2114
 
@@ -2758,4 +3713,9 @@ async def ws_endpoint(ws: WebSocket, token: str = Query(...)):
2758
 
2759
  @app.get("/")
2760
  async def root():
2761
- return FileResponse("static/index.html", headers={"Cache-Control": "no-cache"})
 
 
 
 
 
 
8
  import tempfile
9
  import shutil
10
  import asyncio
11
+ import secrets
12
+ import threading
13
  from typing import Optional, Dict, List, Any, Tuple
14
  from contextlib import asynccontextmanager
15
 
16
  import aiosqlite
17
 
18
+ from fastapi import FastAPI, HTTPException, WebSocket, WebSocketDisconnect, Header, Query, UploadFile, File, Request
19
  from fastapi.staticfiles import StaticFiles
20
  from fastapi.responses import FileResponse, StreamingResponse
21
  from fastapi.middleware.cors import CORSMiddleware
 
27
 
28
  from storage_handler import (
29
  store_file, retrieve_file, delete_file,
30
+ download_database, upload_database, start_db_sync,
31
+ list_backups, create_timestamped_backup, restore_backup, delete_backup,
32
+ start_backup_loop, StorageUnavailableError, _OFFLINE, LOCAL_STORAGE_DIR
33
  )
34
 
35
  # ------------------------------------------------------------------------
36
  # Configuration
37
  # ------------------------------------------------------------------------
38
+ APP_VERSION = "3.0.0"
39
  DATABASE_URL = os.environ.get("DATABASE_URL", "/data/infinitychat.db")
40
  MESSAGE_KEY_B64 = os.environ.get("SECRET_KEY", None)
41
  FILE_ENCRYPTION_KEY_B64 = os.environ.get("FILE_ENCRYPTION_KEY", None)
42
 
43
+ # Admin console credentials. These are intentionally plain-text secrets supplied
44
+ # by the host via secret/environment configuration (ADMIN_USERNAME/ADMIN_PASSWORD).
45
+ ADMIN_USERNAME = (os.environ.get("ADMIN_USERNAME", "") or "").strip()
46
+ ADMIN_PASSWORD = os.environ.get("ADMIN_PASSWORD", "") or ""
47
+ ADMIN_SESSION_TTL = 12 * 60 * 60
48
+ ADMIN_SESSION_LOCK = threading.Lock()
49
+ ADMIN_SESSIONS: Dict[str, Dict[str, Any]] = {}
50
+
51
  logging.basicConfig(level=logging.INFO)
52
  logger = logging.getLogger("InfinityChat")
53
 
 
95
 
96
  async def init_database():
97
  os.makedirs(os.path.dirname(DATABASE_URL), exist_ok=True)
98
+ local_exists = os.path.exists(DATABASE_URL)
99
  db_existed = download_database(DATABASE_URL)
100
  if db_existed:
101
+ logger.info("✅ Restored database from storage")
102
+ elif local_exists:
103
+ logger.info("📁 Using existing local database")
104
  else:
105
  logger.info("🆕 Starting with fresh database")
106
 
 
193
  )
194
  """)
195
 
196
+ # --- v3: social feed / Twitter-like features ---
197
+ await db.execute("""
198
+ CREATE TABLE IF NOT EXISTS social_posts (
199
+ id INTEGER PRIMARY KEY AUTOINCREMENT,
200
+ author_id INTEGER NOT NULL,
201
+ body TEXT NOT NULL DEFAULT '',
202
+ media_json TEXT,
203
+ reply_to_id INTEGER,
204
+ quote_id INTEGER,
205
+ is_deleted INTEGER NOT NULL DEFAULT 0,
206
+ created_at_ms INTEGER NOT NULL,
207
+ edited_at_ms INTEGER,
208
+ FOREIGN KEY(author_id) REFERENCES users(id) ON DELETE CASCADE,
209
+ FOREIGN KEY(reply_to_id) REFERENCES social_posts(id) ON DELETE SET NULL,
210
+ FOREIGN KEY(quote_id) REFERENCES social_posts(id) ON DELETE SET NULL
211
+ )
212
+ """)
213
+ await db.execute("""
214
+ CREATE TABLE IF NOT EXISTS social_likes (
215
+ post_id INTEGER NOT NULL,
216
+ user_id INTEGER NOT NULL,
217
+ created_at_ms INTEGER NOT NULL,
218
+ PRIMARY KEY (post_id, user_id),
219
+ FOREIGN KEY(post_id) REFERENCES social_posts(id) ON DELETE CASCADE,
220
+ FOREIGN KEY(user_id) REFERENCES users(id) ON DELETE CASCADE
221
+ )
222
+ """)
223
+ await db.execute("""
224
+ CREATE TABLE IF NOT EXISTS social_reposts (
225
+ post_id INTEGER NOT NULL,
226
+ user_id INTEGER NOT NULL,
227
+ created_at_ms INTEGER NOT NULL,
228
+ PRIMARY KEY (post_id, user_id),
229
+ FOREIGN KEY(post_id) REFERENCES social_posts(id) ON DELETE CASCADE,
230
+ FOREIGN KEY(user_id) REFERENCES users(id) ON DELETE CASCADE
231
+ )
232
+ """)
233
+ await db.execute("""
234
+ CREATE TABLE IF NOT EXISTS social_bookmarks (
235
+ post_id INTEGER NOT NULL,
236
+ user_id INTEGER NOT NULL,
237
+ created_at_ms INTEGER NOT NULL,
238
+ PRIMARY KEY (post_id, user_id),
239
+ FOREIGN KEY(post_id) REFERENCES social_posts(id) ON DELETE CASCADE,
240
+ FOREIGN KEY(user_id) REFERENCES users(id) ON DELETE CASCADE
241
+ )
242
+ """)
243
+ await db.execute("""
244
+ CREATE TABLE IF NOT EXISTS follows (
245
+ follower_id INTEGER NOT NULL,
246
+ following_id INTEGER NOT NULL,
247
+ created_at_ms INTEGER NOT NULL,
248
+ PRIMARY KEY (follower_id, following_id),
249
+ FOREIGN KEY(follower_id) REFERENCES users(id) ON DELETE CASCADE,
250
+ FOREIGN KEY(following_id) REFERENCES users(id) ON DELETE CASCADE
251
+ )
252
+ """)
253
+ await db.execute("""
254
+ CREATE TABLE IF NOT EXISTS social_notifications (
255
+ id INTEGER PRIMARY KEY AUTOINCREMENT,
256
+ user_id INTEGER NOT NULL,
257
+ actor_id INTEGER NOT NULL,
258
+ type TEXT NOT NULL,
259
+ post_id INTEGER,
260
+ created_at_ms INTEGER NOT NULL,
261
+ read INTEGER NOT NULL DEFAULT 0,
262
+ FOREIGN KEY(user_id) REFERENCES users(id) ON DELETE CASCADE,
263
+ FOREIGN KEY(actor_id) REFERENCES users(id) ON DELETE CASCADE,
264
+ FOREIGN KEY(post_id) REFERENCES social_posts(id) ON DELETE CASCADE
265
+ )
266
+ """)
267
+
268
  # --- Additive migrations for databases created before v2 ---
269
  msg_cols = await _table_columns(db, "messages")
270
  if "conversation_id" not in msg_cols:
 
283
  "avatar_path": "TEXT",
284
  "last_seen": "INTEGER DEFAULT (strftime('%s','now'))",
285
  "status": "TEXT DEFAULT 'offline'",
286
+ "bio": "TEXT NOT NULL DEFAULT ''",
287
+ "location": "TEXT NOT NULL DEFAULT ''",
288
+ "website": "TEXT NOT NULL DEFAULT ''",
289
  }.items():
290
  if col not in usr_cols:
291
  await _ensure_column(db, "users", col, ddl)
 
376
  "ON user_blocks(blocked_id)")
377
  await db.execute("CREATE INDEX IF NOT EXISTS idx_blocks_blocker "
378
  "ON user_blocks(blocker_id)")
379
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_posts_author_time "
380
+ "ON social_posts(author_id, created_at_ms)")
381
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_posts_time "
382
+ "ON social_posts(created_at_ms)")
383
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_posts_reply "
384
+ "ON social_posts(reply_to_id)")
385
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_follows_following "
386
+ "ON follows(following_id)")
387
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_follows_follower "
388
+ "ON follows(follower_id)")
389
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_likes_post "
390
+ "ON social_likes(post_id)")
391
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_reposts_post "
392
+ "ON social_reposts(post_id)")
393
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_bookmarks_user "
394
+ "ON social_bookmarks(user_id, post_id)")
395
+ await db.execute("CREATE INDEX IF NOT EXISTS idx_notifications_user "
396
+ "ON social_notifications(user_id, read, created_at_ms)")
397
  await db.commit()
398
 
399
  upload_database(DATABASE_URL)
400
  start_db_sync(DATABASE_URL)
401
+ start_backup_loop(DATABASE_URL)
402
+ logger.info("✅ Database initialized successfully (schema v3)")
403
 
404
  # Clean up any leftover temp upload dirs from previous runs
405
  _cleanup_temp_dir()
 
517
  db = await get_db()
518
  try:
519
  cursor = await db.execute(
520
+ "SELECT id, username, display_name, avatar_path, status, bio, location, website FROM users WHERE token = ?",
521
  (token,)
522
  )
523
  user = await cursor.fetchone()
 
1102
 
1103
  @app.patch("/api/profile")
1104
  async def update_profile(
1105
+ display_name: str = Query(None, max_length=50),
1106
+ bio: str = Query(None, max_length=200),
1107
+ location: str = Query(None, max_length=80),
1108
+ website: str = Query(None, max_length=160),
1109
  token: str = Header(..., alias="X-Auth-Token")
1110
  ):
1111
  user = await authenticate_user(token)
 
1113
  raise HTTPException(401)
1114
  db = await get_db()
1115
  try:
1116
+ updates = []
1117
+ args = []
1118
+ if display_name is not None:
1119
+ updates.append("display_name = ?")
1120
+ args.append(display_name.strip() or user['username'])
1121
+ if bio is not None:
1122
+ updates.append("bio = ?")
1123
+ args.append(bio.strip()[:200])
1124
+ if location is not None:
1125
+ updates.append("location = ?")
1126
+ args.append(location.strip()[:80])
1127
+ if website is not None:
1128
+ website = (website or "").strip()
1129
+ if website and not re.match(r"^https?://", website, re.I):
1130
+ website = "https://" + website
1131
+ updates.append("website = ?")
1132
+ args.append(website[:160])
1133
+ if updates:
1134
+ args.append(user['id'])
1135
+ await db.execute(f"UPDATE users SET {', '.join(updates)} WHERE id = ?", args)
1136
+ await db.commit()
1137
+ schedule_db_sync()
1138
+ user = await authenticate_user(token)
1139
+ else:
1140
+ raise HTTPException(400, "No profile fields supplied")
1141
+ manager.update_profile(user['id'], display_name=user.get('display_name'),
1142
+ avatar_path=user.get('avatar_path'))
1143
  await manager.broadcast({
1144
  "type": "profile_updated",
1145
  "user_id": user['id'],
1146
  "username": user['username'],
1147
+ "display_name": user.get('display_name'),
1148
  "avatar_path": user.get('avatar_path'),
1149
  })
1150
  return {"user": user}
 
1247
  elif detected == "image/webp":
1248
  ext = ".webp"
1249
  remote_path = f"avatars/{user['username']}_{uuid.uuid4().hex}{ext}"
1250
+ try:
1251
+ await asyncio.to_thread(store_file, remote_path, data)
1252
+ except StorageUnavailableError:
1253
+ raise HTTPException(503, "File storage is unavailable")
1254
  finally:
1255
  os.unlink(tmp_path)
1256
 
 
1929
  await db.close()
1930
 
1931
 
1932
+ # ------------------------------------------------------------------------
1933
+ # Social feed (Twitter/X-like)
1934
+ # ------------------------------------------------------------------------
1935
+ MAX_POST_LENGTH = 28000
1936
+ MAX_POST_MEDIA = 4
1937
+
1938
+
1939
+ async def _blocked_user_ids(db: aiosqlite.Connection, uid: int) -> set:
1940
+ """All users who are in a block relationship with uid (either direction)."""
1941
+ cursor = await db.execute(
1942
+ "SELECT blocker_id, blocked_id FROM user_blocks "
1943
+ "WHERE blocker_id = ? OR blocked_id = ?",
1944
+ (uid, uid)
1945
+ )
1946
+ ids = set()
1947
+ for row in await cursor.fetchall():
1948
+ ids.add(row["blocker_id"])
1949
+ ids.add(row["blocked_id"])
1950
+ ids.discard(uid)
1951
+ return ids
1952
+
1953
+
1954
+ async def _social_user(db: aiosqlite.Connection, uid: int) -> dict:
1955
+ cursor = await db.execute(
1956
+ "SELECT id, username, display_name, avatar_path, bio, location, website "
1957
+ "FROM users WHERE id = ?", (uid,)
1958
+ )
1959
+ row = await cursor.fetchone()
1960
+ if not row:
1961
+ raise HTTPException(404, "User not found")
1962
+ return dict(row)
1963
+
1964
+
1965
+ async def _toggle_row(db: aiosqlite.Connection, table: str, post_id: int, uid: int):
1966
+ """Insert/remove a (post_id,user_id) row; returns True when now present."""
1967
+ cursor = await db.execute(f"SELECT 1 FROM {table} WHERE post_id = ? AND user_id = ?",
1968
+ (post_id, uid))
1969
+ exists = await cursor.fetchone() is not None
1970
+ if exists:
1971
+ await db.execute(f"DELETE FROM {table} WHERE post_id = ? AND user_id = ?", (post_id, uid))
1972
+ return False
1973
+ await db.execute(f"INSERT INTO {table} (post_id, user_id, created_at_ms) VALUES (?,?,?)",
1974
+ (post_id, uid, int(time.time() * 1000)))
1975
+ return True
1976
+
1977
+
1978
+ async def _count_post(db: aiosqlite.Connection, table: str, post_id: int) -> int:
1979
+ cursor = await db.execute(f"SELECT COUNT(*) AS c FROM {table} WHERE post_id = ?", (post_id,))
1980
+ row = await cursor.fetchone()
1981
+ return row["c"] or 0
1982
+
1983
+
1984
+ async def _post_exists(db: aiosqlite.Connection, post_id: int, include_deleted: bool = False) -> bool:
1985
+ sql = "SELECT 1 FROM social_posts WHERE id = ?"
1986
+ args = [post_id]
1987
+ if not include_deleted:
1988
+ sql += " AND is_deleted = 0"
1989
+ cursor = await db.execute(sql, args)
1990
+ return await cursor.fetchone() is not None
1991
+
1992
+
1993
+ async def _post_author_id(db: aiosqlite.Connection, post_id: int) -> Optional[int]:
1994
+ cur = await db.execute("SELECT author_id FROM social_posts WHERE id = ?", (post_id,))
1995
+ row = await cur.fetchone()
1996
+ return row["author_id"] if row else None
1997
+
1998
+
1999
+ async def _ensure_interactable_post(db: aiosqlite.Connection, post_id: int, viewer_id: int) -> None:
2000
+ author = await _post_author_id(db, post_id)
2001
+ if author is None:
2002
+ raise HTTPException(404, "Post not found")
2003
+ if await _block_exists(db, viewer_id, author):
2004
+ raise HTTPException(403, "You cannot interact with this user")
2005
+
2006
+
2007
+ async def _serialize_post(db: aiosqlite.Connection, row: dict, viewer_id: int,
2008
+ reposter_id: int = None) -> dict:
2009
+ author = await _social_user(db, row["author_id"])
2010
+ media = []
2011
+ if row.get("media_json"):
2012
+ try:
2013
+ media = json.loads(row["media_json"])
2014
+ except Exception:
2015
+ media = []
2016
+ quote = None
2017
+ if row.get("quote_id"):
2018
+ qcur = await db.execute("SELECT * FROM social_posts WHERE id = ?", (row["quote_id"],))
2019
+ qrow = await qcur.fetchone()
2020
+ if qrow and not qrow["is_deleted"] and not await _block_exists(db, viewer_id, qrow["author_id"]):
2021
+ quote = await _serialize_post(db, dict(qrow), viewer_id)
2022
+ async def has(table):
2023
+ cur = await db.execute(f"SELECT 1 FROM {table} WHERE post_id = ? AND user_id = ?",
2024
+ (row["id"], viewer_id))
2025
+ return await cur.fetchone() is not None
2026
+ reposter = None
2027
+ if reposter_id:
2028
+ ru = await _social_user(db, reposter_id)
2029
+ reposter = {
2030
+ "id": ru["id"], "username": ru["username"],
2031
+ "display_name": ru["display_name"], "avatar_path": ru["avatar_path"],
2032
+ }
2033
+ return {
2034
+ "id": row["id"],
2035
+ "author": {
2036
+ "id": author["id"],
2037
+ "username": author["username"],
2038
+ "display_name": author["display_name"],
2039
+ "avatar_path": author["avatar_path"],
2040
+ },
2041
+ "body": row["body"],
2042
+ "media": media,
2043
+ "reply_to_id": row["reply_to_id"],
2044
+ "quote_id": row["quote_id"],
2045
+ "quote": quote,
2046
+ "created_at_ms": row["created_at_ms"],
2047
+ "edited_at_ms": row["edited_at_ms"],
2048
+ "reposter_id": reposter_id,
2049
+ "reposter": reposter,
2050
+ "like_count": await _count_post(db, "social_likes", row["id"]),
2051
+ "repost_count": await _count_post(db, "social_reposts", row["id"]),
2052
+ "reply_count": await _count_replies(db, row["id"]),
2053
+ "bookmark_count": await _count_post(db, "social_bookmarks", row["id"]),
2054
+ "liked": await has("social_likes"),
2055
+ "reposted": await has("social_reposts"),
2056
+ "bookmarked": await has("social_bookmarks"),
2057
+ }
2058
+
2059
+
2060
+ async def _count_replies(db: aiosqlite.Connection, post_id: int) -> int:
2061
+ cursor = await db.execute(
2062
+ "SELECT COUNT(*) AS c FROM social_posts WHERE reply_to_id = ? AND is_deleted = 0",
2063
+ (post_id,)
2064
+ )
2065
+ return (await cursor.fetchone())["c"] or 0
2066
+
2067
+
2068
+ async def _serialize_post_rows(db: aiosqlite.Connection, rows: List[dict],
2069
+ viewer_id: int, reposters: List[int] = None) -> List[dict]:
2070
+ out = []
2071
+ blocked = await _blocked_user_ids(db, viewer_id)
2072
+ for i, row in enumerate(rows):
2073
+ if row["author_id"] in blocked:
2074
+ continue
2075
+ if row.get("reposter_id") is not None and row["reposter_id"] in blocked:
2076
+ continue
2077
+ reposter = row.get("reposter_id") if row.get("reposter_id") is not None else ((reposters or [None])[i] if reposters else None)
2078
+ out.append(await _serialize_post(db, row, viewer_id, reposter))
2079
+ return out
2080
+
2081
+
2082
+ async def _extract_mentions(body: str) -> List[str]:
2083
+ return list(dict.fromkeys(re.findall(r"@([A-Za-z0-9_]{1,30})", body or "")))
2084
+
2085
+
2086
+ async def _extract_hashtags(body: str) -> List[str]:
2087
+ return list(dict.fromkeys(re.findall(r"#([A-Za-z0-9_]{1,100})", body or "")))
2088
+
2089
+
2090
+ async def _notify_social(db: aiosqlite.Connection, user_id: int, actor_id: int,
2091
+ ntype: str, post_id: int = None):
2092
+ if user_id == actor_id:
2093
+ return
2094
+ if await _block_exists(db, user_id, actor_id):
2095
+ return
2096
+ await db.execute(
2097
+ "INSERT INTO social_notifications (user_id, actor_id, type, post_id, created_at_ms) "
2098
+ "VALUES (?,?,?,?,?)",
2099
+ (user_id, actor_id, ntype, post_id, int(time.time() * 1000))
2100
+ )
2101
+
2102
+
2103
+ async def _load_post(db: aiosqlite.Connection, post_id: int, viewer_id: int) -> dict:
2104
+ cur = await db.execute("SELECT * FROM social_posts WHERE id = ?", (post_id,))
2105
+ row = await cur.fetchone()
2106
+ if not row or row["is_deleted"]:
2107
+ raise HTTPException(404, "Post not found")
2108
+ if await _block_exists(db, viewer_id, row["author_id"]):
2109
+ raise HTTPException(403, "You cannot view posts from this user")
2110
+ return await _serialize_post(db, dict(row), viewer_id)
2111
+
2112
+
2113
+ async def _social_author_ids(db: aiosqlite.Connection, uid: int) -> List[int]:
2114
+ ids = {uid}
2115
+ cursor = await db.execute("SELECT following_id FROM follows WHERE follower_id = ?", (uid,))
2116
+ for row in await cursor.fetchall():
2117
+ ids.add(row["following_id"])
2118
+ return list(ids)
2119
+
2120
+
2121
+ _SOCIAL_MEDIA_EXTS = (".jpg", ".jpeg", ".png", ".gif", ".webp", ".bmp", ".avif",
2122
+ ".mp4", ".webm", ".mov", ".ogg", ".m4v", ".mpeg", ".mpg")
2123
+
2124
+
2125
+ def _is_social_media(m: dict) -> bool:
2126
+ mtype = (str(m.get("file_type", "") or "")).lower()
2127
+ path = (str(m.get("file_path", "") or "")).lower()
2128
+ return mtype.startswith("image/") or mtype.startswith("video/") or path.endswith(_SOCIAL_MEDIA_EXTS)
2129
+
2130
+
2131
+ @app.post("/api/social/posts")
2132
+ async def create_social_post_rest(
2133
+ body: str = Query("", max_length=MAX_POST_LENGTH),
2134
+ media_json: str = Query(""),
2135
+ reply_to_id: int = Query(None),
2136
+ quote_id: int = Query(None),
2137
+ token: str = Header(..., alias="X-Auth-Token")
2138
+ ):
2139
+ user = await authenticate_user(token)
2140
+ if not user:
2141
+ raise HTTPException(401)
2142
+ body = (body or "").strip()
2143
+ if not body and not media_json:
2144
+ raise HTTPException(400, "Post cannot be empty")
2145
+ media = []
2146
+ if media_json:
2147
+ try:
2148
+ media = json.loads(media_json)
2149
+ except Exception:
2150
+ raise HTTPException(400, "Invalid media payload")
2151
+ if not isinstance(media, list) or len(media) > MAX_POST_MEDIA:
2152
+ raise HTTPException(400, f"A post can contain at most {MAX_POST_MEDIA} media files")
2153
+ for m in media:
2154
+ path = m.get("file_path", "")
2155
+ if not str(path).startswith(f"uploads/{user['username']}/"):
2156
+ raise HTTPException(403, "You can only attach files you uploaded")
2157
+ if not _is_social_media(m):
2158
+ raise HTTPException(400, "Only images and videos can be attached to posts")
2159
+ db = await get_db()
2160
+ try:
2161
+ if reply_to_id:
2162
+ await _ensure_interactable_post(db, reply_to_id, user["id"])
2163
+ if quote_id:
2164
+ await _ensure_interactable_post(db, quote_id, user["id"])
2165
+ cur = await db.execute(
2166
+ "INSERT INTO social_posts (author_id, body, media_json, reply_to_id, quote_id, "
2167
+ "created_at_ms) VALUES (?,?,?,?,?,?)",
2168
+ (user["id"], body, media_json or None, reply_to_id, quote_id,
2169
+ int(time.time() * 1000))
2170
+ )
2171
+ post_id = cur.lastrowid
2172
+ await db.commit()
2173
+ schedule_db_sync()
2174
+ # Notify reply target + mentioned users.
2175
+ if reply_to_id:
2176
+ rt = await db.execute("SELECT author_id FROM social_posts WHERE id = ?", (reply_to_id,))
2177
+ rt_row = await rt.fetchone()
2178
+ if rt_row:
2179
+ await _notify_social(db, rt_row["author_id"], user["id"], "reply", post_id)
2180
+ for handle in await _extract_mentions(body):
2181
+ mc = await db.execute("SELECT id FROM users WHERE lower(username) = lower(?)", (handle,))
2182
+ mrow = await mc.fetchone()
2183
+ if mrow:
2184
+ await _notify_social(db, mrow["id"], user["id"], "mention", post_id)
2185
+ await db.commit()
2186
+ post = await _load_post(db, post_id, user["id"])
2187
+ return {"post": post}
2188
+ finally:
2189
+ await db.close()
2190
+
2191
+
2192
+ @app.get("/api/social/posts/{post_id}")
2193
+ async def get_post_rest(post_id: int, token: str = Header(..., alias="X-Auth-Token")):
2194
+ user = await authenticate_user(token)
2195
+ if not user:
2196
+ raise HTTPException(401)
2197
+ db = await get_db()
2198
+ try:
2199
+ return {"post": await _load_post(db, post_id, user["id"])}
2200
+ finally:
2201
+ await db.close()
2202
+
2203
+
2204
+ @app.get("/api/social/posts/{post_id}/replies")
2205
+ async def get_replies_rest(post_id: int, limit: int = Query(50, le=100),
2206
+ token: str = Header(..., alias="X-Auth-Token")):
2207
+ user = await authenticate_user(token)
2208
+ if not user:
2209
+ raise HTTPException(401)
2210
+ db = await get_db()
2211
+ try:
2212
+ await _ensure_interactable_post(db, post_id, user["id"])
2213
+ cur = await db.execute(
2214
+ "SELECT * FROM social_posts WHERE reply_to_id = ? AND is_deleted = 0 "
2215
+ "ORDER BY created_at_ms ASC LIMIT ?", (post_id, limit))
2216
+ rows = [dict(r) for r in await cur.fetchall()]
2217
+ return {"replies": await _serialize_post_rows(db, rows, user["id"])}
2218
+ finally:
2219
+ await db.close()
2220
+
2221
+
2222
+ @app.delete("/api/social/posts/{post_id}")
2223
+ async def delete_social_post_rest(post_id: int, token: str = Header(..., alias="X-Auth-Token")):
2224
+ user = await authenticate_user(token)
2225
+ if not user:
2226
+ raise HTTPException(401)
2227
+ db = await get_db()
2228
+ try:
2229
+ cur = await db.execute("SELECT author_id FROM social_posts WHERE id = ?", (post_id,))
2230
+ row = await cur.fetchone()
2231
+ if not row:
2232
+ raise HTTPException(404, "Post not found")
2233
+ if row["author_id"] != user["id"]:
2234
+ raise HTTPException(403, "You can only delete your own posts")
2235
+ await db.execute("UPDATE social_posts SET is_deleted = 1 WHERE id = ?", (post_id,))
2236
+ await db.commit()
2237
+ schedule_db_sync()
2238
+ return {"status": "deleted", "post_id": post_id}
2239
+ finally:
2240
+ await db.close()
2241
+
2242
+
2243
+ @app.patch("/api/social/posts/{post_id}")
2244
+ async def edit_social_post_rest(
2245
+ post_id: int,
2246
+ body: str = Query("", max_length=MAX_POST_LENGTH),
2247
+ token: str = Header(..., alias="X-Auth-Token")
2248
+ ):
2249
+ user = await authenticate_user(token)
2250
+ if not user:
2251
+ raise HTTPException(401)
2252
+ body = (body or "").strip()
2253
+ if not body:
2254
+ raise HTTPException(400, "Post body cannot be empty")
2255
+ db = await get_db()
2256
+ try:
2257
+ cur = await db.execute("SELECT author_id, is_deleted FROM social_posts WHERE id = ?", (post_id,))
2258
+ row = await cur.fetchone()
2259
+ if not row or row["is_deleted"]:
2260
+ raise HTTPException(404, "Post not found")
2261
+ if row["author_id"] != user["id"]:
2262
+ raise HTTPException(403, "You can only edit your own posts")
2263
+ await db.execute(
2264
+ "UPDATE social_posts SET body = ?, edited_at_ms = ? WHERE id = ?",
2265
+ (body, int(time.time() * 1000), post_id)
2266
+ )
2267
+ await db.commit()
2268
+ schedule_db_sync()
2269
+ return {"post": await _load_post(db, post_id, user["id"])}
2270
+ finally:
2271
+ await db.close()
2272
+
2273
+
2274
+ @app.post("/api/social/posts/{post_id}/like")
2275
+ async def toggle_like_rest(post_id: int, token: str = Header(..., alias="X-Auth-Token")):
2276
+ user = await authenticate_user(token)
2277
+ if not user:
2278
+ raise HTTPException(401)
2279
+ db = await get_db()
2280
+ try:
2281
+ await _ensure_interactable_post(db, post_id, user["id"])
2282
+ liked = await _toggle_row(db, "social_likes", post_id, user["id"])
2283
+ await db.commit()
2284
+ if liked:
2285
+ cur = await db.execute("SELECT author_id FROM social_posts WHERE id = ?", (post_id,))
2286
+ row = await cur.fetchone()
2287
+ await _notify_social(db, row["author_id"], user["id"], "like", post_id)
2288
+ await db.commit()
2289
+ return {"liked": liked, "like_count": await _count_post(db, "social_likes", post_id)}
2290
+ finally:
2291
+ await db.close()
2292
+
2293
+
2294
+ @app.post("/api/social/posts/{post_id}/repost")
2295
+ async def toggle_repost_rest(post_id: int, token: str = Header(..., alias="X-Auth-Token")):
2296
+ user = await authenticate_user(token)
2297
+ if not user:
2298
+ raise HTTPException(401)
2299
+ db = await get_db()
2300
+ try:
2301
+ await _ensure_interactable_post(db, post_id, user["id"])
2302
+ reposted = await _toggle_row(db, "social_reposts", post_id, user["id"])
2303
+ await db.commit()
2304
+ if reposted:
2305
+ cur = await db.execute("SELECT author_id FROM social_posts WHERE id = ?", (post_id,))
2306
+ row = await cur.fetchone()
2307
+ await _notify_social(db, row["author_id"], user["id"], "repost", post_id)
2308
+ await db.commit()
2309
+ return {"reposted": reposted, "repost_count": await _count_post(db, "social_reposts", post_id)}
2310
+ finally:
2311
+ await db.close()
2312
+
2313
+
2314
+ @app.post("/api/social/posts/{post_id}/bookmark")
2315
+ async def toggle_bookmark_rest(post_id: int, token: str = Header(..., alias="X-Auth-Token")):
2316
+ user = await authenticate_user(token)
2317
+ if not user:
2318
+ raise HTTPException(401)
2319
+ db = await get_db()
2320
+ try:
2321
+ await _ensure_interactable_post(db, post_id, user["id"])
2322
+ bookmarked = await _toggle_row(db, "social_bookmarks", post_id, user["id"])
2323
+ await db.commit()
2324
+ return {"bookmarked": bookmarked}
2325
+ finally:
2326
+ await db.close()
2327
+
2328
+
2329
+ @app.get("/api/social/feed")
2330
+ async def get_social_feed_rest(
2331
+ feed: str = Query("home"),
2332
+ cursor: int = Query(-1),
2333
+ limit: int = Query(30, le=100),
2334
+ token: str = Header(..., alias="X-Auth-Token")
2335
+ ):
2336
+ user = await authenticate_user(token)
2337
+ if not user:
2338
+ raise HTTPException(401)
2339
+ db = await get_db()
2340
+ try:
2341
+ blocked = await _blocked_user_ids(db, user["id"])
2342
+ if feed == "home":
2343
+ author_ids = await _social_author_ids(db, user["id"])
2344
+ rows = []
2345
+ if author_ids:
2346
+ placeholders = ",".join("?" * len(author_ids))
2347
+ cur = await db.execute(
2348
+ f"SELECT * FROM social_posts WHERE author_id IN ({placeholders}) "
2349
+ f"AND is_deleted = 0 ORDER BY created_at_ms DESC LIMIT ?",
2350
+ author_ids + [limit * 3])
2351
+ rows = [dict(r) for r in await cur.fetchall()]
2352
+ # Reposts by followed accounts + self, merged with original posts.
2353
+ reposter_ids = list(author_ids)
2354
+ if reposter_ids:
2355
+ ph = ",".join("?" * len(reposter_ids))
2356
+ cur = await db.execute(
2357
+ f"SELECT p.*, r.user_id AS reposter_id, r.created_at_ms AS reposted_at_ms "
2358
+ f"FROM social_reposts r JOIN social_posts p ON p.id = r.post_id "
2359
+ f"WHERE r.user_id IN ({ph}) AND p.is_deleted = 0",
2360
+ reposter_ids)
2361
+ reposts = [dict(r) for r in await cur.fetchall()]
2362
+ rows = rows + reposts
2363
+ rows = sorted(rows, key=lambda r: r.get("reposted_at_ms") or r["created_at_ms"], reverse=True)
2364
+ elif feed == "explore":
2365
+ cur = await db.execute(
2366
+ "SELECT * FROM social_posts WHERE is_deleted = 0 "
2367
+ "ORDER BY created_at_ms DESC LIMIT ?", (limit * 3,))
2368
+ rows = [dict(r) for r in await cur.fetchall()]
2369
+ else:
2370
+ raise HTTPException(400, "Unknown feed")
2371
+ if cursor >= 0:
2372
+ rows = [r for r in rows if r["created_at_ms"] < cursor]
2373
+ rows = rows[:limit]
2374
+ posts = await _serialize_post_rows(db, rows, user["id"])
2375
+ posts.sort(key=lambda p: p["created_at_ms"], reverse=True)
2376
+ return {"posts": posts, "cursor": posts[-1]["created_at_ms"] if posts else None}
2377
+ finally:
2378
+ await db.close()
2379
+
2380
+
2381
+ @app.get("/api/social/users/{username}/posts")
2382
+ async def get_user_posts_rest(username: str, token: str = Header(..., alias="X-Auth-Token")):
2383
+ user = await authenticate_user(token)
2384
+ if not user:
2385
+ raise HTTPException(401)
2386
+ db = await get_db()
2387
+ try:
2388
+ u = await db.execute("SELECT id FROM users WHERE lower(username) = lower(?)", (username,))
2389
+ urow = await u.fetchone()
2390
+ if not urow:
2391
+ raise HTTPException(404, "User not found")
2392
+ if await _block_exists(db, user["id"], urow["id"]):
2393
+ raise HTTPException(403, "Content is not available")
2394
+ cur = await db.execute(
2395
+ "SELECT * FROM social_posts WHERE author_id = ? AND is_deleted = 0 "
2396
+ "ORDER BY created_at_ms DESC LIMIT 100", (urow["id"],))
2397
+ rows = [dict(r) for r in await cur.fetchall()]
2398
+ return {"posts": await _serialize_post_rows(db, rows, user["id"])}
2399
+ finally:
2400
+ await db.close()
2401
+
2402
+
2403
+ @app.get("/api/social/profile/{username}")
2404
+ async def get_social_profile_rest(username: str, token: str = Header(..., alias="X-Auth-Token")):
2405
+ user = await authenticate_user(token)
2406
+ if not user:
2407
+ raise HTTPException(401)
2408
+ db = await get_db()
2409
+ try:
2410
+ cur = await db.execute(
2411
+ "SELECT id FROM users WHERE lower(username) = lower(?)", (username,))
2412
+ urow = await cur.fetchone()
2413
+ if not urow:
2414
+ raise HTTPException(404, "User not found")
2415
+ target_id = urow["id"]
2416
+ blocked = await _block_exists(db, user["id"], target_id)
2417
+ if blocked:
2418
+ raise HTTPException(403, "Content is not available")
2419
+ prof = await _social_user(db, target_id)
2420
+ followers = await db.execute("SELECT COUNT(*) AS c FROM follows WHERE following_id = ?", (target_id,))
2421
+ following = await db.execute("SELECT COUNT(*) AS c FROM follows WHERE follower_id = ?", (target_id,))
2422
+ is_follow = await db.execute("SELECT 1 FROM follows WHERE follower_id = ? AND following_id = ?",
2423
+ (user["id"], target_id))
2424
+ posts = await db.execute("SELECT COUNT(*) AS c FROM social_posts WHERE author_id = ? AND is_deleted = 0",
2425
+ (target_id,))
2426
+ return {
2427
+ "profile": {
2428
+ **prof,
2429
+ "is_self": target_id == user["id"],
2430
+ "is_following": await is_follow.fetchone() is not None,
2431
+ "followers_count": (await followers.fetchone())["c"],
2432
+ "following_count": (await following.fetchone())["c"],
2433
+ "posts_count": (await posts.fetchone())["c"],
2434
+ }
2435
+ }
2436
+ finally:
2437
+ await db.close()
2438
+
2439
+
2440
+ @app.post("/api/social/users/{username}/follow")
2441
+ async def toggle_follow_rest(username: str, token: str = Header(..., alias="X-Auth-Token")):
2442
+ user = await authenticate_user(token)
2443
+ if not user:
2444
+ raise HTTPException(401)
2445
+ db = await get_db()
2446
+ try:
2447
+ cur = await db.execute("SELECT id FROM users WHERE lower(username) = lower(?)", (username,))
2448
+ row = await cur.fetchone()
2449
+ if not row:
2450
+ raise HTTPException(404, "User not found")
2451
+ target_id = row["id"]
2452
+ if target_id == user["id"]:
2453
+ raise HTTPException(400, "You cannot follow yourself")
2454
+ if await _block_exists(db, user["id"], target_id):
2455
+ raise HTTPException(403, "You cannot interact with this user")
2456
+ now = int(time.time() * 1000)
2457
+ existing = await db.execute("SELECT 1 FROM follows WHERE follower_id = ? AND following_id = ?",
2458
+ (user["id"], target_id))
2459
+ if await existing.fetchone():
2460
+ await db.execute("DELETE FROM follows WHERE follower_id = ? AND following_id = ?",
2461
+ (user["id"], target_id))
2462
+ await db.commit()
2463
+ return {"following": False}
2464
+ await db.execute("INSERT INTO follows (follower_id, following_id, created_at_ms) VALUES (?,?,?)",
2465
+ (user["id"], target_id, now))
2466
+ await _notify_social(db, target_id, user["id"], "follow")
2467
+ await db.commit()
2468
+ return {"following": True}
2469
+ finally:
2470
+ await db.close()
2471
+
2472
+
2473
+ @app.get("/api/social/search")
2474
+ async def social_search_rest(q: str = Query("", max_length=100),
2475
+ token: str = Header(..., alias="X-Auth-Token")):
2476
+ user = await authenticate_user(token)
2477
+ if not user:
2478
+ raise HTTPException(401)
2479
+ q = (q or "").strip()
2480
+ db = await get_db()
2481
+ try:
2482
+ blocked = await _blocked_user_ids(db, user["id"])
2483
+ users = []
2484
+ posts = []
2485
+ if q:
2486
+ like = f"%{q}%"
2487
+ ucur = await db.execute(
2488
+ "SELECT u.id, u.username, u.display_name, u.avatar_path, u.bio, "
2489
+ "(SELECT 1 FROM follows f WHERE f.follower_id = ? AND f.following_id = u.id) AS is_following "
2490
+ "FROM users u WHERE u.id != ? AND (lower(u.username) LIKE ? OR lower(u.display_name) LIKE ?) LIMIT 30",
2491
+ (user["id"], user["id"], like.lower(), like.lower()))
2492
+ for r in await ucur.fetchall():
2493
+ if r["id"] in blocked:
2494
+ continue
2495
+ users.append({
2496
+ "id": r["id"], "username": r["username"], "display_name": r["display_name"],
2497
+ "avatar_path": r["avatar_path"], "bio": r["bio"],
2498
+ "is_following": bool(r["is_following"]),
2499
+ })
2500
+ pcur = await db.execute(
2501
+ "SELECT * FROM social_posts WHERE is_deleted = 0 AND body LIKE ? "
2502
+ "ORDER BY created_at_ms DESC LIMIT 50", (like,))
2503
+ rows = [dict(r) for r in await pcur.fetchall() if r["author_id"] not in blocked]
2504
+ posts = await _serialize_post_rows(db, rows, user["id"])
2505
+ return {"users": users, "posts": posts}
2506
+ finally:
2507
+ await db.close()
2508
+
2509
+
2510
+ @app.get("/api/social/suggestions")
2511
+ async def social_suggestions_rest(limit: int = Query(5, le=20),
2512
+ token: str = Header(..., alias="X-Auth-Token")):
2513
+ user = await authenticate_user(token)
2514
+ if not user:
2515
+ raise HTTPException(401)
2516
+ db = await get_db()
2517
+ try:
2518
+ blocked = await _blocked_user_ids(db, user["id"])
2519
+ cur = await db.execute(
2520
+ "SELECT u.id, u.username, u.display_name, u.avatar_path, u.bio, "
2521
+ "(SELECT COUNT(*) FROM follows f WHERE f.following_id = u.id) AS followers_count "
2522
+ "FROM users u LEFT JOIN follows f ON f.follower_id = ? AND f.following_id = u.id "
2523
+ "WHERE u.id != ? AND u.id NOT IN (SELECT following_id FROM follows WHERE follower_id = ?) "
2524
+ "ORDER BY followers_count DESC, u.username LIMIT ?",
2525
+ (user["id"], user["id"], user["id"], limit)
2526
+ )
2527
+ out = []
2528
+ for r in await cur.fetchall():
2529
+ if r["id"] in blocked:
2530
+ continue
2531
+ out.append({
2532
+ "id": r["id"], "username": r["username"], "display_name": r["display_name"],
2533
+ "avatar_path": r["avatar_path"], "bio": r["bio"], "followers_count": r["followers_count"],
2534
+ "is_following": False,
2535
+ })
2536
+ return {"users": out}
2537
+ finally:
2538
+ await db.close()
2539
+
2540
+
2541
+ @app.get("/api/social/trending")
2542
+ async def social_trending_rest(token: str = Header(..., alias="X-Auth-Token")):
2543
+ user = await authenticate_user(token)
2544
+ if not user:
2545
+ raise HTTPException(401)
2546
+ db = await get_db()
2547
+ try:
2548
+ since = int((time.time() - 7 * 24 * 3600) * 1000)
2549
+ cur = await db.execute(
2550
+ "SELECT body FROM social_posts WHERE is_deleted = 0 AND created_at_ms >= ?", (since,))
2551
+ counts = {}
2552
+ for row in await cur.fetchall():
2553
+ for tag in await _extract_hashtags(row["body"]):
2554
+ t = tag.lower()
2555
+ counts[t] = counts.get(t, 0) + 1
2556
+ top = [{"tag": k, "count": v} for k, v in counts.items()]
2557
+ top.sort(key=lambda x: (-x["count"], x["tag"]))
2558
+ return {"trending": top[:10]}
2559
+ finally:
2560
+ await db.close()
2561
+
2562
+
2563
+ @app.get("/api/social/bookmarks")
2564
+ async def social_bookmarks_rest(limit: int = Query(50, le=100),
2565
+ token: str = Header(..., alias="X-Auth-Token")):
2566
+ user = await authenticate_user(token)
2567
+ if not user:
2568
+ raise HTTPException(401)
2569
+ db = await get_db()
2570
+ try:
2571
+ cur = await db.execute(
2572
+ "SELECT p.* FROM social_bookmarks b JOIN social_posts p ON p.id = b.post_id "
2573
+ "WHERE b.user_id = ? AND p.is_deleted = 0 ORDER BY b.created_at_ms DESC LIMIT ?",
2574
+ (user["id"], limit))
2575
+ rows = [dict(r) for r in await cur.fetchall()]
2576
+ return {"posts": await _serialize_post_rows(db, rows, user["id"])}
2577
+ finally:
2578
+ await db.close()
2579
+
2580
+
2581
+ @app.get("/api/social/notifications")
2582
+ async def social_notifications_rest(token: str = Header(..., alias="X-Auth-Token")):
2583
+ user = await authenticate_user(token)
2584
+ if not user:
2585
+ raise HTTPException(401)
2586
+ db = await get_db()
2587
+ try:
2588
+ blocked = await _blocked_user_ids(db, user["id"])
2589
+ cur = await db.execute(
2590
+ "SELECT n.*, u.username AS actor_username, u.display_name AS actor_display_name, "
2591
+ "u.avatar_path AS actor_avatar_path FROM social_notifications n JOIN users u ON u.id = n.actor_id "
2592
+ "WHERE n.user_id = ? ORDER BY n.created_at_ms DESC LIMIT 100", (user["id"],))
2593
+ out = []
2594
+ for r in await cur.fetchall():
2595
+ if r["actor_id"] in blocked:
2596
+ continue
2597
+ out.append({
2598
+ "id": r["id"], "type": r["type"], "post_id": r["post_id"],
2599
+ "created_at_ms": r["created_at_ms"], "read": bool(r["read"]),
2600
+ "actor": {"id": r["actor_id"], "username": r["actor_username"],
2601
+ "display_name": r["actor_display_name"], "avatar_path": r["actor_avatar_path"]},
2602
+ })
2603
+ return {"notifications": out}
2604
+ finally:
2605
+ await db.close()
2606
+
2607
+
2608
+ @app.post("/api/social/notifications/read")
2609
+ async def mark_notifications_read_rest(token: str = Header(..., alias="X-Auth-Token")):
2610
+ user = await authenticate_user(token)
2611
+ if not user:
2612
+ raise HTTPException(401)
2613
+ db = await get_db()
2614
+ try:
2615
+ await db.execute("UPDATE social_notifications SET read = 1 WHERE user_id = ?", (user["id"],))
2616
+ await db.commit()
2617
+ return {"status": "read"}
2618
+ finally:
2619
+ await db.close()
2620
+
2621
+
2622
+ # ------------------------------------------------------------------------
2623
+ # Admin console (separate /admin page)
2624
+ # ------------------------------------------------------------------------
2625
+ def _admin_configured() -> bool:
2626
+ return bool(ADMIN_USERNAME and ADMIN_PASSWORD)
2627
+
2628
+
2629
+ def _require_admin(admin_token: Optional[str]) -> str:
2630
+ if not _admin_configured():
2631
+ raise HTTPException(503, "Admin login is not configured. Set ADMIN_USERNAME and ADMIN_PASSWORD.")
2632
+ if not admin_token:
2633
+ raise HTTPException(401, "Admin authentication required")
2634
+ now = time.time()
2635
+ with ADMIN_SESSION_LOCK:
2636
+ sess = ADMIN_SESSIONS.get(admin_token)
2637
+ if not sess or sess["expires"] < now:
2638
+ ADMIN_SESSIONS.pop(admin_token, None)
2639
+ raise HTTPException(401, "Admin session expired")
2640
+ return sess.get("username", "admin")
2641
+
2642
+
2643
+ def _cleanup_admin_sessions(now: float = None) -> None:
2644
+ now = now or time.time()
2645
+ with ADMIN_SESSION_LOCK:
2646
+ for key, sess in list(ADMIN_SESSIONS.items()):
2647
+ if sess["expires"] < now:
2648
+ ADMIN_SESSIONS.pop(key, None)
2649
+
2650
+
2651
+ @app.post("/api/admin/login")
2652
+ async def admin_login(request: Request):
2653
+ if not _admin_configured():
2654
+ raise HTTPException(503, "Admin login is not configured. Set ADMIN_USERNAME and ADMIN_PASSWORD.")
2655
+ try:
2656
+ body = await request.json()
2657
+ except Exception:
2658
+ raise HTTPException(400, "Invalid JSON body")
2659
+ username = str(body.get("username", "") or "").strip()
2660
+ password = str(body.get("password", "") or "")
2661
+ if username != ADMIN_USERNAME or password != ADMIN_PASSWORD:
2662
+ raise HTTPException(401, "Invalid admin credentials")
2663
+ _cleanup_admin_sessions()
2664
+ token = secrets.token_urlsafe(32)
2665
+ now = time.time()
2666
+ with ADMIN_SESSION_LOCK:
2667
+ ADMIN_SESSIONS[token] = {
2668
+ "username": username,
2669
+ "created": now,
2670
+ "expires": now + ADMIN_SESSION_TTL,
2671
+ }
2672
+ return {"token": token, "username": username, "expires_in": ADMIN_SESSION_TTL}
2673
+
2674
+
2675
+ @app.post("/api/admin/logout")
2676
+ async def admin_logout(admin_token: str = Header(None, alias="X-Admin-Token")):
2677
+ with ADMIN_SESSION_LOCK:
2678
+ ADMIN_SESSIONS.pop(admin_token, None)
2679
+ return {"status": "logged_out"}
2680
+
2681
+
2682
+ @app.get("/api/admin/status")
2683
+ async def admin_status(admin_token: str = Header(None, alias="X-Admin-Token")):
2684
+ _require_admin(admin_token)
2685
+ backups = await asyncio.to_thread(list_backups)
2686
+ return {
2687
+ "app_version": APP_VERSION,
2688
+ "storage_mode": "local" if _OFFLINE else "bucket",
2689
+ "db_path": DATABASE_URL,
2690
+ "temp_dir": TEMP_DIR,
2691
+ "local_storage_dir": LOCAL_STORAGE_DIR,
2692
+ "backup_count": len(backups),
2693
+ "backup_interval_seconds": 3600,
2694
+ "admin_session_ttl_seconds": ADMIN_SESSION_TTL,
2695
+ }
2696
+
2697
+
2698
+ @app.get("/api/admin/backups")
2699
+ async def admin_backups(admin_token: str = Header(None, alias="X-Admin-Token")):
2700
+ _require_admin(admin_token)
2701
+ return {"backups": await asyncio.to_thread(list_backups), "offline": _OFFLINE}
2702
+
2703
+
2704
+ @app.post("/api/admin/backups/create")
2705
+ async def admin_create_backup(admin_token: str = Header(None, alias="X-Admin-Token")):
2706
+ _require_admin(admin_token)
2707
+ path = await asyncio.to_thread(create_timestamped_backup, DATABASE_URL, "manual")
2708
+ if not path:
2709
+ raise HTTPException(503, "Backup unavailable (storage offline or failure)")
2710
+ return {"status": "created", "path": path}
2711
+
2712
+
2713
+ @app.post("/api/admin/backups/{backup_prefix:path}/restore")
2714
+ async def admin_restore_backup(backup_prefix: str, admin_token: str = Header(None, alias="X-Admin-Token")):
2715
+ _require_admin(admin_token)
2716
+ try:
2717
+ count = await asyncio.to_thread(restore_backup, DATABASE_URL, backup_prefix)
2718
+ except StorageUnavailableError as e:
2719
+ raise HTTPException(503, str(e))
2720
+ except FileNotFoundError as e:
2721
+ raise HTTPException(404, str(e))
2722
+ except Exception as e:
2723
+ raise HTTPException(500, f"Restore failed: {e}")
2724
+ return {"status": "restored", "files": count, "restart_recommended": True}
2725
+
2726
+
2727
+ @app.delete("/api/admin/backups/{backup_prefix:path}")
2728
+ async def admin_delete_backup(backup_prefix: str, admin_token: str = Header(None, alias="X-Admin-Token")):
2729
+ _require_admin(admin_token)
2730
+ try:
2731
+ await asyncio.to_thread(delete_backup, backup_prefix)
2732
+ except ValueError as e:
2733
+ raise HTTPException(400, str(e))
2734
+ except StorageUnavailableError as e:
2735
+ raise HTTPException(503, str(e))
2736
+ except Exception as e:
2737
+ raise HTTPException(500, f"Delete failed: {e}")
2738
+ return {"status": "deleted", "path": backup_prefix}
2739
+
2740
+
2741
  # ------------------------------------------------------------------------
2742
  # Messages (REST fallbacks for reliable edit/delete/receipts)
2743
  # ------------------------------------------------------------------------
 
2998
  with open(assembled_path, 'rb') as f:
2999
  full_data = f.read()
3000
  # Blocking bucket write happens off the event loop
3001
+ try:
3002
+ await asyncio.to_thread(store_file, remote_path, full_data)
3003
+ except StorageUnavailableError:
3004
+ raise HTTPException(503, "File storage is unavailable")
3005
+ finally:
3006
+ del full_data
3007
 
3008
  finally:
3009
  shutil.rmtree(session_dir, ignore_errors=True)
 
3040
  if '..' in file_path or '\\' in file_path:
3041
  raise HTTPException(400, "Invalid path")
3042
 
3043
+ # Avatars are public profile data. Social attachments are linked to a post.
3044
+ # Chat files belong to a conversation, so a user blocked out of it must not
3045
+ # be able to keep downloading those files.
3046
+ social_media = False
3047
  if not file_path.startswith("avatars/"):
3048
  db = await get_db()
3049
  try:
3050
+ # Social-media attachment: any signed-in user may fetch a post's media.
3051
+ sp = await db.execute(
3052
+ "SELECT 1 FROM social_posts WHERE is_deleted = 0 AND media_json LIKE ? LIMIT 1",
3053
+ (f'%{file_path}%',)
3054
  )
3055
+ if await sp.fetchone() is not None:
3056
+ social_media = True
3057
+ if not social_media:
3058
+ cursor = await db.execute(
3059
+ "SELECT conversation_id FROM messages WHERE file_path = ? ORDER BY id DESC LIMIT 1",
3060
+ (file_path,)
3061
+ )
3062
+ row = await cursor.fetchone()
3063
+ if not row:
3064
+ raise HTTPException(404, "File not found")
3065
+ if not await user_in_conversation(db, user['id'], row["conversation_id"]):
3066
+ raise HTTPException(403, "You can no longer access this file")
3067
  finally:
3068
  await db.close()
3069
 
 
3713
 
3714
  @app.get("/")
3715
  async def root():
3716
+ return FileResponse("static/index.html", headers={"Cache-Control": "no-cache"})
3717
+
3718
+
3719
+ @app.get("/admin", include_in_schema=False)
3720
+ async def admin_page():
3721
+ return FileResponse("static/admin.html", headers={"Cache-Control": "no-cache"})