from __future__ import annotations from collections import Counter from datetime import datetime, time as dt_time, timedelta from typing import Any from backend.bootstrap.config import display_compact_date as _display_date from backend.data.numbers import finite_number as _number from backend.features.sentiment.engine import apply_sentiment_to_dashboard from backend.data.providers.tushare_helpers import ( _realtime_market_status, _trading_session_progress, _value_percentile, ) from backend.data.providers.tushare_transport import TushareError class DashboardMixin: def dashboard(self, requested_date: str) -> dict[str, Any]: trade_date, previous_trade_date = self.resolve_trade_context(requested_date) if self.should_use_realtime(requested_date, trade_date): return self._realtime_dashboard( requested_date, trade_date, previous_trade_date, ) daily = self._load_daily(trade_date) if ( not daily and requested_date == datetime.now().astimezone().strftime("%Y%m%d") and trade_date == requested_date and datetime.now().astimezone().time().replace(tzinfo=None) >= dt_time(9, 15) ): return self._realtime_dashboard( requested_date, trade_date, previous_trade_date, ) if not daily: raise TushareError(f"No daily data returned for {trade_date}") notices: list[str] = [] limit_data_source = "official" try: limit_rows = self._load_limit_lists(trade_date) previous_limit_rows = self._load_limit_type(previous_trade_date, "U") if not limit_rows: limit_data_source = "derived" notices.append("涨跌停高级接口当日数据尚未更新,已使用日线数据推算。") limit_rows = self._derive_limits(trade_date, daily) except TushareError as exc: limit_data_source = "derived" notices.append(f"涨跌停高级接口不可用,已使用日线数据推算:{exc}") limit_rows = self._derive_limits(trade_date, daily) previous_daily = self._load_daily(previous_trade_date) previous_limit_rows = [ row for row in self._derive_limits(previous_trade_date, previous_daily) if row.get("limit_type") == "U" ] up_rows = [row for row in limit_rows if row.get("limit_type") == "U"] down_rows = [row for row in limit_rows if row.get("limit_type") == "D"] broken_rows = [row for row in limit_rows if row.get("limit_type") == "Z"] limits = [self._normalize_limit(row, "涨停") for row in up_rows] broken = [self._normalize_limit(row, "炸板") for row in broken_rows] down_limits = [self._normalize_limit(row, "跌停") for row in down_rows] previous_limits = [self._normalize_limit(row, "涨停") for row in previous_limit_rows] yesterday_limits = _build_yesterday_performance( previous_limits, daily, limits, broken, down_limits, ) sectors = _build_sectors(limits) previous_sectors = _build_sectors(previous_limits) dashboard = { "meta": { "requested_date": _display_date(requested_date), "trade_date": _display_date(trade_date), "previous_trade_date": _display_date(previous_trade_date), "source": "tushare", "limit_data_source": limit_data_source, "updated_at": datetime.now().astimezone().isoformat(timespec="seconds"), "notice": ";".join(notices), }, "overview": _build_overview(daily, up_rows, down_rows, broken_rows), "limits": limits, "broken": broken, "down_limits": down_limits, "yesterday_limits": yesterday_limits, "limit_performance": _build_limit_performance(yesterday_limits), "ladders": _build_ladders(limits), "sectors": sectors, "sector_rotation": _build_sector_rotation(sectors, previous_sectors), } return apply_sentiment_to_dashboard(dashboard) @staticmethod def should_use_realtime(requested_date: str, trade_date: str) -> bool: """Use rt_k for today's open market until end-of-day datasets settle.""" now = datetime.now().astimezone() today = now.strftime("%Y%m%d") return ( requested_date == today and trade_date == today and dt_time(9, 15) <= now.time().replace(tzinfo=None) < dt_time(16, 30) ) def _realtime_dashboard( self, requested_date: str, trade_date: str, previous_trade_date: str, ) -> dict[str, Any]: reference = self._load_realtime_reference(trade_date, previous_trade_date) basic_rows = list(reference["basic_rows"]) codes = ",".join( str(row.get("ts_code") or "") for row in basic_rows if row.get("ts_code") ) if not codes: raise TushareError("No active stock codes available for rt_k") quotes = self.query("rt_k", {"ts_code": codes}) if not quotes: raise TushareError(f"No realtime data returned for {trade_date}") basic_map = {str(row.get("ts_code") or ""): row for row in basic_rows} daily: list[dict[str, Any]] = [] for quote in quotes: close = _number(quote.get("close")) previous_close = _number(quote.get("pre_close")) if close <= 0 or previous_close <= 0: continue basic = basic_map.get(str(quote.get("ts_code") or ""), {}) daily.append( { **quote, "trade_date": trade_date, "name": str(quote.get("name") or basic.get("name") or "--").strip(), "industry": basic.get("industry") or "其他", "pct_chg": round((close / previous_close - 1) * 100, 4), "amount_unit": "yuan", } ) with self._realtime_reference_lock: self._latest_realtime_market[trade_date] = { "rows": daily, "updated_at": datetime.now().astimezone().isoformat(timespec="seconds"), } if len(self._latest_realtime_market) > 3: oldest = next(iter(self._latest_realtime_market)) self._latest_realtime_market.pop(oldest, None) limit_rows = self._derive_limits( trade_date, daily, price_limits=list(reference["price_limits"]), basic_rows=basic_rows, previous_limit_rows=list(reference["previous_limit_rows"]), capital_rows=list(reference["capital_rows"]), ) previous_limit_rows = list(reference["previous_limit_rows"]) up_rows = [row for row in limit_rows if row.get("limit_type") == "U"] down_rows = [row for row in limit_rows if row.get("limit_type") == "D"] broken_rows = [row for row in limit_rows if row.get("limit_type") == "Z"] limits = [self._normalize_limit(row, "涨停") for row in up_rows] broken = [self._normalize_limit(row, "炸板") for row in broken_rows] down_limits = [self._normalize_limit(row, "跌停") for row in down_rows] previous_limits = [self._normalize_limit(row, "涨停") for row in previous_limit_rows] yesterday_limits = _build_yesterday_performance( previous_limits, daily, limits, broken, down_limits, ) sectors = _build_sectors(limits) previous_sectors = _build_sectors(previous_limits) now = datetime.now().astimezone() market_status = _realtime_market_status(now.time().replace(tzinfo=None)) dashboard = { "meta": { "requested_date": _display_date(requested_date), "trade_date": _display_date(trade_date), "previous_trade_date": _display_date(previous_trade_date), "source": "tushare", "mode": "realtime", "realtime": True, "market_status": market_status, "refresh_mode": "manual", "auto_refresh": False, "quote_count": len(daily), "updated_at": now.isoformat(timespec="seconds"), "notice": "盘中行情由 Tushare rt_k 实时计算;涨停原因、封板时间和开板次数以盘后榜单校正为准。", }, "overview": _build_overview(daily, up_rows, down_rows, broken_rows), "limits": limits, "broken": broken, "down_limits": down_limits, "yesterday_limits": yesterday_limits, "limit_performance": _build_limit_performance(yesterday_limits), "ladders": _build_ladders(limits), "sectors": sectors, "sector_rotation": _build_sector_rotation(sectors, previous_sectors), } return apply_sentiment_to_dashboard(dashboard) def _load_realtime_reference( self, trade_date: str, previous_trade_date: str, ) -> dict[str, Any]: cache_key = f"{trade_date}:{previous_trade_date}" with self._realtime_reference_lock: cached = self._realtime_reference_cache.get(cache_key) if cached: return cached basic_rows = self.query( "stock_basic", {"exchange": "", "list_status": "L"}, "ts_code,name,industry,market,list_date", ) price_limits = self.query( "stk_limit", {"trade_date": trade_date}, "ts_code,trade_date,up_limit,down_limit", ) previous_limit_rows = self._load_limit_type(previous_trade_date, "U") capital_rows = self.query( "daily_basic", {"trade_date": previous_trade_date}, "ts_code,trade_date,total_share,float_share,free_share,total_mv,circ_mv", ) if not basic_rows or not price_limits: raise TushareError(f"Realtime reference data is incomplete for {trade_date}") result = { "basic_rows": basic_rows, "price_limits": price_limits, "previous_limit_rows": previous_limit_rows, "capital_rows": capital_rows, } with self._realtime_reference_lock: self._realtime_reference_cache[cache_key] = result if len(self._realtime_reference_cache) > 3: oldest = next(iter(self._realtime_reference_cache)) self._realtime_reference_cache.pop(oldest, None) return result def realtime_stock_quote( self, ts_code: str, reference_date: str = "", ) -> dict[str, Any]: rows = self.query("rt_k", {"ts_code": ts_code}) if not rows: raise TushareError(f"No realtime quote returned for {ts_code}") row = rows[0] close = _number(row.get("close")) previous_close = _number(row.get("pre_close")) if close <= 0 or previous_close <= 0: raise TushareError(f"Realtime quote is unavailable for {ts_code}") basic: dict[str, Any] = {} with self._realtime_reference_lock: references = list(self._realtime_reference_cache.values()) for reference in reversed(references): basic = next( ( item for item in reference.get("basic_rows") or [] if str(item.get("ts_code") or "") == ts_code ), {}, ) if basic: break if not basic: basics = self.query( "stock_basic", {"ts_code": ts_code}, "ts_code,name,industry,market,list_date", ) basic = basics[0] if basics else {} capital = self._latest_capital(ts_code, reference_date) float_share = _number(capital.get("float_share")) # rt_k volume is shares; daily_basic float_share is reported in 10k shares. turnover_rate = _number(row.get("vol")) / float_share / 100 if float_share else 0 market_date = reference_date or datetime.now().astimezone().strftime("%Y%m%d") self._ensure_realtime_market_cache(market_date) with self._realtime_reference_lock: market_rows = list((self._latest_realtime_market.get(market_date) or {}).get("rows") or []) references = list(self._realtime_reference_cache.values()) capital_map: dict[str, dict[str, Any]] = {} for reference in reversed(references): capital_map = { str(item.get("ts_code") or ""): item for item in reference.get("capital_rows") or [] } if capital_map: break market_amounts = [_number(item.get("amount")) for item in market_rows if _number(item.get("amount")) > 0] amount_percentile = _value_percentile(_number(row.get("amount")), market_amounts) market_turnovers = [] for item in market_rows: item_capital = capital_map.get(str(item.get("ts_code") or ""), {}) item_float_share = _number(item_capital.get("float_share")) if item_float_share: market_turnovers.append(_number(item.get("vol")) / item_float_share / 100) market_turnover = ( sum(market_turnovers) / len(market_turnovers) if market_turnovers else 0 ) turnover_relative = turnover_rate / market_turnover if market_turnover else 0 activity = self._stock_activity_metrics( ts_code, market_date, _number(row.get("vol")) / 100, ) return { "code": ts_code.split(".")[0], "ts_code": ts_code, "name": str(row.get("name") or basic.get("name") or "--").strip(), "sector": basic.get("industry") or "其他", "price": round(close, 3), "change": round((close / previous_close - 1) * 100, 4), "open": round(_number(row.get("open")), 3), "high": round(_number(row.get("high")), 3), "low": round(_number(row.get("low")), 3), "previous_close": round(previous_close, 3), "amount_billion": round(_number(row.get("amount")) / 100000000, 3), "volume": _number(row.get("vol")), "trade_count": int(_number(row.get("num"))), "turnover_rate": round(turnover_rate, 4), "market_turnover_rate": round(market_turnover, 4), "turnover_relative": round(turnover_relative, 4), "amount_percentile": round(amount_percentile * 100, 2), "volume_activity_ratio": activity.get("volume_activity_ratio", 0), "activity_history_date": activity.get("history_trade_date", ""), "activity_source": activity.get("source", "unavailable"), "float_share_10k": float_share, "capital_trade_date": str(capital.get("trade_date") or ""), "turnover_source": "rt_volume/latest_float_share" if float_share else "unavailable", "data_source": "tushare", "realtime": True, } def _stock_activity_metrics( self, ts_code: str, reference_date: str, current_volume_lots: float, ) -> dict[str, Any]: cache_key = f"{ts_code}:{reference_date}" with self._realtime_reference_lock: history = self._stock_activity_cache.get(cache_key) if history is None: try: end = datetime.strptime(reference_date, "%Y%m%d") except ValueError: end = datetime.now().astimezone().replace(tzinfo=None) rows = self.query( "daily", { "ts_code": ts_code, "start_date": (end - timedelta(days=30)).strftime("%Y%m%d"), "end_date": reference_date, }, "ts_code,trade_date,vol,amount", ) completed = [ item for item in rows if str(item.get("trade_date") or "") < reference_date and _number(item.get("vol")) > 0 ] completed.sort(key=lambda item: str(item.get("trade_date") or "")) recent = completed[-5:] history = { "average_volume_lots": ( sum(_number(item.get("vol")) for item in recent) / len(recent) if recent else 0 ), "history_trade_date": str(recent[-1].get("trade_date") or "") if recent else "", } with self._realtime_reference_lock: self._stock_activity_cache[cache_key] = history if len(self._stock_activity_cache) > 256: oldest = next(iter(self._stock_activity_cache)) self._stock_activity_cache.pop(oldest, None) average_volume = _number(history.get("average_volume_lots")) progress = _trading_session_progress(datetime.now().astimezone().time().replace(tzinfo=None)) expected_volume = average_volume * progress ratio = current_volume_lots / expected_volume if expected_volume else 0 return { **history, "volume_activity_ratio": round(ratio, 4), "session_progress": round(progress, 4), "source": "rt_volume/5d_average_at_same_progress" if expected_volume else "unavailable", } def realtime_factor_snapshot(self, requested_date: str) -> dict[str, Any]: trade_date, previous_trade_date = self.resolve_trade_context(requested_date) reference = self._load_realtime_reference(trade_date, previous_trade_date) codes = [ str(row.get("ts_code") or "") for row in reference.get("basic_rows") or [] if row.get("ts_code") ] quotes = self.query("rt_k", {"ts_code": ",".join(codes)}, "") capital_map = { str(row.get("ts_code") or ""): row for row in reference.get("capital_rows") or [] } rows = [] for quote in quotes: ts_code = str(quote.get("ts_code") or "") close = _number(quote.get("close")) previous_close = _number(quote.get("pre_close")) if not ts_code or close <= 0 or previous_close <= 0: continue capital = capital_map.get(ts_code, {}) float_share = _number(capital.get("float_share")) rows.append( { "ts_code": ts_code, "trade_date": trade_date, "open": _number(quote.get("open")), "high": _number(quote.get("high")), "low": _number(quote.get("low")), "close": close, "pct_chg": (close / previous_close - 1) * 100, "vol": _number(quote.get("vol")) / 100, "amount": _number(quote.get("amount")), "turnover_rate": ( _number(quote.get("vol")) / float_share / 100 if float_share else 0 ), "capital_trade_date": str(capital.get("trade_date") or ""), } ) if not rows: raise TushareError(f"No realtime factor snapshot returned for {trade_date}") return { "trade_date": trade_date, "previous_trade_date": previous_trade_date, "source": "tushare_rt_k", "realtime": True, "rows": rows, } def _ensure_realtime_market_cache(self, requested_date: str) -> list[dict[str, Any]]: with self._realtime_reference_lock: cached = list( (self._latest_realtime_market.get(requested_date) or {}).get("rows") or [] ) if cached: return cached trade_date, previous_trade_date = self.resolve_trade_context(requested_date) if trade_date != requested_date: return [] reference = self._load_realtime_reference(trade_date, previous_trade_date) codes = [ str(row.get("ts_code") or "") for row in reference.get("basic_rows") or [] if row.get("ts_code") ] quotes = self.query("rt_k", {"ts_code": ",".join(codes)}, "") rows = [ row for row in quotes if _number(row.get("close")) > 0 and _number(row.get("pre_close")) > 0 ] with self._realtime_reference_lock: self._latest_realtime_market[trade_date] = { "rows": rows, "updated_at": datetime.now().astimezone().isoformat(timespec="seconds"), } return rows def _latest_capital(self, ts_code: str, reference_date: str = "") -> dict[str, Any]: end_date = reference_date or datetime.now().astimezone().strftime("%Y%m%d") cache_key = f"{ts_code}:{end_date}" with self._realtime_reference_lock: cached = self._capital_cache.get(cache_key) if cached: return cached try: end = datetime.strptime(end_date, "%Y%m%d") except ValueError: end = datetime.now().astimezone().replace(tzinfo=None) end_date = end.strftime("%Y%m%d") start_date = (end - timedelta(days=20)).strftime("%Y%m%d") rows = self.query( "daily_basic", {"ts_code": ts_code, "start_date": start_date, "end_date": end_date}, "ts_code,trade_date,turnover_rate,volume_ratio,total_share,float_share," "free_share,total_mv,circ_mv", ) rows.sort(key=lambda item: str(item.get("trade_date") or "")) result = rows[-1] if rows else {} with self._realtime_reference_lock: self._capital_cache[cache_key] = result if len(self._capital_cache) > 256: oldest = next(iter(self._capital_cache)) self._capital_cache.pop(oldest, None) return result def _build_overview( daily: list[dict[str, Any]], up_rows: list[dict[str, Any]], down_rows: list[dict[str, Any]], broken_rows: list[dict[str, Any]], ) -> dict[str, Any]: up_count = sum(1 for row in daily if _number(row.get("pct_chg")) > 0) down_count = sum(1 for row in daily if _number(row.get("pct_chg")) < 0) flat_count = len(daily) - up_count - down_count amount_billion = sum( _number(row.get("amount")) / (100000000 if row.get("amount_unit") == "yuan" else 100000) for row in daily ) limit_count = len(up_rows) broken_count = len(broken_rows) seal_rate = round(limit_count / max(limit_count + broken_count, 1) * 100, 1) return { "up_count": up_count, "down_count": down_count, "flat_count": flat_count, "limit_up_count": limit_count, "limit_down_count": len(down_rows), "broken_count": broken_count, "amount_billion": round(amount_billion, 1), "seal_rate": seal_rate, } def _build_ladders(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: groups: dict[int, list[dict[str, Any]]] = {} for row in rows: groups.setdefault(int(row.get("streak") or 1), []).append(row) return [ { "level": level, "label": "首板" if level == 1 else f"{level}板", "count": len(stocks), "stocks": sorted(stocks, key=lambda item: item.get("first_time") or "99:99:99"), } for level, stocks in sorted(groups.items(), reverse=True) ] def _build_sectors(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: counts = Counter(row.get("sector") or "其他" for row in rows) result: list[dict[str, Any]] = [] for name, count in counts.most_common(20): stocks = [row for row in rows if (row.get("sector") or "其他") == name] max_streak = max(item.get("streak", 1) for item in stocks) leader = max(stocks, key=lambda item: (item.get("streak", 1), item.get("amount_billion", 0))) result.append( { "name": name, "count": count, "strength": min(100, 44 + count * 8 + max_streak * 5), "amount_billion": round(sum(item.get("amount_billion", 0) for item in stocks), 1), "leader": leader.get("name", "--"), "change": round(sum(item.get("change", 0) for item in stocks) / count, 2), "max_streak": max_streak, } ) return result def _build_yesterday_performance( previous_limits: list[dict[str, Any]], daily: list[dict[str, Any]], current_limits: list[dict[str, Any]], current_broken: list[dict[str, Any]], current_down: list[dict[str, Any]], ) -> list[dict[str, Any]]: daily_map = {str(row.get("ts_code", "")).split(".")[0]: row for row in daily} limit_map = {row["code"]: row for row in current_limits} broken_codes = {row["code"] for row in current_broken} down_codes = {row["code"] for row in current_down} result = [] for previous in previous_limits: code = previous["code"] daily_row = daily_map.get(code, {}) current = limit_map.get(code) if current: outcome = "晋级" elif code in broken_codes: outcome = "炸板" elif code in down_codes: outcome = "跌停" else: outcome = "断板" result.append( { "code": code, "name": previous["name"], "prior_streak": previous.get("streak", 1), "current_streak": current.get("streak", 0) if current else 0, "current_change": _number(daily_row.get("pct_chg")), "current_price": _number(daily_row.get("close")), "sector": previous.get("sector", "其他"), "reason": previous.get("reason", "待补充"), "outcome": outcome, } ) return result def _build_limit_performance(rows: list[dict[str, Any]]) -> list[dict[str, Any]]: result = [] for level in sorted({int(row.get("prior_streak") or 1) for row in rows}, reverse=True): group = [row for row in rows if int(row.get("prior_streak") or 1) == level] advanced = sum(row.get("outcome") == "晋级" for row in group) positive = sum(_number(row.get("current_change")) > 0 for row in group) result.append( { "level": level, "label": "昨日首板" if level == 1 else f"昨日{level}板", "count": len(group), "advanced": advanced, "advance_rate": round(advanced / len(group) * 100, 1), "positive_rate": round(positive / len(group) * 100, 1), "average_change": round(sum(_number(row.get("current_change")) for row in group) / len(group), 2), } ) return result def _build_sector_rotation( current: list[dict[str, Any]], previous: list[dict[str, Any]] ) -> list[dict[str, Any]]: previous_map = {row["name"]: row for row in previous} result = [] for index, sector in enumerate(current, start=1): previous_count = int(previous_map.get(sector["name"], {}).get("count", 0)) delta = int(sector["count"]) - previous_count result.append( { **sector, "rank": index, "previous_count": previous_count, "delta": delta, "trend": "升温" if delta > 0 else "降温" if delta < 0 else "持平", } ) return result