| 12345678910111213141516171819202122232425262728293031323334353637383940414243444546474849505152535455565758596061626364656667686970717273747576777879808182838485868788899091929394959697989910010110210310410510610710810911011111211311411511611711811912012112212312412512612712812913013113213313413513613713813914014114214314414514614714814915015115215315415515615715815916016116216316416516616716816917017117217317417517617717817918018118218318418518618718818919019119219319419519619719819920020120220320420520620720820921021121221321421521621721821922022122222322422522622722822923023123223323423523623723823924024124224324424524624724824925025125225325425525625725825926026126226326426526626726826927027127227327427527627727827928028128228328428528628728828929029129229329429529629729829930030130230330430530630730830931031131231331431531631731831932032132232332432532632732832933033133233333433533633733833934034134234334434534634734834935035135235335435535635735835936036136236336436536636736836937037137237337437537637737837938038138238338438538638738838939039139239339439539639739839940040140240340440540640740840941041141241341441541641741841942042142242342442542642742842943043143243343443543643743843944044144244344444544644744844945045145245345445545645745845946046146246346446546646746846947047147247347447547647747847948048148248348448548648748848949049149249349449549649749849950050150250350450550650750850951051151251351451551651751851952052152252352452552652752852953053153253353453553653753853954054154254354454554654754854955055155255355455555655755855956056156256356456556656756856957057157257357457557657757857958058158258358458558658758858959059159259359459559659759859960060160260360460560660760860961061161261361461561661761861962062162262362462562662762862963063163263363463563663763863964064164264364464564664764864965065165265365465565665765865966066166266366466566666766866967067167267367467567667767867968068168268368468568668768868969069169269369469569669769869970070170270370470570670770870971071171271371471571671771871972072172272372472572672772872973073173273373473573673773873974074174274374474574674774874975075175275375475575675775875976076176276376476576676776876977077177277377477577677777877978078178278378478578678778878979079179279379479579679779879980080180280380480580680780880981081181281381481581681781881982082182282382482582682782882983083183283383483583683783883984084184284384484584684784884985085185285385485585685785885986086186286386486586686786886987087187287387487587687787887988088188288388488588688788888989089189289389489589689789889990090190290390490590690790890991091191291391491591691791891992092192292392492592692792892993093193293393493593693793893994094194294394494594694794894995095195295395495595695795895996096196296396496596696796896997097197297397497597697797897998098198298398498598698798898999099199299399499599699799899910001001100210031004100510061007100810091010101110121013101410151016101710181019102010211022102310241025102610271028102910301031103210331034103510361037103810391040104110421043104410451046104710481049105010511052105310541055105610571058105910601061106210631064106510661067106810691070107110721073107410751076107710781079108010811082108310841085108610871088108910901091109210931094109510961097109810991100110111021103110411051106110711081109111011111112111311141115111611171118111911201121112211231124112511261127112811291130113111321133113411351136113711381139114011411142114311441145114611471148114911501151115211531154115511561157115811591160116111621163116411651166116711681169117011711172117311741175117611771178117911801181118211831184118511861187118811891190119111921193119411951196119711981199120012011202120312041205120612071208120912101211121212131214121512161217121812191220122112221223122412251226122712281229123012311232123312341235123612371238123912401241124212431244124512461247124812491250125112521253125412551256125712581259126012611262126312641265126612671268126912701271127212731274127512761277127812791280128112821283128412851286128712881289129012911292129312941295129612971298129913001301130213031304130513061307130813091310131113121313131413151316131713181319132013211322132313241325132613271328132913301331133213331334133513361337133813391340134113421343134413451346134713481349135013511352135313541355135613571358135913601361136213631364136513661367136813691370137113721373137413751376137713781379138013811382138313841385138613871388138913901391139213931394139513961397139813991400140114021403140414051406140714081409141014111412141314141415141614171418141914201421142214231424142514261427142814291430143114321433143414351436143714381439144014411442144314441445144614471448144914501451145214531454145514561457145814591460146114621463146414651466146714681469147014711472147314741475147614771478147914801481148214831484148514861487148814891490149114921493149414951496149714981499150015011502150315041505150615071508150915101511151215131514151515161517151815191520152115221523152415251526152715281529153015311532153315341535153615371538153915401541154215431544154515461547154815491550155115521553155415551556155715581559156015611562156315641565156615671568156915701571157215731574157515761577157815791580158115821583158415851586158715881589159015911592159315941595159615971598159916001601160216031604160516061607160816091610161116121613161416151616161716181619162016211622162316241625162616271628162916301631163216331634163516361637163816391640164116421643164416451646164716481649165016511652165316541655165616571658165916601661166216631664166516661667166816691670167116721673167416751676167716781679168016811682168316841685168616871688168916901691169216931694169516961697169816991700170117021703170417051706170717081709171017111712171317141715171617171718171917201721172217231724172517261727172817291730173117321733173417351736173717381739174017411742174317441745174617471748174917501751175217531754175517561757175817591760176117621763176417651766176717681769177017711772177317741775177617771778177917801781178217831784178517861787178817891790179117921793179417951796179717981799180018011802180318041805180618071808180918101811181218131814181518161817181818191820182118221823182418251826182718281829183018311832183318341835183618371838183918401841184218431844184518461847184818491850185118521853185418551856185718581859186018611862186318641865186618671868186918701871187218731874187518761877187818791880188118821883188418851886188718881889189018911892189318941895189618971898189919001901190219031904190519061907190819091910191119121913191419151916191719181919192019211922192319241925192619271928192919301931193219331934193519361937193819391940194119421943194419451946194719481949195019511952195319541955195619571958195919601961196219631964196519661967196819691970197119721973197419751976197719781979198019811982198319841985198619871988198919901991199219931994199519961997199819992000200120022003200420052006200720082009201020112012201320142015201620172018201920202021202220232024202520262027202820292030203120322033203420352036203720382039204020412042204320442045204620472048204920502051205220532054205520562057205820592060206120622063206420652066206720682069207020712072207320742075207620772078207920802081208220832084208520862087208820892090209120922093209420952096209720982099210021012102210321042105210621072108210921102111211221132114211521162117211821192120212121222123212421252126212721282129213021312132213321342135213621372138213921402141214221432144214521462147214821492150215121522153215421552156215721582159216021612162216321642165216621672168216921702171217221732174217521762177217821792180218121822183218421852186218721882189219021912192219321942195219621972198219922002201220222032204220522062207220822092210221122122213221422152216221722182219222022212222222322242225222622272228222922302231223222332234223522362237223822392240224122422243224422452246224722482249225022512252225322542255225622572258225922602261226222632264226522662267226822692270227122722273227422752276227722782279228022812282228322842285228622872288228922902291229222932294229522962297229822992300230123022303230423052306230723082309231023112312231323142315231623172318231923202321232223232324232523262327232823292330233123322333233423352336233723382339234023412342234323442345234623472348234923502351235223532354235523562357235823592360236123622363236423652366236723682369237023712372237323742375237623772378237923802381238223832384238523862387238823892390239123922393239423952396239723982399240024012402240324042405240624072408240924102411241224132414241524162417241824192420242124222423242424252426242724282429243024312432243324342435243624372438243924402441244224432444244524462447244824492450245124522453245424552456245724582459246024612462246324642465246624672468246924702471247224732474247524762477247824792480248124822483248424852486248724882489249024912492249324942495249624972498249925002501250225032504250525062507250825092510251125122513251425152516251725182519252025212522252325242525252625272528252925302531253225332534253525362537253825392540254125422543254425452546254725482549255025512552255325542555255625572558255925602561256225632564256525662567256825692570257125722573257425752576257725782579258025812582258325842585258625872588258925902591259225932594259525962597259825992600260126022603260426052606260726082609261026112612261326142615261626172618261926202621262226232624262526262627262826292630263126322633263426352636263726382639264026412642264326442645264626472648264926502651265226532654265526562657265826592660266126622663266426652666266726682669267026712672267326742675267626772678267926802681268226832684268526862687268826892690269126922693269426952696269726982699270027012702270327042705270627072708270927102711271227132714271527162717271827192720272127222723272427252726272727282729273027312732273327342735273627372738273927402741274227432744274527462747274827492750275127522753275427552756275727582759276027612762276327642765276627672768276927702771277227732774277527762777277827792780278127822783278427852786278727882789279027912792279327942795279627972798279928002801280228032804280528062807280828092810281128122813281428152816281728182819282028212822282328242825282628272828282928302831283228332834283528362837283828392840284128422843284428452846284728482849285028512852285328542855285628572858285928602861286228632864286528662867286828692870287128722873287428752876287728782879288028812882288328842885288628872888288928902891289228932894289528962897289828992900290129022903290429052906290729082909291029112912291329142915291629172918291929202921292229232924292529262927292829292930293129322933293429352936293729382939294029412942294329442945294629472948294929502951295229532954295529562957295829592960296129622963296429652966296729682969297029712972297329742975297629772978297929802981298229832984298529862987298829892990299129922993299429952996299729982999300030013002300330043005300630073008300930103011301230133014301530163017301830193020302130223023302430253026302730283029303030313032303330343035303630373038303930403041304230433044304530463047304830493050305130523053305430553056305730583059306030613062306330643065306630673068306930703071307230733074307530763077307830793080308130823083308430853086308730883089309030913092309330943095309630973098309931003101310231033104310531063107310831093110311131123113311431153116311731183119312031213122312331243125312631273128312931303131313231333134313531363137313831393140314131423143314431453146314731483149315031513152315331543155315631573158315931603161316231633164316531663167316831693170317131723173317431753176317731783179318031813182318331843185318631873188318931903191319231933194319531963197319831993200320132023203320432053206320732083209321032113212321332143215321632173218321932203221322232233224322532263227322832293230323132323233323432353236323732383239324032413242324332443245324632473248324932503251325232533254325532563257325832593260326132623263326432653266326732683269327032713272327332743275327632773278327932803281328232833284328532863287328832893290329132923293329432953296329732983299330033013302330333043305330633073308330933103311331233133314331533163317331833193320332133223323332433253326332733283329333033313332333333343335333633373338333933403341334233433344334533463347334833493350335133523353335433553356335733583359336033613362336333643365336633673368336933703371337233733374337533763377337833793380338133823383338433853386338733883389339033913392339333943395339633973398339934003401340234033404340534063407340834093410341134123413341434153416341734183419342034213422342334243425342634273428342934303431343234333434 |
- from __future__ import annotations
- import csv
- import io
- import secrets
- import tempfile
- import traceback
- from contextlib import suppress
- from datetime import UTC, datetime, timedelta
- from pathlib import Path
- from typing import Any
- from aiohttp import ClientError, ClientSession, ClientTimeout, web
- import wbb
- from wbb import MESSAGE_DUMP_CHAT
- from wbb import app as telegram_app
- from wbb.admin.bot_config import (
- BotConfigError,
- create_bot_profile,
- create_bot_role,
- delete_bot_profile,
- delete_bot_role,
- get_bot_profile_secrets,
- store_bot_identity,
- telegram_config_status,
- test_bot_token,
- update_bot_profile,
- update_bot_role,
- update_telegram_config,
- )
- from wbb.admin.security import (
- generate_csrf_token,
- generate_session_token,
- hash_password,
- hash_token,
- validate_new_password,
- verify_password,
- )
- from wbb.admin.telegram_webapp import (
- TelegramWebAppAuthError,
- verify_telegram_webapp_init_data,
- )
- from wbb.services.bot_permissions import api_permission, has_any_permission, has_permission
- from wbb.services.business_assistant import (
- AssistantProviderError,
- FeishuWebhookError,
- runtime_overview,
- )
- from wbb.services.channel_management import (
- create_channel_post,
- delete_channel_post,
- edit_channel_post,
- ensure_channel,
- list_channel_admin_candidates,
- list_channel_admins,
- list_channels,
- public_post,
- remove_channel_admin,
- set_channel_admin,
- update_channel_photo,
- update_channel_profile,
- )
- from wbb.services.chat_management import (
- ChatManagementError,
- apply_automation_settings,
- create_invite_link,
- ensure_permission,
- execute_member_action,
- get_automation_settings,
- get_chat_overview,
- get_invite_links,
- list_accessible_chats,
- list_chat_admins,
- list_recent_members,
- remove_unavailable_chat,
- revoke_invite_link,
- search_chat_members,
- send_announcement,
- update_chat_permissions,
- update_chat_profile,
- )
- from wbb.services.directory import DirectoryServiceError, verify_local_membership
- from wbb.services.external_risk_rules import (
- SOURCES as EXTERNAL_RISK_SOURCES,
- )
- from wbb.services.external_risk_rules import (
- get_policy as get_external_risk_policy,
- )
- from wbb.services.external_risk_rules import (
- list_active_sources as list_active_external_risk_sources,
- )
- from wbb.services.external_risk_rules import (
- list_snapshot_entries as list_external_risk_snapshot_entries,
- )
- from wbb.services.external_risk_rules import (
- list_snapshots as list_external_risk_snapshots,
- )
- from wbb.services.external_risk_rules import (
- publish_snapshot as publish_external_risk_snapshot,
- )
- from wbb.services.external_risk_rules import (
- set_policy as set_external_risk_policy,
- )
- from wbb.services.external_risk_rules import (
- stage_snapshot as stage_external_risk_snapshot,
- )
- from wbb.services.giveaway_eligibility import validate_targets
- from wbb.services.giveaway_templates import (
- create_template,
- list_templates,
- update_template,
- )
- from wbb.services.giveaways import (
- GiveawayServiceError,
- cancel_and_refund_giveaway,
- create_and_publish_giveaway,
- finish_and_publish_giveaway,
- remove_and_optionally_refund_participant,
- reroll_giveaway,
- update_and_refresh_giveaway,
- )
- from wbb.services.interaction_settings import (
- get_interaction_settings,
- save_interaction_settings,
- )
- from wbb.services.point_settings import apply_point_rules
- from wbb.services.technician_reviews import (
- ensure_technician_review_topic,
- publish_approved_review,
- )
- from wbb.utils.dbadmin import (
- create_admin_session,
- dashboard_counts,
- ensure_default_admin,
- get_admin_session,
- get_admin_user,
- list_audit_logs,
- list_member_identity_changes,
- record_audit,
- record_login_failure,
- record_login_success,
- revoke_admin_session,
- update_admin_password,
- )
- from wbb.utils.dbassistant import (
- AssistantDataError,
- clear_conversation,
- close_conversation,
- conversation_detail,
- create_knowledge_entry,
- create_knowledge_source,
- delete_knowledge_entry,
- delete_knowledge_source,
- ensure_assistant_indexes,
- get_account_settings,
- get_business_connection,
- get_knowledge_candidate,
- get_knowledge_source,
- list_business_connections,
- list_conversations,
- list_knowledge_candidates,
- list_knowledge_entries,
- list_knowledge_sources,
- pause_conversation,
- public_account_settings,
- publish_knowledge_candidate,
- reject_knowledge_candidate,
- resume_conversation,
- update_account_settings,
- update_knowledge_entry,
- update_knowledge_source,
- usage_metrics,
- )
- from wbb.utils.dbchannel import list_posts
- from wbb.utils.dbdirectory import (
- APPLICATION_APPROVED,
- DirectoryDataError,
- clear_directory_location,
- decide_teacher_application,
- ensure_directory_indexes,
- get_directory_profile,
- get_directory_settings,
- list_directory_events,
- list_directory_locations,
- list_directory_profiles,
- list_membership_candidates,
- list_teacher_applications,
- save_directory_location,
- set_directory_settings,
- set_teacher_state,
- )
- from wbb.utils.dbfunctions import get_rules, set_chat_rules
- from wbb.utils.dbgiveaway import (
- add_giveaway_ban,
- get_giveaway,
- list_giveaway_bans,
- list_giveaways_page,
- list_participants_page,
- remove_giveaway_ban,
- templatesdb,
- ticket_ordersdb,
- )
- from wbb.utils.dbpoints import (
- SOURCE_ADMIN,
- InsufficientPoints,
- PointsError,
- adjust_points,
- get_point_rules,
- list_point_accounts,
- list_point_transactions,
- set_points,
- )
- from wbb.utils.dbservice import (
- ServiceDataError,
- admin_reveal_order_address,
- create_package_template,
- duplicate_package_template,
- ensure_builtin_package_templates,
- fulfillment_metrics,
- get_active_review_template,
- get_admin_service_order,
- get_service_settings,
- get_technician_self_service_context,
- get_technician_service_profile,
- list_customer_blocks,
- list_leaderboard,
- list_package_templates,
- list_review_qr_records,
- list_reviews,
- list_service_orders,
- list_service_reports,
- moderate_review,
- moderate_service_report,
- publish_review_template,
- publish_technician_profile,
- resolve_service_dispute,
- set_customer_block,
- set_package_template_status,
- set_service_settings,
- set_technician_accepting_requests,
- update_package_template,
- )
- API_PREFIX = "/api/admin/v1"
- TECHNICIAN_API_PREFIX = "/api/technician/v1"
- SESSION_COOKIE = "wbb_admin_session"
- UNSAFE_METHODS = {"POST", "PUT", "PATCH", "DELETE"}
- PUBLIC_API_PATHS = {f"{API_PREFIX}/auth/login", f"{API_PREFIX}/health"}
- BOT_SCOPED_PREFIXES = (
- f"{API_PREFIX}/business-assistant",
- f"{API_PREFIX}/chats",
- f"{API_PREFIX}/channels",
- f"{API_PREFIX}/giveaways",
- f"{API_PREFIX}/points",
- f"{API_PREFIX}/media",
- f"{API_PREFIX}/external-risk",
- )
- TELEGRAM_ID_KEYS = {
- "chat_id",
- "user_id",
- "actor_id",
- "creator_id",
- "customer_id",
- "moderator_id",
- "reporter_id",
- "target_id",
- "technician_id",
- "message_id",
- "removed_by",
- "authorization_chat_id",
- "ops_group_id",
- }
- class ApiProblem(RuntimeError):
- def __init__(
- self,
- code: str,
- message: str,
- *,
- status: int = 400,
- details: Any = None,
- ):
- super().__init__(message)
- self.code = code
- self.status = status
- self.details = details
- def as_utc(value: datetime) -> datetime:
- if value.tzinfo is None:
- return value.replace(tzinfo=UTC)
- return value.astimezone(UTC)
- def jsonable(value: Any, *, key: str = "") -> Any:
- if isinstance(value, datetime):
- return as_utc(value).isoformat().replace("+00:00", "Z")
- if isinstance(value, dict):
- return {
- item_key: jsonable(item_value, key=item_key)
- for item_key, item_value in value.items()
- if item_key not in {"_id", "encrypted_payload", "token_hash"}
- }
- if isinstance(value, (list, tuple)):
- return [jsonable(item) for item in value]
- if isinstance(value, int) and (key in TELEGRAM_ID_KEYS or key.endswith("_telegram_id")):
- return str(value)
- return value
- def success(data: Any, *, status: int = 200) -> web.Response:
- return web.json_response({"data": jsonable(data)}, status=status)
- def error_response(problem: ApiProblem) -> web.Response:
- return web.json_response(
- {
- "error": {
- "code": problem.code,
- "message": str(problem),
- "details": jsonable(problem.details),
- }
- },
- status=problem.status,
- )
- def parse_int(value: Any, name: str, *, minimum: int | None = None) -> int:
- try:
- parsed = int(value)
- except (TypeError, ValueError) as exc:
- raise ApiProblem("invalid_parameter", f"{name} 必须是整数。") from exc
- if minimum is not None and parsed < minimum:
- raise ApiProblem("invalid_parameter", f"{name} 不能小于 {minimum}。")
- return parsed
- def parse_datetime(value: Any, name: str) -> datetime:
- if not isinstance(value, str) or not value.strip():
- raise ApiProblem("invalid_parameter", f"{name} 必须是 UTC ISO 8601 时间。")
- try:
- parsed = datetime.fromisoformat(value.replace("Z", "+00:00"))
- except ValueError as exc:
- raise ApiProblem("invalid_parameter", f"{name} 不是有效时间。") from exc
- if parsed.tzinfo is None:
- raise ApiProblem("invalid_parameter", f"{name} 必须包含时区。")
- return parsed.astimezone(UTC)
- async def json_body(request: web.Request) -> dict[str, Any]:
- try:
- body = await request.json()
- except Exception as exc:
- raise ApiProblem("invalid_json", "请求体必须是 JSON。") from exc
- if not isinstance(body, dict):
- raise ApiProblem("invalid_json", "请求体必须是 JSON 对象。")
- return body
- def require_confirmation(body: dict[str, Any]) -> None:
- if body.get("confirm") is not True:
- raise ApiProblem(
- "confirmation_required", "该操作需要二次确认。", status=409
- )
- def page_params(request: web.Request) -> tuple[int, int]:
- return (
- parse_int(request.query.get("page", 1), "page", minimum=1),
- min(parse_int(request.query.get("page_size", 20), "page_size", minimum=1), 100),
- )
- def chat_id_param(request: web.Request) -> int:
- return parse_int(request.match_info["chat_id"], "chatId")
- def set_audit(
- request: web.Request,
- action: str,
- *,
- chat_id: int | None = None,
- target_id: int | str | None = None,
- summary: str = "",
- metadata: dict[str, Any] | None = None,
- ) -> None:
- request["audit_context"] = {
- "action": action,
- "chat_id": chat_id,
- "target_id": target_id,
- "summary": summary,
- "metadata": metadata or {},
- }
- @web.middleware
- async def api_error_middleware(request: web.Request, handler):
- try:
- response = await handler(request)
- except ApiProblem as exc:
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(exc)
- except (ChatManagementError, GiveawayServiceError, DirectoryServiceError) as exc:
- problem = ApiProblem(
- exc.code, str(exc), status=getattr(exc, "status", 400)
- )
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except (DirectoryDataError, ServiceDataError) as exc:
- problem = ApiProblem(exc.code, str(exc), status=409)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except AssistantDataError as exc:
- status = 404 if exc.code.endswith("_not_found") else 409
- problem = ApiProblem(exc.code, str(exc), status=status)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except (PointsError, InsufficientPoints) as exc:
- problem = ApiProblem("points_error", str(exc), status=409)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except BotConfigError as exc:
- status = {
- "bot_not_found": 404,
- "role_not_found": 404,
- "role_in_use": 409,
- "builtin_role_immutable": 409,
- }.get(exc.code, 400)
- problem = ApiProblem(exc.code, str(exc), status=status)
- await _record_request_audit(request, success_state=False, error=str(exc))
- return error_response(problem)
- except web.HTTPException:
- raise
- except Exception as exc:
- await _record_request_audit(request, success_state=False, error=str(exc))
- wbb.log.error(
- f"Admin API {request.method} {request.path} failed: {exc}\n"
- f"{traceback.format_exc()}"
- )
- return error_response(
- ApiProblem("internal_error", "服务器处理请求失败。", status=500)
- )
- await _record_request_audit(request, success_state=response.status < 400)
- return response
- async def _record_request_audit(
- request: web.Request, *, success_state: bool, error: str = ""
- ) -> None:
- context = request.get("audit_context")
- if not context or request.get("audit_recorded"):
- return
- request["audit_recorded"] = True
- admin = request.get("admin") or {}
- technician = request.get("technician") or {}
- actor_id = admin.get("username") or technician.get("user_id") or "anonymous"
- actor_name = (
- admin.get("username")
- or technician.get("display_name")
- or str(actor_id)
- )
- with suppress(Exception):
- await record_audit(
- source="technician_mini_app" if technician else "web",
- actor_id=actor_id,
- actor_name=actor_name,
- success=success_state,
- error=error,
- **context,
- )
- @web.middleware
- async def technician_authentication_middleware(request: web.Request, handler):
- if not request.path.startswith(TECHNICIAN_API_PREFIX):
- return await handler(request)
- bot_id = str(request.headers.get("X-Telegram-Bot-Id") or "").strip()
- init_data = str(request.headers.get("X-Telegram-Init-Data") or "").strip()
- if not bot_id or not init_data:
- raise ApiProblem(
- "mini_app_auth_required",
- "请从 Telegram Bot 重新打开技师套餐页面。",
- status=401,
- )
- try:
- profile = get_bot_profile_secrets(
- str(getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json")),
- bot_id,
- )
- except BotConfigError as exc:
- raise ApiProblem(
- "mini_app_auth_failed",
- "Telegram Mini App 鉴权失败。",
- status=401,
- ) from exc
- if (
- not profile.get("enabled", True)
- or "teacher_directory.manage" not in set(profile.get("permissions", []))
- ):
- raise ApiProblem(
- "mini_app_bot_unavailable",
- "当前 Bot 未启用技师自助服务。",
- status=403,
- )
- try:
- identity = verify_telegram_webapp_init_data(
- init_data,
- str(profile.get("bot_token") or ""),
- max_age_seconds=max(
- 300,
- min(
- 86400,
- int(
- getattr(
- wbb,
- "SERVICE_MINI_APP_AUTH_MAX_AGE_SECONDS",
- 3600,
- )
- or 3600
- ),
- ),
- ),
- )
- except TelegramWebAppAuthError as exc:
- raise ApiProblem("mini_app_auth_failed", str(exc), status=401) from exc
- request["technician"] = {**identity, "bot_id": bot_id}
- request["telegram_bot_token"] = str(profile.get("bot_token") or "")
- return await handler(request)
- @web.middleware
- async def authentication_middleware(request: web.Request, handler):
- if not request.path.startswith(API_PREFIX) or request.path in PUBLIC_API_PATHS:
- return await handler(request)
- raw_token = request.cookies.get(SESSION_COOKIE, "")
- session = await get_admin_session(hash_token(raw_token)) if raw_token else None
- if not session:
- raise ApiProblem("unauthenticated", "登录已失效,请重新登录。", status=401)
- user = session["user"]
- request["admin"] = {
- "username": user["username"],
- "must_change_password": bool(user.get("must_change_password")),
- "csrf_token": session["csrf_token"],
- "session_token_hash": session["token_hash"],
- }
- allowed_during_password_change = {
- f"{API_PREFIX}/auth/me",
- f"{API_PREFIX}/auth/password",
- f"{API_PREFIX}/auth/logout",
- }
- if user.get("must_change_password") and request.path not in allowed_during_password_change:
- raise ApiProblem(
- "password_change_required", "首次登录必须修改初始密码。", status=428
- )
- if request.method in UNSAFE_METHODS:
- supplied = request.headers.get("X-CSRF-Token", "")
- if not supplied or supplied != session["csrf_token"]:
- raise ApiProblem("csrf_failed", "CSRF 校验失败。", status=403)
- return await handler(request)
- def _selected_bot_id(request: web.Request) -> str:
- bot_id = str(request.headers.get("X-Bot-Id") or "").strip()
- if bot_id:
- return bot_id
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is not None:
- running = [
- key
- for key, value in supervisor.runtimes().items()
- if value.get("state") == "running"
- ]
- if len(running) == 1:
- return running[0]
- raise ApiProblem(
- "bot_selection_required",
- "请先选择要管理的机器人。",
- status=409,
- )
- def _request_bot_token(request: web.Request) -> str:
- direct = str(request.get("telegram_bot_token") or "").strip()
- if direct:
- return direct
- if bool(getattr(wbb, "SUPERVISOR_MODE", False)):
- try:
- profile = get_bot_profile_secrets(
- getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"),
- _selected_bot_id(request),
- )
- except (ApiProblem, BotConfigError):
- profile = {}
- token = str(profile.get("bot_token") or "").strip()
- if token:
- return token
- return str(
- getattr(telegram_app, "bot_token", None)
- or getattr(wbb, "BOT_TOKEN", "")
- or ""
- ).strip()
- def _require_internal_request(request: web.Request) -> None:
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- expected = str(
- getattr(supervisor, "internal_token", "")
- or getattr(wbb, "INTERNAL_TOKEN", "")
- or ""
- )
- supplied = request.headers.get("Authorization", "")
- token = supplied.removeprefix("Bearer ").strip()
- if not expected or not token or not secrets.compare_digest(expected, token):
- raise ApiProblem(
- "internal_auth_failed",
- "内部服务鉴权失败。",
- status=403,
- )
- if request.remote and request.remote not in {"127.0.0.1", "::1"}:
- raise ApiProblem(
- "internal_loopback_required",
- "内部服务只接受本机请求。",
- status=403,
- )
- @web.middleware
- async def bot_role_middleware(request: web.Request, handler):
- required = api_permission(request.method, request.path)
- if not required:
- return await handler(request)
- if bool(getattr(wbb, "SUPERVISOR_MODE", False)):
- profile = get_bot_profile_secrets(
- getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"),
- _selected_bot_id(request),
- )
- permissions = profile.get("permissions", [])
- else:
- permissions = getattr(wbb, "BOT_PERMISSIONS", {"*"})
- allowed = (
- has_any_permission(
- {"channel.profile", "channel.admins", "channel.invites", "channel.posts"},
- permissions,
- )
- if required == "channel.any"
- else has_permission(required, permissions)
- )
- if not allowed:
- raise ApiProblem(
- "bot_role_permission_denied",
- "所选机器人的职责角色不允许执行该操作。",
- status=403,
- details={"required_permission": required},
- )
- return await handler(request)
- @web.middleware
- async def bot_proxy_middleware(request: web.Request, handler):
- if not bool(getattr(wbb, "SUPERVISOR_MODE", False)) or not request.path.startswith(
- BOT_SCOPED_PREFIXES
- ):
- return await handler(request)
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is None:
- raise ApiProblem(
- "bot_supervisor_unavailable",
- "机器人监管服务尚未就绪。",
- status=503,
- )
- bot_id = _selected_bot_id(request)
- endpoint = supervisor.endpoint_for(bot_id)
- if endpoint is None:
- raise ApiProblem(
- "bot_not_running",
- "所选机器人当前未连接 Telegram。",
- status=409,
- )
- target = f"{endpoint}{request.rel_url}"
- headers = {
- key: value
- for key, value in request.headers.items()
- if key.lower() not in {"host", "content-length", "connection"}
- }
- body = await request.read()
- try:
- async with ClientSession(timeout=ClientTimeout(total=90)) as session:
- async with session.request(
- request.method,
- target,
- data=body or None,
- headers=headers,
- allow_redirects=False,
- ) as response:
- response_body = await response.read()
- response_headers = {
- key: value
- for key, value in response.headers.items()
- if key.lower() in {"content-type", "content-disposition"}
- }
- return web.Response(
- body=response_body,
- status=response.status,
- headers=response_headers,
- )
- except (ClientError, TimeoutError) as exc:
- raise ApiProblem(
- "bot_worker_unavailable",
- "机器人工作进程暂时不可用。",
- status=503,
- ) from exc
- class AdminApi:
- def __init__(self) -> None:
- self.username = str(getattr(wbb, "ADMIN_WEB_USERNAME", "admin"))
- self.initial_password = str(
- getattr(wbb, "ADMIN_WEB_INITIAL_PASSWORD", "qwe0.123456")
- )
- self.session_hours = int(getattr(wbb, "ADMIN_WEB_SESSION_HOURS", 12))
- self.cookie_secure = bool(getattr(wbb, "ADMIN_WEB_COOKIE_SECURE", False))
- self.upload_limit = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)) * 1024 * 1024
- self.bot_config_path = Path(
- str(getattr(wbb, "BOT_PROFILES_PATH", "runtime/bot_profiles.json"))
- )
- async def initialize(self) -> None:
- await ensure_default_admin(
- username=self.username,
- password_hash=hash_password(self.initial_password),
- )
- await ensure_directory_indexes()
- await ensure_builtin_package_templates()
- await ensure_assistant_indexes()
- def register(self, application: web.Application) -> None:
- router = application.router
- router.add_get(
- f"{TECHNICIAN_API_PREFIX}/bootstrap",
- self.technician_mini_app_bootstrap,
- )
- router.add_put(
- f"{TECHNICIAN_API_PREFIX}/profile",
- self.technician_mini_app_profile_update,
- )
- router.add_patch(
- f"{TECHNICIAN_API_PREFIX}/availability",
- self.technician_mini_app_availability,
- )
- router.add_post(
- "/api/internal/v1/directory/verify-platform",
- self.internal_directory_verify_platform,
- )
- router.add_post(
- "/api/internal/v1/directory/verify-local",
- self.internal_directory_verify_local,
- )
- router.add_post(
- "/api/internal/v1/directory/notify-local",
- self.internal_directory_notify_local,
- )
- router.add_get(f"{API_PREFIX}/health", self.health)
- router.add_post(f"{API_PREFIX}/auth/login", self.login)
- router.add_get(f"{API_PREFIX}/auth/me", self.me)
- router.add_post(f"{API_PREFIX}/auth/logout", self.logout)
- router.add_put(f"{API_PREFIX}/auth/password", self.change_password)
- router.add_get(f"{API_PREFIX}/dashboard", self.dashboard)
- router.add_get(f"{API_PREFIX}/settings", self.system_settings)
- router.add_get(f"{API_PREFIX}/settings/telegram", self.telegram_settings)
- router.add_put(f"{API_PREFIX}/settings/telegram", self.telegram_settings_update)
- router.add_get(f"{API_PREFIX}/bots", self.bots)
- router.add_post(f"{API_PREFIX}/bots", self.bot_create)
- router.add_put(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_update)
- router.add_delete(f"{API_PREFIX}/bots/{{bot_id}}", self.bot_delete)
- router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/test", self.bot_test)
- router.add_post(f"{API_PREFIX}/bots/{{bot_id}}/restart", self.bot_restart)
- router.add_get(f"{API_PREFIX}/roles", self.roles)
- router.add_post(f"{API_PREFIX}/roles", self.role_create)
- router.add_put(f"{API_PREFIX}/roles/{{role_id}}", self.role_update)
- router.add_delete(f"{API_PREFIX}/roles/{{role_id}}", self.role_delete)
- router.add_get(f"{API_PREFIX}/audit-logs", self.audit_logs)
- router.add_post(f"{API_PREFIX}/media", self.upload_media)
- router.add_get(f"{API_PREFIX}/channels", self.channels)
- router.add_post(f"{API_PREFIX}/channels", self.channel_add)
- router.add_get(f"{API_PREFIX}/channels/{{chat_id}}", self.channel)
- router.add_patch(f"{API_PREFIX}/channels/{{chat_id}}/profile", self.channel_profile)
- router.add_post(f"{API_PREFIX}/channels/{{chat_id}}/photo", self.channel_photo)
- router.add_get(f"{API_PREFIX}/channels/{{chat_id}}/admins", self.channel_admins)
- router.add_get(f"{API_PREFIX}/channels/{{chat_id}}/admins/candidates", self.channel_admin_candidates)
- router.add_put(f"{API_PREFIX}/channels/{{chat_id}}/admins/{{user_id}}", self.channel_admin_update)
- router.add_delete(f"{API_PREFIX}/channels/{{chat_id}}/admins/{{user_id}}", self.channel_admin_remove)
- router.add_get(f"{API_PREFIX}/channels/{{chat_id}}/invites", self.channel_invites)
- router.add_post(f"{API_PREFIX}/channels/{{chat_id}}/invites", self.channel_invite_create)
- router.add_delete(f"{API_PREFIX}/channels/{{chat_id}}/invites", self.channel_invite_revoke)
- router.add_get(f"{API_PREFIX}/channels/{{chat_id}}/posts", self.channel_posts)
- router.add_post(f"{API_PREFIX}/channels/{{chat_id}}/posts", self.channel_post_create)
- router.add_patch(f"{API_PREFIX}/channels/{{chat_id}}/posts/{{post_id}}", self.channel_post_edit)
- router.add_delete(f"{API_PREFIX}/channels/{{chat_id}}/posts/{{post_id}}", self.channel_post_delete)
- router.add_get(
- f"{API_PREFIX}/business-assistant/status",
- self.business_assistant_status,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/connections",
- self.business_assistant_connections,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/settings",
- self.business_assistant_settings,
- )
- router.add_put(
- f"{API_PREFIX}/business-assistant/settings",
- self.business_assistant_settings_update,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/feishu-webhook/test",
- self.business_assistant_feishu_webhook_test,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/knowledge",
- self.business_assistant_knowledge,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/knowledge",
- self.business_assistant_knowledge_create,
- )
- router.add_put(
- f"{API_PREFIX}/business-assistant/knowledge/{{entry_id}}",
- self.business_assistant_knowledge_update,
- )
- router.add_delete(
- f"{API_PREFIX}/business-assistant/knowledge/{{entry_id}}",
- self.business_assistant_knowledge_delete,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/knowledge/test",
- self.business_assistant_knowledge_test,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/knowledge-sources",
- self.business_assistant_knowledge_sources,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/knowledge-sources",
- self.business_assistant_knowledge_source_create,
- )
- router.add_put(
- f"{API_PREFIX}/business-assistant/knowledge-sources/{{source_id}}",
- self.business_assistant_knowledge_source_update,
- )
- router.add_delete(
- f"{API_PREFIX}/business-assistant/knowledge-sources/{{source_id}}",
- self.business_assistant_knowledge_source_delete,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/knowledge-sources/{{source_id}}/actions",
- self.business_assistant_knowledge_source_action,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/knowledge-candidates",
- self.business_assistant_knowledge_candidates,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/knowledge-candidates/{{candidate_id}}/actions",
- self.business_assistant_knowledge_candidate_action,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/conversations",
- self.business_assistant_conversations,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/conversations/{{conversation_id}}",
- self.business_assistant_conversation_detail,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/conversations/{{conversation_id}}/actions",
- self.business_assistant_conversation_action,
- )
- router.add_get(
- f"{API_PREFIX}/business-assistant/usage",
- self.business_assistant_usage,
- )
- router.add_post(
- f"{API_PREFIX}/business-assistant/model/test",
- self.business_assistant_model_test,
- )
- router.add_get(f"{API_PREFIX}/chats", self.chats)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}", self.chat)
- router.add_delete(f"{API_PREFIX}/chats/{{chat_id}}", self.chat_delete)
- router.add_patch(f"{API_PREFIX}/chats/{{chat_id}}/profile", self.chat_profile)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/permissions", self.chat_permissions)
- router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/announcements", self.announcement)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/admins", self.chat_admins)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/search", self.member_search)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/members/recent", self.recent_members)
- router.add_get(
- f"{API_PREFIX}/chats/{{chat_id}}/members/identity-changes",
- self.member_identity_changes,
- )
- router.add_post(
- f"{API_PREFIX}/chats/{{chat_id}}/members/{{user_id}}/actions",
- self.member_action,
- )
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invites)
- router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_create)
- router.add_delete(f"{API_PREFIX}/chats/{{chat_id}}/invites", self.invite_revoke)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/rules", self.rules_update)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/automation", self.automation_update)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings)
- router.add_put(f"{API_PREFIX}/chats/{{chat_id}}/points/settings", self.points_settings_update)
- router.add_get(f"{API_PREFIX}/points/accounts", self.point_accounts)
- router.add_get(f"{API_PREFIX}/points/leaderboard", self.point_leaderboard)
- router.add_get(f"{API_PREFIX}/points/transactions", self.point_transactions)
- router.add_post(f"{API_PREFIX}/points/adjustments", self.point_adjustment)
- router.add_get(f"{API_PREFIX}/points/export", self.points_export)
- router.add_get(f"{API_PREFIX}/directory/settings", self.directory_settings)
- router.add_put(
- f"{API_PREFIX}/directory/settings",
- self.directory_settings_update,
- )
- router.add_get(
- f"{API_PREFIX}/directory/applications",
- self.directory_applications,
- )
- router.add_post(
- f"{API_PREFIX}/directory/applications/{{user_id}}/actions",
- self.directory_application_action,
- )
- router.add_get(f"{API_PREFIX}/directory/profiles", self.directory_profiles)
- router.add_post(
- f"{API_PREFIX}/directory/profiles/{{user_id}}/actions",
- self.directory_profile_action,
- )
- router.add_get(f"{API_PREFIX}/directory/locations", self.directory_locations)
- router.add_put(
- f"{API_PREFIX}/directory/locations/{{user_id}}",
- self.directory_location_update,
- )
- router.add_delete(
- f"{API_PREFIX}/directory/locations/{{user_id}}",
- self.directory_location_delete,
- )
- router.add_get(f"{API_PREFIX}/directory/events", self.directory_events)
- router.add_get(
- f"{API_PREFIX}/directory/service-settings",
- self.service_settings,
- )
- router.add_put(
- f"{API_PREFIX}/directory/service-settings",
- self.service_settings_update,
- )
- router.add_get(
- f"{API_PREFIX}/directory/package-templates",
- self.service_package_templates,
- )
- router.add_post(
- f"{API_PREFIX}/directory/package-templates",
- self.service_package_template_create,
- )
- router.add_put(
- f"{API_PREFIX}/directory/package-templates/{{template_id}}",
- self.service_package_template_update,
- )
- router.add_post(
- f"{API_PREFIX}/directory/package-templates/{{template_id}}/actions",
- self.service_package_template_action,
- )
- router.add_get(
- f"{API_PREFIX}/directory/profiles/{{user_id}}/service-profile",
- self.technician_service_profile,
- )
- router.add_put(
- f"{API_PREFIX}/directory/profiles/{{user_id}}/service-profile",
- self.technician_service_profile_update,
- )
- router.add_get(
- f"{API_PREFIX}/directory/service-orders",
- self.service_orders,
- )
- router.add_get(
- f"{API_PREFIX}/directory/service-orders/{{order_id}}",
- self.service_order_detail,
- )
- router.add_post(
- f"{API_PREFIX}/directory/service-orders/{{order_id}}/actions",
- self.service_order_action,
- )
- router.add_post(
- f"{API_PREFIX}/directory/service-orders/{{order_id}}/address",
- self.service_order_address,
- )
- router.add_post(
- f"{API_PREFIX}/directory/customer-blocks/{{customer_id}}",
- self.service_customer_block,
- )
- router.add_get(
- f"{API_PREFIX}/directory/customer-blocks",
- self.service_customer_blocks,
- )
- router.add_get(f"{API_PREFIX}/directory/reports", self.service_reports)
- router.add_post(
- f"{API_PREFIX}/directory/reports/{{report_id}}/actions",
- self.service_report_action,
- )
- router.add_get(
- f"{API_PREFIX}/directory/review-template",
- self.service_review_template,
- )
- router.add_put(
- f"{API_PREFIX}/directory/review-template",
- self.service_review_template_publish,
- )
- router.add_get(f"{API_PREFIX}/directory/reviews", self.service_reviews)
- router.add_post(
- f"{API_PREFIX}/directory/reviews/{{review_id}}/actions",
- self.service_review_action,
- )
- router.add_get(
- f"{API_PREFIX}/directory/review-qrs",
- self.service_review_qrs,
- )
- router.add_get(
- f"{API_PREFIX}/directory/leaderboard",
- self.service_leaderboard,
- )
- router.add_get(
- f"{API_PREFIX}/directory/fulfillment-metrics",
- self.service_fulfillment_metrics,
- )
- router.add_get(f"{API_PREFIX}/giveaways", self.giveaways)
- router.add_post(f"{API_PREFIX}/giveaways", self.giveaway_create)
- router.add_get(f"{API_PREFIX}/giveaways/templates", self.giveaway_templates)
- router.add_post(f"{API_PREFIX}/giveaways/templates", self.giveaway_template_create)
- router.add_patch(
- f"{API_PREFIX}/giveaways/templates/{{template_id}}", self.giveaway_template_update
- )
- router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}", self.giveaway_detail)
- router.add_patch(f"{API_PREFIX}/giveaways/{{giveaway_id}}", self.giveaway_update)
- router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/finish", self.giveaway_finish)
- router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/cancel", self.giveaway_cancel)
- router.add_post(f"{API_PREFIX}/giveaways/{{giveaway_id}}/reroll", self.giveaway_reroll)
- router.add_get(
- f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants",
- self.giveaway_participants,
- )
- router.add_get(
- f"{API_PREFIX}/giveaways/{{giveaway_id}}/orders", self.giveaway_orders
- )
- router.add_delete(
- f"{API_PREFIX}/giveaways/{{giveaway_id}}/participants/{{user_id}}",
- self.giveaway_participant_remove,
- )
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_bans)
- router.add_post(f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans", self.giveaway_ban_add)
- router.add_delete(
- f"{API_PREFIX}/chats/{{chat_id}}/giveaway-bans/{{user_id}}",
- self.giveaway_ban_remove,
- )
- router.add_get(f"{API_PREFIX}/giveaways/{{giveaway_id}}/export", self.giveaway_export)
- router.add_get(f"{API_PREFIX}/external-risk/sources", self.external_risk_sources)
- router.add_post(f"{API_PREFIX}/external-risk/sources/{{source_id}}/stage", self.external_risk_stage)
- router.add_post(f"{API_PREFIX}/external-risk/snapshots/{{snapshot_id}}/publish", self.external_risk_publish)
- router.add_get(f"{API_PREFIX}/external-risk/snapshots/{{snapshot_id}}/entries", self.external_risk_entries)
- router.add_get(f"{API_PREFIX}/external-risk/chats/{{chat_id}}", self.external_risk_chat_policy)
- router.add_patch(f"{API_PREFIX}/external-risk/chats/{{chat_id}}", self.external_risk_chat_policy_update)
- router.add_get(f"{API_PREFIX}/chats/{{chat_id}}/interaction-settings", self.chat_interaction_settings)
- router.add_patch(f"{API_PREFIX}/chats/{{chat_id}}/interaction-settings", self.chat_interaction_settings_update)
- async def health(self, _: web.Request) -> web.Response:
- return success({"status": "ok"})
- async def internal_directory_verify_platform(
- self,
- request: web.Request,
- ) -> web.Response:
- _require_internal_request(request)
- body = await json_body(request)
- user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
- require_admin = bool(body.get("require_admin"))
- candidates = await list_membership_candidates(user_id)
- if not candidates:
- return success(
- {"allowed": False, "reason": "membership_evidence_required"}
- )
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is None:
- local = [
- int(item["chat_id"])
- for item in candidates
- if str(item.get("bot_id")) == str(getattr(wbb, "BOT_PROFILE_ID", "primary"))
- ]
- return success(
- await verify_local_membership(
- user_id=user_id,
- chat_ids=local,
- require_admin=require_admin,
- )
- )
- candidates_by_bot: dict[str, list[int]] = {}
- for item in candidates:
- candidates_by_bot.setdefault(str(item["bot_id"]), []).append(
- int(item["chat_id"])
- )
- for bot_id, chat_ids in candidates_by_bot.items():
- endpoint = supervisor.endpoint_for(bot_id)
- if endpoint is None:
- continue
- try:
- async with ClientSession(timeout=ClientTimeout(total=15)) as session:
- async with session.post(
- f"{endpoint}/api/internal/v1/directory/verify-local",
- headers={
- "Authorization": f"Bearer {supervisor.internal_token}"
- },
- json={
- "user_id": str(user_id),
- "chat_ids": [str(value) for value in chat_ids],
- "require_admin": require_admin,
- },
- ) as response:
- if response.status != 200:
- continue
- result = (await response.json()).get("data") or {}
- if result.get("allowed"):
- return success(result)
- except (ClientError, TimeoutError, ValueError):
- continue
- return success({"allowed": False})
- async def internal_directory_verify_local(
- self,
- request: web.Request,
- ) -> web.Response:
- _require_internal_request(request)
- body = await json_body(request)
- user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
- raw_chat_ids = body.get("chat_ids")
- if not isinstance(raw_chat_ids, list):
- raise ApiProblem("invalid_parameter", "chat_ids 必须是数组。")
- chat_ids = [parse_int(value, "chat_id") for value in raw_chat_ids[:200]]
- return success(
- await verify_local_membership(
- user_id=user_id,
- chat_ids=chat_ids,
- require_admin=bool(body.get("require_admin")),
- )
- )
- async def internal_directory_notify_local(
- self,
- request: web.Request,
- ) -> web.Response:
- _require_internal_request(request)
- body = await json_body(request)
- user_id = parse_int(body.get("user_id"), "user_id", minimum=1)
- text = str(body.get("text") or "").strip()
- if not text:
- raise ApiProblem("invalid_parameter", "通知内容不能为空。")
- try:
- await telegram_app.send_message(user_id, text)
- except Exception as exc:
- raise ApiProblem(
- "telegram_notification_failed",
- "无法向该用户发送 Telegram 通知。",
- status=409,
- ) from exc
- return success({"sent": True})
- async def _notify_directory_user(
- self,
- *,
- profile: dict[str, Any],
- text: str,
- ) -> None:
- bot_id = str(
- profile.get("application_source_bot_id")
- or getattr(wbb, "BOT_PROFILE_ID", "primary")
- )
- supervisor = getattr(wbb, "BOT_SUPERVISOR", None)
- if supervisor is None:
- with suppress(Exception):
- await telegram_app.send_message(int(profile["user_id"]), text)
- return
- endpoint = supervisor.endpoint_for(bot_id)
- if endpoint is None:
- return
- with suppress(Exception):
- async with ClientSession(timeout=ClientTimeout(total=10)) as session:
- await session.post(
- f"{endpoint}/api/internal/v1/directory/notify-local",
- headers={
- "Authorization": f"Bearer {supervisor.internal_token}"
- },
- json={"user_id": str(profile["user_id"]), "text": text},
- )
- async def directory_settings(self, _: web.Request) -> web.Response:
- return success(await get_directory_settings())
- async def directory_settings_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- values = {
- key: body[key]
- for key in (
- "nearby_radius_km",
- "nearby_radius_options_km",
- )
- if key in body
- }
- saved = await set_directory_settings(values)
- set_audit(
- request,
- "directory.settings.update",
- summary="更新技师目录设置",
- )
- return success(saved)
- async def directory_applications(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_teacher_applications(
- status=str(request.query.get("status") or ""),
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def directory_application_action(
- self,
- request: web.Request,
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- profile = await decide_teacher_application(
- user_id=user_id,
- action=action,
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- source="web",
- reason=reason,
- )
- labels = {"approve": "批准", "reject": "拒绝", "revoke": "撤销"}
- set_audit(
- request,
- f"directory.application.{action}",
- target_id=user_id,
- summary=f"{labels.get(action, action)}技师申请",
- metadata={"reason": reason},
- )
- notification = {
- "approve": "你的技师申请已通过,可以在菜单中自助上榜。",
- "reject": f"你的技师申请未通过。原因:{reason}",
- "revoke": f"你的技师资格已被撤销。原因:{reason}",
- }.get(action)
- if notification:
- await self._notify_directory_user(profile=profile, text=notification)
- return success(profile)
- async def directory_profiles(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- listed_raw = request.query.get("listed")
- listed = (
- listed_raw.lower() in {"1", "true"}
- if listed_raw is not None
- else None
- )
- items, total = await list_directory_profiles(
- query=str(request.query.get("query") or ""),
- application_status=str(
- request.query.get("application_status") or APPLICATION_APPROVED
- ),
- listed=listed,
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def directory_profile_action(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "强制状态操作必须填写原因。")
- profile, _ = await set_teacher_state(
- user_id=user_id,
- action=action,
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- source="web",
- reason=reason,
- )
- set_audit(
- request,
- f"directory.profile.{action}",
- target_id=user_id,
- summary=f"强制修改技师状态:{action}",
- metadata={"reason": reason},
- )
- await self._notify_directory_user(
- profile=profile,
- text=f"管理员已将你的技师状态修改为“{action}”。原因:{reason}",
- )
- return success(profile)
- async def directory_locations(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_directory_locations(
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- safe_items = []
- for item in items:
- profile = item.get("profile") or {}
- safe_items.append(
- {
- "user_id": item["user_id"],
- "longitude": item.get("longitude"),
- "latitude": item.get("latitude"),
- "source": item.get("source"),
- "updated_at": item.get("updated_at"),
- "display_name": profile.get("display_name"),
- "username": profile.get("username"),
- "application_status": profile.get("application_status", "none"),
- }
- )
- return success(
- {
- "items": safe_items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def directory_location_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "修正位置必须填写原因。")
- profile = await get_directory_profile(user_id)
- if not profile:
- raise ApiProblem("profile_not_found", "未找到该成员资料。", status=404)
- saved = await save_directory_location(
- user_id=user_id,
- longitude=body.get("longitude"),
- latitude=body.get("latitude"),
- source="web",
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- )
- set_audit(
- request,
- "directory.location.update",
- target_id=user_id,
- summary="修正成员位置",
- metadata={"reason": reason},
- )
- return success(
- {
- "user_id": user_id,
- "longitude": saved["longitude"],
- "latitude": saved["latitude"],
- "source": saved.get("source"),
- "updated_at": saved.get("updated_at"),
- }
- )
- async def directory_location_delete(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "清除位置必须填写原因。")
- deleted = await clear_directory_location(
- user_id=user_id,
- actor_id=request["admin"]["username"],
- actor_name=request["admin"]["username"],
- source="web",
- reason=reason,
- )
- set_audit(
- request,
- "directory.location.delete",
- target_id=user_id,
- summary="清除成员位置",
- metadata={"reason": reason},
- )
- return success({"deleted": deleted})
- async def directory_events(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_directory_events(
- query=str(request.query.get("query") or ""),
- event_type=str(request.query.get("event_type") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- }
- )
- async def service_settings(self, _: web.Request) -> web.Response:
- return success(await get_service_settings())
- async def service_settings_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- saved = await set_service_settings(body)
- set_audit(request, "service.settings.update", summary="更新技师咨询与评价设置")
- return success(saved)
- async def service_package_templates(self, request: web.Request) -> web.Response:
- items = await list_package_templates(
- include_archived=request.query.get("include_archived", "false").lower()
- == "true"
- )
- return success({"items": items, "total": len(items)})
- async def service_package_template_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- created = await create_package_template(
- body,
- actor_id=request["admin"]["username"],
- )
- set_audit(
- request,
- "service.package_template.create",
- target_id=created["template_id"],
- summary="创建套餐模板",
- )
- return success(created, status=201)
- async def service_package_template_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- template_id = request.match_info["template_id"]
- updated = await update_package_template(
- template_id,
- body,
- actor_id=request["admin"]["username"],
- )
- set_audit(
- request,
- "service.package_template.update",
- target_id=template_id,
- summary="更新套餐模板并创建新版本",
- )
- return success(updated)
- async def service_package_template_action(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- template_id = request.match_info["template_id"]
- action = str(body.get("action") or "")
- if action == "duplicate":
- updated = await duplicate_package_template(
- template_id,
- actor_id=request["admin"]["username"],
- )
- else:
- updated = await set_package_template_status(
- template_id,
- action,
- actor_id=request["admin"]["username"],
- )
- set_audit(
- request,
- f"service.package_template.{action}",
- target_id=template_id,
- summary=f"套餐模板操作:{action}",
- )
- return success(updated)
- async def technician_service_profile(self, request: web.Request) -> web.Response:
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- return success(await get_technician_service_profile(user_id))
- async def technician_service_profile_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- values = body.get("service_profile") or {
- key: value for key, value in body.items() if key != "confirm"
- }
- updated = await publish_technician_profile(
- user_id,
- values,
- actor_id=request["admin"]["username"],
- )
- set_audit(
- request,
- "service.technician_profile.update",
- target_id=user_id,
- summary="更新技师服务资料",
- )
- with suppress(Exception):
- await ensure_technician_review_topic(
- _request_bot_token(request),
- user_id,
- )
- return success(updated)
- async def technician_mini_app_bootstrap(
- self,
- request: web.Request,
- ) -> web.Response:
- identity = request["technician"]
- context = await get_technician_self_service_context(identity["user_id"])
- templates = await list_package_templates()
- public_templates = [
- {
- "template_id": item["template_id"],
- "version": item["version"],
- "name": item["name"],
- "builtin": bool(item.get("builtin")),
- "package": item["package"],
- }
- for item in templates
- if item.get("status") == "enabled"
- ]
- response = success(
- {
- "technician": context,
- "templates": public_templates,
- "limits": {"max_packages": 5},
- }
- )
- response.headers["Cache-Control"] = "no-store"
- return response
- async def technician_mini_app_profile_update(
- self,
- request: web.Request,
- ) -> web.Response:
- identity = request["technician"]
- body = await json_body(request)
- values = body.get("service_profile") or body
- updated = await publish_technician_profile(
- identity["user_id"],
- values,
- actor_id=identity["user_id"],
- )
- set_audit(
- request,
- "service.technician_self_profile.update",
- target_id=identity["user_id"],
- summary="技师通过 Mini App 更新服务资料和套餐",
- )
- with suppress(Exception):
- await ensure_technician_review_topic(
- _request_bot_token(request),
- int(identity["user_id"]),
- )
- return success(
- await get_technician_self_service_context(updated["user_id"])
- )
- async def technician_mini_app_availability(
- self,
- request: web.Request,
- ) -> web.Response:
- identity = request["technician"]
- body = await json_body(request)
- if not isinstance(body.get("accepting_requests"), bool):
- raise ApiProblem(
- "invalid_accepting_requests",
- "接单状态必须是布尔值。",
- )
- updated = await set_technician_accepting_requests(
- identity["user_id"],
- body["accepting_requests"],
- )
- set_audit(
- request,
- "service.technician_self_availability.update",
- target_id=identity["user_id"],
- summary=(
- "开启接单"
- if body["accepting_requests"]
- else "暂停接单"
- ),
- )
- return success(
- await get_technician_self_service_context(updated["user_id"])
- )
- async def service_orders(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_service_orders(
- query=str(request.query.get("query") or ""),
- status=str(request.query.get("status") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def service_order_detail(self, request: web.Request) -> web.Response:
- return success(await get_admin_service_order(request.match_info["order_id"]))
- async def service_order_action(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "争议处理或作废必须填写原因。")
- order_id = request.match_info["order_id"]
- updated = await resolve_service_dispute(
- order_id,
- actor_id=request["admin"]["username"],
- action=action,
- reason=reason,
- )
- set_audit(
- request,
- f"service.order.{action}",
- target_id=order_id,
- summary=f"服务单操作:{action}",
- metadata={"reason": reason},
- )
- return success(updated)
- async def service_order_address(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "查看精确地址必须填写原因。")
- order_id = request.match_info["order_id"]
- address = await admin_reveal_order_address(
- order_id,
- actor_id=request["admin"]["username"],
- reason=reason,
- )
- set_audit(
- request,
- "service.order.address_reveal",
- target_id=order_id,
- summary="管理员查看服务单精确地址",
- metadata={"reason": reason},
- )
- return success(address)
- async def service_customer_block(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- customer_id = parse_int(
- request.match_info["customer_id"], "customer_id", minimum=1
- )
- active = bool(body.get("active", True))
- reason = str(body.get("reason") or "").strip()
- if active and not reason:
- raise ApiProblem("reason_required", "封禁顾客必须填写原因。")
- updated = await set_customer_block(
- customer_id=customer_id,
- scope="global",
- technician_id=None,
- active=active,
- actor_id=request["admin"]["username"],
- reason=reason,
- )
- set_audit(
- request,
- "service.customer.block" if active else "service.customer.unblock",
- target_id=customer_id,
- summary="全局封禁顾客" if active else "解除顾客封禁",
- metadata={"reason": reason},
- )
- return success(updated)
- async def service_customer_blocks(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- active_raw = request.query.get("active")
- active = (
- active_raw.lower() in {"1", "true"}
- if active_raw is not None
- else None
- )
- items, total = await list_customer_blocks(
- scope=str(request.query.get("scope") or "global"),
- active=active,
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def service_reports(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_service_reports(
- status=str(request.query.get("status") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def service_report_action(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "处理举报必须填写原因。")
- report = await moderate_service_report(
- request.match_info["report_id"],
- action=action,
- actor_id=request["admin"]["username"],
- reason=reason,
- )
- set_audit(
- request,
- f"service.report.{action}",
- target_id=report["report_id"],
- summary=f"处理服务举报:{action}",
- metadata={"reason": reason},
- )
- return success(report)
- async def service_review_template(self, _: web.Request) -> web.Response:
- return success(await get_active_review_template())
- async def service_review_template_publish(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- published = await publish_review_template(
- body,
- actor_id=request["admin"]["username"],
- )
- set_audit(
- request,
- "service.review_template.publish",
- target_id=published["template_id"],
- summary="发布新版评价模板",
- )
- return success(published, status=201)
- async def service_reviews(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_reviews(
- status=str(request.query.get("status") or ""),
- source=str(request.query.get("source") or ""),
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def service_review_action(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- action = str(body.get("action") or "")
- reason = str(body.get("reason") or "").strip()
- updated = await moderate_review(
- request.match_info["review_id"],
- action,
- actor_id=request["admin"]["username"],
- reason=reason,
- )
- publication = None
- if action == "approve":
- publication = await publish_approved_review(
- _request_bot_token(request),
- updated,
- )
- updated = {**updated, "topic_publication": publication}
- set_audit(
- request,
- f"service.review.{action}",
- target_id=request.match_info["review_id"],
- summary=f"评价审核:{action}",
- metadata={"reason": reason, "topic_publication": publication},
- )
- with suppress(Exception):
- await telegram_app.send_message(
- int(updated["customer_id"]),
- (
- "你的评价已审核通过。"
- if action == "approve"
- else f"你的评价未通过审核。原因:{reason}"
- ),
- )
- return success(updated)
- async def service_review_qrs(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_review_qr_records(
- status=str(request.query.get("status") or ""),
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def service_leaderboard(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_leaderboard(
- category=str(request.query.get("category") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {
- "items": items,
- "total": total,
- "page": page,
- "page_size": page_size,
- "rule": "v/(v+5) × 技师平均分 + 5/(v+5) × 分类平均分;至少 3 条批准评价。",
- "title": "审核评价排行",
- }
- )
- async def service_fulfillment_metrics(self, _: web.Request) -> web.Response:
- return success(await fulfillment_metrics())
- def _set_session_cookie(self, response: web.Response, token: str) -> None:
- response.set_cookie(
- SESSION_COOKIE,
- token,
- httponly=True,
- secure=self.cookie_secure,
- samesite="Strict",
- max_age=self.session_hours * 3600,
- path="/",
- )
- async def _new_session(self, request: web.Request, username: str) -> tuple[str, str]:
- token = generate_session_token()
- csrf = generate_csrf_token()
- await create_admin_session(
- username=username,
- token_hash=hash_token(token),
- csrf_token=csrf,
- expires_at=datetime.now(UTC) + timedelta(hours=self.session_hours),
- remote_address=request.remote or "",
- user_agent=request.headers.get("User-Agent", ""),
- )
- return token, csrf
- async def login(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- username = str(body.get("username") or "").strip()
- password = str(body.get("password") or "")
- set_audit(request, "auth.login", summary=f"Login as {username}")
- user = await get_admin_user(username)
- if user and int(user.get("failed_login_count", 0)) >= 5:
- last_failed = user.get("last_failed_login_at")
- if last_failed and datetime.now(UTC) - as_utc(last_failed) < timedelta(minutes=15):
- raise ApiProblem(
- "login_rate_limited", "登录失败次数过多,请 15 分钟后再试。", status=429
- )
- if not user or not verify_password(password, user.get("password_hash", "")):
- await record_login_failure(username)
- raise ApiProblem("invalid_credentials", "用户名或密码错误。", status=401)
- await record_login_success(username)
- token, csrf = await self._new_session(request, username)
- response = success(
- {
- "username": username,
- "must_change_password": bool(user.get("must_change_password")),
- "csrf_token": csrf,
- }
- )
- self._set_session_cookie(response, token)
- request["admin"] = {"username": username}
- return response
- async def me(self, request: web.Request) -> web.Response:
- return success(
- {
- "username": request["admin"]["username"],
- "must_change_password": request["admin"]["must_change_password"],
- "csrf_token": request["admin"]["csrf_token"],
- }
- )
- async def logout(self, request: web.Request) -> web.Response:
- set_audit(request, "auth.logout", summary="Logout")
- await revoke_admin_session(request["admin"]["session_token_hash"])
- response = success({"logged_out": True})
- response.del_cookie(SESSION_COOKIE, path="/")
- return response
- async def change_password(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- current = str(body.get("current_password") or "")
- new_password = str(body.get("new_password") or "")
- username = request["admin"]["username"]
- user = await get_admin_user(username)
- if not user or not verify_password(current, user["password_hash"]):
- raise ApiProblem("invalid_current_password", "当前密码错误。", status=409)
- if verify_password(new_password, user["password_hash"]):
- raise ApiProblem("password_unchanged", "新密码不能与当前密码相同。", status=409)
- errors = validate_new_password(new_password, username)
- if errors:
- raise ApiProblem("weak_password", "新密码不符合要求。", details=errors)
- set_audit(request, "auth.password.change", summary="Change administrator password")
- await update_admin_password(username, hash_password(new_password))
- token, csrf = await self._new_session(request, username)
- response = success(
- {"username": username, "must_change_password": False, "csrf_token": csrf}
- )
- self._set_session_cookie(response, token)
- return response
- async def dashboard(self, _: web.Request) -> web.Response:
- counts = await dashboard_counts()
- giveaways, giveaway_total = await list_giveaways_page(
- page=1,
- page_size=5,
- all_bots=bool(getattr(wbb, "SUPERVISOR_MODE", False)),
- )
- accounts, account_total = await list_point_accounts(page=1, page_size=5)
- transactions, transaction_total = await list_point_transactions(page=1, page_size=8)
- return success(
- {
- "counts": {
- **counts,
- "giveaways": giveaway_total,
- "point_accounts": account_total,
- "point_transactions": transaction_total,
- },
- "recent_giveaways": giveaways,
- "top_accounts": accounts,
- "recent_transactions": transactions,
- }
- )
- async def system_settings(self, _: web.Request) -> web.Response:
- bot_connected = bool(getattr(wbb, "TELEGRAM_CONNECTED", False))
- return success(
- {
- "web_address": f"http://127.0.0.1:{getattr(wbb, 'ADMIN_WEB_PORT', 8088)}/admin",
- "bind_host": str(getattr(wbb, "ADMIN_WEB_HOST", "0.0.0.0")),
- "upload_limit_mb": int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20)),
- "supervisor_mode": bool(getattr(wbb, "SUPERVISOR_MODE", False)),
- "bot": {
- "connected": bot_connected,
- "id": str(wbb.BOT_ID) if bot_connected else "",
- "username": wbb.BOT_USERNAME if bot_connected else "",
- "name": wbb.BOT_NAME if bot_connected else "",
- },
- }
- )
- def _supervisor(self):
- return getattr(wbb, "BOT_SUPERVISOR", None)
- def _telegram_status(self) -> dict[str, Any]:
- supervisor = self._supervisor()
- runtimes = supervisor.runtimes() if supervisor is not None else {}
- return telegram_config_status(self.bot_config_path, runtimes=runtimes)
- async def telegram_settings(self, _: web.Request) -> web.Response:
- return success(self._telegram_status())
- async def telegram_settings_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "telegram.settings.update", summary="Update Telegram API credentials")
- body.pop("confirm", None)
- update_telegram_config(self.bot_config_path, body)
- supervisor = self._supervisor()
- if supervisor is not None:
- status = telegram_config_status(self.bot_config_path)
- for profile in status["bots"]:
- if profile["enabled"] and profile["ready_to_connect"]:
- await supervisor.reconcile(profile["bot_id"])
- return success(self._telegram_status())
- async def bots(self, _: web.Request) -> web.Response:
- return success({"items": self._telegram_status()["bots"]})
- async def bot_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "bot.create", summary=str(body.get("label") or ""))
- body.pop("confirm", None)
- profile = create_bot_profile(self.bot_config_path, body)
- supervisor = self._supervisor()
- if supervisor is not None and profile["enabled"] and profile["ready_to_connect"]:
- await supervisor.start(profile["bot_id"])
- current = next(
- item
- for item in self._telegram_status()["bots"]
- if item["bot_id"] == profile["bot_id"]
- )
- return success(current, status=201)
- async def bot_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- set_audit(request, "bot.update", target_id=bot_id, summary="Update Bot profile")
- body.pop("confirm", None)
- update_bot_profile(self.bot_config_path, bot_id, body)
- supervisor = self._supervisor()
- if supervisor is not None:
- await supervisor.reconcile(bot_id)
- current = next(
- item for item in self._telegram_status()["bots"] if item["bot_id"] == bot_id
- )
- return success(current)
- async def bot_delete(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
- set_audit(request, "bot.delete", target_id=bot_id, summary=str(profile.get("label") or ""))
- supervisor = self._supervisor()
- if supervisor is not None:
- await supervisor.stop(bot_id)
- delete_bot_profile(self.bot_config_path, bot_id)
- return success({"deleted": True})
- async def bot_test(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- profile = get_bot_profile_secrets(self.bot_config_path, bot_id)
- set_audit(request, "bot.test", target_id=bot_id, summary="Validate Bot Token")
- identity = await test_bot_token(str(profile.get("bot_token") or ""))
- store_bot_identity(self.bot_config_path, bot_id, identity)
- return success(identity)
- async def bot_restart(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- bot_id = request.match_info["bot_id"]
- set_audit(request, "bot.restart", target_id=bot_id, summary="Restart Bot worker")
- supervisor = self._supervisor()
- if supervisor is None:
- raise ApiProblem(
- "bot_supervisor_unavailable",
- "机器人监管服务尚未就绪。",
- status=503,
- )
- return success(await supervisor.restart(bot_id))
- async def roles(self, _: web.Request) -> web.Response:
- status = self._telegram_status()
- return success(
- {
- "items": status["roles"],
- "permission_catalog": status["permission_catalog"],
- }
- )
- async def role_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- set_audit(
- request,
- "bot.role.create",
- summary=str(body.get("name") or ""),
- )
- body.pop("confirm", None)
- return success(create_bot_role(self.bot_config_path, body), status=201)
- async def role_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- role_id = request.match_info["role_id"]
- set_audit(
- request,
- "bot.role.update",
- target_id=role_id,
- summary=str(body.get("name") or ""),
- )
- body.pop("confirm", None)
- role = update_bot_role(self.bot_config_path, role_id, body)
- supervisor = self._supervisor()
- if supervisor is not None:
- for profile in self._telegram_status()["bots"]:
- if role_id in profile["role_ids"]:
- await supervisor.reconcile(profile["bot_id"])
- return success(role)
- async def role_delete(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- role_id = request.match_info["role_id"]
- set_audit(
- request,
- "bot.role.delete",
- target_id=role_id,
- summary="Delete Bot role",
- )
- delete_bot_role(self.bot_config_path, role_id)
- return success({"deleted": True})
- async def audit_logs(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- items, total = await list_audit_logs(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- action=request.query.get("action") or None,
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def upload_media(self, request: web.Request) -> web.Response:
- if not MESSAGE_DUMP_CHAT:
- raise ApiProblem(
- "message_dump_chat_missing",
- "请先配置媒体中转群。",
- status=409,
- )
- reader = await request.multipart()
- part = await reader.next()
- if not part or part.name != "file":
- raise ApiProblem("file_required", "请选择要上传的文件。")
- filename = Path(part.filename or "upload.bin").name
- content_type = part.headers.get("Content-Type", "application/octet-stream")
- size = 0
- temporary_path: Path | None = None
- try:
- with tempfile.NamedTemporaryFile(delete=False, suffix=Path(filename).suffix) as handle:
- temporary_path = Path(handle.name)
- while True:
- chunk = await part.read_chunk(1024 * 1024)
- if not chunk:
- break
- size += len(chunk)
- if size > self.upload_limit:
- raise ApiProblem("file_too_large", "上传文件超过 20 MB 限制。", status=413)
- handle.write(chunk)
- if content_type == "image/gif":
- sent = await telegram_app.send_animation(MESSAGE_DUMP_CHAT, str(temporary_path))
- media = sent.animation or sent.document
- if not media:
- raise ApiProblem("media_upload_failed", "Telegram 未能保存该动图。", status=502)
- media_type = "animation" if sent.animation else "document"
- file_id = media.file_id
- elif content_type.startswith("image/"):
- media_type = "photo"
- sent = await telegram_app.send_photo(MESSAGE_DUMP_CHAT, str(temporary_path))
- file_id = sent.photo.file_id
- elif content_type.startswith("video/"):
- media_type = "video"
- sent = await telegram_app.send_video(MESSAGE_DUMP_CHAT, str(temporary_path))
- file_id = sent.video.file_id
- else:
- media_type = "document"
- sent = await telegram_app.send_document(
- MESSAGE_DUMP_CHAT, str(temporary_path), file_name=filename
- )
- file_id = sent.document.file_id
- finally:
- if temporary_path:
- temporary_path.unlink(missing_ok=True)
- set_audit(request, "media.upload", summary=filename, target_id=file_id)
- return success(
- {
- "file_id": file_id,
- "type": media_type,
- "filename": filename,
- "size": size,
- "mime_type": content_type,
- },
- status=201,
- )
- async def channels(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_channels(
- query=request.query.get("query", ""), page=page, page_size=page_size
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def channel_add(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = parse_int(body.get("chat_id"), "chat_id")
- set_audit(request, "channel.connect", chat_id=chat_id)
- return success(await ensure_channel(chat_id), status=201)
- async def channel(self, request: web.Request) -> web.Response:
- return success(await ensure_channel(chat_id_param(request), include_photo=True))
- async def channel_profile(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "channel.profile.update", chat_id=chat_id)
- return success(await update_channel_profile(
- chat_id, title=body.get("title"), description=body.get("description")
- ))
- async def channel_photo(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- reader = await request.multipart()
- part = await reader.next()
- if not part or part.name != "photo":
- raise ApiProblem("photo_required", "请选择 JPG 或 PNG 头像。")
- filename = Path(part.filename or "photo").name
- suffix = Path(filename).suffix.lower()
- if suffix not in {".jpg", ".jpeg", ".png"}:
- raise ApiProblem("invalid_photo", "头像只能是 JPG 或 PNG。")
- temporary_path: Path | None = None
- try:
- with tempfile.NamedTemporaryFile(delete=False, suffix=suffix) as handle:
- temporary_path = Path(handle.name)
- size = 0
- while chunk := await part.read_chunk(1024 * 1024):
- size += len(chunk)
- if size > 5_000_000:
- raise ApiProblem("photo_too_large", "头像不能超过 5 MB。", status=413)
- handle.write(chunk)
- from PIL import Image
- try:
- with Image.open(temporary_path) as photo:
- if photo.format not in {"JPEG", "PNG"}:
- raise ValueError("unsupported image format")
- photo.verify()
- except Exception as exc:
- raise ApiProblem("invalid_photo", "头像文件无效。") from exc
- set_audit(request, "channel.photo.update", chat_id=chat_id, summary=filename)
- return success(await update_channel_photo(chat_id, temporary_path))
- finally:
- if temporary_path:
- temporary_path.unlink(missing_ok=True)
- async def channel_admins(self, request: web.Request) -> web.Response:
- return success({"items": await list_channel_admins(chat_id_param(request))})
- async def channel_admin_candidates(self, request: web.Request) -> web.Response:
- limit = parse_int(request.query.get("limit", "30"), "limit", minimum=1)
- return success({"items": await list_channel_admin_candidates(
- chat_id_param(request), query=request.query.get("query", ""), limit=limit,
- )})
- async def channel_admin_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- body = await json_body(request)
- require_confirmation(body)
- privileges = body.get("privileges")
- if not isinstance(privileges, list) or not all(isinstance(item, str) for item in privileges):
- raise ApiProblem("invalid_privileges", "管理员权限须为字符串列表。")
- set_audit(request, "channel.admin.update", chat_id=chat_id, target_id=user_id)
- return success(await set_channel_admin(chat_id, user_id, privileges))
- async def channel_admin_remove(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- user_id = parse_int(request.match_info["user_id"], "user_id", minimum=1)
- require_confirmation(await json_body(request))
- set_audit(request, "channel.admin.remove", chat_id=chat_id, target_id=user_id)
- return success(await remove_channel_admin(chat_id, user_id))
- async def channel_invites(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- await ensure_channel(chat_id)
- return success({"items": await get_invite_links(chat_id)})
- async def channel_invite_create(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- await ensure_channel(chat_id)
- expires_at = parse_datetime(body["expires_at"], "expires_at") if body.get("expires_at") else None
- member_limit = parse_int(body["member_limit"], "member_limit", minimum=1) if body.get("member_limit") else None
- set_audit(request, "channel.invite.create", chat_id=chat_id)
- try:
- result = await create_invite_link(
- chat_id, name=str(body.get("name") or "Admin panel"),
- expires_at=expires_at, member_limit=member_limit,
- )
- except Exception as exc:
- raise ChatManagementError("channel_invite_failed", "Telegram 未能创建频道邀请链接。", status=502) from exc
- return success(result, status=201)
- async def channel_invite_revoke(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- await ensure_channel(chat_id)
- link = str(body.get("invite_link") or "")
- set_audit(request, "channel.invite.revoke", chat_id=chat_id)
- try:
- await revoke_invite_link(chat_id, link)
- except Exception as exc:
- raise ChatManagementError("channel_invite_revoke_failed", "Telegram 未能撤销频道邀请链接。", status=502) from exc
- return success({"revoked": True})
- async def channel_posts(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- await ensure_channel(chat_id)
- page, page_size = page_params(request)
- items, total = await list_posts(chat_id, page=page, page_size=page_size)
- return success({"items": [public_post(item) for item in items], "total": total, "page": page, "page_size": page_size})
- async def channel_post_create(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "channel.post.create", chat_id=chat_id, summary=str(body.get("text") or "")[:200])
- return success(await create_channel_post(chat_id, body), status=201)
- async def channel_post_edit(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "channel.post.edit", chat_id=chat_id, target_id=request.match_info["post_id"])
- return success(await edit_channel_post(chat_id, request.match_info["post_id"], body))
- async def channel_post_delete(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- require_confirmation(await json_body(request))
- set_audit(request, "channel.post.delete", chat_id=chat_id, target_id=request.match_info["post_id"])
- return success(await delete_channel_post(chat_id, request.match_info["post_id"]))
- async def chats(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_accessible_chats(
- query=request.query.get("query", ""), page=page, page_size=page_size
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def chat(self, request: web.Request) -> web.Response:
- return success(await get_chat_overview(chat_id_param(request)))
- async def chat_delete(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- set_audit(
- request,
- "chat.record.remove",
- chat_id=chat_id,
- summary="Remove unavailable chat from management list",
- )
- return success(await remove_unavailable_chat(chat_id))
- async def chat_profile(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.profile.update", chat_id=chat_id, summary="Update group profile")
- return success(
- await update_chat_profile(
- chat_id, title=body.get("title"), description=body.get("description")
- )
- )
- async def chat_permissions(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.permissions.update", chat_id=chat_id, summary="Update default member permissions")
- return success(await update_chat_permissions(chat_id, body.get("permissions") or {}))
- async def announcement(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.announcement.send", chat_id=chat_id, summary=str(body.get("text") or "")[:200])
- return success(
- await send_announcement(
- chat_id,
- text=str(body.get("text") or ""),
- media_type=body.get("media_type"),
- file_id=body.get("file_id"),
- pin=bool(body.get("pin")),
- ),
- status=201,
- )
- async def chat_admins(self, request: web.Request) -> web.Response:
- return success({"items": await list_chat_admins(chat_id_param(request))})
- async def member_search(self, request: web.Request) -> web.Response:
- return success(
- {
- "items": await search_chat_members(
- chat_id_param(request),
- request.query.get("query", ""),
- limit=parse_int(request.query.get("limit", "20"), "limit"),
- )
- }
- )
- async def recent_members(self, request: web.Request) -> web.Response:
- return success(
- {
- "items": await list_recent_members(
- chat_id_param(request),
- limit=parse_int(request.query.get("limit", "30"), "limit"),
- )
- }
- )
- async def member_identity_changes(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_member_identity_changes(
- chat_id_param(request), page=page, page_size=page_size
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def member_action(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- user_id = parse_int(request.match_info["user_id"], "userId")
- body = await json_body(request)
- require_confirmation(body)
- action = str(body.get("action") or "")
- set_audit(
- request,
- f"chat.member.{action}",
- chat_id=chat_id,
- target_id=user_id,
- summary=str(body.get("reason") or ""),
- )
- duration = body.get("duration_seconds")
- return success(
- await execute_member_action(
- chat_id,
- user_id=user_id,
- action=action,
- reason=str(body.get("reason") or ""),
- duration_seconds=parse_int(duration, "duration_seconds", minimum=60) if duration else None,
- privileges=body.get("privileges") or {},
- )
- )
- async def invites(self, request: web.Request) -> web.Response:
- return success({"items": await get_invite_links(chat_id_param(request))})
- async def invite_create(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- expires_at = parse_datetime(body["expires_at"], "expires_at") if body.get("expires_at") else None
- member_limit = parse_int(body["member_limit"], "member_limit", minimum=1) if body.get("member_limit") else None
- set_audit(request, "chat.invite.create", chat_id=chat_id, summary=str(body.get("name") or ""))
- return success(
- await create_invite_link(
- chat_id,
- name=str(body.get("name") or "Admin panel"),
- expires_at=expires_at,
- member_limit=member_limit,
- ),
- status=201,
- )
- async def invite_revoke(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- invite_link = str(body.get("invite_link") or "")
- if not invite_link:
- raise ApiProblem("invite_link_required", "缺少邀请链接。")
- set_audit(request, "chat.invite.revoke", chat_id=chat_id, summary=invite_link)
- await revoke_invite_link(chat_id, invite_link)
- return success({"revoked": True})
- async def rules(self, request: web.Request) -> web.Response:
- return success({"rules": await get_rules(chat_id_param(request))})
- async def rules_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- await ensure_permission(chat_id, "can_change_info")
- rules = str(body.get("rules") or "")
- if len(rules) > 4000:
- raise ApiProblem("rules_too_long", "群规不能超过 4000 个字符。")
- set_audit(request, "chat.rules.update", chat_id=chat_id, summary="Update group rules")
- await set_chat_rules(chat_id, rules)
- return success({"rules": rules})
- async def automation(self, request: web.Request) -> web.Response:
- return success(await get_automation_settings(chat_id_param(request)))
- async def automation_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- set_audit(request, "chat.automation.update", chat_id=chat_id, summary="Update automation rules")
- body.pop("confirm", None)
- return success(await apply_automation_settings(chat_id, body))
- async def points_settings(self, request: web.Request) -> web.Response:
- return success(await get_point_rules(chat_id_param(request)))
- async def points_settings_update(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- await ensure_permission(chat_id, "can_change_info")
- set_audit(request, "points.settings.update", chat_id=chat_id, summary="Update point rules")
- body.pop("confirm", None)
- return success(await apply_point_rules(chat_id, body))
- async def point_accounts(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- items, total = await list_point_accounts(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- query=request.query.get("query", ""),
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def point_leaderboard(self, request: web.Request) -> web.Response:
- if not request.query.get("chat_id"):
- raise ApiProblem("chat_id_required", "排行榜必须选择群组。")
- return await self.point_accounts(request)
- async def point_transactions(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- user_raw = request.query.get("user_id")
- created_from = parse_datetime(request.query["created_from"], "created_from") if request.query.get("created_from") else None
- created_to = parse_datetime(request.query["created_to"], "created_to") if request.query.get("created_to") else None
- items, total = await list_point_transactions(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- user_id=parse_int(user_raw, "user_id") if user_raw else None,
- source=request.query.get("source") or None,
- created_from=created_from,
- created_to=created_to,
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def point_adjustment(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = parse_int(body.get("chat_id"), "chat_id")
- user_id = parse_int(body.get("user_id"), "user_id")
- operation = str(body.get("operation") or "")
- amount = parse_int(body.get("amount"), "amount", minimum=0)
- reason = str(body.get("reason") or "").strip()
- if not reason:
- raise ApiProblem("reason_required", "积分调整必须填写原因。")
- if operation not in {"add", "deduct", "set"}:
- raise ApiProblem(
- "invalid_operation",
- "积分操作必须是增加、扣减或设置余额。",
- )
- try:
- user = await telegram_app.get_users(user_id)
- username, first_name = user.username, user.first_name
- display_name = " ".join(
- value for value in (user.first_name, user.last_name) if value
- )
- except Exception:
- username, first_name = body.get("username"), body.get("first_name")
- display_name = body.get("display_name") or first_name
- request_key = str(body.get("request_id") or generate_session_token())
- set_audit(
- request,
- f"points.adjust.{operation}",
- chat_id=chat_id,
- target_id=user_id,
- summary=reason,
- metadata={"amount": amount},
- )
- if operation == "set":
- account, created = await set_points(
- chat_id=chat_id,
- user_id=user_id,
- balance=amount,
- actor_id=request["admin"]["username"],
- reason=reason,
- idempotency_key=f"web-points:{request_key}",
- username=username,
- first_name=first_name,
- display_name=display_name,
- )
- else:
- account, created = await adjust_points(
- chat_id=chat_id,
- user_id=user_id,
- delta=amount if operation == "add" else -amount,
- source=SOURCE_ADMIN,
- idempotency_key=f"web-points:{request_key}",
- actor_id=request["admin"]["username"],
- reason=reason,
- username=username,
- first_name=first_name,
- display_name=display_name,
- )
- return success({"account": account, "created": created}, status=201 if created else 200)
- async def points_export(self, request: web.Request) -> web.Response:
- chat_raw = request.query.get("chat_id")
- if not chat_raw:
- raise ApiProblem("chat_id_required", "导出积分数据必须选择群组。")
- chat_id = parse_int(chat_raw, "chat_id")
- kind = request.query.get("kind", "accounts")
- output = io.StringIO()
- if kind == "accounts":
- items: list[dict[str, Any]] = []
- page = 1
- while True:
- batch, total = await list_point_accounts(
- chat_id=chat_id, page=page, page_size=100
- )
- items.extend(batch)
- if not batch or len(items) >= total:
- break
- page += 1
- fieldnames = ["chat_id", "user_id", "display_name", "username", "first_name", "balance", "lifetime_earned", "lifetime_spent", "updated_at"]
- elif kind == "transactions":
- items = []
- page = 1
- while True:
- batch, total = await list_point_transactions(
- chat_id=chat_id, page=page, page_size=100
- )
- items.extend(batch)
- if not batch or len(items) >= total:
- break
- page += 1
- fieldnames = ["transaction_id", "chat_id", "user_id", "display_name", "username", "first_name", "delta", "balance_after", "source", "actor_id", "reason", "reference_id", "created_at"]
- else:
- raise ApiProblem(
- "invalid_export_kind",
- "导出类型必须是积分账户或积分流水。",
- )
- writer = csv.DictWriter(output, fieldnames=fieldnames, extrasaction="ignore")
- writer.writeheader()
- for item in items:
- writer.writerow(jsonable(item))
- return web.Response(
- text=output.getvalue(),
- content_type="text/csv",
- headers={"Content-Disposition": f'attachment; filename="points-{chat_id}-{kind}.csv"'},
- )
- async def business_assistant_status(self, _: web.Request) -> web.Response:
- return success(await runtime_overview())
- async def business_assistant_connections(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_business_connections(page=page, page_size=page_size)
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def business_assistant_settings(self, request: web.Request) -> web.Response:
- connection_id = str(request.query.get("connection_id") or "").strip()
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- if not await get_business_connection(connection_id):
- raise AssistantDataError("connection_not_found", "未找到 Business 连接。")
- return success(public_account_settings(await get_account_settings(connection_id)))
- async def business_assistant_settings_update(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- connection_id = str(body.pop("connection_id", "")).strip()
- body.pop("confirm", None)
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- set_audit(
- request,
- "business_assistant.settings.update",
- target_id=connection_id,
- summary="更新智能接待设置",
- )
- saved = await update_account_settings(connection_id, body)
- return success(public_account_settings(saved))
- async def business_assistant_feishu_webhook_test(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- connection_id = str(body.pop("connection_id", "")).strip()
- body.pop("confirm", None)
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- allowed = {"feishu_webhook_url", "feishu_webhook_signing_secret"}
- overrides = {key: value for key, value in body.items() if key in allowed}
- from wbb.modules.business_assistant import get_business_assistant_runtime
- runtime = get_business_assistant_runtime()
- if runtime is None:
- raise ApiProblem(
- "assistant_runtime_unavailable", "智能接待运行时尚未启动。", status=503
- )
- try:
- result = await runtime.test_feishu_webhook(connection_id, overrides)
- except FeishuWebhookError as exc:
- raise ApiProblem("feishu_webhook_error", str(exc), status=502) from exc
- set_audit(
- request,
- "business_assistant.feishu_webhook.test",
- target_id=connection_id,
- summary="测试飞书群机器人 Webhook",
- )
- return success(result)
- async def business_assistant_knowledge(self, request: web.Request) -> web.Response:
- connection_id = str(request.query.get("connection_id") or "").strip()
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- page, page_size = page_params(request)
- items, total = await list_knowledge_entries(
- connection_id,
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def business_assistant_knowledge_create(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- connection_id = str(body.pop("connection_id", "")).strip()
- body.pop("confirm", None)
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- entry = await create_knowledge_entry(connection_id, body)
- set_audit(
- request,
- "business_assistant.knowledge.create",
- target_id=entry["entry_id"],
- summary=str(entry["question"]),
- metadata={"connection_id": connection_id},
- )
- return success(entry, status=201)
- async def business_assistant_knowledge_update(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- body.pop("confirm", None)
- entry_id = request.match_info["entry_id"]
- set_audit(
- request,
- "business_assistant.knowledge.update",
- target_id=entry_id,
- summary="更新知识条目",
- )
- return success(await update_knowledge_entry(entry_id, body))
- async def business_assistant_knowledge_delete(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- entry_id = request.match_info["entry_id"]
- set_audit(
- request,
- "business_assistant.knowledge.delete",
- target_id=entry_id,
- summary="删除知识条目",
- )
- await delete_knowledge_entry(entry_id)
- return success({"deleted": True})
- async def business_assistant_knowledge_test(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- connection_id = str(body.get("connection_id") or "").strip()
- text = str(body.get("text") or "").strip()
- if not connection_id or not text:
- raise ApiProblem("invalid_parameter", "连接和测试问题不能为空。")
- from wbb.modules.business_assistant import get_business_assistant_runtime
- runtime = get_business_assistant_runtime()
- if runtime is None:
- raise ApiProblem(
- "assistant_runtime_unavailable", "智能接待运行时尚未启动。", status=503
- )
- try:
- result = await runtime.preview_answer(connection_id, text)
- except AssistantProviderError as exc:
- raise ApiProblem("assistant_provider_error", str(exc), status=502) from exc
- return success(result)
- async def business_assistant_knowledge_sources(
- self, request: web.Request
- ) -> web.Response:
- connection_id = str(request.query.get("connection_id") or "").strip()
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- page, page_size = page_params(request)
- items, total = await list_knowledge_sources(
- connection_id, page=page, page_size=page_size
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def business_assistant_knowledge_source_create(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- connection_id = str(body.pop("connection_id", "")).strip()
- body.pop("confirm", None)
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- from wbb.modules.business_assistant import validate_knowledge_source
- validated = await validate_knowledge_source(body)
- source = await create_knowledge_source(connection_id, validated)
- set_audit(
- request,
- "business_assistant.knowledge_source.create",
- target_id=source["source_id"],
- summary=str(source["title"]),
- metadata={"connection_id": connection_id, "source_type": source["source_type"]},
- )
- return success(source, status=201)
- async def business_assistant_knowledge_source_update(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- body.pop("confirm", None)
- source_id = request.match_info["source_id"]
- if body.get("enabled") is False:
- from wbb.modules.business_assistant import cancel_knowledge_source_sync
- await cancel_knowledge_source_sync(source_id)
- source = await update_knowledge_source(source_id, body)
- set_audit(
- request,
- "business_assistant.knowledge_source.update",
- target_id=source_id,
- summary=str(source.get("title") or "更新知识来源"),
- )
- return success(source)
- async def business_assistant_knowledge_source_delete(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- source_id = request.match_info["source_id"]
- source = await get_knowledge_source(source_id)
- if not source:
- raise AssistantDataError("source_not_found", "未找到知识来源。")
- set_audit(
- request,
- "business_assistant.knowledge_source.delete",
- target_id=source_id,
- summary=str(source.get("title") or "删除知识来源"),
- )
- from wbb.modules.business_assistant import cancel_knowledge_source_sync
- await cancel_knowledge_source_sync(source_id)
- await delete_knowledge_source(source_id)
- return success({"deleted": True})
- async def business_assistant_knowledge_source_action(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- source_id = request.match_info["source_id"]
- action = str(body.get("action") or "")
- if action in {"enable", "disable"}:
- if action == "disable":
- from wbb.modules.business_assistant import cancel_knowledge_source_sync
- await cancel_knowledge_source_sync(source_id)
- result = await update_knowledge_source(
- source_id, {"enabled": action == "enable"}
- )
- elif action == "sync":
- from wbb.modules.business_assistant import schedule_knowledge_source_sync
- result = await schedule_knowledge_source_sync(source_id)
- else:
- raise ApiProblem("invalid_action", "不支持的知识来源操作。")
- set_audit(
- request,
- f"business_assistant.knowledge_source.{action}",
- target_id=source_id,
- summary=f"知识来源操作:{action}",
- )
- return success(result)
- async def business_assistant_knowledge_candidates(
- self, request: web.Request
- ) -> web.Response:
- connection_id = str(request.query.get("connection_id") or "").strip()
- if not connection_id:
- raise ApiProblem("connection_required", "请先选择一个 Business 连接。")
- page, page_size = page_params(request)
- items, total = await list_knowledge_candidates(
- connection_id,
- status=str(request.query.get("status") or ""),
- source_id=str(request.query.get("source_id") or ""),
- query=str(request.query.get("query") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def business_assistant_knowledge_candidate_action(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- candidate_id = request.match_info["candidate_id"]
- action = str(body.get("action") or "")
- if not await get_knowledge_candidate(candidate_id):
- raise AssistantDataError("candidate_not_found", "未找到知识候选。")
- if action == "approve":
- result: Any = await publish_knowledge_candidate(candidate_id)
- elif action == "reject":
- result = await reject_knowledge_candidate(candidate_id)
- elif action == "regenerate":
- from wbb.modules.business_assistant import regenerate_knowledge_candidate
- try:
- result = await regenerate_knowledge_candidate(candidate_id)
- except AssistantProviderError as exc:
- raise ApiProblem("assistant_provider_error", str(exc), status=502) from exc
- else:
- raise ApiProblem("invalid_action", "不支持的知识候选操作。")
- set_audit(
- request,
- f"business_assistant.knowledge_candidate.{action}",
- target_id=candidate_id,
- summary=f"知识候选操作:{action}",
- )
- return success(result)
- async def business_assistant_conversations(
- self, request: web.Request
- ) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_conversations(
- connection_id=str(request.query.get("connection_id") or ""),
- status=str(request.query.get("status") or ""),
- page=page,
- page_size=page_size,
- )
- return success(
- {"items": items, "total": total, "page": page, "page_size": page_size}
- )
- async def business_assistant_conversation_detail(
- self, request: web.Request
- ) -> web.Response:
- return success(
- await conversation_detail(request.match_info["conversation_id"])
- )
- async def business_assistant_conversation_action(
- self, request: web.Request
- ) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- action = str(body.get("action") or "")
- conversation_id = request.match_info["conversation_id"]
- actions = {
- "pause": pause_conversation,
- "resume": resume_conversation,
- "close": close_conversation,
- "clear": clear_conversation,
- }
- handler = actions.get(action)
- if handler is None:
- raise ApiProblem("invalid_action", "不支持的会话操作。")
- set_audit(
- request,
- f"business_assistant.conversation.{action}",
- target_id=conversation_id,
- summary=f"智能接待会话操作:{action}",
- )
- return success(await handler(conversation_id))
- async def business_assistant_usage(self, request: web.Request) -> web.Response:
- return success(
- await usage_metrics(
- connection_id=str(request.query.get("connection_id") or "")
- )
- )
- async def business_assistant_model_test(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- from wbb.modules.business_assistant import get_business_assistant_runtime
- runtime = get_business_assistant_runtime()
- if runtime is None:
- raise ApiProblem(
- "assistant_runtime_unavailable", "智能接待运行时尚未启动。", status=503
- )
- try:
- result = await runtime.provider.test_connection()
- except AssistantProviderError as exc:
- raise ApiProblem("assistant_provider_error", str(exc), status=502) from exc
- set_audit(
- request,
- "business_assistant.model.test",
- summary="验证 OpenAI 兼容模型",
- )
- return success(result)
- async def giveaways(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- chat_raw = request.query.get("chat_id")
- items, total = await list_giveaways_page(
- chat_id=parse_int(chat_raw, "chat_id") if chat_raw else None,
- status=request.query.get("status") or None,
- query=request.query.get("query", ""),
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def giveaway_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = parse_int(body.get("chat_id"), "chat_id")
- await ensure_permission(chat_id, "can_change_info")
- prizes = body.get("prizes")
- if not isinstance(prizes, list):
- raise ApiProblem("invalid_prizes", "奖项必须是数组。")
- starts_at = parse_datetime(body.get("starts_at"), "starts_at")
- ends_at = parse_datetime(body.get("ends_at"), "ends_at")
- if ends_at <= datetime.now(UTC):
- raise ApiProblem("invalid_schedule", "开奖时间必须晚于当前时间。")
- if starts_at >= ends_at:
- raise ApiProblem("invalid_schedule", "报名开始时间必须早于开奖时间。")
- if str(body.get("eligibility_mode", "all")) not in {"all", "any"}:
- raise ApiProblem("invalid_eligibility_mode", "资格条件组合方式必须是全部或任一。")
- ticket_limit = parse_int(body.get("max_tickets_per_user", 10), "max_tickets_per_user", minimum=1)
- if ticket_limit > 100:
- raise ApiProblem("invalid_ticket_limit", "每人奖票上限不能超过 100。")
- targets = body.get("eligibility_targets") or []
- if not isinstance(targets, list):
- raise ApiProblem("invalid_targets", "资格目标必须是数组。")
- try:
- targets = await validate_targets(targets)
- except ValueError as exc:
- raise ApiProblem("invalid_targets", str(exc)) from exc
- set_audit(request, "giveaway.create", chat_id=chat_id, summary=str(body.get("title") or ""))
- giveaway = await create_and_publish_giveaway(
- chat_id=chat_id,
- creator_id=0,
- creator_name=request["admin"]["username"],
- title=str(body.get("title") or ""),
- description=str(body.get("description") or ""),
- prizes=prizes,
- starts_at=starts_at,
- ends_at=ends_at,
- minimum_points=parse_int(body.get("minimum_points", 0), "minimum_points", minimum=0),
- entry_cost=parse_int(body.get("entry_cost", 0), "entry_cost", minimum=0),
- participation_reward=parse_int(body.get("participation_reward", 0), "participation_reward", minimum=0),
- max_tickets_per_user=ticket_limit,
- eligibility_targets=targets,
- eligibility_mode=str(body.get("eligibility_mode", "all")),
- )
- return success(giveaway, status=201)
- async def giveaway_templates(self, request: web.Request) -> web.Response:
- raw = request.query.get("chat_id")
- return success({"items": await list_templates(int(raw) if raw else None)})
- async def giveaway_template_create(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = parse_int(body.get("chat_id"), "chat_id")
- await ensure_permission(chat_id, "can_change_info")
- try:
- body["eligibility_targets"] = await validate_targets(
- body.get("eligibility_targets") or []
- )
- template = await create_template(body, request["admin"]["username"])
- except (ValueError, KeyError, TypeError) as exc:
- raise ApiProblem("invalid_template", str(exc)) from exc
- set_audit(request, "giveaway.template.create", chat_id=chat_id, target_id=template["template_id"])
- return success(template, status=201)
- async def giveaway_template_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- template_id = request.match_info["template_id"]
- current = await templatesdb.find_one({"template_id": template_id, "bot_id": wbb.BOT_PROFILE_ID})
- if not current:
- raise ApiProblem("template_not_found", "未找到循环抽奖模板。", status=404)
- await ensure_permission(int(current["chat_id"]), "can_change_info")
- if int(body.get("chat_id", current["chat_id"])) != int(current["chat_id"]):
- raise ApiProblem("template_chat_fixed", "循环抽奖模板不能更换群组。")
- merged = {**current, **body}
- try:
- if merged.get("active", True):
- merged["eligibility_targets"] = await validate_targets(
- merged.get("eligibility_targets") or []
- )
- updated = await update_template(template_id, merged)
- except (ValueError, KeyError, TypeError) as exc:
- raise ApiProblem("invalid_template", str(exc)) from exc
- set_audit(request, "giveaway.template.update", chat_id=int(current["chat_id"]), target_id=template_id)
- return success(updated)
- async def external_risk_sources(self, request: web.Request) -> web.Response:
- snapshots = await list_external_risk_snapshots()
- return success({
- "sources": EXTERNAL_RISK_SOURCES, "snapshots": snapshots,
- "active": await list_active_external_risk_sources(),
- })
- async def external_risk_stage(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- source_id = request.match_info["source_id"]
- try:
- snapshot = await stage_external_risk_snapshot(source_id)
- except ValueError as exc:
- raise ApiProblem("external_risk_invalid_source", str(exc)) from exc
- except (ClientError, TimeoutError) as exc:
- raise ApiProblem("external_risk_source_unavailable", "规则源暂不可访问,请稍后重试。", status=502) from exc
- set_audit(request, "external_risk.stage", target_id=source_id)
- return success(snapshot, status=201)
- async def external_risk_publish(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- try:
- snapshot = await publish_external_risk_snapshot(request.match_info["snapshot_id"])
- except ValueError as exc:
- raise ApiProblem("external_risk_snapshot_invalid", str(exc)) from exc
- set_audit(request, "external_risk.publish", target_id=snapshot["snapshot_id"])
- return success(snapshot)
- async def external_risk_entries(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- try:
- items, total = await list_external_risk_snapshot_entries(
- request.match_info["snapshot_id"], page=page, page_size=page_size,
- query=request.query.get("query", ""),
- )
- except ValueError as exc:
- raise ApiProblem("external_risk_snapshot_not_found", str(exc), status=404) from exc
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def external_risk_chat_policy(self, request: web.Request) -> web.Response:
- return success(await get_external_risk_policy(chat_id_param(request)))
- async def external_risk_chat_policy_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = chat_id_param(request)
- await ensure_permission(chat_id, "can_change_info")
- if body.get("mode") == "delete":
- await ensure_permission(chat_id, "can_delete_messages")
- try:
- result = await set_external_risk_policy(chat_id, str(body.get("mode")))
- except ValueError as exc:
- raise ApiProblem("external_risk_policy_invalid", str(exc)) from exc
- set_audit(request, "external_risk.policy", chat_id=chat_id, summary=result["mode"])
- return success(result)
- async def chat_interaction_settings(self, request: web.Request) -> web.Response:
- return success(await get_interaction_settings(chat_id_param(request)))
- async def chat_interaction_settings_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- chat_id = chat_id_param(request)
- await ensure_permission(chat_id, "can_change_info")
- try:
- saved = await save_interaction_settings(chat_id, body)
- except (TypeError, ValueError) as exc:
- raise ApiProblem("invalid_interaction_settings", str(exc)) from exc
- set_audit(request, "chat.interaction_settings.update", chat_id=chat_id)
- return success(saved)
- async def giveaway_detail(self, request: web.Request) -> web.Response:
- giveaway = await get_giveaway(request.match_info["giveaway_id"])
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- return success(giveaway)
- async def giveaway_update(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- starts_at = parse_datetime(body.get("starts_at"), "starts_at")
- ends_at = parse_datetime(body.get("ends_at"), "ends_at")
- expected_updated_at = parse_datetime(
- body.get("expected_updated_at"), "expected_updated_at"
- )
- set_audit(
- request,
- "giveaway.update",
- chat_id=int(giveaway["chat_id"]),
- target_id=giveaway_id,
- summary=str(body.get("title") or ""),
- metadata={
- "starts_at": starts_at,
- "ends_at": ends_at,
- },
- )
- return success(
- await update_and_refresh_giveaway(
- giveaway_id,
- title=str(body.get("title") or ""),
- description=str(body.get("description") or ""),
- starts_at=starts_at,
- ends_at=ends_at,
- expected_updated_at=expected_updated_at,
- )
- )
- async def giveaway_finish(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.finish", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Manual draw")
- ok, message, current = await finish_and_publish_giveaway(giveaway_id)
- if not ok:
- raise ApiProblem("giveaway_not_running", message, status=409)
- return success(current)
- async def giveaway_cancel(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.cancel", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary="Cancel and refund")
- ok, message, current = await cancel_and_refund_giveaway(
- giveaway_id, chat_id=int(giveaway["chat_id"])
- )
- if not ok:
- raise ApiProblem("giveaway_not_running", message, status=409)
- return success(current)
- async def giveaway_reroll(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.reroll", chat_id=int(giveaway["chat_id"]), target_id=giveaway_id, summary=str(body.get("tier_name") or "All tiers"))
- winners, reroll_id = await reroll_giveaway(
- giveaway_id=giveaway_id,
- moderator_id=0,
- tier_name=body.get("tier_name") or None,
- reroll_id=body.get("request_id") or None,
- )
- return success({"winners": winners, "reroll_id": reroll_id})
- async def giveaway_participants(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_participants_page(
- giveaway_id=request.match_info["giveaway_id"],
- active_only=request.query.get("active_only", "false").lower() == "true",
- page=page,
- page_size=page_size,
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def giveaway_orders(self, request: web.Request) -> web.Response:
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- page, page_size = page_params(request)
- query = {"bot_id": wbb.BOT_PROFILE_ID, "giveaway_id": giveaway_id}
- total = await ticket_ordersdb.count_documents(query)
- cursor = ticket_ordersdb.find(query).sort("created_at", -1).skip((page - 1) * page_size).limit(page_size)
- return success({"items": [item async for item in cursor], "total": total, "page": page, "page_size": page_size})
- async def giveaway_participant_remove(self, request: web.Request) -> web.Response:
- body = await json_body(request)
- require_confirmation(body)
- giveaway_id = request.match_info["giveaway_id"]
- user_id = parse_int(request.match_info["user_id"], "user_id")
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- await ensure_permission(int(giveaway["chat_id"]), "can_change_info")
- set_audit(request, "giveaway.participant.remove", chat_id=int(giveaway["chat_id"]), target_id=user_id, summary=str(body.get("reason") or ""))
- removed = await remove_and_optionally_refund_participant(
- giveaway_id=giveaway_id,
- user_id=user_id,
- moderator_id=0,
- reason=str(body.get("reason") or "Removed in admin panel"),
- refund=bool(body.get("refund", True)),
- )
- if not removed:
- raise ApiProblem("participant_not_found", "参与者不存在或已被移除。", status=404)
- return success({"removed": True, "refunded": bool(body.get("refund", True))})
- async def giveaway_bans(self, request: web.Request) -> web.Response:
- page, page_size = page_params(request)
- items, total = await list_giveaway_bans(
- chat_id=chat_id_param(request), page=page, page_size=page_size
- )
- return success({"items": items, "total": total, "page": page, "page_size": page_size})
- async def giveaway_ban_add(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- body = await json_body(request)
- require_confirmation(body)
- user_id = parse_int(body.get("user_id"), "user_id")
- await ensure_permission(chat_id, "can_change_info")
- set_audit(request, "giveaway.ban.add", chat_id=chat_id, target_id=user_id, summary=str(body.get("reason") or ""))
- await add_giveaway_ban(chat_id=chat_id, user_id=user_id, moderator_id=0, reason=str(body.get("reason") or ""))
- return success({"created": True}, status=201)
- async def giveaway_ban_remove(self, request: web.Request) -> web.Response:
- chat_id = chat_id_param(request)
- user_id = parse_int(request.match_info["user_id"], "user_id")
- body = await json_body(request)
- require_confirmation(body)
- await ensure_permission(chat_id, "can_change_info")
- set_audit(request, "giveaway.ban.remove", chat_id=chat_id, target_id=user_id)
- return success({"removed": await remove_giveaway_ban(chat_id, user_id)})
- async def giveaway_export(self, request: web.Request) -> web.Response:
- giveaway_id = request.match_info["giveaway_id"]
- giveaway = await get_giveaway(giveaway_id)
- if not giveaway:
- raise ApiProblem("giveaway_not_found", "未找到该抽奖。", status=404)
- participants: list[dict[str, Any]] = []
- page = 1
- while True:
- batch, total = await list_participants_page(
- giveaway_id=giveaway_id, page=page, page_size=100
- )
- participants.extend(batch)
- if not batch or len(participants) >= total:
- break
- page += 1
- output = io.StringIO()
- fields = ["giveaway_id", "chat_id", "user_id", "display_name", "username", "first_name", "active", "entry_cost", "joined_at", "removed_at", "refunded_at"]
- writer = csv.DictWriter(output, fieldnames=fields, extrasaction="ignore")
- writer.writeheader()
- for participant in participants:
- writer.writerow(jsonable(participant))
- return web.Response(
- text=output.getvalue(),
- content_type="text/csv",
- headers={"Content-Disposition": f'attachment; filename="giveaway-{giveaway_id}.csv"'},
- )
- def build_admin_application() -> web.Application:
- max_upload_mb = int(getattr(wbb, "ADMIN_WEB_UPLOAD_MAX_MB", 20))
- application = web.Application(
- middlewares=[
- api_error_middleware,
- technician_authentication_middleware,
- authentication_middleware,
- bot_role_middleware,
- bot_proxy_middleware,
- ],
- client_max_size=(max_upload_mb + 1) * 1024 * 1024,
- )
- admin_api = AdminApi()
- application["admin_api"] = admin_api
- admin_api.register(application)
- return application
|