feat: expand market discovery and auction workflow

This commit is contained in:
leefer
2026-07-24 17:32:28 +08:00
parent fde2728a86
commit 2d2a3aa5e5
46 changed files with 32992 additions and 168 deletions
+235 -65
View File
@@ -9,7 +9,7 @@ import re
import secrets
import threading
import time
from datetime import date, datetime, timedelta, timezone
from datetime import date, datetime, time as dt_time, timedelta, timezone
from http import HTTPStatus
from http.cookies import SimpleCookie
from http.server import BaseHTTPRequestHandler, ThreadingHTTPServer
@@ -19,6 +19,7 @@ from urllib.parse import parse_qs, unquote, urlparse
from alert_service import AlertService
from assistant_agent import ReviewAssistantError, stream_review_assistant
from api_access import required_role
from chart_data_provider import ChartDataError, EastmoneyChartClient
from app_config import (
DATA_DIR,
MENTOR_SKILLS_DIR,
@@ -51,6 +52,7 @@ from heaven_engine import (
)
from llm_strategy import LLMCompilerError, compile_strategy_with_llm, test_llm_connection
from mentor_agent import MentorAgentError, MentorSkillRegistry, stream_with_mentor
from market_insights import MarketInsightsService
from realtime_aggregator import WebRealtimeAggregator
from screener import (
FACTOR_FIELDS,
@@ -138,6 +140,7 @@ class DashboardService:
self.trade_journal = TradeJournalService(self.database)
self.mentor_skills = MentorSkillRegistry(MENTOR_SKILLS_DIR, PRIVATE_MENTOR_SKILLS_DIR)
self.realtime_aggregator = WebRealtimeAggregator()
self.chart_data = EastmoneyChartClient()
self.screener.ensure_builtin_strategies()
self._background_stop = threading.Event()
self._background_thread = threading.Thread(
@@ -1066,10 +1069,30 @@ class DashboardService:
sector = validate_text(sector, "板块名称", 50)
return self.realtime_aggregator.health_snapshot(sector)
def _market_insights(self) -> MarketInsightsService:
if not self.configured:
raise ValueError("行情数据尚未配置。")
return MarketInsightsService(self.database, TushareClient(self.token))
def auction_center(self, trade_date: str, force: bool = False) -> dict[str, Any]:
return self._market_insights().auction_center(
normalize_date(trade_date), force, self.current_user_id
)
def theme_library(self, trade_date: str, force: bool = False) -> dict[str, Any]:
return self._market_insights().theme_library(normalize_date(trade_date), force)
def theme_detail(self, code: str, trade_date: str) -> dict[str, Any]:
return self._market_insights().theme_detail(code, normalize_date(trade_date))
def popularity(self, trade_date: str, force: bool = False) -> dict[str, Any]:
return self._market_insights().popularity(normalize_date(trade_date), force)
def screener_setup(self, trade_date: str) -> dict[str, Any]:
normalized_date = normalize_date(trade_date)
regime = self.screener.detect_regime(normalized_date)
factor_dates = self.database.factor_dates(normalized_date, 100)
auction_dates = self.database.auction_factor_dates(normalized_date, 100)
return {
"trade_date": normalized_date,
"regime": regime,
@@ -1081,6 +1104,8 @@ class DashboardService:
"start_date": factor_dates[0] if factor_dates else "",
"end_date": factor_dates[-1] if factor_dates else "",
"ready": len(factor_dates) >= 21,
"auction_date_count": len(auction_dates),
"auction_ready": bool(auction_dates and auction_dates[-1] == factor_dates[-1]) if factor_dates else False,
},
"llm": {
"configured": self.llm_configured,
@@ -3160,6 +3185,51 @@ class DashboardService:
raise ValueError("未找到对应的板块或题材。")
return self._ths_search_detail(basic, normalized_date)
def get_intraday_chart(
self, entity_type: str, identifier: str
) -> dict[str, Any]:
entity_type = str(entity_type or "").strip().lower()
identifier = str(identifier or "").strip().upper()
if entity_type == "stock":
code = validate_stock_code(identifier)
chart = self.chart_data.stock_intraday(code)
type_label = SEARCH_TYPE_LABELS["stock"]
elif entity_type == "index":
basic = next((item for item in SEARCH_INDEXES if item["id"] == identifier), None)
if not basic:
raise ValueError("暂不支持该指数分时行情。")
chart = self.chart_data.index_intraday(identifier)
type_label = SEARCH_TYPE_LABELS["index"]
elif entity_type in {"sector", "theme"}:
basic = next(
(
item for item in self._search_market_directory()
if item.get("id") == identifier and item.get("type") == entity_type
),
None,
)
if not basic:
raise ValueError("未找到对应的板块或题材。")
chart = self.chart_data.board_intraday(identifier, str(basic.get("name") or ""))
type_label = SEARCH_TYPE_LABELS[entity_type]
else:
raise ValueError("分时行情类型不支持。")
return {
"meta": {
"trade_date": str(chart.get("trade_date") or ""),
"previous_close": float(chart.get("previous_close") or 0),
},
"entity": {
"id": identifier,
"code": str(chart.get("code") or identifier),
"name": str(chart.get("name") or ""),
"type": entity_type,
"type_label": type_label,
},
"points": list(chart.get("points") or []),
}
def _ths_search_detail(
self, basic: dict[str, Any], trade_date: str
) -> dict[str, Any]:
@@ -3309,8 +3379,9 @@ class DashboardService:
if not force:
cached = self.database.get_data_snapshot("stock_detail", cache_key)
if cached and str((cached.get("meta") or {}).get("source") or "") != "demo":
cached["meta"] = {**cached.get("meta", {}), "cached": True}
return self._enrich_stock_detail(cached)
if not self._stock_detail_cache_needs_refresh(cached, normalized_date):
cached["meta"] = {**cached.get("meta", {}), "cached": True}
return self._prepare_stock_detail(cached, code, normalized_date)
name, sector = self._stock_identity(code, normalized_date)
source = "tushare"
@@ -3333,7 +3404,7 @@ class DashboardService:
"cached": True,
"notice": "最新行情暂不可用,已沿用最近真实收盘数据。",
}
return self._enrich_stock_detail(payload)
return self._prepare_stock_detail(payload, code, normalized_date)
else:
payload = self.database.get_latest_data_snapshot(
"stock_detail", f"{code}:", cache_key, exclude_source="demo"
@@ -3346,11 +3417,96 @@ class DashboardService:
"cached": True,
"notice": "公共行情尚未配置,已沿用最近真实收盘数据。",
}
return self._enrich_stock_detail(payload)
return self._prepare_stock_detail(payload, code, normalized_date)
payload["meta"]["source"] = source
payload["meta"]["cached"] = False
self.database.save_data_snapshot("stock_detail", cache_key, source, payload)
return self._enrich_stock_detail(payload)
return self._prepare_stock_detail(payload, code, normalized_date)
@staticmethod
def _stock_detail_bar_date(payload: dict[str, Any]) -> str:
prices = list(payload.get("prices") or [])
return str((prices[-1] if prices else {}).get("trade_date") or "").replace("-", "")
def _stock_detail_cache_needs_refresh(
self, payload: dict[str, Any], requested_date: str
) -> bool:
now = datetime.now().astimezone()
return (
requested_date == now.strftime("%Y%m%d")
and now.time().replace(tzinfo=None) >= dt_time(15, 0)
and self._stock_detail_bar_date(payload) < requested_date
)
def _prepare_stock_detail(
self, payload: dict[str, Any], code: str, requested_date: str
) -> dict[str, Any]:
result = copy.deepcopy(payload)
actual_date = self._stock_detail_bar_date(result)
if actual_date:
result["meta"] = {
**(result.get("meta") or {}),
"trade_date": f"{actual_date[:4]}-{actual_date[4:6]}-{actual_date[6:]}",
}
if self.configured:
client = TushareClient(self.token)
now = datetime.now().astimezone()
today = now.strftime("%Y%m%d")
should_merge = (
requested_date == today
and actual_date < today
and now.time().replace(tzinfo=None) >= dt_time(9, 15)
)
if should_merge:
try:
resolved_date, _ = client.resolve_trade_context(requested_date)
if resolved_date == today:
quote = client.realtime_stock_quote(tushare_code(code), requested_date)
self._merge_realtime_stock_detail(result, quote, requested_date)
except TushareError:
pass
return self._enrich_stock_detail(result)
@staticmethod
def _merge_realtime_stock_detail(
payload: dict[str, Any], quote: dict[str, Any], trade_date: str
) -> None:
display_date = f"{trade_date[:4]}-{trade_date[4:6]}-{trade_date[6:]}"
realtime_bar = {
"trade_date": display_date,
"open": quote["open"],
"high": quote["high"],
"low": quote["low"],
"close": quote["price"],
"change": quote["change"],
"volume": quote["volume"] / 100,
"amount_billion": quote["amount_billion"],
"realtime": True,
}
prices = list(payload.get("prices") or [])
if prices and str(prices[-1].get("trade_date") or "").replace("-", "") == trade_date:
prices[-1] = realtime_bar
else:
prices.append(realtime_bar)
payload["prices"] = prices[-90:]
stock = dict(payload.get("stock") or {})
stock.update(
{
"name": quote["name"],
"industry": quote["sector"],
"price": quote["price"],
"change": quote["change"],
"amount_billion": quote["amount_billion"],
"turnover_rate": quote["turnover_rate"],
}
)
payload["stock"] = stock
payload["meta"] = {
**(payload.get("meta") or {}),
"trade_date": display_date,
"realtime": True,
"updated_at": datetime.now().astimezone().isoformat(timespec="seconds"),
}
def get_stock_preview(
self, code: str, trade_date: str, force: bool = False
@@ -3359,75 +3515,30 @@ class DashboardService:
detail = self.get_stock_detail(code, trade_date, force)
detail_meta = detail.get("meta") or {}
resolved_date = str(detail_meta.get("trade_date") or trade_date)
compact_date = normalize_date(resolved_date)
intraday_points: list[dict[str, Any]] = []
intraday_status = "unavailable"
intraday_notice = "分时行情暂不可用。"
if self.configured:
cache_key = f"{code}:{compact_date}"
cached = None if force else self.database.get_data_snapshot("stock_intraday", cache_key)
if cached and cached.get("points"):
intraday_points = list(cached["points"])
intraday_trade_date = ""
intraday_previous_close = 0.0
try:
intraday = self.chart_data.stock_intraday(code)
intraday_points = list(intraday.get("points") or [])
intraday_trade_date = str(intraday.get("trade_date") or "")
intraday_previous_close = float(intraday.get("previous_close") or 0)
if intraday_points:
intraday_status = "available"
intraday_notice = ""
else:
try:
intraday = TushareClient(self.token).stock_intraday(
tushare_code(code), compact_date
)
intraday_points = list(intraday.get("points") or [])
if intraday_points:
intraday_status = "available"
intraday_notice = ""
self.database.save_data_snapshot(
"stock_intraday", cache_key, "tushare", intraday
)
else:
intraday_status = "empty"
intraday_notice = "该交易日暂无分时数据。"
except TushareError as exc:
intraday_status = "unavailable"
intraday_notice = "分时行情暂不可用,请稍后重试。"
intraday_status = "empty"
intraday_notice = "最近交易日暂无分时数据。"
except ChartDataError:
intraday_status = "unavailable"
intraday_notice = "分时行情暂不可用,请稍后重试。"
prices = list(detail.get("prices") or [])[-60:]
stock = dict(detail.get("stock") or {"code": code})
realtime = False
if self.configured and compact_date == date.today().strftime("%Y%m%d"):
try:
quote = TushareClient(self.token).realtime_stock_quote(
tushare_code(code),
compact_date,
)
realtime_bar = {
"trade_date": f"{compact_date[:4]}-{compact_date[4:6]}-{compact_date[6:]}",
"open": quote["open"],
"high": quote["high"],
"low": quote["low"],
"close": quote["price"],
"change": quote["change"],
"volume": quote["volume"] / 100,
"amount_billion": quote["amount_billion"],
"realtime": True,
}
if prices and str(prices[-1].get("trade_date") or "").replace("-", "") == compact_date:
prices[-1] = realtime_bar
else:
prices.append(realtime_bar)
prices = prices[-60:]
stock.update(
{
"name": quote["name"],
"industry": quote["sector"],
"price": quote["price"],
"change": quote["change"],
"amount_billion": quote["amount_billion"],
"turnover_rate": quote["turnover_rate"],
}
)
realtime = True
except TushareError:
realtime = False
realtime = bool(detail_meta.get("realtime"))
return {
"meta": {
"trade_date": resolved_date,
@@ -3435,6 +3546,8 @@ class DashboardService:
"notice": detail_meta.get("notice") or "",
"intraday_status": intraday_status,
"intraday_notice": intraday_notice,
"intraday_trade_date": intraday_trade_date,
"intraday_previous_close": intraday_previous_close,
"realtime": realtime,
"refresh_interval_seconds": 10 if realtime else 0,
},
@@ -3751,6 +3864,54 @@ class RequestHandler(BaseHTTPRequestHandler):
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:
@@ -3814,6 +3975,15 @@ class RequestHandler(BaseHTTPRequestHandler):
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)