api.py 107 KB

1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006100710081009101010111012101310141015101610171018101910201021102210231024102510261027102810291030103110321033103410351036103710381039104010411042104310441045104610471048104910501051105210531054105510561057105810591060106110621063106410651066106710681069107010711072107310741075107610771078107910801081108210831084108510861087108810891090109110921093109410951096109710981099110011011102110311041105110611071108110911101111111211131114111511161117111811191120112111221123112411251126112711281129113011311132113311341135113611371138113911401141114211431144114511461147114811491150115111521153115411551156115711581159116011611162116311641165116611671168116911701171117211731174117511761177117811791180118111821183118411851186118711881189119011911192119311941195119611971198119912001201120212031204120512061207120812091210121112121213121412151216121712181219122012211222122312241225122612271228122912301231123212331234123512361237123812391240124112421243124412451246124712481249125012511252125312541255125612571258125912601261126212631264126512661267126812691270127112721273127412751276127712781279128012811282128312841285128612871288128912901291129212931294129512961297129812991300130113021303130413051306130713081309131013111312131313141315131613171318131913201321132213231324132513261327132813291330133113321333133413351336133713381339134013411342134313441345134613471348134913501351135213531354135513561357135813591360136113621363136413651366136713681369137013711372137313741375137613771378137913801381138213831384138513861387138813891390139113921393139413951396139713981399140014011402140314041405140614071408140914101411141214131414141514161417141814191420142114221423142414251426142714281429143014311432143314341435143614371438143914401441144214431444144514461447144814491450145114521453145414551456145714581459146014611462146314641465146614671468146914701471147214731474147514761477147814791480148114821483148414851486148714881489149014911492149314941495149614971498149915001501150215031504150515061507150815091510151115121513151415151516151715181519152015211522152315241525152615271528152915301531153215331534153515361537153815391540154115421543154415451546154715481549155015511552155315541555155615571558155915601561156215631564156515661567156815691570157115721573157415751576157715781579158015811582158315841585158615871588158915901591159215931594159515961597159815991600160116021603160416051606160716081609161016111612161316141615161616171618161916201621162216231624162516261627162816291630163116321633163416351636163716381639164016411642164316441645164616471648164916501651165216531654165516561657165816591660166116621663166416651666166716681669167016711672167316741675167616771678167916801681168216831684168516861687168816891690169116921693169416951696169716981699170017011702170317041705170617071708170917101711171217131714171517161717171817191720172117221723172417251726172717281729173017311732173317341735173617371738173917401741174217431744174517461747174817491750175117521753175417551756175717581759176017611762176317641765176617671768176917701771177217731774177517761777177817791780178117821783178417851786178717881789179017911792179317941795179617971798179918001801180218031804180518061807180818091810181118121813181418151816181718181819182018211822182318241825182618271828182918301831183218331834183518361837183818391840184118421843184418451846184718481849185018511852185318541855185618571858185918601861186218631864186518661867186818691870187118721873187418751876187718781879188018811882188318841885188618871888188918901891189218931894189518961897189818991900190119021903190419051906190719081909191019111912191319141915191619171918191919201921192219231924192519261927192819291930193119321933193419351936193719381939194019411942194319441945194619471948194919501951195219531954195519561957195819591960196119621963196419651966196719681969197019711972197319741975197619771978197919801981198219831984198519861987198819891990199119921993199419951996199719981999200020012002200320042005200620072008200920102011201220132014201520162017201820192020202120222023202420252026202720282029203020312032203320342035203620372038203920402041204220432044204520462047204820492050205120522053205420552056205720582059206020612062206320642065206620672068206920702071207220732074207520762077207820792080208120822083208420852086208720882089209020912092209320942095209620972098209921002101210221032104210521062107210821092110211121122113211421152116211721182119212021212122212321242125212621272128212921302131213221332134213521362137213821392140214121422143214421452146214721482149215021512152215321542155215621572158215921602161216221632164216521662167216821692170217121722173217421752176217721782179218021812182218321842185218621872188218921902191219221932194219521962197219821992200220122022203220422052206220722082209221022112212221322142215221622172218221922202221222222232224222522262227222822292230223122322233223422352236223722382239224022412242224322442245224622472248224922502251225222532254225522562257225822592260226122622263226422652266226722682269227022712272227322742275227622772278227922802281228222832284228522862287228822892290229122922293229422952296229722982299230023012302230323042305230623072308230923102311231223132314231523162317231823192320232123222323232423252326232723282329233023312332233323342335233623372338233923402341234223432344234523462347234823492350235123522353235423552356235723582359236023612362236323642365236623672368236923702371237223732374237523762377237823792380238123822383238423852386238723882389239023912392239323942395239623972398239924002401240224032404240524062407240824092410241124122413241424152416241724182419242024212422242324242425242624272428242924302431243224332434243524362437243824392440244124422443244424452446244724482449245024512452245324542455245624572458245924602461246224632464246524662467246824692470247124722473247424752476247724782479248024812482248324842485248624872488248924902491249224932494249524962497249824992500250125022503250425052506250725082509251025112512251325142515251625172518251925202521252225232524252525262527252825292530253125322533253425352536253725382539254025412542254325442545254625472548254925502551255225532554255525562557255825592560256125622563256425652566256725682569257025712572257325742575257625772578257925802581258225832584258525862587258825892590259125922593259425952596259725982599260026012602260326042605260626072608260926102611261226132614261526162617261826192620262126222623262426252626262726282629263026312632263326342635263626372638263926402641264226432644264526462647264826492650265126522653265426552656265726582659266026612662266326642665266626672668266926702671267226732674267526762677267826792680268126822683268426852686268726882689269026912692269326942695269626972698269927002701270227032704270527062707270827092710271127122713271427152716271727182719272027212722272327242725272627272728272927302731273227332734273527362737273827392740274127422743274427452746274727482749275027512752275327542755275627572758275927602761
  1. from __future__ import annotations
  2. import csv
  3. import io
  4. import secrets
  5. import tempfile
  6. import traceback
  7. from contextlib import suppress
  8. from datetime import UTC, datetime, timedelta
  9. from pathlib import Path
  10. from typing import Any
  11. from aiohttp import ClientError, ClientSession, ClientTimeout, web
  12. import wbb
  13. from wbb import MESSAGE_DUMP_CHAT
  14. from wbb import app as telegram_app
  15. from wbb.admin.bot_config import (
  16. BotConfigError,
  17. create_bot_profile,
  18. create_bot_role,
  19. delete_bot_profile,
  20. delete_bot_role,
  21. get_bot_profile_secrets,
  22. store_bot_identity,
  23. telegram_config_status,
  24. test_bot_token,
  25. update_bot_profile,
  26. update_bot_role,
  27. update_telegram_config,
  28. )
  29. from wbb.admin.security import (
  30. generate_csrf_token,
  31. generate_session_token,
  32. hash_password,
  33. hash_token,
  34. validate_new_password,
  35. verify_password,
  36. )
  37. from wbb.admin.telegram_webapp import (
  38. TelegramWebAppAuthError,
  39. verify_telegram_webapp_init_data,
  40. )
  41. from wbb.services.bot_permissions import api_permission, has_permission
  42. from wbb.services.business_assistant import (
  43. AssistantProviderError,
  44. runtime_overview,
  45. )
  46. from wbb.services.chat_management import (
  47. ChatManagementError,
  48. apply_automation_settings,
  49. create_invite_link,
  50. ensure_permission,
  51. execute_member_action,
  52. get_automation_settings,
  53. get_chat_overview,
  54. get_invite_links,
  55. list_accessible_chats,
  56. list_chat_admins,
  57. list_recent_members,
  58. revoke_invite_link,
  59. search_chat_members,
  60. send_announcement,
  61. update_chat_permissions,
  62. update_chat_profile,
  63. )
  64. from wbb.services.directory import DirectoryServiceError, verify_local_membership
  65. from wbb.services.giveaways import (
  66. GiveawayServiceError,
  67. cancel_and_refund_giveaway,
  68. create_and_publish_giveaway,
  69. finish_and_publish_giveaway,
  70. remove_and_optionally_refund_participant,
  71. reroll_giveaway,
  72. )
  73. from wbb.services.point_settings import apply_point_rules
  74. from wbb.services.technician_reviews import (
  75. ensure_technician_review_topic,
  76. publish_approved_review,
  77. )
  78. from wbb.utils.dbadmin import (
  79. create_admin_session,
  80. dashboard_counts,
  81. ensure_default_admin,
  82. get_admin_session,
  83. get_admin_user,
  84. list_audit_logs,
  85. list_member_identity_changes,
  86. record_audit,
  87. record_login_failure,
  88. record_login_success,
  89. revoke_admin_session,
  90. update_admin_password,
  91. )
  92. from wbb.utils.dbassistant import (
  93. AssistantDataError,
  94. clear_conversation,
  95. close_conversation,
  96. conversation_detail,
  97. create_knowledge_entry,
  98. delete_knowledge_entry,
  99. ensure_assistant_indexes,
  100. get_account_settings,
  101. get_business_connection,
  102. list_business_connections,
  103. list_conversations,
  104. list_knowledge_entries,
  105. pause_conversation,
  106. resume_conversation,
  107. update_account_settings,
  108. update_knowledge_entry,
  109. usage_metrics,
  110. )
  111. from wbb.utils.dbdirectory import (
  112. APPLICATION_APPROVED,
  113. DirectoryDataError,
  114. clear_directory_location,
  115. decide_teacher_application,
  116. ensure_directory_indexes,
  117. get_directory_profile,
  118. get_directory_settings,
  119. list_directory_events,
  120. list_directory_locations,
  121. list_directory_profiles,
  122. list_membership_candidates,
  123. list_teacher_applications,
  124. save_directory_location,
  125. set_directory_settings,
  126. set_teacher_state,
  127. )
  128. from wbb.utils.dbfunctions import get_rules, set_chat_rules
  129. from wbb.utils.dbgiveaway import (
  130. add_giveaway_ban,
  131. get_giveaway,
  132. list_giveaway_bans,
  133. list_giveaways_page,
  134. list_participants_page,
  135. remove_giveaway_ban,
  136. )
  137. from wbb.utils.dbpoints import (
  138. SOURCE_ADMIN,
  139. InsufficientPoints,
  140. PointsError,
  141. adjust_points,
  142. get_point_rules,
  143. list_point_accounts,
  144. list_point_transactions,
  145. set_points,
  146. )
  147. from wbb.utils.dbservice import (
  148. ServiceDataError,
  149. admin_reveal_order_address,
  150. create_package_template,
  151. duplicate_package_template,
  152. ensure_builtin_package_templates,
  153. fulfillment_metrics,
  154. get_active_review_template,
  155. get_admin_service_order,
  156. get_service_settings,
  157. get_technician_self_service_context,
  158. get_technician_service_profile,
  159. list_customer_blocks,
  160. list_leaderboard,
  161. list_package_templates,
  162. list_review_qr_records,
  163. list_reviews,
  164. list_service_orders,
  165. list_service_reports,
  166. moderate_review,
  167. moderate_service_report,
  168. publish_review_template,
  169. publish_technician_profile,
  170. resolve_service_dispute,
  171. set_customer_block,
  172. set_package_template_status,
  173. set_service_settings,
  174. set_technician_accepting_requests,
  175. update_package_template,
  176. )
  177. API_PREFIX = "/api/admin/v1"
  178. TECHNICIAN_API_PREFIX = "/api/technician/v1"
  179. SESSION_COOKIE = "wbb_admin_session"
  180. UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"}
  181. PUBLIC_API_PATHS = {f"{API_PREFIX}/auth/login", f"{API_PREFIX}/health"}
  182. BOT_SCOPED_PREFIXES = (
  183. f"{API_PREFIX}/business-assistant",
  184. f"{API_PREFIX}/chats",
  185. f"{API_PREFIX}/giveaways",
  186. f"{API_PREFIX}/points",
  187. f"{API_PREFIX}/media",
  188. )
  189. TELEGRAM_ID_KEYS = {
  190. "chat_id",
  191. "user_id",
  192. "actor_id",
  193. "creator_id",
  194. "customer_id",
  195. "moderator_id",
  196. "reporter_id",
  197. "target_id",
  198. "technician_id",
  199. "message_id",
  200. "removed_by",
  201. "authorization_chat_id",
  202. "ops_group_id",
  203. }
  204. class ApiProblem(RuntimeError):
  205. def __init__(
  206. self,
  207. code: str,
  208. message: str,
  209. *,
  210. status: int = 400,
  211. details: Any = None,
  212. ):
  213. super().__init__(message)
  214. self.code = code
  215. self.status = status
  216. self.details = details
  217. def as_utc(value: datetime) -> datetime:
  218. if value.tzinfo is None:
  219. return value.replace(tzinfo=UTC)
  220. return value.astimezone(UTC)
  221. def jsonable(value: Any, *, key: str = "") -> Any:
  222. if isinstance(value, datetime):
  223. return as_utc(value).isoformat().replace("+00:00", "Z")
  224. if isinstance(value, dict):
  225. return {
  226. item_key: jsonable(item_value, key=item_key)
  227. for item_key, item_value in value.items()
  228. if item_key not in {"_id", "encrypted_payload", "token_hash"}
  229. }
  230. if isinstance(value, (list, tuple)):
  231. return [jsonable(item) for item in value]
  232. if isinstance(value, int) and (key in TELEGRAM_ID_KEYS or key.endswith("_telegram_id")):
  233. return str(value)
  234. return value
  235. def success(data: Any, *, status: int = 200) -> web.Response:
  236. return web.json_response({"data": jsonable(data)}, status=status)
  237. def error_response(problem: ApiProblem) -> web.Response:
  238. return web.json_response(
  239. {
  240. "error": {
  241. "code": problem.code,
  242. "message": str(problem),
  243. "details": jsonable(problem.details),
  244. }
  245. },
  246. status=problem.status,
  247. )
  248. def parse_int(value: Any, name: str, *, minimum: int | None = None) -> int:
  249. try:
  250. parsed = int(value)
  251. except (TypeError, ValueError) as exc:
  252. raise ApiProblem("invalid_parameter", f"{name} 必须是整数。") from exc
  253. if minimum is not None and parsed < minimum:
  254. raise ApiProblem("invalid_parameter", f"{name} 不能小于 {minimum}。")
  255. return parsed
  256. def parse_datetime(value: Any, name: str) -> datetime:
  257. if not isinstance(value, str) or not value.strip():
  258. raise ApiProblem("invalid_parameter", f"{name} 必须是 UTC ISO 8601 时间。")
  259. try:
  260. parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
  261. except ValueError as exc:
  262. raise ApiProblem("invalid_parameter", f"{name} 不是有效时间。") from exc
  263. if parsed.tzinfo is None:
  264. raise ApiProblem("invalid_parameter", f"{name} 必须包含时区。")
  265. return parsed.astimezone(UTC)
  266. async def json_body(request: web.Request) -> dict[str, Any]:
  267. try:
  268. body = await request.json()
  269. except Exception as exc:
  270. raise ApiProblem("invalid_json", "请求体必须是 JSON。") from exc
  271. if not isinstance(body, dict):
  272. raise ApiProblem("invalid_json", "请求体必须是 JSON 对象。")
  273. return body
  274. def require_confirmation(body: dict[str, Any]) -> None:
  275. if body.get("confirm") is not True:
  276. raise ApiProblem(
  277. "confirmation_required", "该操作需要二次确认。", status=409
  278. )
  279. def page_params(request: web.Request) -> tuple[int, int]:
  280. return (
  281. parse_int(request.query.get("page", 1), "page", minimum=1),
  282. min(parse_int(request.query.get("page_size", 20), "page_size", minimum=1), 100),
  283. )
  284. def chat_id_param(request: web.Request) -> int:
  285. return parse_int(request.match_info["chat_id"], "chatId")
  286. def set_audit(
  287. request: web.Request,
  288. action: str,
  289. *,
  290. chat_id: int | None = None,
  291. target_id: int | str | None = None,
  292. summary: str = "",
  293. metadata: dict[str, Any] | None = None,
  294. ) -> None:
  295. request["audit_context"] = {
  296. "action": action,
  297. "chat_id": chat_id,
  298. "target_id": target_id,
  299. "summary": summary,
  300. "metadata": metadata or {},
  301. }
  302. @web.middleware
  303. async def api_error_middleware(request: web.Request, handler):
  304. try:
  305. response = await handler(request)
  306. except ApiProblem as exc:
  307. await _record_request_audit(request, success_state=False, error=str(exc))
  308. return error_response(exc)
  309. except (ChatManagementError, GiveawayServiceError, DirectoryServiceError) as exc:
  310. problem = ApiProblem(
  311. exc.code, str(exc), status=getattr(exc, "status", 400)
  312. )
  313. await _record_request_audit(request, success_state=False, error=str(exc))
  314. return error_response(problem)
  315. except (DirectoryDataError, ServiceDataError) as exc:
  316. problem = ApiProblem(exc.code, str(exc), status=409)
  317. await _record_request_audit(request, success_state=False, error=str(exc))
  318. return error_response(problem)
  319. except AssistantDataError as exc:
  320. status = 404 if exc.code.endswith("_not_found") else 409
  321. problem = ApiProblem(exc.code, str(exc), status=status)
  322. await _record_request_audit(request, success_state=False, error=str(exc))
  323. return error_response(problem)
  324. except (PointsError, InsufficientPoints) as exc:
  325. problem = ApiProblem("points_error", str(exc), status=409)
  326. await _record_request_audit(request, success_state=False, error=str(exc))
  327. return error_response(problem)
  328. except BotConfigError as exc:
  329. status = {
  330. "bot_not_found": 404,
  331. "role_not_found": 404,
  332. "role_in_use": 409,
  333. "builtin_role_immutable": 409,
  334. }.get(exc.code, 400)
  335. problem = ApiProblem(exc.code, str(exc), status=status)
  336. await _record_request_audit(request, success_state=False, error=str(exc))
  337. return error_response(problem)
  338. except web.HTTPException:
  339. raise
  340. except Exception as exc:
  341. await _record_request_audit(request, success_state=False, error=str(exc))
  342. wbb.log.error(
  343. f"Admin API {request.method} {request.path} failed: {exc}\n"
  344. f"{traceback.format_exc()}"
  345. )
  346. return error_response(
  347. ApiProblem("internal_error", "服务器处理请求失败。", status=500)
  348. )
  349. await _record_request_audit(request, success_state=response.status < 400)
  350. return response
  351. async def _record_request_audit(
  352. request: web.Request, *, success_state: bool, error: str = ""
  353. ) -> None:
  354. context = request.get("audit_context")
  355. if not context or request.get("audit_recorded"):
  356. return
  357. request["audit_recorded"] = True
  358. admin = request.get("admin") or {}
  359. technician = request.get("technician") or {}
  360. actor_id = admin.get("username") or technician.get("user_id") or "anonymous"
  361. actor_name = (
  362. admin.get("username")
  363. or technician.get("display_name")
  364. or str(actor_id)
  365. )
  366. with suppress(Exception):
  367. await record_audit(
  368. source="technician_mini_app" if technician else "web",
  369. actor_id=actor_id,
  370. actor_name=actor_name,
  371. success=success_state,
  372. error=error,
  373. **context,
  374. )
  375. @web.middleware
  376. async def technician_authentication_middleware(request: web.Request, handler):
  377. if not request.path.startswith(TECHNICIAN_API_PREFIX):
  378. return await handler(request)
  379. bot_id = str(request.headers.get("X-Telegram-Bot-Id") or "").strip()
  380. init_data = str(request.headers.get("X-Telegram-Init-Data") or "").strip()
  381. if not bot_id or not init_data:
  382. raise ApiProblem(
  383. "mini_app_auth_required",
  384. "请从 Telegram Bot 重新打开技师套餐页面。",
  385. status=401,
  386. )
  387. try:
  388. profile = get_bot_profile_secrets(
  389. str(getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json")),
  390. bot_id,
  391. )
  392. except BotConfigError as exc:
  393. raise ApiProblem(
  394. "mini_app_auth_failed",
  395. "Telegram Mini App 鉴权失败。",
  396. status=401,
  397. ) from exc
  398. if (
  399. not profile.get("enabled", True)
  400. or "teacher_directory.manage" not in set(profile.get("permissions", []))
  401. ):
  402. raise ApiProblem(
  403. "mini_app_bot_unavailable",
  404. "当前 Bot 未启用技师自助服务。",
  405. status=403,
  406. )
  407. try:
  408. identity = verify_telegram_webapp_init_data(
  409. init_data,
  410. str(profile.get("bot_token") or ""),
  411. max_age_seconds=max(
  412. 300,
  413. min(
  414. 86400,
  415. int(
  416. getattr(
  417. wbb,
  418. "SERVICE_MINI_APP_AUTH_MAX_AGE_SECONDS",
  419. 3600,
  420. )
  421. or 3600
  422. ),
  423. ),
  424. ),
  425. )
  426. except TelegramWebAppAuthError as exc:
  427. raise ApiProblem("mini_app_auth_failed", str(exc), status=401) from exc
  428. request["technician"] = {**identity, "bot_id": bot_id}
  429. request["telegram_bot_token"] = str(profile.get("bot_token") or "")
  430. return await handler(request)
  431. @web.middleware
  432. async def authentication_middleware(request: web.Request, handler):
  433. if not request.path.startswith(API_PREFIX) or request.path in PUBLIC_API_PATHS:
  434. return await handler(request)
  435. raw_token = request.cookies.get(SESSION_COOKIE, "")
  436. session = await get_admin_session(hash_token(raw_token)) if raw_token else None
  437. if not session:
  438. raise ApiProblem("unauthenticated", "登录已失效,请重新登录。", status=401)
  439. user = session["user"]
  440. request["admin"] = {
  441. "username": user["username"],
  442. "must_change_password": bool(user.get("must_change_password")),
  443. "csrf_token": session["csrf_token"],
  444. "session_token_hash": session["token_hash"],
  445. }
  446. allowed_during_password_change = {
  447. f"{API_PREFIX}/auth/me",
  448. f"{API_PREFIX}/auth/password",
  449. f"{API_PREFIX}/auth/logout",
  450. }
  451. if user.get("must_change_password") and request.path not in allowed_during_password_change:
  452. raise ApiProblem(
  453. "password_change_required", "首次登录必须修改初始密码。", status=428
  454. )
  455. if request.method in UNSAFE_METHODS:
  456. supplied = request.headers.get("X-CSRF-Token", "")
  457. if not supplied or supplied != session["csrf_token"]:
  458. raise ApiProblem("csrf_failed", "CSRF 校验失败。", status=403)
  459. return await handler(request)
  460. def _selected_bot_id(request: web.Request) -> str:
  461. bot_id = str(request.headers.get("X-Bot-Id") or "").strip()
  462. if bot_id:
  463. return bot_id
  464. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  465. if supervisor is not None:
  466. running = [
  467. key
  468. for key, value in supervisor.runtimes().items()
  469. if value.get("state") == "running"
  470. ]
  471. if len(running) == 1:
  472. return running[0]
  473. raise ApiProblem(
  474. "bot_selection_required",
  475. "请先选择要管理的机器人。",
  476. status=409,
  477. )
  478. def _request_bot_token(request: web.Request) -> str:
  479. direct = str(request.get("telegram_bot_token") or "").strip()
  480. if direct:
  481. return direct
  482. if bool(getattr(wbb, "SUPERVISOR_MODE", False)):
  483. try:
  484. profile = get_bot_profile_secrets(
  485. getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"),
  486. _selected_bot_id(request),
  487. )
  488. except (ApiProblem, BotConfigError):
  489. profile = {}
  490. token = str(profile.get("bot_token") or "").strip()
  491. if token:
  492. return token
  493. return str(
  494. getattr(telegram_app, "bot_token", None)
  495. or getattr(wbb, "BOT_TOKEN", "")
  496. or ""
  497. ).strip()
  498. def _require_internal_request(request: web.Request) -> None:
  499. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  500. expected = str(
  501. getattr(supervisor, "internal_token", "")
  502. or getattr(wbb, "INTERNAL_TOKEN", "")
  503. or ""
  504. )
  505. supplied = request.headers.get("Authorization", "")
  506. token = supplied.removeprefix("Bearer ").strip()
  507. if not expected or not token or not secrets.compare_digest(expected, token):
  508. raise ApiProblem(
  509. "internal_auth_failed",
  510. "内部服务鉴权失败。",
  511. status=403,
  512. )
  513. if request.remote and request.remote not in {"127.0.0.1", "::1"}:
  514. raise ApiProblem(
  515. "internal_loopback_required",
  516. "内部服务只接受本机请求。",
  517. status=403,
  518. )
  519. @web.middleware
  520. async def bot_role_middleware(request: web.Request, handler):
  521. required = api_permission(request.method, request.path)
  522. if not required:
  523. return await handler(request)
  524. if bool(getattr(wbb, "SUPERVISOR_MODE", False)):
  525. profile = get_bot_profile_secrets(
  526. getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"),
  527. _selected_bot_id(request),
  528. )
  529. permissions = profile.get("permissions", [])
  530. else:
  531. permissions = getattr(wbb, "BOT_PERMISSIONS", {"*"})
  532. if not has_permission(required, permissions):
  533. raise ApiProblem(
  534. "bot_role_permission_denied",
  535. "所选机器人的职责角色不允许执行该操作。",
  536. status=403,
  537. details={"required_permission": required},
  538. )
  539. return await handler(request)
  540. @web.middleware
  541. async def bot_proxy_middleware(request: web.Request, handler):
  542. if not bool(getattr(wbb, "SUPERVISOR_MODE", False)) or not request.path.startswith(
  543. BOT_SCOPED_PREFIXES
  544. ):
  545. return await handler(request)
  546. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  547. if supervisor is None:
  548. raise ApiProblem(
  549. "bot_supervisor_unavailable",
  550. "机器人监管服务尚未就绪。",
  551. status=503,
  552. )
  553. bot_id = _selected_bot_id(request)
  554. endpoint = supervisor.endpoint_for(bot_id)
  555. if endpoint is None:
  556. raise ApiProblem(
  557. "bot_not_running",
  558. "所选机器人当前未连接 Telegram。",
  559. status=409,
  560. )
  561. target = f"{endpoint}{request.rel_url}"
  562. headers = {
  563. key: value
  564. for key, value in request.headers.items()
  565. if key.lower() not in {"host", "content-length", "connection"}
  566. }
  567. body = await request.read()
  568. try:
  569. async with ClientSession(timeout=ClientTimeout(total=90)) as session:
  570. async with session.request(
  571. request.method,
  572. target,
  573. data=body or None,
  574. headers=headers,
  575. allow_redirects=False,
  576. ) as response:
  577. response_body = await response.read()
  578. response_headers = {
  579. key: value
  580. for key, value in response.headers.items()
  581. if key.lower() in {"content-type", "content-disposition"}
  582. }
  583. return web.Response(
  584. body=response_body,
  585. status=response.status,
  586. headers=response_headers,
  587. )
  588. except (ClientError, TimeoutError) as exc:
  589. raise ApiProblem(
  590. "bot_worker_unavailable",
  591. "机器人工作进程暂时不可用。",
  592. status=503,
  593. ) from exc
  594. class AdminApi:
  595. def __init__(self) -> None:
  596. self.username = str(getattr(wbb, "ADMIN_WEB_USERNAME", "admin"))
  597. self.initial_password = str(
  598. getattr(wbb, "ADMIN_WEB_INITIAL_PASSWORD", "qwe0.123456")
  599. )
  600. self.session_hours = int(getattr(wbb, "ADMIN_WEB_SESSION_HOURS", 12))
  601. self.cookie_secure = bool(getattr(wbb, "ADMIN_WEB_COOKIE_SECURE", False))
  602. self.upload_limit = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)) * 1024 * 1024
  603. self.bot_config_path = Path(
  604. str(getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"))
  605. )
  606. async def initialize(self) -> None:
  607. await ensure_default_admin(
  608. username=self.username,
  609. password_hash=hash_password(self.initial_password),
  610. )
  611. await ensure_directory_indexes()
  612. await ensure_builtin_package_templates()
  613. await ensure_assistant_indexes()
  614. def register(self, application: web.Application) -> None:
  615. router = application.router
  616. router.add_get(
  617. f"{TECHNICIAN_API_PREFIX}/bootstrap",
  618. self.technician_mini_app_bootstrap,
  619. )
  620. router.add_put(
  621. f"{TECHNICIAN_API_PREFIX}/profile",
  622. self.technician_mini_app_profile_update,
  623. )
  624. router.add_patch(
  625. f"{TECHNICIAN_API_PREFIX}/availability",
  626. self.technician_mini_app_availability,
  627. )
  628. router.add_post(
  629. "/api/internal/v1/directory/verify-platform",
  630. self.internal_directory_verify_platform,
  631. )
  632. router.add_post(
  633. "/api/internal/v1/directory/verify-local",
  634. self.internal_directory_verify_local,
  635. )
  636. router.add_post(
  637. "/api/internal/v1/directory/notify-local",
  638. self.internal_directory_notify_local,
  639. )
  640. router.add_get(f"{API_PREFIX}/health", self.health)
  641. router.add_post(f"{API_PREFIX}/auth/login", self.login)
  642. router.add_get(f"{API_PREFIX}/auth/me", self.me)
  643. router.add_post(f"{API_PREFIX}/auth/logout", self.logout)
  644. router.add_put(f"{API_PREFIX}/auth/password", self.change_password)
  645. router.add_get(f"{API_PREFIX}/dashboard", self.dashboard)
  646. router.add_get(f"{API_PREFIX}/settings", self.system_settings)
  647. router.add_get(f"{API_PREFIX}/settings/telegram", self.telegram_settings)
  648. router.add_put(f"{API_PREFIX}/settings/telegram", self.telegram_settings_update)
  649. router.add_get(f"{API_PREFIX}/bots", self.bots)
  650. router.add_post(f"{API_PREFIX}/bots", self.bot_create)
  651. router.add_put(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_update)
  652. router.add_delete(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_delete)
  653. router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/test", self.bot_test)
  654. router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/restart", self.bot_restart)
  655. router.add_get(f"{API_PREFIX}/roles", self.roles)
  656. router.add_post(f"{API_PREFIX}/roles", self.role_create)
  657. router.add_put(f"{API_PREFIX}/roles/{{role_id}}", self.role_update)
  658. router.add_delete(f"{API_PREFIX}/roles/{{role_id}}", self.role_delete)
  659. router.add_get(f"{API_PREFIX}/audit-logs", self.audit_logs)
  660. router.add_post(f"{API_PREFIX}/media", self.upload_media)
  661. router.add_get(
  662. f"{API_PREFIX}/business-assistant/status",
  663. self.business_assistant_status,
  664. )
  665. router.add_get(
  666. f"{API_PREFIX}/business-assistant/connections",
  667. self.business_assistant_connections,
  668. )
  669. router.add_get(
  670. f"{API_PREFIX}/business-assistant/settings",
  671. self.business_assistant_settings,
  672. )
  673. router.add_put(
  674. f"{API_PREFIX}/business-assistant/settings",
  675. self.business_assistant_settings_update,
  676. )
  677. router.add_get(
  678. f"{API_PREFIX}/business-assistant/knowledge",
  679. self.business_assistant_knowledge,
  680. )
  681. router.add_post(
  682. f"{API_PREFIX}/business-assistant/knowledge",
  683. self.business_assistant_knowledge_create,
  684. )
  685. router.add_put(
  686. f"{API_PREFIX}/business-assistant/knowledge/{{entry_id}}",
  687. self.business_assistant_knowledge_update,
  688. )
  689. router.add_delete(
  690. f"{API_PREFIX}/business-assistant/knowledge/{{entry_id}}",
  691. self.business_assistant_knowledge_delete,
  692. )
  693. router.add_post(
  694. f"{API_PREFIX}/business-assistant/knowledge/test",
  695. self.business_assistant_knowledge_test,
  696. )
  697. router.add_get(
  698. f"{API_PREFIX}/business-assistant/conversations",
  699. self.business_assistant_conversations,
  700. )
  701. router.add_get(
  702. f"{API_PREFIX}/business-assistant/conversations/{{conversation_id}}",
  703. self.business_assistant_conversation_detail,
  704. )
  705. router.add_post(
  706. f"{API_PREFIX}/business-assistant/conversations/{{conversation_id}}/actions",
  707. self.business_assistant_conversation_action,
  708. )
  709. router.add_get(
  710. f"{API_PREFIX}/business-assistant/usage",
  711. self.business_assistant_usage,
  712. )
  713. router.add_post(
  714. f"{API_PREFIX}/business-assistant/model/test",
  715. self.business_assistant_model_test,
  716. )
  717. router.add_get(f"{API_PREFIX}/chats", self.chats)
  718. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}", self.chat)
  719. router.add_patch(f"{API_PREFIX}/chats/{{chat_id}}/profile", self.chat_profile)
  720. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/permissions", self.chat_permissions)
  721. router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/announcements", self.announcement)
  722. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/admins", self.chat_admins)
  723. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/search", self.member_search)
  724. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/recent", self.recent_members)
  725. router.add_get(
  726. f"{API_PREFIX}/chats/{{chat_id}}/members/identity-changes",
  727. self.member_identity_changes,
  728. )
  729. router.add_post(
  730. f"{API_PREFIX}/chats/{{chat_id}}/members/{{user_id}}/actions",
  731. self.member_action,
  732. )
  733. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invites)
  734. router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_create)
  735. router.add_delete(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_revoke)
  736. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules)
  737. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules_update)
  738. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation)
  739. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation_update)
  740. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings)
  741. router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings_update)
  742. router.add_get(f"{API_PREFIX}/points/accounts", self.point_accounts)
  743. router.add_get(f"{API_PREFIX}/points/leaderboard", self.point_leaderboard)
  744. router.add_get(f"{API_PREFIX}/points/transactions", self.point_transactions)
  745. router.add_post(f"{API_PREFIX}/points/adjustments", self.point_adjustment)
  746. router.add_get(f"{API_PREFIX}/points/export", self.points_export)
  747. router.add_get(f"{API_PREFIX}/directory/settings", self.directory_settings)
  748. router.add_put(
  749. f"{API_PREFIX}/directory/settings",
  750. self.directory_settings_update,
  751. )
  752. router.add_get(
  753. f"{API_PREFIX}/directory/applications",
  754. self.directory_applications,
  755. )
  756. router.add_post(
  757. f"{API_PREFIX}/directory/applications/{{user_id}}/actions",
  758. self.directory_application_action,
  759. )
  760. router.add_get(f"{API_PREFIX}/directory/profiles", self.directory_profiles)
  761. router.add_post(
  762. f"{API_PREFIX}/directory/profiles/{{user_id}}/actions",
  763. self.directory_profile_action,
  764. )
  765. router.add_get(f"{API_PREFIX}/directory/locations", self.directory_locations)
  766. router.add_put(
  767. f"{API_PREFIX}/directory/locations/{{user_id}}",
  768. self.directory_location_update,
  769. )
  770. router.add_delete(
  771. f"{API_PREFIX}/directory/locations/{{user_id}}",
  772. self.directory_location_delete,
  773. )
  774. router.add_get(f"{API_PREFIX}/directory/events", self.directory_events)
  775. router.add_get(
  776. f"{API_PREFIX}/directory/service-settings",
  777. self.service_settings,
  778. )
  779. router.add_put(
  780. f"{API_PREFIX}/directory/service-settings",
  781. self.service_settings_update,
  782. )
  783. router.add_get(
  784. f"{API_PREFIX}/directory/package-templates",
  785. self.service_package_templates,
  786. )
  787. router.add_post(
  788. f"{API_PREFIX}/directory/package-templates",
  789. self.service_package_template_create,
  790. )
  791. router.add_put(
  792. f"{API_PREFIX}/directory/package-templates/{{template_id}}",
  793. self.service_package_template_update,
  794. )
  795. router.add_post(
  796. f"{API_PREFIX}/directory/package-templates/{{template_id}}/actions",
  797. self.service_package_template_action,
  798. )
  799. router.add_get(
  800. f"{API_PREFIX}/directory/profiles/{{user_id}}/service-profile",
  801. self.technician_service_profile,
  802. )
  803. router.add_put(
  804. f"{API_PREFIX}/directory/profiles/{{user_id}}/service-profile",
  805. self.technician_service_profile_update,
  806. )
  807. router.add_get(
  808. f"{API_PREFIX}/directory/service-orders",
  809. self.service_orders,
  810. )
  811. router.add_get(
  812. f"{API_PREFIX}/directory/service-orders/{{order_id}}",
  813. self.service_order_detail,
  814. )
  815. router.add_post(
  816. f"{API_PREFIX}/directory/service-orders/{{order_id}}/actions",
  817. self.service_order_action,
  818. )
  819. router.add_post(
  820. f"{API_PREFIX}/directory/service-orders/{{order_id}}/address",
  821. self.service_order_address,
  822. )
  823. router.add_post(
  824. f"{API_PREFIX}/directory/customer-blocks/{{customer_id}}",
  825. self.service_customer_block,
  826. )
  827. router.add_get(
  828. f"{API_PREFIX}/directory/customer-blocks",
  829. self.service_customer_blocks,
  830. )
  831. router.add_get(f"{API_PREFIX}/directory/reports", self.service_reports)
  832. router.add_post(
  833. f"{API_PREFIX}/directory/reports/{{report_id}}/actions",
  834. self.service_report_action,
  835. )
  836. router.add_get(
  837. f"{API_PREFIX}/directory/review-template",
  838. self.service_review_template,
  839. )
  840. router.add_put(
  841. f"{API_PREFIX}/directory/review-template",
  842. self.service_review_template_publish,
  843. )
  844. router.add_get(f"{API_PREFIX}/directory/reviews", self.service_reviews)
  845. router.add_post(
  846. f"{API_PREFIX}/directory/reviews/{{review_id}}/actions",
  847. self.service_review_action,
  848. )
  849. router.add_get(
  850. f"{API_PREFIX}/directory/review-qrs",
  851. self.service_review_qrs,
  852. )
  853. router.add_get(
  854. f"{API_PREFIX}/directory/leaderboard",
  855. self.service_leaderboard,
  856. )
  857. router.add_get(
  858. f"{API_PREFIX}/directory/fulfillment-metrics",
  859. self.service_fulfillment_metrics,
  860. )
  861. router.add_get(f"{API_PREFIX}/giveaways", self.giveaways)
  862. router.add_post(f"{API_PREFIX}/giveaways", self.giveaway_create)
  863. router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}", self.giveaway_detail)
  864. router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/finish", self.giveaway_finish)
  865. router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/cancel", self.giveaway_cancel)
  866. router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/reroll", self.giveaway_reroll)
  867. router.add_get(
  868. f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants",
  869. self.giveaway_participants,
  870. )
  871. router.add_delete(
  872. f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants/{{user_id}}",
  873. self.giveaway_participant_remove,
  874. )
  875. router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_bans)
  876. router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_ban_add)
  877. router.add_delete(
  878. f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans/{{user_id}}",
  879. self.giveaway_ban_remove,
  880. )
  881. router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}/export", self.giveaway_export)
  882. async def health(self, _: web.Request) -> web.Response:
  883. return success({"status": "ok"})
  884. async def internal_directory_verify_platform(
  885. self,
  886. request: web.Request,
  887. ) -> web.Response:
  888. _require_internal_request(request)
  889. body = await json_body(request)
  890. user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
  891. require_admin = bool(body.get("require_admin"))
  892. candidates = await list_membership_candidates(user_id)
  893. if not candidates:
  894. return success(
  895. {"allowed": False, "reason": "membership_evidence_required"}
  896. )
  897. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  898. if supervisor is None:
  899. local = [
  900. int(item["chat_id"])
  901. for item in candidates
  902. if str(item.get("bot_id")) == str(getattr(wbb, "BOT_PROFILE_ID", "primary"))
  903. ]
  904. return success(
  905. await verify_local_membership(
  906. user_id=user_id,
  907. chat_ids=local,
  908. require_admin=require_admin,
  909. )
  910. )
  911. candidates_by_bot: dict[str, list[int]] = {}
  912. for item in candidates:
  913. candidates_by_bot.setdefault(str(item["bot_id"]), []).append(
  914. int(item["chat_id"])
  915. )
  916. for bot_id, chat_ids in candidates_by_bot.items():
  917. endpoint = supervisor.endpoint_for(bot_id)
  918. if endpoint is None:
  919. continue
  920. try:
  921. async with ClientSession(timeout=ClientTimeout(total=15)) as session:
  922. async with session.post(
  923. f"{endpoint}/api/internal/v1/directory/verify-local",
  924. headers={
  925. "Authorization": f"Bearer {supervisor.internal_token}"
  926. },
  927. json={
  928. "user_id": str(user_id),
  929. "chat_ids": [str(value) for value in chat_ids],
  930. "require_admin": require_admin,
  931. },
  932. ) as response:
  933. if response.status != 200:
  934. continue
  935. result = (await response.json()).get("data") or {}
  936. if result.get("allowed"):
  937. return success(result)
  938. except (ClientError, TimeoutError, ValueError):
  939. continue
  940. return success({"allowed": False})
  941. async def internal_directory_verify_local(
  942. self,
  943. request: web.Request,
  944. ) -> web.Response:
  945. _require_internal_request(request)
  946. body = await json_body(request)
  947. user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
  948. raw_chat_ids = body.get("chat_ids")
  949. if not isinstance(raw_chat_ids, list):
  950. raise ApiProblem("invalid_parameter", "chat_ids 必须是数组。")
  951. chat_ids = [parse_int(value, "chat_id") for value in raw_chat_ids[:200]]
  952. return success(
  953. await verify_local_membership(
  954. user_id=user_id,
  955. chat_ids=chat_ids,
  956. require_admin=bool(body.get("require_admin")),
  957. )
  958. )
  959. async def internal_directory_notify_local(
  960. self,
  961. request: web.Request,
  962. ) -> web.Response:
  963. _require_internal_request(request)
  964. body = await json_body(request)
  965. user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
  966. text = str(body.get("text") or "").strip()
  967. if not text:
  968. raise ApiProblem("invalid_parameter", "通知内容不能为空。")
  969. try:
  970. await telegram_app.send_message(user_id, text)
  971. except Exception as exc:
  972. raise ApiProblem(
  973. "telegram_notification_failed",
  974. "无法向该用户发送 Telegram 通知。",
  975. status=409,
  976. ) from exc
  977. return success({"sent": True})
  978. async def _notify_directory_user(
  979. self,
  980. *,
  981. profile: dict[str, Any],
  982. text: str,
  983. ) -> None:
  984. bot_id = str(
  985. profile.get("application_source_bot_id")
  986. or getattr(wbb, "BOT_PROFILE_ID", "primary")
  987. )
  988. supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
  989. if supervisor is None:
  990. with suppress(Exception):
  991. await telegram_app.send_message(int(profile["user_id"]), text)
  992. return
  993. endpoint = supervisor.endpoint_for(bot_id)
  994. if endpoint is None:
  995. return
  996. with suppress(Exception):
  997. async with ClientSession(timeout=ClientTimeout(total=10)) as session:
  998. await session.post(
  999. f"{endpoint}/api/internal/v1/directory/notify-local",
  1000. headers={
  1001. "Authorization": f"Bearer {supervisor.internal_token}"
  1002. },
  1003. json={"user_id": str(profile["user_id"]), "text": text},
  1004. )
  1005. async def directory_settings(self, _: web.Request) -> web.Response:
  1006. return success(await get_directory_settings())
  1007. async def directory_settings_update(self, request: web.Request) -> web.Response:
  1008. body = await json_body(request)
  1009. require_confirmation(body)
  1010. values = {
  1011. key: body[key]
  1012. for key in (
  1013. "nearby_radius_km",
  1014. "nearby_radius_options_km",
  1015. )
  1016. if key in body
  1017. }
  1018. saved = await set_directory_settings(values)
  1019. set_audit(
  1020. request,
  1021. "directory.settings.update",
  1022. summary="更新技师目录设置",
  1023. )
  1024. return success(saved)
  1025. async def directory_applications(self, request: web.Request) -> web.Response:
  1026. page, page_size = page_params(request)
  1027. items, total = await list_teacher_applications(
  1028. status=str(request.query.get("status") or ""),
  1029. query=str(request.query.get("query") or ""),
  1030. page=page,
  1031. page_size=page_size,
  1032. )
  1033. return success(
  1034. {
  1035. "items": items,
  1036. "total": total,
  1037. "page": page,
  1038. "page_size": page_size,
  1039. }
  1040. )
  1041. async def directory_application_action(
  1042. self,
  1043. request: web.Request,
  1044. ) -> web.Response:
  1045. body = await json_body(request)
  1046. require_confirmation(body)
  1047. user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
  1048. action = str(body.get("action") or "")
  1049. reason = str(body.get("reason") or "").strip()
  1050. profile = await decide_teacher_application(
  1051. user_id=user_id,
  1052. action=action,
  1053. actor_id=request["admin"]["username"],
  1054. actor_name=request["admin"]["username"],
  1055. source="web",
  1056. reason=reason,
  1057. )
  1058. labels = {"approve": "批准", "reject": "拒绝", "revoke": "撤销"}
  1059. set_audit(
  1060. request,
  1061. f"directory.application.{action}",
  1062. target_id=user_id,
  1063. summary=f"{labels.get(action, action)}技师申请",
  1064. metadata={"reason": reason},
  1065. )
  1066. notification = {
  1067. "approve": "你的技师申请已通过,可以在菜单中自助上榜。",
  1068. "reject": f"你的技师申请未通过。原因:{reason}",
  1069. "revoke": f"你的技师资格已被撤销。原因:{reason}",
  1070. }.get(action)
  1071. if notification:
  1072. await self._notify_directory_user(profile=profile, text=notification)
  1073. return success(profile)
  1074. async def directory_profiles(self, request: web.Request) -> web.Response:
  1075. page, page_size = page_params(request)
  1076. listed_raw = request.query.get("listed")
  1077. listed = (
  1078. listed_raw.lower() in {"1", "true"}
  1079. if listed_raw is not None
  1080. else None
  1081. )
  1082. items, total = await list_directory_profiles(
  1083. query=str(request.query.get("query") or ""),
  1084. application_status=str(
  1085. request.query.get("application_status") or APPLICATION_APPROVED
  1086. ),
  1087. listed=listed,
  1088. page=page,
  1089. page_size=page_size,
  1090. )
  1091. return success(
  1092. {
  1093. "items": items,
  1094. "total": total,
  1095. "page": page,
  1096. "page_size": page_size,
  1097. }
  1098. )
  1099. async def directory_profile_action(self, request: web.Request) -> web.Response:
  1100. body = await json_body(request)
  1101. require_confirmation(body)
  1102. user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
  1103. action = str(body.get("action") or "")
  1104. reason = str(body.get("reason") or "").strip()
  1105. if not reason:
  1106. raise ApiProblem("reason_required", "强制状态操作必须填写原因。")
  1107. profile, _ = await set_teacher_state(
  1108. user_id=user_id,
  1109. action=action,
  1110. actor_id=request["admin"]["username"],
  1111. actor_name=request["admin"]["username"],
  1112. source="web",
  1113. reason=reason,
  1114. )
  1115. set_audit(
  1116. request,
  1117. f"directory.profile.{action}",
  1118. target_id=user_id,
  1119. summary=f"强制修改技师状态:{action}",
  1120. metadata={"reason": reason},
  1121. )
  1122. await self._notify_directory_user(
  1123. profile=profile,
  1124. text=f"管理员已将你的技师状态修改为“{action}”。原因:{reason}",
  1125. )
  1126. return success(profile)
  1127. async def directory_locations(self, request: web.Request) -> web.Response:
  1128. page, page_size = page_params(request)
  1129. items, total = await list_directory_locations(
  1130. query=str(request.query.get("query") or ""),
  1131. page=page,
  1132. page_size=page_size,
  1133. )
  1134. safe_items = []
  1135. for item in items:
  1136. profile = item.get("profile") or {}
  1137. safe_items.append(
  1138. {
  1139. "user_id": item["user_id"],
  1140. "longitude": item.get("longitude"),
  1141. "latitude": item.get("latitude"),
  1142. "source": item.get("source"),
  1143. "updated_at": item.get("updated_at"),
  1144. "display_name": profile.get("display_name"),
  1145. "username": profile.get("username"),
  1146. "application_status": profile.get("application_status", "none"),
  1147. }
  1148. )
  1149. return success(
  1150. {
  1151. "items": safe_items,
  1152. "total": total,
  1153. "page": page,
  1154. "page_size": page_size,
  1155. }
  1156. )
  1157. async def directory_location_update(self, request: web.Request) -> web.Response:
  1158. body = await json_body(request)
  1159. require_confirmation(body)
  1160. user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
  1161. reason = str(body.get("reason") or "").strip()
  1162. if not reason:
  1163. raise ApiProblem("reason_required", "修正位置必须填写原因。")
  1164. profile = await get_directory_profile(user_id)
  1165. if not profile:
  1166. raise ApiProblem("profile_not_found", "未找到该成员资料。", status=404)
  1167. saved = await save_directory_location(
  1168. user_id=user_id,
  1169. longitude=body.get("longitude"),
  1170. latitude=body.get("latitude"),
  1171. source="web",
  1172. actor_id=request["admin"]["username"],
  1173. actor_name=request["admin"]["username"],
  1174. )
  1175. set_audit(
  1176. request,
  1177. "directory.location.update",
  1178. target_id=user_id,
  1179. summary="修正成员位置",
  1180. metadata={"reason": reason},
  1181. )
  1182. return success(
  1183. {
  1184. "user_id": user_id,
  1185. "longitude": saved["longitude"],
  1186. "latitude": saved["latitude"],
  1187. "source": saved.get("source"),
  1188. "updated_at": saved.get("updated_at"),
  1189. }
  1190. )
  1191. async def directory_location_delete(self, request: web.Request) -> web.Response:
  1192. body = await json_body(request)
  1193. require_confirmation(body)
  1194. user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
  1195. reason = str(body.get("reason") or "").strip()
  1196. if not reason:
  1197. raise ApiProblem("reason_required", "清除位置必须填写原因。")
  1198. deleted = await clear_directory_location(
  1199. user_id=user_id,
  1200. actor_id=request["admin"]["username"],
  1201. actor_name=request["admin"]["username"],
  1202. source="web",
  1203. reason=reason,
  1204. )
  1205. set_audit(
  1206. request,
  1207. "directory.location.delete",
  1208. target_id=user_id,
  1209. summary="清除成员位置",
  1210. metadata={"reason": reason},
  1211. )
  1212. return success({"deleted": deleted})
  1213. async def directory_events(self, request: web.Request) -> web.Response:
  1214. page, page_size = page_params(request)
  1215. items, total = await list_directory_events(
  1216. query=str(request.query.get("query") or ""),
  1217. event_type=str(request.query.get("event_type") or ""),
  1218. page=page,
  1219. page_size=page_size,
  1220. )
  1221. return success(
  1222. {
  1223. "items": items,
  1224. "total": total,
  1225. "page": page,
  1226. "page_size": page_size,
  1227. }
  1228. )
  1229. async def service_settings(self, _: web.Request) -> web.Response:
  1230. return success(await get_service_settings())
  1231. async def service_settings_update(self, request: web.Request) -> web.Response:
  1232. body = await json_body(request)
  1233. require_confirmation(body)
  1234. saved = await set_service_settings(body)
  1235. set_audit(request, "service.settings.update", summary="更新技师咨询与评价设置")
  1236. return success(saved)
  1237. async def service_package_templates(self, request: web.Request) -> web.Response:
  1238. items = await list_package_templates(
  1239. include_archived=request.query.get("include_archived", "false").lower()
  1240. == "true"
  1241. )
  1242. return success({"items": items, "total": len(items)})
  1243. async def service_package_template_create(self, request: web.Request) -> web.Response:
  1244. body = await json_body(request)
  1245. require_confirmation(body)
  1246. created = await create_package_template(
  1247. body,
  1248. actor_id=request["admin"]["username"],
  1249. )
  1250. set_audit(
  1251. request,
  1252. "service.package_template.create",
  1253. target_id=created["template_id"],
  1254. summary="创建套餐模板",
  1255. )
  1256. return success(created, status=201)
  1257. async def service_package_template_update(self, request: web.Request) -> web.Response:
  1258. body = await json_body(request)
  1259. require_confirmation(body)
  1260. template_id = request.match_info["template_id"]
  1261. updated = await update_package_template(
  1262. template_id,
  1263. body,
  1264. actor_id=request["admin"]["username"],
  1265. )
  1266. set_audit(
  1267. request,
  1268. "service.package_template.update",
  1269. target_id=template_id,
  1270. summary="更新套餐模板并创建新版本",
  1271. )
  1272. return success(updated)
  1273. async def service_package_template_action(self, request: web.Request) -> web.Response:
  1274. body = await json_body(request)
  1275. require_confirmation(body)
  1276. template_id = request.match_info["template_id"]
  1277. action = str(body.get("action") or "")
  1278. if action == "duplicate":
  1279. updated = await duplicate_package_template(
  1280. template_id,
  1281. actor_id=request["admin"]["username"],
  1282. )
  1283. else:
  1284. updated = await set_package_template_status(
  1285. template_id,
  1286. action,
  1287. actor_id=request["admin"]["username"],
  1288. )
  1289. set_audit(
  1290. request,
  1291. f"service.package_template.{action}",
  1292. target_id=template_id,
  1293. summary=f"套餐模板操作:{action}",
  1294. )
  1295. return success(updated)
  1296. async def technician_service_profile(self, request: web.Request) -> web.Response:
  1297. user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
  1298. return success(await get_technician_service_profile(user_id))
  1299. async def technician_service_profile_update(self, request: web.Request) -> web.Response:
  1300. body = await json_body(request)
  1301. require_confirmation(body)
  1302. user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
  1303. values = body.get("service_profile") or {
  1304. key: value for key, value in body.items() if key != "confirm"
  1305. }
  1306. updated = await publish_technician_profile(
  1307. user_id,
  1308. values,
  1309. actor_id=request["admin"]["username"],
  1310. )
  1311. set_audit(
  1312. request,
  1313. "service.technician_profile.update",
  1314. target_id=user_id,
  1315. summary="更新技师服务资料",
  1316. )
  1317. with suppress(Exception):
  1318. await ensure_technician_review_topic(
  1319. _request_bot_token(request),
  1320. user_id,
  1321. )
  1322. return success(updated)
  1323. async def technician_mini_app_bootstrap(
  1324. self,
  1325. request: web.Request,
  1326. ) -> web.Response:
  1327. identity = request["technician"]
  1328. context = await get_technician_self_service_context(identity["user_id"])
  1329. templates = await list_package_templates()
  1330. public_templates = [
  1331. {
  1332. "template_id": item["template_id"],
  1333. "version": item["version"],
  1334. "name": item["name"],
  1335. "builtin": bool(item.get("builtin")),
  1336. "package": item["package"],
  1337. }
  1338. for item in templates
  1339. if item.get("status") == "enabled"
  1340. ]
  1341. response = success(
  1342. {
  1343. "technician": context,
  1344. "templates": public_templates,
  1345. "limits": {"max_packages": 5},
  1346. }
  1347. )
  1348. response.headers["Cache-Control"] = "no-store"
  1349. return response
  1350. async def technician_mini_app_profile_update(
  1351. self,
  1352. request: web.Request,
  1353. ) -> web.Response:
  1354. identity = request["technician"]
  1355. body = await json_body(request)
  1356. values = body.get("service_profile") or body
  1357. updated = await publish_technician_profile(
  1358. identity["user_id"],
  1359. values,
  1360. actor_id=identity["user_id"],
  1361. )
  1362. set_audit(
  1363. request,
  1364. "service.technician_self_profile.update",
  1365. target_id=identity["user_id"],
  1366. summary="技师通过 Mini App 更新服务资料和套餐",
  1367. )
  1368. with suppress(Exception):
  1369. await ensure_technician_review_topic(
  1370. _request_bot_token(request),
  1371. int(identity["user_id"]),
  1372. )
  1373. return success(
  1374. await get_technician_self_service_context(updated["user_id"])
  1375. )
  1376. async def technician_mini_app_availability(
  1377. self,
  1378. request: web.Request,
  1379. ) -> web.Response:
  1380. identity = request["technician"]
  1381. body = await json_body(request)
  1382. if not isinstance(body.get("accepting_requests"), bool):
  1383. raise ApiProblem(
  1384. "invalid_accepting_requests",
  1385. "接单状态必须是布尔值。",
  1386. )
  1387. updated = await set_technician_accepting_requests(
  1388. identity["user_id"],
  1389. body["accepting_requests"],
  1390. )
  1391. set_audit(
  1392. request,
  1393. "service.technician_self_availability.update",
  1394. target_id=identity["user_id"],
  1395. summary=(
  1396. "开启接单"
  1397. if body["accepting_requests"]
  1398. else "暂停接单"
  1399. ),
  1400. )
  1401. return success(
  1402. await get_technician_self_service_context(updated["user_id"])
  1403. )
  1404. async def service_orders(self, request: web.Request) -> web.Response:
  1405. page, page_size = page_params(request)
  1406. items, total = await list_service_orders(
  1407. query=str(request.query.get("query") or ""),
  1408. status=str(request.query.get("status") or ""),
  1409. page=page,
  1410. page_size=page_size,
  1411. )
  1412. return success(
  1413. {"items": items, "total": total, "page": page, "page_size": page_size}
  1414. )
  1415. async def service_order_detail(self, request: web.Request) -> web.Response:
  1416. return success(await get_admin_service_order(request.match_info["order_id"]))
  1417. async def service_order_action(self, request: web.Request) -> web.Response:
  1418. body = await json_body(request)
  1419. require_confirmation(body)
  1420. action = str(body.get("action") or "")
  1421. reason = str(body.get("reason") or "").strip()
  1422. if not reason:
  1423. raise ApiProblem("reason_required", "争议处理或作废必须填写原因。")
  1424. order_id = request.match_info["order_id"]
  1425. updated = await resolve_service_dispute(
  1426. order_id,
  1427. actor_id=request["admin"]["username"],
  1428. action=action,
  1429. reason=reason,
  1430. )
  1431. set_audit(
  1432. request,
  1433. f"service.order.{action}",
  1434. target_id=order_id,
  1435. summary=f"服务单操作:{action}",
  1436. metadata={"reason": reason},
  1437. )
  1438. return success(updated)
  1439. async def service_order_address(self, request: web.Request) -> web.Response:
  1440. body = await json_body(request)
  1441. require_confirmation(body)
  1442. reason = str(body.get("reason") or "").strip()
  1443. if not reason:
  1444. raise ApiProblem("reason_required", "查看精确地址必须填写原因。")
  1445. order_id = request.match_info["order_id"]
  1446. address = await admin_reveal_order_address(
  1447. order_id,
  1448. actor_id=request["admin"]["username"],
  1449. reason=reason,
  1450. )
  1451. set_audit(
  1452. request,
  1453. "service.order.address_reveal",
  1454. target_id=order_id,
  1455. summary="管理员查看服务单精确地址",
  1456. metadata={"reason": reason},
  1457. )
  1458. return success(address)
  1459. async def service_customer_block(self, request: web.Request) -> web.Response:
  1460. body = await json_body(request)
  1461. require_confirmation(body)
  1462. customer_id = parse_int(
  1463. request.match_info["customer_id"], "customer_id", minimum=1
  1464. )
  1465. active = bool(body.get("active", True))
  1466. reason = str(body.get("reason") or "").strip()
  1467. if active and not reason:
  1468. raise ApiProblem("reason_required", "封禁顾客必须填写原因。")
  1469. updated = await set_customer_block(
  1470. customer_id=customer_id,
  1471. scope="global",
  1472. technician_id=None,
  1473. active=active,
  1474. actor_id=request["admin"]["username"],
  1475. reason=reason,
  1476. )
  1477. set_audit(
  1478. request,
  1479. "service.customer.block" if active else "service.customer.unblock",
  1480. target_id=customer_id,
  1481. summary="全局封禁顾客" if active else "解除顾客封禁",
  1482. metadata={"reason": reason},
  1483. )
  1484. return success(updated)
  1485. async def service_customer_blocks(self, request: web.Request) -> web.Response:
  1486. page, page_size = page_params(request)
  1487. active_raw = request.query.get("active")
  1488. active = (
  1489. active_raw.lower() in {"1", "true"}
  1490. if active_raw is not None
  1491. else None
  1492. )
  1493. items, total = await list_customer_blocks(
  1494. scope=str(request.query.get("scope") or "global"),
  1495. active=active,
  1496. page=page,
  1497. page_size=page_size,
  1498. )
  1499. return success(
  1500. {"items": items, "total": total, "page": page, "page_size": page_size}
  1501. )
  1502. async def service_reports(self, request: web.Request) -> web.Response:
  1503. page, page_size = page_params(request)
  1504. items, total = await list_service_reports(
  1505. status=str(request.query.get("status") or ""),
  1506. page=page,
  1507. page_size=page_size,
  1508. )
  1509. return success(
  1510. {"items": items, "total": total, "page": page, "page_size": page_size}
  1511. )
  1512. async def service_report_action(self, request: web.Request) -> web.Response:
  1513. body = await json_body(request)
  1514. require_confirmation(body)
  1515. action = str(body.get("action") or "")
  1516. reason = str(body.get("reason") or "").strip()
  1517. if not reason:
  1518. raise ApiProblem("reason_required", "处理举报必须填写原因。")
  1519. report = await moderate_service_report(
  1520. request.match_info["report_id"],
  1521. action=action,
  1522. actor_id=request["admin"]["username"],
  1523. reason=reason,
  1524. )
  1525. set_audit(
  1526. request,
  1527. f"service.report.{action}",
  1528. target_id=report["report_id"],
  1529. summary=f"处理服务举报:{action}",
  1530. metadata={"reason": reason},
  1531. )
  1532. return success(report)
  1533. async def service_review_template(self, _: web.Request) -> web.Response:
  1534. return success(await get_active_review_template())
  1535. async def service_review_template_publish(self, request: web.Request) -> web.Response:
  1536. body = await json_body(request)
  1537. require_confirmation(body)
  1538. published = await publish_review_template(
  1539. body,
  1540. actor_id=request["admin"]["username"],
  1541. )
  1542. set_audit(
  1543. request,
  1544. "service.review_template.publish",
  1545. target_id=published["template_id"],
  1546. summary="发布新版评价模板",
  1547. )
  1548. return success(published, status=201)
  1549. async def service_reviews(self, request: web.Request) -> web.Response:
  1550. page, page_size = page_params(request)
  1551. items, total = await list_reviews(
  1552. status=str(request.query.get("status") or ""),
  1553. source=str(request.query.get("source") or ""),
  1554. query=str(request.query.get("query") or ""),
  1555. page=page,
  1556. page_size=page_size,
  1557. )
  1558. return success(
  1559. {"items": items, "total": total, "page": page, "page_size": page_size}
  1560. )
  1561. async def service_review_action(self, request: web.Request) -> web.Response:
  1562. body = await json_body(request)
  1563. require_confirmation(body)
  1564. action = str(body.get("action") or "")
  1565. reason = str(body.get("reason") or "").strip()
  1566. updated = await moderate_review(
  1567. request.match_info["review_id"],
  1568. action,
  1569. actor_id=request["admin"]["username"],
  1570. reason=reason,
  1571. )
  1572. publication = None
  1573. if action == "approve":
  1574. publication = await publish_approved_review(
  1575. _request_bot_token(request),
  1576. updated,
  1577. )
  1578. updated = {**updated, "topic_publication": publication}
  1579. set_audit(
  1580. request,
  1581. f"service.review.{action}",
  1582. target_id=request.match_info["review_id"],
  1583. summary=f"评价审核:{action}",
  1584. metadata={"reason": reason, "topic_publication": publication},
  1585. )
  1586. with suppress(Exception):
  1587. await telegram_app.send_message(
  1588. int(updated["customer_id"]),
  1589. (
  1590. "你的评价已审核通过。"
  1591. if action == "approve"
  1592. else f"你的评价未通过审核。原因:{reason}"
  1593. ),
  1594. )
  1595. return success(updated)
  1596. async def service_review_qrs(self, request: web.Request) -> web.Response:
  1597. page, page_size = page_params(request)
  1598. items, total = await list_review_qr_records(
  1599. status=str(request.query.get("status") or ""),
  1600. query=str(request.query.get("query") or ""),
  1601. page=page,
  1602. page_size=page_size,
  1603. )
  1604. return success(
  1605. {"items": items, "total": total, "page": page, "page_size": page_size}
  1606. )
  1607. async def service_leaderboard(self, request: web.Request) -> web.Response:
  1608. page, page_size = page_params(request)
  1609. items, total = await list_leaderboard(
  1610. category=str(request.query.get("category") or ""),
  1611. page=page,
  1612. page_size=page_size,
  1613. )
  1614. return success(
  1615. {
  1616. "items": items,
  1617. "total": total,
  1618. "page": page,
  1619. "page_size": page_size,
  1620. "rule": "v/(v+5) × 技师平均分 + 5/(v+5) × 分类平均分;至少 3 条批准评价。",
  1621. "title": "审核评价排行",
  1622. }
  1623. )
  1624. async def service_fulfillment_metrics(self, _: web.Request) -> web.Response:
  1625. return success(await fulfillment_metrics())
  1626. def _set_session_cookie(self, response: web.Response, token: str) -> None:
  1627. response.set_cookie(
  1628. SESSION_COOKIE,
  1629. token,
  1630. httponly=True,
  1631. secure=self.cookie_secure,
  1632. samesite="Strict",
  1633. max_age=self.session_hours * 3600,
  1634. path="/",
  1635. )
  1636. async def _new_session(self, request: web.Request, username: str) -> tuple[str, str]:
  1637. token = generate_session_token()
  1638. csrf = generate_csrf_token()
  1639. await create_admin_session(
  1640. username=username,
  1641. token_hash=hash_token(token),
  1642. csrf_token=csrf,
  1643. expires_at=datetime.now(UTC) + timedelta(hours=self.session_hours),
  1644. remote_address=request.remote or "",
  1645. user_agent=request.headers.get("User-Agent", ""),
  1646. )
  1647. return token, csrf
  1648. async def login(self, request: web.Request) -> web.Response:
  1649. body = await json_body(request)
  1650. username = str(body.get("username") or "").strip()
  1651. password = str(body.get("password") or "")
  1652. set_audit(request, "auth.login", summary=f"Login as {username}")
  1653. user = await get_admin_user(username)
  1654. if user and int(user.get("failed_login_count", 0)) >= 5:
  1655. last_failed = user.get("last_failed_login_at")
  1656. if last_failed and datetime.now(UTC) - as_utc(last_failed) < timedelta(minutes=15):
  1657. raise ApiProblem(
  1658. "login_rate_limited", "登录失败次数过多,请 15 分钟后再试。", status=429
  1659. )
  1660. if not user or not verify_password(password, user.get("password_hash", "")):
  1661. await record_login_failure(username)
  1662. raise ApiProblem("invalid_credentials", "用户名或密码错误。", status=401)
  1663. await record_login_success(username)
  1664. token, csrf = await self._new_session(request, username)
  1665. response = success(
  1666. {
  1667. "username": username,
  1668. "must_change_password": bool(user.get("must_change_password")),
  1669. "csrf_token": csrf,
  1670. }
  1671. )
  1672. self._set_session_cookie(response, token)
  1673. request["admin"] = {"username": username}
  1674. return response
  1675. async def me(self, request: web.Request) -> web.Response:
  1676. return success(
  1677. {
  1678. "username": request["admin"]["username"],
  1679. "must_change_password": request["admin"]["must_change_password"],
  1680. "csrf_token": request["admin"]["csrf_token"],
  1681. }
  1682. )
  1683. async def logout(self, request: web.Request) -> web.Response:
  1684. set_audit(request, "auth.logout", summary="Logout")
  1685. await revoke_admin_session(request["admin"]["session_token_hash"])
  1686. response = success({"logged_out": True})
  1687. response.del_cookie(SESSION_COOKIE, path="/")
  1688. return response
  1689. async def change_password(self, request: web.Request) -> web.Response:
  1690. body = await json_body(request)
  1691. current = str(body.get("current_password") or "")
  1692. new_password = str(body.get("new_password") or "")
  1693. username = request["admin"]["username"]
  1694. user = await get_admin_user(username)
  1695. if not user or not verify_password(current, user["password_hash"]):
  1696. raise ApiProblem("invalid_current_password", "当前密码错误。", status=409)
  1697. if verify_password(new_password, user["password_hash"]):
  1698. raise ApiProblem("password_unchanged", "新密码不能与当前密码相同。", status=409)
  1699. errors = validate_new_password(new_password, username)
  1700. if errors:
  1701. raise ApiProblem("weak_password", "新密码不符合要求。", details=errors)
  1702. set_audit(request, "auth.password.change", summary="Change administrator password")
  1703. await update_admin_password(username, hash_password(new_password))
  1704. token, csrf = await self._new_session(request, username)
  1705. response = success(
  1706. {"username": username, "must_change_password": False, "csrf_token": csrf}
  1707. )
  1708. self._set_session_cookie(response, token)
  1709. return response
  1710. async def dashboard(self, _: web.Request) -> web.Response:
  1711. counts = await dashboard_counts()
  1712. giveaways, giveaway_total = await list_giveaways_page(
  1713. page=1,
  1714. page_size=5,
  1715. all_bots=bool(getattr(wbb, "SUPERVISOR_MODE", False)),
  1716. )
  1717. accounts, account_total = await list_point_accounts(page=1, page_size=5)
  1718. transactions, transaction_total = await list_point_transactions(page=1, page_size=8)
  1719. return success(
  1720. {
  1721. "counts": {
  1722. **counts,
  1723. "giveaways": giveaway_total,
  1724. "point_accounts": account_total,
  1725. "point_transactions": transaction_total,
  1726. },
  1727. "recent_giveaways": giveaways,
  1728. "top_accounts": accounts,
  1729. "recent_transactions": transactions,
  1730. }
  1731. )
  1732. async def system_settings(self, _: web.Request) -> web.Response:
  1733. bot_connected = bool(getattr(wbb, "TELEGRAM_CONNECTED", False))
  1734. return success(
  1735. {
  1736. "web_address": f"http://127.0.0.1:{getattr(wbb, 'ADMIN_WEB_PORT', 8088)}/admin",
  1737. "bind_host": str(getattr(wbb, "ADMIN_WEB_HOST", "0.0.0.0")),
  1738. "upload_limit_mb": int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)),
  1739. "supervisor_mode": bool(getattr(wbb, "SUPERVISOR_MODE", False)),
  1740. "bot": {
  1741. "connected": bot_connected,
  1742. "id": str(wbb.BOT_ID) if bot_connected else "",
  1743. "username": wbb.BOT_USERNAME if bot_connected else "",
  1744. "name": wbb.BOT_NAME if bot_connected else "",
  1745. },
  1746. }
  1747. )
  1748. def _supervisor(self):
  1749. return getattr(wbb, "BOT_SUPERVISOR", None)
  1750. def _telegram_status(self) -> dict[str, Any]:
  1751. supervisor = self._supervisor()
  1752. runtimes = supervisor.runtimes() if supervisor is not None else {}
  1753. return telegram_config_status(self.bot_config_path, runtimes=runtimes)
  1754. async def telegram_settings(self, _: web.Request) -> web.Response:
  1755. return success(self._telegram_status())
  1756. async def telegram_settings_update(self, request: web.Request) -> web.Response:
  1757. body = await json_body(request)
  1758. require_confirmation(body)
  1759. set_audit(request, "telegram.settings.update", summary="Update Telegram API credentials")
  1760. body.pop("confirm", None)
  1761. update_telegram_config(self.bot_config_path, body)
  1762. supervisor = self._supervisor()
  1763. if supervisor is not None:
  1764. status = telegram_config_status(self.bot_config_path)
  1765. for profile in status["bots"]:
  1766. if profile["enabled"] and profile["ready_to_connect"]:
  1767. await supervisor.reconcile(profile["bot_id"])
  1768. return success(self._telegram_status())
  1769. async def bots(self, _: web.Request) -> web.Response:
  1770. return success({"items": self._telegram_status()["bots"]})
  1771. async def bot_create(self, request: web.Request) -> web.Response:
  1772. body = await json_body(request)
  1773. require_confirmation(body)
  1774. set_audit(request, "bot.create", summary=str(body.get("label") or ""))
  1775. body.pop("confirm", None)
  1776. profile = create_bot_profile(self.bot_config_path, body)
  1777. supervisor = self._supervisor()
  1778. if supervisor is not None and profile["enabled"] and profile["ready_to_connect"]:
  1779. await supervisor.start(profile["bot_id"])
  1780. current = next(
  1781. item
  1782. for item in self._telegram_status()["bots"]
  1783. if item["bot_id"] == profile["bot_id"]
  1784. )
  1785. return success(current, status=201)
  1786. async def bot_update(self, request: web.Request) -> web.Response:
  1787. body = await json_body(request)
  1788. require_confirmation(body)
  1789. bot_id = request.match_info["bot_id"]
  1790. set_audit(request, "bot.update", target_id=bot_id, summary="Update Bot profile")
  1791. body.pop("confirm", None)
  1792. update_bot_profile(self.bot_config_path, bot_id, body)
  1793. supervisor = self._supervisor()
  1794. if supervisor is not None:
  1795. await supervisor.reconcile(bot_id)
  1796. current = next(
  1797. item for item in self._telegram_status()["bots"] if item["bot_id"] == bot_id
  1798. )
  1799. return success(current)
  1800. async def bot_delete(self, request: web.Request) -> web.Response:
  1801. body = await json_body(request)
  1802. require_confirmation(body)
  1803. bot_id = request.match_info["bot_id"]
  1804. profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
  1805. set_audit(request, "bot.delete", target_id=bot_id, summary=str(profile.get("label") or ""))
  1806. supervisor = self._supervisor()
  1807. if supervisor is not None:
  1808. await supervisor.stop(bot_id)
  1809. delete_bot_profile(self.bot_config_path, bot_id)
  1810. return success({"deleted": True})
  1811. async def bot_test(self, request: web.Request) -> web.Response:
  1812. body = await json_body(request)
  1813. require_confirmation(body)
  1814. bot_id = request.match_info["bot_id"]
  1815. profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
  1816. set_audit(request, "bot.test", target_id=bot_id, summary="Validate Bot Token")
  1817. identity = await test_bot_token(str(profile.get("bot_token") or ""))
  1818. store_bot_identity(self.bot_config_path, bot_id, identity)
  1819. return success(identity)
  1820. async def bot_restart(self, request: web.Request) -> web.Response:
  1821. body = await json_body(request)
  1822. require_confirmation(body)
  1823. bot_id = request.match_info["bot_id"]
  1824. set_audit(request, "bot.restart", target_id=bot_id, summary="Restart Bot worker")
  1825. supervisor = self._supervisor()
  1826. if supervisor is None:
  1827. raise ApiProblem(
  1828. "bot_supervisor_unavailable",
  1829. "机器人监管服务尚未就绪。",
  1830. status=503,
  1831. )
  1832. return success(await supervisor.restart(bot_id))
  1833. async def roles(self, _: web.Request) -> web.Response:
  1834. status = self._telegram_status()
  1835. return success(
  1836. {
  1837. "items": status["roles"],
  1838. "permission_catalog": status["permission_catalog"],
  1839. }
  1840. )
  1841. async def role_create(self, request: web.Request) -> web.Response:
  1842. body = await json_body(request)
  1843. require_confirmation(body)
  1844. set_audit(
  1845. request,
  1846. "bot.role.create",
  1847. summary=str(body.get("name") or ""),
  1848. )
  1849. body.pop("confirm", None)
  1850. return success(create_bot_role(self.bot_config_path, body), status=201)
  1851. async def role_update(self, request: web.Request) -> web.Response:
  1852. body = await json_body(request)
  1853. require_confirmation(body)
  1854. role_id = request.match_info["role_id"]
  1855. set_audit(
  1856. request,
  1857. "bot.role.update",
  1858. target_id=role_id,
  1859. summary=str(body.get("name") or ""),
  1860. )
  1861. body.pop("confirm", None)
  1862. role = update_bot_role(self.bot_config_path, role_id, body)
  1863. supervisor = self._supervisor()
  1864. if supervisor is not None:
  1865. for profile in self._telegram_status()["bots"]:
  1866. if role_id in profile["role_ids"]:
  1867. await supervisor.reconcile(profile["bot_id"])
  1868. return success(role)
  1869. async def role_delete(self, request: web.Request) -> web.Response:
  1870. body = await json_body(request)
  1871. require_confirmation(body)
  1872. role_id = request.match_info["role_id"]
  1873. set_audit(
  1874. request,
  1875. "bot.role.delete",
  1876. target_id=role_id,
  1877. summary="Delete Bot role",
  1878. )
  1879. delete_bot_role(self.bot_config_path, role_id)
  1880. return success({"deleted": True})
  1881. async def audit_logs(self, request: web.Request) -> web.Response:
  1882. page, page_size = page_params(request)
  1883. chat_raw = request.query.get("chat_id")
  1884. items, total = await list_audit_logs(
  1885. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  1886. action=request.query.get("action") or None,
  1887. page=page,
  1888. page_size=page_size,
  1889. )
  1890. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  1891. async def upload_media(self, request: web.Request) -> web.Response:
  1892. if not MESSAGE_DUMP_CHAT:
  1893. raise ApiProblem(
  1894. "message_dump_chat_missing",
  1895. "请先配置媒体中转群。",
  1896. status=409,
  1897. )
  1898. reader = await request.multipart()
  1899. part = await reader.next()
  1900. if not part or part.name != "file":
  1901. raise ApiProblem("file_required", "请选择要上传的文件。")
  1902. filename = Path(part.filename or "upload.bin").name
  1903. content_type = part.headers.get("Content-Type", "application/octet-stream")
  1904. size = 0
  1905. temporary_path: Path | None = None
  1906. try:
  1907. with tempfile.NamedTemporaryFile(delete=False, suffix=Path(filename).suffix) as handle:
  1908. temporary_path = Path(handle.name)
  1909. while True:
  1910. chunk = await part.read_chunk(1024 * 1024)
  1911. if not chunk:
  1912. break
  1913. size += len(chunk)
  1914. if size > self.upload_limit:
  1915. raise ApiProblem("file_too_large", "上传文件超过 20 MB 限制。", status=413)
  1916. handle.write(chunk)
  1917. if content_type == "image/gif":
  1918. media_type = "animation"
  1919. sent = await telegram_app.send_animation(MESSAGE_DUMP_CHAT, str(temporary_path))
  1920. file_id = sent.animation.file_id
  1921. elif content_type.startswith("image/"):
  1922. media_type = "photo"
  1923. sent = await telegram_app.send_photo(MESSAGE_DUMP_CHAT, str(temporary_path))
  1924. file_id = sent.photo.file_id
  1925. elif content_type.startswith("video/"):
  1926. media_type = "video"
  1927. sent = await telegram_app.send_video(MESSAGE_DUMP_CHAT, str(temporary_path))
  1928. file_id = sent.video.file_id
  1929. else:
  1930. media_type = "document"
  1931. sent = await telegram_app.send_document(
  1932. MESSAGE_DUMP_CHAT, str(temporary_path), file_name=filename
  1933. )
  1934. file_id = sent.document.file_id
  1935. finally:
  1936. if temporary_path:
  1937. temporary_path.unlink(missing_ok=True)
  1938. set_audit(request, "media.upload", summary=filename, target_id=file_id)
  1939. return success(
  1940. {
  1941. "file_id": file_id,
  1942. "type": media_type,
  1943. "filename": filename,
  1944. "size": size,
  1945. },
  1946. status=201,
  1947. )
  1948. async def chats(self, request: web.Request) -> web.Response:
  1949. page, page_size = page_params(request)
  1950. items, total = await list_accessible_chats(
  1951. query=request.query.get("query", ""), page=page, page_size=page_size
  1952. )
  1953. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  1954. async def chat(self, request: web.Request) -> web.Response:
  1955. return success(await get_chat_overview(chat_id_param(request)))
  1956. async def chat_profile(self, request: web.Request) -> web.Response:
  1957. chat_id = chat_id_param(request)
  1958. body = await json_body(request)
  1959. require_confirmation(body)
  1960. set_audit(request, "chat.profile.update", chat_id=chat_id, summary="Update group profile")
  1961. return success(
  1962. await update_chat_profile(
  1963. chat_id, title=body.get("title"), description=body.get("description")
  1964. )
  1965. )
  1966. async def chat_permissions(self, request: web.Request) -> web.Response:
  1967. chat_id = chat_id_param(request)
  1968. body = await json_body(request)
  1969. require_confirmation(body)
  1970. set_audit(request, "chat.permissions.update", chat_id=chat_id, summary="Update default member permissions")
  1971. return success(await update_chat_permissions(chat_id, body.get("permissions") or {}))
  1972. async def announcement(self, request: web.Request) -> web.Response:
  1973. chat_id = chat_id_param(request)
  1974. body = await json_body(request)
  1975. require_confirmation(body)
  1976. set_audit(request, "chat.announcement.send", chat_id=chat_id, summary=str(body.get("text") or "")[:200])
  1977. return success(
  1978. await send_announcement(
  1979. chat_id,
  1980. text=str(body.get("text") or ""),
  1981. media_type=body.get("media_type"),
  1982. file_id=body.get("file_id"),
  1983. pin=bool(body.get("pin")),
  1984. ),
  1985. status=201,
  1986. )
  1987. async def chat_admins(self, request: web.Request) -> web.Response:
  1988. return success({"items": await list_chat_admins(chat_id_param(request))})
  1989. async def member_search(self, request: web.Request) -> web.Response:
  1990. return success(
  1991. {
  1992. "items": await search_chat_members(
  1993. chat_id_param(request),
  1994. request.query.get("query", ""),
  1995. limit=parse_int(request.query.get("limit", "20"), "limit"),
  1996. )
  1997. }
  1998. )
  1999. async def recent_members(self, request: web.Request) -> web.Response:
  2000. return success(
  2001. {
  2002. "items": await list_recent_members(
  2003. chat_id_param(request),
  2004. limit=parse_int(request.query.get("limit", "30"), "limit"),
  2005. )
  2006. }
  2007. )
  2008. async def member_identity_changes(self, request: web.Request) -> web.Response:
  2009. page, page_size = page_params(request)
  2010. items, total = await list_member_identity_changes(
  2011. chat_id_param(request), page=page, page_size=page_size
  2012. )
  2013. return success(
  2014. {"items": items, "total": total, "page": page, "page_size": page_size}
  2015. )
  2016. async def member_action(self, request: web.Request) -> web.Response:
  2017. chat_id = chat_id_param(request)
  2018. user_id = parse_int(request.match_info["user_id"], "userId")
  2019. body = await json_body(request)
  2020. require_confirmation(body)
  2021. action = str(body.get("action") or "")
  2022. set_audit(
  2023. request,
  2024. f"chat.member.{action}",
  2025. chat_id=chat_id,
  2026. target_id=user_id,
  2027. summary=str(body.get("reason") or ""),
  2028. )
  2029. duration = body.get("duration_seconds")
  2030. return success(
  2031. await execute_member_action(
  2032. chat_id,
  2033. user_id=user_id,
  2034. action=action,
  2035. reason=str(body.get("reason") or ""),
  2036. duration_seconds=parse_int(duration, "duration_seconds", minimum=60) if duration else None,
  2037. privileges=body.get("privileges") or {},
  2038. )
  2039. )
  2040. async def invites(self, request: web.Request) -> web.Response:
  2041. return success({"items": await get_invite_links(chat_id_param(request))})
  2042. async def invite_create(self, request: web.Request) -> web.Response:
  2043. chat_id = chat_id_param(request)
  2044. body = await json_body(request)
  2045. require_confirmation(body)
  2046. expires_at = parse_datetime(body["expires_at"], "expires_at") if body.get("expires_at") else None
  2047. member_limit = parse_int(body["member_limit"], "member_limit", minimum=1) if body.get("member_limit") else None
  2048. set_audit(request, "chat.invite.create", chat_id=chat_id, summary=str(body.get("name") or ""))
  2049. return success(
  2050. await create_invite_link(
  2051. chat_id,
  2052. name=str(body.get("name") or "Admin panel"),
  2053. expires_at=expires_at,
  2054. member_limit=member_limit,
  2055. ),
  2056. status=201,
  2057. )
  2058. async def invite_revoke(self, request: web.Request) -> web.Response:
  2059. chat_id = chat_id_param(request)
  2060. body = await json_body(request)
  2061. require_confirmation(body)
  2062. invite_link = str(body.get("invite_link") or "")
  2063. if not invite_link:
  2064. raise ApiProblem("invite_link_required", "缺少邀请链接。")
  2065. set_audit(request, "chat.invite.revoke", chat_id=chat_id, summary=invite_link)
  2066. await revoke_invite_link(chat_id, invite_link)
  2067. return success({"revoked": True})
  2068. async def rules(self, request: web.Request) -> web.Response:
  2069. return success({"rules": await get_rules(chat_id_param(request))})
  2070. async def rules_update(self, request: web.Request) -> web.Response:
  2071. chat_id = chat_id_param(request)
  2072. body = await json_body(request)
  2073. require_confirmation(body)
  2074. await ensure_permission(chat_id, "can_change_info")
  2075. rules = str(body.get("rules") or "")
  2076. if len(rules) > 4000:
  2077. raise ApiProblem("rules_too_long", "群规不能超过 4000 个字符。")
  2078. set_audit(request, "chat.rules.update", chat_id=chat_id, summary="Update group rules")
  2079. await set_chat_rules(chat_id, rules)
  2080. return success({"rules": rules})
  2081. async def automation(self, request: web.Request) -> web.Response:
  2082. return success(await get_automation_settings(chat_id_param(request)))
  2083. async def automation_update(self, request: web.Request) -> web.Response:
  2084. chat_id = chat_id_param(request)
  2085. body = await json_body(request)
  2086. require_confirmation(body)
  2087. set_audit(request, "chat.automation.update", chat_id=chat_id, summary="Update automation rules")
  2088. body.pop("confirm", None)
  2089. return success(await apply_automation_settings(chat_id, body))
  2090. async def points_settings(self, request: web.Request) -> web.Response:
  2091. return success(await get_point_rules(chat_id_param(request)))
  2092. async def points_settings_update(self, request: web.Request) -> web.Response:
  2093. chat_id = chat_id_param(request)
  2094. body = await json_body(request)
  2095. require_confirmation(body)
  2096. await ensure_permission(chat_id, "can_change_info")
  2097. set_audit(request, "points.settings.update", chat_id=chat_id, summary="Update point rules")
  2098. body.pop("confirm", None)
  2099. return success(await apply_point_rules(chat_id, body))
  2100. async def point_accounts(self, request: web.Request) -> web.Response:
  2101. page, page_size = page_params(request)
  2102. chat_raw = request.query.get("chat_id")
  2103. items, total = await list_point_accounts(
  2104. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  2105. query=request.query.get("query", ""),
  2106. page=page,
  2107. page_size=page_size,
  2108. )
  2109. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  2110. async def point_leaderboard(self, request: web.Request) -> web.Response:
  2111. if not request.query.get("chat_id"):
  2112. raise ApiProblem("chat_id_required", "排行榜必须选择群组。")
  2113. return await self.point_accounts(request)
  2114. async def point_transactions(self, request: web.Request) -> web.Response:
  2115. page, page_size = page_params(request)
  2116. chat_raw = request.query.get("chat_id")
  2117. user_raw = request.query.get("user_id")
  2118. created_from = parse_datetime(request.query["created_from"], "created_from") if request.query.get("created_from") else None
  2119. created_to = parse_datetime(request.query["created_to"], "created_to") if request.query.get("created_to") else None
  2120. items, total = await list_point_transactions(
  2121. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  2122. user_id=parse_int(user_raw, "user_id") if user_raw else None,
  2123. source=request.query.get("source") or None,
  2124. created_from=created_from,
  2125. created_to=created_to,
  2126. page=page,
  2127. page_size=page_size,
  2128. )
  2129. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  2130. async def point_adjustment(self, request: web.Request) -> web.Response:
  2131. body = await json_body(request)
  2132. require_confirmation(body)
  2133. chat_id = parse_int(body.get("chat_id"), "chat_id")
  2134. user_id = parse_int(body.get("user_id"), "user_id")
  2135. operation = str(body.get("operation") or "")
  2136. amount = parse_int(body.get("amount"), "amount", minimum=0)
  2137. reason = str(body.get("reason") or "").strip()
  2138. if not reason:
  2139. raise ApiProblem("reason_required", "积分调整必须填写原因。")
  2140. if operation not in {"add", "deduct", "set"}:
  2141. raise ApiProblem(
  2142. "invalid_operation",
  2143. "积分操作必须是增加、扣减或设置余额。",
  2144. )
  2145. try:
  2146. user = await telegram_app.get_users(user_id)
  2147. username, first_name = user.username, user.first_name
  2148. display_name = " ".join(
  2149. value for value in (user.first_name, user.last_name) if value
  2150. )
  2151. except Exception:
  2152. username, first_name = body.get("username"), body.get("first_name")
  2153. display_name = body.get("display_name") or first_name
  2154. request_key = str(body.get("request_id") or generate_session_token())
  2155. set_audit(
  2156. request,
  2157. f"points.adjust.{operation}",
  2158. chat_id=chat_id,
  2159. target_id=user_id,
  2160. summary=reason,
  2161. metadata={"amount": amount},
  2162. )
  2163. if operation == "set":
  2164. account, created = await set_points(
  2165. chat_id=chat_id,
  2166. user_id=user_id,
  2167. balance=amount,
  2168. actor_id=request["admin"]["username"],
  2169. reason=reason,
  2170. idempotency_key=f"web-points:{request_key}",
  2171. username=username,
  2172. first_name=first_name,
  2173. display_name=display_name,
  2174. )
  2175. else:
  2176. account, created = await adjust_points(
  2177. chat_id=chat_id,
  2178. user_id=user_id,
  2179. delta=amount if operation == "add" else -amount,
  2180. source=SOURCE_ADMIN,
  2181. idempotency_key=f"web-points:{request_key}",
  2182. actor_id=request["admin"]["username"],
  2183. reason=reason,
  2184. username=username,
  2185. first_name=first_name,
  2186. display_name=display_name,
  2187. )
  2188. return success({"account": account, "created": created}, status=201 if created else 200)
  2189. async def points_export(self, request: web.Request) -> web.Response:
  2190. chat_raw = request.query.get("chat_id")
  2191. if not chat_raw:
  2192. raise ApiProblem("chat_id_required", "导出积分数据必须选择群组。")
  2193. chat_id = parse_int(chat_raw, "chat_id")
  2194. kind = request.query.get("kind", "accounts")
  2195. output = io.StringIO()
  2196. if kind == "accounts":
  2197. items: list[dict[str, Any]] = []
  2198. page = 1
  2199. while True:
  2200. batch, total = await list_point_accounts(
  2201. chat_id=chat_id, page=page, page_size=100
  2202. )
  2203. items.extend(batch)
  2204. if not batch or len(items) >= total:
  2205. break
  2206. page += 1
  2207. fieldnames = ["chat_id", "user_id", "display_name", "username", "first_name", "balance", "lifetime_earned", "lifetime_spent", "updated_at"]
  2208. elif kind == "transactions":
  2209. items = []
  2210. page = 1
  2211. while True:
  2212. batch, total = await list_point_transactions(
  2213. chat_id=chat_id, page=page, page_size=100
  2214. )
  2215. items.extend(batch)
  2216. if not batch or len(items) >= total:
  2217. break
  2218. page += 1
  2219. fieldnames = ["transaction_id", "chat_id", "user_id", "display_name", "username", "first_name", "delta", "balance_after", "source", "actor_id", "reason", "reference_id", "created_at"]
  2220. else:
  2221. raise ApiProblem(
  2222. "invalid_export_kind",
  2223. "导出类型必须是积分账户或积分流水。",
  2224. )
  2225. writer = csv.DictWriter(output, fieldnames=fieldnames, extrasaction="ignore")
  2226. writer.writeheader()
  2227. for item in items:
  2228. writer.writerow(jsonable(item))
  2229. return web.Response(
  2230. text=output.getvalue(),
  2231. content_type="text/csv",
  2232. headers={"Content-Disposition": f'attachment; filename="points-{chat_id}-{kind}.csv"'},
  2233. )
  2234. async def business_assistant_status(self, _: web.Request) -> web.Response:
  2235. return success(await runtime_overview())
  2236. async def business_assistant_connections(self, request: web.Request) -> web.Response:
  2237. page, page_size = page_params(request)
  2238. items, total = await list_business_connections(page=page, page_size=page_size)
  2239. return success(
  2240. {"items": items, "total": total, "page": page, "page_size": page_size}
  2241. )
  2242. async def business_assistant_settings(self, request: web.Request) -> web.Response:
  2243. connection_id = str(request.query.get("connection_id") or "").strip()
  2244. if not connection_id:
  2245. raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
  2246. if not await get_business_connection(connection_id):
  2247. raise AssistantDataError("connection_not_found", "未找到 Business 连接。")
  2248. return success(await get_account_settings(connection_id))
  2249. async def business_assistant_settings_update(
  2250. self, request: web.Request
  2251. ) -> web.Response:
  2252. body = await json_body(request)
  2253. require_confirmation(body)
  2254. connection_id = str(body.pop("connection_id", "")).strip()
  2255. body.pop("confirm", None)
  2256. if not connection_id:
  2257. raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
  2258. set_audit(
  2259. request,
  2260. "business_assistant.settings.update",
  2261. target_id=connection_id,
  2262. summary="更新智能接待设置",
  2263. )
  2264. return success(await update_account_settings(connection_id, body))
  2265. async def business_assistant_knowledge(self, request: web.Request) -> web.Response:
  2266. connection_id = str(request.query.get("connection_id") or "").strip()
  2267. if not connection_id:
  2268. raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
  2269. page, page_size = page_params(request)
  2270. items, total = await list_knowledge_entries(
  2271. connection_id,
  2272. query=str(request.query.get("query") or ""),
  2273. page=page,
  2274. page_size=page_size,
  2275. )
  2276. return success(
  2277. {"items": items, "total": total, "page": page, "page_size": page_size}
  2278. )
  2279. async def business_assistant_knowledge_create(
  2280. self, request: web.Request
  2281. ) -> web.Response:
  2282. body = await json_body(request)
  2283. require_confirmation(body)
  2284. connection_id = str(body.pop("connection_id", "")).strip()
  2285. body.pop("confirm", None)
  2286. if not connection_id:
  2287. raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
  2288. entry = await create_knowledge_entry(connection_id, body)
  2289. set_audit(
  2290. request,
  2291. "business_assistant.knowledge.create",
  2292. target_id=entry["entry_id"],
  2293. summary=str(entry["question"]),
  2294. metadata={"connection_id": connection_id},
  2295. )
  2296. return success(entry, status=201)
  2297. async def business_assistant_knowledge_update(
  2298. self, request: web.Request
  2299. ) -> web.Response:
  2300. body = await json_body(request)
  2301. require_confirmation(body)
  2302. body.pop("confirm", None)
  2303. entry_id = request.match_info["entry_id"]
  2304. set_audit(
  2305. request,
  2306. "business_assistant.knowledge.update",
  2307. target_id=entry_id,
  2308. summary="更新知识条目",
  2309. )
  2310. return success(await update_knowledge_entry(entry_id, body))
  2311. async def business_assistant_knowledge_delete(
  2312. self, request: web.Request
  2313. ) -> web.Response:
  2314. body = await json_body(request)
  2315. require_confirmation(body)
  2316. entry_id = request.match_info["entry_id"]
  2317. set_audit(
  2318. request,
  2319. "business_assistant.knowledge.delete",
  2320. target_id=entry_id,
  2321. summary="删除知识条目",
  2322. )
  2323. await delete_knowledge_entry(entry_id)
  2324. return success({"deleted": True})
  2325. async def business_assistant_knowledge_test(
  2326. self, request: web.Request
  2327. ) -> web.Response:
  2328. body = await json_body(request)
  2329. connection_id = str(body.get("connection_id") or "").strip()
  2330. text = str(body.get("text") or "").strip()
  2331. if not connection_id or not text:
  2332. raise ApiProblem("invalid_parameter", "连接和测试问题不能为空。")
  2333. from wbb.modules.business_assistant import get_business_assistant_runtime
  2334. runtime = get_business_assistant_runtime()
  2335. if runtime is None:
  2336. raise ApiProblem(
  2337. "assistant_runtime_unavailable", "智能接待运行时尚未启动。", status=503
  2338. )
  2339. try:
  2340. result = await runtime.preview_answer(connection_id, text)
  2341. except AssistantProviderError as exc:
  2342. raise ApiProblem("assistant_provider_error", str(exc), status=502) from exc
  2343. return success(result)
  2344. async def business_assistant_conversations(
  2345. self, request: web.Request
  2346. ) -> web.Response:
  2347. page, page_size = page_params(request)
  2348. items, total = await list_conversations(
  2349. connection_id=str(request.query.get("connection_id") or ""),
  2350. status=str(request.query.get("status") or ""),
  2351. page=page,
  2352. page_size=page_size,
  2353. )
  2354. return success(
  2355. {"items": items, "total": total, "page": page, "page_size": page_size}
  2356. )
  2357. async def business_assistant_conversation_detail(
  2358. self, request: web.Request
  2359. ) -> web.Response:
  2360. return success(
  2361. await conversation_detail(request.match_info["conversation_id"])
  2362. )
  2363. async def business_assistant_conversation_action(
  2364. self, request: web.Request
  2365. ) -> web.Response:
  2366. body = await json_body(request)
  2367. require_confirmation(body)
  2368. action = str(body.get("action") or "")
  2369. conversation_id = request.match_info["conversation_id"]
  2370. actions = {
  2371. "pause": pause_conversation,
  2372. "resume": resume_conversation,
  2373. "close": close_conversation,
  2374. "clear": clear_conversation,
  2375. }
  2376. handler = actions.get(action)
  2377. if handler is None:
  2378. raise ApiProblem("invalid_action", "不支持的会话操作。")
  2379. set_audit(
  2380. request,
  2381. f"business_assistant.conversation.{action}",
  2382. target_id=conversation_id,
  2383. summary=f"智能接待会话操作:{action}",
  2384. )
  2385. return success(await handler(conversation_id))
  2386. async def business_assistant_usage(self, request: web.Request) -> web.Response:
  2387. return success(
  2388. await usage_metrics(
  2389. connection_id=str(request.query.get("connection_id") or "")
  2390. )
  2391. )
  2392. async def business_assistant_model_test(self, request: web.Request) -> web.Response:
  2393. body = await json_body(request)
  2394. require_confirmation(body)
  2395. from wbb.modules.business_assistant import get_business_assistant_runtime
  2396. runtime = get_business_assistant_runtime()
  2397. if runtime is None:
  2398. raise ApiProblem(
  2399. "assistant_runtime_unavailable", "智能接待运行时尚未启动。", status=503
  2400. )
  2401. try:
  2402. result = await runtime.provider.test_connection()
  2403. except AssistantProviderError as exc:
  2404. raise ApiProblem("assistant_provider_error", str(exc), status=502) from exc
  2405. set_audit(
  2406. request,
  2407. "business_assistant.model.test",
  2408. summary="验证 OpenAI 兼容模型",
  2409. )
  2410. return success(result)
  2411. async def giveaways(self, request: web.Request) -> web.Response:
  2412. page, page_size = page_params(request)
  2413. chat_raw = request.query.get("chat_id")
  2414. items, total = await list_giveaways_page(
  2415. chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
  2416. status=request.query.get("status") or None,
  2417. query=request.query.get("query", ""),
  2418. page=page,
  2419. page_size=page_size,
  2420. )
  2421. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  2422. async def giveaway_create(self, request: web.Request) -> web.Response:
  2423. body = await json_body(request)
  2424. require_confirmation(body)
  2425. chat_id = parse_int(body.get("chat_id"), "chat_id")
  2426. await ensure_permission(chat_id, "can_change_info")
  2427. prizes = body.get("prizes")
  2428. if not isinstance(prizes, list):
  2429. raise ApiProblem("invalid_prizes", "奖项必须是数组。")
  2430. set_audit(request, "giveaway.create", chat_id=chat_id, summary=str(body.get("title") or ""))
  2431. giveaway = await create_and_publish_giveaway(
  2432. chat_id=chat_id,
  2433. creator_id=0,
  2434. creator_name=request["admin"]["username"],
  2435. title=str(body.get("title") or ""),
  2436. description=str(body.get("description") or ""),
  2437. prizes=prizes,
  2438. ends_at=parse_datetime(body.get("ends_at"), "ends_at"),
  2439. minimum_points=parse_int(body.get("minimum_points", 0), "minimum_points", minimum=0),
  2440. entry_cost=parse_int(body.get("entry_cost", 0), "entry_cost", minimum=0),
  2441. participation_reward=parse_int(body.get("participation_reward", 0), "participation_reward", minimum=0),
  2442. )
  2443. return success(giveaway, status=201)
  2444. async def giveaway_detail(self, request: web.Request) -> web.Response:
  2445. giveaway = await get_giveaway(request.match_info["giveaway_id"])
  2446. if not giveaway:
  2447. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  2448. return success(giveaway)
  2449. async def giveaway_finish(self, request: web.Request) -> web.Response:
  2450. body = await json_body(request)
  2451. require_confirmation(body)
  2452. giveaway_id = request.match_info["giveaway_id"]
  2453. giveaway = await get_giveaway(giveaway_id)
  2454. if not giveaway:
  2455. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  2456. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  2457. set_audit(request, "giveaway.finish", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Manual draw")
  2458. ok, message, current = await finish_and_publish_giveaway(giveaway_id)
  2459. if not ok:
  2460. raise ApiProblem("giveaway_not_running", message, status=409)
  2461. return success(current)
  2462. async def giveaway_cancel(self, request: web.Request) -> web.Response:
  2463. body = await json_body(request)
  2464. require_confirmation(body)
  2465. giveaway_id = request.match_info["giveaway_id"]
  2466. giveaway = await get_giveaway(giveaway_id)
  2467. if not giveaway:
  2468. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  2469. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  2470. set_audit(request, "giveaway.cancel", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Cancel and refund")
  2471. ok, message, current = await cancel_and_refund_giveaway(
  2472. giveaway_id, chat_id=int(giveaway["chat_id"])
  2473. )
  2474. if not ok:
  2475. raise ApiProblem("giveaway_not_running", message, status=409)
  2476. return success(current)
  2477. async def giveaway_reroll(self, request: web.Request) -> web.Response:
  2478. body = await json_body(request)
  2479. require_confirmation(body)
  2480. giveaway_id = request.match_info["giveaway_id"]
  2481. giveaway = await get_giveaway(giveaway_id)
  2482. if not giveaway:
  2483. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  2484. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  2485. set_audit(request, "giveaway.reroll", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary=str(body.get("tier_name") or "All tiers"))
  2486. winners, reroll_id = await reroll_giveaway(
  2487. giveaway_id=giveaway_id,
  2488. moderator_id=0,
  2489. tier_name=body.get("tier_name") or None,
  2490. reroll_id=body.get("request_id") or None,
  2491. )
  2492. return success({"winners": winners, "reroll_id": reroll_id})
  2493. async def giveaway_participants(self, request: web.Request) -> web.Response:
  2494. page, page_size = page_params(request)
  2495. items, total = await list_participants_page(
  2496. giveaway_id=request.match_info["giveaway_id"],
  2497. active_only=request.query.get("active_only", "false").lower() == "true",
  2498. page=page,
  2499. page_size=page_size,
  2500. )
  2501. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  2502. async def giveaway_participant_remove(self, request: web.Request) -> web.Response:
  2503. body = await json_body(request)
  2504. require_confirmation(body)
  2505. giveaway_id = request.match_info["giveaway_id"]
  2506. user_id = parse_int(request.match_info["user_id"], "user_id")
  2507. giveaway = await get_giveaway(giveaway_id)
  2508. if not giveaway:
  2509. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  2510. await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
  2511. set_audit(request, "giveaway.participant.remove", chat_id=int(giveaway["chat_id"]), target_id=user_id, summary=str(body.get("reason") or ""))
  2512. removed = await remove_and_optionally_refund_participant(
  2513. giveaway_id=giveaway_id,
  2514. user_id=user_id,
  2515. moderator_id=0,
  2516. reason=str(body.get("reason") or "Removed in admin panel"),
  2517. refund=bool(body.get("refund", True)),
  2518. )
  2519. if not removed:
  2520. raise ApiProblem("participant_not_found", "参与者不存在或已被移除。", status=404)
  2521. return success({"removed": True, "refunded": bool(body.get("refund", True))})
  2522. async def giveaway_bans(self, request: web.Request) -> web.Response:
  2523. page, page_size = page_params(request)
  2524. items, total = await list_giveaway_bans(
  2525. chat_id=chat_id_param(request), page=page, page_size=page_size
  2526. )
  2527. return success({"items": items, "total": total, "page": page, "page_size": page_size})
  2528. async def giveaway_ban_add(self, request: web.Request) -> web.Response:
  2529. chat_id = chat_id_param(request)
  2530. body = await json_body(request)
  2531. require_confirmation(body)
  2532. user_id = parse_int(body.get("user_id"), "user_id")
  2533. await ensure_permission(chat_id, "can_change_info")
  2534. set_audit(request, "giveaway.ban.add", chat_id=chat_id, target_id=user_id, summary=str(body.get("reason") or ""))
  2535. await add_giveaway_ban(chat_id=chat_id, user_id=user_id, moderator_id=0, reason=str(body.get("reason") or ""))
  2536. return success({"created": True}, status=201)
  2537. async def giveaway_ban_remove(self, request: web.Request) -> web.Response:
  2538. chat_id = chat_id_param(request)
  2539. user_id = parse_int(request.match_info["user_id"], "user_id")
  2540. body = await json_body(request)
  2541. require_confirmation(body)
  2542. await ensure_permission(chat_id, "can_change_info")
  2543. set_audit(request, "giveaway.ban.remove", chat_id=chat_id, target_id=user_id)
  2544. return success({"removed": await remove_giveaway_ban(chat_id, user_id)})
  2545. async def giveaway_export(self, request: web.Request) -> web.Response:
  2546. giveaway_id = request.match_info["giveaway_id"]
  2547. giveaway = await get_giveaway(giveaway_id)
  2548. if not giveaway:
  2549. raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
  2550. participants: list[dict[str, Any]] = []
  2551. page = 1
  2552. while True:
  2553. batch, total = await list_participants_page(
  2554. giveaway_id=giveaway_id, page=page, page_size=100
  2555. )
  2556. participants.extend(batch)
  2557. if not batch or len(participants) >= total:
  2558. break
  2559. page += 1
  2560. output = io.StringIO()
  2561. fields = ["giveaway_id", "chat_id", "user_id", "display_name", "username", "first_name", "active", "entry_cost", "joined_at", "removed_at", "refunded_at"]
  2562. writer = csv.DictWriter(output, fieldnames=fields, extrasaction="ignore")
  2563. writer.writeheader()
  2564. for participant in participants:
  2565. writer.writerow(jsonable(participant))
  2566. return web.Response(
  2567. text=output.getvalue(),
  2568. content_type="text/csv",
  2569. headers={"Content-Disposition": f'attachment; filename="giveaway-{giveaway_id}.csv"'},
  2570. )
  2571. def build_admin_application() -> web.Application:
  2572. max_upload_mb = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20))
  2573. application = web.Application(
  2574. middlewares=[
  2575. api_error_middleware,
  2576. technician_authentication_middleware,
  2577. authentication_middleware,
  2578. bot_role_middleware,
  2579. bot_proxy_middleware,
  2580. ],
  2581. client_max_size=(max_upload_mb + 1) * 1024 * 1024,
  2582. )
  2583. admin_api = AdminApi()
  2584. application["admin_api"] = admin_api
  2585. admin_api.register(application)
  2586. return application