dbpoints.py 27 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837
  1. from __future__ import annotations
  2. import asyncio
  3. import hashlib
  4. import re
  5. from collections import defaultdict
  6. from datetime import UTC, datetime, timedelta
  7. from secrets import token_hex
  8. from typing import Any
  9. from zoneinfo import ZoneInfo, ZoneInfoNotFoundError
  10. from pymongo import ASCENDING, DESCENDING
  11. from pymongo.errors import DuplicateKeyError
  12. from wbb import db
  13. SOURCE_CHECKIN = "checkin"
  14. SOURCE_ACTIVITY = "activity"
  15. SOURCE_UPVOTE = "upvote"
  16. SOURCE_ADMIN = "admin_adjustment"
  17. SOURCE_GIVEAWAY_ENTRY = "giveaway_entry"
  18. SOURCE_GIVEAWAY_REFUND = "giveaway_refund"
  19. SOURCE_GIVEAWAY_PARTICIPATION = "giveaway_participation"
  20. SOURCE_GIVEAWAY_WINNER = "giveaway_winner"
  21. DEFAULT_CHECKIN_SUCCESS_MESSAGE = (
  22. "签到成功:+{points} 积分,当前余额:{balance}。"
  23. )
  24. DEFAULT_CHECKIN_ALREADY_MESSAGE = (
  25. "你今天已经签到过了,当前余额:{balance}。"
  26. )
  27. CHECKIN_MESSAGE_MAX_LENGTH = 500
  28. POINT_SOURCES = {
  29. SOURCE_CHECKIN,
  30. SOURCE_ACTIVITY,
  31. SOURCE_UPVOTE,
  32. SOURCE_ADMIN,
  33. SOURCE_GIVEAWAY_ENTRY,
  34. SOURCE_GIVEAWAY_REFUND,
  35. SOURCE_GIVEAWAY_PARTICIPATION,
  36. SOURCE_GIVEAWAY_WINNER,
  37. }
  38. DEFAULT_POINT_RULES: dict[str, Any] = {
  39. "enabled": False,
  40. "timezone": "Asia/Shanghai",
  41. "checkin_enabled": True,
  42. "checkin_button_enabled": False,
  43. "checkin_points": 10,
  44. "checkin_success_message": DEFAULT_CHECKIN_SUCCESS_MESSAGE,
  45. "checkin_already_message": DEFAULT_CHECKIN_ALREADY_MESSAGE,
  46. "activity_enabled": True,
  47. "activity_points": 1,
  48. "activity_cooldown_seconds": 300,
  49. "activity_daily_cap": 10,
  50. "upvote_enabled": True,
  51. "upvote_points": 2,
  52. "upvote_pair_cooldown_seconds": 86400,
  53. "upvote_daily_cap": 10,
  54. }
  55. accountsdb = db.point_accounts
  56. transactionsdb = db.point_transactions
  57. rulesdb = db.point_rules
  58. rate_limitsdb = db.point_rate_limits
  59. _index_lock = asyncio.Lock()
  60. _indexes_ready = False
  61. _account_locks: defaultdict[tuple[int, int], asyncio.Lock] = defaultdict(asyncio.Lock)
  62. _rate_locks: defaultdict[str, asyncio.Lock] = defaultdict(asyncio.Lock)
  63. class PointsError(RuntimeError):
  64. pass
  65. class InsufficientPoints(PointsError):
  66. pass
  67. def utc_now() -> datetime:
  68. return datetime.now(UTC)
  69. def _as_utc(value: datetime) -> datetime:
  70. if value.tzinfo is None:
  71. return value.replace(tzinfo=UTC)
  72. return value.astimezone(UTC)
  73. def _local_day(now: datetime, timezone_name: str) -> str:
  74. try:
  75. zone = ZoneInfo(timezone_name)
  76. except ZoneInfoNotFoundError:
  77. zone = ZoneInfo("Asia/Shanghai")
  78. return _as_utc(now).astimezone(zone).date().isoformat()
  79. def _safe_int(value: Any, default: int, minimum: int, maximum: int) -> int:
  80. try:
  81. parsed = int(value)
  82. except (TypeError, ValueError):
  83. return default
  84. return max(minimum, min(parsed, maximum))
  85. def _safe_bool(value: Any, default: bool) -> bool:
  86. if isinstance(value, bool):
  87. return value
  88. if isinstance(value, (int, float)):
  89. return value != 0
  90. if isinstance(value, str):
  91. normalized = value.strip().lower()
  92. if normalized in {
  93. "true",
  94. "1",
  95. "yes",
  96. "on",
  97. "enable",
  98. "enabled",
  99. "开启",
  100. "打开",
  101. }:
  102. return True
  103. if normalized in {
  104. "false",
  105. "0",
  106. "no",
  107. "off",
  108. "disable",
  109. "disabled",
  110. "关闭",
  111. }:
  112. return False
  113. return default
  114. def _safe_text(value: Any, default: str, maximum: int) -> str:
  115. text = str(value or "").strip()
  116. if not text:
  117. return default
  118. return text[:maximum]
  119. def normalize_point_rules(raw: dict[str, Any] | None) -> dict[str, Any]:
  120. rules = dict(DEFAULT_POINT_RULES)
  121. if raw:
  122. rules.update({key: value for key, value in raw.items() if key in rules})
  123. for key in (
  124. "enabled",
  125. "checkin_enabled",
  126. "checkin_button_enabled",
  127. "activity_enabled",
  128. "upvote_enabled",
  129. ):
  130. rules[key] = _safe_bool(rules[key], DEFAULT_POINT_RULES[key])
  131. timezone_name = str(rules.get("timezone") or "Asia/Shanghai")
  132. try:
  133. ZoneInfo(timezone_name)
  134. except ZoneInfoNotFoundError:
  135. timezone_name = "Asia/Shanghai"
  136. rules["timezone"] = timezone_name
  137. rules["checkin_points"] = _safe_int(rules["checkin_points"], 10, 0, 100000)
  138. rules["checkin_success_message"] = _safe_text(
  139. rules["checkin_success_message"],
  140. DEFAULT_CHECKIN_SUCCESS_MESSAGE,
  141. CHECKIN_MESSAGE_MAX_LENGTH,
  142. )
  143. rules["checkin_already_message"] = _safe_text(
  144. rules["checkin_already_message"],
  145. DEFAULT_CHECKIN_ALREADY_MESSAGE,
  146. CHECKIN_MESSAGE_MAX_LENGTH,
  147. )
  148. rules["activity_points"] = _safe_int(rules["activity_points"], 1, 0, 100000)
  149. rules["activity_cooldown_seconds"] = _safe_int(
  150. rules["activity_cooldown_seconds"], 300, 10, 86400
  151. )
  152. rules["activity_daily_cap"] = _safe_int(
  153. rules["activity_daily_cap"], 10, 0, 1000000
  154. )
  155. rules["upvote_points"] = _safe_int(rules["upvote_points"], 2, 0, 100000)
  156. rules["upvote_pair_cooldown_seconds"] = _safe_int(
  157. rules["upvote_pair_cooldown_seconds"], 86400, 60, 2592000
  158. )
  159. rules["upvote_daily_cap"] = _safe_int(
  160. rules["upvote_daily_cap"], 10, 0, 1000000
  161. )
  162. return rules
  163. async def ensure_point_indexes() -> None:
  164. global _indexes_ready
  165. if _indexes_ready:
  166. return
  167. async with _index_lock:
  168. if _indexes_ready:
  169. return
  170. await accountsdb.create_index(
  171. [("chat_id", ASCENDING), ("user_id", ASCENDING)], unique=True
  172. )
  173. await accountsdb.create_index(
  174. [("chat_id", ASCENDING), ("balance", DESCENDING), ("user_id", ASCENDING)]
  175. )
  176. await transactionsdb.create_index([("idempotency_key", ASCENDING)], unique=True)
  177. await transactionsdb.create_index(
  178. [("chat_id", ASCENDING), ("user_id", ASCENDING), ("created_at", DESCENDING)]
  179. )
  180. await transactionsdb.create_index(
  181. [("chat_id", ASCENDING), ("source", ASCENDING), ("created_at", DESCENDING)]
  182. )
  183. await rulesdb.create_index([("chat_id", ASCENDING)], unique=True)
  184. await rate_limitsdb.create_index([("rate_key", ASCENDING)], unique=True)
  185. await rate_limitsdb.create_index("expires_at", expireAfterSeconds=0)
  186. _indexes_ready = True
  187. async def get_point_rules(chat_id: int) -> dict[str, Any]:
  188. await ensure_point_indexes()
  189. doc = await rulesdb.find_one({"chat_id": chat_id})
  190. return normalize_point_rules(doc)
  191. async def set_point_rules(chat_id: int, values: dict[str, Any]) -> dict[str, Any]:
  192. await ensure_point_indexes()
  193. current = await get_point_rules(chat_id)
  194. current.update({key: value for key, value in values.items() if key in DEFAULT_POINT_RULES})
  195. normalized = normalize_point_rules(current)
  196. now = utc_now()
  197. await rulesdb.update_one(
  198. {"chat_id": chat_id},
  199. {"$set": {**normalized, "updated_at": now}, "$setOnInsert": {"created_at": now}},
  200. upsert=True,
  201. )
  202. return normalized
  203. async def list_checkin_keyboard_migrations(
  204. keyboard_version: int,
  205. ) -> list[int]:
  206. await ensure_point_indexes()
  207. cursor = rulesdb.find(
  208. {
  209. "enabled": True,
  210. "checkin_enabled": True,
  211. "checkin_button_enabled": True,
  212. "$or": [
  213. {"checkin_keyboard_version": {"$exists": False}},
  214. {"checkin_keyboard_version": {"$lt": int(keyboard_version)}},
  215. ],
  216. },
  217. {"chat_id": 1},
  218. )
  219. documents = await cursor.to_list(length=10000)
  220. return [int(document["chat_id"]) for document in documents]
  221. async def mark_checkin_keyboard_version(
  222. chat_id: int,
  223. keyboard_version: int,
  224. ) -> None:
  225. await ensure_point_indexes()
  226. await rulesdb.update_one(
  227. {"chat_id": int(chat_id)},
  228. {
  229. "$set": {
  230. "checkin_keyboard_version": int(keyboard_version),
  231. "updated_at": utc_now(),
  232. }
  233. },
  234. )
  235. async def _reconcile_account_locked(chat_id: int, user_id: int) -> dict[str, Any]:
  236. cached = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  237. pipeline = [
  238. {"$match": {"chat_id": chat_id, "user_id": user_id}},
  239. {
  240. "$group": {
  241. "_id": None,
  242. "balance": {"$sum": "$delta"},
  243. "lifetime_earned": {
  244. "$sum": {"$cond": [{"$gt": ["$delta", 0]}, "$delta", 0]}
  245. },
  246. "lifetime_spent": {
  247. "$sum": {
  248. "$cond": [
  249. {"$lt": ["$delta", 0]},
  250. {"$multiply": ["$delta", -1]},
  251. 0,
  252. ]
  253. }
  254. },
  255. }
  256. },
  257. ]
  258. totals = [doc async for doc in transactionsdb.aggregate(pipeline)]
  259. latest = await transactionsdb.find_one(
  260. {"chat_id": chat_id, "user_id": user_id},
  261. sort=[("created_at", DESCENDING), ("_id", DESCENDING)],
  262. )
  263. total = totals[0] if totals else {}
  264. now = utc_now()
  265. account = {
  266. "chat_id": chat_id,
  267. "user_id": user_id,
  268. "balance": int(total.get("balance", 0)),
  269. "lifetime_earned": int(total.get("lifetime_earned", 0)),
  270. "lifetime_spent": int(total.get("lifetime_spent", 0)),
  271. "last_transaction_id": latest.get("transaction_id") if latest else None,
  272. "updated_at": now,
  273. }
  274. if latest:
  275. account["username"] = (
  276. cached.get("username") if cached and "username" in cached else latest.get("username")
  277. )
  278. account["first_name"] = (
  279. cached.get("first_name")
  280. if cached and "first_name" in cached
  281. else latest.get("first_name")
  282. )
  283. account["display_name"] = (
  284. cached.get("display_name")
  285. if cached and "display_name" in cached
  286. else latest.get("display_name") or latest.get("first_name")
  287. )
  288. await accountsdb.update_one(
  289. {"chat_id": chat_id, "user_id": user_id},
  290. {"$set": account, "$setOnInsert": {"created_at": now}},
  291. upsert=True,
  292. )
  293. return account
  294. async def reconcile_account(chat_id: int, user_id: int) -> dict[str, Any]:
  295. await ensure_point_indexes()
  296. async with _account_locks[(chat_id, user_id)]:
  297. return await _reconcile_account_locked(chat_id, user_id)
  298. async def get_point_account(chat_id: int, user_id: int) -> dict[str, Any]:
  299. await ensure_point_indexes()
  300. account = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  301. latest = await transactionsdb.find_one(
  302. {"chat_id": chat_id, "user_id": user_id},
  303. projection={"transaction_id": 1},
  304. sort=[("created_at", DESCENDING), ("_id", DESCENDING)],
  305. )
  306. if latest and (
  307. not account or account.get("last_transaction_id") != latest.get("transaction_id")
  308. ):
  309. return await reconcile_account(chat_id, user_id)
  310. if account:
  311. return account
  312. return {
  313. "chat_id": chat_id,
  314. "user_id": user_id,
  315. "balance": 0,
  316. "lifetime_earned": 0,
  317. "lifetime_spent": 0,
  318. "last_transaction_id": None,
  319. }
  320. async def update_point_account_identity(
  321. *,
  322. chat_id: int,
  323. user_id: int,
  324. username: str | None,
  325. first_name: str | None,
  326. last_name: str | None,
  327. display_name: str,
  328. ) -> None:
  329. await ensure_point_indexes()
  330. await accountsdb.update_one(
  331. {"chat_id": int(chat_id), "user_id": int(user_id)},
  332. {
  333. "$set": {
  334. "username": username,
  335. "first_name": first_name,
  336. "last_name": last_name,
  337. "display_name": display_name or first_name,
  338. "identity_updated_at": utc_now(),
  339. }
  340. },
  341. )
  342. async def _write_points_locked(
  343. *,
  344. chat_id: int,
  345. user_id: int,
  346. source: str,
  347. idempotency_key: str,
  348. delta: int | None = None,
  349. target_balance: int | None = None,
  350. actor_id: int | str | None = None,
  351. reason: str = "",
  352. reference_id: str | None = None,
  353. username: str | None = None,
  354. first_name: str | None = None,
  355. display_name: str | None = None,
  356. ) -> tuple[dict[str, Any], bool]:
  357. existing = await transactionsdb.find_one({"idempotency_key": idempotency_key})
  358. if existing:
  359. account = await _reconcile_account_locked(chat_id, user_id)
  360. return account, False
  361. account = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  362. latest = await transactionsdb.find_one(
  363. {"chat_id": chat_id, "user_id": user_id},
  364. projection={"transaction_id": 1},
  365. sort=[("created_at", DESCENDING), ("_id", DESCENDING)],
  366. )
  367. if not account or (
  368. latest and account.get("last_transaction_id") != latest.get("transaction_id")
  369. ):
  370. account = await _reconcile_account_locked(chat_id, user_id)
  371. balance_before = int(account.get("balance", 0))
  372. actual_delta = (
  373. int(target_balance) - balance_before
  374. if target_balance is not None
  375. else int(delta or 0)
  376. )
  377. balance_after = balance_before + actual_delta
  378. if balance_after < 0:
  379. raise InsufficientPoints("积分余额不足。")
  380. now = utc_now()
  381. transaction = {
  382. "transaction_id": token_hex(10),
  383. "idempotency_key": idempotency_key,
  384. "chat_id": chat_id,
  385. "user_id": user_id,
  386. "username": username,
  387. "first_name": first_name,
  388. "display_name": display_name or first_name,
  389. "delta": actual_delta,
  390. "balance_before": balance_before,
  391. "balance_after": balance_after,
  392. "source": source,
  393. "actor_id": actor_id,
  394. "reason": reason.strip(),
  395. "reference_id": reference_id,
  396. "created_at": now,
  397. }
  398. try:
  399. await transactionsdb.insert_one(transaction)
  400. except DuplicateKeyError:
  401. account = await _reconcile_account_locked(chat_id, user_id)
  402. return account, False
  403. update = {
  404. "$set": {
  405. "balance": balance_after,
  406. "last_transaction_id": transaction["transaction_id"],
  407. "username": username or account.get("username"),
  408. "first_name": first_name or account.get("first_name"),
  409. "display_name": display_name or first_name or account.get("display_name"),
  410. "updated_at": now,
  411. },
  412. "$setOnInsert": {"created_at": now},
  413. "$inc": {
  414. "lifetime_earned": max(actual_delta, 0),
  415. "lifetime_spent": max(-actual_delta, 0),
  416. },
  417. }
  418. await accountsdb.update_one(
  419. {"chat_id": chat_id, "user_id": user_id}, update, upsert=True
  420. )
  421. account = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  422. return account or {"chat_id": chat_id, "user_id": user_id, "balance": balance_after}, True
  423. async def adjust_points(
  424. *,
  425. chat_id: int,
  426. user_id: int,
  427. delta: int,
  428. source: str,
  429. idempotency_key: str,
  430. actor_id: int | str | None = None,
  431. reason: str = "",
  432. reference_id: str | None = None,
  433. username: str | None = None,
  434. first_name: str | None = None,
  435. display_name: str | None = None,
  436. ) -> tuple[dict[str, Any], bool]:
  437. await ensure_point_indexes()
  438. if source not in POINT_SOURCES:
  439. raise PointsError(f"不支持的积分来源:{source}")
  440. if not idempotency_key.strip():
  441. raise PointsError("缺少请求幂等键。")
  442. delta = int(delta)
  443. if delta == 0:
  444. account = await get_point_account(chat_id, user_id)
  445. return account, False
  446. async with _account_locks[(chat_id, user_id)]:
  447. return await _write_points_locked(
  448. chat_id=chat_id,
  449. user_id=user_id,
  450. delta=delta,
  451. source=source,
  452. idempotency_key=idempotency_key,
  453. actor_id=actor_id,
  454. reason=reason,
  455. reference_id=reference_id,
  456. username=username,
  457. first_name=first_name,
  458. display_name=display_name,
  459. )
  460. async def set_points(
  461. *,
  462. chat_id: int,
  463. user_id: int,
  464. balance: int,
  465. actor_id: int | str,
  466. reason: str,
  467. idempotency_key: str,
  468. username: str | None = None,
  469. first_name: str | None = None,
  470. display_name: str | None = None,
  471. ) -> tuple[dict[str, Any], bool]:
  472. await ensure_point_indexes()
  473. if balance < 0:
  474. raise PointsError("积分余额不能小于 0。")
  475. if not idempotency_key.strip():
  476. raise PointsError("缺少请求幂等键。")
  477. async with _account_locks[(chat_id, user_id)]:
  478. return await _write_points_locked(
  479. chat_id=chat_id,
  480. user_id=user_id,
  481. target_balance=balance,
  482. source=SOURCE_ADMIN,
  483. idempotency_key=idempotency_key,
  484. actor_id=actor_id,
  485. reason=reason,
  486. username=username,
  487. first_name=first_name,
  488. display_name=display_name,
  489. )
  490. async def get_point_transaction_by_key(idempotency_key: str) -> dict[str, Any] | None:
  491. await ensure_point_indexes()
  492. return await transactionsdb.find_one({"idempotency_key": idempotency_key})
  493. async def list_point_accounts(
  494. *,
  495. chat_id: int | None = None,
  496. query: str = "",
  497. page: int = 1,
  498. page_size: int = 20,
  499. ) -> tuple[list[dict[str, Any]], int]:
  500. await ensure_point_indexes()
  501. filters: dict[str, Any] = {}
  502. if chat_id is not None:
  503. filters["chat_id"] = chat_id
  504. query = query.strip()
  505. if query:
  506. if query.lstrip("-").isdigit():
  507. filters["user_id"] = int(query)
  508. else:
  509. value = re.escape(query.lstrip("@"))
  510. filters["$or"] = [
  511. {"username": {"$regex": value, "$options": "i"}},
  512. {"first_name": {"$regex": value, "$options": "i"}},
  513. {"display_name": {"$regex": value, "$options": "i"}},
  514. ]
  515. page = max(1, page)
  516. page_size = max(1, min(page_size, 100))
  517. total = await accountsdb.count_documents(filters)
  518. cursor = (
  519. accountsdb.find(filters)
  520. .sort([("balance", DESCENDING), ("user_id", ASCENDING)])
  521. .skip((page - 1) * page_size)
  522. .limit(page_size)
  523. )
  524. return [doc async for doc in cursor], total
  525. async def list_point_transactions(
  526. *,
  527. chat_id: int | None = None,
  528. user_id: int | None = None,
  529. source: str | None = None,
  530. created_from: datetime | None = None,
  531. created_to: datetime | None = None,
  532. page: int = 1,
  533. page_size: int = 20,
  534. ) -> tuple[list[dict[str, Any]], int]:
  535. await ensure_point_indexes()
  536. filters: dict[str, Any] = {}
  537. if chat_id is not None:
  538. filters["chat_id"] = chat_id
  539. if user_id is not None:
  540. filters["user_id"] = user_id
  541. if source:
  542. if source not in POINT_SOURCES:
  543. return [], 0
  544. filters["source"] = source
  545. if created_from or created_to:
  546. created_filter: dict[str, datetime] = {}
  547. if created_from:
  548. created_filter["$gte"] = _as_utc(created_from)
  549. if created_to:
  550. created_filter["$lte"] = _as_utc(created_to)
  551. filters["created_at"] = created_filter
  552. page = max(1, page)
  553. page_size = max(1, min(page_size, 100))
  554. total = await transactionsdb.count_documents(filters)
  555. cursor = (
  556. transactionsdb.find(filters)
  557. .sort([("created_at", DESCENDING), ("_id", DESCENDING)])
  558. .skip((page - 1) * page_size)
  559. .limit(page_size)
  560. )
  561. return [doc async for doc in cursor], total
  562. async def reconcile_all_accounts(
  563. *, chat_id: int | None = None, limit: int = 1000
  564. ) -> dict[str, int]:
  565. """Rebuild cached accounts from the immutable ledger."""
  566. await ensure_point_indexes()
  567. match: dict[str, Any] = {}
  568. if chat_id is not None:
  569. match["chat_id"] = int(chat_id)
  570. pipeline: list[dict[str, Any]] = []
  571. if match:
  572. pipeline.append({"$match": match})
  573. pipeline.extend(
  574. [
  575. {"$group": {"_id": {"chat_id": "$chat_id", "user_id": "$user_id"}}},
  576. {"$limit": max(1, min(int(limit), 10000))},
  577. ]
  578. )
  579. repaired = 0
  580. async for item in transactionsdb.aggregate(pipeline):
  581. await reconcile_account(int(item["_id"]["chat_id"]), int(item["_id"]["user_id"]))
  582. repaired += 1
  583. return {"repaired": repaired}
  584. async def award_checkin(
  585. *,
  586. chat_id: int,
  587. user_id: int,
  588. username: str | None,
  589. first_name: str | None,
  590. display_name: str | None = None,
  591. now: datetime | None = None,
  592. ) -> tuple[dict[str, Any], bool]:
  593. rules = await get_point_rules(chat_id)
  594. if not rules["enabled"] or not rules["checkin_enabled"]:
  595. raise PointsError("本群尚未开启积分签到。")
  596. now = _as_utc(now or utc_now())
  597. day = _local_day(now, rules["timezone"])
  598. return await adjust_points(
  599. chat_id=chat_id,
  600. user_id=user_id,
  601. delta=rules["checkin_points"],
  602. source=SOURCE_CHECKIN,
  603. idempotency_key=f"checkin:{chat_id}:{user_id}:{day}",
  604. reason=f"每日签到({day})",
  605. username=username,
  606. first_name=first_name,
  607. display_name=display_name,
  608. )
  609. async def award_activity(
  610. *,
  611. chat_id: int,
  612. user_id: int,
  613. message_id: int,
  614. content: str,
  615. username: str | None,
  616. first_name: str | None,
  617. display_name: str | None = None,
  618. now: datetime | None = None,
  619. ) -> tuple[dict[str, Any] | None, bool]:
  620. rules = await get_point_rules(chat_id)
  621. if not rules["enabled"] or not rules["activity_enabled"]:
  622. return None, False
  623. points = int(rules["activity_points"])
  624. cap = int(rules["activity_daily_cap"])
  625. if points <= 0 or cap <= 0:
  626. return None, False
  627. now = _as_utc(now or utc_now())
  628. day = _local_day(now, rules["timezone"])
  629. rate_key = f"activity:{chat_id}:{user_id}:{day}"
  630. content_hash = hashlib.sha256(content.strip().lower().encode("utf-8")).hexdigest()
  631. async with _rate_locks[rate_key]:
  632. state = await rate_limitsdb.find_one({"rate_key": rate_key}) or {}
  633. last_awarded = state.get("last_awarded_at")
  634. if last_awarded and (
  635. now - _as_utc(last_awarded)
  636. ).total_seconds() < rules["activity_cooldown_seconds"]:
  637. return None, False
  638. if int(state.get("awarded_points", 0)) >= cap:
  639. return None, False
  640. recent_hashes = [
  641. item
  642. for item in state.get("recent_hashes", [])
  643. if now - _as_utc(item["created_at"]) < timedelta(hours=1)
  644. ]
  645. if any(item["hash"] == content_hash for item in recent_hashes):
  646. return None, False
  647. award = min(points, cap - int(state.get("awarded_points", 0)))
  648. account, created = await adjust_points(
  649. chat_id=chat_id,
  650. user_id=user_id,
  651. delta=award,
  652. source=SOURCE_ACTIVITY,
  653. idempotency_key=f"activity:{chat_id}:{message_id}:{user_id}",
  654. reference_id=str(message_id),
  655. reason="活跃消息奖励",
  656. username=username,
  657. first_name=first_name,
  658. display_name=display_name,
  659. )
  660. if created:
  661. recent_hashes.append({"hash": content_hash, "created_at": now})
  662. await rate_limitsdb.update_one(
  663. {"rate_key": rate_key},
  664. {
  665. "$set": {
  666. "chat_id": chat_id,
  667. "user_id": user_id,
  668. "rule": SOURCE_ACTIVITY,
  669. "last_awarded_at": now,
  670. "recent_hashes": recent_hashes[-20:],
  671. "expires_at": now + timedelta(days=3),
  672. },
  673. "$inc": {"awarded_points": award},
  674. "$setOnInsert": {"created_at": now},
  675. },
  676. upsert=True,
  677. )
  678. return account, created
  679. async def award_upvote(
  680. *,
  681. chat_id: int,
  682. voter_id: int,
  683. target_id: int,
  684. message_id: int,
  685. username: str | None,
  686. first_name: str | None,
  687. display_name: str | None = None,
  688. now: datetime | None = None,
  689. ) -> tuple[dict[str, Any] | None, bool]:
  690. if voter_id == target_id:
  691. return None, False
  692. rules = await get_point_rules(chat_id)
  693. if not rules["enabled"] or not rules["upvote_enabled"]:
  694. return None, False
  695. points = int(rules["upvote_points"])
  696. cap = int(rules["upvote_daily_cap"])
  697. if points <= 0 or cap <= 0:
  698. return None, False
  699. now = _as_utc(now or utc_now())
  700. day = _local_day(now, rules["timezone"])
  701. daily_key = f"upvote-target:{chat_id}:{target_id}:{day}"
  702. pair_key = f"upvote-pair:{chat_id}:{voter_id}:{target_id}"
  703. async with _rate_locks[pair_key]:
  704. async with _rate_locks[daily_key]:
  705. pair = await rate_limitsdb.find_one({"rate_key": pair_key}) or {}
  706. last_awarded = pair.get("last_awarded_at")
  707. if last_awarded and (
  708. now - _as_utc(last_awarded)
  709. ).total_seconds() < rules["upvote_pair_cooldown_seconds"]:
  710. return None, False
  711. daily = await rate_limitsdb.find_one({"rate_key": daily_key}) or {}
  712. already = int(daily.get("awarded_points", 0))
  713. if already >= cap:
  714. return None, False
  715. award = min(points, cap - already)
  716. account, created = await adjust_points(
  717. chat_id=chat_id,
  718. user_id=target_id,
  719. delta=award,
  720. source=SOURCE_UPVOTE,
  721. idempotency_key=f"upvote:{chat_id}:{message_id}:{voter_id}:{target_id}",
  722. actor_id=voter_id,
  723. reference_id=str(message_id),
  724. reason="收到有效点赞",
  725. username=username,
  726. first_name=first_name,
  727. display_name=display_name,
  728. )
  729. if created:
  730. expires_at = now + timedelta(days=35)
  731. await rate_limitsdb.update_one(
  732. {"rate_key": pair_key},
  733. {
  734. "$set": {
  735. "chat_id": chat_id,
  736. "user_id": target_id,
  737. "voter_id": voter_id,
  738. "rule": "upvote_pair",
  739. "last_awarded_at": now,
  740. "expires_at": expires_at,
  741. },
  742. "$setOnInsert": {"created_at": now},
  743. },
  744. upsert=True,
  745. )
  746. await rate_limitsdb.update_one(
  747. {"rate_key": daily_key},
  748. {
  749. "$set": {
  750. "chat_id": chat_id,
  751. "user_id": target_id,
  752. "rule": SOURCE_UPVOTE,
  753. "last_awarded_at": now,
  754. "expires_at": now + timedelta(days=3),
  755. },
  756. "$inc": {"awarded_points": award},
  757. "$setOnInsert": {"created_at": now},
  758. },
  759. upsert=True,
  760. )
  761. return account, created