dbadmin.py 18 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603
  1. from __future__ import annotations
  2. import asyncio
  3. import re
  4. from datetime import UTC, datetime
  5. from secrets import token_hex
  6. from typing import Any
  7. from pymongo import ASCENDING, DESCENDING
  8. from wbb import BOT_PROFILE_ID, control_db, db
  9. admin_usersdb = control_db.admin_users
  10. admin_sessionsdb = control_db.admin_sessions
  11. auditdb = db.admin_audit_logs
  12. managed_chatsdb = db.managed_chats
  13. managed_chat_settingsdb = db.managed_chat_settings
  14. invite_linksdb = db.admin_invite_links
  15. recent_chat_membersdb = db.recent_chat_members
  16. member_identity_changesdb = db.member_identity_changes
  17. _index_lock = asyncio.Lock()
  18. _indexes_ready = False
  19. def utc_now() -> datetime:
  20. return datetime.now(UTC)
  21. async def ensure_admin_indexes() -> None:
  22. global _indexes_ready
  23. if _indexes_ready:
  24. return
  25. async with _index_lock:
  26. if _indexes_ready:
  27. return
  28. await admin_usersdb.create_index([("username", ASCENDING)], unique=True)
  29. await admin_sessionsdb.create_index([("token_hash", ASCENDING)], unique=True)
  30. await admin_sessionsdb.create_index("expires_at", expireAfterSeconds=0)
  31. await auditdb.create_index([("created_at", DESCENDING)])
  32. await auditdb.create_index(
  33. [("chat_id", ASCENDING), ("created_at", DESCENDING)]
  34. )
  35. await managed_chatsdb.create_index(
  36. [("bot_id", ASCENDING), ("chat_id", ASCENDING)], unique=True
  37. )
  38. await managed_chatsdb.create_index([("last_seen_at", DESCENDING)])
  39. await managed_chat_settingsdb.create_index(
  40. [("bot_id", ASCENDING), ("chat_id", ASCENDING)], unique=True
  41. )
  42. await invite_linksdb.create_index(
  43. [("bot_id", ASCENDING), ("chat_id", ASCENDING), ("created_at", DESCENDING)]
  44. )
  45. await recent_chat_membersdb.create_index(
  46. [
  47. ("bot_id", ASCENDING),
  48. ("chat_id", ASCENDING),
  49. ("user_id", ASCENDING),
  50. ],
  51. unique=True,
  52. )
  53. await recent_chat_membersdb.create_index(
  54. [("bot_id", ASCENDING), ("chat_id", ASCENDING), ("last_seen_at", DESCENDING)]
  55. )
  56. await member_identity_changesdb.create_index(
  57. [("bot_id", ASCENDING), ("chat_id", ASCENDING), ("observed_at", DESCENDING)]
  58. )
  59. await member_identity_changesdb.create_index(
  60. [
  61. ("bot_id", ASCENDING),
  62. ("chat_id", ASCENDING),
  63. ("user_id", ASCENDING),
  64. ("observed_at", DESCENDING),
  65. ]
  66. )
  67. _indexes_ready = True
  68. async def ensure_default_admin(
  69. *, username: str, password_hash: str
  70. ) -> dict[str, Any]:
  71. await ensure_admin_indexes()
  72. now = utc_now()
  73. await admin_usersdb.update_one(
  74. {"username": username},
  75. {
  76. "$setOnInsert": {
  77. "username": username,
  78. "password_hash": password_hash,
  79. "must_change_password": True,
  80. "failed_login_count": 0,
  81. "created_at": now,
  82. "updated_at": now,
  83. }
  84. },
  85. upsert=True,
  86. )
  87. return await admin_usersdb.find_one({"username": username})
  88. async def get_admin_user(username: str) -> dict[str, Any] | None:
  89. await ensure_admin_indexes()
  90. return await admin_usersdb.find_one({"username": username})
  91. async def record_login_failure(username: str) -> None:
  92. await ensure_admin_indexes()
  93. await admin_usersdb.update_one(
  94. {"username": username},
  95. {
  96. "$inc": {"failed_login_count": 1},
  97. "$set": {"last_failed_login_at": utc_now(), "updated_at": utc_now()},
  98. },
  99. )
  100. async def record_login_success(username: str) -> None:
  101. await ensure_admin_indexes()
  102. await admin_usersdb.update_one(
  103. {"username": username},
  104. {
  105. "$set": {
  106. "failed_login_count": 0,
  107. "last_login_at": utc_now(),
  108. "updated_at": utc_now(),
  109. }
  110. },
  111. )
  112. async def update_admin_password(username: str, password_hash: str) -> bool:
  113. await ensure_admin_indexes()
  114. result = await admin_usersdb.update_one(
  115. {"username": username},
  116. {
  117. "$set": {
  118. "password_hash": password_hash,
  119. "must_change_password": False,
  120. "password_updated_at": utc_now(),
  121. "updated_at": utc_now(),
  122. }
  123. },
  124. )
  125. if result.modified_count:
  126. await admin_sessionsdb.delete_many({"username": username})
  127. return result.modified_count == 1
  128. async def create_admin_session(
  129. *,
  130. username: str,
  131. token_hash: str,
  132. csrf_token: str,
  133. expires_at: datetime,
  134. remote_address: str,
  135. user_agent: str,
  136. ) -> None:
  137. await ensure_admin_indexes()
  138. now = utc_now()
  139. await admin_sessionsdb.insert_one(
  140. {
  141. "username": username,
  142. "token_hash": token_hash,
  143. "csrf_token": csrf_token,
  144. "remote_address": remote_address[:200],
  145. "user_agent": user_agent[:500],
  146. "created_at": now,
  147. "last_seen_at": now,
  148. "expires_at": expires_at,
  149. }
  150. )
  151. async def get_admin_session(token_hash: str) -> dict[str, Any] | None:
  152. await ensure_admin_indexes()
  153. session = await admin_sessionsdb.find_one(
  154. {"token_hash": token_hash, "expires_at": {"$gt": utc_now()}}
  155. )
  156. if not session:
  157. return None
  158. user = await get_admin_user(session["username"])
  159. if not user:
  160. return None
  161. await admin_sessionsdb.update_one(
  162. {"_id": session["_id"]}, {"$set": {"last_seen_at": utc_now()}}
  163. )
  164. return {**session, "user": user}
  165. async def revoke_admin_session(token_hash: str) -> None:
  166. await ensure_admin_indexes()
  167. await admin_sessionsdb.delete_one({"token_hash": token_hash})
  168. async def revoke_admin_sessions(username: str) -> None:
  169. await ensure_admin_indexes()
  170. await admin_sessionsdb.delete_many({"username": username})
  171. async def upsert_managed_chat(
  172. *,
  173. chat_id: int,
  174. title: str | None,
  175. username: str | None,
  176. chat_type: str,
  177. member_count: int | None = None,
  178. accessible: bool | None = True,
  179. bot_status: str | None = None,
  180. bot_privileges: list[str] | None = None,
  181. ) -> None:
  182. await ensure_admin_indexes()
  183. now = utc_now()
  184. values: dict[str, Any] = {
  185. "title": title or str(chat_id),
  186. "username": username,
  187. "type": chat_type,
  188. "last_seen_at": now,
  189. "updated_at": now,
  190. }
  191. if accessible is not None:
  192. values["accessible"] = accessible
  193. if member_count is not None:
  194. values["member_count"] = int(member_count)
  195. if bot_status is not None:
  196. values["bot_status"] = bot_status
  197. if bot_privileges is not None:
  198. values["bot_privileges"] = bot_privileges
  199. await managed_chatsdb.update_one(
  200. {"bot_id": BOT_PROFILE_ID, "chat_id": chat_id},
  201. {
  202. "$set": values,
  203. "$setOnInsert": {
  204. "bot_id": BOT_PROFILE_ID,
  205. "created_at": now,
  206. **({"accessible": True} if accessible is None else {}),
  207. },
  208. },
  209. upsert=True,
  210. )
  211. async def mark_managed_chat_unavailable(chat_id: int, reason: str) -> None:
  212. await ensure_admin_indexes()
  213. await managed_chatsdb.update_one(
  214. {"bot_id": BOT_PROFILE_ID, "chat_id": chat_id},
  215. {
  216. "$set": {
  217. "accessible": False,
  218. "last_error": reason[:500],
  219. "updated_at": utc_now(),
  220. }
  221. },
  222. upsert=True,
  223. )
  224. async def delete_unavailable_managed_chat(chat_id: int) -> bool:
  225. """Remove an unavailable chat from this bot's management index only."""
  226. await ensure_admin_indexes()
  227. result = await managed_chatsdb.delete_one(
  228. {
  229. "bot_id": BOT_PROFILE_ID,
  230. "chat_id": int(chat_id),
  231. "accessible": False,
  232. }
  233. )
  234. return result.deleted_count == 1
  235. async def managed_chat_exists(chat_id: int) -> bool:
  236. await ensure_admin_indexes()
  237. return (
  238. await managed_chatsdb.find_one(
  239. {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id)}, {"_id": 1}
  240. )
  241. is not None
  242. )
  243. async def list_managed_chats(
  244. *, query: str = "", page: int = 1, page_size: int = 20,
  245. chat_type: str | None = None,
  246. ) -> tuple[list[dict[str, Any]], int]:
  247. await ensure_admin_indexes()
  248. filters: dict[str, Any] = {"bot_id": BOT_PROFILE_ID}
  249. if chat_type == "channel":
  250. filters["type"] = "channel"
  251. filters["bot_status"] = {"$in": ["owner", "administrator"]}
  252. elif chat_type == "group":
  253. filters["type"] = {"$in": ["group", "supergroup"]}
  254. query = query.strip()
  255. if query:
  256. if query.lstrip("-").isdigit():
  257. filters["chat_id"] = int(query)
  258. else:
  259. value = re.escape(query.lstrip("@"))
  260. filters["$or"] = [
  261. {"title": {"$regex": value, "$options": "i"}},
  262. {"username": {"$regex": value, "$options": "i"}},
  263. ]
  264. page = max(1, page)
  265. page_size = max(1, min(page_size, 100))
  266. total = await managed_chatsdb.count_documents(filters)
  267. cursor = (
  268. managed_chatsdb.find(filters)
  269. .sort([("accessible", DESCENDING), ("last_seen_at", DESCENDING)])
  270. .skip((page - 1) * page_size)
  271. .limit(page_size)
  272. )
  273. return [doc async for doc in cursor], total
  274. async def record_audit(
  275. *,
  276. source: str,
  277. actor_id: int | str | None,
  278. actor_name: str,
  279. action: str,
  280. chat_id: int | None = None,
  281. target_id: int | str | None = None,
  282. summary: str = "",
  283. success: bool = True,
  284. error: str = "",
  285. metadata: dict[str, Any] | None = None,
  286. ) -> None:
  287. await ensure_admin_indexes()
  288. await auditdb.insert_one(
  289. {
  290. "source": source,
  291. "actor_id": actor_id,
  292. "actor_name": actor_name,
  293. "action": action,
  294. "chat_id": chat_id,
  295. "target_id": target_id,
  296. "summary": summary[:1000],
  297. "success": bool(success),
  298. "error": error[:1000],
  299. "metadata": metadata or {},
  300. "bot_id": BOT_PROFILE_ID,
  301. "created_at": utc_now(),
  302. }
  303. )
  304. async def list_audit_logs(
  305. *,
  306. chat_id: int | None = None,
  307. action: str | None = None,
  308. page: int = 1,
  309. page_size: int = 20,
  310. ) -> tuple[list[dict[str, Any]], int]:
  311. await ensure_admin_indexes()
  312. filters: dict[str, Any] = {}
  313. if chat_id is not None:
  314. filters["chat_id"] = chat_id
  315. if action:
  316. filters["action"] = action
  317. page = max(1, page)
  318. page_size = max(1, min(page_size, 100))
  319. total = await auditdb.count_documents(filters)
  320. cursor = (
  321. auditdb.find(filters)
  322. .sort("created_at", DESCENDING)
  323. .skip((page - 1) * page_size)
  324. .limit(page_size)
  325. )
  326. return [doc async for doc in cursor], total
  327. async def get_managed_chat_settings(chat_id: int) -> dict[str, Any]:
  328. await ensure_admin_indexes()
  329. doc = await managed_chat_settingsdb.find_one(
  330. {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id)}
  331. )
  332. if doc is not None:
  333. doc.setdefault("captcha_enabled", False)
  334. doc.setdefault("interaction_settings", {
  335. "cleanup_seconds": 10,
  336. "commands": {key: True for key in ("giveaway", "checkin", "points", "history", "shortcuts")},
  337. })
  338. return doc
  339. return {
  340. "chat_id": int(chat_id),
  341. "auto_replies": [],
  342. "blacklist_words": [],
  343. "welcome": {"enabled": False, "text": "", "media": None},
  344. "identity_monitor": {"enabled": False, "notify_in_chat": False},
  345. "captcha_enabled": False,
  346. "chatbot_enabled": False,
  347. "antiflood_enabled": True,
  348. "interaction_settings": {
  349. "cleanup_seconds": 10,
  350. "commands": {key: True for key in ("giveaway", "checkin", "points", "history", "shortcuts")},
  351. },
  352. }
  353. async def update_managed_chat_settings(
  354. chat_id: int, values: dict[str, Any]
  355. ) -> dict[str, Any]:
  356. await ensure_admin_indexes()
  357. allowed = {
  358. "auto_replies",
  359. "blacklist_words",
  360. "risk_control",
  361. "risk_rules",
  362. "identity_monitor",
  363. "welcome",
  364. "captcha_enabled",
  365. "chatbot_enabled",
  366. "antiflood_enabled",
  367. "interaction_settings",
  368. }
  369. update = {key: value for key, value in values.items() if key in allowed}
  370. now = utc_now()
  371. await managed_chat_settingsdb.update_one(
  372. {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id)},
  373. {
  374. "$set": {**update, "updated_at": now},
  375. "$setOnInsert": {
  376. "bot_id": BOT_PROFILE_ID,
  377. "chat_id": int(chat_id),
  378. "created_at": now,
  379. },
  380. },
  381. upsert=True,
  382. )
  383. return await get_managed_chat_settings(chat_id)
  384. async def upsert_recent_chat_member(
  385. *,
  386. chat_id: int,
  387. user_id: int,
  388. username: str | None,
  389. first_name: str | None,
  390. last_name: str | None,
  391. is_bot: bool,
  392. ) -> None:
  393. await ensure_admin_indexes()
  394. now = utc_now()
  395. await recent_chat_membersdb.update_one(
  396. {
  397. "bot_id": BOT_PROFILE_ID,
  398. "chat_id": int(chat_id),
  399. "user_id": int(user_id),
  400. },
  401. {
  402. "$set": {
  403. "username": username,
  404. "first_name": first_name,
  405. "last_name": last_name,
  406. "is_bot": bool(is_bot),
  407. "last_seen_at": now,
  408. },
  409. "$inc": {"message_count": 1},
  410. "$setOnInsert": {
  411. "bot_id": BOT_PROFILE_ID,
  412. "chat_id": int(chat_id),
  413. "user_id": int(user_id),
  414. "created_at": now,
  415. },
  416. },
  417. upsert=True,
  418. )
  419. async def get_recent_chat_member(chat_id: int, user_id: int) -> dict[str, Any] | None:
  420. await ensure_admin_indexes()
  421. return await recent_chat_membersdb.find_one(
  422. {
  423. "bot_id": BOT_PROFILE_ID,
  424. "chat_id": int(chat_id),
  425. "user_id": int(user_id),
  426. }
  427. )
  428. async def record_member_identity_change(
  429. *,
  430. chat_id: int,
  431. user_id: int,
  432. before: dict[str, Any],
  433. after: dict[str, Any],
  434. changed_fields: list[str],
  435. ) -> dict[str, Any]:
  436. await ensure_admin_indexes()
  437. event = {
  438. "change_id": token_hex(10),
  439. "bot_id": BOT_PROFILE_ID,
  440. "chat_id": int(chat_id),
  441. "user_id": int(user_id),
  442. "before": before,
  443. "after": after,
  444. "changed_fields": changed_fields,
  445. "observed_at": utc_now(),
  446. }
  447. await member_identity_changesdb.insert_one(event)
  448. return event
  449. async def list_member_identity_changes(
  450. chat_id: int,
  451. *,
  452. page: int = 1,
  453. page_size: int = 20,
  454. ) -> tuple[list[dict[str, Any]], int]:
  455. await ensure_admin_indexes()
  456. filters = {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id)}
  457. page = max(1, int(page))
  458. page_size = max(1, min(int(page_size), 100))
  459. total = await member_identity_changesdb.count_documents(filters)
  460. cursor = (
  461. member_identity_changesdb.find(filters)
  462. .sort("observed_at", DESCENDING)
  463. .skip((page - 1) * page_size)
  464. .limit(page_size)
  465. )
  466. return [item async for item in cursor], total
  467. async def list_recent_chat_members(
  468. chat_id: int,
  469. *,
  470. query: str = "",
  471. limit: int = 30,
  472. ) -> list[dict[str, Any]]:
  473. await ensure_admin_indexes()
  474. filters: dict[str, Any] = {
  475. "bot_id": BOT_PROFILE_ID,
  476. "chat_id": int(chat_id),
  477. "is_bot": {"$ne": True},
  478. }
  479. normalized_query = query.strip().lstrip("@")
  480. if normalized_query:
  481. if normalized_query.isdigit():
  482. filters["user_id"] = int(normalized_query)
  483. else:
  484. pattern = re.escape(normalized_query)
  485. filters["$or"] = [
  486. {"username": {"$regex": pattern, "$options": "i"}},
  487. {"first_name": {"$regex": pattern, "$options": "i"}},
  488. {"last_name": {"$regex": pattern, "$options": "i"}},
  489. ]
  490. cursor = (
  491. recent_chat_membersdb.find(filters)
  492. .sort("last_seen_at", DESCENDING)
  493. .limit(max(1, min(int(limit), 50)))
  494. )
  495. return [item async for item in cursor]
  496. async def store_invite_link(
  497. *, chat_id: int, invite_link: str, name: str, expires_at: datetime | None
  498. ) -> None:
  499. await ensure_admin_indexes()
  500. await invite_linksdb.update_one(
  501. {
  502. "bot_id": BOT_PROFILE_ID,
  503. "chat_id": int(chat_id),
  504. "invite_link": invite_link,
  505. },
  506. {
  507. "$set": {
  508. "name": name,
  509. "expires_at": expires_at,
  510. "revoked": False,
  511. "updated_at": utc_now(),
  512. },
  513. "$setOnInsert": {"bot_id": BOT_PROFILE_ID, "created_at": utc_now()},
  514. },
  515. upsert=True,
  516. )
  517. async def revoke_stored_invite_link(chat_id: int, invite_link: str) -> None:
  518. await ensure_admin_indexes()
  519. await invite_linksdb.update_one(
  520. {
  521. "bot_id": BOT_PROFILE_ID,
  522. "chat_id": int(chat_id),
  523. "invite_link": invite_link,
  524. },
  525. {"$set": {"revoked": True, "revoked_at": utc_now()}},
  526. )
  527. async def list_stored_invite_links(chat_id: int) -> list[dict[str, Any]]:
  528. await ensure_admin_indexes()
  529. cursor = invite_linksdb.find(
  530. {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id)}
  531. ).sort(
  532. "created_at", DESCENDING
  533. )
  534. return [doc async for doc in cursor]
  535. async def dashboard_counts() -> dict[str, int]:
  536. await ensure_admin_indexes()
  537. return {
  538. "chats": await managed_chatsdb.count_documents({"accessible": {"$ne": False}}),
  539. "audit_events": await auditdb.count_documents({}),
  540. }