api.py 139 KB

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