channel_management.py 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621
  1. from __future__ import annotations
  2. import base64
  3. from datetime import UTC, datetime, timedelta
  4. from pathlib import Path
  5. from typing import Any
  6. from uuid import uuid4
  7. from aiohttp import ClientSession, ClientTimeout
  8. from pymongo.errors import DuplicateKeyError
  9. from pyrogram.enums import ChatMemberStatus, ChatType, ParseMode
  10. from pyrogram.errors import RPCError
  11. from pyrogram.types import (
  12. InputMediaAnimation,
  13. InputMediaDocument,
  14. InputMediaPhoto,
  15. InputMediaVideo,
  16. )
  17. import wbb
  18. from wbb import BOT_ID, BOT_PERMISSIONS, BOT_PROFILE_ID, app, log
  19. from wbb.services.bot_permissions import has_permission
  20. from wbb.services.chat_management import (
  21. ChatManagementError,
  22. _privilege_list,
  23. ensure_permission,
  24. serialize_member,
  25. sync_managed_chat,
  26. update_chat_profile,
  27. )
  28. from wbb.utils.dbadmin import list_managed_chats, record_audit, upsert_managed_chat
  29. from wbb.utils.dbchannel import (
  30. claim_due_posts,
  31. create_post,
  32. ensure_channel_indexes,
  33. get_post,
  34. mark_stale_sends_uncertain,
  35. now_utc,
  36. postsdb,
  37. telegram_key,
  38. update_scheduled_post,
  39. )
  40. MEDIA_METHODS = {
  41. "photo": ("send_photo", "photo", InputMediaPhoto),
  42. "animation": ("send_animation", "animation", InputMediaAnimation),
  43. "video": ("send_video", "video", InputMediaVideo),
  44. "document": ("send_document", "document", InputMediaDocument),
  45. }
  46. CHANNEL_ADMIN_PRIVILEGES = {
  47. "can_change_info",
  48. "can_post_messages",
  49. "can_edit_messages",
  50. "can_delete_messages",
  51. "can_invite_users",
  52. "can_promote_members",
  53. "can_manage_video_chats",
  54. }
  55. def telegram_error_detail(exc: Exception) -> str:
  56. return f"{type(exc).__name__}: {str(exc)[:180]}" if str(exc) else type(exc).__name__
  57. async def channel_bot_api(method: str, **params: Any) -> Any:
  58. token = str(getattr(wbb, "BOT_TOKEN", "") or "")
  59. if not token:
  60. raise ChatManagementError("bot_token_missing", "机器人令牌不可用。", status=503)
  61. data = {
  62. key: ("true" if value else "false") if isinstance(value, bool) else str(value)
  63. for key, value in params.items()
  64. }
  65. try:
  66. async with ClientSession(timeout=ClientTimeout(total=15)) as session:
  67. async with session.post(f"https://api.telegram.org/bot{token}/{method}", data=data) as response:
  68. payload = await response.json(content_type=None)
  69. except Exception as exc:
  70. raise ChatManagementError(
  71. "channel_telegram_unavailable", "Telegram 管理员接口暂时不可用。", status=502
  72. ) from exc
  73. if not isinstance(payload, dict) or not payload.get("ok"):
  74. detail = str(payload.get("description") or "请求失败")[:180] if isinstance(payload, dict) else "响应无效"
  75. raise ChatManagementError(
  76. "channel_telegram_failed", f"Telegram 拒绝管理员操作:{detail}", status=502
  77. )
  78. return payload.get("result")
  79. def serialize_bot_api_member(member: dict[str, Any]) -> dict[str, Any]:
  80. user = member.get("user") or {}
  81. return {
  82. "user": {
  83. "id": str(user.get("id") or ""),
  84. "username": user.get("username"),
  85. "first_name": user.get("first_name"),
  86. "last_name": user.get("last_name"),
  87. "is_bot": bool(user.get("is_bot")),
  88. "is_deleted": bool(user.get("is_deleted")),
  89. },
  90. "status": "owner" if member.get("status") == "creator" else member.get("status"),
  91. "custom_title": member.get("custom_title"),
  92. "privileges": [name for name in sorted(CHANNEL_ADMIN_PRIVILEGES) if member.get(name)],
  93. "until_date": member.get("until_date"),
  94. }
  95. def public_post(post: dict[str, Any]) -> dict[str, Any]:
  96. return {key: value for key, value in post.items() if key not in {"_id", "telegram_key"}}
  97. async def warm_channel_peer(chat_id: int) -> None:
  98. try:
  99. resolver = getattr(app, "resolve_peer", None)
  100. if resolver:
  101. await resolver(chat_id)
  102. else:
  103. await app.get_chat(chat_id)
  104. return
  105. except ValueError as exc:
  106. if "Peer id invalid" not in str(exc):
  107. raise
  108. token = str(getattr(wbb, "BOT_TOKEN", "") or "")
  109. if not token:
  110. raise ChatManagementError("channel_peer_unavailable", "Bot 尚未识别该频道,请等待频道新帖后重试。", status=409)
  111. try:
  112. async with ClientSession(timeout=ClientTimeout(total=10)) as session:
  113. async with session.post(
  114. f"https://api.telegram.org/bot{token}/getChat", data={"chat_id": str(chat_id)}
  115. ) as response:
  116. data = await response.json()
  117. chat = data.get("result") if data.get("ok") else None
  118. if not chat or int(chat.get("id", 0)) != chat_id or chat.get("type") != "channel":
  119. raise ChatManagementError("channel_unavailable", "Telegram 未能确认该频道及 Bot 的访问权限。", status=404)
  120. username = chat.get("username")
  121. if not username:
  122. raise ChatManagementError(
  123. "channel_peer_unavailable", "Bot 尚未缓存该私有频道;请在频道发布新帖后重试。", status=409
  124. )
  125. await app.get_chat(f"@{username}")
  126. except ChatManagementError:
  127. raise
  128. except Exception as exc:
  129. raise ChatManagementError("channel_peer_unavailable", "Bot 暂时无法解析频道 ID,请稍后重试。", status=502) from exc
  130. async def ensure_channel(
  131. chat_id: int, *, admin: bool = True, include_photo: bool = False
  132. ) -> dict[str, Any]:
  133. await warm_channel_peer(chat_id)
  134. overview = await sync_managed_chat(chat_id)
  135. if overview["type"] != "channel":
  136. raise ChatManagementError("not_channel", "目标不是频道。", status=400)
  137. if admin and overview["bot_status"] not in {"owner", "administrator"}:
  138. raise ChatManagementError("bot_not_admin", "机器人不是该频道管理员。", status=403)
  139. chat = await app.get_chat(chat_id)
  140. photo = getattr(chat, "photo", None)
  141. overview["photo_file_id"] = getattr(photo, "small_file_id", None)
  142. if include_photo and overview["photo_file_id"]:
  143. try:
  144. content = await app.download_media(overview["photo_file_id"], in_memory=True)
  145. overview["photo_data_url"] = "data:image/jpeg;base64," + base64.b64encode(
  146. content.getvalue()
  147. ).decode("ascii")
  148. except Exception:
  149. overview["photo_data_url"] = None
  150. return overview
  151. async def list_channels(*, query: str, page: int, page_size: int) -> tuple[list[dict], int]:
  152. return await list_managed_chats(query=query, page=page, page_size=page_size, chat_type="channel")
  153. async def update_channel_profile(chat_id: int, *, title: str | None, description: str | None) -> dict:
  154. await ensure_channel(chat_id)
  155. try:
  156. return await update_chat_profile(chat_id, title=title, description=description)
  157. except ChatManagementError as exc:
  158. raise ChatManagementError(exc.code, str(exc).replace("群", "频道"), status=exc.status) from exc
  159. async def update_channel_photo(chat_id: int, photo_path: Path) -> dict:
  160. await ensure_channel(chat_id)
  161. await ensure_permission(chat_id, "can_change_info")
  162. try:
  163. await app.set_chat_photo(chat_id, photo=str(photo_path))
  164. except RPCError as exc:
  165. raise ChatManagementError(
  166. "channel_photo_failed", f"Telegram 拒绝更新频道头像:{telegram_error_detail(exc)}", status=502
  167. ) from exc
  168. except Exception as exc:
  169. raise ChatManagementError("channel_photo_failed", "Telegram 未能更新频道头像。", status=502) from exc
  170. return await ensure_channel(chat_id, include_photo=True)
  171. async def list_channel_admins(chat_id: int) -> list[dict]:
  172. await ensure_channel(chat_id)
  173. from pyrogram.enums import ChatMembersFilter
  174. try:
  175. return [
  176. serialize_member(member)
  177. async for member in app.get_chat_members(chat_id, filter=ChatMembersFilter.ADMINISTRATORS)
  178. ]
  179. except Exception as exc:
  180. raise ChatManagementError("admin_list_failed", "Telegram 未能读取频道管理员。", status=502) from exc
  181. async def list_channel_admin_candidates(chat_id: int, *, query: str = "", limit: int = 30) -> list[dict]:
  182. await ensure_channel(chat_id)
  183. await ensure_permission(chat_id, "can_promote_members")
  184. from pyrogram.enums import ChatMembersFilter
  185. query = query.strip()
  186. if len(query) > 64:
  187. raise ChatManagementError("member_query_too_long", "搜索内容不能超过 64 个字符。")
  188. limit = max(1, min(limit, 30))
  189. if query.isdigit():
  190. member = await channel_bot_api("getChatMember", chat_id=chat_id, user_id=int(query))
  191. if member.get("status") != "member" or (member.get("user") or {}).get("is_bot"):
  192. return []
  193. return [serialize_bot_api_member(member)]
  194. results: list[dict] = []
  195. try:
  196. members = app.get_chat_members(
  197. chat_id,
  198. query=query.lstrip("@"),
  199. limit=limit,
  200. filter=ChatMembersFilter.SEARCH if query else ChatMembersFilter.RECENT,
  201. )
  202. async for member in members:
  203. if member.status != ChatMemberStatus.MEMBER or member.user.is_bot:
  204. continue
  205. results.append(serialize_member(member))
  206. except RPCError as exc:
  207. raise ChatManagementError(
  208. "member_list_failed", f"Telegram 拒绝读取频道订阅者:{telegram_error_detail(exc)}", status=502
  209. ) from exc
  210. except Exception as exc:
  211. raise ChatManagementError("member_list_failed", "Telegram 未能读取频道订阅者。", status=502) from exc
  212. return results
  213. async def set_channel_admin(chat_id: int, user_id: int, privileges: list[str]) -> dict:
  214. await ensure_channel(chat_id)
  215. await ensure_permission(chat_id, "can_promote_members")
  216. requested = set(privileges)
  217. if not requested or requested - CHANNEL_ADMIN_PRIVILEGES:
  218. raise ChatManagementError("invalid_privileges", "请选择有效的频道管理员权限。")
  219. try:
  220. bot = await app.get_chat_member(chat_id, BOT_ID)
  221. except Exception as exc:
  222. raise ChatManagementError("bot_permission_unavailable", "Telegram 未能读取机器人权限。", status=502) from exc
  223. allowed = set(_privilege_list(bot))
  224. if bot.status != ChatMemberStatus.OWNER and requested - allowed:
  225. raise ChatManagementError("privilege_exceeds_bot", "不能授予超出机器人自身的权限。", status=403)
  226. target = await channel_bot_api("getChatMember", chat_id=chat_id, user_id=user_id)
  227. if target.get("status") in {"creator", "owner"}:
  228. raise ChatManagementError("owner_protected", "不能修改频道所有者。", status=403)
  229. values = {key: key in requested for key in CHANNEL_ADMIN_PRIVILEGES}
  230. await channel_bot_api("promoteChatMember", chat_id=chat_id, user_id=user_id, **values)
  231. return serialize_bot_api_member(
  232. await channel_bot_api("getChatMember", chat_id=chat_id, user_id=user_id)
  233. )
  234. async def remove_channel_admin(chat_id: int, user_id: int) -> dict:
  235. await ensure_channel(chat_id)
  236. await ensure_permission(chat_id, "can_promote_members")
  237. target = await channel_bot_api("getChatMember", chat_id=chat_id, user_id=user_id)
  238. if target.get("status") in {"creator", "owner"}:
  239. raise ChatManagementError("owner_protected", "不能移除频道所有者。", status=403)
  240. if target.get("status") != "administrator":
  241. raise ChatManagementError("not_admin", "该用户不是频道管理员。", status=409)
  242. await channel_bot_api(
  243. "promoteChatMember", chat_id=chat_id, user_id=user_id,
  244. **{key: False for key in CHANNEL_ADMIN_PRIVILEGES},
  245. )
  246. return {"user_id": str(user_id), "removed": True}
  247. def validate_post_values(
  248. values: dict[str, Any], *, scheduled: bool = False, keep_publish_at: bool = False
  249. ) -> dict[str, Any]:
  250. text = str(values.get("text") or "").strip()
  251. media_type = values.get("media_type") or None
  252. file_id = values.get("file_id") or None
  253. if media_type not in (None, *MEDIA_METHODS) or bool(media_type) != bool(file_id):
  254. raise ChatManagementError("invalid_media", "媒体类型与文件 ID 不匹配。")
  255. if not text and not file_id:
  256. raise ChatManagementError("empty_post", "帖子文字或媒体至少填写一项。")
  257. if len(text) > (1024 if file_id else 4096):
  258. raise ChatManagementError("post_too_long", "文字超过 Telegram 长度限制。")
  259. media_filename = str(values.get("media_filename") or "").strip() if file_id else ""
  260. media_mime_type = str(values.get("media_mime_type") or "").strip() if file_id else ""
  261. if len(media_filename) > 255 or len(media_mime_type) > 100:
  262. raise ChatManagementError("invalid_media_metadata", "媒体文件信息过长。")
  263. try:
  264. media_size = int(values["media_size"]) if file_id and values.get("media_size") is not None else None
  265. except (TypeError, ValueError) as exc:
  266. raise ChatManagementError("invalid_media_metadata", "媒体文件大小无效。") from exc
  267. if media_size is not None and media_size < 0:
  268. raise ChatManagementError("invalid_media_metadata", "媒体文件大小无效。")
  269. publish_at = values.get("publish_at")
  270. if publish_at is not None:
  271. if isinstance(publish_at, str):
  272. try:
  273. publish_at = datetime.fromisoformat(publish_at.replace("Z", "+00:00"))
  274. except ValueError as exc:
  275. raise ChatManagementError("invalid_publish_at", "定时时间无效。") from exc
  276. if publish_at.tzinfo is None:
  277. raise ChatManagementError("invalid_publish_at", "定时时间须包含时区。")
  278. if not isinstance(publish_at, datetime):
  279. raise ChatManagementError("invalid_publish_at", "定时时间须包含时区。")
  280. if publish_at.tzinfo is None:
  281. publish_at = publish_at.replace(tzinfo=UTC)
  282. publish_at = publish_at.astimezone(UTC)
  283. if not keep_publish_at and publish_at <= now_utc() + timedelta(seconds=5):
  284. raise ChatManagementError("invalid_publish_at", "定时时间至少晚于当前 5 秒。")
  285. elif scheduled:
  286. raise ChatManagementError("invalid_publish_at", "请填写定时时间。")
  287. return {
  288. "text": text,
  289. "media_type": media_type,
  290. "file_id": file_id,
  291. "media_filename": media_filename or None,
  292. "media_size": media_size,
  293. "media_mime_type": media_mime_type or None,
  294. "pin": bool(values.get("pin", False)),
  295. "publish_at": publish_at,
  296. }
  297. def same_publish_at(value: Any, current: Any) -> bool:
  298. if isinstance(value, str):
  299. try:
  300. value = datetime.fromisoformat(value.replace("Z", "+00:00"))
  301. except ValueError:
  302. return False
  303. if not isinstance(value, datetime) or not isinstance(current, datetime):
  304. return False
  305. value = value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC)
  306. current = current.replace(tzinfo=UTC) if current.tzinfo is None else current.astimezone(UTC)
  307. return value == current
  308. async def publish_post(post: dict[str, Any]) -> dict:
  309. if not has_permission("channel.posts", BOT_PERMISSIONS):
  310. raise ChatManagementError("bot_role_permission_denied", "所选机器人的职责角色不允许发布频道帖子。", status=403)
  311. chat_id = int(post["chat_id"])
  312. await ensure_channel(chat_id)
  313. await ensure_permission(chat_id, "can_post_messages")
  314. if post.get("pin"):
  315. await ensure_permission(chat_id, "can_edit_messages")
  316. try:
  317. if post.get("file_id"):
  318. method_name, keyword, _ = MEDIA_METHODS[post["media_type"]]
  319. sent = await getattr(app, method_name)(
  320. chat_id, **{keyword: post["file_id"]}, caption=post["text"] or None,
  321. parse_mode=ParseMode.DISABLED,
  322. )
  323. else:
  324. sent = await app.send_message(chat_id, post["text"], parse_mode=ParseMode.DISABLED)
  325. except RPCError as exc:
  326. detail = telegram_error_detail(exc)
  327. await postsdb.update_one(
  328. {"post_id": post["post_id"], "status": "sending"},
  329. {"$set": {"status": "failed", "last_error": detail, "updated_at": now_utc()}},
  330. )
  331. raise ChatManagementError("post_send_failed", f"Telegram 拒绝发布:{detail}", status=502) from exc
  332. except Exception as exc:
  333. await postsdb.update_one(
  334. {"post_id": post["post_id"], "status": "sending"},
  335. {"$set": {"status": "uncertain", "last_error": "Telegram 发送返回异常,结果需要人工核实。", "updated_at": now_utc()}},
  336. )
  337. raise ChatManagementError(
  338. "post_delivery_uncertain", "Telegram 发送结果不明,已停止自动重试;请在频道核实。", status=502
  339. ) from exc
  340. key = telegram_key(chat_id, int(sent.id))
  341. try:
  342. for attempt in range(3):
  343. await postsdb.delete_one({"telegram_key": key, "source": "telegram"})
  344. try:
  345. await postsdb.update_one(
  346. {"post_id": post["post_id"]},
  347. {"$set": {
  348. "telegram_key": key,
  349. "message_id": int(sent.id),
  350. "status": "published",
  351. "published_at": now_utc(),
  352. "updated_at": now_utc(),
  353. }},
  354. )
  355. break
  356. except DuplicateKeyError:
  357. if attempt == 2:
  358. raise
  359. except Exception as exc:
  360. raise ChatManagementError(
  361. "post_delivery_uncertain", "消息已发出,但记录更新失败;请在频道核实,系统不会自动重发。", status=502
  362. ) from exc
  363. if post.get("pin"):
  364. try:
  365. await app.pin_chat_message(chat_id, int(sent.id))
  366. except Exception:
  367. await postsdb.update_one(
  368. {"post_id": post["post_id"]},
  369. {"$set": {"pin_error": "Telegram 未能置顶该消息。"}},
  370. )
  371. return public_post(await get_post(chat_id, post["post_id"]))
  372. async def create_channel_post(chat_id: int, values: dict[str, Any]) -> dict:
  373. await ensure_channel(chat_id)
  374. await ensure_permission(chat_id, "can_post_messages")
  375. data = validate_post_values(values)
  376. if data["pin"]:
  377. await ensure_permission(chat_id, "can_edit_messages")
  378. post = await create_post(chat_id, **data)
  379. if data["publish_at"]:
  380. return public_post(post)
  381. try:
  382. return await publish_post(post)
  383. except ChatManagementError as exc:
  384. if exc.code != "post_delivery_uncertain":
  385. await postsdb.update_one(
  386. {"post_id": post["post_id"], "status": "sending"},
  387. {"$set": {"status": "failed", "last_error": str(exc)[:300], "updated_at": now_utc()}},
  388. )
  389. raise
  390. async def edit_channel_post(chat_id: int, post_id: str, values: dict[str, Any]) -> dict:
  391. await ensure_channel(chat_id)
  392. current = await get_post(chat_id, post_id)
  393. if not current:
  394. raise ChatManagementError("post_not_found", "帖子不存在。", status=404)
  395. merged = {key: current.get(key) for key in (
  396. "text", "media_type", "file_id", "media_filename", "media_size", "media_mime_type", "pin", "publish_at"
  397. )}
  398. if current["status"] != "scheduled":
  399. merged["publish_at"] = None
  400. merged.update(values)
  401. if values.get("file_id") and values["file_id"] != current.get("file_id"):
  402. for key in ("media_filename", "media_size", "media_mime_type"):
  403. if key not in values:
  404. merged[key] = None
  405. data = validate_post_values(
  406. merged,
  407. scheduled=current["status"] == "scheduled",
  408. keep_publish_at=current["status"] == "scheduled" and (
  409. "publish_at" not in values or same_publish_at(values["publish_at"], current.get("publish_at"))
  410. ),
  411. )
  412. if current["status"] == "scheduled":
  413. await ensure_permission(chat_id, "can_post_messages")
  414. updated = await update_scheduled_post(chat_id, post_id, data)
  415. if not updated:
  416. raise ChatManagementError("post_in_progress", "帖子已进入发送流程,不能编辑。", status=409)
  417. return public_post(updated)
  418. if current["status"] != "published":
  419. raise ChatManagementError("post_not_editable", "当前状态不能编辑帖子。", status=409)
  420. if "publish_at" in values or "pin" in values:
  421. raise ChatManagementError("invalid_edit", "已发布帖子不能修改定时或置顶设置。")
  422. await ensure_permission(chat_id, "can_edit_messages")
  423. message_id = int(current["message_id"])
  424. try:
  425. if data["media_type"]:
  426. if not current.get("media_type"):
  427. raise ChatManagementError("media_change_unsupported", "文字帖不能改为媒体帖。")
  428. if data["file_id"] != current.get("file_id") or data["media_type"] != current.get("media_type"):
  429. _, _, media_class = MEDIA_METHODS[data["media_type"]]
  430. await app.edit_message_media(
  431. chat_id, message_id, media_class(
  432. data["file_id"], caption=data["text"], parse_mode=ParseMode.DISABLED
  433. )
  434. )
  435. else:
  436. await app.edit_message_caption(
  437. chat_id, message_id, data["text"], parse_mode=ParseMode.DISABLED
  438. )
  439. else:
  440. if current.get("media_type"):
  441. raise ChatManagementError("media_change_unsupported", "媒体帖不能改为纯文字帖。")
  442. await app.edit_message_text(
  443. chat_id, message_id, data["text"], parse_mode=ParseMode.DISABLED
  444. )
  445. except ChatManagementError:
  446. raise
  447. except RPCError as exc:
  448. raise ChatManagementError(
  449. "post_edit_failed", f"Telegram 拒绝编辑帖子:{telegram_error_detail(exc)}", status=502
  450. ) from exc
  451. except Exception as exc:
  452. raise ChatManagementError("post_edit_failed", "Telegram 未能编辑帖子;请检查 Bot 权限与消息状态。", status=502) from exc
  453. await postsdb.update_one(
  454. {"post_id": post_id, "bot_id": BOT_PROFILE_ID},
  455. {"$set": {key: data[key] for key in (
  456. "text", "media_type", "file_id", "media_filename", "media_size", "media_mime_type"
  457. )} | {"updated_at": now_utc()}},
  458. )
  459. return public_post(await get_post(chat_id, post_id))
  460. async def delete_channel_post(chat_id: int, post_id: str) -> dict:
  461. await ensure_channel(chat_id)
  462. current = await get_post(chat_id, post_id)
  463. if not current:
  464. raise ChatManagementError("post_not_found", "帖子不存在。", status=404)
  465. if current["status"] == "scheduled":
  466. updated = await update_scheduled_post(chat_id, post_id, {"status": "canceled"})
  467. if not updated:
  468. raise ChatManagementError("post_in_progress", "帖子已进入发送流程,不能取消。", status=409)
  469. return public_post(updated)
  470. if current["status"] != "published":
  471. raise ChatManagementError("post_not_deletable", "当前状态不能删除帖子。", status=409)
  472. await ensure_permission(chat_id, "can_delete_messages")
  473. try:
  474. deleted = await app.delete_messages(chat_id, int(current["message_id"]))
  475. if deleted == 0:
  476. raise RuntimeError("Telegram did not delete the message")
  477. except RPCError as exc:
  478. raise ChatManagementError(
  479. "post_delete_failed", f"Telegram 拒绝删除帖子:{telegram_error_detail(exc)}", status=502
  480. ) from exc
  481. except Exception as exc:
  482. raise ChatManagementError("post_delete_failed", "Telegram 未能删除帖子;请检查 Bot 权限与消息状态。", status=502) from exc
  483. await postsdb.update_one(
  484. {"post_id": post_id, "bot_id": BOT_PROFILE_ID},
  485. {"$set": {"status": "deleted", "updated_at": now_utc()}},
  486. )
  487. return public_post(await get_post(chat_id, post_id))
  488. async def observe_channel_post(message: Any) -> None:
  489. if getattr(message.chat, "type", None) != ChatType.CHANNEL:
  490. return
  491. chat_id, message_id = int(message.chat.id), int(message.id)
  492. bot = await app.get_chat_member(chat_id, BOT_ID)
  493. if bot.status not in {ChatMemberStatus.OWNER, ChatMemberStatus.ADMINISTRATOR}:
  494. return
  495. await upsert_managed_chat(
  496. chat_id=chat_id,
  497. title=message.chat.title,
  498. username=message.chat.username,
  499. chat_type="channel",
  500. bot_status=str(getattr(bot.status, "value", bot.status)),
  501. bot_privileges=_privilege_list(bot),
  502. )
  503. media_type = next((name for name in MEDIA_METHODS if getattr(message, name, None)), None)
  504. media = getattr(message, media_type, None) if media_type else None
  505. key = telegram_key(chat_id, message_id)
  506. values = {
  507. "text": getattr(message, "text", None) or getattr(message, "caption", None) or "",
  508. "media_type": media_type or ("other" if getattr(message, "media", None) else None),
  509. "file_id": getattr(media, "file_id", None),
  510. "media_filename": getattr(media, "file_name", None),
  511. "media_size": getattr(media, "file_size", None),
  512. "media_mime_type": getattr(media, "mime_type", None),
  513. "updated_at": now_utc(),
  514. }
  515. await ensure_channel_indexes()
  516. update = {
  517. "$set": values,
  518. "$setOnInsert": {
  519. "post_id": uuid4().hex,
  520. "bot_id": BOT_PROFILE_ID,
  521. "chat_id": chat_id,
  522. "telegram_key": key,
  523. "message_id": message_id,
  524. "source": "telegram",
  525. "status": "published",
  526. "pin": False,
  527. "created_at": getattr(message, "date", None) or now_utc(),
  528. "published_at": getattr(message, "date", None) or now_utc(),
  529. },
  530. }
  531. try:
  532. await postsdb.update_one({"telegram_key": key}, update, upsert=True)
  533. except DuplicateKeyError:
  534. await postsdb.update_one({"telegram_key": key}, {"$set": values})
  535. async def observe_channel_deletion(chat_id: int, message_id: int) -> None:
  536. await ensure_channel_indexes()
  537. await postsdb.update_one(
  538. {"telegram_key": telegram_key(chat_id, message_id)},
  539. {"$set": {"status": "deleted", "updated_at": now_utc()}},
  540. )
  541. async def sweep_channel_posts() -> None:
  542. await mark_stale_sends_uncertain()
  543. for post in await claim_due_posts():
  544. try:
  545. await publish_post(post)
  546. except Exception as exc:
  547. if isinstance(exc, ChatManagementError) and exc.code != "post_delivery_uncertain":
  548. await postsdb.update_one(
  549. {"post_id": post["post_id"], "status": "sending"},
  550. {"$set": {"status": "failed", "last_error": str(exc)[:300], "updated_at": now_utc()}},
  551. )
  552. try:
  553. await record_audit(
  554. source="system", actor_id=None, actor_name="频道定时发布",
  555. action="channel.post.publish", chat_id=int(post["chat_id"]),
  556. target_id=post["post_id"], summary="定时发布失败或结果不明",
  557. success=False, error=str(exc)[:300],
  558. )
  559. except Exception as audit_exc:
  560. log.error(f"频道定时发布审计写入失败:{audit_exc}")
  561. else:
  562. try:
  563. await record_audit(
  564. source="system", actor_id=None, actor_name="频道定时发布",
  565. action="channel.post.publish", chat_id=int(post["chat_id"]),
  566. target_id=post["post_id"], summary="定时帖子已发布",
  567. )
  568. except Exception as audit_exc:
  569. log.error(f"频道定时发布审计写入失败:{audit_exc}")