| 123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134135136137138139140141142143144145146147148149150151152153154155156157158159160161162163164165166167168169170171172173174175176177178179180181182183184185186187188189190191192193194195196197198199200201202203204205206207208209210211212213214215216217218219220221222223224225226227228229230231232233234235236237238239240241242243244245246247248249250251252253254255256257258259260261262263264265266267268269270271272273274275276277278279280281282283284285286287288289290291292293294295296297298299300301302303304305306307308309310311312313314315316317318319320321322323324325326327328329330331332333334335336337338339340341342343344345346347348349350351352353354355356357358359360361362363364365366367368369370371372373374375376377378379380381382383384385386387388389390391392393394395396397398399400401402403404405406407408409410411412413414415416417418419420421422423424425426427428429430431432433434435436437438439440441442443444445446447448449450451452453454455456457458459460461462463464465466467468469470471472473474475476477478479480481482483484485486487488489490491492493494495496497498499500501502503504505506507508509510511512513514515516517518519520521522523524525526527528529530531532533534535536537538539540541542543544545546547548549550551552553554555556557558559560561562563564565566567568569570571572573574575576577578579580581582583584585586587588589590591592593594595596597598599600601602603604605606607608609610611612613614615616617618619620621622623624625626627628629630631632633634635636637638639640641642643644645646647648649650651652653654655656657658659660661662663664665666667668669670671672673674675676677678679680681682683684685686687688689690691692693694695696697698699700701702703704705706707708709710711712713714715716717718719720721722723724725726727728729730731732733734735736737738739740741742743744745746747748749750751752753754755756757758759760761762763764765766767768769770771772773774775776777778779780781782783784785786787788789790791792793794795796797798799800801802803804805806807808809810811812813814815816817818819820821822823824825826827828829830831832833834835836837838839840841842843844845846847848849850851852853854855856857858859860861862863864865866867868869870871872873874875876877878879880881882883884885886887888889890891892893894895896897898899900901902903904905906907908909910911912913914915916917918919920921922923924925926927928929930931932933934935936937938939940941942943944945946947948949950951952953954955956957958959960961962963964965966967968969970971972973974975976977978979980981982983984985986987988989990991992993994995996997998999100010011002100310041005100610071008100910101011101210131014101510161017101810191020102110221023102410251026102710281029103010311032103310341035103610371038103910401041104210431044104510461047104810491050105110521053105410551056105710581059106010611062106310641065106610671068106910701071107210731074107510761077107810791080108110821083108410851086108710881089109010911092109310941095109610971098109911001101110211031104110511061107110811091110111111121113111411151116111711181119112011211122112311241125112611271128112911301131113211331134113511361137113811391140114111421143114411451146114711481149115011511152115311541155115611571158115911601161116211631164116511661167116811691170117111721173117411751176117711781179118011811182118311841185118611871188118911901191119211931194119511961197119811991200120112021203120412051206120712081209121012111212121312141215121612171218121912201221122212231224122512261227122812291230123112321233123412351236123712381239124012411242124312441245124612471248124912501251125212531254125512561257125812591260126112621263126412651266126712681269127012711272127312741275127612771278127912801281128212831284128512861287128812891290129112921293129412951296129712981299130013011302130313041305130613071308130913101311131213131314131513161317131813191320132113221323132413251326132713281329133013311332133313341335133613371338133913401341134213431344134513461347134813491350135113521353135413551356135713581359136013611362136313641365136613671368136913701371137213731374137513761377137813791380138113821383138413851386138713881389139013911392139313941395139613971398139914001401140214031404140514061407140814091410141114121413141414151416141714181419142014211422142314241425142614271428142914301431143214331434143514361437143814391440144114421443144414451446144714481449145014511452145314541455145614571458145914601461146214631464146514661467146814691470147114721473147414751476147714781479148014811482148314841485148614871488148914901491149214931494149514961497149814991500150115021503150415051506150715081509151015111512151315141515151615171518151915201521152215231524152515261527152815291530153115321533153415351536153715381539154015411542154315441545154615471548154915501551155215531554155515561557155815591560156115621563156415651566156715681569157015711572157315741575157615771578157915801581158215831584158515861587158815891590159115921593159415951596159715981599160016011602160316041605160616071608160916101611161216131614161516161617161816191620162116221623162416251626162716281629163016311632163316341635163616371638163916401641164216431644164516461647164816491650165116521653165416551656165716581659166016611662166316641665166616671668166916701671167216731674167516761677167816791680168116821683168416851686168716881689169016911692169316941695169616971698169917001701170217031704170517061707170817091710171117121713171417151716171717181719172017211722172317241725172617271728172917301731173217331734173517361737173817391740174117421743174417451746174717481749175017511752175317541755175617571758175917601761176217631764176517661767176817691770177117721773177417751776177717781779178017811782178317841785178617871788178917901791179217931794179517961797179817991800180118021803180418051806180718081809181018111812181318141815181618171818181918201821182218231824182518261827182818291830183118321833183418351836183718381839184018411842184318441845184618471848184918501851185218531854185518561857185818591860186118621863186418651866186718681869187018711872187318741875 |
- from __future__ import annotations
- import csv
- import io
- import secrets
- import tempfile
- import traceback
- from contextlib import suppress
- from datetime import UTC, datetime, timedelta
- from pathlib import Path
- from typing import Any
- from aiohttp import ClientError, ClientSession, ClientTimeout, web
- import wbb
- from wbb import MESSAGE_DUMP_CHAT
- from wbb import app as telegram_app
- from wbb.admin.bot_config import (
- BotConfigError,
- create_bot_profile,
- create_bot_role,
- delete_bot_profile,
- delete_bot_role,
- get_bot_profile_secrets,
- store_bot_identity,
- telegram_config_status,
- test_bot_token,
- update_bot_profile,
- update_bot_role,
- update_telegram_config,
- )
- from wbb.admin.security import (
- generate_csrf_token,
- generate_session_token,
- hash_password,
- hash_token,
- validate_new_password,
- verify_password,
- )
- from wbb.services.bot_permissions import api_permission, has_permission
- from wbb.services.chat_management import (
- ChatManagementError,
- apply_automation_settings,
- create_invite_link,
- ensure_permission,
- execute_member_action,
- get_automation_settings,
- get_chat_overview,
- get_invite_links,
- list_accessible_chats,
- list_chat_admins,
- list_recent_members,
- revoke_invite_link,
- search_chat_members,
- send_announcement,
- update_chat_permissions,
- update_chat_profile,
- )
- from wbb.services.directory import DirectoryServiceError, verify_local_membership
- from wbb.services.directory_geocoding import (
- start_directory_geocoder,
- stop_directory_geocoder,
- )
- from wbb.services.giveaways import (
- GiveawayServiceError,
- cancel_and_refund_giveaway,
- create_and_publish_giveaway,
- finish_and_publish_giveaway,
- remove_and_optionally_refund_participant,
- reroll_giveaway,
- )
- from wbb.services.point_settings import apply_point_rules
- from wbb.utils.dbadmin import (
- create_admin_session,
- dashboard_counts,
- ensure_default_admin,
- get_admin_session,
- get_admin_user,
- list_audit_logs,
- list_member_identity_changes,
- record_audit,
- record_login_failure,
- record_login_success,
- revoke_admin_session,
- update_admin_password,
- )
- from wbb.utils.dbdirectory import (
- APPLICATION_APPROVED,
- DirectoryDataError,
- clear_directory_location,
- decide_teacher_application,
- ensure_directory_indexes,
- get_directory_location,
- get_directory_profile,
- get_directory_settings,
- list_directory_events,
- list_directory_locations,
- list_directory_profiles,
- list_membership_candidates,
- list_teacher_applications,
- save_directory_location,
- set_directory_settings,
- set_teacher_state,
- )
- from wbb.utils.dbfunctions import get_rules, set_chat_rules
- from wbb.utils.dbgiveaway import (
- add_giveaway_ban,
- get_giveaway,
- list_giveaway_bans,
- list_giveaways_page,
- list_participants_page,
- remove_giveaway_ban,
- )
- from wbb.utils.dbpoints import (
- SOURCE_ADMIN,
- InsufficientPoints,
- PointsError,
- adjust_points,
- get_point_rules,
- list_point_accounts,
- list_point_transactions,
- set_points,
- )
- API_PREFIX = "/api/admin/v1"
- SESSION_COOKIE = "wbb_admin_session"
- UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"}
- PUBLIC_API_PATHS = {f"{API_PREFIX}/auth/login", f"{API_PREFIX}/health"}
- BOT_SCOPED_PREFIXES = (
- f"{API_PREFIX}/chats",
- f"{API_PREFIX}/giveaways",
- f"{API_PREFIX}/points",
- f"{API_PREFIX}/media",
- )
- TELEGRAM_ID_KEYS = {
- "chat_id",
- "user_id",
- "actor_id",
- "creator_id",
- "moderator_id",
- "target_id",
- "message_id",
- "removed_by",
- "authorization_chat_id",
- }
- class ApiProblem(RuntimeError):
- def __init__(
- self,
- code: str,
- message: str,
- *,
- status: int = 400,
- details: Any = None,
- ):
- super().__init__(message)
- self.code = code
- self.status = status
- self.details = details
- def as_utc(value: datetime) -> datetime:
- if value.tzinfo is None:
- return value.replace(tzinfo=UTC)
- return value.astimezone(UTC)
- def jsonable(value: Any, *, key: str = "") -> Any:
- if isinstance(value, datetime):
- return as_utc(value).isoformat().replace("+00:00", "Z")
- if isinstance(value, dict):
- return {
- item_key: jsonable(item_value, key=item_key)
- for item_key, item_value in value.items()
- if item_key != "_id"
- }
- if isinstance(value, (list, tuple)):
- return [jsonable(item) for item in value]
- if isinstance(value, int) and (key in TELEGRAM_ID_KEYS or key.endswith("_telegram_id")):
- return str(value)
- return value
- def success(data: Any, *, status: int = 200) -> web.Response:
- return web.json_response({"data": jsonable(data)}, status=status)
- def error_response(problem: ApiProblem) -> web.Response:
- return web.json_response(
- {
- "error": {
- "code": problem.code,
- "message": str(problem),
- "details": jsonable(problem.details),
- }
- },
- status=problem.status,
- )
- def parse_int(value: Any, name: str, *, minimum: int | None = None) -> int:
- try:
- parsed = int(value)
- except (TypeError, ValueError) as exc:
- raise ApiProblem("invalid_parameter", f"{name} 必须是整数。") from exc
- if minimum is not None and parsed < minimum:
- raise ApiProblem("invalid_parameter", f"{name} 不能小于 {minimum}。")
- return parsed
- def parse_datetime(value: Any, name: str) -> datetime:
- if not isinstance(value, str) or not value.strip():
- raise ApiProblem("invalid_parameter", f"{name} 必须是 UTC ISO 8601 时间。")
- try:
- parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
- except ValueError as exc:
- raise ApiProblem("invalid_parameter", f"{name} 不是有效时间。") from exc
- if parsed.tzinfo is None:
- raise ApiProblem("invalid_parameter", f"{name} 必须包含时区。")
- return parsed.astimezone(UTC)
- async def json_body(request: web.Request) -> dict[str, Any]:
- try:
- body = await request.json()
- except Exception as exc:
- raise ApiProblem("invalid_json", "请求体必须是 JSON。") from exc
- if not isinstance(body, dict):
- raise ApiProblem("invalid_json", "请求体必须是 JSON 对象。")
- return body
- def require_confirmation(body: dict[str, Any]) -> None:
- if body.get("confirm") is not True:
- raise ApiProblem(
- "confirmation_required", "该操作需要二次确认。", status=409
- )
- def page_params(request: web.Request) -> tuple[int, int]:
- return (
- parse_int(request.query.get("page", 1), "page", minimum=1),
- min(parse_int(request.query.get("page_size", 20), "page_size", minimum=1), 100),
- )
- def chat_id_param(request: web.Request) -> int:
- return parse_int(request.match_info["chat_id"], "chatId")
- def set_audit(
- request: web.Request,
- action: str,
- *,
- chat_id: int | None = None,
- target_id: int | str | None = None,
- summary: str = "",
- metadata: dict[str, Any] | None = None,
- ) -> None:
- request["audit_context"] = {
- "action": action,
- "chat_id": chat_id,
- "target_id": target_id,
- "summary": summary,
- "metadata": metadata or {},
- }
- @web.middleware
- async def api_error_middleware(request: web.Request, handler):
- try:
- response = await handler(request)
- except ApiProblem as exc:
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(exc)
- except (ChatManagementError, GiveawayServiceError, DirectoryServiceError) as exc:
- problem = ApiProblem(
- exc.code, str(exc), status=getattr(exc, "status", 400)
- )
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except DirectoryDataError as exc:
- problem = ApiProblem(exc.code, str(exc), status=409)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except (PointsError, InsufficientPoints) as exc:
- problem = ApiProblem("points_error", str(exc), status=409)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except BotConfigError as exc:
- status = {
- "bot_not_found": 404,
- "role_not_found": 404,
- "role_in_use": 409,
- "builtin_role_immutable": 409,
- }.get(exc.code, 400)
- problem = ApiProblem(exc.code, str(exc), status=status)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except web.HTTPException:
- raise
- except Exception as exc:
- await _record_request_audit(request, success_state=False, error=str(exc))
- wbb.log.error(
- f"Admin API {request.method} {request.path} failed: {exc}\n"
- f"{traceback.format_exc()}"
- )
- return error_response(
- ApiProblem("internal_error", "服务器处理请求失败。", status=500)
- )
- await _record_request_audit(request, success_state=response.status < 400)
- return response
- async def _record_request_audit(
- request: web.Request, *, success_state: bool, error: str = ""
- ) -> None:
- context = request.get("audit_context")
- if not context or request.get("audit_recorded"):
- return
- request["audit_recorded"] = True
- admin = request.get("admin") or {}
- with suppress(Exception):
- await record_audit(
- source="web",
- actor_id=admin.get("username", "anonymous"),
- actor_name=admin.get("username", "anonymous"),
- success=success_state,
- error=error,
- **context,
- )
- @web.middleware
- async def authentication_middleware(request: web.Request, handler):
- if not request.path.startswith(API_PREFIX) or request.path in PUBLIC_API_PATHS:
- return await handler(request)
- raw_token = request.cookies.get(SESSION_COOKIE, "")
- session = await get_admin_session(hash_token(raw_token)) if raw_token else None
- if not session:
- raise ApiProblem("unauthenticated", "登录已失效,请重新登录。", status=401)
- user = session["user"]
- request["admin"] = {
- "username": user["username"],
- "must_change_password": bool(user.get("must_change_password")),
- "csrf_token": session["csrf_token"],
- "session_token_hash": session["token_hash"],
- }
- allowed_during_password_change = {
- f"{API_PREFIX}/auth/me",
- f"{API_PREFIX}/auth/password",
- f"{API_PREFIX}/auth/logout",
- }
- if user.get("must_change_password") and request.path not in allowed_during_password_change:
- raise ApiProblem(
- "password_change_required", "首次登录必须修改初始密码。", status=428
- )
- if request.method in UNSAFE_METHODS:
- supplied = request.headers.get("X-CSRF-Token", "")
- if not supplied or supplied != session["csrf_token"]:
- raise ApiProblem("csrf_failed", "CSRF 校验失败。", status=403)
- return await handler(request)
- def _selected_bot_id(request: web.Request) -> str:
- bot_id = str(request.headers.get("X-Bot-Id") or "").strip()
- if bot_id:
- return bot_id
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is not None:
- running = [
- key
- for key, value in supervisor.runtimes().items()
- if value.get("state") == "running"
- ]
- if len(running) == 1:
- return running[0]
- raise ApiProblem(
- "bot_selection_required",
- "请先选择要管理的机器人。",
- status=409,
- )
- def _require_internal_request(request: web.Request) -> None:
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- expected = str(
- getattr(supervisor, "internal_token", "")
- or getattr(wbb, "INTERNAL_TOKEN", "")
- or ""
- )
- supplied = request.headers.get("Authorization", "")
- token = supplied.removeprefix("Bearer ").strip()
- if not expected or not token or not secrets.compare_digest(expected, token):
- raise ApiProblem(
- "internal_auth_failed",
- "内部服务鉴权失败。",
- status=403,
- )
- if request.remote and request.remote not in {"127.0.0.1", "::1"}:
- raise ApiProblem(
- "internal_loopback_required",
- "内部服务只接受本机请求。",
- status=403,
- )
- @web.middleware
- async def bot_role_middleware(request: web.Request, handler):
- required = api_permission(request.method, request.path)
- if not required:
- return await handler(request)
- if bool(getattr(wbb, "SUPERVISOR_MODE", False)):
- profile = get_bot_profile_secrets(
- getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"),
- _selected_bot_id(request),
- )
- permissions = profile.get("permissions", [])
- else:
- permissions = getattr(wbb, "BOT_PERMISSIONS", {"*"})
- if not has_permission(required, permissions):
- raise ApiProblem(
- "bot_role_permission_denied",
- "所选机器人的职责角色不允许执行该操作。",
- status=403,
- details={"required_permission": required},
- )
- return await handler(request)
- @web.middleware
- async def bot_proxy_middleware(request: web.Request, handler):
- if not bool(getattr(wbb, "SUPERVISOR_MODE", False)) or not request.path.startswith(
- BOT_SCOPED_PREFIXES
- ):
- return await handler(request)
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is None:
- raise ApiProblem(
- "bot_supervisor_unavailable",
- "机器人监管服务尚未就绪。",
- status=503,
- )
- bot_id = _selected_bot_id(request)
- endpoint = supervisor.endpoint_for(bot_id)
- if endpoint is None:
- raise ApiProblem(
- "bot_not_running",
- "所选机器人当前未连接 Telegram。",
- status=409,
- )
- target = f"{endpoint}{request.rel_url}"
- headers = {
- key: value
- for key, value in request.headers.items()
- if key.lower() not in {"host", "content-length", "connection"}
- }
- body = await request.read()
- try:
- async with ClientSession(timeout=ClientTimeout(total=90)) as session:
- async with session.request(
- request.method,
- target,
- data=body or None,
- headers=headers,
- allow_redirects=False,
- ) as response:
- response_body = await response.read()
- response_headers = {
- key: value
- for key, value in response.headers.items()
- if key.lower() in {"content-type", "content-disposition"}
- }
- return web.Response(
- body=response_body,
- status=response.status,
- headers=response_headers,
- )
- except (ClientError, TimeoutError) as exc:
- raise ApiProblem(
- "bot_worker_unavailable",
- "机器人工作进程暂时不可用。",
- status=503,
- ) from exc
- class AdminApi:
- def __init__(self) -> None:
- self.username = str(getattr(wbb, "ADMIN_WEB_USERNAME", "admin"))
- self.initial_password = str(
- getattr(wbb, "ADMIN_WEB_INITIAL_PASSWORD", "qwe0.123456")
- )
- self.session_hours = int(getattr(wbb, "ADMIN_WEB_SESSION_HOURS", 12))
- self.cookie_secure = bool(getattr(wbb, "ADMIN_WEB_COOKIE_SECURE", False))
- self.upload_limit = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)) * 1024 * 1024
- self.bot_config_path = Path(
- str(getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"))
- )
- async def initialize(self) -> None:
- await ensure_default_admin(
- username=self.username,
- password_hash=hash_password(self.initial_password),
- )
- await ensure_directory_indexes()
- def register(self, application: web.Application) -> None:
- router = application.router
- router.add_post(
- "/api/internal/v1/directory/verify-platform",
- self.internal_directory_verify_platform,
- )
- router.add_post(
- "/api/internal/v1/directory/verify-local",
- self.internal_directory_verify_local,
- )
- router.add_post(
- "/api/internal/v1/directory/notify-local",
- self.internal_directory_notify_local,
- )
- router.add_get(f"{API_PREFIX}/health", self.health)
- router.add_post(f"{API_PREFIX}/auth/login", self.login)
- router.add_get(f"{API_PREFIX}/auth/me", self.me)
- router.add_post(f"{API_PREFIX}/auth/logout", self.logout)
- router.add_put(f"{API_PREFIX}/auth/password", self.change_password)
- router.add_get(f"{API_PREFIX}/dashboard", self.dashboard)
- router.add_get(f"{API_PREFIX}/settings", self.system_settings)
- router.add_get(f"{API_PREFIX}/settings/telegram", self.telegram_settings)
- router.add_put(f"{API_PREFIX}/settings/telegram", self.telegram_settings_update)
- router.add_get(f"{API_PREFIX}/bots", self.bots)
- router.add_post(f"{API_PREFIX}/bots", self.bot_create)
- router.add_put(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_update)
- router.add_delete(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_delete)
- router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/test", self.bot_test)
- router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/restart", self.bot_restart)
- router.add_get(f"{API_PREFIX}/roles", self.roles)
- router.add_post(f"{API_PREFIX}/roles", self.role_create)
- router.add_put(f"{API_PREFIX}/roles/{{role_id}}", self.role_update)
- router.add_delete(f"{API_PREFIX}/roles/{{role_id}}", self.role_delete)
- router.add_get(f"{API_PREFIX}/audit-logs", self.audit_logs)
- router.add_post(f"{API_PREFIX}/media", self.upload_media)
- router.add_get(f"{API_PREFIX}/chats", self.chats)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}", self.chat)
- router.add_patch(f"{API_PREFIX}/chats/{{chat_id}}/profile", self.chat_profile)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/permissions", self.chat_permissions)
- router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/announcements", self.announcement)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/admins", self.chat_admins)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/search", self.member_search)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/recent", self.recent_members)
- router.add_get(
- f"{API_PREFIX}/chats/{{chat_id}}/members/identity-changes",
- self.member_identity_changes,
- )
- router.add_post(
- f"{API_PREFIX}/chats/{{chat_id}}/members/{{user_id}}/actions",
- self.member_action,
- )
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invites)
- router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_create)
- router.add_delete(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_revoke)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules_update)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation_update)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings_update)
- router.add_get(f"{API_PREFIX}/points/accounts", self.point_accounts)
- router.add_get(f"{API_PREFIX}/points/leaderboard", self.point_leaderboard)
- router.add_get(f"{API_PREFIX}/points/transactions", self.point_transactions)
- router.add_post(f"{API_PREFIX}/points/adjustments", self.point_adjustment)
- router.add_get(f"{API_PREFIX}/points/export", self.points_export)
- router.add_get(f"{API_PREFIX}/directory/settings", self.directory_settings)
- router.add_put(
- f"{API_PREFIX}/directory/settings",
- self.directory_settings_update,
- )
- router.add_get(
- f"{API_PREFIX}/directory/applications",
- self.directory_applications,
- )
- router.add_post(
- f"{API_PREFIX}/directory/applications/{{user_id}}/actions",
- self.directory_application_action,
- )
- router.add_get(f"{API_PREFIX}/directory/profiles", self.directory_profiles)
- router.add_post(
- f"{API_PREFIX}/directory/profiles/{{user_id}}/actions",
- self.directory_profile_action,
- )
- router.add_get(f"{API_PREFIX}/directory/locations", self.directory_locations)
- router.add_put(
- f"{API_PREFIX}/directory/locations/{{user_id}}",
- self.directory_location_update,
- )
- router.add_delete(
- f"{API_PREFIX}/directory/locations/{{user_id}}",
- self.directory_location_delete,
- )
- router.add_get(f"{API_PREFIX}/directory/events", self.directory_events)
- router.add_get(f"{API_PREFIX}/giveaways", self.giveaways)
- router.add_post(f"{API_PREFIX}/giveaways", self.giveaway_create)
- router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}", self.giveaway_detail)
- router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/finish", self.giveaway_finish)
- router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/cancel", self.giveaway_cancel)
- router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/reroll", self.giveaway_reroll)
- router.add_get(
- f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants",
- self.giveaway_participants,
- )
- router.add_delete(
- f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants/{{user_id}}",
- self.giveaway_participant_remove,
- )
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_bans)
- router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_ban_add)
- router.add_delete(
- f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans/{{user_id}}",
- self.giveaway_ban_remove,
- )
- router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}/export", self.giveaway_export)
- async def health(self, _: web.Request) -> web.Response:
- return success({"status": "ok"})
- async def internal_directory_verify_platform(
- self,
- request: web.Request,
- ) -> web.Response:
- _require_internal_request(request)
- body = await json_body(request)
- user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
- require_admin = bool(body.get("require_admin"))
- candidates = await list_membership_candidates(user_id)
- if not candidates:
- return success(
- {"allowed": False, "reason": "membership_evidence_required"}
- )
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is None:
- local = [
- int(item["chat_id"])
- for item in candidates
- if str(item.get("bot_id")) == str(getattr(wbb, "BOT_PROFILE_ID", "primary"))
- ]
- return success(
- await verify_local_membership(
- user_id=user_id,
- chat_ids=local,
- require_admin=require_admin,
- )
- )
- candidates_by_bot: dict[str, list[int]] = {}
- for item in candidates:
- candidates_by_bot.setdefault(str(item["bot_id"]), []).append(
- int(item["chat_id"])
- )
- for bot_id, chat_ids in candidates_by_bot.items():
- endpoint = supervisor.endpoint_for(bot_id)
- if endpoint is None:
- continue
- try:
- async with ClientSession(timeout=ClientTimeout(total=15)) as session:
- async with session.post(
- f"{endpoint}/api/internal/v1/directory/verify-local",
- headers={
- "Authorization": f"Bearer {supervisor.internal_token}"
- },
- json={
- "user_id": str(user_id),
- "chat_ids": [str(value) for value in chat_ids],
- "require_admin": require_admin,
- },
- ) as response:
- if response.status != 200:
- continue
- result = (await response.json()).get("data") or {}
- if result.get("allowed"):
- return success(result)
- except (ClientError, TimeoutError, ValueError):
- continue
- return success({"allowed": False})
- async def internal_directory_verify_local(
- self,
- request: web.Request,
- ) -> web.Response:
- _require_internal_request(request)
- body = await json_body(request)
- user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
- raw_chat_ids = body.get("chat_ids")
- if not isinstance(raw_chat_ids, list):
- raise ApiProblem("invalid_parameter", "chat_ids 必须是数组。")
- chat_ids = [parse_int(value, "chat_id") for value in raw_chat_ids[:200]]
- return success(
- await verify_local_membership(
- user_id=user_id,
- chat_ids=chat_ids,
- require_admin=bool(body.get("require_admin")),
- )
- )
- async def internal_directory_notify_local(
- self,
- request: web.Request,
- ) -> web.Response:
- _require_internal_request(request)
- body = await json_body(request)
- user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
- text = str(body.get("text") or "").strip()
- if not text:
- raise ApiProblem("invalid_parameter", "通知内容不能为空。")
- try:
- await telegram_app.send_message(user_id, text)
- except Exception as exc:
- raise ApiProblem(
- "telegram_notification_failed",
- "无法向该用户发送 Telegram 通知。",
- status=409,
- ) from exc
- return success({"sent": True})
- async def _notify_directory_user(
- self,
- *,
- profile: dict[str, Any],
- text: str,
- ) -> None:
- bot_id = str(
- profile.get("application_source_bot_id")
- or getattr(wbb, "BOT_PROFILE_ID", "primary")
- )
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is None:
- with suppress(Exception):
- await telegram_app.send_message(int(profile["user_id"]), text)
- return
- endpoint = supervisor.endpoint_for(bot_id)
- if endpoint is None:
- return
- with suppress(Exception):
- async with ClientSession(timeout=ClientTimeout(total=10)) as session:
- await session.post(
- f"{endpoint}/api/internal/v1/directory/notify-local",
- headers={
- "Authorization": f"Bearer {supervisor.internal_token}"
- },
- json={"user_id": str(profile["user_id"]), "text": text},
- )
- async def directory_settings(self, _: web.Request) -> web.Response:
- settings = await get_directory_settings()
- settings["nominatim_ready"] = bool(
- settings.get("nominatim_enabled")
- and str(settings.get("nominatim_contact") or "").strip()
- )
- settings["attribution"] = "© OpenStreetMap contributors"
- settings["attribution_url"] = "https://www.openstreetmap.org/copyright"
- return success(settings)
- async def directory_settings_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- values = {
- key: body[key]
- for key in (
- "nominatim_enabled",
- "nominatim_endpoint",
- "nominatim_contact",
- "nominatim_user_agent",
- "nearby_radius_km",
- "nearby_radius_options_km",
- )
- if key in body
- }
- saved = await set_directory_settings(values)
- set_audit(
- request,
- "directory.settings.update",
- summary="更新师生目录和 Nominatim 设置",
- )
- return success(
- {
- **saved,
- "nominatim_ready": bool(
- saved.get("nominatim_enabled")
- and str(saved.get("nominatim_contact") or "").strip()
- ),
- "attribution": "© OpenStreetMap contributors",
- "attribution_url": "https://www.openstreetmap.org/copyright",
- }
- )
- async def directory_applications(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_teacher_applications(
- status=str(request.query.get("status") or ""),
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def directory_application_action(
- self,
- request: web.Request,
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- profile = await decide_teacher_application(
- user_id=user_id,
- action=action,
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- source="web",
- reason=reason,
- )
- labels = {"approve": "批准", "reject": "拒绝", "revoke": "撤销"}
- set_audit(
- request,
- f"directory.application.{action}",
- target_id=user_id,
- summary=f"{labels.get(action, action)}老师申请",
- metadata={"reason": reason},
- )
- notification = {
- "approve": "你的老师申请已通过,可以在菜单中自助上榜。",
- "reject": f"你的老师申请未通过。原因:{reason}",
- "revoke": f"你的老师资格已被撤销。原因:{reason}",
- }.get(action)
- if notification:
- await self._notify_directory_user(profile=profile, text=notification)
- return success(profile)
- async def directory_profiles(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- listed_raw = request.query.get("listed")
- listed = (
- listed_raw.lower() in {"1", "true"}
- if listed_raw is not None
- else None
- )
- items, total = await list_directory_profiles(
- query=str(request.query.get("query") or ""),
- application_status=str(
- request.query.get("application_status") or APPLICATION_APPROVED
- ),
- listed=listed,
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def directory_profile_action(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "强制状态操作必须填写原因。")
- profile, _ = await set_teacher_state(
- user_id=user_id,
- action=action,
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- source="web",
- reason=reason,
- )
- set_audit(
- request,
- f"directory.profile.{action}",
- target_id=user_id,
- summary=f"强制修改老师状态:{action}",
- metadata={"reason": reason},
- )
- await self._notify_directory_user(
- profile=profile,
- text=f"管理员已将你的老师状态修改为“{action}”。原因:{reason}",
- )
- return success(profile)
- async def directory_locations(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_directory_locations(
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- safe_items = []
- for item in items:
- profile = item.get("profile") or {}
- safe_items.append(
- {
- "user_id": item["user_id"],
- "longitude": item.get("longitude"),
- "latitude": item.get("latitude"),
- "region": item.get("region"),
- "geocode_consent": item.get("geocode_consent"),
- "geocode_status": item.get("geocode_status"),
- "source": item.get("source"),
- "updated_at": item.get("updated_at"),
- "display_name": profile.get("display_name"),
- "username": profile.get("username"),
- "application_status": profile.get("application_status", "none"),
- }
- )
- return success(
- {
- "items": safe_items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def directory_location_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "修正位置必须填写原因。")
- profile = await get_directory_profile(user_id)
- if not profile:
- raise ApiProblem("profile_not_found", "未找到该成员资料。", status=404)
- previous = await get_directory_location(user_id)
- saved = await save_directory_location(
- user_id=user_id,
- longitude=body.get("longitude"),
- latitude=body.get("latitude"),
- source="web",
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- geocode_consent=(
- previous.get("geocode_consent") if previous else None
- ),
- )
- set_audit(
- request,
- "directory.location.update",
- target_id=user_id,
- summary="修正成员位置",
- metadata={"reason": reason},
- )
- return success(
- {
- "user_id": user_id,
- "longitude": saved["longitude"],
- "latitude": saved["latitude"],
- "region": saved.get("region"),
- "geocode_status": saved.get("geocode_status"),
- "geocode_consent": saved.get("geocode_consent"),
- "updated_at": saved.get("updated_at"),
- }
- )
- async def directory_location_delete(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "清除位置必须填写原因。")
- deleted = await clear_directory_location(
- user_id=user_id,
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- source="web",
- reason=reason,
- )
- set_audit(
- request,
- "directory.location.delete",
- target_id=user_id,
- summary="清除成员位置",
- metadata={"reason": reason},
- )
- return success({"deleted": deleted})
- async def directory_events(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_directory_events(
- query=str(request.query.get("query") or ""),
- event_type=str(request.query.get("event_type") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- def _set_session_cookie(self, response: web.Response, token: str) -> None:
- response.set_cookie(
- SESSION_COOKIE,
- token,
- httponly=True,
- secure=self.cookie_secure,
- samesite="Strict",
- max_age=self.session_hours * 3600,
- path="/",
- )
- async def _new_session(self, request: web.Request, username: str) -> tuple[str, str]:
- token = generate_session_token()
- csrf = generate_csrf_token()
- await create_admin_session(
- username=username,
- token_hash=hash_token(token),
- csrf_token=csrf,
- expires_at=datetime.now(UTC) + timedelta(hours=self.session_hours),
- remote_address=request.remote or "",
- user_agent=request.headers.get("User-Agent", ""),
- )
- return token, csrf
- async def login(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- username = str(body.get("username") or "").strip()
- password = str(body.get("password") or "")
- set_audit(request, "auth.login", summary=f"Login as {username}")
- user = await get_admin_user(username)
- if user and int(user.get("failed_login_count", 0)) >= 5:
- last_failed = user.get("last_failed_login_at")
- if last_failed and datetime.now(UTC) - as_utc(last_failed) < timedelta(minutes=15):
- raise ApiProblem(
- "login_rate_limited", "登录失败次数过多,请 15 分钟后再试。", status=429
- )
- if not user or not verify_password(password, user.get("password_hash", "")):
- await record_login_failure(username)
- raise ApiProblem("invalid_credentials", "用户名或密码错误。", status=401)
- await record_login_success(username)
- token, csrf = await self._new_session(request, username)
- response = success(
- {
- "username": username,
- "must_change_password": bool(user.get("must_change_password")),
- "csrf_token": csrf,
- }
- )
- self._set_session_cookie(response, token)
- request["admin"] = {"username": username}
- return response
- async def me(self, request: web.Request) -> web.Response:
- return success(
- {
- "username": request["admin"]["username"],
- "must_change_password": request["admin"]["must_change_password"],
- "csrf_token": request["admin"]["csrf_token"],
- }
- )
- async def logout(self, request: web.Request) -> web.Response:
- set_audit(request, "auth.logout", summary="Logout")
- await revoke_admin_session(request["admin"]["session_token_hash"])
- response = success({"logged_out": True})
- response.del_cookie(SESSION_COOKIE, path="/")
- return response
- async def change_password(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- current = str(body.get("current_password") or "")
- new_password = str(body.get("new_password") or "")
- username = request["admin"]["username"]
- user = await get_admin_user(username)
- if not user or not verify_password(current, user["password_hash"]):
- raise ApiProblem("invalid_current_password", "当前密码错误。", status=409)
- if verify_password(new_password, user["password_hash"]):
- raise ApiProblem("password_unchanged", "新密码不能与当前密码相同。", status=409)
- errors = validate_new_password(new_password, username)
- if errors:
- raise ApiProblem("weak_password", "新密码不符合要求。", details=errors)
- set_audit(request, "auth.password.change", summary="Change administrator password")
- await update_admin_password(username, hash_password(new_password))
- token, csrf = await self._new_session(request, username)
- response = success(
- {"username": username, "must_change_password": False, "csrf_token": csrf}
- )
- self._set_session_cookie(response, token)
- return response
- async def dashboard(self, _: web.Request) -> web.Response:
- counts = await dashboard_counts()
- giveaways, giveaway_total = await list_giveaways_page(
- page=1,
- page_size=5,
- all_bots=bool(getattr(wbb, "SUPERVISOR_MODE", False)),
- )
- accounts, account_total = await list_point_accounts(page=1, page_size=5)
- transactions, transaction_total = await list_point_transactions(page=1, page_size=8)
- return success(
- {
- "counts": {
- **counts,
- "giveaways": giveaway_total,
- "point_accounts": account_total,
- "point_transactions": transaction_total,
- },
- "recent_giveaways": giveaways,
- "top_accounts": accounts,
- "recent_transactions": transactions,
- }
- )
- async def system_settings(self, _: web.Request) -> web.Response:
- bot_connected = bool(getattr(wbb, "TELEGRAM_CONNECTED", False))
- return success(
- {
- "web_address": f"http://127.0.0.1:{getattr(wbb, 'ADMIN_WEB_PORT', 8088)}/admin",
- "bind_host": str(getattr(wbb, "ADMIN_WEB_HOST", "0.0.0.0")),
- "upload_limit_mb": int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)),
- "supervisor_mode": bool(getattr(wbb, "SUPERVISOR_MODE", False)),
- "bot": {
- "connected": bot_connected,
- "id": str(wbb.BOT_ID) if bot_connected else "",
- "username": wbb.BOT_USERNAME if bot_connected else "",
- "name": wbb.BOT_NAME if bot_connected else "",
- },
- }
- )
- def _supervisor(self):
- return getattr(wbb, "BOT_SUPERVISOR", None)
- def _telegram_status(self) -> dict[str, Any]:
- supervisor = self._supervisor()
- runtimes = supervisor.runtimes() if supervisor is not None else {}
- return telegram_config_status(self.bot_config_path, runtimes=runtimes)
- async def telegram_settings(self, _: web.Request) -> web.Response:
- return success(self._telegram_status())
- async def telegram_settings_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "telegram.settings.update", summary="Update Telegram API credentials")
- body.pop("confirm", None)
- update_telegram_config(self.bot_config_path, body)
- supervisor = self._supervisor()
- if supervisor is not None:
- status = telegram_config_status(self.bot_config_path)
- for profile in status["bots"]:
- if profile["enabled"] and profile["ready_to_connect"]:
- await supervisor.reconcile(profile["bot_id"])
- return success(self._telegram_status())
- async def bots(self, _: web.Request) -> web.Response:
- return success({"items": self._telegram_status()["bots"]})
- async def bot_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "bot.create", summary=str(body.get("label") or ""))
- body.pop("confirm", None)
- profile = create_bot_profile(self.bot_config_path, body)
- supervisor = self._supervisor()
- if supervisor is not None and profile["enabled"] and profile["ready_to_connect"]:
- await supervisor.start(profile["bot_id"])
- current = next(
- item
- for item in self._telegram_status()["bots"]
- if item["bot_id"] == profile["bot_id"]
- )
- return success(current, status=201)
- async def bot_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- set_audit(request, "bot.update", target_id=bot_id, summary="Update Bot profile")
- body.pop("confirm", None)
- update_bot_profile(self.bot_config_path, bot_id, body)
- supervisor = self._supervisor()
- if supervisor is not None:
- await supervisor.reconcile(bot_id)
- current = next(
- item for item in self._telegram_status()["bots"] if item["bot_id"] == bot_id
- )
- return success(current)
- async def bot_delete(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
- set_audit(request, "bot.delete", target_id=bot_id, summary=str(profile.get("label") or ""))
- supervisor = self._supervisor()
- if supervisor is not None:
- await supervisor.stop(bot_id)
- delete_bot_profile(self.bot_config_path, bot_id)
- return success({"deleted": True})
- async def bot_test(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
- set_audit(request, "bot.test", target_id=bot_id, summary="Validate Bot Token")
- identity = await test_bot_token(str(profile.get("bot_token") or ""))
- store_bot_identity(self.bot_config_path, bot_id, identity)
- return success(identity)
- async def bot_restart(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- set_audit(request, "bot.restart", target_id=bot_id, summary="Restart Bot worker")
- supervisor = self._supervisor()
- if supervisor is None:
- raise ApiProblem(
- "bot_supervisor_unavailable",
- "机器人监管服务尚未就绪。",
- status=503,
- )
- return success(await supervisor.restart(bot_id))
- async def roles(self, _: web.Request) -> web.Response:
- status = self._telegram_status()
- return success(
- {
- "items": status["roles"],
- "permission_catalog": status["permission_catalog"],
- }
- )
- async def role_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- set_audit(
- request,
- "bot.role.create",
- summary=str(body.get("name") or ""),
- )
- body.pop("confirm", None)
- return success(create_bot_role(self.bot_config_path, body), status=201)
- async def role_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- role_id = request.match_info["role_id"]
- set_audit(
- request,
- "bot.role.update",
- target_id=role_id,
- summary=str(body.get("name") or ""),
- )
- body.pop("confirm", None)
- role = update_bot_role(self.bot_config_path, role_id, body)
- supervisor = self._supervisor()
- if supervisor is not None:
- for profile in self._telegram_status()["bots"]:
- if role_id in profile["role_ids"]:
- await supervisor.reconcile(profile["bot_id"])
- return success(role)
- async def role_delete(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- role_id = request.match_info["role_id"]
- set_audit(
- request,
- "bot.role.delete",
- target_id=role_id,
- summary="Delete Bot role",
- )
- delete_bot_role(self.bot_config_path, role_id)
- return success({"deleted": True})
- async def audit_logs(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- items, total = await list_audit_logs(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- action=request.query.get("action") or None,
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def upload_media(self, request: web.Request) -> web.Response:
- if not MESSAGE_DUMP_CHAT:
- raise ApiProblem(
- "message_dump_chat_missing",
- "请先配置媒体中转群。",
- status=409,
- )
- reader = await request.multipart()
- part = await reader.next()
- if not part or part.name != "file":
- raise ApiProblem("file_required", "请选择要上传的文件。")
- filename = Path(part.filename or "upload.bin").name
- content_type = part.headers.get("Content-Type", "application/octet-stream")
- size = 0
- temporary_path: Path | None = None
- try:
- with tempfile.NamedTemporaryFile(delete=False, suffix=Path(filename).suffix) as handle:
- temporary_path = Path(handle.name)
- while True:
- chunk = await part.read_chunk(1024 * 1024)
- if not chunk:
- break
- size += len(chunk)
- if size > self.upload_limit:
- raise ApiProblem("file_too_large", "上传文件超过 20 MB 限制。", status=413)
- handle.write(chunk)
- if content_type == "image/gif":
- media_type = "animation"
- sent = await telegram_app.send_animation(MESSAGE_DUMP_CHAT, str(temporary_path))
- file_id = sent.animation.file_id
- elif content_type.startswith("image/"):
- media_type = "photo"
- sent = await telegram_app.send_photo(MESSAGE_DUMP_CHAT, str(temporary_path))
- file_id = sent.photo.file_id
- elif content_type.startswith("video/"):
- media_type = "video"
- sent = await telegram_app.send_video(MESSAGE_DUMP_CHAT, str(temporary_path))
- file_id = sent.video.file_id
- else:
- media_type = "document"
- sent = await telegram_app.send_document(
- MESSAGE_DUMP_CHAT, str(temporary_path), file_name=filename
- )
- file_id = sent.document.file_id
- finally:
- if temporary_path:
- temporary_path.unlink(missing_ok=True)
- set_audit(request, "media.upload", summary=filename, target_id=file_id)
- return success(
- {
- "file_id": file_id,
- "type": media_type,
- "filename": filename,
- "size": size,
- },
- status=201,
- )
- async def chats(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_accessible_chats(
- query=request.query.get("query", ""), page=page, page_size=page_size
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def chat(self, request: web.Request) -> web.Response:
- return success(await get_chat_overview(chat_id_param(request)))
- async def chat_profile(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.profile.update", chat_id=chat_id, summary="Update group profile")
- return success(
- await update_chat_profile(
- chat_id, title=body.get("title"), description=body.get("description")
- )
- )
- async def chat_permissions(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.permissions.update", chat_id=chat_id, summary="Update default member permissions")
- return success(await update_chat_permissions(chat_id, body.get("permissions") or {}))
- async def announcement(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.announcement.send", chat_id=chat_id, summary=str(body.get("text") or "")[:200])
- return success(
- await send_announcement(
- chat_id,
- text=str(body.get("text") or ""),
- media_type=body.get("media_type"),
- file_id=body.get("file_id"),
- pin=bool(body.get("pin")),
- ),
- status=201,
- )
- async def chat_admins(self, request: web.Request) -> web.Response:
- return success({"items": await list_chat_admins(chat_id_param(request))})
- async def member_search(self, request: web.Request) -> web.Response:
- return success(
- {
- "items": await search_chat_members(
- chat_id_param(request),
- request.query.get("query", ""),
- limit=parse_int(request.query.get("limit", "20"), "limit"),
- )
- }
- )
- async def recent_members(self, request: web.Request) -> web.Response:
- return success(
- {
- "items": await list_recent_members(
- chat_id_param(request),
- limit=parse_int(request.query.get("limit", "30"), "limit"),
- )
- }
- )
- async def member_identity_changes(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_member_identity_changes(
- chat_id_param(request), page=page, page_size=page_size
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def member_action(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- user_id = parse_int(request.match_info["user_id"], "userId")
- body = await json_body(request)
- require_confirmation(body)
- action = str(body.get("action") or "")
- set_audit(
- request,
- f"chat.member.{action}",
- chat_id=chat_id,
- target_id=user_id,
- summary=str(body.get("reason") or ""),
- )
- duration = body.get("duration_seconds")
- return success(
- await execute_member_action(
- chat_id,
- user_id=user_id,
- action=action,
- reason=str(body.get("reason") or ""),
- duration_seconds=parse_int(duration, "duration_seconds", minimum=60) if duration else None,
- privileges=body.get("privileges") or {},
- )
- )
- async def invites(self, request: web.Request) -> web.Response:
- return success({"items": await get_invite_links(chat_id_param(request))})
- async def invite_create(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- expires_at = parse_datetime(body["expires_at"], "expires_at") if body.get("expires_at") else None
- member_limit = parse_int(body["member_limit"], "member_limit", minimum=1) if body.get("member_limit") else None
- set_audit(request, "chat.invite.create", chat_id=chat_id, summary=str(body.get("name") or ""))
- return success(
- await create_invite_link(
- chat_id,
- name=str(body.get("name") or "Admin panel"),
- expires_at=expires_at,
- member_limit=member_limit,
- ),
- status=201,
- )
- async def invite_revoke(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- invite_link = str(body.get("invite_link") or "")
- if not invite_link:
- raise ApiProblem("invite_link_required", "缺少邀请链接。")
- set_audit(request, "chat.invite.revoke", chat_id=chat_id, summary=invite_link)
- await revoke_invite_link(chat_id, invite_link)
- return success({"revoked": True})
- async def rules(self, request: web.Request) -> web.Response:
- return success({"rules": await get_rules(chat_id_param(request))})
- async def rules_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- await ensure_permission(chat_id, "can_change_info")
- rules = str(body.get("rules") or "")
- if len(rules) > 4000:
- raise ApiProblem("rules_too_long", "群规不能超过 4000 个字符。")
- set_audit(request, "chat.rules.update", chat_id=chat_id, summary="Update group rules")
- await set_chat_rules(chat_id, rules)
- return success({"rules": rules})
- async def automation(self, request: web.Request) -> web.Response:
- return success(await get_automation_settings(chat_id_param(request)))
- async def automation_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.automation.update", chat_id=chat_id, summary="Update automation rules")
- body.pop("confirm", None)
- return success(await apply_automation_settings(chat_id, body))
- async def points_settings(self, request: web.Request) -> web.Response:
- return success(await get_point_rules(chat_id_param(request)))
- async def points_settings_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- await ensure_permission(chat_id, "can_change_info")
- set_audit(request, "points.settings.update", chat_id=chat_id, summary="Update point rules")
- body.pop("confirm", None)
- return success(await apply_point_rules(chat_id, body))
- async def point_accounts(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- items, total = await list_point_accounts(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- query=request.query.get("query", ""),
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def point_leaderboard(self, request: web.Request) -> web.Response:
- if not request.query.get("chat_id"):
- raise ApiProblem("chat_id_required", "排行榜必须选择群组。")
- return await self.point_accounts(request)
- async def point_transactions(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- user_raw = request.query.get("user_id")
- created_from = parse_datetime(request.query["created_from"], "created_from") if request.query.get("created_from") else None
- created_to = parse_datetime(request.query["created_to"], "created_to") if request.query.get("created_to") else None
- items, total = await list_point_transactions(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- user_id=parse_int(user_raw, "user_id") if user_raw else None,
- source=request.query.get("source") or None,
- created_from=created_from,
- created_to=created_to,
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def point_adjustment(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = parse_int(body.get("chat_id"), "chat_id")
- user_id = parse_int(body.get("user_id"), "user_id")
- operation = str(body.get("operation") or "")
- amount = parse_int(body.get("amount"), "amount", minimum=0)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "积分调整必须填写原因。")
- if operation not in {"add", "deduct", "set"}:
- raise ApiProblem(
- "invalid_operation",
- "积分操作必须是增加、扣减或设置余额。",
- )
- try:
- user = await telegram_app.get_users(user_id)
- username, first_name = user.username, user.first_name
- display_name = " ".join(
- value for value in (user.first_name, user.last_name) if value
- )
- except Exception:
- username, first_name = body.get("username"), body.get("first_name")
- display_name = body.get("display_name") or first_name
- request_key = str(body.get("request_id") or generate_session_token())
- set_audit(
- request,
- f"points.adjust.{operation}",
- chat_id=chat_id,
- target_id=user_id,
- summary=reason,
- metadata={"amount": amount},
- )
- if operation == "set":
- account, created = await set_points(
- chat_id=chat_id,
- user_id=user_id,
- balance=amount,
- actor_id=request["admin"]["username"],
- reason=reason,
- idempotency_key=f"web-points:{request_key}",
- username=username,
- first_name=first_name,
- display_name=display_name,
- )
- else:
- account, created = await adjust_points(
- chat_id=chat_id,
- user_id=user_id,
- delta=amount if operation == "add" else -amount,
- source=SOURCE_ADMIN,
- idempotency_key=f"web-points:{request_key}",
- actor_id=request["admin"]["username"],
- reason=reason,
- username=username,
- first_name=first_name,
- display_name=display_name,
- )
- return success({"account": account, "created": created}, status=201 if created else 200)
- async def points_export(self, request: web.Request) -> web.Response:
- chat_raw = request.query.get("chat_id")
- if not chat_raw:
- raise ApiProblem("chat_id_required", "导出积分数据必须选择群组。")
- chat_id = parse_int(chat_raw, "chat_id")
- kind = request.query.get("kind", "accounts")
- output = io.StringIO()
- if kind == "accounts":
- items: list[dict[str, Any]] = []
- page = 1
- while True:
- batch, total = await list_point_accounts(
- chat_id=chat_id, page=page, page_size=100
- )
- items.extend(batch)
- if not batch or len(items) >= total:
- break
- page += 1
- fieldnames = ["chat_id", "user_id", "display_name", "username", "first_name", "balance", "lifetime_earned", "lifetime_spent", "updated_at"]
- elif kind == "transactions":
- items = []
- page = 1
- while True:
- batch, total = await list_point_transactions(
- chat_id=chat_id, page=page, page_size=100
- )
- items.extend(batch)
- if not batch or len(items) >= total:
- break
- page += 1
- fieldnames = ["transaction_id", "chat_id", "user_id", "display_name", "username", "first_name", "delta", "balance_after", "source", "actor_id", "reason", "reference_id", "created_at"]
- else:
- raise ApiProblem(
- "invalid_export_kind",
- "导出类型必须是积分账户或积分流水。",
- )
- writer = csv.DictWriter(output, fieldnames=fieldnames, extrasaction="ignore")
- writer.writeheader()
- for item in items:
- writer.writerow(jsonable(item))
- return web.Response(
- text=output.getvalue(),
- content_type="text/csv",
- headers={"Content-Disposition": f'attachment; filename="points-{chat_id}-{kind}.csv"'},
- )
- async def giveaways(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- items, total = await list_giveaways_page(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- status=request.query.get("status") or None,
- query=request.query.get("query", ""),
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def giveaway_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = parse_int(body.get("chat_id"), "chat_id")
- await ensure_permission(chat_id, "can_change_info")
- prizes = body.get("prizes")
- if not isinstance(prizes, list):
- raise ApiProblem("invalid_prizes", "奖项必须是数组。")
- set_audit(request, "giveaway.create", chat_id=chat_id, summary=str(body.get("title") or ""))
- giveaway = await create_and_publish_giveaway(
- chat_id=chat_id,
- creator_id=0,
- creator_name=request["admin"]["username"],
- title=str(body.get("title") or ""),
- description=str(body.get("description") or ""),
- prizes=prizes,
- ends_at=parse_datetime(body.get("ends_at"), "ends_at"),
- minimum_points=parse_int(body.get("minimum_points", 0), "minimum_points", minimum=0),
- entry_cost=parse_int(body.get("entry_cost", 0), "entry_cost", minimum=0),
- participation_reward=parse_int(body.get("participation_reward", 0), "participation_reward", minimum=0),
- )
- return success(giveaway, status=201)
- async def giveaway_detail(self, request: web.Request) -> web.Response:
- giveaway = await get_giveaway(request.match_info["giveaway_id"])
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- return success(giveaway)
- async def giveaway_finish(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.finish", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Manual draw")
- ok, message, current = await finish_and_publish_giveaway(giveaway_id)
- if not ok:
- raise ApiProblem("giveaway_not_running", message, status=409)
- return success(current)
- async def giveaway_cancel(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.cancel", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Cancel and refund")
- ok, message, current = await cancel_and_refund_giveaway(
- giveaway_id, chat_id=int(giveaway["chat_id"])
- )
- if not ok:
- raise ApiProblem("giveaway_not_running", message, status=409)
- return success(current)
- async def giveaway_reroll(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.reroll", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary=str(body.get("tier_name") or "All tiers"))
- winners, reroll_id = await reroll_giveaway(
- giveaway_id=giveaway_id,
- moderator_id=0,
- tier_name=body.get("tier_name") or None,
- reroll_id=body.get("request_id") or None,
- )
- return success({"winners": winners, "reroll_id": reroll_id})
- async def giveaway_participants(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_participants_page(
- giveaway_id=request.match_info["giveaway_id"],
- active_only=request.query.get("active_only", "false").lower() == "true",
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def giveaway_participant_remove(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- user_id = parse_int(request.match_info["user_id"], "user_id")
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.participant.remove", chat_id=int(giveaway["chat_id"]), target_id=user_id, summary=str(body.get("reason") or ""))
- removed = await remove_and_optionally_refund_participant(
- giveaway_id=giveaway_id,
- user_id=user_id,
- moderator_id=0,
- reason=str(body.get("reason") or "Removed in admin panel"),
- refund=bool(body.get("refund", True)),
- )
- if not removed:
- raise ApiProblem("participant_not_found", "参与者不存在或已被移除。", status=404)
- return success({"removed": True, "refunded": bool(body.get("refund", True))})
- async def giveaway_bans(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_giveaway_bans(
- chat_id=chat_id_param(request), page=page, page_size=page_size
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def giveaway_ban_add(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(body.get("user_id"), "user_id")
- await ensure_permission(chat_id, "can_change_info")
- set_audit(request, "giveaway.ban.add", chat_id=chat_id, target_id=user_id, summary=str(body.get("reason") or ""))
- await add_giveaway_ban(chat_id=chat_id, user_id=user_id, moderator_id=0, reason=str(body.get("reason") or ""))
- return success({"created": True}, status=201)
- async def giveaway_ban_remove(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- user_id = parse_int(request.match_info["user_id"], "user_id")
- body = await json_body(request)
- require_confirmation(body)
- await ensure_permission(chat_id, "can_change_info")
- set_audit(request, "giveaway.ban.remove", chat_id=chat_id, target_id=user_id)
- return success({"removed": await remove_giveaway_ban(chat_id, user_id)})
- async def giveaway_export(self, request: web.Request) -> web.Response:
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- participants: list[dict[str, Any]] = []
- page = 1
- while True:
- batch, total = await list_participants_page(
- giveaway_id=giveaway_id, page=page, page_size=100
- )
- participants.extend(batch)
- if not batch or len(participants) >= total:
- break
- page += 1
- output = io.StringIO()
- fields = ["giveaway_id", "chat_id", "user_id", "display_name", "username", "first_name", "active", "entry_cost", "joined_at", "removed_at", "refunded_at"]
- writer = csv.DictWriter(output, fieldnames=fields, extrasaction="ignore")
- writer.writeheader()
- for participant in participants:
- writer.writerow(jsonable(participant))
- return web.Response(
- text=output.getvalue(),
- content_type="text/csv",
- headers={"Content-Disposition": f'attachment; filename="giveaway-{giveaway_id}.csv"'},
- )
- async def _start_global_admin_services(_: web.Application) -> None:
- if not bool(getattr(wbb, "BOT_WORKER_MODE", False)):
- await start_directory_geocoder()
- async def _stop_global_admin_services(_: web.Application) -> None:
- if not bool(getattr(wbb, "BOT_WORKER_MODE", False)):
- await stop_directory_geocoder()
- def build_admin_application() -> web.Application:
- max_upload_mb = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20))
- application = web.Application(
- middlewares=[
- api_error_middleware,
- authentication_middleware,
- bot_role_middleware,
- bot_proxy_middleware,
- ],
- client_max_size=(max_upload_mb + 1) * 1024 * 1024,
- )
- admin_api = AdminApi()
- application["admin_api"] = admin_api
- admin_api.register(application)
- application.on_startup.append(_start_global_admin_services)
- application.on_cleanup.append(_stop_global_admin_services)
- return application
|