| 1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003100410051006 |
- from __future__ import annotations
- import asyncio
- import math
- import re
- from datetime import UTC, datetime, timedelta
- from typing import Any
- from uuid import uuid4
- from pymongo import ASCENDING, DESCENDING, GEOSPHERE, ReturnDocument
- from pymongo.errors import DuplicateKeyError, OperationFailure
- from wbb import BOT_PROFILE_ID, control_db
- profilesdb = control_db.directory_profiles
- locationsdb = control_db.directory_locations
- membershipsdb = control_db.directory_memberships
- eventsdb = control_db.directory_events
- settingsdb = control_db.directory_settings
- APPLICATION_NONE = "none"
- APPLICATION_PENDING = "pending"
- APPLICATION_APPROVED = "approved"
- APPLICATION_REJECTED = "rejected"
- APPLICATION_REVOKED = "revoked"
- APPLICATION_STATUSES = {
- APPLICATION_NONE,
- APPLICATION_PENDING,
- APPLICATION_APPROVED,
- APPLICATION_REJECTED,
- APPLICATION_REVOKED,
- }
- PRESENCE_ONLINE = "online"
- PRESENCE_OFFLINE = "offline"
- DEFAULT_DIRECTORY_SETTINGS: dict[str, Any] = {
- "online_duration_hours": 24,
- "nearby_radius_km": 20,
- "nearby_radius_options_km": [5, 10, 20, 50],
- }
- _index_lock = asyncio.Lock()
- _indexes_ready = False
- class DirectoryDataError(ValueError):
- def __init__(self, code: str, message: str):
- super().__init__(message)
- self.code = code
- def utc_now() -> datetime:
- return datetime.now(UTC)
- def aware_utc(value: datetime | None) -> datetime | None:
- if value is None:
- return None
- if value.tzinfo is None:
- return value.replace(tzinfo=UTC)
- return value.astimezone(UTC)
- def display_name(
- first_name: str | None,
- last_name: str | None,
- username: str | None = None,
- user_id: int | None = None,
- ) -> str:
- name = " ".join(
- value.strip()
- for value in (first_name, last_name)
- if value and value.strip()
- )
- if name:
- return name
- if username:
- return f"@{username}"
- return f"用户 {user_id}" if user_id is not None else "未知用户"
- async def ensure_directory_indexes() -> None:
- global _indexes_ready
- if _indexes_ready:
- return
- async with _index_lock:
- if _indexes_ready:
- return
- await profilesdb.create_index([("user_id", ASCENDING)], unique=True)
- await profilesdb.create_index(
- [
- ("application_status", ASCENDING),
- ("listed", ASCENDING),
- ("updated_at", DESCENDING),
- ]
- )
- await profilesdb.create_index(
- [("presence_status", ASCENDING), ("online_until", ASCENDING)]
- )
- await locationsdb.create_index([("user_id", ASCENDING)], unique=True)
- await locationsdb.create_index([("point", GEOSPHERE)])
- await locationsdb.create_index([("updated_at", DESCENDING)])
- await membershipsdb.create_index(
- [
- ("bot_id", ASCENDING),
- ("chat_id", ASCENDING),
- ("user_id", ASCENDING),
- ],
- unique=True,
- )
- await membershipsdb.create_index(
- [("user_id", ASCENDING), ("active", ASCENDING), ("verified_at", DESCENDING)]
- )
- await eventsdb.create_index([("event_id", ASCENDING)], unique=True)
- await eventsdb.create_index(
- [("idempotency_key", ASCENDING)],
- unique=True,
- sparse=True,
- )
- await eventsdb.create_index([("created_at", DESCENDING)])
- await eventsdb.create_index([("user_id", ASCENDING), ("created_at", DESCENDING)])
- await settingsdb.create_index([("settings_id", ASCENDING)], unique=True)
- _indexes_ready = True
- async def get_directory_settings() -> dict[str, Any]:
- await ensure_directory_indexes()
- stored = await settingsdb.find_one({"settings_id": "global"}) or {}
- return {
- **DEFAULT_DIRECTORY_SETTINGS,
- **{
- key: stored[key]
- for key in DEFAULT_DIRECTORY_SETTINGS
- if key in stored
- },
- }
- async def set_directory_settings(values: dict[str, Any]) -> dict[str, Any]:
- await ensure_directory_indexes()
- current = await get_directory_settings()
- target = {**current, **values}
- try:
- radius_options = sorted(
- {
- int(value)
- for value in target.get("nearby_radius_options_km", [])
- if int(value) > 0
- }
- )
- default_radius = int(target.get("nearby_radius_km") or 0)
- except (TypeError, ValueError) as exc:
- raise DirectoryDataError(
- "invalid_radius_options", "附近范围必须是正整数。"
- ) from exc
- if not radius_options:
- raise DirectoryDataError("invalid_radius_options", "附近范围选项不能为空。")
- if default_radius not in radius_options:
- raise DirectoryDataError(
- "invalid_default_radius", "默认附近范围必须属于可选范围。"
- )
- target.update(
- {
- "online_duration_hours": 24,
- "nearby_radius_km": default_radius,
- "nearby_radius_options_km": radius_options,
- "updated_at": utc_now(),
- }
- )
- await settingsdb.update_one(
- {"settings_id": "global"},
- {"$set": target, "$setOnInsert": {"created_at": utc_now()}},
- upsert=True,
- )
- return await get_directory_settings()
- async def record_directory_event(
- *,
- event_type: str,
- user_id: int,
- actor_id: int | str | None,
- actor_name: str = "",
- source: str,
- reason: str = "",
- metadata: dict[str, Any] | None = None,
- idempotency_key: str | None = None,
- bot_id: str | None = None,
- chat_id: int | None = None,
- ) -> dict[str, Any] | None:
- await ensure_directory_indexes()
- document = {
- "event_id": uuid4().hex,
- "event_type": event_type,
- "user_id": int(user_id),
- "actor_id": actor_id,
- "actor_name": actor_name,
- "source": source,
- "reason": reason,
- "metadata": metadata or {},
- "bot_id": bot_id or str(BOT_PROFILE_ID),
- "chat_id": int(chat_id) if chat_id is not None else None,
- "created_at": utc_now(),
- }
- if idempotency_key:
- document["idempotency_key"] = idempotency_key
- try:
- await eventsdb.insert_one(document)
- except DuplicateKeyError:
- return None
- return document
- async def upsert_directory_identity(
- *,
- user_id: int,
- username: str | None,
- first_name: str | None,
- last_name: str | None,
- source: str = "telegram",
- ) -> dict[str, Any]:
- await ensure_directory_indexes()
- normalized_username = str(username or "").lstrip("@").strip() or None
- now = utc_now()
- previous = await profilesdb.find_one({"user_id": int(user_id)})
- identity = {
- "username": normalized_username,
- "first_name": str(first_name or "").strip() or None,
- "last_name": str(last_name or "").strip() or None,
- "display_name": display_name(
- first_name,
- last_name,
- normalized_username,
- int(user_id),
- ),
- }
- updates: dict[str, Any] = {**identity, "updated_at": now}
- username_removed = bool(previous and previous.get("username") and not normalized_username)
- if username_removed and previous.get("application_status") == APPLICATION_APPROVED:
- updates.update(
- {
- "listed": False,
- "presence_status": PRESENCE_OFFLINE,
- "online_until": None,
- }
- )
- await profilesdb.update_one(
- {"user_id": int(user_id)},
- {
- "$set": updates,
- "$setOnInsert": {
- "application_status": APPLICATION_NONE,
- "listed": False,
- "presence_status": PRESENCE_OFFLINE,
- "created_at": now,
- },
- },
- upsert=True,
- )
- if username_removed:
- await record_directory_event(
- event_type="teacher_username_removed",
- user_id=int(user_id),
- actor_id=int(user_id),
- actor_name=identity["display_name"],
- source=source,
- reason="技师移除了 Telegram 用户名,系统自动下榜并下线。",
- )
- return await get_directory_profile(int(user_id)) or {}
- async def get_directory_profile(user_id: int) -> dict[str, Any] | None:
- await ensure_directory_indexes()
- return await profilesdb.find_one({"user_id": int(user_id)})
- async def get_directory_location(user_id: int) -> dict[str, Any] | None:
- await ensure_directory_indexes()
- return await locationsdb.find_one({"user_id": int(user_id)})
- def validate_coordinates(longitude: float, latitude: float) -> tuple[float, float]:
- try:
- lon = float(longitude)
- lat = float(latitude)
- except (TypeError, ValueError) as exc:
- raise DirectoryDataError("invalid_coordinates", "经纬度必须是数字。") from exc
- if not math.isfinite(lon) or not -180 <= lon <= 180:
- raise DirectoryDataError("invalid_longitude", "经度必须在 -180 到 180 之间。")
- if not math.isfinite(lat) or not -90 <= lat <= 90:
- raise DirectoryDataError("invalid_latitude", "纬度必须在 -90 到 90 之间。")
- return lon, lat
- async def save_directory_location(
- *,
- user_id: int,
- longitude: float,
- latitude: float,
- source: str,
- actor_id: int | str,
- actor_name: str,
- bot_id: str | None = None,
- ) -> dict[str, Any]:
- await ensure_directory_indexes()
- lon, lat = validate_coordinates(longitude, latitude)
- now = utc_now()
- await locationsdb.update_one(
- {"user_id": int(user_id)},
- {
- "$set": {
- "point": {"type": "Point", "coordinates": [lon, lat]},
- "longitude": lon,
- "latitude": lat,
- "source": source,
- "source_bot_id": bot_id or str(BOT_PROFILE_ID),
- "updated_at": now,
- },
- "$unset": {
- "coordinate_key": "",
- "geocode_consent": "",
- "geocode_status": "",
- "region": "",
- "geocoded_at": "",
- "geocode_error": "",
- },
- "$setOnInsert": {"created_at": now},
- },
- upsert=True,
- )
- await profilesdb.update_one(
- {"user_id": int(user_id)},
- {
- "$set": {"location_updated_at": now, "updated_at": now},
- "$setOnInsert": {
- "application_status": APPLICATION_NONE,
- "listed": False,
- "presence_status": PRESENCE_OFFLINE,
- "created_at": now,
- },
- },
- upsert=True,
- )
- await record_directory_event(
- event_type="location_updated",
- user_id=int(user_id),
- actor_id=actor_id,
- actor_name=actor_name,
- source=source,
- reason="更新位置",
- metadata={"location_source": source},
- bot_id=bot_id,
- )
- return await get_directory_location(int(user_id)) or {}
- async def clear_directory_location(
- *,
- user_id: int,
- actor_id: int | str,
- actor_name: str,
- source: str,
- reason: str,
- ) -> bool:
- await ensure_directory_indexes()
- result = await locationsdb.delete_one({"user_id": int(user_id)})
- await profilesdb.update_one(
- {"user_id": int(user_id)},
- {
- "$set": {
- "listed": False,
- "presence_status": PRESENCE_OFFLINE,
- "online_until": None,
- "location_updated_at": None,
- "updated_at": utc_now(),
- }
- },
- )
- if result.deleted_count:
- await record_directory_event(
- event_type="location_cleared",
- user_id=int(user_id),
- actor_id=actor_id,
- actor_name=actor_name,
- source=source,
- reason=reason,
- )
- return bool(result.deleted_count)
- async def submit_teacher_application(
- *,
- user_id: int,
- source: str,
- bot_id: str,
- chat_id: int,
- chat_title: str = "",
- ) -> tuple[dict[str, Any], bool]:
- await ensure_directory_indexes()
- profile = await get_directory_profile(int(user_id))
- location = await get_directory_location(int(user_id))
- if not profile or not profile.get("username"):
- raise DirectoryDataError(
- "username_required", "申请技师前必须先设置 Telegram 用户名。"
- )
- if not location:
- raise DirectoryDataError("location_required", "申请技师前必须先更新位置。")
- if profile.get("application_status") == APPLICATION_APPROVED:
- raise DirectoryDataError("already_teacher", "你已经是技师,无需重复申请。")
- if profile.get("application_status") == APPLICATION_REVOKED:
- raise DirectoryDataError(
- "teacher_revoked", "技师资格已被撤销,只能由管理员重新批准。"
- )
- if profile.get("application_status") == APPLICATION_PENDING:
- return profile, False
- now = utc_now()
- await profilesdb.update_one(
- {"user_id": int(user_id)},
- {
- "$set": {
- "application_status": APPLICATION_PENDING,
- "application_reason": "",
- "application_source_bot_id": str(bot_id),
- "application_source_chat_id": int(chat_id),
- "application_source_chat_title": str(chat_title or "").strip(),
- "applied_at": now,
- "decided_at": None,
- "decided_by": None,
- "updated_at": now,
- }
- },
- )
- await record_directory_event(
- event_type="teacher_application_submitted",
- user_id=int(user_id),
- actor_id=int(user_id),
- actor_name=str(profile.get("display_name") or user_id),
- source=source,
- reason="申请成为技师",
- bot_id=bot_id,
- chat_id=int(chat_id),
- )
- return await get_directory_profile(int(user_id)) or {}, True
- async def decide_teacher_application(
- *,
- user_id: int,
- action: str,
- actor_id: int | str,
- actor_name: str,
- source: str,
- reason: str = "",
- authorization_chat_id: int | None = None,
- ) -> dict[str, Any]:
- await ensure_directory_indexes()
- profile = await get_directory_profile(int(user_id))
- if not profile:
- raise DirectoryDataError("profile_not_found", "未找到该成员资料。")
- action_status = {
- "approve": APPLICATION_APPROVED,
- "reject": APPLICATION_REJECTED,
- "revoke": APPLICATION_REVOKED,
- }
- if action not in action_status:
- raise DirectoryDataError("invalid_decision", "不支持的审批操作。")
- if action in {"reject", "revoke"} and not reason.strip():
- raise DirectoryDataError("reason_required", "拒绝或撤销时必须填写原因。")
- if action == "approve":
- if not profile.get("username"):
- raise DirectoryDataError(
- "username_required", "该成员没有 Telegram 用户名,无法批准。"
- )
- if not await get_directory_location(int(user_id)):
- raise DirectoryDataError("location_required", "该成员尚未更新位置。")
- now = utc_now()
- updates: dict[str, Any] = {
- "application_status": action_status[action],
- "application_reason": reason.strip(),
- "decided_at": now,
- "decided_by": actor_id,
- "decided_by_name": actor_name,
- "authorization_chat_id": authorization_chat_id,
- "updated_at": now,
- }
- if action != "approve":
- updates.update(
- {
- "listed": False,
- "presence_status": PRESENCE_OFFLINE,
- "online_until": None,
- }
- )
- await profilesdb.update_one({"user_id": int(user_id)}, {"$set": updates})
- await record_directory_event(
- event_type=f"teacher_application_{action}",
- user_id=int(user_id),
- actor_id=actor_id,
- actor_name=actor_name,
- source=source,
- reason=reason,
- metadata={"authorization_chat_id": authorization_chat_id},
- chat_id=authorization_chat_id,
- )
- return await get_directory_profile(int(user_id)) or {}
- async def set_teacher_state(
- *,
- user_id: int,
- action: str,
- actor_id: int | str,
- actor_name: str,
- source: str,
- reason: str = "",
- idempotency_key: str | None = None,
- chat_id: int | None = None,
- ) -> tuple[dict[str, Any], bool]:
- await ensure_directory_indexes()
- profile = await get_directory_profile(int(user_id))
- if not profile or profile.get("application_status") != APPLICATION_APPROVED:
- raise DirectoryDataError("teacher_required", "该成员尚未取得技师资格。")
- location = await get_directory_location(int(user_id))
- if action in {"list", "online"}:
- if not profile.get("username"):
- raise DirectoryDataError("username_required", "请先设置 Telegram 用户名。")
- if not location:
- raise DirectoryDataError("location_required", "请先更新位置。")
- if action == "online" and not profile.get("listed"):
- raise DirectoryDataError("listing_required", "请先上榜,再切换为上线状态。")
- if action not in {"list", "unlist", "online", "offline"}:
- raise DirectoryDataError("invalid_teacher_action", "不支持的技师状态操作。")
- event = await record_directory_event(
- event_type=f"teacher_{action}",
- user_id=int(user_id),
- actor_id=actor_id,
- actor_name=actor_name,
- source=source,
- reason=reason or f"技师{action}",
- idempotency_key=idempotency_key,
- chat_id=chat_id,
- )
- applied = event is not None
- now = utc_now()
- updates: dict[str, Any] = {"updated_at": now}
- if action == "list":
- updates["listed"] = True
- updates["listed_at"] = now
- elif action == "unlist":
- updates.update(
- {
- "listed": False,
- "presence_status": PRESENCE_OFFLINE,
- "online_until": None,
- }
- )
- elif action == "online":
- updates.update(
- {
- "presence_status": PRESENCE_ONLINE,
- "online_until": now + timedelta(hours=24),
- "last_online_at": now,
- }
- )
- else:
- updates.update(
- {"presence_status": PRESENCE_OFFLINE, "online_until": None}
- )
- await profilesdb.update_one({"user_id": int(user_id)}, {"$set": updates})
- return await get_directory_profile(int(user_id)) or {}, applied
- async def expire_teacher_presence() -> int:
- await ensure_directory_indexes()
- now = utc_now()
- expired = 0
- while True:
- item = await profilesdb.find_one_and_update(
- {
- "presence_status": PRESENCE_ONLINE,
- "online_until": {"$lte": now},
- },
- {
- "$set": {
- "presence_status": PRESENCE_OFFLINE,
- "online_until": None,
- "updated_at": now,
- }
- },
- sort=[("online_until", ASCENDING)],
- return_document=ReturnDocument.BEFORE,
- )
- if not item:
- break
- expired += 1
- previous_until = aware_utc(item.get("online_until"))
- await record_directory_event(
- event_type="teacher_auto_offline",
- user_id=int(item["user_id"]),
- actor_id="system",
- actor_name="系统",
- source="scheduler",
- reason="上线状态已超过 24 小时。",
- idempotency_key=(
- f"teacher-auto-offline:{item['user_id']}:"
- f"{previous_until.isoformat() if previous_until else 'unknown'}"
- ),
- )
- return expired
- def teacher_is_online(profile: dict[str, Any], *, now: datetime | None = None) -> bool:
- current = now or utc_now()
- until = aware_utc(profile.get("online_until"))
- return bool(
- profile.get("presence_status") == PRESENCE_ONLINE
- and until
- and until > current
- )
- def _haversine_meters(
- longitude_a: float,
- latitude_a: float,
- longitude_b: float,
- latitude_b: float,
- ) -> float:
- radius = 6_371_008.8
- lon_a, lat_a, lon_b, lat_b = map(
- math.radians,
- (longitude_a, latitude_a, longitude_b, latitude_b),
- )
- delta_lon = lon_b - lon_a
- delta_lat = lat_b - lat_a
- value = (
- math.sin(delta_lat / 2) ** 2
- + math.cos(lat_a) * math.cos(lat_b) * math.sin(delta_lon / 2) ** 2
- )
- return 2 * radius * math.asin(math.sqrt(value))
- async def _list_teachers_fallback(
- *,
- longitude: float,
- latitude: float,
- max_distance_meters: float | None,
- query: str,
- page: int,
- page_size: int,
- ) -> tuple[list[dict[str, Any]], int]:
- profile_filter: dict[str, Any] = {
- "application_status": APPLICATION_APPROVED,
- "listed": True,
- "username": {"$nin": [None, ""]},
- }
- if query:
- pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
- profile_filter["$or"] = [
- {"username": pattern},
- {"display_name": pattern},
- {"user_id": int(query)} if query.isdigit() else {"user_id": -1},
- ]
- profiles = {
- int(item["user_id"]): item
- async for item in profilesdb.find(profile_filter)
- }
- values: list[dict[str, Any]] = []
- if profiles:
- async for location in locationsdb.find(
- {"user_id": {"$in": list(profiles)}}
- ):
- point = location.get("point", {}).get("coordinates") or []
- if len(point) != 2:
- continue
- distance = _haversine_meters(
- longitude,
- latitude,
- float(point[0]),
- float(point[1]),
- )
- if max_distance_meters is not None and distance > max_distance_meters:
- continue
- values.append(
- {
- **profiles[int(location["user_id"])],
- "distance_meters": distance,
- "location_updated_at": location.get("updated_at"),
- }
- )
- values.sort(key=lambda item: (item["distance_meters"], int(item["user_id"])))
- total = len(values)
- start = (page - 1) * page_size
- return values[start : start + page_size], total
- async def list_directory_teachers(
- *,
- longitude: float,
- latitude: float,
- max_distance_meters: float | None = None,
- query: str = "",
- page: int = 1,
- page_size: int = 10,
- ) -> tuple[list[dict[str, Any]], int]:
- await ensure_directory_indexes()
- lon, lat = validate_coordinates(longitude, latitude)
- profile_match: dict[str, Any] = {
- "profile.application_status": APPLICATION_APPROVED,
- "profile.listed": True,
- "profile.username": {"$nin": [None, ""]},
- }
- if query:
- pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
- profile_match["$or"] = [
- {"profile.username": pattern},
- {"profile.display_name": pattern},
- {"profile.user_id": int(query)} if query.isdigit() else {"profile.user_id": -1},
- ]
- geo_near: dict[str, Any] = {
- "near": {"type": "Point", "coordinates": [lon, lat]},
- "distanceField": "distance_meters",
- "spherical": True,
- "key": "point",
- }
- if max_distance_meters is not None:
- geo_near["maxDistance"] = float(max_distance_meters)
- pipeline = [
- {"$geoNear": geo_near},
- {
- "$lookup": {
- "from": profilesdb.name,
- "localField": "user_id",
- "foreignField": "user_id",
- "as": "profile",
- }
- },
- {"$unwind": "$profile"},
- {"$match": profile_match},
- {"$sort": {"distance_meters": 1, "user_id": 1}},
- {
- "$facet": {
- "items": [
- {"$skip": (max(1, page) - 1) * max(1, page_size)},
- {"$limit": max(1, page_size)},
- ],
- "count": [{"$count": "total"}],
- }
- },
- ]
- try:
- results = await locationsdb.aggregate(pipeline).to_list(length=1)
- result = results[0] if results else {"items": [], "count": []}
- items = [
- {
- **item["profile"],
- "distance_meters": item["distance_meters"],
- "location_updated_at": item.get("updated_at"),
- }
- for item in result.get("items", [])
- ]
- count = result.get("count") or []
- return items, int(count[0]["total"]) if count else 0
- except (NotImplementedError, OperationFailure):
- return await _list_teachers_fallback(
- longitude=lon,
- latitude=lat,
- max_distance_meters=max_distance_meters,
- query=query,
- page=page,
- page_size=page_size,
- )
- async def list_teacher_applications(
- *,
- status: str = "",
- query: str = "",
- source_bot_id: str = "",
- source_chat_ids: list[int] | None = None,
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- await ensure_directory_indexes()
- filters: dict[str, Any] = {}
- if status:
- if status not in APPLICATION_STATUSES:
- raise DirectoryDataError("invalid_application_status", "申请状态无效。")
- filters["application_status"] = status
- else:
- filters["application_status"] = {"$ne": APPLICATION_NONE}
- if source_bot_id:
- filters["application_source_bot_id"] = str(source_bot_id)
- if source_chat_ids is not None:
- filters["application_source_chat_id"] = {
- "$in": [int(value) for value in source_chat_ids]
- }
- if query:
- pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
- filters["$or"] = [
- {"username": pattern},
- {"display_name": pattern},
- {"user_id": int(query)} if query.isdigit() else {"user_id": -1},
- ]
- total = await profilesdb.count_documents(filters)
- items = await (
- profilesdb.find(filters)
- .sort([("applied_at", DESCENDING), ("updated_at", DESCENDING)])
- .skip((page - 1) * page_size)
- .limit(page_size)
- .to_list(length=page_size)
- )
- location_map = {
- int(item["user_id"]): item
- async for item in locationsdb.find(
- {"user_id": {"$in": [int(item["user_id"]) for item in items]}}
- )
- }
- for item in items:
- location = location_map.get(int(item["user_id"]))
- item["location"] = (
- {
- "source": location.get("source"),
- "updated_at": location.get("updated_at"),
- }
- if location
- else None
- )
- item["online"] = teacher_is_online(item)
- return items, total
- async def list_directory_profiles(
- *,
- query: str = "",
- application_status: str = "",
- listed: bool | None = None,
- source_bot_id: str = "",
- source_chat_ids: list[int] | None = None,
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- filters: dict[str, Any] = {}
- if application_status:
- filters["application_status"] = application_status
- if listed is not None:
- filters["listed"] = listed
- if source_bot_id:
- filters["application_source_bot_id"] = str(source_bot_id)
- if source_chat_ids is not None:
- filters["application_source_chat_id"] = {
- "$in": [int(value) for value in source_chat_ids]
- }
- if query:
- pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
- filters["$or"] = [
- {"username": pattern},
- {"display_name": pattern},
- {"user_id": int(query)} if query.isdigit() else {"user_id": -1},
- ]
- total = await profilesdb.count_documents(filters)
- items = await (
- profilesdb.find(filters)
- .sort([("updated_at", DESCENDING)])
- .skip((page - 1) * page_size)
- .limit(page_size)
- .to_list(length=page_size)
- )
- locations = {
- int(item["user_id"]): item
- async for item in locationsdb.find(
- {"user_id": {"$in": [int(item["user_id"]) for item in items]}}
- )
- }
- for item in items:
- location = locations.get(int(item["user_id"]))
- item["location"] = (
- {
- "source": location.get("source"),
- "updated_at": location.get("updated_at"),
- }
- if location
- else None
- )
- item["online"] = teacher_is_online(item)
- return items, total
- async def list_directory_locations(
- *,
- query: str = "",
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- await ensure_directory_indexes()
- filters: dict[str, Any] = {}
- if query:
- pattern = re.compile(re.escape(query.lstrip("@")), re.IGNORECASE)
- profile_filters = {
- "$or": [
- {"username": pattern},
- {"display_name": pattern},
- {"user_id": int(query)} if query.isdigit() else {"user_id": -1},
- ]
- }
- user_ids = [
- int(item["user_id"])
- async for item in profilesdb.find(profile_filters, {"user_id": 1})
- ]
- filters["user_id"] = {"$in": user_ids}
- total = await locationsdb.count_documents(filters)
- items = await (
- locationsdb.find(filters)
- .sort("updated_at", DESCENDING)
- .skip((page - 1) * page_size)
- .limit(page_size)
- .to_list(length=page_size)
- )
- profiles = {
- int(item["user_id"]): item
- async for item in profilesdb.find(
- {"user_id": {"$in": [int(item["user_id"]) for item in items]}}
- )
- }
- for item in items:
- item["profile"] = profiles.get(int(item["user_id"]), {})
- return items, total
- async def list_directory_events(
- *,
- query: str = "",
- event_type: str = "",
- page: int = 1,
- page_size: int = 20,
- ) -> tuple[list[dict[str, Any]], int]:
- filters: dict[str, Any] = {}
- if event_type:
- filters["event_type"] = event_type
- if query:
- pattern = re.compile(re.escape(query), re.IGNORECASE)
- filters["$or"] = [
- {"actor_name": pattern},
- {"reason": pattern},
- {"user_id": int(query)} if query.isdigit() else {"user_id": -1},
- ]
- total = await eventsdb.count_documents(filters)
- items = await (
- eventsdb.find(filters)
- .sort("created_at", DESCENDING)
- .skip((page - 1) * page_size)
- .limit(page_size)
- .to_list(length=page_size)
- )
- return items, total
- async def upsert_directory_membership(
- *,
- bot_id: str,
- chat_id: int,
- user_id: int,
- status: str,
- active: bool,
- chat_title: str = "",
- username: str | None = None,
- display_name_value: str = "",
- verified: bool = False,
- ) -> dict[str, Any]:
- await ensure_directory_indexes()
- now = utc_now()
- values: dict[str, Any] = {
- "status": status,
- "active": bool(active),
- "username": username,
- "display_name": display_name_value,
- "observed_at": now,
- "updated_at": now,
- }
- if chat_title:
- values["chat_title"] = chat_title
- if verified:
- values["verified_at"] = now
- await membershipsdb.update_one(
- {
- "bot_id": str(bot_id),
- "chat_id": int(chat_id),
- "user_id": int(user_id),
- },
- {"$set": values, "$setOnInsert": {"created_at": now}},
- upsert=True,
- )
- return await membershipsdb.find_one(
- {
- "bot_id": str(bot_id),
- "chat_id": int(chat_id),
- "user_id": int(user_id),
- }
- ) or {}
- async def list_membership_candidates(user_id: int) -> list[dict[str, Any]]:
- await ensure_directory_indexes()
- return await membershipsdb.find(
- {"user_id": int(user_id), "active": True}
- ).sort("verified_at", DESCENDING).to_list(length=200)
|