giveaways.py 43 KB

12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106
  1. from __future__ import annotations
  2. import asyncio
  3. import secrets
  4. from collections import defaultdict
  5. from contextlib import suppress
  6. from datetime import UTC, datetime
  7. from html import escape
  8. from typing import Any
  9. from zoneinfo import ZoneInfo
  10. from pymongo.errors import DuplicateKeyError
  11. from pyrogram.enums import ParseMode
  12. from pyrogram.types import InlineKeyboardButton, InlineKeyboardMarkup
  13. from wbb import BOT_PROFILE_ID, app, log
  14. from wbb.services.giveaway_eligibility import EligibilityUnavailable, check_eligibility
  15. from wbb.utils.dbadmin import record_audit
  16. from wbb.utils.dbgiveaway import (
  17. STATUS_CANCELED,
  18. STATUS_CANCELING,
  19. STATUS_DRAWING,
  20. STATUS_FINISHED,
  21. STATUS_RUNNING,
  22. add_participant,
  23. as_utc,
  24. attach_giveaway_message,
  25. claim_giveaway_cancel,
  26. claim_giveaway_draw,
  27. count_participants,
  28. create_giveaway,
  29. finalize_giveaway_cancel,
  30. finish_giveaway,
  31. get_giveaway,
  32. get_participant,
  33. giveawaysdb,
  34. is_giveaway_banned,
  35. list_participants,
  36. mark_participant_refunded,
  37. participantsdb,
  38. record_reroll,
  39. remove_participant,
  40. save_pending_winners,
  41. ticket_ordersdb,
  42. update_running_giveaway,
  43. utc_now,
  44. )
  45. from wbb.utils.dbpoints import (
  46. SOURCE_GIVEAWAY_ENTRY,
  47. SOURCE_GIVEAWAY_PARTICIPATION,
  48. SOURCE_GIVEAWAY_REFUND,
  49. SOURCE_GIVEAWAY_WINNER,
  50. InsufficientPoints,
  51. adjust_points,
  52. get_point_account,
  53. get_point_transaction_by_key,
  54. )
  55. from wbb.utils.i18n import giveaway_status_label
  56. _RANDOM = secrets.SystemRandom()
  57. _giveaway_locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock)
  58. GIVEAWAY_TIMEZONE = ZoneInfo("Asia/Shanghai")
  59. class GiveawayServiceError(RuntimeError):
  60. def __init__(self, code: str, message: str, *, status: int = 400):
  61. super().__init__(message)
  62. self.code = code
  63. self.status = status
  64. def _mongo_datetime_millis(value: datetime) -> int:
  65. return int(as_utc(value).timestamp() * 1000)
  66. def parse_giveaway_time(raw: str) -> datetime | None:
  67. try:
  68. local_time = datetime.strptime(raw.strip(), "%Y-%m-%d %H:%M")
  69. except ValueError:
  70. return None
  71. return local_time.replace(tzinfo=GIVEAWAY_TIMEZONE).astimezone(UTC)
  72. def format_giveaway_time(value: datetime) -> str:
  73. return as_utc(value).astimezone(GIVEAWAY_TIMEZONE).strftime("%Y-%m-%d %H:%M")
  74. def giveaway_starts_at(giveaway: dict[str, Any]) -> datetime:
  75. return as_utc(giveaway.get("starts_at") or giveaway["created_at"])
  76. def _participant_mention(participant: dict[str, Any]) -> str:
  77. user_id = int(participant["user_id"])
  78. label = (
  79. participant.get("display_name")
  80. or ("@" + str(participant["username"]) if participant.get("username") else None)
  81. or participant.get("first_name")
  82. or str(user_id)
  83. )
  84. return f'<a href="tg://user?id={user_id}">{escape(str(label))}</a>'
  85. def _format_prizes(prizes: list[dict[str, Any]]) -> str:
  86. lines: list[str] = []
  87. for prize in prizes:
  88. points = int(prize.get("points_reward", 0))
  89. suffix = f"(每人 +{points} 积分)" if points else ""
  90. lines.append(
  91. f"- {escape(str(prize['name']))}:{int(prize['count'])} 人{suffix}"
  92. )
  93. return "\n".join(lines)
  94. async def render_giveaway(giveaway: dict[str, Any]) -> str:
  95. participants = await count_participants(giveaway["giveaway_id"])
  96. status_label = giveaway_status_label(giveaway["status"])
  97. if giveaway["status"] == STATUS_RUNNING and giveaway_starts_at(giveaway) > utc_now():
  98. status_label = "待报名"
  99. text = [
  100. f"<b>抽奖 #{escape(giveaway['giveaway_id'])}</b>",
  101. f"<b>状态:</b>{escape(status_label)}",
  102. f"<b>标题:</b>{escape(giveaway['title'])}",
  103. ]
  104. if giveaway.get("description"):
  105. text.append(f"<b>说明:</b>{escape(giveaway['description'])}")
  106. text.extend(
  107. [
  108. "<b>奖项:</b>",
  109. _format_prizes(giveaway["prizes"]),
  110. f"<b>报名开始:</b>{format_giveaway_time(giveaway_starts_at(giveaway))}",
  111. f"<b>开奖时间:</b>{format_giveaway_time(giveaway['ends_at'])}",
  112. "<b>时区:</b>北京时间",
  113. f"<b>参与人数:</b>{participants}",
  114. ]
  115. )
  116. minimum = int(giveaway.get("minimum_points", 0))
  117. cost = int(giveaway.get("entry_cost", 0))
  118. reward = int(giveaway.get("participation_reward", 0))
  119. if minimum:
  120. text.append(f"<b>最低积分:</b>{minimum}")
  121. if cost:
  122. text.append(f"<b>每张奖票:</b>{cost} 积分;每人最多 {int(giveaway.get('max_tickets_per_user', 10))} 张")
  123. else:
  124. text.append("<b>免费抽奖:</b>每人限 1 张奖票")
  125. targets = giveaway.get("eligibility_targets") or []
  126. if targets:
  127. condition = "全部满足" if giveaway.get("eligibility_mode", "all") == "all" else "满足任意一个"
  128. labels = [
  129. "@" + str(item["username"]) if item.get("username") else str(item.get("title") or item["chat_id"])
  130. for item in targets
  131. ]
  132. text.append(f"<b>报名资格({condition}):</b>{escape('、'.join(labels))}")
  133. if reward:
  134. text.append(f"<b>参与奖励:</b>+{reward}")
  135. return "\n".join(text)
  136. def join_markup(giveaway_id: str, entry_cost: int = 0) -> InlineKeyboardMarkup:
  137. rows = [
  138. [
  139. InlineKeyboardButton(
  140. "立即报名" if not entry_cost else f"购买 1 张({entry_cost} 积分)",
  141. callback_data=f"giveaway_join:{giveaway_id}"
  142. )
  143. ],
  144. ]
  145. if entry_cost:
  146. rows.append([
  147. InlineKeyboardButton("购买 5 张", callback_data=f"giveaway_buy:{giveaway_id}:5"),
  148. ])
  149. rows.extend([
  150. [
  151. InlineKeyboardButton(
  152. "查看参与者",
  153. callback_data=f"giveaway_participants:{giveaway_id}",
  154. )
  155. ],
  156. ])
  157. return InlineKeyboardMarkup(rows)
  158. async def refresh_giveaway_message(giveaway: dict[str, Any]) -> bool:
  159. message_id = giveaway.get("message_id")
  160. if not message_id:
  161. return True
  162. try:
  163. await app.edit_message_text(
  164. int(giveaway["chat_id"]),
  165. int(message_id),
  166. await render_giveaway(giveaway),
  167. parse_mode=ParseMode.HTML,
  168. reply_markup=join_markup(giveaway["giveaway_id"], int(giveaway.get("entry_cost", 0))),
  169. )
  170. except Exception:
  171. return False
  172. return True
  173. async def update_and_refresh_giveaway(
  174. giveaway_id: str,
  175. *,
  176. title: str,
  177. description: str,
  178. starts_at: datetime,
  179. ends_at: datetime,
  180. expected_updated_at: datetime,
  181. ) -> dict[str, Any]:
  182. async with _giveaway_locks[giveaway_id]:
  183. giveaway = await get_giveaway(giveaway_id)
  184. if not giveaway:
  185. raise GiveawayServiceError(
  186. "giveaway_not_found", "未找到该抽奖。", status=404
  187. )
  188. if giveaway["status"] != STATUS_RUNNING:
  189. raise GiveawayServiceError(
  190. "giveaway_not_running", "只有进行中的抽奖可以编辑。", status=409
  191. )
  192. if _mongo_datetime_millis(giveaway["updated_at"]) != _mongo_datetime_millis(
  193. expected_updated_at
  194. ):
  195. raise GiveawayServiceError(
  196. "giveaway_stale",
  197. "抽奖已被其他操作更新,请刷新页面后重试。",
  198. status=409,
  199. )
  200. current_starts_at = giveaway_starts_at(giveaway)
  201. if current_starts_at <= utc_now() and _mongo_datetime_millis(
  202. starts_at
  203. ) != _mongo_datetime_millis(current_starts_at):
  204. raise GiveawayServiceError(
  205. "giveaway_registration_started",
  206. "报名开始后不能修改报名开始时间。",
  207. status=409,
  208. )
  209. try:
  210. updated = await update_running_giveaway(
  211. giveaway_id,
  212. title=title,
  213. description=description,
  214. starts_at=starts_at,
  215. ends_at=ends_at,
  216. expected_updated_at=giveaway["updated_at"],
  217. )
  218. except ValueError as exc:
  219. raise GiveawayServiceError("invalid_giveaway", str(exc)) from exc
  220. if not updated:
  221. raise GiveawayServiceError(
  222. "giveaway_stale",
  223. "抽奖状态或内容已变化,请刷新页面后重试。",
  224. status=409,
  225. )
  226. if not await refresh_giveaway_message(updated):
  227. raise GiveawayServiceError(
  228. "giveaway_message_refresh_failed",
  229. "抽奖内容已保存,但群消息刷新失败,请刷新页面后重试保存。",
  230. status=502,
  231. )
  232. return updated
  233. def pick_winners(
  234. participants: list[dict[str, Any]], prizes: list[dict[str, Any]]
  235. ) -> list[dict[str, Any]]:
  236. winners: list[dict[str, Any]] = []
  237. picked_user_ids: set[int] = set()
  238. for tier_index, prize in enumerate(prizes):
  239. remaining = [
  240. participant
  241. for participant in participants
  242. if int(participant["user_id"]) not in picked_user_ids
  243. ]
  244. picked = []
  245. for _ in range(min(int(prize["count"]), len(remaining))):
  246. total = sum(max(1, int(item.get("paid_ticket_count", item.get("ticket_count", 1)))) for item in remaining)
  247. position = _RANDOM.randrange(total)
  248. for index, item in enumerate(remaining):
  249. position -= max(1, int(item.get("paid_ticket_count", item.get("ticket_count", 1))))
  250. if position < 0:
  251. picked.append(remaining.pop(index))
  252. break
  253. for participant in picked:
  254. user_id = int(participant["user_id"])
  255. picked_user_ids.add(user_id)
  256. winners.append(
  257. {
  258. "user_id": user_id,
  259. "username": participant.get("username"),
  260. "first_name": participant.get("first_name"),
  261. "display_name": participant.get("display_name"),
  262. "tier": str(prize["name"]),
  263. "tier_index": tier_index,
  264. "points_reward": int(prize.get("points_reward", 0)),
  265. }
  266. )
  267. return winners
  268. def render_winners(
  269. giveaway: dict[str, Any],
  270. winners: list[dict[str, Any]],
  271. *,
  272. reroll: bool = False,
  273. tier_name: str | None = None,
  274. ) -> str:
  275. title = "抽奖重抽结果" if reroll else "抽奖结果"
  276. lines = [
  277. f"<b>{title} #{escape(giveaway['giveaway_id'])}</b>",
  278. f"<b>标题:</b>{escape(giveaway['title'])}",
  279. ]
  280. if tier_name:
  281. lines.append(f"<b>奖项:</b>{escape(tier_name)}")
  282. if not winners:
  283. lines.append("没有符合条件的参与者。")
  284. return "\n".join(lines)
  285. lines.append("<b>中奖者:</b>")
  286. for winner in winners:
  287. points = int(winner.get("points_reward", 0))
  288. suffix = f"(+{points} 积分)" if points else ""
  289. lines.append(
  290. f"- {escape(str(winner['tier']))}:{_participant_mention(winner)}{suffix}"
  291. )
  292. return "\n".join(lines)
  293. async def create_and_publish_giveaway(**values: Any) -> dict[str, Any]:
  294. giveaway = await create_giveaway(**values)
  295. try:
  296. sent = await app.send_message(
  297. giveaway["chat_id"],
  298. await render_giveaway(giveaway),
  299. parse_mode=ParseMode.HTML,
  300. reply_markup=join_markup(giveaway["giveaway_id"], int(giveaway.get("entry_cost", 0))),
  301. disable_web_page_preview=True,
  302. )
  303. except Exception:
  304. await claim_giveaway_cancel(giveaway["giveaway_id"], giveaway["chat_id"])
  305. await finalize_giveaway_cancel(giveaway["giveaway_id"])
  306. raise
  307. await attach_giveaway_message(giveaway["giveaway_id"], sent.chat.id, sent.id)
  308. giveaway["message_id"] = sent.id
  309. return giveaway
  310. async def join_giveaway(
  311. giveaway_id: str,
  312. *,
  313. user_id: int,
  314. username: str | None,
  315. first_name: str | None,
  316. display_name: str | None = None,
  317. verify_membership: bool = True,
  318. ) -> tuple[str, dict[str, Any] | None]:
  319. async with _giveaway_locks[giveaway_id]:
  320. giveaway = await get_giveaway(giveaway_id)
  321. if not giveaway:
  322. return "missing", None
  323. if giveaway["status"] != STATUS_RUNNING:
  324. return "closed", giveaway
  325. now = utc_now()
  326. if giveaway_starts_at(giveaway) > now:
  327. return "not_started", giveaway
  328. if as_utc(giveaway["ends_at"]) <= now:
  329. return "ended", giveaway
  330. existing = await get_participant(giveaway_id, user_id)
  331. if existing:
  332. return "duplicate", giveaway
  333. if verify_membership:
  334. try:
  335. eligible, failed = await check_eligibility(giveaway, user_id)
  336. except EligibilityUnavailable:
  337. return "eligibility_unavailable", giveaway
  338. if not eligible:
  339. return (
  340. "not_member" if int(giveaway["chat_id"]) in failed else "not_eligible",
  341. giveaway,
  342. )
  343. minimum = int(giveaway.get("minimum_points", 0))
  344. entry_cost = int(giveaway.get("entry_cost", 0))
  345. account = await get_point_account(int(giveaway["chat_id"]), user_id)
  346. debit_key: str | None = None
  347. refund_key: str | None = None
  348. pending_debit = False
  349. if entry_cost:
  350. debit_key, refund_key, pending_debit = await _entry_attempt_keys(
  351. giveaway_id, user_id
  352. )
  353. available_balance = int(account.get("balance", 0)) + (
  354. entry_cost if pending_debit else 0
  355. )
  356. if available_balance < max(minimum, entry_cost):
  357. return "insufficient_points", giveaway
  358. if entry_cost and debit_key:
  359. try:
  360. await adjust_points(
  361. chat_id=int(giveaway["chat_id"]),
  362. user_id=user_id,
  363. delta=-entry_cost,
  364. source=SOURCE_GIVEAWAY_ENTRY,
  365. idempotency_key=debit_key,
  366. reference_id=giveaway_id,
  367. reason=f"抽奖 #{giveaway_id} 报名消耗",
  368. username=username,
  369. first_name=first_name,
  370. display_name=display_name,
  371. )
  372. except InsufficientPoints:
  373. return "insufficient_points", giveaway
  374. try:
  375. result = await add_participant(
  376. giveaway_id=giveaway_id,
  377. user_id=user_id,
  378. username=username,
  379. first_name=first_name,
  380. display_name=display_name,
  381. )
  382. except Exception:
  383. if entry_cost:
  384. await _refund_entry(
  385. giveaway,
  386. user_id=user_id,
  387. username=username,
  388. first_name=first_name,
  389. display_name=display_name,
  390. key_suffix="join-failed",
  391. idempotency_key=refund_key,
  392. )
  393. raise
  394. if result != "ok" and entry_cost:
  395. await _refund_entry(
  396. giveaway,
  397. user_id=user_id,
  398. username=username,
  399. first_name=first_name,
  400. display_name=display_name,
  401. key_suffix=f"join-{result}",
  402. idempotency_key=refund_key,
  403. )
  404. if result == "ok":
  405. await refresh_giveaway_message(giveaway)
  406. return result, giveaway
  407. async def buy_additional_tickets(
  408. giveaway_id: str,
  409. *,
  410. user_id: int,
  411. quantity: int,
  412. request_id: str,
  413. username: str | None = None,
  414. first_name: str | None = None,
  415. display_name: str | None = None,
  416. ) -> tuple[str, dict[str, Any] | None]:
  417. """Reserve capacity before charging; replaying a request never charges twice."""
  418. if quantity < 1 or quantity > 100 or not request_id or len(request_id) > 100:
  419. return "invalid_quantity", None
  420. order_id = f"{giveaway_id}:{user_id}:{request_id}"
  421. async with _giveaway_locks[giveaway_id]:
  422. giveaway = await get_giveaway(giveaway_id)
  423. if not giveaway:
  424. return "missing", None
  425. if int(giveaway.get("entry_cost", 0)) <= 0:
  426. return "free_single_ticket", giveaway
  427. order = await ticket_ordersdb.find_one({"order_id": order_id})
  428. if order and order.get("status") == "complete":
  429. return "duplicate", giveaway
  430. if order and order.get("status") == "rejected":
  431. return str(order.get("failure_code") or "ticket_limit"), giveaway
  432. if not order:
  433. now = utc_now()
  434. if giveaway["status"] != STATUS_RUNNING or as_utc(giveaway["ends_at"]) <= now:
  435. return "ended", giveaway
  436. if giveaway_starts_at(giveaway) > now:
  437. return "not_started", giveaway
  438. participant = await get_participant(giveaway_id, user_id)
  439. if participant and participant.get("active") is False:
  440. return "not_joined", giveaway
  441. if await is_giveaway_banned(int(giveaway["chat_id"]), user_id):
  442. return "banned", giveaway
  443. try:
  444. eligible, _ = await check_eligibility(giveaway, user_id)
  445. except EligibilityUnavailable:
  446. return "eligibility_unavailable", giveaway
  447. if not eligible:
  448. return "not_eligible", giveaway
  449. if not participant:
  450. account = await get_point_account(int(giveaway["chat_id"]), user_id)
  451. if int(account.get("balance", 0)) < int(giveaway.get("minimum_points", 0)):
  452. return "insufficient_points", giveaway
  453. try:
  454. await participantsdb.insert_one({
  455. "giveaway_id": giveaway_id,
  456. "chat_id": int(giveaway["chat_id"]),
  457. "user_id": int(user_id),
  458. "username": username,
  459. "first_name": first_name,
  460. "display_name": display_name or first_name,
  461. "entry_cost": int(giveaway["entry_cost"]),
  462. "ticket_count": 0,
  463. "paid_ticket_count": 0,
  464. "pending": True,
  465. "active": True,
  466. "joined_at": now,
  467. })
  468. except DuplicateKeyError:
  469. pass
  470. order = {
  471. "order_id": order_id,
  472. "bot_id": BOT_PROFILE_ID,
  473. "request_id": request_id,
  474. "giveaway_id": giveaway_id,
  475. "chat_id": int(giveaway["chat_id"]),
  476. "user_id": int(user_id),
  477. "quantity": quantity,
  478. "amount": int(giveaway["entry_cost"]) * quantity,
  479. "username": username,
  480. "first_name": first_name,
  481. "display_name": display_name,
  482. "status": "pending",
  483. "created_at": now,
  484. "updated_at": now,
  485. }
  486. try:
  487. await ticket_ordersdb.insert_one(order)
  488. except Exception:
  489. order = await ticket_ordersdb.find_one({"order_id": order_id})
  490. if not order:
  491. raise
  492. quantity = int(order["quantity"])
  493. participant = await get_participant(giveaway_id, user_id)
  494. if participant and "ticket_count" not in participant:
  495. await participantsdb.update_one(
  496. {"giveaway_id": giveaway_id, "user_id": user_id, "ticket_count": {"$exists": False}},
  497. {"$set": {"ticket_count": 1, "paid_ticket_count": 1}},
  498. )
  499. if not participant or participant.get("active") is False or giveaway["status"] in {STATUS_CANCELED, STATUS_CANCELING, STATUS_FINISHED}:
  500. debit_key = f"giveaway-ticket:{order_id}"
  501. if await get_point_transaction_by_key(debit_key):
  502. await adjust_points(
  503. chat_id=int(giveaway["chat_id"]), user_id=user_id,
  504. delta=int(order["amount"]), source=SOURCE_GIVEAWAY_REFUND,
  505. idempotency_key=f"giveaway-ticket-refund:{order_id}",
  506. reference_id=giveaway_id, reason=f"抽奖 #{giveaway_id} 奖票购买取消退款",
  507. username=username, first_name=first_name, display_name=display_name,
  508. )
  509. if participant and order_id in participant.get("ticket_order_ids", []):
  510. await participantsdb.update_one(
  511. {"giveaway_id": giveaway_id, "user_id": user_id, "ticket_order_ids": order_id},
  512. {"$inc": {"ticket_count": -quantity}, "$pull": {"ticket_order_ids": order_id}},
  513. )
  514. await ticket_ordersdb.update_one(
  515. {"order_id": order_id},
  516. {"$set": {"status": "rejected", "failure_code": "not_joined", "updated_at": utc_now()}},
  517. )
  518. return "not_joined", giveaway
  519. participant_filter = {
  520. "giveaway_id": giveaway_id,
  521. "user_id": int(user_id),
  522. "active": {"$ne": False},
  523. }
  524. reserved = await participantsdb.update_one(
  525. {
  526. **participant_filter,
  527. "ticket_order_ids": {"$ne": order_id},
  528. "ticket_count": {
  529. "$lte": int(giveaway.get("max_tickets_per_user", 10)) - quantity
  530. },
  531. },
  532. {"$inc": {"ticket_count": quantity}, "$addToSet": {"ticket_order_ids": order_id}},
  533. )
  534. if not reserved.modified_count:
  535. participant = await get_participant(giveaway_id, user_id)
  536. if not participant or order_id not in participant.get("ticket_order_ids", []):
  537. await ticket_ordersdb.update_one(
  538. {"order_id": order_id}, {"$set": {"status": "rejected", "failure_code": "ticket_limit", "updated_at": utc_now()}}
  539. )
  540. return "ticket_limit", giveaway
  541. try:
  542. await adjust_points(
  543. chat_id=int(giveaway["chat_id"]), user_id=user_id,
  544. delta=-int(order["amount"]), source=SOURCE_GIVEAWAY_ENTRY,
  545. idempotency_key=f"giveaway-ticket:{order_id}", reference_id=giveaway_id,
  546. reason=f"抽奖 #{giveaway_id} 购买 {quantity} 张奖票",
  547. username=username, first_name=first_name, display_name=display_name,
  548. )
  549. except InsufficientPoints:
  550. await participantsdb.update_one(
  551. {**participant_filter, "ticket_order_ids": order_id},
  552. {"$inc": {"ticket_count": -quantity}, "$pull": {"ticket_order_ids": order_id}},
  553. )
  554. await ticket_ordersdb.update_one(
  555. {"order_id": order_id}, {"$set": {"status": "rejected", "failure_code": "insufficient_points", "updated_at": utc_now()}}
  556. )
  557. return "insufficient_points", giveaway
  558. await participantsdb.update_one(
  559. {**participant_filter, "completed_ticket_order_ids": {"$ne": order_id}},
  560. {
  561. "$inc": {"paid_ticket_count": quantity},
  562. "$addToSet": {"completed_ticket_order_ids": order_id},
  563. "$set": {"pending": False},
  564. },
  565. )
  566. await ticket_ordersdb.update_one(
  567. {"order_id": order_id}, {"$set": {"status": "complete", "updated_at": utc_now()}}
  568. )
  569. await refresh_giveaway_message(giveaway)
  570. return "ok", giveaway
  571. async def reconcile_pending_ticket_orders(limit: int = 50) -> None:
  572. cursor = ticket_ordersdb.find({"bot_id": BOT_PROFILE_ID, "status": "pending"}).sort("created_at", 1).limit(limit)
  573. for order in [item async for item in cursor]:
  574. try:
  575. await buy_additional_tickets(
  576. str(order["giveaway_id"]), user_id=int(order["user_id"]),
  577. quantity=int(order["quantity"]),
  578. request_id=str(order.get("request_id") or order["order_id"].rsplit(":", 1)[-1]),
  579. username=order.get("username"), first_name=order.get("first_name"),
  580. display_name=order.get("display_name"),
  581. )
  582. except Exception as exc:
  583. log.error(f"奖票订单 {order['order_id']} 恢复失败:{exc}")
  584. async def _entry_attempt_keys(
  585. giveaway_id: str, user_id: int
  586. ) -> tuple[str, str, bool]:
  587. base = f"giveaway-entry:{giveaway_id}:{user_id}"
  588. for attempt in range(100):
  589. suffix = "" if attempt == 0 else f":{attempt}"
  590. debit_key = f"{base}{suffix}"
  591. refund_key = f"giveaway-refund:{giveaway_id}:{user_id}:entry-attempt:{attempt}"
  592. debit = await get_point_transaction_by_key(debit_key)
  593. if not debit:
  594. return debit_key, refund_key, False
  595. if not await get_point_transaction_by_key(refund_key):
  596. return debit_key, refund_key, True
  597. raise GiveawayServiceError(
  598. "entry_retry_limit", "报名重试次数过多,请稍后再试。", status=409
  599. )
  600. async def _refund_entry(
  601. giveaway: dict[str, Any],
  602. *,
  603. user_id: int,
  604. username: str | None,
  605. first_name: str | None,
  606. display_name: str | None,
  607. key_suffix: str,
  608. idempotency_key: str | None = None,
  609. amount: int | None = None,
  610. ) -> bool:
  611. entry_cost = int(giveaway.get("entry_cost", 0)) if amount is None else int(amount)
  612. if entry_cost <= 0:
  613. return False
  614. _, created = await adjust_points(
  615. chat_id=int(giveaway["chat_id"]),
  616. user_id=user_id,
  617. delta=entry_cost,
  618. source=SOURCE_GIVEAWAY_REFUND,
  619. idempotency_key=idempotency_key
  620. or f"giveaway-refund:{giveaway['giveaway_id']}:{user_id}:{key_suffix}",
  621. reference_id=giveaway["giveaway_id"],
  622. reason=f"抽奖 #{giveaway['giveaway_id']} 报名退款",
  623. username=username,
  624. first_name=first_name,
  625. display_name=display_name,
  626. )
  627. if created:
  628. await mark_participant_refunded(
  629. giveaway["giveaway_id"], user_id, key_suffix
  630. )
  631. return created
  632. async def _award_draw_points(
  633. giveaway: dict[str, Any],
  634. participants: list[dict[str, Any]],
  635. winners: list[dict[str, Any]],
  636. ) -> None:
  637. participation_reward = int(giveaway.get("participation_reward", 0))
  638. if participation_reward:
  639. for participant in participants:
  640. await adjust_points(
  641. chat_id=int(giveaway["chat_id"]),
  642. user_id=int(participant["user_id"]),
  643. delta=participation_reward,
  644. source=SOURCE_GIVEAWAY_PARTICIPATION,
  645. idempotency_key=(
  646. f"giveaway-participation:{giveaway['giveaway_id']}:"
  647. f"{participant['user_id']}"
  648. ),
  649. reference_id=giveaway["giveaway_id"],
  650. reason=f"抽奖 #{giveaway['giveaway_id']} 参与奖励",
  651. username=participant.get("username"),
  652. first_name=participant.get("first_name"),
  653. display_name=participant.get("display_name"),
  654. )
  655. for winner in winners:
  656. reward = int(winner.get("points_reward", 0))
  657. if not reward:
  658. continue
  659. await adjust_points(
  660. chat_id=int(giveaway["chat_id"]),
  661. user_id=int(winner["user_id"]),
  662. delta=reward,
  663. source=SOURCE_GIVEAWAY_WINNER,
  664. idempotency_key=(
  665. f"giveaway-winner:{giveaway['giveaway_id']}:"
  666. f"{winner.get('tier_index', 0)}:{winner['user_id']}"
  667. ),
  668. reference_id=giveaway["giveaway_id"],
  669. reason=(
  670. f"抽奖 #{giveaway['giveaway_id']} 中奖奖励"
  671. f"({winner['tier']})"
  672. ),
  673. username=winner.get("username"),
  674. first_name=winner.get("first_name"),
  675. display_name=winner.get("display_name"),
  676. )
  677. async def finish_and_publish_giveaway(
  678. giveaway_id: str, *, publish: bool = True
  679. ) -> tuple[bool, str, dict[str, Any] | None]:
  680. async with _giveaway_locks[giveaway_id]:
  681. current = await get_giveaway(giveaway_id)
  682. if (
  683. current
  684. and current["status"] in {STATUS_RUNNING, STATUS_DRAWING}
  685. and current.get("template_id")
  686. and not current.get("message_id")
  687. ):
  688. return False, "报名公告未发布,本期不能开奖。", current
  689. giveaway = await claim_giveaway_draw(giveaway_id)
  690. if not giveaway:
  691. existing = await get_giveaway(giveaway_id)
  692. if not existing:
  693. return False, "未找到该抽奖。", None
  694. if existing["status"] == STATUS_FINISHED:
  695. return False, "该抽奖已经完成开奖。", existing
  696. return False, "该抽奖当前不在进行中。", existing
  697. pending_order = await ticket_ordersdb.find_one(
  698. {"giveaway_id": giveaway_id, "status": "pending"}
  699. )
  700. if pending_order:
  701. raise GiveawayServiceError(
  702. "ticket_order_pending", "奖票订单尚未结算,开奖将稍后重试。", status=503
  703. )
  704. participants = await list_participants(giveaway_id)
  705. eligible_participants = [
  706. participant for participant in participants
  707. if participant.get("eligibility_status") == "eligible"
  708. ]
  709. if giveaway.get("pending_winners") is not None and not any(
  710. "eligibility_status" in item for item in participants
  711. ) and not giveaway.get("eligibility_targets"):
  712. eligible_participants = participants
  713. if giveaway.get("pending_winners") is None:
  714. eligible_participants = []
  715. try:
  716. for participant in participants:
  717. eligible, failed = await check_eligibility(
  718. giveaway, int(participant["user_id"])
  719. )
  720. await participantsdb.update_one(
  721. {"giveaway_id": giveaway_id, "user_id": participant["user_id"]},
  722. {"$set": {
  723. "eligibility_status": "eligible" if eligible else "ineligible",
  724. "eligibility_failed_chat_ids": failed,
  725. "eligibility_checked_at": utc_now(),
  726. }},
  727. )
  728. if eligible:
  729. eligible_participants.append(participant)
  730. except EligibilityUnavailable as exc:
  731. changed = await giveawaysdb.update_one(
  732. {
  733. "giveaway_id": giveaway_id,
  734. "status": STATUS_DRAWING,
  735. "eligibility_error": {"$ne": str(exc)},
  736. },
  737. {"$set": {"eligibility_error": str(exc), "updated_at": utc_now()}},
  738. )
  739. if changed.modified_count:
  740. with suppress(Exception):
  741. await record_audit(
  742. source="system", actor_id=None, actor_name="抽奖定时任务",
  743. action="giveaway.eligibility_paused",
  744. chat_id=int(giveaway["chat_id"]), target_id=giveaway_id,
  745. summary="开奖资格核验失败,已暂停并等待重试",
  746. success=False, error=str(exc),
  747. )
  748. raise GiveawayServiceError(
  749. "eligibility_unavailable", "资格核验暂不可用,开奖已暂停并将自动重试。", status=503
  750. ) from exc
  751. await giveawaysdb.update_one(
  752. {"giveaway_id": giveaway_id}, {"$unset": {"eligibility_error": ""}}
  753. )
  754. pending = pick_winners(eligible_participants, giveaway["prizes"])
  755. giveaway = await save_pending_winners(giveaway_id, pending) or giveaway
  756. winners = list(giveaway.get("pending_winners", []))
  757. await _award_draw_points(giveaway, eligible_participants, winners)
  758. marked = await finish_giveaway(giveaway_id, winners)
  759. if not marked:
  760. current = await get_giveaway(giveaway_id)
  761. if not current or current["status"] != STATUS_FINISHED:
  762. return False, "抽奖无法完成,请稍后重试。", current
  763. giveaway = current
  764. winners = list(current.get("winners", []))
  765. else:
  766. giveaway = {**giveaway, "status": STATUS_FINISHED, "winners": winners}
  767. if publish:
  768. await publish_finished(giveaway, winners)
  769. return True, f"抽奖 #{giveaway_id} 已完成开奖。", giveaway
  770. async def publish_finished(
  771. giveaway: dict[str, Any], winners: list[dict[str, Any]]
  772. ) -> None:
  773. giveaway = await get_giveaway(giveaway["giveaway_id"]) or giveaway
  774. chat_id = int(giveaway["chat_id"])
  775. message_id = giveaway.get("message_id")
  776. text = render_winners(giveaway, winners)
  777. result_message = None
  778. if giveaway.get("result_message_id"):
  779. from types import SimpleNamespace
  780. result_message = SimpleNamespace(id=int(giveaway["result_message_id"]))
  781. if message_id:
  782. with suppress(Exception):
  783. await app.edit_message_text(
  784. chat_id,
  785. int(message_id),
  786. await render_giveaway({**giveaway, "status": STATUS_FINISHED}),
  787. parse_mode=ParseMode.HTML,
  788. reply_markup=None,
  789. )
  790. if result_message is None:
  791. with suppress(Exception):
  792. result_message = await app.send_message(
  793. chat_id, text, parse_mode=ParseMode.HTML,
  794. reply_to_message_id=int(message_id),
  795. )
  796. if result_message is None:
  797. with suppress(Exception):
  798. result_message = await app.send_message(
  799. chat_id, text, parse_mode=ParseMode.HTML
  800. )
  801. if result_message is None:
  802. log.error(f"抽奖 {giveaway['giveaway_id']} 结果发送失败,等待重试")
  803. return
  804. await giveawaysdb.update_one(
  805. {"giveaway_id": giveaway["giveaway_id"]},
  806. {"$set": {"result_message_id": int(result_message.id), "updated_at": utc_now()}},
  807. )
  808. result_pinned = bool(giveaway.get("result_pinned"))
  809. if not result_pinned:
  810. try:
  811. await app.pin_chat_message(
  812. chat_id,
  813. int(result_message.id),
  814. disable_notification=True,
  815. )
  816. result_pinned = True
  817. except Exception as exc:
  818. log.error(
  819. f"抽奖 {giveaway['giveaway_id']} 结果置顶失败:{exc}"
  820. )
  821. if result_pinned and message_id and int(message_id) != int(result_message.id):
  822. try:
  823. await app.delete_messages(chat_id, int(message_id))
  824. except Exception:
  825. # If Telegram refuses deletion, at least remove the stale pinned entry.
  826. with suppress(Exception):
  827. await app.unpin_chat_message(chat_id, int(message_id))
  828. await giveawaysdb.update_one(
  829. {"giveaway_id": giveaway["giveaway_id"]},
  830. {"$set": {
  831. "result_pinned": result_pinned,
  832. "publication_pending": not result_pinned,
  833. "updated_at": utc_now(),
  834. }},
  835. )
  836. if not result_pinned:
  837. return
  838. for winner in winners:
  839. with suppress(Exception):
  840. await app.send_message(
  841. int(winner["user_id"]),
  842. (
  843. f"恭喜你在 <b>{escape(giveaway['title'])}</b> 中获得"
  844. f"<b>{escape(str(winner['tier']))}</b>。"
  845. ),
  846. parse_mode=ParseMode.HTML,
  847. )
  848. async def cancel_and_refund_giveaway(
  849. giveaway_id: str, *, chat_id: int, publish: bool = True
  850. ) -> tuple[bool, str, dict[str, Any] | None]:
  851. await reconcile_pending_ticket_orders()
  852. async with _giveaway_locks[giveaway_id]:
  853. giveaway = await claim_giveaway_cancel(giveaway_id, chat_id)
  854. if not giveaway:
  855. existing = await get_giveaway(giveaway_id)
  856. if not existing:
  857. return False, "未找到该抽奖。", None
  858. if existing["status"] == STATUS_CANCELED:
  859. return False, "该抽奖已经取消。", existing
  860. return False, "该抽奖当前不在进行中。", existing
  861. participants = await list_participants(giveaway_id)
  862. for participant in participants:
  863. await _refund_entry(
  864. giveaway,
  865. user_id=int(participant["user_id"]),
  866. username=participant.get("username"),
  867. first_name=participant.get("first_name"),
  868. display_name=participant.get("display_name"),
  869. key_suffix="cancel",
  870. amount=int(giveaway.get("entry_cost", 0)) * int(
  871. participant.get("paid_ticket_count", participant.get("ticket_count", 1))
  872. ),
  873. )
  874. await finalize_giveaway_cancel(giveaway_id)
  875. giveaway = {**giveaway, "status": STATUS_CANCELED}
  876. if publish and giveaway.get("message_id"):
  877. with suppress(Exception):
  878. await app.edit_message_text(
  879. int(giveaway["chat_id"]),
  880. int(giveaway["message_id"]),
  881. await render_giveaway(giveaway),
  882. parse_mode=ParseMode.HTML,
  883. reply_markup=None,
  884. )
  885. return True, f"抽奖 #{giveaway_id} 已取消,报名积分已退还。", giveaway
  886. async def remove_and_optionally_refund_participant(
  887. *,
  888. giveaway_id: str,
  889. user_id: int,
  890. moderator_id: int,
  891. reason: str,
  892. refund: bool = True,
  893. ) -> bool:
  894. giveaway = await get_giveaway(giveaway_id)
  895. if not giveaway or giveaway["status"] != STATUS_RUNNING:
  896. raise GiveawayServiceError(
  897. "giveaway_not_running", "该抽奖当前不在进行中。"
  898. )
  899. participant = await remove_participant(
  900. giveaway_id=giveaway_id,
  901. user_id=user_id,
  902. moderator_id=moderator_id,
  903. reason=reason,
  904. refund=refund,
  905. )
  906. if not participant:
  907. return False
  908. if refund:
  909. await _refund_entry(
  910. giveaway,
  911. user_id=user_id,
  912. username=participant.get("username"),
  913. first_name=participant.get("first_name"),
  914. display_name=participant.get("display_name"),
  915. key_suffix="participant-removed",
  916. amount=int(giveaway.get("entry_cost", 0)) * int(
  917. participant.get("paid_ticket_count", participant.get("ticket_count", 1))
  918. ),
  919. )
  920. await refresh_giveaway_message(giveaway)
  921. return True
  922. async def reroll_giveaway(
  923. *,
  924. giveaway_id: str,
  925. moderator_id: int,
  926. tier_name: str | None = None,
  927. reroll_id: str | None = None,
  928. publish: bool = True,
  929. ) -> tuple[list[dict[str, Any]], str]:
  930. reroll_id = reroll_id or secrets.token_hex(8)
  931. async with _giveaway_locks[giveaway_id]:
  932. giveaway = await get_giveaway(giveaway_id)
  933. if not giveaway or giveaway["status"] != STATUS_FINISHED:
  934. raise GiveawayServiceError(
  935. "giveaway_not_finished", "只有已开奖的抽奖可以重抽。"
  936. )
  937. for previous in giveaway.get("rerolls", []):
  938. if previous.get("reroll_id") == reroll_id:
  939. return list(previous.get("winners", [])), reroll_id
  940. prizes = list(giveaway["prizes"])
  941. if tier_name:
  942. prizes = [
  943. prize
  944. for prize in prizes
  945. if str(prize["name"]).strip().lower() == tier_name.strip().lower()
  946. ]
  947. if not prizes:
  948. raise GiveawayServiceError("tier_not_found", "未找到该奖项。")
  949. excluded = {int(winner["user_id"]) for winner in giveaway.get("winners", [])}
  950. for previous in giveaway.get("rerolls", []):
  951. excluded.update(
  952. int(winner["user_id"]) for winner in previous.get("winners", [])
  953. )
  954. participants = [
  955. participant
  956. for participant in await list_participants(giveaway_id)
  957. if int(participant["user_id"]) not in excluded
  958. ]
  959. eligible_participants = []
  960. try:
  961. for participant in participants:
  962. eligible, _ = await check_eligibility(giveaway, int(participant["user_id"]))
  963. if eligible:
  964. eligible_participants.append(participant)
  965. except EligibilityUnavailable as exc:
  966. raise GiveawayServiceError(
  967. "eligibility_unavailable", "资格核验暂不可用,重抽已暂停。", status=503
  968. ) from exc
  969. winners = pick_winners(eligible_participants, prizes)
  970. await record_reroll(
  971. giveaway_id,
  972. winners,
  973. tier_name,
  974. moderator_id,
  975. reroll_id=reroll_id,
  976. )
  977. for winner in winners:
  978. reward = int(winner.get("points_reward", 0))
  979. if reward:
  980. await adjust_points(
  981. chat_id=int(giveaway["chat_id"]),
  982. user_id=int(winner["user_id"]),
  983. delta=reward,
  984. source=SOURCE_GIVEAWAY_WINNER,
  985. idempotency_key=(
  986. f"giveaway-reroll-winner:{giveaway_id}:{reroll_id}:"
  987. f"{winner.get('tier_index', 0)}:{winner['user_id']}"
  988. ),
  989. reference_id=giveaway_id,
  990. reason=f"抽奖 #{giveaway_id} 重抽中奖奖励",
  991. username=winner.get("username"),
  992. first_name=winner.get("first_name"),
  993. display_name=winner.get("display_name"),
  994. )
  995. if publish:
  996. with suppress(Exception):
  997. await app.send_message(
  998. int(giveaway["chat_id"]),
  999. render_winners(giveaway, winners, reroll=True, tier_name=tier_name),
  1000. parse_mode=ParseMode.HTML,
  1001. disable_web_page_preview=True,
  1002. )
  1003. return winners, reroll_id
  1004. async def resume_pending_giveaway(giveaway: dict[str, Any]) -> None:
  1005. try:
  1006. if (
  1007. giveaway["status"] == STATUS_RUNNING
  1008. and giveaway.get("template_id")
  1009. and not giveaway.get("message_id")
  1010. and as_utc(giveaway["ends_at"]) <= utc_now()
  1011. ):
  1012. canceled, _, _ = await cancel_and_refund_giveaway(
  1013. giveaway["giveaway_id"], chat_id=int(giveaway["chat_id"]),
  1014. publish=False,
  1015. )
  1016. if canceled:
  1017. await giveawaysdb.update_one(
  1018. {"giveaway_id": giveaway["giveaway_id"], "status": STATUS_CANCELED},
  1019. {"$set": {
  1020. "cancellation_reason": "announcement_not_published",
  1021. "updated_at": utc_now(),
  1022. }},
  1023. )
  1024. with suppress(Exception):
  1025. await record_audit(
  1026. source="system", actor_id=None, actor_name="抽奖定时任务",
  1027. action="giveaway.announcement_missed",
  1028. chat_id=int(giveaway["chat_id"]),
  1029. target_id=giveaway["giveaway_id"],
  1030. summary="报名公告未发布,开奖时间已到,本期自动取消",
  1031. success=False,
  1032. )
  1033. return
  1034. if giveaway["status"] in {STATUS_RUNNING, STATUS_DRAWING}:
  1035. await finish_and_publish_giveaway(giveaway["giveaway_id"])
  1036. elif giveaway["status"] == STATUS_CANCELING:
  1037. await cancel_and_refund_giveaway(
  1038. giveaway["giveaway_id"], chat_id=int(giveaway["chat_id"])
  1039. )
  1040. except Exception as exc:
  1041. log.error(
  1042. f"抽奖 {giveaway.get('giveaway_id')} 恢复任务失败:{exc}"
  1043. )