channel_management.py 25 KB

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