|
|
@@ -1,8 +1,6 @@
|
|
|
from __future__ import annotations
|
|
|
|
|
|
import asyncio
|
|
|
-import hashlib
|
|
|
-import json
|
|
|
import math
|
|
|
import re
|
|
|
from datetime import UTC, datetime, timedelta
|
|
|
@@ -19,8 +17,6 @@ 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"
|
|
|
@@ -39,10 +35,6 @@ 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],
|
|
|
@@ -129,11 +121,6 @@ async def ensure_directory_indexes() -> None:
|
|
|
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
|
|
|
|
|
|
|
|
|
@@ -143,9 +130,9 @@ async def get_directory_settings() -> dict[str, Any]:
|
|
|
return {
|
|
|
**DEFAULT_DIRECTORY_SETTINGS,
|
|
|
**{
|
|
|
- key: value
|
|
|
- for key, value in stored.items()
|
|
|
- if key not in {"_id", "settings_id"}
|
|
|
+ key: stored[key]
|
|
|
+ for key in DEFAULT_DIRECTORY_SETTINGS
|
|
|
+ if key in stored
|
|
|
},
|
|
|
}
|
|
|
|
|
|
@@ -154,12 +141,6 @@ 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(
|
|
|
{
|
|
|
@@ -179,26 +160,8 @@ async def set_directory_settings(values: dict[str, Any]) -> dict[str, Any]:
|
|
|
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,
|
|
|
@@ -330,15 +293,6 @@ def validate_coordinates(longitude: float, latitude: float) -> tuple[float, floa
|
|
|
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,
|
|
|
@@ -347,20 +301,11 @@ async def save_directory_location(
|
|
|
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)},
|
|
|
{
|
|
|
@@ -368,14 +313,18 @@ async def save_directory_location(
|
|
|
"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,
|
|
|
},
|
|
|
+ "$unset": {
|
|
|
+ "coordinate_key": "",
|
|
|
+ "geocode_consent": "",
|
|
|
+ "geocode_status": "",
|
|
|
+ "region": "",
|
|
|
+ "geocoded_at": "",
|
|
|
+ "geocode_error": "",
|
|
|
+ },
|
|
|
"$setOnInsert": {"created_at": now},
|
|
|
},
|
|
|
upsert=True,
|
|
|
@@ -400,45 +349,9 @@ async def save_directory_location(
|
|
|
actor_name=actor_name,
|
|
|
source=source,
|
|
|
reason="更新位置",
|
|
|
- metadata={"geocode_consent": consent},
|
|
|
+ metadata={"location_source": source},
|
|
|
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 {}
|
|
|
|
|
|
|
|
|
@@ -452,7 +365,6 @@ async def clear_directory_location(
|
|
|
) -> 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)},
|
|
|
{
|
|
|
@@ -482,7 +394,9 @@ async def submit_teacher_application(
|
|
|
user_id: int,
|
|
|
source: str,
|
|
|
bot_id: str,
|
|
|
-) -> dict[str, Any]:
|
|
|
+ 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))
|
|
|
@@ -499,7 +413,7 @@ async def submit_teacher_application(
|
|
|
"teacher_revoked", "老师资格已被撤销,只能由管理员重新批准。"
|
|
|
)
|
|
|
if profile.get("application_status") == APPLICATION_PENDING:
|
|
|
- return profile
|
|
|
+ return profile, False
|
|
|
now = utc_now()
|
|
|
await profilesdb.update_one(
|
|
|
{"user_id": int(user_id)},
|
|
|
@@ -508,6 +422,8 @@ async def submit_teacher_application(
|
|
|
"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,
|
|
|
@@ -523,8 +439,9 @@ async def submit_teacher_application(
|
|
|
source=source,
|
|
|
reason="申请成为老师",
|
|
|
bot_id=bot_id,
|
|
|
+ chat_id=int(chat_id),
|
|
|
)
|
|
|
- return await get_directory_profile(int(user_id)) or {}
|
|
|
+ return await get_directory_profile(int(user_id)) or {}, True
|
|
|
|
|
|
|
|
|
async def decide_teacher_application(
|
|
|
@@ -768,7 +685,6 @@ async def _list_teachers_fallback(
|
|
|
{
|
|
|
**profiles[int(location["user_id"])],
|
|
|
"distance_meters": distance,
|
|
|
- "region": location.get("region"),
|
|
|
"location_updated_at": location.get("updated_at"),
|
|
|
}
|
|
|
)
|
|
|
@@ -839,7 +755,6 @@ async def list_directory_teachers(
|
|
|
{
|
|
|
**item["profile"],
|
|
|
"distance_meters": item["distance_meters"],
|
|
|
- "region": item.get("region"),
|
|
|
"location_updated_at": item.get("updated_at"),
|
|
|
}
|
|
|
for item in result.get("items", [])
|
|
|
@@ -861,6 +776,8 @@ 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]:
|
|
|
@@ -872,6 +789,12 @@ async def list_teacher_applications(
|
|
|
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"] = [
|
|
|
@@ -897,8 +820,7 @@ async def list_teacher_applications(
|
|
|
location = location_map.get(int(item["user_id"]))
|
|
|
item["location"] = (
|
|
|
{
|
|
|
- "region": location.get("region"),
|
|
|
- "geocode_status": location.get("geocode_status"),
|
|
|
+ "source": location.get("source"),
|
|
|
"updated_at": location.get("updated_at"),
|
|
|
}
|
|
|
if location
|
|
|
@@ -913,6 +835,8 @@ 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]:
|
|
|
@@ -921,6 +845,12 @@ async def list_directory_profiles(
|
|
|
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"] = [
|
|
|
@@ -946,8 +876,7 @@ async def list_directory_profiles(
|
|
|
location = locations.get(int(item["user_id"]))
|
|
|
item["location"] = (
|
|
|
{
|
|
|
- "region": location.get("region"),
|
|
|
- "geocode_status": location.get("geocode_status"),
|
|
|
+ "source": location.get("source"),
|
|
|
"updated_at": location.get("updated_at"),
|
|
|
}
|
|
|
if location
|
|
|
@@ -1075,153 +1004,3 @@ async def list_membership_candidates(user_id: int) -> list[dict[str, Any]]:
|
|
|
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(),
|
|
|
- }
|
|
|
- },
|
|
|
- )
|