|
|
@@ -0,0 +1,1227 @@
|
|
|
+from __future__ import annotations
|
|
|
+
|
|
|
+import asyncio
|
|
|
+import hashlib
|
|
|
+import json
|
|
|
+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
|
|
|
+geocode_cachedb = control_db.directory_geocode_cache
|
|
|
+geocode_jobsdb = control_db.directory_geocode_jobs
|
|
|
+
|
|
|
+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] = {
|
|
|
+ "nominatim_enabled": True,
|
|
|
+ "nominatim_endpoint": "https://nominatim.openstreetmap.org/reverse",
|
|
|
+ "nominatim_contact": "",
|
|
|
+ "nominatim_user_agent": "TelegramBotAdmin/1.0",
|
|
|
+ "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)
|
|
|
+ await geocode_cachedb.create_index([("coordinate_key", ASCENDING)], unique=True)
|
|
|
+ await geocode_jobsdb.create_index([("user_id", ASCENDING)], unique=True)
|
|
|
+ await geocode_jobsdb.create_index(
|
|
|
+ [("status", ASCENDING), ("next_attempt_at", ASCENDING)]
|
|
|
+ )
|
|
|
+ _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: value
|
|
|
+ for key, value in stored.items()
|
|
|
+ if key not in {"_id", "settings_id"}
|
|
|
+ },
|
|
|
+ }
|
|
|
+
|
|
|
+
|
|
|
+async def set_directory_settings(values: dict[str, Any]) -> dict[str, Any]:
|
|
|
+ await ensure_directory_indexes()
|
|
|
+ current = await get_directory_settings()
|
|
|
+ target = {**current, **values}
|
|
|
+ endpoint = str(target.get("nominatim_endpoint") or "").strip()
|
|
|
+ if not endpoint.startswith(("https://", "http://127.0.0.1", "http://localhost")):
|
|
|
+ raise DirectoryDataError(
|
|
|
+ "invalid_nominatim_endpoint",
|
|
|
+ "Nominatim 地址必须使用 HTTPS;仅本机自建服务可以使用 HTTP。",
|
|
|
+ )
|
|
|
+ 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", "默认附近范围必须属于可选范围。"
|
|
|
+ )
|
|
|
+ contact = str(target.get("nominatim_contact") or "").strip()
|
|
|
+ user_agent = str(
|
|
|
+ target.get("nominatim_user_agent") or "TelegramBotAdmin/1.0"
|
|
|
+ ).strip()
|
|
|
+ if "\n" in contact or "\r" in contact or "\n" in user_agent or "\r" in user_agent:
|
|
|
+ raise DirectoryDataError(
|
|
|
+ "invalid_nominatim_identity",
|
|
|
+ "Nominatim 联系信息和 User-Agent 不能包含换行符。",
|
|
|
+ )
|
|
|
+ if contact and "@" not in contact and not contact.startswith(("https://", "http://")):
|
|
|
+ raise DirectoryDataError(
|
|
|
+ "invalid_nominatim_contact",
|
|
|
+ "Nominatim 联系信息必须是邮箱或 URL。",
|
|
|
+ )
|
|
|
+ target.update(
|
|
|
+ {
|
|
|
+ "nominatim_enabled": bool(target.get("nominatim_enabled")),
|
|
|
+ "nominatim_endpoint": endpoint,
|
|
|
+ "nominatim_contact": contact,
|
|
|
+ "nominatim_user_agent": user_agent,
|
|
|
+ "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
|
|
|
+
|
|
|
+
|
|
|
+def coordinate_key(longitude: float, latitude: float) -> str:
|
|
|
+ raw = json.dumps(
|
|
|
+ [float(longitude), float(latitude)],
|
|
|
+ ensure_ascii=True,
|
|
|
+ separators=(",", ":"),
|
|
|
+ )
|
|
|
+ return hashlib.sha256(raw.encode("ascii")).hexdigest()
|
|
|
+
|
|
|
+
|
|
|
+async def save_directory_location(
|
|
|
+ *,
|
|
|
+ user_id: int,
|
|
|
+ longitude: float,
|
|
|
+ latitude: float,
|
|
|
+ source: str,
|
|
|
+ actor_id: int | str,
|
|
|
+ actor_name: str,
|
|
|
+ geocode_consent: bool | None = None,
|
|
|
+ bot_id: str | None = None,
|
|
|
+) -> dict[str, Any]:
|
|
|
+ await ensure_directory_indexes()
|
|
|
+ lon, lat = validate_coordinates(longitude, latitude)
|
|
|
+ previous = await get_directory_location(int(user_id))
|
|
|
+ consent = (
|
|
|
+ geocode_consent
|
|
|
+ if geocode_consent is not None
|
|
|
+ else previous.get("geocode_consent") if previous else None
|
|
|
+ )
|
|
|
+ key = coordinate_key(lon, lat)
|
|
|
+ now = utc_now()
|
|
|
+ status = "pending" if consent is True else "disabled" if consent is False else "consent_required"
|
|
|
+ await locationsdb.update_one(
|
|
|
+ {"user_id": int(user_id)},
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "point": {"type": "Point", "coordinates": [lon, lat]},
|
|
|
+ "longitude": lon,
|
|
|
+ "latitude": lat,
|
|
|
+ "coordinate_key": key,
|
|
|
+ "geocode_consent": consent,
|
|
|
+ "geocode_status": status,
|
|
|
+ "region": None,
|
|
|
+ "source": source,
|
|
|
+ "source_bot_id": bot_id or str(BOT_PROFILE_ID),
|
|
|
+ "updated_at": now,
|
|
|
+ },
|
|
|
+ "$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={"geocode_consent": consent},
|
|
|
+ bot_id=bot_id,
|
|
|
+ )
|
|
|
+ if consent is True:
|
|
|
+ await enqueue_geocode_job(int(user_id))
|
|
|
+ return await get_directory_location(int(user_id)) or {}
|
|
|
+
|
|
|
+
|
|
|
+async def set_geocode_consent(
|
|
|
+ *,
|
|
|
+ user_id: int,
|
|
|
+ consent: bool,
|
|
|
+ source: str,
|
|
|
+ actor_id: int | str,
|
|
|
+ actor_name: str,
|
|
|
+) -> dict[str, Any]:
|
|
|
+ location = await get_directory_location(int(user_id))
|
|
|
+ if not location:
|
|
|
+ raise DirectoryDataError("location_required", "请先更新位置。")
|
|
|
+ updates: dict[str, Any] = {
|
|
|
+ "geocode_consent": bool(consent),
|
|
|
+ "geocode_status": "pending" if consent else "disabled",
|
|
|
+ "updated_at": utc_now(),
|
|
|
+ }
|
|
|
+ if not consent:
|
|
|
+ updates["region"] = None
|
|
|
+ await geocode_jobsdb.delete_one({"user_id": int(user_id)})
|
|
|
+ await locationsdb.update_one({"user_id": int(user_id)}, {"$set": updates})
|
|
|
+ await record_directory_event(
|
|
|
+ event_type="geocode_consent_updated",
|
|
|
+ user_id=int(user_id),
|
|
|
+ actor_id=actor_id,
|
|
|
+ actor_name=actor_name,
|
|
|
+ source=source,
|
|
|
+ reason="同意区域解析" if consent else "关闭区域解析",
|
|
|
+ metadata={"geocode_consent": bool(consent)},
|
|
|
+ )
|
|
|
+ if consent:
|
|
|
+ await enqueue_geocode_job(int(user_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 geocode_jobsdb.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,
|
|
|
+) -> dict[str, Any]:
|
|
|
+ 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
|
|
|
+ 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),
|
|
|
+ "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,
|
|
|
+ )
|
|
|
+ return await get_directory_profile(int(user_id)) or {}
|
|
|
+
|
|
|
+
|
|
|
+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,
|
|
|
+ "region": location.get("region"),
|
|
|
+ "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"],
|
|
|
+ "region": item.get("region"),
|
|
|
+ "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 = "",
|
|
|
+ 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 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"] = (
|
|
|
+ {
|
|
|
+ "region": location.get("region"),
|
|
|
+ "geocode_status": location.get("geocode_status"),
|
|
|
+ "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,
|
|
|
+ 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 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"] = (
|
|
|
+ {
|
|
|
+ "region": location.get("region"),
|
|
|
+ "geocode_status": location.get("geocode_status"),
|
|
|
+ "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)
|
|
|
+
|
|
|
+
|
|
|
+async def enqueue_geocode_job(user_id: int) -> dict[str, Any] | None:
|
|
|
+ await ensure_directory_indexes()
|
|
|
+ location = await get_directory_location(int(user_id))
|
|
|
+ if not location or location.get("geocode_consent") is not True:
|
|
|
+ return None
|
|
|
+ key = str(location["coordinate_key"])
|
|
|
+ cached = await geocode_cachedb.find_one({"coordinate_key": key})
|
|
|
+ if cached:
|
|
|
+ await locationsdb.update_one(
|
|
|
+ {"user_id": int(user_id), "coordinate_key": key},
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "region": cached.get("region"),
|
|
|
+ "geocode_status": "completed",
|
|
|
+ "geocoded_at": cached.get("created_at"),
|
|
|
+ "updated_at": utc_now(),
|
|
|
+ }
|
|
|
+ },
|
|
|
+ )
|
|
|
+ return None
|
|
|
+ now = utc_now()
|
|
|
+ await geocode_jobsdb.update_one(
|
|
|
+ {"user_id": int(user_id)},
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "coordinate_key": key,
|
|
|
+ "longitude": float(location["longitude"]),
|
|
|
+ "latitude": float(location["latitude"]),
|
|
|
+ "status": "pending",
|
|
|
+ "attempts": 0,
|
|
|
+ "next_attempt_at": now,
|
|
|
+ "updated_at": now,
|
|
|
+ },
|
|
|
+ "$setOnInsert": {"created_at": now},
|
|
|
+ },
|
|
|
+ upsert=True,
|
|
|
+ )
|
|
|
+ return await geocode_jobsdb.find_one({"user_id": int(user_id)})
|
|
|
+
|
|
|
+
|
|
|
+async def claim_geocode_job() -> dict[str, Any] | None:
|
|
|
+ await ensure_directory_indexes()
|
|
|
+ now = utc_now()
|
|
|
+ return await geocode_jobsdb.find_one_and_update(
|
|
|
+ {
|
|
|
+ "status": "pending",
|
|
|
+ "$or": [
|
|
|
+ {"next_attempt_at": {"$lte": now}},
|
|
|
+ {"next_attempt_at": {"$exists": False}},
|
|
|
+ ],
|
|
|
+ },
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "status": "processing",
|
|
|
+ "claimed_at": now,
|
|
|
+ "updated_at": now,
|
|
|
+ },
|
|
|
+ "$inc": {"attempts": 1},
|
|
|
+ },
|
|
|
+ sort=[("created_at", ASCENDING)],
|
|
|
+ return_document=ReturnDocument.AFTER,
|
|
|
+ )
|
|
|
+
|
|
|
+
|
|
|
+async def complete_geocode_job(
|
|
|
+ *,
|
|
|
+ job: dict[str, Any],
|
|
|
+ region: dict[str, Any],
|
|
|
+) -> None:
|
|
|
+ now = utc_now()
|
|
|
+ key = str(job["coordinate_key"])
|
|
|
+ await geocode_cachedb.update_one(
|
|
|
+ {"coordinate_key": key},
|
|
|
+ {
|
|
|
+ "$setOnInsert": {
|
|
|
+ "coordinate_key": key,
|
|
|
+ "region": region,
|
|
|
+ "provider": "nominatim",
|
|
|
+ "created_at": now,
|
|
|
+ }
|
|
|
+ },
|
|
|
+ upsert=True,
|
|
|
+ )
|
|
|
+ await locationsdb.update_one(
|
|
|
+ {
|
|
|
+ "user_id": int(job["user_id"]),
|
|
|
+ "coordinate_key": key,
|
|
|
+ "geocode_consent": True,
|
|
|
+ },
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "region": region,
|
|
|
+ "geocode_status": "completed",
|
|
|
+ "geocoded_at": now,
|
|
|
+ "updated_at": now,
|
|
|
+ }
|
|
|
+ },
|
|
|
+ )
|
|
|
+ await geocode_jobsdb.delete_one(
|
|
|
+ {"user_id": int(job["user_id"]), "coordinate_key": key}
|
|
|
+ )
|
|
|
+
|
|
|
+
|
|
|
+async def fail_geocode_job(
|
|
|
+ *,
|
|
|
+ job: dict[str, Any],
|
|
|
+ error: str,
|
|
|
+ retry_after_seconds: int = 60,
|
|
|
+) -> None:
|
|
|
+ attempts = int(job.get("attempts") or 1)
|
|
|
+ key_filter = {
|
|
|
+ "user_id": int(job["user_id"]),
|
|
|
+ "coordinate_key": str(job["coordinate_key"]),
|
|
|
+ }
|
|
|
+ if attempts >= 3:
|
|
|
+ await geocode_jobsdb.update_one(
|
|
|
+ key_filter,
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "status": "failed",
|
|
|
+ "error": error[:500],
|
|
|
+ "updated_at": utc_now(),
|
|
|
+ }
|
|
|
+ },
|
|
|
+ )
|
|
|
+ await locationsdb.update_one(
|
|
|
+ key_filter,
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "geocode_status": "failed",
|
|
|
+ "geocode_error": error[:500],
|
|
|
+ "updated_at": utc_now(),
|
|
|
+ }
|
|
|
+ },
|
|
|
+ )
|
|
|
+ return
|
|
|
+ await geocode_jobsdb.update_one(
|
|
|
+ key_filter,
|
|
|
+ {
|
|
|
+ "$set": {
|
|
|
+ "status": "pending",
|
|
|
+ "error": error[:500],
|
|
|
+ "next_attempt_at": utc_now()
|
|
|
+ + timedelta(seconds=max(1, retry_after_seconds)),
|
|
|
+ "updated_at": utc_now(),
|
|
|
+ }
|
|
|
+ },
|
|
|
+ )
|