directory.py 17 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528
  1. from __future__ import annotations
  2. import asyncio
  3. from contextlib import suppress
  4. from typing import Any
  5. from aiohttp import ClientError, ClientSession, ClientTimeout
  6. from pyrogram.enums import ChatMemberStatus
  7. import wbb
  8. from wbb import BOT_PROFILE_ID, SUDOERS, app, log
  9. from wbb.utils.dbdirectory import (
  10. DirectoryDataError,
  11. aware_utc,
  12. clear_directory_location,
  13. decide_teacher_application,
  14. expire_teacher_presence,
  15. get_directory_location,
  16. get_directory_profile,
  17. list_directory_teachers,
  18. list_membership_candidates,
  19. save_directory_location,
  20. set_teacher_state,
  21. submit_teacher_application,
  22. teacher_is_online,
  23. upsert_directory_identity,
  24. upsert_directory_membership,
  25. )
  26. ACTIVE_MEMBER_STATUSES = {
  27. ChatMemberStatus.OWNER,
  28. ChatMemberStatus.ADMINISTRATOR,
  29. ChatMemberStatus.MEMBER,
  30. ChatMemberStatus.RESTRICTED,
  31. }
  32. ADMIN_MEMBER_STATUSES = {
  33. ChatMemberStatus.OWNER,
  34. ChatMemberStatus.ADMINISTRATOR,
  35. }
  36. _presence_task: asyncio.Task[None] | None = None
  37. class DirectoryServiceError(RuntimeError):
  38. def __init__(self, code: str, message: str, *, status: int = 400):
  39. super().__init__(message)
  40. self.code = code
  41. self.status = status
  42. def status_value(value: Any) -> str:
  43. return str(getattr(value, "value", value or "")).lower()
  44. def is_sudoer(user_id: int) -> bool:
  45. try:
  46. return int(user_id) in SUDOERS
  47. except (TypeError, ValueError):
  48. return False
  49. def user_identity(user: Any) -> dict[str, Any]:
  50. return {
  51. "user_id": int(user.id),
  52. "username": getattr(user, "username", None),
  53. "first_name": getattr(user, "first_name", None),
  54. "last_name": getattr(user, "last_name", None),
  55. "display_name": " ".join(
  56. value.strip()
  57. for value in (
  58. getattr(user, "first_name", None),
  59. getattr(user, "last_name", None),
  60. )
  61. if value and value.strip()
  62. )
  63. or (
  64. f"@{user.username}"
  65. if getattr(user, "username", None)
  66. else f"用户 {user.id}"
  67. ),
  68. }
  69. async def observe_directory_user(user: Any, *, source: str = "telegram") -> dict[str, Any]:
  70. identity = user_identity(user)
  71. return await upsert_directory_identity(
  72. user_id=identity["user_id"],
  73. username=identity["username"],
  74. first_name=identity["first_name"],
  75. last_name=identity["last_name"],
  76. source=source,
  77. )
  78. async def observe_group_member(
  79. *,
  80. chat_id: int,
  81. chat_title: str,
  82. user: Any,
  83. status: str = "member",
  84. active: bool = True,
  85. verified: bool = False,
  86. ) -> dict[str, Any]:
  87. identity = user_identity(user)
  88. await observe_directory_user(user)
  89. return await upsert_directory_membership(
  90. bot_id=str(BOT_PROFILE_ID),
  91. chat_id=int(chat_id),
  92. user_id=identity["user_id"],
  93. status=status,
  94. active=active,
  95. chat_title=chat_title,
  96. username=identity["username"],
  97. display_name_value=identity["display_name"],
  98. verified=verified,
  99. )
  100. async def verify_local_membership(
  101. *,
  102. user_id: int,
  103. chat_ids: list[int],
  104. require_admin: bool = False,
  105. ) -> dict[str, Any]:
  106. accepted = ADMIN_MEMBER_STATUSES if require_admin else ACTIVE_MEMBER_STATUSES
  107. for chat_id in dict.fromkeys(int(value) for value in chat_ids):
  108. try:
  109. member = await app.get_chat_member(int(chat_id), int(user_id))
  110. except Exception:
  111. continue
  112. status = getattr(member, "status", None)
  113. active = status in ACTIVE_MEMBER_STATUSES
  114. user = getattr(member, "user", None)
  115. if user is not None:
  116. with suppress(Exception):
  117. await observe_group_member(
  118. chat_id=int(chat_id),
  119. chat_title="",
  120. user=user,
  121. status=status_value(status),
  122. active=active,
  123. verified=True,
  124. )
  125. if status in accepted:
  126. return {
  127. "allowed": True,
  128. "bot_id": str(BOT_PROFILE_ID),
  129. "chat_id": int(chat_id),
  130. "status": status_value(status),
  131. }
  132. return {"allowed": False}
  133. async def verify_platform_membership(
  134. *,
  135. user_id: int,
  136. require_admin: bool = False,
  137. ) -> dict[str, Any]:
  138. if require_admin and is_sudoer(int(user_id)):
  139. return {
  140. "allowed": True,
  141. "bot_id": str(BOT_PROFILE_ID),
  142. "chat_id": None,
  143. "status": "sudoer",
  144. }
  145. candidates = await list_membership_candidates(int(user_id))
  146. if not candidates:
  147. return {"allowed": False, "reason": "membership_evidence_required"}
  148. internal_url = str(
  149. getattr(wbb, "SUPERVISOR_INTERNAL_URL", "") or ""
  150. ).rstrip("/")
  151. internal_token = str(getattr(wbb, "INTERNAL_TOKEN", "") or "")
  152. if internal_url and internal_token:
  153. try:
  154. async with ClientSession(timeout=ClientTimeout(total=20)) as session:
  155. async with session.post(
  156. f"{internal_url}/api/internal/v1/directory/verify-platform",
  157. headers={"Authorization": f"Bearer {internal_token}"},
  158. json={
  159. "user_id": str(user_id),
  160. "require_admin": bool(require_admin),
  161. },
  162. ) as response:
  163. if response.status != 200:
  164. return {"allowed": False, "reason": "verification_unavailable"}
  165. payload = await response.json()
  166. return payload.get("data") or {"allowed": False}
  167. except (ClientError, TimeoutError, ValueError) as exc:
  168. log.error(f"跨 Bot 成员身份校验失败:{exc}")
  169. return {"allowed": False, "reason": "verification_unavailable"}
  170. local_chat_ids = [
  171. int(item["chat_id"])
  172. for item in candidates
  173. if str(item.get("bot_id")) == str(BOT_PROFILE_ID)
  174. ]
  175. if not local_chat_ids:
  176. return {"allowed": False, "reason": "verification_unavailable"}
  177. return await verify_local_membership(
  178. user_id=int(user_id),
  179. chat_ids=local_chat_ids,
  180. require_admin=require_admin,
  181. )
  182. async def require_platform_membership(
  183. user_id: int,
  184. *,
  185. require_admin: bool = False,
  186. ) -> dict[str, Any]:
  187. result = await verify_platform_membership(
  188. user_id=int(user_id),
  189. require_admin=require_admin,
  190. )
  191. if result.get("allowed"):
  192. return result
  193. reason = result.get("reason")
  194. if reason == "membership_evidence_required":
  195. raise DirectoryServiceError(
  196. "membership_verification_required",
  197. "请先在任一受管群发送 /directory 完成群成员身份验证。",
  198. status=403,
  199. )
  200. if reason == "verification_unavailable":
  201. raise DirectoryServiceError(
  202. "membership_verification_unavailable",
  203. "暂时无法通过对应机器人核验群身份,请稍后重试。",
  204. status=503,
  205. )
  206. raise DirectoryServiceError(
  207. "managed_group_member_required",
  208. "该功能仅对受管群的当前成员开放。",
  209. status=403,
  210. )
  211. async def update_user_location(
  212. *,
  213. user: Any,
  214. longitude: float,
  215. latitude: float,
  216. source: str = "telegram_location",
  217. ) -> dict[str, Any]:
  218. await require_platform_membership(int(user.id))
  219. identity = user_identity(user)
  220. await observe_directory_user(user)
  221. try:
  222. return await save_directory_location(
  223. user_id=int(user.id),
  224. longitude=longitude,
  225. latitude=latitude,
  226. source=source,
  227. actor_id=int(user.id),
  228. actor_name=identity["display_name"],
  229. bot_id=str(BOT_PROFILE_ID),
  230. )
  231. except DirectoryDataError as exc:
  232. raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
  233. async def clear_user_location(*, user: Any) -> bool:
  234. identity = user_identity(user)
  235. return await clear_directory_location(
  236. user_id=int(user.id),
  237. actor_id=int(user.id),
  238. actor_name=identity["display_name"],
  239. source="telegram_private",
  240. reason="用户主动清除位置",
  241. )
  242. async def apply_as_teacher(
  243. *,
  244. user: Any,
  245. chat_id: int | None = None,
  246. chat_title: str = "",
  247. source: str = "telegram_private",
  248. ) -> tuple[dict[str, Any], bool]:
  249. authorization = await require_platform_membership(int(user.id))
  250. await observe_directory_user(user)
  251. application_chat_id = int(chat_id or authorization.get("chat_id") or 0)
  252. if not application_chat_id:
  253. raise DirectoryServiceError(
  254. "application_chat_required",
  255. "请先在要申请技师的群里打开附近技师菜单。",
  256. status=409,
  257. )
  258. application_chat_title = str(chat_title or "").strip()
  259. if not application_chat_title:
  260. with suppress(Exception):
  261. chat = await app.get_chat(application_chat_id)
  262. application_chat_title = str(getattr(chat, "title", "") or "").strip()
  263. try:
  264. return await submit_teacher_application(
  265. user_id=int(user.id),
  266. source=source,
  267. bot_id=str(BOT_PROFILE_ID),
  268. chat_id=application_chat_id,
  269. chat_title=application_chat_title,
  270. )
  271. except DirectoryDataError as exc:
  272. raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
  273. async def change_own_teacher_state(
  274. *,
  275. user: Any,
  276. action: str,
  277. source: str = "telegram_private",
  278. idempotency_key: str | None = None,
  279. chat_id: int | None = None,
  280. ) -> tuple[dict[str, Any], bool]:
  281. await require_platform_membership(int(user.id))
  282. identity = user_identity(user)
  283. await observe_directory_user(user)
  284. try:
  285. return await set_teacher_state(
  286. user_id=int(user.id),
  287. action=action,
  288. actor_id=int(user.id),
  289. actor_name=identity["display_name"],
  290. source=source,
  291. reason="技师自助操作",
  292. idempotency_key=idempotency_key,
  293. chat_id=chat_id,
  294. )
  295. except DirectoryDataError as exc:
  296. raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
  297. async def admin_source_chat_ids(actor: Any) -> list[int] | None:
  298. if is_sudoer(int(actor.id)):
  299. return None
  300. candidates = await list_membership_candidates(int(actor.id))
  301. chat_ids: list[int] = []
  302. for item in candidates:
  303. if str(item.get("bot_id")) != str(BOT_PROFILE_ID):
  304. continue
  305. chat_id = int(item["chat_id"])
  306. try:
  307. member = await app.get_chat_member(chat_id, int(actor.id))
  308. except Exception:
  309. continue
  310. if getattr(member, "status", None) in ADMIN_MEMBER_STATUSES:
  311. chat_ids.append(chat_id)
  312. values = list(dict.fromkeys(chat_ids))
  313. if not values:
  314. raise DirectoryServiceError(
  315. "managed_group_admin_required",
  316. "只有申请来源群的管理员可以审批技师。",
  317. status=403,
  318. )
  319. return values
  320. async def require_teacher_admin(actor: Any, user_id: int) -> dict[str, Any]:
  321. profile = await get_directory_profile(int(user_id))
  322. if not profile:
  323. raise DirectoryServiceError("profile_not_found", "未找到该成员资料。", status=404)
  324. if is_sudoer(int(actor.id)):
  325. return profile
  326. source_chat_id = profile.get("application_source_chat_id")
  327. source_bot_id = profile.get("application_source_bot_id")
  328. if not source_chat_id:
  329. await require_platform_membership(int(actor.id), require_admin=True)
  330. return profile
  331. if source_bot_id and str(source_bot_id) != str(BOT_PROFILE_ID):
  332. raise DirectoryServiceError(
  333. "source_bot_required",
  334. "请使用接收技师申请提醒的机器人完成审批。",
  335. status=403,
  336. )
  337. try:
  338. member = await app.get_chat_member(int(source_chat_id), int(actor.id))
  339. except Exception as exc:
  340. raise DirectoryServiceError(
  341. "admin_verification_unavailable",
  342. "暂时无法核验申请来源群的管理员身份。",
  343. status=503,
  344. ) from exc
  345. if getattr(member, "status", None) not in ADMIN_MEMBER_STATUSES:
  346. raise DirectoryServiceError(
  347. "source_chat_admin_required",
  348. "只有申请来源群的管理员可以执行此操作。",
  349. status=403,
  350. )
  351. return profile
  352. async def admin_decide_teacher(
  353. *,
  354. actor: Any,
  355. user_id: int,
  356. action: str,
  357. reason: str = "",
  358. ) -> dict[str, Any]:
  359. profile = await require_teacher_admin(actor, int(user_id))
  360. identity = user_identity(actor)
  361. try:
  362. return await decide_teacher_application(
  363. user_id=int(user_id),
  364. action=action,
  365. actor_id=int(actor.id),
  366. actor_name=identity["display_name"],
  367. source="telegram_private",
  368. reason=reason,
  369. authorization_chat_id=profile.get("application_source_chat_id"),
  370. )
  371. except DirectoryDataError as exc:
  372. raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
  373. async def admin_change_teacher_state(
  374. *,
  375. actor: Any,
  376. user_id: int,
  377. action: str,
  378. reason: str,
  379. ) -> dict[str, Any]:
  380. await require_teacher_admin(actor, int(user_id))
  381. if not reason.strip():
  382. raise DirectoryServiceError("reason_required", "必须填写操作原因。")
  383. identity = user_identity(actor)
  384. try:
  385. profile, _ = await set_teacher_state(
  386. user_id=int(user_id),
  387. action=action,
  388. actor_id=int(actor.id),
  389. actor_name=identity["display_name"],
  390. source="telegram_private",
  391. reason=reason,
  392. )
  393. return profile
  394. except DirectoryDataError as exc:
  395. raise DirectoryServiceError(exc.code, str(exc), status=409) from exc
  396. async def directory_for_user(
  397. *,
  398. user: Any,
  399. radius_km: int | None,
  400. page: int = 1,
  401. page_size: int = 10,
  402. ) -> tuple[list[dict[str, Any]], int]:
  403. await require_platform_membership(int(user.id))
  404. await observe_directory_user(user)
  405. location = await get_directory_location(int(user.id))
  406. if not location:
  407. raise DirectoryServiceError(
  408. "location_required",
  409. "请先使用 Telegram 定位或手动选择位置,再查看技师。",
  410. status=409,
  411. )
  412. try:
  413. return await list_directory_teachers(
  414. longitude=float(location["longitude"]),
  415. latitude=float(location["latitude"]),
  416. max_distance_meters=(
  417. float(radius_km) * 1000 if radius_km is not None else None
  418. ),
  419. page=max(1, int(page)),
  420. page_size=max(1, min(int(page_size), 20)),
  421. )
  422. except DirectoryDataError as exc:
  423. raise DirectoryServiceError(exc.code, str(exc)) from exc
  424. def public_teacher(profile: dict[str, Any]) -> dict[str, Any]:
  425. return {
  426. "user_id": int(profile["user_id"]),
  427. "username": profile.get("username"),
  428. "display_name": profile.get("display_name") or f"用户 {profile['user_id']}",
  429. "online": teacher_is_online(profile),
  430. "online_until": aware_utc(profile.get("online_until")),
  431. "distance_meters": float(profile.get("distance_meters") or 0),
  432. "location_updated_at": aware_utc(profile.get("location_updated_at")),
  433. }
  434. def profile_summary(profile: dict[str, Any], location: dict[str, Any] | None) -> dict[str, Any]:
  435. return {
  436. "user_id": int(profile["user_id"]),
  437. "username": profile.get("username"),
  438. "display_name": profile.get("display_name") or f"用户 {profile['user_id']}",
  439. "application_status": profile.get("application_status", "none"),
  440. "listed": bool(profile.get("listed")),
  441. "online": teacher_is_online(profile),
  442. "online_until": aware_utc(profile.get("online_until")),
  443. "has_location": bool(location),
  444. "location_source": location.get("source") if location else None,
  445. "location_updated_at": aware_utc(location.get("updated_at")) if location else None,
  446. }
  447. def format_distance(distance_meters: float) -> str:
  448. if distance_meters < 1000:
  449. return f"{max(0, round(distance_meters))} 米"
  450. return f"{distance_meters / 1000:.1f} 公里"
  451. async def _presence_sweeper() -> None:
  452. while True:
  453. try:
  454. await expire_teacher_presence()
  455. except Exception as exc:
  456. log.error(f"技师在线状态过期任务失败:{exc}")
  457. await asyncio.sleep(60)
  458. def start_presence_sweeper() -> None:
  459. global _presence_task
  460. if _presence_task and not _presence_task.done():
  461. return
  462. try:
  463. _presence_task = asyncio.create_task(
  464. _presence_sweeper(),
  465. name="teacher-presence-sweeper",
  466. )
  467. except RuntimeError:
  468. _presence_task = None
  469. async def stop_presence_sweeper() -> None:
  470. global _presence_task
  471. if _presence_task:
  472. _presence_task.cancel()
  473. with suppress(asyncio.CancelledError):
  474. await _presence_task
  475. _presence_task = None