diff --git a/CLAUDE.md b/CLAUDE.md index 671034a8..d06cefc5 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -38,8 +38,8 @@ Users (Web / Telegram) → Next.js Frontend (Vercel) → FastAPI /web/app.py Payment Layer (Intent + Event + Confirm Loop) ``` -- **Backend**: FastAPI on port 8000 (`web/app.py` → `web/core.py` + `web/routes.py` + `web/analysis_service.py`) -- **Frontend**: Next.js 15 + React 19 + TypeScript + Tailwind CSS 3 on port 3000 (dev) +- **Backend**: FastAPI on port 8000 (`web/app.py` → `web/app_factory.py` → `web/routers/` (8 route modules) + `web/services/` (14 service modules) + `web/core.py`) +- **Frontend**: Next.js 15 + React 19 + TypeScript + Tailwind CSS 3 + shadcn/ui (new-york style) on port 3000 (dev) - **Bot**: Telegram bot via `bot_listener.py` → `src/bot/` - **Shared analysis core** in `src/` is used by both web API and bot - **Scan Terminal**: Real-time city opportunity scanning (`web/scan_terminal_service.py` and `frontend/components/dashboard/scan-terminal/`) @@ -54,9 +54,12 @@ Users (Web / Telegram) → Next.js Frontend (Vercel) → FastAPI /web/app.py ```bash cd frontend npm ci -npm run dev # Next.js dev server -npm run build # Production build -npm run lint # ESLint via next lint +npm run dev # Next.js dev server (runs sync-next-server-chunks.mjs first) +npm run build # Production build (runs sync-next-server-chunks.mjs after) +npm run start # Production server +npm run lint # ESLint via next lint +npm run typecheck # tsc --noEmit +npm run test:business # Business state tests (also runs in CI) ``` ### Backend (dev on port 8000) @@ -73,10 +76,9 @@ python run.py ### Docker (production-like stack) ```bash -docker compose up -d --build # bot + web API -docker compose --profile workers up -d # + prewarm worker -docker compose --profile monitoring up -d # + Prometheus/Grafana/Alertmanager +docker compose up -d --build # bot + web API (polyweather + polyweather_web) ``` +The compose file defines two services: `polyweather` (bot) and `polyweather_web` (FastAPI on :8000). Prewarm worker and monitoring profiles were removed in v1.6.0. ### Python tests ```bash @@ -101,21 +103,26 @@ curl http://127.0.0.1:8000/metrics | Directory | Purpose | |-----------|---------| -| `src/data_collection/` | Weather sources (METAR, TAF, Open-Meteo, JMA, KMA, MGM, NMC, Russia stations, settlement sources), city registry (52 cities), Polymarket readonly layer. Also: `madis_sources.py` (NOAA 5-min NetCDF), `amos_station_sources.py` (Korean runway sensors), `country_networks.py` (per-country provider routing + `_airport_primary_from_raw`) | | `src/analysis/` | DEB algorithm, trend engine, probability calibration (EMOS/LGBM), market alert engine, settlement rounding | -| `src/models/` | LightGBM daily-high model training and feature engineering | -| `src/payments/` | Onchain checkout, event listener, confirm loop, contract audit | +| `src/auth/` | Supabase entitlement checks, Telegram group pricing | | `src/bot/` | Telegram bot handlers and orchestrator | | `src/database/` | SQLite-based runtime state, DB manager, daily/truth/training feature repositories | -| `web/` | FastAPI app, routes (~65K), analysis service (~130K), scan terminal service (~56K), AI scan modules | +| `src/data_collection/` | Weather sources (METAR, TAF, Open-Meteo, JMA, KMA, MGM, NMC, Russia stations, settlement sources), city registry (52 cities), Polymarket readonly layer. Also: `madis_sources.py` (NOAA 5-min NetCDF), `amos_station_sources.py` (Korean runway sensors), `country_networks.py` (per-country provider routing + `_airport_primary_from_raw`) | +| `src/data_mining/` | Historical data fetch utilities | +| `src/models/` | LightGBM daily-high model training and feature engineering | +| `src/onchain/` | Polygon wallet watcher, Polymarket wallet activity watcher | +| `src/payments/` | Onchain checkout, event listener, confirm loop, contract audit | +| `src/strategy/` | Trading strategy modules | +| `src/trading/` | Trading execution modules | +| `src/utils/` | Shared utilities: config loader, logging, metrics, Telegram push, chat ID helpers | +| `web/` | FastAPI app (`app.py` → `app_factory.py`), `routers/` (8 route modules), `services/` (14 service modules), `core.py`, scan terminal modules (AI fallback, AI prompts, METAR gate, city rows, ranker, cache) | | `frontend/app/` | Next.js App Router pages (dashboard, account, auth, docs, ops, probabilities, scan) | | `frontend/components/dashboard/` | Dashboard UI components (map, sidebar, detail panel, modals, charts, scan terminal). `scan-root-styles.ts` is the CSS Module barrel, combining 22 module roots into one pre-composed className. `monitoring/` subdirectory: `MonitorPanel`, `monitor-temperature.ts` (temp resolution chain), `monitor-refresh-policy.ts`. | | `frontend/lib/` | Shared client logic: types (`dashboard-types.ts`, including `AirportCurrentConditions`, `CityDetail`), API client, chart utils, i18n, `source-freshness.ts` (per-source freshness with `expected_next_update_at`), dashboard utils | | `frontend/hooks/` | React hooks: dashboard store (global state), Leaflet map, chart helper | -| `scripts/` | Operational scripts: probability calibration training, backfills, payment reconciliation, prewarm worker | +| `scripts/` | Operational scripts: probability calibration training, backfills, payment reconciliation. `supabase/` subdirectory: DB schema and migration SQL. | | `config/` | YAML config (city list, weather settings, logging) | | `docs/` | Bilingual product & technical docs | -| `monitoring/` | Prometheus/Grafana/Alertmanager configs | ## Key Technical Details @@ -124,6 +131,7 @@ curl http://127.0.0.1:8000/metrics - **Frontend package manager**: npm - **State storage**: SQLite primary path (set via `POLYWEATHER_STATE_STORAGE_MODE=sqlite` + `POLYWEATHER_DB_PATH`). Legacy JSON/JSONL files are migration/fallback only. - **Runtime data**: External dir recommended (`POLYWEATHER_RUNTIME_DATA_DIR=/var/lib/polyweather`) to avoid git conflicts +- **Configuration**: `.env.example` is the comprehensive reference (8 config sections: runtime, Telegram, weather cache, auth, ops, frontend, optional modules, Polygon monitor). Copy to `.env` and fill in secrets. - **Auth gating** (frontend middleware): Token-based (`POLYWEATHER_DASHBOARD_ACCESS_TOKEN`) or Supabase session-based (`POLYWEATHER_AUTH_ENABLED`). Local dev hosts bypass auth. - **CORS**: Allowed origins from `WEB_CORS_ORIGINS` env var (defaults: localhost:3000, polyweather-pro.vercel.app) - **EMOS/CRPS calibration**: Trainable but production should use `legacy` or `emos_shadow` engine; `emos_primary` only after local evaluation + manual rollout diff --git a/frontend/app/api/auth/telegram/bind-by-token/route.ts b/frontend/app/api/auth/telegram/bind-by-token/route.ts new file mode 100644 index 00000000..0ffb738c --- /dev/null +++ b/frontend/app/api/auth/telegram/bind-by-token/route.ts @@ -0,0 +1,44 @@ +import { NextRequest, NextResponse } from "next/server"; +import { + applyAuthResponseCookies, + buildBackendRequestHeaders, +} from "@/lib/backend-auth"; +import { + buildProxyExceptionResponse, + buildUpstreamErrorResponse, +} from "@/lib/api-proxy"; + +const API_BASE = process.env.POLYWEATHER_API_BASE_URL; + +export async function POST(req: NextRequest) { + if (!API_BASE) { + return NextResponse.json( + { error: "POLYWEATHER_API_BASE_URL is not configured" }, + { status: 500 }, + ); + } + try { + const body = await req.text(); + const auth = await buildBackendRequestHeaders(req); + const headers = new Headers(auth.headers); + headers.set("Content-Type", "application/json"); + const res = await fetch(`${API_BASE}/api/auth/telegram/bind-by-token`, { + method: "POST", + headers, + body, + cache: "no-store", + }); + if (!res.ok) { + const raw = await res.text(); + const response = buildUpstreamErrorResponse(res.status, raw); + return applyAuthResponseCookies(response, auth.response); + } + const data = await res.json(); + const response = NextResponse.json(data); + return applyAuthResponseCookies(response, auth.response); + } catch (error) { + return buildProxyExceptionResponse(error, { + publicMessage: "Failed to bind Telegram account", + }); + } +} diff --git a/frontend/components/account/AccountCenter.tsx b/frontend/components/account/AccountCenter.tsx index c0d0f6c6..20dd5b33 100644 --- a/frontend/components/account/AccountCenter.tsx +++ b/frontend/components/account/AccountCenter.tsx @@ -1,6 +1,6 @@ "use client"; -import { useCallback, useEffect, useMemo, useRef, useState } from "react"; +import { useCallback, useEffect, useMemo, useState } from "react"; import Link from "next/link"; import { useRouter } from "next/navigation"; import type { User } from "@supabase/supabase-js"; @@ -80,16 +80,6 @@ type TelegramPricing = { pricing_source?: string; }; -type TelegramAuthPayload = { - id: number; - first_name?: string; - last_name?: string; - username?: string; - photo_url?: string; - auth_date: number; - hash: string; -}; - type PaymentPlan = { plan_code: string; plan_id: number; @@ -250,11 +240,6 @@ const TELEGRAM_GROUP_URL = String( const TELEGRAM_BOT_URL = String( process.env.NEXT_PUBLIC_TELEGRAM_BOT_URL || "https://t.me/WeatherQuant_bot", ).trim(); -const TELEGRAM_LOGIN_BOT_USERNAME = String( - process.env.NEXT_PUBLIC_TELEGRAM_LOGIN_BOT_USERNAME || "WeatherQuant_bot", -) - .replace(/^@/, "") - .trim(); const TELEGRAM_MARKET_CHANNEL_URL = "https://t.me/+hGAk7JsjtdhiOTUx"; const TELEGRAM_TOPICS_GROUP_URL = "https://t.me/+8vel7rwjZagxODUx"; const SUBSCRIPTION_HELP_HREF = "/subscription-help"; @@ -898,8 +883,6 @@ export function AccountCenter() { const [lastTxHash, setLastTxHash] = useState(""); const [manualPayment, setManualPayment] = useState(null); const [manualTxHash, setManualTxHash] = useState(""); - const [telegramLoginBusy, setTelegramLoginBusy] = useState(false); - const telegramLoginWidgetRef = useRef(null); const [lastPaymentStartedAt, setLastPaymentStartedAt] = useState(0); const [showSecondarySections, setShowSecondarySections] = useState(false); const [reconcileBusy, setReconcileBusy] = useState(false); @@ -1510,64 +1493,44 @@ export function AccountCenter() { const userId = backend?.user_id || user?.id || ""; const isAuthenticated = Boolean(userId); const email = backend?.email || user?.email || ""; + // Handle ?bind_token=xxx from Telegram bot /bind deep link useEffect(() => { - const container = telegramLoginWidgetRef.current; - if (!container || !TELEGRAM_LOGIN_BOT_USERNAME || !isAuthenticated) return; - container.innerHTML = ""; - const callbackName = "onPolyWeatherTelegramAuth"; - (window as unknown as Record)[callbackName] = async ( - payload: TelegramAuthPayload, - ) => { - setTelegramLoginBusy(true); + if (!isAuthenticated) return; + const params = new URLSearchParams(window.location.search); + const token = params.get("bind_token"); + if (!token) return; + // Remove token from URL so refresh doesn't retry + const url = new URL(window.location.href); + url.searchParams.delete("bind_token"); + window.history.replaceState(null, "", url.toString()); + (async () => { setPaymentError(""); setPaymentInfo(""); try { const authHeaders = await buildAuthedHeaders(true, false); - const res = await fetch("/api/auth/telegram/login", { + const res = await fetch("/api/auth/telegram/bind-by-token", { method: "POST", headers: authHeaders, - body: JSON.stringify(payload), + body: JSON.stringify({ token }), }); if (!res.ok) { const raw = (await res.text()).slice(0, 350); - throw new Error(`telegram login failed: ${raw}`); + throw new Error(`bind failed: ${raw}`); } const data = (await res.json()) as { telegram_pricing?: TelegramPricing | null; }; - const amount = data.telegram_pricing?.amount_usdc || "10"; - setPaymentInfo( - data.telegram_pricing?.is_group_member - ? `Telegram 群成员验证成功,当前会员价 ${amount}U。` - : `Telegram 已登录,但未检测到群成员身份,当前价格 ${amount}U。`, - ); + if (data.telegram_pricing?.is_group_member) { + const amount = data.telegram_pricing.amount_usdc || "5"; + setPaymentInfo(`Telegram 群成员验证成功,当前会员价 ${amount}U。`); + } await loadSnapshot(); await loadPaymentSnapshot(); } catch (error) { setPaymentError(normalizePaymentError(error).message); - } finally { - setTelegramLoginBusy(false); } - }; - const script = document.createElement("script"); - script.src = "https://telegram.org/js/telegram-widget.js?22"; - script.async = true; - script.setAttribute("data-telegram-login", TELEGRAM_LOGIN_BOT_USERNAME); - script.setAttribute("data-size", "large"); - script.setAttribute("data-radius", "12"); - script.setAttribute("data-userpic", "false"); - script.setAttribute("data-request-access", "write"); - script.setAttribute("data-onauth", `${callbackName}(user)`); - container.appendChild(script); - return () => { - container.innerHTML = ""; - }; - }, [ - buildAuthedHeaders, - isAuthenticated, - loadPaymentSnapshot, - loadSnapshot, - ]); + })(); + }, [isAuthenticated, buildAuthedHeaders, loadPaymentSnapshot, loadSnapshot]); const displayName = String(user?.user_metadata?.full_name || "").trim() || (email ? String(email).split("@")[0] : "") || @@ -3011,28 +2974,21 @@ export function AccountCenter() {

{copy.telegramHint}

-
-

- Telegram 群内价格 -

-

- 登录 Telegram 后由后端检查群成员身份:群成员 5U,非群成员 10U。 -

-
-
- {telegramLoginBusy ? ( - - 正在验证 Telegram... - - ) : null} - {backend?.telegram_pricing?.amount_usdc ? ( + {backend?.telegram_pricing?.is_group_member ? ( +
+

+ Telegram 群成员价格 +

+

+ 已验证群成员身份,当前会员价 {backend.telegram_pricing.amount_usdc ?? "5"}U。 +

+
- 当前价格: {backend.telegram_pricing.amount_usdc}U - {backend.telegram_pricing.is_group_member ? " · 群成员" : " · 普通价"} + 当前价格: {backend.telegram_pricing.amount_usdc ?? "5"}U · 群成员 - ) : null} +
-
+ ) : null} {/* Channel Upgrade / Migration Notice — Pro users only */} {isSubscribed ? ( diff --git a/src/bot/handlers/basic.py b/src/bot/handlers/basic.py index 47652a11..574a6b10 100644 --- a/src/bot/handlers/basic.py +++ b/src/bot/handlers/basic.py @@ -1,5 +1,6 @@ from __future__ import annotations +import os from typing import Any from typing import Callable @@ -7,6 +8,7 @@ from src.bot.command_parser import extract_command_name from src.bot.io_layer import BotIOLayer from src.bot.observability import CommandTrace from src.bot.runtime_coordinator import RuntimeStatus, render_runtime_status_html +from src.auth.telegram_group_pricing import TelegramGroupPricing, TELEGRAM_MEMBER_STATUSES _BASIC_COMMANDS = {"start", "help", "id", "top", "diag", "bind", "unbind"} _BASIC_COMMANDS = {"start", "help", "id", "top", "diag", "bind", "unbind", "markets"} @@ -140,27 +142,48 @@ class BasicCommandHandler: trace = CommandTrace("/bind", message) try: parts = (message.text or "").split(maxsplit=2) + user = message.from_user + + # No-args mode: generate a one-time bind token for group members if len(parts) < 2: + pricing = TelegramGroupPricing() + if not pricing.configured: + self.bot.reply_to( + message, + "⚠️ 机器人未配置群组定价,请联系管理员。", + ) + trace.set_status("error", "pricing_not_configured") + return + member_status = pricing.get_member_status(user.id) + if not member_status or member_status not in TELEGRAM_MEMBER_STATUSES: + self.bot.reply_to( + message, + "🔒 此功能仅限内部群成员使用。", + ) + trace.set_status("blocked", f"not_group_member:{member_status or 'none'}") + return + token = self.io_layer.db.create_bind_token(user.id, ttl_minutes=10) + app_url = str(os.getenv("POLYWEATHER_APP_URL") or "https://polyweather-pro.vercel.app").rstrip("/") + bind_url = f"{app_url}/account?bind_token={token}" self.bot.reply_to( message, ( - "❌ 用法:\n" - "/bind <supabase_user_id> [email]\n\n" - "示例:\n" - "/bind 11111111-2222-3333-4444-555555555555 user@example.com" + "🔗 点击以下链接绑定网页账户(10 分钟内有效):\n" + f"{bind_url}\n\n" + "打开链接后系统将自动验证群成员身份并绑定 Telegram。" ), - parse_mode="HTML", + disable_web_page_preview=False, ) - trace.set_status("bad_request", "missing_supabase_user_id") + trace.set_status("ok", "token_generated") return + # Args mode: manual bind with supabase_user_id supabase_user_id = str(parts[1] or "").strip() if len(supabase_user_id) < 8: self.bot.reply_to(message, "❌ supabase_user_id 格式不正确。") trace.set_status("bad_request", "invalid_supabase_user_id") return supabase_email = str(parts[2] or "").strip() if len(parts) >= 3 else "" - user = message.from_user self.io_layer.db.upsert_user(user.id, self.io_layer.display_name(user)) result = self.io_layer.db.bind_supabase_identity( telegram_id=user.id, diff --git a/src/database/db_manager.py b/src/database/db_manager.py index e3f29cf7..dd8a20e1 100644 --- a/src/database/db_manager.py +++ b/src/database/db_manager.py @@ -2,6 +2,7 @@ import sqlite3 import os import hashlib import json +import secrets import threading from datetime import datetime, timedelta from typing import Optional, Dict, Any, List, Set @@ -364,6 +365,17 @@ class DBManager: conn.execute( "CREATE INDEX IF NOT EXISTS idx_supabase_bindings_telegram_id ON supabase_bindings(telegram_id)" ) + conn.execute(""" + CREATE TABLE IF NOT EXISTS telegram_bind_tokens ( + token TEXT PRIMARY KEY, + telegram_id INTEGER NOT NULL, + created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, + expires_at TIMESTAMP NOT NULL + ) + """) + conn.execute( + "CREATE INDEX IF NOT EXISTS idx_telegram_bind_tokens_expires ON telegram_bind_tokens(expires_at)" + ) conn.execute(""" CREATE TABLE IF NOT EXISTS airport_obs_log ( id INTEGER PRIMARY KEY AUTOINCREMENT, @@ -1368,6 +1380,58 @@ class DBManager: ) return {"ok": True, "reason": "unbound", "previous_supabase_user_id": current_uid} + def create_bind_token(self, telegram_id: int, ttl_minutes: int = 10) -> str: + token = secrets.token_urlsafe(16) + now = datetime.now() + expires_at = now + timedelta(minutes=max(1, int(ttl_minutes))) + with self._get_connection() as conn: + conn.execute( + """ + DELETE FROM telegram_bind_tokens + WHERE telegram_id = ? OR expires_at < ? + """, + (int(telegram_id), now.isoformat()), + ) + conn.execute( + """ + INSERT INTO telegram_bind_tokens (token, telegram_id, expires_at) + VALUES (?, ?, ?) + """, + (token, int(telegram_id), expires_at.isoformat()), + ) + conn.commit() + return token + + def consume_bind_token(self, token: str) -> Optional[int]: + token = str(token or "").strip() + if not token: + return None + now = datetime.now() + with self._get_connection() as conn: + conn.row_factory = sqlite3.Row + row = conn.execute( + """ + SELECT telegram_id, expires_at + FROM telegram_bind_tokens + WHERE token = ? + LIMIT 1 + """, + (token,), + ).fetchone() + if not row: + return None + try: + expires_at = datetime.fromisoformat(row["expires_at"]) + except Exception: + expires_at = now + if now > expires_at: + conn.execute("DELETE FROM telegram_bind_tokens WHERE token = ?", (token,)) + conn.commit() + return None + conn.execute("DELETE FROM telegram_bind_tokens WHERE token = ?", (token,)) + conn.commit() + return int(row["telegram_id"]) + def add_message_activity( self, telegram_id: int, diff --git a/web/core.py b/web/core.py index 1b5a68aa..d6e0b702 100644 --- a/web/core.py +++ b/web/core.py @@ -411,6 +411,10 @@ class TelegramLoginRequest(BaseModel): hash: str = Field(..., min_length=10) +class TelegramBindTokenRequest(BaseModel): + token: str = Field(..., min_length=8) + + class AnalyticsEventRequest(BaseModel): event_type: str = Field(..., min_length=3, max_length=64) client_id: Optional[str] = Field(default=None, max_length=128) diff --git a/web/routers/auth.py b/web/routers/auth.py index 395722cb..7df41dc5 100644 --- a/web/routers/auth.py +++ b/web/routers/auth.py @@ -2,8 +2,12 @@ from fastapi import APIRouter, Request -from web.core import TelegramLoginRequest -from web.services.auth_api import get_auth_me_payload, login_with_telegram +from web.core import TelegramBindTokenRequest, TelegramLoginRequest +from web.services.auth_api import ( + bind_telegram_by_token, + get_auth_me_payload, + login_with_telegram, +) router = APIRouter(tags=["auth"]) @@ -16,3 +20,8 @@ async def auth_me(request: Request): @router.post("/api/auth/telegram/login") async def auth_telegram_login(request: Request, body: TelegramLoginRequest): return login_with_telegram(request, body) + + +@router.post("/api/auth/telegram/bind-by-token") +async def auth_telegram_bind_by_token(request: Request, body: TelegramBindTokenRequest): + return bind_telegram_by_token(request, body) diff --git a/web/services/auth_api.py b/web/services/auth_api.py index 7a443d33..00e1864f 100644 --- a/web/services/auth_api.py +++ b/web/services/auth_api.py @@ -157,3 +157,43 @@ def login_with_telegram(request: Request, body: TelegramLoginRequest) -> Dict[st "binding": bind_result, "telegram_pricing": price, } + + +def bind_telegram_by_token(request: Request, body) -> Dict[str, Any]: + """Bind Telegram identity using a one-time token from the bot /bind command.""" + legacy_routes._assert_entitlement(request) + identity = legacy_routes._require_supabase_identity(request) + + token = str(getattr(body, "token", "") or "").strip() + if not token: + raise HTTPException(status_code=400, detail="bind_token is required") + + db = DBManager() + telegram_id = db.consume_bind_token(token) + if telegram_id is None: + raise HTTPException(status_code=400, detail="invalid or expired bind token") + + pricing = TelegramGroupPricing() + member_status = pricing.get_member_status(telegram_id) if pricing.configured else None + if not member_status or member_status not in ("creator", "administrator", "member"): + raise HTTPException(status_code=403, detail="not a group member") + + db.upsert_user(telegram_id, "") + bind_result = db.bind_supabase_identity( + telegram_id=telegram_id, + supabase_user_id=identity["user_id"], + supabase_email=identity.get("email") or "", + ) + if not bind_result.get("ok"): + raise HTTPException( + status_code=409, + detail=str(bind_result.get("reason") or "telegram bind failed"), + ) + + price = pricing.resolve_price_for_telegram_id(telegram_id) + return { + "ok": True, + "telegram_id": telegram_id, + "binding": bind_result, + "telegram_pricing": price, + }