Files
xiaobaifupan/app/backend/application.py
T

1093 lines
47 KiB
Python

from __future__ import annotations
import json
import re
import secrets
import threading
import time
from datetime import date, datetime, time as dt_time, timedelta, timezone
from http import HTTPStatus
from http.server import BaseHTTPRequestHandler
from typing import Any
from urllib.parse import parse_qs, unquote, urlparse
from api_access import ROUTES
from backend.bootstrap.container import build_application_container
from backend.bootstrap.settings import load_runtime_settings
from backend.http import HttpTransportMixin
from backend.llm import LLMGateway, LLMGatewayError
from backend.llm.http import LLMHttpMixin
from backend.llm.service import LLMServiceMixin
from backend.features.market import ChartDataError, MarketServiceMixin
from backend.features.heaven import HeavenHttpMixin, HeavenServiceMixin, build_personal_field
from backend.features.alerts import AlertHttpMixin, AlertServiceMixin
from backend.features.review import ReviewHttpMixin, ReviewServiceMixin
from backend.bootstrap.config import (
DATA_DIR,
MENTOR_SKILLS_DIR,
PRIVATE_MENTOR_SKILLS_DIR,
TOKEN_PATTERN,
normalize_date,
validate_text,
)
from database import ReviewDatabase
from backend.features.accounts.http import AccountHttpMixin
from backend.features.accounts.security import SecretVault
from backend.features.accounts.service import AccountService
from backend.features.auction import AuctionServiceMixin
from backend.features.dragon_tiger import DragonTigerServiceMixin
from backend.features.mentor import MentorHttpMixin, MentorServiceMixin
from backend.features.pools import PoolServiceMixin
from backend.features.popularity import PopularityServiceMixin
from backend.features.rotation import RotationServiceMixin
from backend.features.screener.service import (
SCREENER_LIBRARY_VERSION,
ScreenerServiceMixin,
automatic_screener_jobs,
)
from backend.features.sentiment import SentimentServiceMixin
from backend.features.system import SystemHttpMixin
from backend.features.themes import ThemeServiceMixin
from backend.data.providers.tushare_client import TushareError
LEGACY_SECRET_KEYS = {
"TUSHARE_TOKEN",
"IFIND_REFRESH_TOKEN",
"IFIND_ACCESS_TOKEN",
"LLM_API_KEY",
"LLM_BASE_URL",
"LLM_MODEL",
"LLM_PRIMARY_API_KEY",
"LLM_PRIMARY_BASE_URL",
"LLM_PRIMARY_MODEL",
"LLM_FALLBACK_API_KEY",
"LLM_FALLBACK_BASE_URL",
"LLM_FALLBACK_MODEL",
}
class DashboardService(
MarketServiceMixin,
SentimentServiceMixin,
PoolServiceMixin,
RotationServiceMixin,
AuctionServiceMixin,
ThemeServiceMixin,
PopularityServiceMixin,
DragonTigerServiceMixin,
ScreenerServiceMixin,
MentorServiceMixin,
HeavenServiceMixin,
AlertServiceMixin,
ReviewServiceMixin,
LLMServiceMixin,
):
def __init__(self) -> None:
runtime = load_runtime_settings()
self.vault = SecretVault(runtime.encryption_key)
self.database = ReviewDatabase(DATA_DIR / "review.db")
self.sync_lock = threading.Lock()
self.auth_lock = threading.Lock()
self.system_lock = threading.Lock()
self.auto_screener_lock = threading.Lock()
self._auto_screener_last_attempt: dict[str, datetime] = {}
self._ifind_event_lock = threading.Lock()
self._request_context = threading.local()
self.accounts = AccountService(
database=self.database,
vault=self.vault,
current_user_supplier=lambda: self.current_user_id,
access_supplier=lambda: getattr(self._request_context, "access", {}),
bind_user=self.bind_user,
personal_field_builder=build_personal_field,
auth_lock=self.auth_lock,
)
self._system_credentials = self._load_system_credentials(runtime.initial_credentials)
self.container = build_application_container(
self.database,
self._system_credentials,
MENTOR_SKILLS_DIR,
PRIVATE_MENTOR_SKILLS_DIR,
lambda: self.token,
)
self.data_gateway = self.container.data_gateway
self.ifind = self.container.ifind
self.screener = self.container.screener
self.strategy_tracking = self.container.strategy_tracking
self.alert_service = self.container.alert_service
self.trade_journal = self.container.trade_journal
self.mentor_skills = self.container.mentor_skills
self.realtime_aggregator = self.container.realtime_aggregator
self.chart_data = self.container.chart_data
self.jobs = self.container.jobs
self.llm_gateway = LLMGateway(
database=self.database,
user_id_supplier=lambda: self.current_user_id,
membership_supplier=self.membership,
settings_supplier=lambda: self._system_credentials,
profile_supplier=self._resolved_llm_profile,
)
self.screener.ensure_builtin_strategies()
def start_background_jobs(self) -> threading.Thread:
return self.jobs.start_scheduler(
self._background_refresh_tick,
interval_seconds=5,
initial_delay_seconds=3,
)
def stop_background_jobs(self, timeout_seconds: float = 5) -> bool:
scheduler_stopped = self.jobs.stop_scheduler(timeout_seconds)
workers_stopped = self.jobs.wait_for_idle(timeout_seconds)
return scheduler_stopped and workers_stopped
def _load_system_credentials(self, environment: dict[str, str]) -> dict[str, Any]:
encrypted = self.database.get_system_setting("credentials")
current = self.vault.decrypt_json(encrypted) if encrypted else {}
changed = False
first_user_id = self.database.first_user_id()
first_personal: dict[str, Any] = {}
if first_user_id:
first_encrypted = self.database.get_user_credentials(first_user_id)
first_personal = self.vault.decrypt_json(first_encrypted) if first_encrypted else {}
defaults = {
"tushare_token": environment.get("tushare_token") or first_personal.get("tushare_token") or "",
"ifind_refresh_token": environment.get("ifind_refresh_token") or "",
"ifind_access_token": environment.get("ifind_access_token") or "",
"platform_llm_primary_api_key": environment.get("platform_llm_primary_api_key") or first_personal.get("llm_primary_api_key") or "",
"platform_llm_primary_base_url": environment.get("platform_llm_primary_base_url") or first_personal.get("llm_primary_base_url") or "https://api.openai.com/v1",
"platform_llm_primary_model": environment.get("platform_llm_primary_model") or first_personal.get("llm_primary_model") or "",
"platform_llm_fallback_api_key": environment.get("platform_llm_fallback_api_key") or first_personal.get("llm_fallback_api_key") or "",
"platform_llm_fallback_base_url": environment.get("platform_llm_fallback_base_url") or first_personal.get("llm_fallback_base_url") or "",
"platform_llm_fallback_model": environment.get("platform_llm_fallback_model") or first_personal.get("llm_fallback_model") or "",
"member_daily_limit": 50,
"background_refresh_enabled": True,
}
for key, value in defaults.items():
if key not in current:
current[key] = value
changed = True
if not isinstance(current.get("llm_models"), list):
migrated_models: list[dict[str, str]] = []
for role, label in (("primary", "原主模型"), ("fallback", "原辅助模型")):
profile = {
"api_key": str(current.get(f"platform_llm_{role}_api_key") or ""),
"base_url": str(current.get(f"platform_llm_{role}_base_url") or ""),
"model": str(current.get(f"platform_llm_{role}_model") or ""),
}
if profile["api_key"] or profile["model"]:
model_id = f"migrated-{role}"
migrated_models.append(
{"id": model_id, "name": label, **profile}
)
current[f"{role}_model_id"] = model_id
current["llm_models"] = migrated_models
current.setdefault("primary_model_id", "")
current.setdefault("fallback_model_id", "")
changed = True
if changed or not encrypted:
self.database.save_system_setting("credentials", self.vault.encrypt_json(current))
for row in self.database.list_user_credentials():
personal = self.vault.decrypt_json(str(row.get("encrypted_payload") or ""))
if "tushare_token" in personal:
personal.pop("tushare_token", None)
self.database.save_user_credentials(
int(row["user_id"]), self.vault.encrypt_json(personal)
)
return current
def _save_system_credentials(self, credentials: dict[str, Any]) -> None:
with self.system_lock:
self.database.save_system_setting("credentials", self.vault.encrypt_json(credentials))
self._system_credentials = dict(credentials)
if hasattr(self, "ifind"):
self.ifind.set_credentials(
str(credentials.get("ifind_refresh_token") or ""),
str(credentials.get("ifind_access_token") or ""),
)
@property
def configured(self) -> bool:
return bool(self.token)
def bind_user(self, user_id: int) -> None:
self._request_context.user_id = int(user_id)
encrypted = self.database.get_user_credentials(int(user_id))
self._request_context.credentials = self.vault.decrypt_json(encrypted) if encrypted else {}
self._request_context.access = self.database.user_access(int(user_id)) or {}
@property
def current_user_id(self) -> int:
user_id = getattr(self._request_context, "user_id", 0)
if not user_id:
raise ValueError("当前请求尚未绑定账号。")
return int(user_id)
def _credentials(self) -> dict[str, str]:
credentials = getattr(self._request_context, "credentials", {})
return {
"llm_primary_api_key": str(credentials.get("llm_primary_api_key") or ""),
"llm_primary_base_url": str(
credentials.get("llm_primary_base_url") or "https://api.openai.com/v1"
),
"llm_primary_model": str(credentials.get("llm_primary_model") or ""),
"llm_fallback_api_key": str(credentials.get("llm_fallback_api_key") or ""),
"llm_fallback_base_url": str(credentials.get("llm_fallback_base_url") or ""),
"llm_fallback_model": str(credentials.get("llm_fallback_model") or ""),
}
def _save_credentials(self, credentials: dict[str, str]) -> None:
self.database.save_user_credentials(
self.current_user_id,
self.vault.encrypt_json(credentials),
)
self._request_context.credentials = dict(credentials)
@property
def token(self) -> str:
return str(self._system_credentials.get("tushare_token") or "")
def membership(self) -> dict[str, Any]:
return self.accounts.membership()
def system_status(self) -> dict[str, Any]:
platform = self._platform_llm_profile()
model_pool = []
for item in self._system_credentials.get("llm_models") or []:
if not isinstance(item, dict):
continue
profile = {
"api_key": str(item.get("api_key") or ""),
"base_url": str(item.get("base_url") or ""),
"model": str(item.get("model") or ""),
}
model_pool.append(
{
"id": str(item.get("id") or ""),
"name": str(item.get("name") or ""),
"base_url": profile["base_url"],
"model": profile["model"],
"configured": self._profile_configured(profile),
}
)
return {
"data": {
"configured": self.configured,
"ifind": self.ifind.status(),
"background_refresh_enabled": bool(
self._system_credentials.get("background_refresh_enabled", True)
),
**self.database.status(),
"jobs": self.jobs.repository.recent(12),
},
"llm": {
"primary_configured": self._profile_configured(platform["primary"]),
"fallback_configured": self._profile_configured(platform["fallback"]),
"models": model_pool,
"primary_model_id": str(self._system_credentials.get("primary_model_id") or ""),
"fallback_model_id": str(self._system_credentials.get("fallback_model_id") or ""),
},
"membership": {
"member_daily_limit": max(
1, int(self._system_credentials.get("member_daily_limit") or 50)
)
},
}
def save_system_settings(self, payload: dict[str, Any]) -> dict[str, Any]:
current = dict(self._system_credentials)
token = str(payload.get("tushare_token") or current.get("tushare_token") or "").strip()
if token and not TOKEN_PATTERN.fullmatch(token):
raise ValueError("Tushare Token 格式不正确。")
ifind_refresh_token = str(
payload.get("ifind_refresh_token")
or current.get("ifind_refresh_token")
or ""
).strip()
if ifind_refresh_token and (
len(ifind_refresh_token) > 2048
or any(character.isspace() for character in ifind_refresh_token)
):
raise ValueError("iFinD Refresh Token 格式不正确。")
existing_models = {
str(item.get("id") or ""): item
for item in current.get("llm_models") or []
if isinstance(item, dict) and item.get("id")
}
raw_models = payload.get("models")
models: list[dict[str, str]] = []
if raw_models is not None:
if not isinstance(raw_models, list) or len(raw_models) > 20:
raise ValueError("模型池格式不正确,最多可保存 20 个模型。")
seen_ids: set[str] = set()
seen_names: set[str] = set()
for index, raw in enumerate(raw_models, start=1):
if not isinstance(raw, dict):
raise ValueError("模型池条目格式不正确。")
model_id = str(raw.get("id") or f"model-{secrets.token_hex(6)}").strip()
if not re.fullmatch(r"[A-Za-z0-9_-]{3,80}", model_id) or model_id in seen_ids:
raise ValueError("模型 ID 不正确或重复。")
name = validate_text(raw.get("name"), f"模型 {index} 名称", 50, required=True)
normalized_name = name.casefold()
if normalized_name in seen_names:
raise ValueError("模型名称不能重复。")
profile = self._validate_llm_profile(
raw,
existing_models.get(model_id) or {},
required=True,
label=name,
)
models.append({"id": model_id, "name": name, **profile})
seen_ids.add(model_id)
seen_names.add(normalized_name)
else:
models = [dict(item) for item in existing_models.values()]
model_ids = {item["id"] for item in models}
primary_model_id = str(
payload.get("primary_model_id", current.get("primary_model_id") or "") or ""
).strip()
fallback_model_id = str(
payload.get("fallback_model_id", current.get("fallback_model_id") or "") or ""
).strip()
if models and primary_model_id not in model_ids:
raise ValueError("请从模型池选择主模型。")
if not models:
primary_model_id = ""
fallback_model_id = ""
if fallback_model_id and fallback_model_id not in model_ids:
raise ValueError("辅助模型不在模型池中。")
if fallback_model_id and fallback_model_id == primary_model_id:
raise ValueError("主模型与辅助模型不能相同。")
try:
daily_limit = max(
1,
min(
1000,
int(payload.get("member_daily_limit", current.get("member_daily_limit") or 50)),
),
)
except (TypeError, ValueError) as exc:
raise ValueError("会员每日额度应为 1 至 1000。") from exc
current.update(
{
"tushare_token": token,
"ifind_refresh_token": ifind_refresh_token,
"llm_models": models,
"primary_model_id": primary_model_id,
"fallback_model_id": fallback_model_id,
"member_daily_limit": daily_limit,
"background_refresh_enabled": bool(
payload.get(
"background_refresh_enabled",
current.get("background_refresh_enabled", True),
)
),
}
)
self._save_system_credentials(current)
return self.system_status()
def admin_users(self) -> list[dict[str, Any]]:
return self.accounts.admin_users(self._platform_usage_today_for_user)
def update_membership(self, payload: dict[str, Any]) -> None:
self.accounts.update_membership(payload)
def request_background_sync(self, trade_date: str) -> bool:
normalized = normalize_date(trade_date)
key = f"manual:{normalized}:{time.time_ns()}"
return self.jobs.submit(
"market.refresh",
key,
lambda: self.sync_dashboard(normalized),
{"trade_date": normalized, "trigger": "administrator"},
)
def _background_refresh_tick(self) -> None:
if not (
self.configured
and self._system_credentials.get("background_refresh_enabled", True)
):
return
today = date.today().strftime("%Y%m%d")
snapshot = self.database.get_snapshot(today) or {}
if self._realtime_snapshot_due(today, snapshot):
bucket = int(time.time() // 5)
self.jobs.submit(
"market.refresh",
f"realtime:{today}:{bucket}",
lambda: self.sync_dashboard(today),
{"trade_date": today, "trigger": "realtime-poll"},
)
self._schedule_automatic_screeners(today, snapshot)
def register_account(self, username: str, password: str) -> dict[str, Any]:
return self.accounts.register(username, password)
def login_account(self, username: str, password: str) -> dict[str, Any]:
return self.accounts.login(username, password)
def change_password(self, current_password: str, new_password: str) -> None:
self.accounts.change_password(current_password, new_password)
def create_account_session(self, user: dict[str, Any]) -> dict[str, Any]:
return self.accounts.create_session(user)
@staticmethod
def _validate_account_input(username: str, password: str) -> None:
AccountService.validate_input(username, password)
def save_birth_profile(self, payload: dict[str, Any]) -> dict[str, Any]:
return self.accounts.save_birth_profile(payload)
def stored_birth_profile(self) -> dict[str, str] | None:
return self.accounts.stored_birth_profile()
def account_personal_field(
self,
current_date: str,
current_field: dict[str, Any],
public: bool = False,
) -> dict[str, Any] | None:
return self.accounts.personal_field(current_date, current_field, public)
@staticmethod
def _public_personal_profile(personal: dict[str, Any]) -> dict[str, Any]:
return AccountService.public_personal_profile(personal)
def status(self) -> dict[str, Any]:
llm_access = self.llm_access_status()
return {
"configured": self.configured,
"mode": "tushare" if self.configured else "unavailable",
"llm_configured": self.llm_configured,
"llm_model": self.llm_primary_model if self.llm_configured else "",
"llm_fallback_configured": self.llm_fallback_configured,
"llm_fallback_model": self.llm_fallback_model if self.llm_fallback_configured else "",
"llm_access": llm_access,
"birth_profile_configured": bool(self.stored_birth_profile()),
"birth_profile": self.stored_birth_profile(),
**self.database.status(),
}
SERVICE = DashboardService()
PUBLIC_POST_HANDLERS = {
"/api/auth/register": "auth_register",
"/api/auth/login": "auth_login",
}
AUTHENTICATED_POST_HANDLERS = {
"/api/auth/logout": "auth_logout",
"/api/account/birth-profile": "save_birth_profile",
"/api/account/password": "change_password",
"/api/alerts": "save_alert",
"/api/trades": "save_trade_entry",
"/api/assistant/chat": "stream_assistant_chat",
"/api/admin/settings": "save_system_settings",
"/api/admin/settings/test": "test_system_llm_settings",
"/api/admin/membership": "save_membership",
"/api/admin/refresh": "start_background_refresh",
"/api/watchlist": "save_watchlist",
"/api/notes": "save_note",
"/api/reasons": "save_reason",
"/api/seat-aliases": "save_seat_alias",
"/api/heaven/sector-phases": "save_sector_phase_override",
"/api/backfill": "backfill_data",
"/api/screener/sync": "sync_screener_data",
"/api/screener/compile": "compile_screener_strategy",
"/api/screener/strategies": "save_screener_strategy",
"/api/screener/run": "run_screener",
"/api/screener/tracking/refresh": "refresh_screener_tracking",
"/api/mentors/chat": "stream_mentor_chat",
"/api/heaven/hexagram": "heaven_hexagram",
"/api/heaven/personal": "heaven_personal",
"/api/heaven/interpret": "heaven_interpret",
}
class RequestHandler(
AccountHttpMixin,
SystemHttpMixin,
MentorHttpMixin,
HeavenHttpMixin,
AlertHttpMixin,
ReviewHttpMixin,
LLMHttpMixin,
HttpTransportMixin,
BaseHTTPRequestHandler,
):
server_version = "XiaobaiReviewWeb/0.8"
application_service = SERVICE
route_registry = ROUTES
def _dispatch_named_handler(self, path: str, handlers: dict[str, str]) -> bool:
handler_name = handlers.get(path)
if handler_name is None:
return False
getattr(self, handler_name)()
return True
def do_GET(self) -> None:
parsed = urlparse(self.path)
if parsed.path == "/api/health":
self.send_json(
{
"ok": True,
"storage": "sqlite",
"account_required": True,
"time": datetime.now().astimezone().isoformat(timespec="seconds"),
}
)
return
if parsed.path == "/api/auth/me":
self.auth_me()
return
if parsed.path.startswith("/api/"):
if not self.require_auth():
return
if not self.require_access("GET", parsed.path):
return
if parsed.path == "/api/admin/settings":
self.send_json(
{"ok": True, **SERVICE.system_status(), "users": SERVICE.admin_users()}
)
return
if parsed.path == "/api/account/status":
self.send_json({"ok": True, **SERVICE.status()})
return
if parsed.path == "/api/alerts":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.alert_center(
query.get("status", ["all"])[0],
query.get("as_of", [date.today().isoformat()])[0],
)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/trades":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.trade_entries(
query.get("start_date", [""])[0],
query.get("end_date", [""])[0],
query.get("code", [""])[0],
)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/assistant/messages":
self.send_json({"items": SERVICE.assistant_messages()})
return
if parsed.path == "/api/dashboard":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
self.send_json(SERVICE.get_dashboard(trade_date, False))
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
except Exception as exc:
self.send_json({"error": f"数据加载失败:{exc}"}, HTTPStatus.INTERNAL_SERVER_ERROR)
return
if parsed.path == "/api/auction":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.auction_center(
query.get("trade_date", [date.today().isoformat()])[0],
query.get("force", ["0"])[0] == "1",
)
)
except (ValueError, TushareError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/themes":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.theme_library(
query.get("trade_date", [date.today().isoformat()])[0],
query.get("force", ["0"])[0] == "1",
)
)
except (ValueError, TushareError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/themes/detail":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.theme_detail(
query.get("code", [""])[0],
query.get("trade_date", [date.today().isoformat()])[0],
)
)
except (ValueError, TushareError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/popularity":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.popularity(
query.get("trade_date", [date.today().isoformat()])[0],
query.get("force", ["0"])[0] == "1",
)
)
except (ValueError, TushareError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/realtime-aggregate/health":
query = parse_qs(parsed.query)
try:
self.send_json(
{
"ok": True,
"aggregate": SERVICE.realtime_aggregate_health(
query.get("sector", [""])[0]
),
}
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/sentiment/history":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
limit = int(query.get("limit", ["20"])[0])
self.send_json(SERVICE.sentiment_history(trade_date, limit))
except (TypeError, ValueError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/rotation/history":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
self.send_json(SERVICE.rotation_history(trade_date, 9))
except (TypeError, ValueError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/rotation/members":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.rotation_sector_members(
query.get("trade_date", [date.today().isoformat()])[0],
query.get("sector", [""])[0],
)
)
except (TypeError, ValueError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/dragon-tiger":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
force = query.get("force", ["0"])[0] == "1"
try:
self.send_json(SERVICE.get_dragon_tiger(trade_date, force))
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/dragon-tiger/profiles":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.get_hot_money_profiles(
query.get("force", ["0"])[0] == "1"
)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/search":
query = parse_qs(parsed.query)
search_query = query.get("q", [""])[0]
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
self.send_json(SERVICE.search_entities(search_query, trade_date))
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/search/detail":
query = parse_qs(parsed.query)
entity_type = query.get("type", [""])[0]
identifier = query.get("id", [""])[0]
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
self.send_json(
SERVICE.get_search_detail(entity_type, identifier, trade_date)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
except TushareError as exc:
self.send_json({"error": f"行情加载失败:{exc}"}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/chart/intraday":
query = parse_qs(parsed.query)
entity_type = query.get("type", [""])[0]
identifier = query.get("id", [""])[0]
try:
self.send_json(SERVICE.get_intraday_chart(entity_type, identifier))
except (ValueError, ChartDataError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
stock_preview_match = re.fullmatch(r"/api/stock/(\d{6})/preview", parsed.path)
if stock_preview_match:
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
force = query.get("force", ["0"])[0] == "1"
try:
self.send_json(
SERVICE.get_stock_preview(stock_preview_match.group(1), trade_date, force)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
stock_match = re.fullmatch(r"/api/stock/(\d{6})", parsed.path)
if stock_match:
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
force = query.get("force", ["0"])[0] == "1"
try:
self.send_json(SERVICE.get_stock_detail(stock_match.group(1), trade_date, force))
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/watchlist":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.review_watchlist(
query.get("trade_date", [date.today().isoformat()])[0]
)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/notes":
query = parse_qs(parsed.query)
code = query.get("code", [""])[0]
trade_date = query.get("trade_date", [""])[0].replace("-", "")
scope = query.get("scope", ["all"])[0]
if scope not in {"all", "daily", "stock"}:
self.send_json({"error": "复盘记录范围不支持。"}, HTTPStatus.BAD_REQUEST)
return
self.send_json(
{
"items": SERVICE.database.list_notes(
SERVICE.current_user_id, code, trade_date, scope
)
}
)
return
if parsed.path == "/api/seat-aliases":
self.send_json({"items": SERVICE.database.list_seat_aliases()})
return
if parsed.path == "/api/screener/setup":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
self.send_json(SERVICE.screener_setup(trade_date))
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/screener/tracking":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.screener_tracking(int(query.get("limit", ["12"])[0]))
)
except (TypeError, ValueError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/mentors/setup":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
try:
self.send_json(SERVICE.mentor_setup(trade_date))
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/mentors/messages":
query = parse_qs(parsed.query)
try:
self.send_json(
{
"items": SERVICE.mentor_messages(
query.get("mentor_id", [""])[0],
query.get("trade_date", [date.today().isoformat()])[0],
)
}
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/heaven/readings":
query = parse_qs(parsed.query)
try:
self.send_json(
SERVICE.heaven_readings(
query.get("mode", [""])[0],
query.get("context_date", [""])[0],
int(query.get("limit", ["100"])[0]),
)
)
except (TypeError, ValueError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/heaven/setup":
query = parse_qs(parsed.query)
trade_date = query.get("trade_date", [date.today().isoformat()])[0]
sector_name = query.get("sector", [""])[0]
stock_code = query.get("stock_code", [""])[0]
manual_data = None
manual_text = query.get("manual_data", [""])[0]
if manual_text:
try:
manual_data = json.loads(manual_text)
except json.JSONDecodeError:
self.send_json({"error": "六爻补录数据格式不正确。"}, HTTPStatus.BAD_REQUEST)
return
try:
self.send_json(
SERVICE.heaven_setup(
trade_date,
sector_name,
stock_code,
manual_data,
)
)
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
self.serve_static(parsed.path)
def do_POST(self) -> None:
parsed = urlparse(self.path)
if self._dispatch_named_handler(parsed.path, PUBLIC_POST_HANDLERS):
return
if not self.require_auth() or not self.require_csrf():
return
if not self.require_access("POST", parsed.path):
return
if self._dispatch_named_handler(parsed.path, AUTHENTICATED_POST_HANDLERS):
return
alert_read_match = re.fullmatch(r"/api/alerts/(\d+)/read", parsed.path)
if alert_read_match:
self.send_json(
{"ok": True, **SERVICE.mark_alert_read(int(alert_read_match.group(1)))}
)
return
if parsed.path == "/api/alerts/read-all":
body = self.read_json_body(True)
self.send_json(
{"ok": True, **SERVICE.mark_all_alerts_read(str(body.get("as_of") or ""))}
)
return
if parsed.path == "/api/screener/tracking":
try:
result = SERVICE.add_screener_tracking(self.read_json_body())
self.send_json({"ok": True, **result})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
if parsed.path == "/api/mentors/preferences":
try:
result = SERVICE.save_mentor_preferences(self.read_json_body())
self.send_json({"ok": True, **result})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
self.send_json({"error": "Not found"}, HTTPStatus.NOT_FOUND)
def do_DELETE(self) -> None:
parsed = urlparse(self.path)
if not self.require_auth() or not self.require_csrf():
return
if not self.require_access("DELETE", parsed.path):
return
if parsed.path == "/api/account/birth-profile":
deleted = SERVICE.database.delete_user_birth_profile(SERVICE.current_user_id)
self.send_json({"ok": True, "deleted": deleted})
return
if parsed.path == "/api/assistant/messages":
deleted = SERVICE.clear_assistant_messages()
self.send_json({"ok": True, "deleted": deleted})
return
if parsed.path == "/api/mentors/messages":
query = parse_qs(parsed.query)
try:
deleted = SERVICE.clear_mentor_messages(
query.get("mentor_id", [""])[0],
query.get("trade_date", [date.today().isoformat()])[0],
)
self.send_json({"ok": True, "deleted": deleted})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
strategy_match = re.fullmatch(r"/api/screener/strategies/(\d+)", parsed.path)
if strategy_match:
try:
result = SERVICE.delete_screener_strategy(int(strategy_match.group(1)))
self.send_json({"ok": True, **result})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
return
tracking_match = re.fullmatch(r"/api/screener/tracking/(\d+)", parsed.path)
if tracking_match:
result = SERVICE.remove_screener_tracking(int(tracking_match.group(1)))
self.send_json({"ok": True, **result})
return
watchlist_match = re.fullmatch(r"/api/watchlist/(\d{6})", parsed.path)
if watchlist_match:
deleted = SERVICE.database.delete_watchlist(
SERVICE.current_user_id, watchlist_match.group(1)
)
self.send_json({"ok": True, "deleted": deleted})
return
note_match = re.fullmatch(r"/api/notes/(\d+)", parsed.path)
if note_match:
deleted = SERVICE.database.delete_note(
SERVICE.current_user_id, int(note_match.group(1))
)
self.send_json({"ok": True, "deleted": deleted})
return
alert_match = re.fullmatch(r"/api/alerts/(\d+)", parsed.path)
if alert_match:
self.send_json(
{"ok": True, **SERVICE.delete_alert(int(alert_match.group(1)))}
)
return
trade_match = re.fullmatch(r"/api/trades/(\d+)", parsed.path)
if trade_match:
self.send_json(
{"ok": True, **SERVICE.delete_trade_entry(int(trade_match.group(1)))}
)
return
heaven_reading_match = re.fullmatch(r"/api/heaven/readings/(\d+)", parsed.path)
if heaven_reading_match:
deleted = SERVICE.database.delete_heaven_reading(
SERVICE.current_user_id, int(heaven_reading_match.group(1))
)
self.send_json({"ok": True, "deleted": deleted})
return
sector_phase_match = re.fullmatch(r"/api/heaven/sector-phases/(.+)", parsed.path)
if sector_phase_match:
name = unquote(sector_phase_match.group(1)).strip()
deleted = SERVICE.database.delete_sector_phase_override(name)
self.send_json({"ok": True, "deleted": deleted})
return
self.send_json({"error": "Not found"}, HTTPStatus.NOT_FOUND)
def save_reason(self) -> None:
try:
body = self.read_json_body()
SERVICE.save_reason(
str(body.get("trade_date") or ""),
str(body.get("code") or ""),
str(body.get("reason") or ""),
)
self.send_json({"ok": True})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def save_seat_alias(self) -> None:
try:
body = self.read_json_body()
seat_name = validate_text(body.get("seat_name"), "席位名称", 200, required=True)
alias = validate_text(body.get("alias"), "席位别名", 50, required=True)
SERVICE.database.save_seat_alias(seat_name, alias)
self.send_json({"ok": True})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def save_sector_phase_override(self) -> None:
try:
body = self.read_json_body()
name = validate_text(body.get("name"), "行业或题材名称", 50, required=True)
element = str(body.get("element") or "").strip()
if element not in {"木", "火", "土", "金", "水"}:
raise ValueError("五行归类必须是木、火、土、金或水。")
SERVICE.database.save_sector_phase_override(name, element)
self.send_json({"ok": True})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def backfill_data(self) -> None:
try:
body = self.read_json_body()
results = SERVICE.backfill(
str(body.get("start_date") or ""),
str(body.get("end_date") or ""),
)
self.send_json({"ok": True, "results": results})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
except Exception as exc:
self.send_json({"error": f"历史回补失败:{exc}"}, HTTPStatus.INTERNAL_SERVER_ERROR)
def sync_screener_data(self) -> None:
try:
body = self.read_json_body()
result = SERVICE.sync_screener_data(
str(body.get("trade_date") or date.today().isoformat()),
int(body.get("lookback") or 45),
)
self.send_json({"ok": True, "result": result})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
except Exception as exc:
self.send_json({"error": f"因子数据同步失败:{exc}"}, HTTPStatus.INTERNAL_SERVER_ERROR)
def compile_screener_strategy(self) -> None:
try:
body = self.read_json_body()
result = SERVICE.compile_screener_strategy(
str(body.get("prompt") or ""), str(body.get("regime") or "")
)
self.send_json({"ok": True, "strategy": result})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def save_screener_strategy(self) -> None:
try:
body = self.read_json_body()
result = SERVICE.save_screener_strategy(body)
self.send_json({"ok": True, **result})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
def run_screener(self) -> None:
try:
body = self.read_json_body()
result = SERVICE.run_screener(body)
self.send_json({"ok": True, "result": result})
except ValueError as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
except Exception as exc:
self.send_json({"error": f"选股执行失败:{exc}"}, HTTPStatus.INTERNAL_SERVER_ERROR)
def refresh_screener_tracking(self) -> None:
try:
body = self.read_json_body(True)
trade_date = str(body.get("trade_date") or date.today().isoformat())
self.send_json({"ok": True, **SERVICE.refresh_screener_tracking(trade_date)})
except (ValueError, json.JSONDecodeError) as exc:
self.send_json({"error": str(exc)}, HTTPStatus.BAD_REQUEST)
except Exception as exc:
self.send_json({"error": f"跟踪刷新失败:{exc}"}, HTTPStatus.INTERNAL_SERVER_ERROR)