feat: add intraday realtime quality metadata (#1538)

Refs #1386

Refs #1389
This commit is contained in:
Alfred
2026-05-31 21:05:12 +08:00
committed by GitHub
parent 9178e03762
commit 9f14850265
14 changed files with 636 additions and 30 deletions
+97 -4
View File
@@ -19,7 +19,7 @@ import random
import time
from threading import BoundedSemaphore, RLock, Thread
from abc import ABC, abstractmethod
from datetime import datetime
from datetime import datetime, timezone
from typing import Callable, Optional, List, Tuple, Dict, Any
import pandas as pd
@@ -1427,6 +1427,78 @@ class DataFetcherManager:
except Exception as e:
logger.error(f"[预取] 批量预取异常: {e}")
return 0
@staticmethod
def _utc_now_iso() -> str:
return datetime.now(timezone.utc).isoformat()
@staticmethod
def _parse_realtime_timestamp(value: Any) -> Optional[datetime]:
if value in (None, ""):
return None
if isinstance(value, datetime):
parsed = value
else:
text = str(value).strip()
if not text:
return None
if text.endswith("Z"):
text = text[:-1] + "+00:00"
try:
parsed = datetime.fromisoformat(text)
except ValueError:
return None
if parsed.tzinfo is None:
return parsed.replace(tzinfo=timezone.utc)
return parsed.astimezone(timezone.utc)
@staticmethod
def _realtime_fetcher_token(fetcher_name: str, **kw) -> str:
if fetcher_name == "AkshareFetcher" and kw.get("source") == "hk":
return "akshare_hk"
mapping = {
"LongbridgeFetcher": "longbridge",
"YfinanceFetcher": "yfinance",
"AkshareFetcher": "akshare",
"FinnhubFetcher": "finnhub",
"AlphaVantageFetcher": "alphavantage",
"EfinanceFetcher": "efinance",
"TushareFetcher": "tushare",
}
return mapping.get(fetcher_name, fetcher_name.replace("Fetcher", "").lower())
def _enrich_realtime_quote(
self,
quote,
*,
fallback_from: Optional[str] = None,
realtime_cache_ttl: Optional[int] = None,
):
"""Attach runtime metadata without inventing provider-side timestamps."""
if quote is None:
return None
fetched_at = self._utc_now_iso()
setattr(quote, "fetched_at", fetched_at)
if fallback_from:
setattr(quote, "fallback_from", str(fallback_from))
provider_dt = self._parse_realtime_timestamp(
getattr(quote, "provider_timestamp", None)
)
if provider_dt is None:
setattr(quote, "provider_timestamp", None)
setattr(quote, "stale_seconds", None)
setattr(quote, "is_stale", None)
return quote
setattr(quote, "provider_timestamp", provider_dt.isoformat())
fetched_dt = self._parse_realtime_timestamp(fetched_at) or datetime.now(timezone.utc)
stale_seconds = max(0, int((fetched_dt - provider_dt).total_seconds()))
ttl = realtime_cache_ttl if realtime_cache_ttl is not None else 600
setattr(quote, "stale_seconds", stale_seconds)
setattr(quote, "is_stale", stale_seconds > int(ttl))
return quote
def get_realtime_quote(self, stock_code: str, *, log_final_failure: bool = True):
"""
@@ -1488,7 +1560,9 @@ class DataFetcherManager:
primary_kw = {"source": "hk"} if primary_src == "AkshareFetcher" else {}
secondary_kw = {"source": "hk"} if secondary_src == "AkshareFetcher" else {}
primary_token = self._realtime_fetcher_token(primary_src, **primary_kw)
primary_quote = self._try_fetcher_quote(stock_code, primary_src, **primary_kw)
fallback_from = primary_token if primary_quote is None else None
if primary_quote is not None:
logger.info(f"[实时行情] {market_label} {stock_code} 成功获取 (来源: {primary_src})")
primary_quote = self._supplement_quote(
@@ -1501,7 +1575,11 @@ class DataFetcherManager:
stock_code, primary_quote, extra_src,
)
if primary_quote is not None:
return primary_quote
return self._enrich_realtime_quote(
primary_quote,
fallback_from=fallback_from,
realtime_cache_ttl=getattr(config, "realtime_cache_ttl", None),
)
if log_final_failure:
logger.info(f"[实时行情] {market_label} {stock_code} 无可用数据源")
return None
@@ -1514,9 +1592,11 @@ class DataFetcherManager:
]
errors = []
failed_sources: List[str] = []
# primary_quote holds the first successful result; we may supplement
# missing fields (volume_ratio, turnover_rate, etc.) from later sources.
primary_quote = None
primary_fallback_from: Optional[str] = None
for source_index, source in enumerate(source_priority):
attempt_start = time.time()
@@ -1565,10 +1645,15 @@ class DataFetcherManager:
if primary_quote is None:
# First successful source becomes primary
primary_quote = quote
primary_fallback_from = failed_sources[0] if failed_sources else None
logger.info(f"[实时行情] {stock_code} 成功获取 (来源: {source})")
# If all key supplementary fields are present, return early
if not self._quote_needs_supplement(primary_quote):
return primary_quote
return self._enrich_realtime_quote(
primary_quote,
fallback_from=primary_fallback_from,
realtime_cache_ttl=getattr(config, "realtime_cache_ttl", None),
)
# Otherwise, continue to try later sources for missing fields
logger.debug(f"[实时行情] {stock_code} 部分字段缺失,尝试从后续数据源补充")
supplement_attempts = 0
@@ -1596,6 +1681,8 @@ class DataFetcherManager:
fallback_to=fallback_to,
record_count=0,
)
if primary_quote is None:
failed_sources.append(source)
except Exception as e:
error_msg = f"[{source}] 失败: {str(e)}"
@@ -1612,11 +1699,17 @@ class DataFetcherManager:
)
logger.info(f"[实时行情] {stock_code} {error_msg},继续尝试下一个数据源")
errors.append(error_msg)
if primary_quote is None:
failed_sources.append(source)
continue
# Return primary even if some fields are still missing
if primary_quote is not None:
return primary_quote
return self._enrich_realtime_quote(
primary_quote,
fallback_from=primary_fallback_from,
realtime_cache_ttl=getattr(config, "realtime_cache_ttl", None),
)
# 所有数据源都失败,返回 None(降级兜底)
if log_final_failure:
+9
View File
@@ -118,6 +118,13 @@ class UnifiedRealtimeQuote:
code: str
name: str = ""
source: RealtimeSource = RealtimeSource.FALLBACK
# === 数据质量元数据(由 DataFetcherManager 统一补齐)===
fetched_at: Optional[str] = None # 本系统获取时间(ISO 8601 datetime)
provider_timestamp: Optional[str] = None # Provider 真实行情时间(ISO 8601 datetime)
is_stale: Optional[bool] = None # provider_timestamp 超过最小 TTL 阈值时为 True
stale_seconds: Optional[int] = None # provider_timestamp 距 fetched_at 的秒数
fallback_from: Optional[str] = None # 整源 fallback 的失败首选源 token
# === 核心价格数据(几乎所有源都有)===
price: Optional[float] = None # 最新价
@@ -157,6 +164,8 @@ class UnifiedRealtimeQuote:
}
# 只添加非 None 的字段
optional_fields = [
'fetched_at', 'provider_timestamp', 'is_stale', 'stale_seconds',
'fallback_from',
'price', 'change_pct', 'change_amount', 'volume', 'amount',
'volume_ratio', 'turnover_rate', 'amplitude',
'open_price', 'high', 'low', 'pre_close',
+1
View File
@@ -31,6 +31,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
- [测试] 补充 ETF 日线数据源路由、输入变体、fallback 与 MA 字段回归覆盖。
- [改进] 优化 Web 报告详情页信息层级,将输入数据块和运行诊断下移为主体内容后的折叠辅助信息。
- [新功能] 市场阶段低敏摘要接入历史详情、同步分析响应和 completed 任务状态的 report metadata。
- [改进] 盘中分析补齐实时行情获取时间、provider 时间、stale、fallback 与 partial/estimated 标记,供 AnalysisContextPack 映射输入数据限制。
## [3.19.0] - 2026-05-29
+3 -3
View File
@@ -82,10 +82,10 @@ P2 block 组装边界:
- `subject` 仍只写 `code`、`stock_name`、`market` 三字段,不扩 `AnalysisSubject`。
- `phase` 只接收传入的 `MarketPhaseContext.to_dict()` 产物,不从 `enhanced_context` 反推。
- `quote` 从 `realtime_quote` 组装;缺失为 `missing`;`source=fallback` 映射为 `fallback`;`fallback_from` 只在 artifact/metadata 显式提供时填写,否则只记录稳定 warning code,不伪造 provider 链。
- `quote` stale 只透传 `price_stale`、`quote_stale`、`quote_stale_seconds` 等显式 marker;builder 不推断新鲜度。
- `quote` 从 `realtime_quote` 组装;缺失为 `missing`;`source=fallback` 或显式 `fallback_from` 映射为 `fallback`,但 `source` 保留真实成功源;`fallback_from` 只在 artifact/metadata 显式提供时填写,否则只记录稳定 warning code,不伪造 provider 链。
- `quote` 会透传 #1386 P3 的 `fetched_at`、`provider_timestamp`、`is_stale`、`stale_seconds`、`fallback_from`。状态优先级固定为 `STALE > FALLBACK > AVAILABLE`:`is_stale=True`、`price_stale`、`quote_stale`、`quote_stale_seconds` 等显式 marker 标为 `stale`;`stale_seconds` 且 `is_stale=False` 只是元数据,不单独推断 stale。builder 只映射上游 artifact,不做质量评分。
- `daily_bars` 只表达完整日线窗口,优先读 `base_context.today`、`base_context.yesterday`、`base_context.date`、`base_context.data_missing`;date-only 放入 `value` 或 `metadata`,不写入 `timestamp`。
- `enhanced_context.today.data_source` 为 `realtime:*` 时,只影响 `technical`:block 标 `partial`,相关 item 标 `estimated`,warning 使用 `intraday_realtime_overlay`。
- `enhanced_context.today` 上的 `is_partial_bar`、`is_estimated`、`estimated_fields` 优先进入 `technical`;缺失时仍兼容 `enhanced_context.today.data_source` 为 `realtime:*` 的旧 heuristic。partial/estimated 只进入 `technical`,`daily_bars` 不承载 partial/estimated,warning 使用 `intraday_realtime_overlay`。
- `technical` 优先复用 `trend_result.to_dict()`;无 trend artifact 时为 `missing`。
- `chip` 复用 `chip_data.to_dict()`;无 chip artifact 默认 `missing`,只有输入 metadata/artifact 明确 not_supported 时才标 `not_supported`。
- `fundamentals` 只读 `fundamental_context` 参数;`ok` 映射为 `available`,`not_supported` 映射为 `not_supported`,`partial` 映射为 `partial`,`failed` 映射为 `missing` + 稳定 reason code;不写入 `errors[]` 原文。
+6
View File
@@ -764,6 +764,12 @@ P2-min 开始在已获得 `market_phase_context` 的分析路径中,把运行
P2-min 仍不新增 API/Web/Bot 参数,不写入 history/task status/report metadata,不改变报告 JSON schema,也不引入完整 quote freshness、fallback、stale 或 data_quality 契约。Bot/API 直连 Agent 若未经过 P1a pipeline 构建 `market_phase_context`,仍保持旧行为;入口透传和可见展示留给后续 P4+。
#### 盘中数据包与实时质量控制(Issue #1386 P3)
P3 补齐普通分析主路径使用的实时行情质量元数据,但仍不新增 `analysis_phase` 参数,不改 API/Web/Bot 阶段入口,不改变报告 JSON schema,也不做 #1389 P5 数据质量评分或模型置信度限制。实时 quote 会带上 `fetched_at`、`provider_timestamp`、`is_stale`、`stale_seconds`、`fallback_from`;其中 `fetched_at` 是系统获取时间,`provider_timestamp` 只在 provider 真实提供行情时间时填写。缺少 provider 时间时不会伪造 fresh,`stale_seconds` 和 `is_stale` 保持空值。
整源 fallback 的语义固定为:`source` 保留实际成功的数据源 token,`fallback_from` 记录本轮失败的最高优先级整源 token;首选源成功后只从后续源补字段时不写 `fallback_from`。`AnalysisContextBuilder` 只映射这些上游 artifact,不重新取数、不做质量评分;quote block 状态按 `STALE > FALLBACK > AVAILABLE` 归并。盘中实时价覆盖 `today` 时会标记 `is_partial_bar`、`is_estimated`、`estimated_fields`、`realtime_source` 和 quote 元数据;`daily_bars` block 仍表示 storage 中完整日线窗口,partial/estimated 只进入 technical block。freshness scoring、盘中 cache TTL 分级、Agent 工具级复用和 API/Web 展示留给后续阶段。
#### AnalysisContextPack Prompt 摘要(Issue #1389 P3)
P3 在普通分析和 Agent 初始上下文中接入 `AnalysisContextPack` 低敏摘要。Pipeline 会用已获取的行情、日线、趋势、筹码、基本面、新闻和市场阶段 artifacts 组装 pack,再把 `analysis_context_pack_summary` 插入 Prompt;在这个新增的 pack 摘要区块中,LLM 只看到 subject、版本、各数据块的状态/来源/warning/missing reason 和新闻结果数,不会通过该区块看到完整 `news.content`、`trend_result`、筹码或基本面原始 payload。既有 `news_context`、Agent pre-fetched JSON 和 `enhanced_context` 原始数据通道保持 P3 前行为,不由本摘要替代或脱敏。
+6
View File
@@ -641,6 +641,12 @@ P2-min starts rendering the runtime market phase into an LLM-readable prompt sec
P2-min still does not add API/Web/Bot parameters, persist phase into history/task status/report metadata, change report JSON schemas, or introduce the full quote freshness, fallback, stale, or data-quality contract. Bot/API direct Agent entrypoints that do not go through the P1a pipeline to build `market_phase_context` keep their previous behavior; entrypoint propagation and visible labels are left to later P4+ work.
### Intraday Data Packet and Realtime Quality Control (Issue #1386 P3)
P3 adds realtime quote quality metadata for the regular analysis path, but still does not add an `analysis_phase` parameter, change API/Web/Bot phase entrypoints, change report JSON schemas, or implement #1389 P5 data-quality scoring or model confidence limits. Realtime quotes may carry `fetched_at`, `provider_timestamp`, `is_stale`, `stale_seconds`, and `fallback_from`; `fetched_at` is the system fetch time, while `provider_timestamp` is only populated when the provider actually returns a quote timestamp. If provider time is unavailable, the system does not fabricate freshness, and `stale_seconds` / `is_stale` stay empty.
Whole-source fallback semantics are fixed: `source` keeps the actual successful provider token, while `fallback_from` records the highest-priority whole source that failed in the current attempt; if the primary source succeeds and later providers only supplement missing fields, `fallback_from` is not set. `AnalysisContextBuilder` only maps these upstream artifacts, performs no extra fetches, and does no quality scoring; quote block status collapses as `STALE > FALLBACK > AVAILABLE`. When realtime price overlays `today`, the pipeline marks `is_partial_bar`, `is_estimated`, `estimated_fields`, `realtime_source`, and quote metadata. The `daily_bars` block still represents the complete daily-bar window in storage; partial/estimated markers only enter the technical block. Freshness scoring, intraday cache TTL tiers, Agent tool-level reuse, and API/Web display remain follow-ups.
### AnalysisContextPack Prompt Summary (Issue #1389 P3)
P3 injects a low-sensitivity `AnalysisContextPack` summary into regular analysis and Agent initial prompts. The pipeline builds the pack from already-fetched quote, daily-bar, trend, chip, fundamentals, news, and market-phase artifacts, then passes `analysis_context_pack_summary` downstream; in this new pack-summary section, the LLM only sees subject, version, data-block status/source/warnings/missing reason, and news result count, not full `news.content`, `trend_result`, chip, or fundamentals raw payloads through that section. Existing `news_context`, Agent pre-fetched JSON, and `enhanced_context` raw-payload channels keep their pre-P3 behavior and are not replaced or sanitized by this summary.
+36 -2
View File
@@ -512,6 +512,7 @@ class StockAnalysisPipeline:
trend_result,
stock_name, # 传入股票名称
fundamental_context,
market_phase_context=market_phase_context_dict,
)
enhanced_context["market_phase_context"] = market_phase_context_dict
@@ -657,7 +658,8 @@ class StockAnalysisPipeline:
chip_data: Optional[ChipDistribution],
trend_result: Optional[TrendAnalysisResult],
stock_name: str = "",
fundamental_context: Optional[Dict[str, Any]] = None
fundamental_context: Optional[Dict[str, Any]] = None,
market_phase_context: Optional[Dict[str, Any]] = None,
) -> Dict[str, Any]:
"""
增强分析上下文
@@ -670,6 +672,7 @@ class StockAnalysisPipeline:
chip_data: 筹码分布数据
trend_result: 趋势分析结果
stock_name: 股票名称
market_phase_context: 已构建的市场阶段上下文,用于标记盘中 partial bar
Returns:
增强后的上下文
@@ -690,6 +693,9 @@ class StockAnalysisPipeline:
if realtime_quote:
# 使用 getattr 安全获取字段,缺失字段返回 None 或默认值
volume_ratio = getattr(realtime_quote, 'volume_ratio', None)
quote_source = getattr(realtime_quote, 'source', None)
quote_source_name = getattr(quote_source, 'value', quote_source)
quote_source_name = str(quote_source_name) if quote_source_name is not None else None
enhanced['realtime'] = {
'name': getattr(realtime_quote, 'name', ''),
'price': getattr(realtime_quote, 'price', None),
@@ -702,7 +708,12 @@ class StockAnalysisPipeline:
'total_mv': getattr(realtime_quote, 'total_mv', None),
'circ_mv': getattr(realtime_quote, 'circ_mv', None),
'change_60d': getattr(realtime_quote, 'change_60d', None),
'source': getattr(realtime_quote, 'source', None),
'source': quote_source_name,
'fetched_at': getattr(realtime_quote, 'fetched_at', None),
'provider_timestamp': getattr(realtime_quote, 'provider_timestamp', None),
'is_stale': getattr(realtime_quote, 'is_stale', None),
'stale_seconds': getattr(realtime_quote, 'stale_seconds', None),
'fallback_from': getattr(realtime_quote, 'fallback_from', None),
}
# 移除 None 值以减少上下文大小
enhanced['realtime'] = {k: v for k, v in enhanced['realtime'].items() if v is not None}
@@ -757,6 +768,9 @@ class StockAnalysisPipeline:
vol = getattr(realtime_quote, 'volume', None)
amt = getattr(realtime_quote, 'amount', None)
pct = getattr(realtime_quote, 'change_pct', None)
fetched_at = getattr(realtime_quote, 'fetched_at', None)
provider_timestamp = getattr(realtime_quote, 'provider_timestamp', None)
fallback_from = getattr(realtime_quote, 'fallback_from', None)
realtime_today = {
'close': price,
'open': open_p,
@@ -768,18 +782,38 @@ class StockAnalysisPipeline:
'date': market_today,
'data_source': f"realtime:{source_name}",
'realtime_source': source_name,
'is_estimated': True,
}
estimated_fields = [
'close', 'open', 'high', 'low', 'ma5', 'ma10', 'ma20',
]
if vol is not None:
realtime_today['volume'] = vol
estimated_fields.append('volume')
if amt is not None:
realtime_today['amount'] = amt
estimated_fields.append('amount')
if pct is not None:
realtime_today['pct_chg'] = pct
estimated_fields.append('pct_chg')
realtime_today['estimated_fields'] = estimated_fields
if isinstance(market_phase_context, dict) and "is_partial_bar" in market_phase_context:
realtime_today['is_partial_bar'] = market_phase_context.get("is_partial_bar")
if fetched_at is not None:
realtime_today['fetched_at'] = fetched_at
if provider_timestamp is not None:
realtime_today['provider_timestamp'] = provider_timestamp
if fallback_from is not None:
realtime_today['fallback_from'] = fallback_from
realtime_owned_fields = {
'open', 'high', 'low', 'close',
'volume', 'amount', 'pct_chg', 'pctChg',
'date', 'data_source', 'dataSource', 'source',
'realtime_source', 'realtimeSource',
'is_partial_bar', 'isPartialBar', 'is_estimated',
'isEstimated', 'estimated_fields', 'estimatedFields',
'fetched_at', 'fetchedAt', 'provider_timestamp',
'providerTimestamp', 'fallback_from', 'fallbackFrom',
}
for k, v in orig_today.items():
if k not in realtime_today and k not in realtime_owned_fields and v is not None:
+125 -17
View File
@@ -5,6 +5,7 @@ from __future__ import annotations
from collections.abc import Mapping
from dataclasses import dataclass
from datetime import datetime
from typing import Any, Dict, List, Optional, Sequence
from src.schemas.analysis_context_pack import (
@@ -107,11 +108,13 @@ def _build_quote_block(artifacts: PipelineAnalysisArtifacts) -> AnalysisContextB
"realtime_fallback_from",
"fallback_from",
)
timestamp = _quote_timestamp(artifacts, quote)
is_fallback = fallback_from is not None or source == "fallback"
if _has_explicit_quote_stale_marker(artifacts, quote):
status = ContextFieldStatus.STALE
warnings.append("quote_stale")
elif source == "fallback":
elif is_fallback:
status = ContextFieldStatus.FALLBACK
if fallback_from is None:
warnings.append(_REALTIME_FALLBACK_WARNING)
@@ -121,7 +124,8 @@ def _build_quote_block(artifacts: PipelineAnalysisArtifacts) -> AnalysisContextB
status=status,
value=value,
source=source,
fallback_from=fallback_from if status == ContextFieldStatus.FALLBACK else None,
timestamp=timestamp,
fallback_from=fallback_from if is_fallback else None,
warnings=list(warnings),
)
for key, value in quote.items()
@@ -131,6 +135,7 @@ def _build_quote_block(artifacts: PipelineAnalysisArtifacts) -> AnalysisContextB
status=status,
items=items,
source=source,
timestamp=timestamp,
warnings=warnings,
metadata=_quote_metadata(artifacts, quote),
)
@@ -217,7 +222,12 @@ def _build_technical_block(
[],
)
has_realtime_overlay = _has_realtime_overlay(artifacts.enhanced_context)
explicit_intraday_overlay = _has_explicit_intraday_overlay(
artifacts.enhanced_context
)
has_realtime_overlay = explicit_intraday_overlay or _has_realtime_overlay(
artifacts.enhanced_context
)
warnings = [_REALTIME_OVERLAY_WARNING] if has_realtime_overlay else []
block_status = (
ContextFieldStatus.PARTIAL
@@ -244,7 +254,24 @@ def _build_technical_block(
items=items,
warnings=warnings,
metadata={
"overlay_source": _realtime_overlay_source(artifacts.enhanced_context)
key: value
for key, value in {
"overlay_source": _realtime_overlay_source(
artifacts.enhanced_context
),
"is_partial_bar": _today_metadata_value(
artifacts.enhanced_context, "is_partial_bar", "isPartialBar"
),
"is_estimated": _today_metadata_value(
artifacts.enhanced_context, "is_estimated", "isEstimated"
),
"estimated_fields": _today_metadata_value(
artifacts.enhanced_context,
"estimated_fields",
"estimatedFields",
),
}.items()
if value is not None
},
),
warnings,
@@ -423,19 +450,63 @@ def _metadata_value(metadata: Dict[str, Any], *keys: str) -> Optional[str]:
return None
def _metadata_iso_datetime_value(metadata: Dict[str, Any], *keys: str) -> Optional[str]:
for key in keys:
value = (metadata or {}).get(key)
if value in (None, ""):
continue
if isinstance(value, datetime):
return value.isoformat()
text = str(value).strip()
if not text:
continue
if "T" not in text:
continue
normalized = text[:-1] + "+00:00" if text.endswith("Z") else text
try:
datetime.fromisoformat(normalized)
except ValueError:
continue
return text
return None
def _quote_timestamp(
artifacts: PipelineAnalysisArtifacts,
quote: Dict[str, Any],
) -> Optional[str]:
return _metadata_iso_datetime_value(
quote,
"provider_timestamp",
"quote_timestamp",
) or _metadata_iso_datetime_value(
artifacts.metadata,
"provider_timestamp",
"quote_timestamp",
"realtime_provider_timestamp",
) or _metadata_iso_datetime_value(
quote,
"fetched_at",
"realtime_fetched_at",
) or _metadata_iso_datetime_value(
artifacts.metadata,
"fetched_at",
"realtime_fetched_at",
)
def _has_explicit_quote_stale_marker(
artifacts: PipelineAnalysisArtifacts,
quote: Dict[str, Any],
) -> bool:
metadata = artifacts.metadata or {}
for key in (
"price_stale",
"quote_stale",
"quote_stale_seconds",
"stale_seconds",
):
for key in ("price_stale", "quote_stale", "is_stale"):
if bool(metadata.get(key)) or bool(quote.get(key)):
return True
if bool(metadata.get("quote_stale_seconds")) or bool(
quote.get("quote_stale_seconds")
):
return True
return False
@@ -448,27 +519,64 @@ def _quote_metadata(
"price_stale",
"quote_stale",
"quote_stale_seconds",
"is_stale",
"stale_seconds",
"fetched_at",
"provider_timestamp",
"fallback_from",
):
value = (artifacts.metadata or {}).get(key)
if value is None:
value = quote.get(key)
if key in {"fetched_at", "provider_timestamp"}:
value = _metadata_iso_datetime_value(artifacts.metadata or {}, key)
if value is None:
value = _metadata_iso_datetime_value(quote, key)
else:
value = (artifacts.metadata or {}).get(key)
if value is None:
value = quote.get(key)
if value is not None:
metadata[key] = value
return metadata
def _has_realtime_overlay(enhanced_context: Dict[str, Any]) -> bool:
def _today_dict(enhanced_context: Dict[str, Any]) -> Optional[Dict[str, Any]]:
today = (enhanced_context or {}).get("today")
if not isinstance(today, dict):
return today if isinstance(today, dict) else None
def _today_metadata_value(enhanced_context: Dict[str, Any], *keys: str) -> Any:
today = _today_dict(enhanced_context)
if today is None:
return None
for key in keys:
value = today.get(key)
if value is not None:
return value
return None
def _has_explicit_intraday_overlay(enhanced_context: Dict[str, Any]) -> bool:
today = _today_dict(enhanced_context)
if today is None:
return False
if bool(today.get("is_partial_bar")) or bool(today.get("isPartialBar")):
return True
if bool(today.get("is_estimated")) or bool(today.get("isEstimated")):
return True
estimated_fields = today.get("estimated_fields") or today.get("estimatedFields")
return bool(estimated_fields)
def _has_realtime_overlay(enhanced_context: Dict[str, Any]) -> bool:
today = _today_dict(enhanced_context)
if today is None:
return False
data_source = today.get("data_source") or today.get("dataSource")
return isinstance(data_source, str) and data_source.startswith("realtime:")
def _realtime_overlay_source(enhanced_context: Dict[str, Any]) -> Optional[str]:
today = (enhanced_context or {}).get("today")
if not isinstance(today, dict):
today = _today_dict(enhanced_context)
if today is None:
return None
value = today.get("data_source") or today.get("dataSource")
return value if isinstance(value, str) and value else None
+111 -1
View File
@@ -44,7 +44,10 @@ class _InvalidTrend:
return ["not", "a", "mapping"]
def _quote(source: RealtimeSource = RealtimeSource.AKSHARE_EM) -> UnifiedRealtimeQuote:
def _quote(
source: RealtimeSource = RealtimeSource.AKSHARE_EM,
**overrides,
) -> UnifiedRealtimeQuote:
return UnifiedRealtimeQuote(
code="600519",
name="贵州茅台",
@@ -53,6 +56,7 @@ def _quote(source: RealtimeSource = RealtimeSource.AKSHARE_EM) -> UnifiedRealtim
change_pct=1.2,
volume_ratio=1.3,
turnover_rate=0.5,
**overrides,
)
@@ -143,6 +147,92 @@ def test_quote_block_maps_available_missing_fallback_and_explicit_stale() -> Non
assert "quote_stale" in stale.warnings
def test_quote_block_maps_realtime_metadata_and_status_priority() -> None:
fallback = AnalysisContextBuilder.build(
_artifacts(
realtime_quote=_quote(
fetched_at="2026-05-31T10:00:05+00:00",
provider_timestamp="2026-05-31T10:00:00+00:00",
is_stale=False,
stale_seconds=5,
fallback_from="efinance",
)
)
).blocks["quote"]
assert fallback.status == ContextFieldStatus.FALLBACK
assert fallback.source == "akshare_em"
assert fallback.timestamp == "2026-05-31T10:00:00+00:00"
assert fallback.items["price"].timestamp == "2026-05-31T10:00:00+00:00"
assert fallback.items["price"].fallback_from == "efinance"
assert fallback.metadata["fetched_at"] == "2026-05-31T10:00:05+00:00"
assert fallback.metadata["provider_timestamp"] == "2026-05-31T10:00:00+00:00"
assert fallback.metadata["is_stale"] is False
assert fallback.metadata["stale_seconds"] == 5
assert fallback.metadata["fallback_from"] == "efinance"
stale = AnalysisContextBuilder.build(
_artifacts(
realtime_quote=_quote(
fetched_at="2026-05-31T10:15:00+00:00",
provider_timestamp="2026-05-31T10:00:00+00:00",
is_stale=True,
stale_seconds=900,
fallback_from="efinance",
)
)
).blocks["quote"]
assert stale.status == ContextFieldStatus.STALE
assert stale.source == "akshare_em"
assert stale.items["price"].fallback_from == "efinance"
assert "quote_stale" in stale.warnings
def test_quote_block_ignores_invalid_or_legacy_timestamp_metadata() -> None:
block = AnalysisContextBuilder.build(
_artifacts(
realtime_quote={
"source": "akshare_em",
"price": 1870.0,
"provider_timestamp": "not-a-date",
"fetched_at": "2026-05-31T10:00:05+00:00",
"timestamp": 0,
}
)
).blocks["quote"]
assert block.status == ContextFieldStatus.AVAILABLE
assert block.timestamp == "2026-05-31T10:00:05+00:00"
assert block.items["price"].timestamp == "2026-05-31T10:00:05+00:00"
assert block.metadata["fetched_at"] == "2026-05-31T10:00:05+00:00"
assert "provider_timestamp" not in block.metadata
legacy_timestamp_only = AnalysisContextBuilder.build(
_artifacts(
realtime_quote={
"source": "akshare_em",
"price": 1870.0,
"timestamp": 0,
}
)
).blocks["quote"]
assert legacy_timestamp_only.timestamp is None
assert legacy_timestamp_only.items["price"].timestamp is None
space_separated_timestamp = AnalysisContextBuilder.build(
_artifacts(
realtime_quote={
"source": "akshare_em",
"price": 1870.0,
"provider_timestamp": "2026-05-31 10:00:00",
}
)
).blocks["quote"]
assert space_separated_timestamp.timestamp is None
assert space_separated_timestamp.items["price"].timestamp is None
def test_daily_bars_uses_base_context_and_keeps_dates_out_of_timestamp() -> None:
pack = AnalysisContextBuilder.build(
_artifacts(
@@ -221,6 +311,26 @@ def test_technical_missing_and_realtime_overlay_statuses_are_explicit() -> None:
assert "intraday_realtime_overlay" in block.warnings
assert "intraday_realtime_overlay" in pack.data_quality.warnings
explicit_pack = AnalysisContextBuilder.build(
_artifacts(
enhanced_context={
"today": {
"close": 1880.0,
"is_partial_bar": True,
"is_estimated": True,
"estimated_fields": ["close", "ma5"],
}
}
)
)
explicit_block = explicit_pack.blocks["technical"]
assert explicit_block.status == ContextFieldStatus.PARTIAL
assert explicit_block.items["intraday_overlay"].status == ContextFieldStatus.ESTIMATED
assert explicit_block.metadata["is_partial_bar"] is True
assert explicit_block.metadata["is_estimated"] is True
assert explicit_block.metadata["estimated_fields"] == ["close", "ma5"]
def test_chip_missing_defaults_to_missing_and_explicit_not_supported() -> None:
missing = AnalysisContextBuilder.build(_artifacts(chip_data=None)).blocks["chip"]
+28
View File
@@ -6,6 +6,8 @@ from pathlib import Path
PROJECT_ROOT = Path(__file__).resolve().parents[1]
DOC_PATH = PROJECT_ROOT / "docs" / "analysis-context-pack.md"
FULL_GUIDE_PATH = PROJECT_ROOT / "docs" / "full-guide.md"
FULL_GUIDE_EN_PATH = PROJECT_ROOT / "docs" / "full-guide_EN.md"
def _read_doc() -> str:
@@ -226,6 +228,32 @@ def test_analysis_context_pack_doc_defines_p2_builder_boundaries() -> None:
assert token in section
def test_analysis_context_pack_docs_record_issue_1386_p3_quality_boundaries() -> None:
section = _section(_read_doc(), "P2 Builder 契约")
for token in (
"`fetched_at`",
"`provider_timestamp`",
"`is_stale`",
"`stale_seconds`",
"`fallback_from`",
"`STALE > FALLBACK > AVAILABLE`",
"builder 只映射上游 artifact,不做质量评分",
"`is_partial_bar`、`is_estimated`、`estimated_fields`",
"`daily_bars` 不承载 partial/estimated",
):
assert token in section
full_guide = FULL_GUIDE_PATH.read_text(encoding="utf-8")
full_guide_en = FULL_GUIDE_EN_PATH.read_text(encoding="utf-8")
assert "盘中数据包与实时质量控制(Issue #1386 P3)" in full_guide
assert "source` 保留实际成功的数据源 token" in full_guide
assert "`AnalysisContextBuilder` 只映射这些上游 artifact" in full_guide
assert "daily_bars` block 仍表示 storage 中完整日线窗口" in full_guide
assert "Intraday Data Packet and Realtime Quality Control (Issue #1386 P3)" in full_guide_en
assert "source` keeps the actual successful provider token" in full_guide_en
def test_analysis_context_pack_doc_defines_p3_runtime_consumption_boundaries() -> None:
section = _section(_read_doc(), "P3 Runtime Consumption")
+30
View File
@@ -172,6 +172,36 @@ class TestFetcherSourceOptimization(unittest.TestCase):
yfinance.get_realtime_quote.assert_called_once_with("AAPL")
longbridge.get_realtime_quote.assert_not_called()
@patch("src.config.get_config")
def test_us_realtime_route_marks_longbridge_fallback_when_secondary_succeeds(self, mock_get_config):
mock_get_config.return_value = SimpleNamespace(
enable_realtime_quote=True,
realtime_source_priority="efinance,akshare_em,tushare",
realtime_cache_ttl=600,
)
longbridge = MagicMock()
longbridge.name = "LongbridgeFetcher"
longbridge.priority = 5
longbridge.is_available_for_request.return_value = True
longbridge.get_realtime_quote.return_value = None
yfinance_quote = _make_quote("AAPL")
yfinance = MagicMock()
yfinance.name = "YfinanceFetcher"
yfinance.priority = 4
yfinance.get_realtime_quote.return_value = yfinance_quote
manager = DataFetcherManager(fetchers=[longbridge, yfinance])
quote = manager.get_realtime_quote("AAPL")
self.assertIs(quote, yfinance_quote)
self.assertEqual(quote.fallback_from, "longbridge")
self.assertIsNotNone(quote.fetched_at)
longbridge.get_realtime_quote.assert_called_once_with("AAPL")
yfinance.get_realtime_quote.assert_called_once_with("AAPL")
@patch("src.config.get_config")
def test_us_daily_route_skips_temporarily_unavailable_longbridge(self, mock_get_config):
mock_get_config.return_value = SimpleNamespace(
@@ -31,6 +31,7 @@ def _make_realtime_quote(
volume: int = 13995600,
amount: float = None,
change_pct: float = 0.96,
**overrides,
) -> UnifiedRealtimeQuote:
return UnifiedRealtimeQuote(
code="600519",
@@ -43,6 +44,7 @@ def _make_realtime_quote(
volume=volume,
amount=amount,
change_pct=change_pct,
**overrides,
)
@@ -279,6 +281,70 @@ class TestEnhanceContextRealtimeOverride(unittest.TestCase):
self.assertEqual(enhanced["today"]["realtime_source"], "tencent")
self.assertNotIn("dataSource", enhanced["today"])
@patch("src.core.pipeline.get_market_now")
@patch("src.core.pipeline.get_market_for_stock", return_value="cn")
def test_realtime_metadata_and_partial_estimated_fields_are_propagated(
self, _mock_market, mock_now
) -> None:
today = date.today()
mock_now.return_value = datetime(
today.year, today.month, today.day, 10, 0, tzinfo=timezone.utc
)
context = {
"code": "600519",
"date": (today - timedelta(days=1)).isoformat(),
"today": {
"close": 15.0,
"amount": 999999,
"date": (today - timedelta(days=1)).isoformat(),
"dataSource": "AkshareFetcher",
},
"yesterday": {"close": 14.5, "volume": 1000000},
}
quote = _make_realtime_quote(
price=15.72,
amount=None,
fetched_at="2026-05-31T10:00:05+00:00",
provider_timestamp="2026-05-31T10:00:00+00:00",
is_stale=False,
stale_seconds=5,
fallback_from="efinance",
)
trend = TrendAnalysisResult(
code="600519",
trend_status=TrendStatus.BULL,
ma5=15.5,
ma10=15.2,
ma20=14.9,
)
enhanced = self.pipeline._enhance_context(
context,
quote,
None,
trend,
"贵州茅台",
market_phase_context={"is_partial_bar": True},
)
self.assertEqual(enhanced["realtime"]["source"], "tencent")
self.assertEqual(enhanced["realtime"]["fetched_at"], "2026-05-31T10:00:05+00:00")
self.assertEqual(enhanced["realtime"]["provider_timestamp"], "2026-05-31T10:00:00+00:00")
self.assertIs(enhanced["realtime"]["is_stale"], False)
self.assertEqual(enhanced["realtime"]["stale_seconds"], 5)
self.assertEqual(enhanced["realtime"]["fallback_from"], "efinance")
self.assertTrue(enhanced["today"]["is_partial_bar"])
self.assertTrue(enhanced["today"]["is_estimated"])
self.assertEqual(
enhanced["today"]["estimated_fields"],
["close", "open", "high", "low", "ma5", "ma10", "ma20", "volume", "pct_chg"],
)
self.assertEqual(enhanced["today"]["fetched_at"], "2026-05-31T10:00:05+00:00")
self.assertEqual(enhanced["today"]["provider_timestamp"], "2026-05-31T10:00:00+00:00")
self.assertEqual(enhanced["today"]["fallback_from"], "efinance")
self.assertNotIn("amount", enhanced["today"])
self.assertNotIn("dataSource", enhanced["today"])
@patch("src.core.pipeline.get_market_now")
@patch("src.core.pipeline.get_market_for_stock", return_value="cn")
def test_realtime_today_does_not_backfill_historical_amount_or_source(
+82 -2
View File
@@ -40,13 +40,19 @@ class _DummyFetcher:
return self._result
def _make_quote(code: str = "600519", name: str = "贵州茅台") -> UnifiedRealtimeQuote:
def _make_quote(
code: str = "600519",
name: str = "贵州茅台",
source: RealtimeSource = RealtimeSource.AKSHARE_EM,
**overrides,
) -> UnifiedRealtimeQuote:
return UnifiedRealtimeQuote(
code=code,
name=name,
source=RealtimeSource.AKSHARE_EM,
source=source,
price=1688.0,
change_pct=1.2,
**overrides,
)
@@ -106,10 +112,84 @@ def test_manager_does_not_warn_when_fallback_source_succeeds(mock_get_config, ca
assert quote is not None
assert quote.name == "贵州茅台"
assert quote.fetched_at is not None
assert quote.fallback_from == "efinance"
assert not [record for record in caplog.records if record.levelno >= logging.WARNING]
assert "所有数据源均不可用" not in caplog.text
@patch("src.config.get_config")
def test_manager_supplement_does_not_mark_fallback_from(mock_get_config):
mock_get_config.return_value = SimpleNamespace(
enable_realtime_quote=True,
realtime_source_priority="efinance,akshare_em",
)
primary = _make_quote(source=RealtimeSource.EFINANCE)
supplement = _make_quote(source=RealtimeSource.AKSHARE_EM, volume_ratio=1.7)
manager = DataFetcherManager(
fetchers=[
_DummyFetcher("EfinanceFetcher", 0, result=primary),
_DummyFetcher("AkshareFetcher", 1, result=supplement),
]
)
quote = manager.get_realtime_quote("600519")
assert quote is primary
assert quote.fetched_at is not None
assert quote.fallback_from is None
assert quote.source == RealtimeSource.EFINANCE
assert quote.volume_ratio == 1.7
@patch("src.config.get_config")
def test_manager_fallback_from_records_highest_priority_failed_source(mock_get_config):
mock_get_config.return_value = SimpleNamespace(
enable_realtime_quote=True,
realtime_source_priority="efinance,tushare,akshare_em",
)
manager = DataFetcherManager(
fetchers=[
_DummyFetcher("EfinanceFetcher", 0, error=RuntimeError("efinance timeout")),
_DummyFetcher("TushareFetcher", 1, error=RuntimeError("tushare timeout")),
_DummyFetcher("AkshareFetcher", 2, result=_make_quote()),
]
)
quote = manager.get_realtime_quote("600519")
assert quote is not None
assert quote.source == RealtimeSource.AKSHARE_EM
assert quote.fallback_from == "efinance"
assert quote.fetched_at is not None
@patch("src.config.get_config")
def test_manager_drops_invalid_provider_timestamp_before_return(mock_get_config):
mock_get_config.return_value = SimpleNamespace(
enable_realtime_quote=True,
realtime_source_priority="efinance",
realtime_cache_ttl=600,
)
raw_quote = _make_quote(
source=RealtimeSource.EFINANCE,
provider_timestamp="not-a-date",
)
manager = DataFetcherManager(
fetchers=[
_DummyFetcher("EfinanceFetcher", 0, result=raw_quote),
]
)
quote = manager.get_realtime_quote("600519")
assert quote is raw_quote
assert quote.fetched_at is not None
assert quote.provider_timestamp is None
assert quote.stale_seconds is None
assert quote.is_stale is None
def test_pipeline_warns_once_when_all_realtime_sources_fail(caplog):
pipeline = _make_pipeline(enable_realtime_quote=True, realtime_quote=None)
+36 -1
View File
@@ -5,7 +5,42 @@ import threading
import time
import unittest
from data_provider.realtime_types import CircuitBreaker
from data_provider.realtime_types import CircuitBreaker, RealtimeSource, UnifiedRealtimeQuote
class UnifiedRealtimeQuoteMetadataTestCase(unittest.TestCase):
def test_metadata_defaults_are_filtered_from_to_dict(self):
quote = UnifiedRealtimeQuote(code="600519", source=RealtimeSource.AKSHARE_EM)
data = quote.to_dict()
self.assertEqual(data["code"], "600519")
self.assertEqual(data["source"], "akshare_em")
self.assertNotIn("fetched_at", data)
self.assertNotIn("provider_timestamp", data)
self.assertNotIn("is_stale", data)
self.assertNotIn("stale_seconds", data)
self.assertNotIn("fallback_from", data)
def test_metadata_is_included_when_present(self):
quote = UnifiedRealtimeQuote(
code="600519",
source=RealtimeSource.TENCENT,
price=1688.0,
fetched_at="2026-05-31T10:00:05+00:00",
provider_timestamp="2026-05-31T10:00:00+00:00",
is_stale=False,
stale_seconds=5,
fallback_from="efinance",
)
data = quote.to_dict()
self.assertEqual(data["fetched_at"], "2026-05-31T10:00:05+00:00")
self.assertEqual(data["provider_timestamp"], "2026-05-31T10:00:00+00:00")
self.assertIs(data["is_stale"], False)
self.assertEqual(data["stale_seconds"], 5)
self.assertEqual(data["fallback_from"], "efinance")
class CircuitBreakerConcurrencyTestCase(unittest.TestCase):