dbpoints.py 25 KB

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