修复网站 API 阻塞与积分同步问题
- uvicorn 改用 import string 格式启动 4 workers,防止数据采集阻塞 event loop - prewarm 去掉 --force-refresh,仅依赖缓存预热避免每 5 分钟全量采集 - 积分变动时同步写入 Supabase user_metadata,避免前端回退路径失败时显示 0 积分 - /api/auth/me 积分解析增加 Supabase email 回退路径
This commit is contained in:
+1
-1
@@ -39,7 +39,7 @@ services:
|
|||||||
container_name: polyweather_prewarm
|
container_name: polyweather_prewarm
|
||||||
restart: unless-stopped
|
restart: unless-stopped
|
||||||
profiles: ["workers"]
|
profiles: ["workers"]
|
||||||
command: python scripts/prewarm_dashboard_worker.py --include-detail --include-market --force-refresh
|
command: python scripts/prewarm_dashboard_worker.py --include-detail --include-market
|
||||||
volumes:
|
volumes:
|
||||||
- ${POLYWEATHER_RUNTIME_DATA_DIR:-/var/lib/polyweather}:/var/lib/polyweather
|
- ${POLYWEATHER_RUNTIME_DATA_DIR:-/var/lib/polyweather}:/var/lib/polyweather
|
||||||
- ${POLYWEATHER_RUNTIME_DATA_DIR:-/var/lib/polyweather}:/app/data
|
- ${POLYWEATHER_RUNTIME_DATA_DIR:-/var/lib/polyweather}:/app/data
|
||||||
|
|||||||
@@ -42,6 +42,75 @@ class DBManager:
|
|||||||
"Prefer": "return=minimal",
|
"Prefer": "return=minimal",
|
||||||
}
|
}
|
||||||
|
|
||||||
|
def _supabase_admin_users_endpoint(self) -> str:
|
||||||
|
supabase_url = str(os.getenv("SUPABASE_URL") or "").strip().rstrip("/")
|
||||||
|
if not supabase_url:
|
||||||
|
return ""
|
||||||
|
return f"{supabase_url}/auth/v1/admin/users"
|
||||||
|
|
||||||
|
def _sync_points_to_supabase_user_metadata(self, telegram_id: int) -> bool:
|
||||||
|
supabase_url = str(os.getenv("SUPABASE_URL") or "").strip().rstrip("/")
|
||||||
|
if not supabase_url:
|
||||||
|
return False
|
||||||
|
headers = self._supabase_service_headers()
|
||||||
|
if not headers:
|
||||||
|
return False
|
||||||
|
endpoint = self._supabase_admin_users_endpoint()
|
||||||
|
if not endpoint:
|
||||||
|
return False
|
||||||
|
|
||||||
|
supabase_user_id = None
|
||||||
|
points = 0
|
||||||
|
with self._get_connection() as conn:
|
||||||
|
conn.row_factory = sqlite3.Row
|
||||||
|
row = conn.execute(
|
||||||
|
"SELECT supabase_user_id FROM supabase_bindings WHERE telegram_id = ? LIMIT 1",
|
||||||
|
(int(telegram_id),),
|
||||||
|
).fetchone()
|
||||||
|
if row and row["supabase_user_id"]:
|
||||||
|
supabase_user_id = str(row["supabase_user_id"]).strip()
|
||||||
|
if not supabase_user_id:
|
||||||
|
row = conn.execute(
|
||||||
|
"SELECT supabase_user_id FROM users WHERE telegram_id = ? LIMIT 1",
|
||||||
|
(int(telegram_id),),
|
||||||
|
).fetchone()
|
||||||
|
if row and row["supabase_user_id"]:
|
||||||
|
supabase_user_id = str(row["supabase_user_id"]).strip()
|
||||||
|
if not supabase_user_id:
|
||||||
|
return False
|
||||||
|
pts_row = conn.execute(
|
||||||
|
"SELECT points FROM users WHERE telegram_id = ? LIMIT 1",
|
||||||
|
(int(telegram_id),),
|
||||||
|
).fetchone()
|
||||||
|
if pts_row:
|
||||||
|
points = max(0, int(pts_row["points"] or 0))
|
||||||
|
|
||||||
|
try:
|
||||||
|
resp = requests.patch(
|
||||||
|
f"{endpoint}/{supabase_user_id}",
|
||||||
|
json={"user_metadata": {"points": points}},
|
||||||
|
headers={**headers, "Prefer": "return=minimal"},
|
||||||
|
timeout=8,
|
||||||
|
)
|
||||||
|
if resp.status_code not in (200, 204):
|
||||||
|
logger.warning(
|
||||||
|
"supabase points sync failed tg={} suid={} status={} body={}",
|
||||||
|
telegram_id,
|
||||||
|
supabase_user_id,
|
||||||
|
resp.status_code,
|
||||||
|
(resp.text or "")[:200],
|
||||||
|
)
|
||||||
|
return False
|
||||||
|
return True
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning(
|
||||||
|
"supabase points sync error tg={} suid={}: {}",
|
||||||
|
telegram_id,
|
||||||
|
supabase_user_id,
|
||||||
|
exc,
|
||||||
|
)
|
||||||
|
return False
|
||||||
|
|
||||||
def _sync_supabase_profile_telegram_fields(
|
def _sync_supabase_profile_telegram_fields(
|
||||||
self,
|
self,
|
||||||
*,
|
*,
|
||||||
@@ -1010,6 +1079,36 @@ class DBManager:
|
|||||||
except Exception:
|
except Exception:
|
||||||
return 0
|
return 0
|
||||||
|
|
||||||
|
def get_points_by_supabase_email(self, supabase_email: str) -> int:
|
||||||
|
email = str(supabase_email or "").strip().lower()
|
||||||
|
if not email:
|
||||||
|
return 0
|
||||||
|
with self._get_connection() as conn:
|
||||||
|
conn.row_factory = sqlite3.Row
|
||||||
|
row = conn.execute(
|
||||||
|
"""
|
||||||
|
SELECT points
|
||||||
|
FROM users
|
||||||
|
WHERE lower(trim(COALESCE(supabase_email, ''))) = ?
|
||||||
|
LIMIT 1
|
||||||
|
""",
|
||||||
|
(email,),
|
||||||
|
).fetchone()
|
||||||
|
if not row:
|
||||||
|
row = conn.execute(
|
||||||
|
"""
|
||||||
|
SELECT u.points
|
||||||
|
FROM users u
|
||||||
|
JOIN supabase_bindings b ON b.telegram_id = u.telegram_id
|
||||||
|
WHERE lower(trim(COALESCE(b.supabase_email, ''))) = ?
|
||||||
|
LIMIT 1
|
||||||
|
""",
|
||||||
|
(email,),
|
||||||
|
).fetchone()
|
||||||
|
if row:
|
||||||
|
return max(0, int(row["points"] or 0))
|
||||||
|
return 0
|
||||||
|
|
||||||
def grant_points_by_supabase_email(
|
def grant_points_by_supabase_email(
|
||||||
self,
|
self,
|
||||||
supabase_email: str,
|
supabase_email: str,
|
||||||
@@ -1048,6 +1147,7 @@ class DBManager:
|
|||||||
(after, telegram_id),
|
(after, telegram_id),
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
|
self._sync_points_to_supabase_user_metadata(telegram_id)
|
||||||
return {
|
return {
|
||||||
"ok": True,
|
"ok": True,
|
||||||
"telegram_id": telegram_id,
|
"telegram_id": telegram_id,
|
||||||
@@ -1421,6 +1521,7 @@ class DBManager:
|
|||||||
points=weekly_points + total_added,
|
points=weekly_points + total_added,
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
|
self._sync_points_to_supabase_user_metadata(telegram_id)
|
||||||
return {
|
return {
|
||||||
"awarded": True,
|
"awarded": True,
|
||||||
"reason": "ok",
|
"reason": "ok",
|
||||||
@@ -1492,6 +1593,7 @@ class DBManager:
|
|||||||
(new_balance, telegram_id),
|
(new_balance, telegram_id),
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
|
self._sync_points_to_supabase_user_metadata(telegram_id)
|
||||||
return {"ok": True, "balance": new_balance, "spent": amount}
|
return {"ok": True, "balance": new_balance, "spent": amount}
|
||||||
|
|
||||||
def spend_points_by_supabase_user_id(self, supabase_user_id: str, amount: int) -> Dict[str, Any]:
|
def spend_points_by_supabase_user_id(self, supabase_user_id: str, amount: int) -> Dict[str, Any]:
|
||||||
@@ -1529,6 +1631,7 @@ class DBManager:
|
|||||||
(new_balance, telegram_id),
|
(new_balance, telegram_id),
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
|
self._sync_points_to_supabase_user_metadata(telegram_id)
|
||||||
return {"ok": True, "balance": new_balance, "spent": amount}
|
return {"ok": True, "balance": new_balance, "spent": amount}
|
||||||
|
|
||||||
def set_premium(self, telegram_id: int, plan: str, months: int = 1):
|
def set_premium(self, telegram_id: int, plan: str, months: int = 1):
|
||||||
@@ -1839,6 +1942,8 @@ class DBManager:
|
|||||||
),
|
),
|
||||||
)
|
)
|
||||||
conn.commit()
|
conn.commit()
|
||||||
|
if bonus > 0:
|
||||||
|
self._sync_points_to_supabase_user_metadata(int(telegram_id))
|
||||||
return True
|
return True
|
||||||
|
|
||||||
def append_airport_obs(
|
def append_airport_obs(
|
||||||
|
|||||||
+6
-1
@@ -35,4 +35,9 @@ __all__ = [
|
|||||||
if __name__ == "__main__":
|
if __name__ == "__main__":
|
||||||
import uvicorn
|
import uvicorn
|
||||||
|
|
||||||
uvicorn.run(app, host="0.0.0.0", port=8000)
|
uvicorn.run(
|
||||||
|
"web.app:app",
|
||||||
|
host="0.0.0.0",
|
||||||
|
port=8000,
|
||||||
|
workers=int(os.getenv("UVICORN_WORKERS", "4")),
|
||||||
|
)
|
||||||
|
|||||||
+21
-11
@@ -239,19 +239,29 @@ def _resolve_auth_points(request: Request) -> int:
|
|||||||
points = max(0, int(raw_points or 0))
|
points = max(0, int(raw_points or 0))
|
||||||
except Exception:
|
except Exception:
|
||||||
points = 0
|
points = 0
|
||||||
if points > 0:
|
|
||||||
return points
|
|
||||||
|
|
||||||
user_id = str(getattr(request.state, "auth_user_id", "") or "").strip()
|
user_id = str(getattr(request.state, "auth_user_id", "") or "").strip()
|
||||||
if not user_id:
|
|
||||||
return points
|
if user_id:
|
||||||
try:
|
try:
|
||||||
db_points = _account_db.get_points_by_supabase_user_id(user_id)
|
db_points = _account_db.get_points_by_supabase_user_id(user_id)
|
||||||
if db_points > points:
|
if db_points > points:
|
||||||
request.state.auth_points = db_points
|
request.state.auth_points = db_points
|
||||||
return db_points
|
points = db_points
|
||||||
except Exception as exc:
|
except Exception as exc:
|
||||||
logger.warning(f"auth points fallback failed user_id={user_id}: {exc}")
|
logger.warning(f"auth points fallback failed user_id={user_id}: {exc}")
|
||||||
|
|
||||||
|
if points <= 0:
|
||||||
|
email = str(getattr(request.state, "auth_email", "") or "").strip().lower()
|
||||||
|
if email:
|
||||||
|
try:
|
||||||
|
email_points = _account_db.get_points_by_supabase_email(email)
|
||||||
|
if email_points > points:
|
||||||
|
request.state.auth_points = email_points
|
||||||
|
points = email_points
|
||||||
|
except Exception as exc:
|
||||||
|
logger.warning(f"auth points email fallback failed email={email}: {exc}")
|
||||||
|
|
||||||
return points
|
return points
|
||||||
|
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user