api.py 57 KB

123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380
  1. from __future__ import annotations
  2. import csv
  3. import io
  4. import tempfile
  5. import traceback
  6. from contextlib import suppress
  7. from datetime import UTC, datetime, timedelta
  8. from pathlib import Path
  9. from typing import Any
  10. from aiohttp import ClientError, ClientSession, ClientTimeout, web
  11. import wbb
  12. from wbb import MESSAGE_DUMP_CHAT
  13. from wbb import app as telegram_app
  14. from wbb.admin.bot_config import (
  15. BotConfigError,
  16. create_bot_profile,
  17. create_bot_role,
  18. delete_bot_profile,
  19. delete_bot_role,
  20. get_bot_profile_secrets,
  21. store_bot_identity,
  22. telegram_config_status,
  23. test_bot_token,
  24. update_bot_profile,
  25. update_bot_role,
  26. update_telegram_config,
  27. )
  28. from wbb.admin.security import (
  29. generate_csrf_token,
  30. generate_session_token,
  31. hash_password,
  32. hash_token,
  33. validate_new_password,
  34. verify_password,
  35. )
  36. from wbb.services.bot_permissions import api_permission, has_permission
  37. from wbb.services.chat_management import (
  38. ChatManagementError,
  39. apply_automation_settings,
  40. create_invite_link,
  41. ensure_permission,
  42. execute_member_action,
  43. get_automation_settings,
  44. get_chat_overview,
  45. get_invite_links,
  46. list_accessible_chats,
  47. list_chat_admins,
  48. list_recent_members,
  49. revoke_invite_link,
  50. search_chat_members,
  51. send_announcement,
  52. update_chat_permissions,
  53. update_chat_profile,
  54. )
  55. from wbb.services.giveaways import (
  56. GiveawayServiceError,
  57. cancel_and_refund_giveaway,
  58. create_and_publish_giveaway,
  59. finish_and_publish_giveaway,
  60. remove_and_optionally_refund_participant,
  61. reroll_giveaway,
  62. )
  63. from wbb.services.point_settings import apply_point_rules
  64. from wbb.utils.dbadmin import (
  65. create_admin_session,
  66. dashboard_counts,
  67. ensure_default_admin,
  68. get_admin_session,
  69. get_admin_user,
  70. list_audit_logs,
  71. list_member_identity_changes,
  72. record_audit,
  73. record_login_failure,
  74. record_login_success,
  75. revoke_admin_session,
  76. update_admin_password,
  77. )
  78. from wbb.utils.dbfunctions import get_rules, set_chat_rules
  79. from wbb.utils.dbgiveaway import (
  80. add_giveaway_ban,
  81. get_giveaway,
  82. list_giveaway_bans,
  83. list_giveaways_page,
  84. list_participants_page,
  85. remove_giveaway_ban,
  86. )
  87. from wbb.utils.dbpoints import (
  88. SOURCE_ADMIN,
  89. InsufficientPoints,
  90. PointsError,
  91. adjust_points,
  92. get_point_rules,
  93. list_point_accounts,
  94. list_point_transactions,
  95. set_points,
  96. )
  97. API_PREFIX = "/api/admin/v1"
  98. SESSION_COOKIE = "wbb_admin_session"
  99. UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"}
  100. PUBLIC_API_PATHS = {f"{API_PREFIX}/auth/login", f"{API_PREFIX}/health"}
  101. BOT_SCOPED_PREFIXES = (
  102. f"{API_PREFIX}/chats",
  103. f"{API_PREFIX}/giveaways",
  104. f"{API_PREFIX}/points",
  105. f"{API_PREFIX}/media",
  106. )
  107. TELEGRAM_ID_KEYS = {
  108. "chat_id",
  109. "user_id",
  110. "actor_id",
  111. "creator_id",
  112. "moderator_id",
  113. "target_id",
  114. "message_id",
  115. "removed_by",
  116. }
  117. class ApiProblem(RuntimeError):
  118. def __init__(
  119. self,
  120. code: str,
  121. message: str,
  122. *,
  123. status: int = 400,
  124. details: Any = None,
  125. ):
  126. super().__init__(message)
  127. self.code = code
  128. self.status = status
  129. self.details = details
  130. def as_utc(value: datetime) -> datetime:
  131. if value.tzinfo is None:
  132. return value.replace(tzinfo=UTC)
  133. return value.astimezone(UTC)
  134. def jsonable(value: Any, *, key: str = "") -> Any:
  135. if isinstance(value, datetime):
  136. return as_utc(value).isoformat().replace("+00:00", "Z")
  137. if isinstance(value, dict):
  138. return {
  139. item_key: jsonable(item_value, key=item_key)
  140. for item_key, item_value in value.items()
  141. if item_key != "_id"
  142. }
  143. if isinstance(value, (list, tuple)):
  144. return [jsonable(item) for item in value]
  145. if isinstance(value, int) and (key in TELEGRAM_ID_KEYS or key.endswith("_telegram_id")):
  146. return str(value)
  147. return value
  148. def success(data: Any, *, status: int = 200) -> web.Response:
  149. return web.json_response({"data": jsonable(data)}, status=status)
  150. def error_response(problem: ApiProblem) -> web.Response:
  151. return web.json_response(
  152. {
  153. "error": {
  154. "code": problem.code,
  155. "message": str(problem),
  156. "details": jsonable(problem.details),
  157. }
  158. },
  159. status=problem.status,
  160. )
  161. def parse_int(value: Any, name: str, *, minimum: int | None = None) -> int:
  162. try:
  163. parsed = int(value)
  164. except (TypeError, ValueError) as exc:
  165. raise ApiProblem("invalid_parameter", f"{name} 必须是整数。") from exc
  166. if minimum is not None and parsed < minimum:
  167. raise ApiProblem("invalid_parameter", f"{name} 不能小于 {minimum}。")
  168. return parsed
  169. def parse_datetime(value: Any, name: str) -> datetime:
  170. if not isinstance(value, str) or not value.strip():
  171. raise ApiProblem("invalid_parameter", f"{name} 必须是 UTC ISO 8601 时间。")
  172. try:
  173. parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
  174. except ValueError as exc:
  175. raise ApiProblem("invalid_parameter", f"{name} 不是有效时间。") from exc
  176. if parsed.tzinfo is None:
  177. raise ApiProblem("invalid_parameter", f"{name} 必须包含时区。")
  178. return parsed.astimezone(UTC)
  179. async def json_body(request: web.Request) -> dict[str, Any]:
  180. try:
  181. body = await request.json()
  182. except Exception as exc:
  183. raise ApiProblem("invalid_json", "请求体必须是 JSON。") from exc
  184. if not isinstance(body, dict):
  185. raise ApiProblem("invalid_json", "请求体必须是 JSON 对象。")
  186. return body
  187. def require_confirmation(body: dict[str, Any]) -> None:
  188. if body.get("confirm") is not True:
  189. raise ApiProblem(
  190. "confirmation_required", "该操作需要二次确认。", status=409
  191. )
  192. def page_params(request: web.Request) -> tuple[int, int]:
  193. return (
  194. parse_int(request.query.get("page", 1), "page", minimum=1),
  195. min(parse_int(request.query.get("page_size", 20), "page_size", minimum=1), 100),
  196. )
  197. def chat_id_param(request: web.Request) -> int:
  198. return parse_int(request.match_info["chat_id"], "chatId")
  199. def set_audit(
  200. request: web.Request,
  201. action: str,
  202. *,
  203. chat_id: int | None = None,
  204. target_id: int | str | None = None,
  205. summary: str = "",
  206. metadata: dict[str, Any] | None = None,
  207. ) -> None:
  208. request["audit_context"] = {
  209. "action": action,
  210. "chat_id": chat_id,
  211. "target_id": target_id,
  212. "summary": summary,
  213. "metadata": metadata or {},
  214. }
  215. @web.middleware
  216. async def api_error_middleware(request: web.Request, handler):
  217. try:
  218. response = await handler(request)
  219. except ApiProblem as exc:
  220. await _record_request_audit(request, success_state=False, error=str(exc))
  221. return error_response(exc)
  222. except (ChatManagementError, GiveawayServiceError) as exc:
  223. problem = ApiProblem(
  224. exc.code, str(exc), status=getattr(exc, "status", 400)
  225. )
  226. await _record_request_audit(request, success_state=False, error=str(exc))
  227. return error_response(problem)
  228. except (PointsError, InsufficientPoints) as exc:
  229. problem = ApiProblem("points_error", str(exc), status=409)
  230. await _record_request_audit(request, success_state=False, error=str(exc))
  231. return error_response(problem)
  232. except BotConfigError as exc:
  233. status = {
  234. "bot_not_found": 404,
  235. "role_not_found": 404,
  236. "role_in_use": 409,
  237. "builtin_role_immutable": 409,
  238. }.get(exc.code, 400)
  239. problem = ApiProblem(exc.code, str(exc), status=status)
  240. await _record_request_audit(request, success_state=False, error=str(exc))
  241. return error_response(problem)
  242. except web.HTTPException:
  243. raise
  244. except Exception as exc:
  245. await _record_request_audit(request, success_state=False, error=str(exc))
  246. wbb.log.error(
  247. f"Admin API {request.method} {request.path} failed: {exc}\n"
  248. f"{traceback.format_exc()}"
  249. )
  250. return error_response(
  251. ApiProblem("internal_error", "服务器处理请求失败。", status=500)
  252. )
  253. await _record_request_audit(request, success_state=response.status < 400)
  254. return response
  255. async def _record_request_audit(
  256. request: web.Request, *, success_state: bool, error: str = ""
  257. ) -> None:
  258. context = request.get("audit_context")
  259. if not context or request.get("audit_recorded"):
  260. return
  261. request["audit_recorded"] = True
  262. admin = request.get("admin") or {}
  263. with suppress(Exception):
  264. await record_audit(
  265. source="web",
  266. actor_id=admin.get("username", "anonymous"),
  267. actor_name=admin.get("username", "anonymous"),
  268. success=success_state,
  269. error=error,
  270. **context,
  271. )
  272. @web.middleware
  273. async def authentication_middleware(request: web.Request, handler):
  274. if not request.path.startswith(API_PREFIX) or request.path in PUBLIC_API_PATHS:
  275. return await handler(request)
  276. raw_token = request.cookies.get(SESSION_COOKIE, "")
  277. session = await get_admin_session(hash_token(raw_token)) if raw_token else None
  278. if not session:
  279. raise ApiProblem("unauthenticated", "登录已失效,请重新登录。", status=401)
  280. user = session["user"]
  281. request["admin"] = {
  282. "username": user["username"],
  283. "must_change_password": bool(user.get("must_change_password")),
  284. "csrf_token": session["csrf_token"],
  285. "session_token_hash": session["token_hash"],
  286. }
  287. allowed_during_password_change = {
  288. f"{API_PREFIX}/auth/me",
  289. f"{API_PREFIX}/auth/password",
  290. f"{API_PREFIX}/auth/logout",
  291. }
  292. if user.get("must_change_password") and request.path not in allowed_during_password_change:
  293. raise ApiProblem(
  294. "password_change_required", "首次登录必须修改初始密码。", status=428
  295. )
  296. if request.method in UNSAFE_METHODS:
  297. supplied = request.headers.get("X-CSRF-Token", "")
  298. if not supplied or supplied != session["csrf_token"]:
  299. raise ApiProblem("csrf_failed", "CSRF 校验失败。", status=403)
  300. return await handler(request)
  301. def _selected_bot_id(request: web.Request) -> str:
  302. bot_id = str(request.headers.get("X-Bot-Id") or "").strip()
  303. if bot_id:
  304. return bot_id
  305. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  306. if supervisor is not None:
  307. running = [
  308. key
  309. for key, value in supervisor.runtimes().items()
  310. if value.get("state") == "running"
  311. ]
  312. if len(running) == 1:
  313. return running[0]
  314. raise ApiProblem(
  315. "bot_selection_required",
  316. "请先选择要管理的机器人。",
  317. status=409,
  318. )
  319. @web.middleware
  320. async def bot_role_middleware(request: web.Request, handler):
  321. required = api_permission(request.method, request.path)
  322. if not required:
  323. return await handler(request)
  324. if bool(getattr(wbb, "SUPERVISOR_MODE", False)):
  325. profile = get_bot_profile_secrets(
  326. getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"),
  327. _selected_bot_id(request),
  328. )
  329. permissions = profile.get("permissions", [])
  330. else:
  331. permissions = getattr(wbb, "BOT_PERMISSIONS", {"*"})
  332. if not has_permission(required, permissions):
  333. raise ApiProblem(
  334. "bot_role_permission_denied",
  335. "所选机器人的职责角色不允许执行该操作。",
  336. status=403,
  337. details={"required_permission": required},
  338. )
  339. return await handler(request)
  340. @web.middleware
  341. async def bot_proxy_middleware(request: web.Request, handler):
  342. if not bool(getattr(wbb, "SUPERVISOR_MODE", False)) or not request.path.startswith(
  343. BOT_SCOPED_PREFIXES
  344. ):
  345. return await handler(request)
  346. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  347. if supervisor is None:
  348. raise ApiProblem(
  349. "bot_supervisor_unavailable",
  350. "机器人监管服务尚未就绪。",
  351. status=503,
  352. )
  353. bot_id = _selected_bot_id(request)
  354. endpoint = supervisor.endpoint_for(bot_id)
  355. if endpoint is None:
  356. raise ApiProblem(
  357. "bot_not_running",
  358. "所选机器人当前未连接 Telegram。",
  359. status=409,
  360. )
  361. target = f"{endpoint}{request.rel_url}"
  362. headers = {
  363. key: value
  364. for key, value in request.headers.items()
  365. if key.lower() not in {"host", "content-length", "connection"}
  366. }
  367. body = await request.read()
  368. try:
  369. async with ClientSession(timeout=ClientTimeout(total=90)) as session:
  370. async with session.request(
  371. request.method,
  372. target,
  373. data=body or None,
  374. headers=headers,
  375. allow_redirects=False,
  376. ) as response:
  377. response_body = await response.read()
  378. response_headers = {
  379. key: value
  380. for key, value in response.headers.items()
  381. if key.lower() in {"content-type", "content-disposition"}
  382. }
  383. return web.Response(
  384. body=response_body,
  385. status=response.status,
  386. headers=response_headers,
  387. )
  388. except (ClientError, TimeoutError) as exc:
  389. raise ApiProblem(
  390. "bot_worker_unavailable",
  391. "机器人工作进程暂时不可用。",
  392. status=503,
  393. ) from exc
  394. class AdminApi:
  395. def __init__(self) -> None:
  396. self.username = str(getattr(wbb, "ADMIN_WEB_USERNAME", "admin"))
  397. self.initial_password = str(
  398. getattr(wbb, "ADMIN_WEB_INITIAL_PASSWORD", "qwe0.123456")
  399. )
  400. self.session_hours = int(getattr(wbb, "ADMIN_WEB_SESSION_HOURS", 12))
  401. self.cookie_secure = bool(getattr(wbb, "ADMIN_WEB_COOKIE_SECURE", False))
  402. self.upload_limit = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)) * 1024 * 1024
  403. self.bot_config_path = Path(
  404. str(getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"))
  405. )
  406. async def initialize(self) -> None:
  407. await ensure_default_admin(
  408. username=self.username,
  409. password_hash=hash_password(self.initial_password),
  410. )
  411. def register(self, application: web.Application) -> None:
  412. router = application.router
  413. router.add_get(f"{API_PREFIX}/health", self.health)
  414. router.add_post(f"{API_PREFIX}/auth/login", self.login)
  415. router.add_get(f"{API_PREFIX}/auth/me", self.me)
  416. router.add_post(f"{API_PREFIX}/auth/logout", self.logout)
  417. router.add_put(f"{API_PREFIX}/auth/password", self.change_password)
  418. router.add_get(f"{API_PREFIX}/dashboard", self.dashboard)
  419. router.add_get(f"{API_PREFIX}/settings", self.system_settings)
  420. router.add_get(f"{API_PREFIX}/settings/telegram", self.telegram_settings)
  421. router.add_put(f"{API_PREFIX}/settings/telegram", self.telegram_settings_update)
  422. router.add_get(f"{API_PREFIX}/bots", self.bots)
  423. router.add_post(f"{API_PREFIX}/bots", self.bot_create)
  424. router.add_put(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_update)
  425. router.add_delete(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_delete)
  426. router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/test", self.bot_test)
  427. router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/restart", self.bot_restart)
  428. router.add_get(f"{API_PREFIX}/roles", self.roles)
  429. router.add_post(f"{API_PREFIX}/roles", self.role_create)
  430. router.add_put(f"{API_PREFIX}/roles/{{role_id}}", self.role_update)
  431. router.add_delete(f"{API_PREFIX}/roles/{{role_id}}", self.role_delete)
  432. router.add_get(f"{API_PREFIX}/audit-logs", self.audit_logs)
  433. router.add_post(f"{API_PREFIX}/media", self.upload_media)
  434. router.add_get(f"{API_PREFIX}/chats", self.chats)
  435. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}", self.chat)
  436. router.add_patch(f"{API_PREFIX}/chats/{{chat_id}}/profile", self.chat_profile)
  437. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/permissions", self.chat_permissions)
  438. router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/announcements", self.announcement)
  439. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/admins", self.chat_admins)
  440. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/search", self.member_search)
  441. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/recent", self.recent_members)
  442. router.add_get(
  443. f"{API_PREFIX}/chats/{{chat_id}}/members/identity-changes",
  444. self.member_identity_changes,
  445. )
  446. router.add_post(
  447. f"{API_PREFIX}/chats/{{chat_id}}/members/{{user_id}}/actions",
  448. self.member_action,
  449. )
  450. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invites)
  451. router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_create)
  452. router.add_delete(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_revoke)
  453. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules)
  454. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules_update)
  455. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation)
  456. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation_update)
  457. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings)
  458. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings_update)
  459. router.add_get(f"{API_PREFIX}/points/accounts", self.point_accounts)
  460. router.add_get(f"{API_PREFIX}/points/leaderboard", self.point_leaderboard)
  461. router.add_get(f"{API_PREFIX}/points/transactions", self.point_transactions)
  462. router.add_post(f"{API_PREFIX}/points/adjustments", self.point_adjustment)
  463. router.add_get(f"{API_PREFIX}/points/export", self.points_export)
  464. router.add_get(f"{API_PREFIX}/giveaways", self.giveaways)
  465. router.add_post(f"{API_PREFIX}/giveaways", self.giveaway_create)
  466. router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}", self.giveaway_detail)
  467. router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/finish", self.giveaway_finish)
  468. router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/cancel", self.giveaway_cancel)
  469. router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/reroll", self.giveaway_reroll)
  470. router.add_get(
  471. f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants",
  472. self.giveaway_participants,
  473. )
  474. router.add_delete(
  475. f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants/{{user_id}}",
  476. self.giveaway_participant_remove,
  477. )
  478. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_bans)
  479. router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_ban_add)
  480. router.add_delete(
  481. f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans/{{user_id}}",
  482. self.giveaway_ban_remove,
  483. )
  484. router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}/export", self.giveaway_export)
  485. async def health(self, _: web.Request) -> web.Response:
  486. return success({"status": "ok"})
  487. def _set_session_cookie(self, response: web.Response, token: str) -> None:
  488. response.set_cookie(
  489. SESSION_COOKIE,
  490. token,
  491. httponly=True,
  492. secure=self.cookie_secure,
  493. samesite="Strict",
  494. max_age=self.session_hours * 3600,
  495. path="/",
  496. )
  497. async def _new_session(self, request: web.Request, username: str) -> tuple[str, str]:
  498. token = generate_session_token()
  499. csrf = generate_csrf_token()
  500. await create_admin_session(
  501. username=username,
  502. token_hash=hash_token(token),
  503. csrf_token=csrf,
  504. expires_at=datetime.now(UTC) + timedelta(hours=self.session_hours),
  505. remote_address=request.remote or "",
  506. user_agent=request.headers.get("User-Agent", ""),
  507. )
  508. return token, csrf
  509. async def login(self, request: web.Request) -> web.Response:
  510. body = await json_body(request)
  511. username = str(body.get("username") or "").strip()
  512. password = str(body.get("password") or "")
  513. set_audit(request, "auth.login", summary=f"Login as {username}")
  514. user = await get_admin_user(username)
  515. if user and int(user.get("failed_login_count", 0)) >= 5:
  516. last_failed = user.get("last_failed_login_at")
  517. if last_failed and datetime.now(UTC) - as_utc(last_failed) < timedelta(minutes=15):
  518. raise ApiProblem(
  519. "login_rate_limited", "登录失败次数过多,请 15 分钟后再试。", status=429
  520. )
  521. if not user or not verify_password(password, user.get("password_hash", "")):
  522. await record_login_failure(username)
  523. raise ApiProblem("invalid_credentials", "用户名或密码错误。", status=401)
  524. await record_login_success(username)
  525. token, csrf = await self._new_session(request, username)
  526. response = success(
  527. {
  528. "username": username,
  529. "must_change_password": bool(user.get("must_change_password")),
  530. "csrf_token": csrf,
  531. }
  532. )
  533. self._set_session_cookie(response, token)
  534. request["admin"] = {"username": username}
  535. return response
  536. async def me(self, request: web.Request) -> web.Response:
  537. return success(
  538. {
  539. "username": request["admin"]["username"],
  540. "must_change_password": request["admin"]["must_change_password"],
  541. "csrf_token": request["admin"]["csrf_token"],
  542. }
  543. )
  544. async def logout(self, request: web.Request) -> web.Response:
  545. set_audit(request, "auth.logout", summary="Logout")
  546. await revoke_admin_session(request["admin"]["session_token_hash"])
  547. response = success({"logged_out": True})
  548. response.del_cookie(SESSION_COOKIE, path="/")
  549. return response
  550. async def change_password(self, request: web.Request) -> web.Response:
  551. body = await json_body(request)
  552. current = str(body.get("current_password") or "")
  553. new_password = str(body.get("new_password") or "")
  554. username = request["admin"]["username"]
  555. user = await get_admin_user(username)
  556. if not user or not verify_password(current, user["password_hash"]):
  557. raise ApiProblem("invalid_current_password", "当前密码错误。", status=409)
  558. if verify_password(new_password, user["password_hash"]):
  559. raise ApiProblem("password_unchanged", "新密码不能与当前密码相同。", status=409)
  560. errors = validate_new_password(new_password, username)
  561. if errors:
  562. raise ApiProblem("weak_password", "新密码不符合要求。", details=errors)
  563. set_audit(request, "auth.password.change", summary="Change administrator password")
  564. await update_admin_password(username, hash_password(new_password))
  565. token, csrf = await self._new_session(request, username)
  566. response = success(
  567. {"username": username, "must_change_password": False, "csrf_token": csrf}
  568. )
  569. self._set_session_cookie(response, token)
  570. return response
  571. async def dashboard(self, _: web.Request) -> web.Response:
  572. counts = await dashboard_counts()
  573. giveaways, giveaway_total = await list_giveaways_page(
  574. page=1,
  575. page_size=5,
  576. all_bots=bool(getattr(wbb, "SUPERVISOR_MODE", False)),
  577. )
  578. accounts, account_total = await list_point_accounts(page=1, page_size=5)
  579. transactions, transaction_total = await list_point_transactions(page=1, page_size=8)
  580. return success(
  581. {
  582. "counts": {
  583. **counts,
  584. "giveaways": giveaway_total,
  585. "point_accounts": account_total,
  586. "point_transactions": transaction_total,
  587. },
  588. "recent_giveaways": giveaways,
  589. "top_accounts": accounts,
  590. "recent_transactions": transactions,
  591. }
  592. )
  593. async def system_settings(self, _: web.Request) -> web.Response:
  594. bot_connected = bool(getattr(wbb, "TELEGRAM_CONNECTED", False))
  595. return success(
  596. {
  597. "web_address": f"http://127.0.0.1:{getattr(wbb, 'ADMIN_WEB_PORT', 8088)}/admin",
  598. "bind_host": str(getattr(wbb, "ADMIN_WEB_HOST", "0.0.0.0")),
  599. "upload_limit_mb": int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)),
  600. "supervisor_mode": bool(getattr(wbb, "SUPERVISOR_MODE", False)),
  601. "bot": {
  602. "connected": bot_connected,
  603. "id": str(wbb.BOT_ID) if bot_connected else "",
  604. "username": wbb.BOT_USERNAME if bot_connected else "",
  605. "name": wbb.BOT_NAME if bot_connected else "",
  606. },
  607. }
  608. )
  609. def _supervisor(self):
  610. return getattr(wbb, "BOT_SUPERVISOR", None)
  611. def _telegram_status(self) -> dict[str, Any]:
  612. supervisor = self._supervisor()
  613. runtimes = supervisor.runtimes() if supervisor is not None else {}
  614. return telegram_config_status(self.bot_config_path, runtimes=runtimes)
  615. async def telegram_settings(self, _: web.Request) -> web.Response:
  616. return success(self._telegram_status())
  617. async def telegram_settings_update(self, request: web.Request) -> web.Response:
  618. body = await json_body(request)
  619. require_confirmation(body)
  620. set_audit(request, "telegram.settings.update", summary="Update Telegram API credentials")
  621. body.pop("confirm", None)
  622. update_telegram_config(self.bot_config_path, body)
  623. supervisor = self._supervisor()
  624. if supervisor is not None:
  625. status = telegram_config_status(self.bot_config_path)
  626. for profile in status["bots"]:
  627. if profile["enabled"] and profile["ready_to_connect"]:
  628. await supervisor.reconcile(profile["bot_id"])
  629. return success(self._telegram_status())
  630. async def bots(self, _: web.Request) -> web.Response:
  631. return success({"items": self._telegram_status()["bots"]})
  632. async def bot_create(self, request: web.Request) -> web.Response:
  633. body = await json_body(request)
  634. require_confirmation(body)
  635. set_audit(request, "bot.create", summary=str(body.get("label") or ""))
  636. body.pop("confirm", None)
  637. profile = create_bot_profile(self.bot_config_path, body)
  638. supervisor = self._supervisor()
  639. if supervisor is not None and profile["enabled"] and profile["ready_to_connect"]:
  640. await supervisor.start(profile["bot_id"])
  641. current = next(
  642. item
  643. for item in self._telegram_status()["bots"]
  644. if item["bot_id"] == profile["bot_id"]
  645. )
  646. return success(current, status=201)
  647. async def bot_update(self, request: web.Request) -> web.Response:
  648. body = await json_body(request)
  649. require_confirmation(body)
  650. bot_id = request.match_info["bot_id"]
  651. set_audit(request, "bot.update", target_id=bot_id, summary="Update Bot profile")
  652. body.pop("confirm", None)
  653. update_bot_profile(self.bot_config_path, bot_id, body)
  654. supervisor = self._supervisor()
  655. if supervisor is not None:
  656. await supervisor.reconcile(bot_id)
  657. current = next(
  658. item for item in self._telegram_status()["bots"] if item["bot_id"] == bot_id
  659. )
  660. return success(current)
  661. async def bot_delete(self, request: web.Request) -> web.Response:
  662. body = await json_body(request)
  663. require_confirmation(body)
  664. bot_id = request.match_info["bot_id"]
  665. profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
  666. set_audit(request, "bot.delete", target_id=bot_id, summary=str(profile.get("label") or ""))
  667. supervisor = self._supervisor()
  668. if supervisor is not None:
  669. await supervisor.stop(bot_id)
  670. delete_bot_profile(self.bot_config_path, bot_id)
  671. return success({"deleted": True})
  672. async def bot_test(self, request: web.Request) -> web.Response:
  673. body = await json_body(request)
  674. require_confirmation(body)
  675. bot_id = request.match_info["bot_id"]
  676. profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
  677. set_audit(request, "bot.test", target_id=bot_id, summary="Validate Bot Token")
  678. identity = await test_bot_token(str(profile.get("bot_token") or ""))
  679. store_bot_identity(self.bot_config_path, bot_id, identity)
  680. return success(identity)
  681. async def bot_restart(self, request: web.Request) -> web.Response:
  682. body = await json_body(request)
  683. require_confirmation(body)
  684. bot_id = request.match_info["bot_id"]
  685. set_audit(request, "bot.restart", target_id=bot_id, summary="Restart Bot worker")
  686. supervisor = self._supervisor()
  687. if supervisor is None:
  688. raise ApiProblem(
  689. "bot_supervisor_unavailable",
  690. "机器人监管服务尚未就绪。",
  691. status=503,
  692. )
  693. return success(await supervisor.restart(bot_id))
  694. async def roles(self, _: web.Request) -> web.Response:
  695. status = self._telegram_status()
  696. return success(
  697. {
  698. "items": status["roles"],
  699. "permission_catalog": status["permission_catalog"],
  700. }
  701. )
  702. async def role_create(self, request: web.Request) -> web.Response:
  703. body = await json_body(request)
  704. require_confirmation(body)
  705. set_audit(
  706. request,
  707. "bot.role.create",
  708. summary=str(body.get("name") or ""),
  709. )
  710. body.pop("confirm", None)
  711. return success(create_bot_role(self.bot_config_path, body), status=201)
  712. async def role_update(self, request: web.Request) -> web.Response:
  713. body = await json_body(request)
  714. require_confirmation(body)
  715. role_id = request.match_info["role_id"]
  716. set_audit(
  717. request,
  718. "bot.role.update",
  719. target_id=role_id,
  720. summary=str(body.get("name") or ""),
  721. )
  722. body.pop("confirm", None)
  723. role = update_bot_role(self.bot_config_path, role_id, body)
  724. supervisor = self._supervisor()
  725. if supervisor is not None:
  726. for profile in self._telegram_status()["bots"]:
  727. if role_id in profile["role_ids"]:
  728. await supervisor.reconcile(profile["bot_id"])
  729. return success(role)
  730. async def role_delete(self, request: web.Request) -> web.Response:
  731. body = await json_body(request)
  732. require_confirmation(body)
  733. role_id = request.match_info["role_id"]
  734. set_audit(
  735. request,
  736. "bot.role.delete",
  737. target_id=role_id,
  738. summary="Delete Bot role",
  739. )
  740. delete_bot_role(self.bot_config_path, role_id)
  741. return success({"deleted": True})
  742. async def audit_logs(self, request: web.Request) -> web.Response:
  743. page, page_size = page_params(request)
  744. chat_raw = request.query.get("chat_id")
  745. items, total = await list_audit_logs(
  746. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  747. action=request.query.get("action") or None,
  748. page=page,
  749. page_size=page_size,
  750. )
  751. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  752. async def upload_media(self, request: web.Request) -> web.Response:
  753. if not MESSAGE_DUMP_CHAT:
  754. raise ApiProblem(
  755. "message_dump_chat_missing",
  756. "请先配置媒体中转群。",
  757. status=409,
  758. )
  759. reader = await request.multipart()
  760. part = await reader.next()
  761. if not part or part.name != "file":
  762. raise ApiProblem("file_required", "请选择要上传的文件。")
  763. filename = Path(part.filename or "upload.bin").name
  764. content_type = part.headers.get("Content-Type", "application/octet-stream")
  765. size = 0
  766. temporary_path: Path | None = None
  767. try:
  768. with tempfile.NamedTemporaryFile(delete=False, suffix=Path(filename).suffix) as handle:
  769. temporary_path = Path(handle.name)
  770. while True:
  771. chunk = await part.read_chunk(1024 * 1024)
  772. if not chunk:
  773. break
  774. size += len(chunk)
  775. if size > self.upload_limit:
  776. raise ApiProblem("file_too_large", "上传文件超过 20 MB 限制。", status=413)
  777. handle.write(chunk)
  778. if content_type == "image/gif":
  779. media_type = "animation"
  780. sent = await telegram_app.send_animation(MESSAGE_DUMP_CHAT, str(temporary_path))
  781. file_id = sent.animation.file_id
  782. elif content_type.startswith("image/"):
  783. media_type = "photo"
  784. sent = await telegram_app.send_photo(MESSAGE_DUMP_CHAT, str(temporary_path))
  785. file_id = sent.photo.file_id
  786. elif content_type.startswith("video/"):
  787. media_type = "video"
  788. sent = await telegram_app.send_video(MESSAGE_DUMP_CHAT, str(temporary_path))
  789. file_id = sent.video.file_id
  790. else:
  791. media_type = "document"
  792. sent = await telegram_app.send_document(
  793. MESSAGE_DUMP_CHAT, str(temporary_path), file_name=filename
  794. )
  795. file_id = sent.document.file_id
  796. finally:
  797. if temporary_path:
  798. temporary_path.unlink(missing_ok=True)
  799. set_audit(request, "media.upload", summary=filename, target_id=file_id)
  800. return success(
  801. {
  802. "file_id": file_id,
  803. "type": media_type,
  804. "filename": filename,
  805. "size": size,
  806. },
  807. status=201,
  808. )
  809. async def chats(self, request: web.Request) -> web.Response:
  810. page, page_size = page_params(request)
  811. items, total = await list_accessible_chats(
  812. query=request.query.get("query", ""), page=page, page_size=page_size
  813. )
  814. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  815. async def chat(self, request: web.Request) -> web.Response:
  816. return success(await get_chat_overview(chat_id_param(request)))
  817. async def chat_profile(self, request: web.Request) -> web.Response:
  818. chat_id = chat_id_param(request)
  819. body = await json_body(request)
  820. require_confirmation(body)
  821. set_audit(request, "chat.profile.update", chat_id=chat_id, summary="Update group profile")
  822. return success(
  823. await update_chat_profile(
  824. chat_id, title=body.get("title"), description=body.get("description")
  825. )
  826. )
  827. async def chat_permissions(self, request: web.Request) -> web.Response:
  828. chat_id = chat_id_param(request)
  829. body = await json_body(request)
  830. require_confirmation(body)
  831. set_audit(request, "chat.permissions.update", chat_id=chat_id, summary="Update default member permissions")
  832. return success(await update_chat_permissions(chat_id, body.get("permissions") or {}))
  833. async def announcement(self, request: web.Request) -> web.Response:
  834. chat_id = chat_id_param(request)
  835. body = await json_body(request)
  836. require_confirmation(body)
  837. set_audit(request, "chat.announcement.send", chat_id=chat_id, summary=str(body.get("text") or "")[:200])
  838. return success(
  839. await send_announcement(
  840. chat_id,
  841. text=str(body.get("text") or ""),
  842. media_type=body.get("media_type"),
  843. file_id=body.get("file_id"),
  844. pin=bool(body.get("pin")),
  845. ),
  846. status=201,
  847. )
  848. async def chat_admins(self, request: web.Request) -> web.Response:
  849. return success({"items": await list_chat_admins(chat_id_param(request))})
  850. async def member_search(self, request: web.Request) -> web.Response:
  851. return success(
  852. {
  853. "items": await search_chat_members(
  854. chat_id_param(request),
  855. request.query.get("query", ""),
  856. limit=parse_int(request.query.get("limit", "20"), "limit"),
  857. )
  858. }
  859. )
  860. async def recent_members(self, request: web.Request) -> web.Response:
  861. return success(
  862. {
  863. "items": await list_recent_members(
  864. chat_id_param(request),
  865. limit=parse_int(request.query.get("limit", "30"), "limit"),
  866. )
  867. }
  868. )
  869. async def member_identity_changes(self, request: web.Request) -> web.Response:
  870. page, page_size = page_params(request)
  871. items, total = await list_member_identity_changes(
  872. chat_id_param(request), page=page, page_size=page_size
  873. )
  874. return success(
  875. {"items": items, "total": total, "page": page, "page_size": page_size}
  876. )
  877. async def member_action(self, request: web.Request) -> web.Response:
  878. chat_id = chat_id_param(request)
  879. user_id = parse_int(request.match_info["user_id"], "userId")
  880. body = await json_body(request)
  881. require_confirmation(body)
  882. action = str(body.get("action") or "")
  883. set_audit(
  884. request,
  885. f"chat.member.{action}",
  886. chat_id=chat_id,
  887. target_id=user_id,
  888. summary=str(body.get("reason") or ""),
  889. )
  890. duration = body.get("duration_seconds")
  891. return success(
  892. await execute_member_action(
  893. chat_id,
  894. user_id=user_id,
  895. action=action,
  896. reason=str(body.get("reason") or ""),
  897. duration_seconds=parse_int(duration, "duration_seconds", minimum=60) if duration else None,
  898. privileges=body.get("privileges") or {},
  899. )
  900. )
  901. async def invites(self, request: web.Request) -> web.Response:
  902. return success({"items": await get_invite_links(chat_id_param(request))})
  903. async def invite_create(self, request: web.Request) -> web.Response:
  904. chat_id = chat_id_param(request)
  905. body = await json_body(request)
  906. require_confirmation(body)
  907. expires_at = parse_datetime(body["expires_at"], "expires_at") if body.get("expires_at") else None
  908. member_limit = parse_int(body["member_limit"], "member_limit", minimum=1) if body.get("member_limit") else None
  909. set_audit(request, "chat.invite.create", chat_id=chat_id, summary=str(body.get("name") or ""))
  910. return success(
  911. await create_invite_link(
  912. chat_id,
  913. name=str(body.get("name") or "Admin panel"),
  914. expires_at=expires_at,
  915. member_limit=member_limit,
  916. ),
  917. status=201,
  918. )
  919. async def invite_revoke(self, request: web.Request) -> web.Response:
  920. chat_id = chat_id_param(request)
  921. body = await json_body(request)
  922. require_confirmation(body)
  923. invite_link = str(body.get("invite_link") or "")
  924. if not invite_link:
  925. raise ApiProblem("invite_link_required", "缺少邀请链接。")
  926. set_audit(request, "chat.invite.revoke", chat_id=chat_id, summary=invite_link)
  927. await revoke_invite_link(chat_id, invite_link)
  928. return success({"revoked": True})
  929. async def rules(self, request: web.Request) -> web.Response:
  930. return success({"rules": await get_rules(chat_id_param(request))})
  931. async def rules_update(self, request: web.Request) -> web.Response:
  932. chat_id = chat_id_param(request)
  933. body = await json_body(request)
  934. require_confirmation(body)
  935. await ensure_permission(chat_id, "can_change_info")
  936. rules = str(body.get("rules") or "")
  937. if len(rules) > 4000:
  938. raise ApiProblem("rules_too_long", "群规不能超过 4000 个字符。")
  939. set_audit(request, "chat.rules.update", chat_id=chat_id, summary="Update group rules")
  940. await set_chat_rules(chat_id, rules)
  941. return success({"rules": rules})
  942. async def automation(self, request: web.Request) -> web.Response:
  943. return success(await get_automation_settings(chat_id_param(request)))
  944. async def automation_update(self, request: web.Request) -> web.Response:
  945. chat_id = chat_id_param(request)
  946. body = await json_body(request)
  947. require_confirmation(body)
  948. set_audit(request, "chat.automation.update", chat_id=chat_id, summary="Update automation rules")
  949. body.pop("confirm", None)
  950. return success(await apply_automation_settings(chat_id, body))
  951. async def points_settings(self, request: web.Request) -> web.Response:
  952. return success(await get_point_rules(chat_id_param(request)))
  953. async def points_settings_update(self, request: web.Request) -> web.Response:
  954. chat_id = chat_id_param(request)
  955. body = await json_body(request)
  956. require_confirmation(body)
  957. await ensure_permission(chat_id, "can_change_info")
  958. set_audit(request, "points.settings.update", chat_id=chat_id, summary="Update point rules")
  959. body.pop("confirm", None)
  960. return success(await apply_point_rules(chat_id, body))
  961. async def point_accounts(self, request: web.Request) -> web.Response:
  962. page, page_size = page_params(request)
  963. chat_raw = request.query.get("chat_id")
  964. items, total = await list_point_accounts(
  965. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  966. query=request.query.get("query", ""),
  967. page=page,
  968. page_size=page_size,
  969. )
  970. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  971. async def point_leaderboard(self, request: web.Request) -> web.Response:
  972. if not request.query.get("chat_id"):
  973. raise ApiProblem("chat_id_required", "排行榜必须选择群组。")
  974. return await self.point_accounts(request)
  975. async def point_transactions(self, request: web.Request) -> web.Response:
  976. page, page_size = page_params(request)
  977. chat_raw = request.query.get("chat_id")
  978. user_raw = request.query.get("user_id")
  979. created_from = parse_datetime(request.query["created_from"], "created_from") if request.query.get("created_from") else None
  980. created_to = parse_datetime(request.query["created_to"], "created_to") if request.query.get("created_to") else None
  981. items, total = await list_point_transactions(
  982. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  983. user_id=parse_int(user_raw, "user_id") if user_raw else None,
  984. source=request.query.get("source") or None,
  985. created_from=created_from,
  986. created_to=created_to,
  987. page=page,
  988. page_size=page_size,
  989. )
  990. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  991. async def point_adjustment(self, request: web.Request) -> web.Response:
  992. body = await json_body(request)
  993. require_confirmation(body)
  994. chat_id = parse_int(body.get("chat_id"), "chat_id")
  995. user_id = parse_int(body.get("user_id"), "user_id")
  996. operation = str(body.get("operation") or "")
  997. amount = parse_int(body.get("amount"), "amount", minimum=0)
  998. reason = str(body.get("reason") or "").strip()
  999. if not reason:
  1000. raise ApiProblem("reason_required", "积分调整必须填写原因。")
  1001. if operation not in {"add", "deduct", "set"}:
  1002. raise ApiProblem(
  1003. "invalid_operation",
  1004. "积分操作必须是增加、扣减或设置余额。",
  1005. )
  1006. try:
  1007. user = await telegram_app.get_users(user_id)
  1008. username, first_name = user.username, user.first_name
  1009. display_name = " ".join(
  1010. value for value in (user.first_name, user.last_name) if value
  1011. )
  1012. except Exception:
  1013. username, first_name = body.get("username"), body.get("first_name")
  1014. display_name = body.get("display_name") or first_name
  1015. request_key = str(body.get("request_id") or generate_session_token())
  1016. set_audit(
  1017. request,
  1018. f"points.adjust.{operation}",
  1019. chat_id=chat_id,
  1020. target_id=user_id,
  1021. summary=reason,
  1022. metadata={"amount": amount},
  1023. )
  1024. if operation == "set":
  1025. account, created = await set_points(
  1026. chat_id=chat_id,
  1027. user_id=user_id,
  1028. balance=amount,
  1029. actor_id=request["admin"]["username"],
  1030. reason=reason,
  1031. idempotency_key=f"web-points:{request_key}",
  1032. username=username,
  1033. first_name=first_name,
  1034. display_name=display_name,
  1035. )
  1036. else:
  1037. account, created = await adjust_points(
  1038. chat_id=chat_id,
  1039. user_id=user_id,
  1040. delta=amount if operation == "add" else -amount,
  1041. source=SOURCE_ADMIN,
  1042. idempotency_key=f"web-points:{request_key}",
  1043. actor_id=request["admin"]["username"],
  1044. reason=reason,
  1045. username=username,
  1046. first_name=first_name,
  1047. display_name=display_name,
  1048. )
  1049. return success({"account": account, "created": created}, status=201 if created else 200)
  1050. async def points_export(self, request: web.Request) -> web.Response:
  1051. chat_raw = request.query.get("chat_id")
  1052. if not chat_raw:
  1053. raise ApiProblem("chat_id_required", "导出积分数据必须选择群组。")
  1054. chat_id = parse_int(chat_raw, "chat_id")
  1055. kind = request.query.get("kind", "accounts")
  1056. output = io.StringIO()
  1057. if kind == "accounts":
  1058. items: list[dict[str, Any]] = []
  1059. page = 1
  1060. while True:
  1061. batch, total = await list_point_accounts(
  1062. chat_id=chat_id, page=page, page_size=100
  1063. )
  1064. items.extend(batch)
  1065. if not batch or len(items) >= total:
  1066. break
  1067. page += 1
  1068. fieldnames = ["chat_id", "user_id", "display_name", "username", "first_name", "balance", "lifetime_earned", "lifetime_spent", "updated_at"]
  1069. elif kind == "transactions":
  1070. items = []
  1071. page = 1
  1072. while True:
  1073. batch, total = await list_point_transactions(
  1074. chat_id=chat_id, page=page, page_size=100
  1075. )
  1076. items.extend(batch)
  1077. if not batch or len(items) >= total:
  1078. break
  1079. page += 1
  1080. fieldnames = ["transaction_id", "chat_id", "user_id", "display_name", "username", "first_name", "delta", "balance_after", "source", "actor_id", "reason", "reference_id", "created_at"]
  1081. else:
  1082. raise ApiProblem(
  1083. "invalid_export_kind",
  1084. "导出类型必须是积分账户或积分流水。",
  1085. )
  1086. writer = csv.DictWriter(output, fieldnames=fieldnames, extrasaction="ignore")
  1087. writer.writeheader()
  1088. for item in items:
  1089. writer.writerow(jsonable(item))
  1090. return web.Response(
  1091. text=output.getvalue(),
  1092. content_type="text/csv",
  1093. headers={"Content-Disposition": f'attachment; filename="points-{chat_id}-{kind}.csv"'},
  1094. )
  1095. async def giveaways(self, request: web.Request) -> web.Response:
  1096. page, page_size = page_params(request)
  1097. chat_raw = request.query.get("chat_id")
  1098. items, total = await list_giveaways_page(
  1099. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  1100. status=request.query.get("status") or None,
  1101. query=request.query.get("query", ""),
  1102. page=page,
  1103. page_size=page_size,
  1104. )
  1105. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  1106. async def giveaway_create(self, request: web.Request) -> web.Response:
  1107. body = await json_body(request)
  1108. require_confirmation(body)
  1109. chat_id = parse_int(body.get("chat_id"), "chat_id")
  1110. await ensure_permission(chat_id, "can_change_info")
  1111. prizes = body.get("prizes")
  1112. if not isinstance(prizes, list):
  1113. raise ApiProblem("invalid_prizes", "奖项必须是数组。")
  1114. set_audit(request, "giveaway.create", chat_id=chat_id, summary=str(body.get("title") or ""))
  1115. giveaway = await create_and_publish_giveaway(
  1116. chat_id=chat_id,
  1117. creator_id=0,
  1118. creator_name=request["admin"]["username"],
  1119. title=str(body.get("title") or ""),
  1120. description=str(body.get("description") or ""),
  1121. prizes=prizes,
  1122. ends_at=parse_datetime(body.get("ends_at"), "ends_at"),
  1123. minimum_points=parse_int(body.get("minimum_points", 0), "minimum_points", minimum=0),
  1124. entry_cost=parse_int(body.get("entry_cost", 0), "entry_cost", minimum=0),
  1125. participation_reward=parse_int(body.get("participation_reward", 0), "participation_reward", minimum=0),
  1126. )
  1127. return success(giveaway, status=201)
  1128. async def giveaway_detail(self, request: web.Request) -> web.Response:
  1129. giveaway = await get_giveaway(request.match_info["giveaway_id"])
  1130. if not giveaway:
  1131. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  1132. return success(giveaway)
  1133. async def giveaway_finish(self, request: web.Request) -> web.Response:
  1134. body = await json_body(request)
  1135. require_confirmation(body)
  1136. giveaway_id = request.match_info["giveaway_id"]
  1137. giveaway = await get_giveaway(giveaway_id)
  1138. if not giveaway:
  1139. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  1140. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  1141. set_audit(request, "giveaway.finish", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Manual draw")
  1142. ok, message, current = await finish_and_publish_giveaway(giveaway_id)
  1143. if not ok:
  1144. raise ApiProblem("giveaway_not_running", message, status=409)
  1145. return success(current)
  1146. async def giveaway_cancel(self, request: web.Request) -> web.Response:
  1147. body = await json_body(request)
  1148. require_confirmation(body)
  1149. giveaway_id = request.match_info["giveaway_id"]
  1150. giveaway = await get_giveaway(giveaway_id)
  1151. if not giveaway:
  1152. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  1153. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  1154. set_audit(request, "giveaway.cancel", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Cancel and refund")
  1155. ok, message, current = await cancel_and_refund_giveaway(
  1156. giveaway_id, chat_id=int(giveaway["chat_id"])
  1157. )
  1158. if not ok:
  1159. raise ApiProblem("giveaway_not_running", message, status=409)
  1160. return success(current)
  1161. async def giveaway_reroll(self, request: web.Request) -> web.Response:
  1162. body = await json_body(request)
  1163. require_confirmation(body)
  1164. giveaway_id = request.match_info["giveaway_id"]
  1165. giveaway = await get_giveaway(giveaway_id)
  1166. if not giveaway:
  1167. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  1168. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  1169. set_audit(request, "giveaway.reroll", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary=str(body.get("tier_name") or "All tiers"))
  1170. winners, reroll_id = await reroll_giveaway(
  1171. giveaway_id=giveaway_id,
  1172. moderator_id=0,
  1173. tier_name=body.get("tier_name") or None,
  1174. reroll_id=body.get("request_id") or None,
  1175. )
  1176. return success({"winners": winners, "reroll_id": reroll_id})
  1177. async def giveaway_participants(self, request: web.Request) -> web.Response:
  1178. page, page_size = page_params(request)
  1179. items, total = await list_participants_page(
  1180. giveaway_id=request.match_info["giveaway_id"],
  1181. active_only=request.query.get("active_only", "false").lower() == "true",
  1182. page=page,
  1183. page_size=page_size,
  1184. )
  1185. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  1186. async def giveaway_participant_remove(self, request: web.Request) -> web.Response:
  1187. body = await json_body(request)
  1188. require_confirmation(body)
  1189. giveaway_id = request.match_info["giveaway_id"]
  1190. user_id = parse_int(request.match_info["user_id"], "user_id")
  1191. giveaway = await get_giveaway(giveaway_id)
  1192. if not giveaway:
  1193. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  1194. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  1195. set_audit(request, "giveaway.participant.remove", chat_id=int(giveaway["chat_id"]), target_id=user_id, summary=str(body.get("reason") or ""))
  1196. removed = await remove_and_optionally_refund_participant(
  1197. giveaway_id=giveaway_id,
  1198. user_id=user_id,
  1199. moderator_id=0,
  1200. reason=str(body.get("reason") or "Removed in admin panel"),
  1201. refund=bool(body.get("refund", True)),
  1202. )
  1203. if not removed:
  1204. raise ApiProblem("participant_not_found", "参与者不存在或已被移除。", status=404)
  1205. return success({"removed": True, "refunded": bool(body.get("refund", True))})
  1206. async def giveaway_bans(self, request: web.Request) -> web.Response:
  1207. page, page_size = page_params(request)
  1208. items, total = await list_giveaway_bans(
  1209. chat_id=chat_id_param(request), page=page, page_size=page_size
  1210. )
  1211. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  1212. async def giveaway_ban_add(self, request: web.Request) -> web.Response:
  1213. chat_id = chat_id_param(request)
  1214. body = await json_body(request)
  1215. require_confirmation(body)
  1216. user_id = parse_int(body.get("user_id"), "user_id")
  1217. await ensure_permission(chat_id, "can_change_info")
  1218. set_audit(request, "giveaway.ban.add", chat_id=chat_id, target_id=user_id, summary=str(body.get("reason") or ""))
  1219. await add_giveaway_ban(chat_id=chat_id, user_id=user_id, moderator_id=0, reason=str(body.get("reason") or ""))
  1220. return success({"created": True}, status=201)
  1221. async def giveaway_ban_remove(self, request: web.Request) -> web.Response:
  1222. chat_id = chat_id_param(request)
  1223. user_id = parse_int(request.match_info["user_id"], "user_id")
  1224. body = await json_body(request)
  1225. require_confirmation(body)
  1226. await ensure_permission(chat_id, "can_change_info")
  1227. set_audit(request, "giveaway.ban.remove", chat_id=chat_id, target_id=user_id)
  1228. return success({"removed": await remove_giveaway_ban(chat_id, user_id)})
  1229. async def giveaway_export(self, request: web.Request) -> web.Response:
  1230. giveaway_id = request.match_info["giveaway_id"]
  1231. giveaway = await get_giveaway(giveaway_id)
  1232. if not giveaway:
  1233. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  1234. participants: list[dict[str, Any]] = []
  1235. page = 1
  1236. while True:
  1237. batch, total = await list_participants_page(
  1238. giveaway_id=giveaway_id, page=page, page_size=100
  1239. )
  1240. participants.extend(batch)
  1241. if not batch or len(participants) >= total:
  1242. break
  1243. page += 1
  1244. output = io.StringIO()
  1245. fields = ["giveaway_id", "chat_id", "user_id", "display_name", "username", "first_name", "active", "entry_cost", "joined_at", "removed_at", "refunded_at"]
  1246. writer = csv.DictWriter(output, fieldnames=fields, extrasaction="ignore")
  1247. writer.writeheader()
  1248. for participant in participants:
  1249. writer.writerow(jsonable(participant))
  1250. return web.Response(
  1251. text=output.getvalue(),
  1252. content_type="text/csv",
  1253. headers={"Content-Disposition": f'attachment; filename="giveaway-{giveaway_id}.csv"'},
  1254. )
  1255. def build_admin_application() -> web.Application:
  1256. max_upload_mb = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20))
  1257. application = web.Application(
  1258. middlewares=[
  1259. api_error_middleware,
  1260. authentication_middleware,
  1261. bot_role_middleware,
  1262. bot_proxy_middleware,
  1263. ],
  1264. client_max_size=(max_upload_mb + 1) * 1024 * 1024,
  1265. )
  1266. admin_api = AdminApi()
  1267. application["admin_api"] = admin_api
  1268. admin_api.register(application)
  1269. return application