dbchannel.py 4.0 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138
  1. from __future__ import annotations
  2. from datetime import UTC, datetime, timedelta
  3. from uuid import uuid4
  4. from pymongo import ASCENDING, DESCENDING, ReturnDocument
  5. from wbb import BOT_PROFILE_ID, db
  6. postsdb = db.channel_posts
  7. _indexes_ready = False
  8. def now_utc() -> datetime:
  9. return datetime.now(UTC)
  10. async def ensure_channel_indexes() -> None:
  11. global _indexes_ready
  12. if _indexes_ready:
  13. return
  14. await postsdb.create_index("post_id", unique=True)
  15. await postsdb.create_index("telegram_key", unique=True, sparse=True)
  16. await postsdb.create_index(
  17. [("bot_id", ASCENDING), ("chat_id", ASCENDING), ("created_at", DESCENDING)]
  18. )
  19. await postsdb.create_index(
  20. [("bot_id", ASCENDING), ("status", ASCENDING), ("publish_at", ASCENDING)]
  21. )
  22. _indexes_ready = True
  23. def telegram_key(chat_id: int, message_id: int) -> str:
  24. return f"{BOT_PROFILE_ID}:{chat_id}:{message_id}"
  25. async def create_post(
  26. chat_id: int,
  27. *,
  28. text: str,
  29. media_type: str | None,
  30. file_id: str | None,
  31. media_filename: str | None,
  32. media_size: int | None,
  33. media_mime_type: str | None,
  34. pin: bool,
  35. publish_at: datetime | None,
  36. ) -> dict:
  37. await ensure_channel_indexes()
  38. created = now_utc()
  39. post = {
  40. "post_id": uuid4().hex,
  41. "bot_id": BOT_PROFILE_ID,
  42. "chat_id": int(chat_id),
  43. "source": "admin",
  44. "status": "scheduled" if publish_at else "sending",
  45. "text": text,
  46. "media_type": media_type,
  47. "file_id": file_id,
  48. "media_filename": media_filename,
  49. "media_size": media_size,
  50. "media_mime_type": media_mime_type,
  51. "pin": bool(pin),
  52. "publish_at": publish_at,
  53. "message_id": None,
  54. "created_at": created,
  55. "updated_at": created,
  56. }
  57. await postsdb.insert_one(post)
  58. return post
  59. async def get_post(chat_id: int, post_id: str) -> dict | None:
  60. await ensure_channel_indexes()
  61. return await postsdb.find_one(
  62. {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id), "post_id": post_id}
  63. )
  64. async def list_posts(chat_id: int, *, page: int, page_size: int) -> tuple[list[dict], int]:
  65. await ensure_channel_indexes()
  66. query = {"bot_id": BOT_PROFILE_ID, "chat_id": int(chat_id)}
  67. total = await postsdb.count_documents(query)
  68. cursor = (
  69. postsdb.find(query)
  70. .sort("created_at", DESCENDING)
  71. .skip((page - 1) * page_size)
  72. .limit(page_size)
  73. )
  74. return [item async for item in cursor], total
  75. async def update_scheduled_post(chat_id: int, post_id: str, values: dict) -> dict | None:
  76. await ensure_channel_indexes()
  77. return await postsdb.find_one_and_update(
  78. {
  79. "bot_id": BOT_PROFILE_ID,
  80. "chat_id": int(chat_id),
  81. "post_id": post_id,
  82. "status": "scheduled",
  83. },
  84. {"$set": {**values, "updated_at": now_utc()}},
  85. return_document=ReturnDocument.AFTER,
  86. )
  87. async def claim_due_posts(limit: int = 20) -> list[dict]:
  88. await ensure_channel_indexes()
  89. deadline = now_utc()
  90. candidates = postsdb.find(
  91. {
  92. "bot_id": BOT_PROFILE_ID,
  93. "status": "scheduled",
  94. "publish_at": {"$lte": deadline},
  95. }
  96. ).sort("publish_at", ASCENDING).limit(limit)
  97. claimed: list[dict] = []
  98. async for candidate in candidates:
  99. item = await postsdb.find_one_and_update(
  100. {"post_id": candidate["post_id"], "status": "scheduled"},
  101. {"$set": {"status": "sending", "updated_at": now_utc()}},
  102. return_document=ReturnDocument.AFTER,
  103. )
  104. if item:
  105. claimed.append(item)
  106. return claimed
  107. async def mark_stale_sends_uncertain() -> None:
  108. await ensure_channel_indexes()
  109. await postsdb.update_many(
  110. {
  111. "bot_id": BOT_PROFILE_ID,
  112. "status": "sending",
  113. "updated_at": {"$lt": now_utc() - timedelta(minutes=5)},
  114. },
  115. {"$set": {"status": "uncertain", "updated_at": now_utc()}},
  116. )