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(), } }, )