dbpoints.py 26 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810
  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 list_checkin_keyboard_migrations(
  180. keyboard_version: int,
  181. ) -> list[int]:
  182. await ensure_point_indexes()
  183. cursor = rulesdb.find(
  184. {
  185. "enabled": True,
  186. "checkin_enabled": True,
  187. "checkin_button_enabled": True,
  188. "$or": [
  189. {"checkin_keyboard_version": {"$exists": False}},
  190. {"checkin_keyboard_version": {"$lt": int(keyboard_version)}},
  191. ],
  192. },
  193. {"chat_id": 1},
  194. )
  195. documents = await cursor.to_list(length=10000)
  196. return [int(document["chat_id"]) for document in documents]
  197. async def mark_checkin_keyboard_version(
  198. chat_id: int,
  199. keyboard_version: int,
  200. ) -> None:
  201. await ensure_point_indexes()
  202. await rulesdb.update_one(
  203. {"chat_id": int(chat_id)},
  204. {
  205. "$set": {
  206. "checkin_keyboard_version": int(keyboard_version),
  207. "updated_at": utc_now(),
  208. }
  209. },
  210. )
  211. async def _reconcile_account_locked(chat_id: int, user_id: int) -> dict[str, Any]:
  212. cached = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  213. pipeline = [
  214. {"$match": {"chat_id": chat_id, "user_id": user_id}},
  215. {
  216. "$group": {
  217. "_id": None,
  218. "balance": {"$sum": "$delta"},
  219. "lifetime_earned": {
  220. "$sum": {"$cond": [{"$gt": ["$delta", 0]}, "$delta", 0]}
  221. },
  222. "lifetime_spent": {
  223. "$sum": {
  224. "$cond": [
  225. {"$lt": ["$delta", 0]},
  226. {"$multiply": ["$delta", -1]},
  227. 0,
  228. ]
  229. }
  230. },
  231. }
  232. },
  233. ]
  234. totals = [doc async for doc in transactionsdb.aggregate(pipeline)]
  235. latest = await transactionsdb.find_one(
  236. {"chat_id": chat_id, "user_id": user_id},
  237. sort=[("created_at", DESCENDING), ("_id", DESCENDING)],
  238. )
  239. total = totals[0] if totals else {}
  240. now = utc_now()
  241. account = {
  242. "chat_id": chat_id,
  243. "user_id": user_id,
  244. "balance": int(total.get("balance", 0)),
  245. "lifetime_earned": int(total.get("lifetime_earned", 0)),
  246. "lifetime_spent": int(total.get("lifetime_spent", 0)),
  247. "last_transaction_id": latest.get("transaction_id") if latest else None,
  248. "updated_at": now,
  249. }
  250. if latest:
  251. account["username"] = (
  252. cached.get("username") if cached and "username" in cached else latest.get("username")
  253. )
  254. account["first_name"] = (
  255. cached.get("first_name")
  256. if cached and "first_name" in cached
  257. else latest.get("first_name")
  258. )
  259. account["display_name"] = (
  260. cached.get("display_name")
  261. if cached and "display_name" in cached
  262. else latest.get("display_name") or latest.get("first_name")
  263. )
  264. await accountsdb.update_one(
  265. {"chat_id": chat_id, "user_id": user_id},
  266. {"$set": account, "$setOnInsert": {"created_at": now}},
  267. upsert=True,
  268. )
  269. return account
  270. async def reconcile_account(chat_id: int, user_id: int) -> dict[str, Any]:
  271. await ensure_point_indexes()
  272. async with _account_locks[(chat_id, user_id)]:
  273. return await _reconcile_account_locked(chat_id, user_id)
  274. async def get_point_account(chat_id: int, user_id: int) -> dict[str, Any]:
  275. await ensure_point_indexes()
  276. account = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  277. latest = await transactionsdb.find_one(
  278. {"chat_id": chat_id, "user_id": user_id},
  279. projection={"transaction_id": 1},
  280. sort=[("created_at", DESCENDING), ("_id", DESCENDING)],
  281. )
  282. if latest and (
  283. not account or account.get("last_transaction_id") != latest.get("transaction_id")
  284. ):
  285. return await reconcile_account(chat_id, user_id)
  286. if account:
  287. return account
  288. return {
  289. "chat_id": chat_id,
  290. "user_id": user_id,
  291. "balance": 0,
  292. "lifetime_earned": 0,
  293. "lifetime_spent": 0,
  294. "last_transaction_id": None,
  295. }
  296. async def update_point_account_identity(
  297. *,
  298. chat_id: int,
  299. user_id: int,
  300. username: str | None,
  301. first_name: str | None,
  302. last_name: str | None,
  303. display_name: str,
  304. ) -> None:
  305. await ensure_point_indexes()
  306. await accountsdb.update_one(
  307. {"chat_id": int(chat_id), "user_id": int(user_id)},
  308. {
  309. "$set": {
  310. "username": username,
  311. "first_name": first_name,
  312. "last_name": last_name,
  313. "display_name": display_name or first_name,
  314. "identity_updated_at": utc_now(),
  315. }
  316. },
  317. )
  318. async def _write_points_locked(
  319. *,
  320. chat_id: int,
  321. user_id: int,
  322. source: str,
  323. idempotency_key: str,
  324. delta: int | None = None,
  325. target_balance: int | None = None,
  326. actor_id: int | str | None = None,
  327. reason: str = "",
  328. reference_id: str | None = None,
  329. username: str | None = None,
  330. first_name: str | None = None,
  331. display_name: str | None = None,
  332. ) -> tuple[dict[str, Any], bool]:
  333. existing = await transactionsdb.find_one({"idempotency_key": idempotency_key})
  334. if existing:
  335. account = await _reconcile_account_locked(chat_id, user_id)
  336. return account, False
  337. account = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  338. latest = await transactionsdb.find_one(
  339. {"chat_id": chat_id, "user_id": user_id},
  340. projection={"transaction_id": 1},
  341. sort=[("created_at", DESCENDING), ("_id", DESCENDING)],
  342. )
  343. if not account or (
  344. latest and account.get("last_transaction_id") != latest.get("transaction_id")
  345. ):
  346. account = await _reconcile_account_locked(chat_id, user_id)
  347. balance_before = int(account.get("balance", 0))
  348. actual_delta = (
  349. int(target_balance) - balance_before
  350. if target_balance is not None
  351. else int(delta or 0)
  352. )
  353. balance_after = balance_before + actual_delta
  354. if balance_after < 0:
  355. raise InsufficientPoints("积分余额不足。")
  356. now = utc_now()
  357. transaction = {
  358. "transaction_id": token_hex(10),
  359. "idempotency_key": idempotency_key,
  360. "chat_id": chat_id,
  361. "user_id": user_id,
  362. "username": username,
  363. "first_name": first_name,
  364. "display_name": display_name or first_name,
  365. "delta": actual_delta,
  366. "balance_before": balance_before,
  367. "balance_after": balance_after,
  368. "source": source,
  369. "actor_id": actor_id,
  370. "reason": reason.strip(),
  371. "reference_id": reference_id,
  372. "created_at": now,
  373. }
  374. try:
  375. await transactionsdb.insert_one(transaction)
  376. except DuplicateKeyError:
  377. account = await _reconcile_account_locked(chat_id, user_id)
  378. return account, False
  379. update = {
  380. "$set": {
  381. "balance": balance_after,
  382. "last_transaction_id": transaction["transaction_id"],
  383. "username": username or account.get("username"),
  384. "first_name": first_name or account.get("first_name"),
  385. "display_name": display_name or first_name or account.get("display_name"),
  386. "updated_at": now,
  387. },
  388. "$setOnInsert": {"created_at": now},
  389. "$inc": {
  390. "lifetime_earned": max(actual_delta, 0),
  391. "lifetime_spent": max(-actual_delta, 0),
  392. },
  393. }
  394. await accountsdb.update_one(
  395. {"chat_id": chat_id, "user_id": user_id}, update, upsert=True
  396. )
  397. account = await accountsdb.find_one({"chat_id": chat_id, "user_id": user_id})
  398. return account or {"chat_id": chat_id, "user_id": user_id, "balance": balance_after}, True
  399. async def adjust_points(
  400. *,
  401. chat_id: int,
  402. user_id: int,
  403. delta: int,
  404. source: str,
  405. idempotency_key: str,
  406. actor_id: int | str | None = None,
  407. reason: str = "",
  408. reference_id: str | None = None,
  409. username: str | None = None,
  410. first_name: str | None = None,
  411. display_name: str | None = None,
  412. ) -> tuple[dict[str, Any], bool]:
  413. await ensure_point_indexes()
  414. if source not in POINT_SOURCES:
  415. raise PointsError(f"不支持的积分来源:{source}")
  416. if not idempotency_key.strip():
  417. raise PointsError("缺少请求幂等键。")
  418. delta = int(delta)
  419. if delta == 0:
  420. account = await get_point_account(chat_id, user_id)
  421. return account, False
  422. async with _account_locks[(chat_id, user_id)]:
  423. return await _write_points_locked(
  424. chat_id=chat_id,
  425. user_id=user_id,
  426. delta=delta,
  427. source=source,
  428. idempotency_key=idempotency_key,
  429. actor_id=actor_id,
  430. reason=reason,
  431. reference_id=reference_id,
  432. username=username,
  433. first_name=first_name,
  434. display_name=display_name,
  435. )
  436. async def set_points(
  437. *,
  438. chat_id: int,
  439. user_id: int,
  440. balance: int,
  441. actor_id: int | str,
  442. reason: str,
  443. idempotency_key: str,
  444. username: str | None = None,
  445. first_name: str | None = None,
  446. display_name: str | None = None,
  447. ) -> tuple[dict[str, Any], bool]:
  448. await ensure_point_indexes()
  449. if balance < 0:
  450. raise PointsError("积分余额不能小于 0。")
  451. if not idempotency_key.strip():
  452. raise PointsError("缺少请求幂等键。")
  453. async with _account_locks[(chat_id, user_id)]:
  454. return await _write_points_locked(
  455. chat_id=chat_id,
  456. user_id=user_id,
  457. target_balance=balance,
  458. source=SOURCE_ADMIN,
  459. idempotency_key=idempotency_key,
  460. actor_id=actor_id,
  461. reason=reason,
  462. username=username,
  463. first_name=first_name,
  464. display_name=display_name,
  465. )
  466. async def get_point_transaction_by_key(idempotency_key: str) -> dict[str, Any] | None:
  467. await ensure_point_indexes()
  468. return await transactionsdb.find_one({"idempotency_key": idempotency_key})
  469. async def list_point_accounts(
  470. *,
  471. chat_id: int | None = None,
  472. query: str = "",
  473. page: int = 1,
  474. page_size: int = 20,
  475. ) -> tuple[list[dict[str, Any]], int]:
  476. await ensure_point_indexes()
  477. filters: dict[str, Any] = {}
  478. if chat_id is not None:
  479. filters["chat_id"] = chat_id
  480. query = query.strip()
  481. if query:
  482. if query.lstrip("-").isdigit():
  483. filters["user_id"] = int(query)
  484. else:
  485. value = re.escape(query.lstrip("@"))
  486. filters["$or"] = [
  487. {"username": {"$regex": value, "$options": "i"}},
  488. {"first_name": {"$regex": value, "$options": "i"}},
  489. {"display_name": {"$regex": value, "$options": "i"}},
  490. ]
  491. page = max(1, page)
  492. page_size = max(1, min(page_size, 100))
  493. total = await accountsdb.count_documents(filters)
  494. cursor = (
  495. accountsdb.find(filters)
  496. .sort([("balance", DESCENDING), ("user_id", ASCENDING)])
  497. .skip((page - 1) * page_size)
  498. .limit(page_size)
  499. )
  500. return [doc async for doc in cursor], total
  501. async def list_point_transactions(
  502. *,
  503. chat_id: int | None = None,
  504. user_id: int | None = None,
  505. source: str | None = None,
  506. created_from: datetime | None = None,
  507. created_to: datetime | None = None,
  508. page: int = 1,
  509. page_size: int = 20,
  510. ) -> tuple[list[dict[str, Any]], int]:
  511. await ensure_point_indexes()
  512. filters: dict[str, Any] = {}
  513. if chat_id is not None:
  514. filters["chat_id"] = chat_id
  515. if user_id is not None:
  516. filters["user_id"] = user_id
  517. if source:
  518. if source not in POINT_SOURCES:
  519. return [], 0
  520. filters["source"] = source
  521. if created_from or created_to:
  522. created_filter: dict[str, datetime] = {}
  523. if created_from:
  524. created_filter["$gte"] = _as_utc(created_from)
  525. if created_to:
  526. created_filter["$lte"] = _as_utc(created_to)
  527. filters["created_at"] = created_filter
  528. page = max(1, page)
  529. page_size = max(1, min(page_size, 100))
  530. total = await transactionsdb.count_documents(filters)
  531. cursor = (
  532. transactionsdb.find(filters)
  533. .sort([("created_at", DESCENDING), ("_id", DESCENDING)])
  534. .skip((page - 1) * page_size)
  535. .limit(page_size)
  536. )
  537. return [doc async for doc in cursor], total
  538. async def reconcile_all_accounts(
  539. *, chat_id: int | None = None, limit: int = 1000
  540. ) -> dict[str, int]:
  541. """Rebuild cached accounts from the immutable ledger."""
  542. await ensure_point_indexes()
  543. match: dict[str, Any] = {}
  544. if chat_id is not None:
  545. match["chat_id"] = int(chat_id)
  546. pipeline: list[dict[str, Any]] = []
  547. if match:
  548. pipeline.append({"$match": match})
  549. pipeline.extend(
  550. [
  551. {"$group": {"_id": {"chat_id": "$chat_id", "user_id": "$user_id"}}},
  552. {"$limit": max(1, min(int(limit), 10000))},
  553. ]
  554. )
  555. repaired = 0
  556. async for item in transactionsdb.aggregate(pipeline):
  557. await reconcile_account(int(item["_id"]["chat_id"]), int(item["_id"]["user_id"]))
  558. repaired += 1
  559. return {"repaired": repaired}
  560. async def award_checkin(
  561. *,
  562. chat_id: int,
  563. user_id: int,
  564. username: str | None,
  565. first_name: str | None,
  566. display_name: str | None = None,
  567. now: datetime | None = None,
  568. ) -> tuple[dict[str, Any], bool]:
  569. rules = await get_point_rules(chat_id)
  570. if not rules["enabled"] or not rules["checkin_enabled"]:
  571. raise PointsError("本群尚未开启积分签到。")
  572. now = _as_utc(now or utc_now())
  573. day = _local_day(now, rules["timezone"])
  574. return await adjust_points(
  575. chat_id=chat_id,
  576. user_id=user_id,
  577. delta=rules["checkin_points"],
  578. source=SOURCE_CHECKIN,
  579. idempotency_key=f"checkin:{chat_id}:{user_id}:{day}",
  580. reason=f"每日签到({day})",
  581. username=username,
  582. first_name=first_name,
  583. display_name=display_name,
  584. )
  585. async def award_activity(
  586. *,
  587. chat_id: int,
  588. user_id: int,
  589. message_id: int,
  590. content: str,
  591. username: str | None,
  592. first_name: str | None,
  593. display_name: str | None = None,
  594. now: datetime | None = None,
  595. ) -> tuple[dict[str, Any] | None, bool]:
  596. rules = await get_point_rules(chat_id)
  597. if not rules["enabled"] or not rules["activity_enabled"]:
  598. return None, False
  599. points = int(rules["activity_points"])
  600. cap = int(rules["activity_daily_cap"])
  601. if points <= 0 or cap <= 0:
  602. return None, False
  603. now = _as_utc(now or utc_now())
  604. day = _local_day(now, rules["timezone"])
  605. rate_key = f"activity:{chat_id}:{user_id}:{day}"
  606. content_hash = hashlib.sha256(content.strip().lower().encode("utf-8")).hexdigest()
  607. async with _rate_locks[rate_key]:
  608. state = await rate_limitsdb.find_one({"rate_key": rate_key}) or {}
  609. last_awarded = state.get("last_awarded_at")
  610. if last_awarded and (
  611. now - _as_utc(last_awarded)
  612. ).total_seconds() < rules["activity_cooldown_seconds"]:
  613. return None, False
  614. if int(state.get("awarded_points", 0)) >= cap:
  615. return None, False
  616. recent_hashes = [
  617. item
  618. for item in state.get("recent_hashes", [])
  619. if now - _as_utc(item["created_at"]) < timedelta(hours=1)
  620. ]
  621. if any(item["hash"] == content_hash for item in recent_hashes):
  622. return None, False
  623. award = min(points, cap - int(state.get("awarded_points", 0)))
  624. account, created = await adjust_points(
  625. chat_id=chat_id,
  626. user_id=user_id,
  627. delta=award,
  628. source=SOURCE_ACTIVITY,
  629. idempotency_key=f"activity:{chat_id}:{message_id}:{user_id}",
  630. reference_id=str(message_id),
  631. reason="活跃消息奖励",
  632. username=username,
  633. first_name=first_name,
  634. display_name=display_name,
  635. )
  636. if created:
  637. recent_hashes.append({"hash": content_hash, "created_at": now})
  638. await rate_limitsdb.update_one(
  639. {"rate_key": rate_key},
  640. {
  641. "$set": {
  642. "chat_id": chat_id,
  643. "user_id": user_id,
  644. "rule": SOURCE_ACTIVITY,
  645. "last_awarded_at": now,
  646. "recent_hashes": recent_hashes[-20:],
  647. "expires_at": now + timedelta(days=3),
  648. },
  649. "$inc": {"awarded_points": award},
  650. "$setOnInsert": {"created_at": now},
  651. },
  652. upsert=True,
  653. )
  654. return account, created
  655. async def award_upvote(
  656. *,
  657. chat_id: int,
  658. voter_id: int,
  659. target_id: int,
  660. message_id: int,
  661. username: str | None,
  662. first_name: str | None,
  663. display_name: str | None = None,
  664. now: datetime | None = None,
  665. ) -> tuple[dict[str, Any] | None, bool]:
  666. if voter_id == target_id:
  667. return None, False
  668. rules = await get_point_rules(chat_id)
  669. if not rules["enabled"] or not rules["upvote_enabled"]:
  670. return None, False
  671. points = int(rules["upvote_points"])
  672. cap = int(rules["upvote_daily_cap"])
  673. if points <= 0 or cap <= 0:
  674. return None, False
  675. now = _as_utc(now or utc_now())
  676. day = _local_day(now, rules["timezone"])
  677. daily_key = f"upvote-target:{chat_id}:{target_id}:{day}"
  678. pair_key = f"upvote-pair:{chat_id}:{voter_id}:{target_id}"
  679. async with _rate_locks[pair_key]:
  680. async with _rate_locks[daily_key]:
  681. pair = await rate_limitsdb.find_one({"rate_key": pair_key}) or {}
  682. last_awarded = pair.get("last_awarded_at")
  683. if last_awarded and (
  684. now - _as_utc(last_awarded)
  685. ).total_seconds() < rules["upvote_pair_cooldown_seconds"]:
  686. return None, False
  687. daily = await rate_limitsdb.find_one({"rate_key": daily_key}) or {}
  688. already = int(daily.get("awarded_points", 0))
  689. if already >= cap:
  690. return None, False
  691. award = min(points, cap - already)
  692. account, created = await adjust_points(
  693. chat_id=chat_id,
  694. user_id=target_id,
  695. delta=award,
  696. source=SOURCE_UPVOTE,
  697. idempotency_key=f"upvote:{chat_id}:{message_id}:{voter_id}:{target_id}",
  698. actor_id=voter_id,
  699. reference_id=str(message_id),
  700. reason="收到有效点赞",
  701. username=username,
  702. first_name=first_name,
  703. display_name=display_name,
  704. )
  705. if created:
  706. expires_at = now + timedelta(days=35)
  707. await rate_limitsdb.update_one(
  708. {"rate_key": pair_key},
  709. {
  710. "$set": {
  711. "chat_id": chat_id,
  712. "user_id": target_id,
  713. "voter_id": voter_id,
  714. "rule": "upvote_pair",
  715. "last_awarded_at": now,
  716. "expires_at": expires_at,
  717. },
  718. "$setOnInsert": {"created_at": now},
  719. },
  720. upsert=True,
  721. )
  722. await rate_limitsdb.update_one(
  723. {"rate_key": daily_key},
  724. {
  725. "$set": {
  726. "chat_id": chat_id,
  727. "user_id": target_id,
  728. "rule": SOURCE_UPVOTE,
  729. "last_awarded_at": now,
  730. "expires_at": now + timedelta(days=3),
  731. },
  732. "$inc": {"awarded_points": award},
  733. "$setOnInsert": {"created_at": now},
  734. },
  735. upsert=True,
  736. )
  737. return account, created