mirror of
https://github.com/ZhuLinsen/daily_stock_analysis.git
synced 2026-10-06 14:33:11 +08:00
feat: add data capability contract (#2289)
* feat: add data capability contract * fix(review-feedback-2289): preserve unknown status until availability is checked and Make * fix(review-feedback-2289): Aggregate daily quality across supported markets and add kline * fix(review-feedback-2289): Honor daily-source circuit breakers in quality selection and Make * fix(review-feedback-2289): Scope market-overview quality by market and Do not select a news * fix: align data capability with runtime routes * fix: align data capability runtime coverage * fix: align source and index capability routes * fix: preserve runtime capability uncertainty * fix: include US index capability routes * fix: align realtime and monitor capabilities * fix: remove unsupported Tushare index capability * fix: align capabilities with runtime routes * fix: model realtime and breaker routes * fix: align US realtime request priority * fix: align realtime circuit coverage * fix: filter unavailable daily priorities * fix: align US realtime capability claims * fix: align executable US data routes
This commit is contained in:
@@ -14,6 +14,7 @@ from api.v1.endpoints import (
|
||||
history,
|
||||
stocks,
|
||||
backtest,
|
||||
data,
|
||||
system_config,
|
||||
auth,
|
||||
agent,
|
||||
@@ -29,6 +30,7 @@ __all__ = [
|
||||
"history",
|
||||
"stocks",
|
||||
"backtest",
|
||||
"data",
|
||||
"system_config",
|
||||
"auth",
|
||||
"agent",
|
||||
|
||||
@@ -0,0 +1,72 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Data capability and quality endpoints."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import logging
|
||||
|
||||
from fastapi import APIRouter, Depends, HTTPException, Request
|
||||
|
||||
from api.deps import get_config_dep
|
||||
from api.v1.schemas.common import ErrorResponse
|
||||
from api.v1.schemas.data_capability import DataCapabilityOverviewResponse
|
||||
from src.config import Config
|
||||
from src.services.data_capability_service import DataCapabilityService
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
|
||||
def _overview_response(config: Config, *, runtime_scheduler: object = None) -> DataCapabilityOverviewResponse:
|
||||
try:
|
||||
payload = DataCapabilityService(
|
||||
config=config,
|
||||
runtime_scheduler=runtime_scheduler,
|
||||
).get_overview()
|
||||
return DataCapabilityOverviewResponse.model_validate(payload)
|
||||
except Exception as exc: # noqa: BLE001 - keep diagnostics fail-open at API boundary.
|
||||
logger.error("Failed to build data capability overview: %s", exc, exc_info=True)
|
||||
raise HTTPException(
|
||||
status_code=500,
|
||||
detail={
|
||||
"error": "internal_error",
|
||||
"message": "Failed to build data capability overview",
|
||||
},
|
||||
)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/overview",
|
||||
response_model=DataCapabilityOverviewResponse,
|
||||
responses={500: {"model": ErrorResponse}},
|
||||
summary="Get data capability overview",
|
||||
description="Return provider capabilities, dataset quality, and source priority without exposing secrets.",
|
||||
)
|
||||
def get_data_overview(
|
||||
request: Request,
|
||||
config: Config = Depends(get_config_dep),
|
||||
) -> DataCapabilityOverviewResponse:
|
||||
"""Return the canonical read-only data overview."""
|
||||
return _overview_response(
|
||||
config,
|
||||
runtime_scheduler=getattr(request.app.state, "runtime_scheduler_service", None),
|
||||
)
|
||||
|
||||
|
||||
@router.get(
|
||||
"/capabilities",
|
||||
response_model=DataCapabilityOverviewResponse,
|
||||
responses={500: {"model": ErrorResponse}},
|
||||
summary="Get data provider capabilities",
|
||||
description="Alias of /data/overview for clients that only need capability metadata.",
|
||||
)
|
||||
def get_data_capabilities(
|
||||
request: Request,
|
||||
config: Config = Depends(get_config_dep),
|
||||
) -> DataCapabilityOverviewResponse:
|
||||
"""Return the data overview under the capability-oriented alias."""
|
||||
return _overview_response(
|
||||
config,
|
||||
runtime_scheduler=getattr(request.app.state, "runtime_scheduler_service", None),
|
||||
)
|
||||
@@ -18,6 +18,7 @@ from api.v1.endpoints import (
|
||||
analysis,
|
||||
auth,
|
||||
backtest,
|
||||
data,
|
||||
decision_signals,
|
||||
health,
|
||||
history,
|
||||
@@ -104,6 +105,12 @@ router.include_router(
|
||||
tags=["Screening"]
|
||||
)
|
||||
|
||||
router.include_router(
|
||||
data.router,
|
||||
prefix="/data",
|
||||
tags=["Data"]
|
||||
)
|
||||
|
||||
router.include_router(
|
||||
intelligence.router,
|
||||
prefix="/intelligence",
|
||||
|
||||
@@ -47,6 +47,14 @@ from api.v1.schemas.backtest import (
|
||||
BacktestResultsResponse,
|
||||
PerformanceMetrics,
|
||||
)
|
||||
from api.v1.schemas.data_capability import (
|
||||
DataCapabilityOverviewResponse,
|
||||
DataDatasetQuality,
|
||||
DataPriorityView,
|
||||
DataProviderCapability,
|
||||
DatasetQualityStatus,
|
||||
ProviderCapabilityStatus,
|
||||
)
|
||||
from api.v1.schemas.system_config import (
|
||||
SystemConfigFieldSchema,
|
||||
SystemConfigCategorySchema,
|
||||
@@ -165,6 +173,13 @@ __all__ = [
|
||||
"BacktestResultItem",
|
||||
"BacktestResultsResponse",
|
||||
"PerformanceMetrics",
|
||||
# data capability
|
||||
"DataCapabilityOverviewResponse",
|
||||
"DataDatasetQuality",
|
||||
"DataPriorityView",
|
||||
"DataProviderCapability",
|
||||
"DatasetQualityStatus",
|
||||
"ProviderCapabilityStatus",
|
||||
# system config
|
||||
"SystemConfigFieldSchema",
|
||||
"SystemConfigCategorySchema",
|
||||
|
||||
@@ -0,0 +1,79 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Data source capability and dataset quality schemas."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
from typing import Any, Dict, List, Literal, Optional
|
||||
|
||||
from pydantic import BaseModel, Field
|
||||
|
||||
|
||||
ProviderCapabilityStatus = Literal[
|
||||
"ok",
|
||||
"partial",
|
||||
"unconfigured",
|
||||
"unavailable",
|
||||
"unknown",
|
||||
]
|
||||
DatasetQualityStatus = Literal[
|
||||
"ok",
|
||||
"degraded",
|
||||
"partial",
|
||||
"unconfigured",
|
||||
"unavailable",
|
||||
"unknown",
|
||||
"stale",
|
||||
]
|
||||
|
||||
|
||||
class DataProviderCapability(BaseModel):
|
||||
"""Visible capability metadata for one data provider."""
|
||||
|
||||
name: str = Field(..., description="Stable provider token")
|
||||
label: str = Field("", description="Human-readable provider name")
|
||||
enabled: bool = Field(..., description="Whether the provider is enabled for runtime routing")
|
||||
configured: bool = Field(..., description="Whether required configuration is present")
|
||||
status: ProviderCapabilityStatus = Field(..., description="Configuration/runtime availability summary")
|
||||
priority: Optional[int] = Field(None, description="Runtime fetcher priority when available")
|
||||
markets: List[str] = Field(default_factory=list, description="Supported markets")
|
||||
datasets: List[str] = Field(default_factory=list, description="Supported dataset identifiers")
|
||||
dataset_markets: Dict[str, List[str]] = Field(
|
||||
default_factory=dict,
|
||||
description="Exact supported markets for each dataset; consumers must not infer a markets × datasets cross-product",
|
||||
)
|
||||
warnings: List[str] = Field(default_factory=list, description="Stable warning codes")
|
||||
last_error: Optional[str] = Field(None, description="Last known non-sensitive error summary")
|
||||
cooldown: Optional[bool] = Field(None, description="Whether the provider is currently in cooldown")
|
||||
|
||||
|
||||
class DataDatasetQuality(BaseModel):
|
||||
"""Dataset-level quality view consumed by dashboards and data center."""
|
||||
|
||||
dataset: str = Field(..., description="Stable dataset identifier")
|
||||
status: DatasetQualityStatus = Field(..., description="Current quality status")
|
||||
source: Optional[str] = Field(None, description="Selected source token when known")
|
||||
stale: Optional[bool] = Field(None, description="Whether the data is known stale")
|
||||
last_success: Optional[str] = Field(None, description="Last successful load timestamp")
|
||||
last_error: Optional[str] = Field(None, description="Last non-sensitive error summary")
|
||||
fallback_from: List[str] = Field(default_factory=list, description="Earlier priority sources skipped before selected source")
|
||||
coverage: Optional[Dict[str, Any]] = Field(None, description="Optional dataset coverage summary")
|
||||
warnings: List[str] = Field(default_factory=list, description="Stable warning codes")
|
||||
|
||||
|
||||
class DataPriorityView(BaseModel):
|
||||
"""Configured provider order for one usage scenario."""
|
||||
|
||||
scenario: str = Field(..., description="Stable scenario identifier")
|
||||
providers: List[str] = Field(default_factory=list, description="Configured provider/source tokens")
|
||||
source: str = Field("", description="Where this priority list comes from")
|
||||
warnings: List[str] = Field(default_factory=list, description="Stable warning codes")
|
||||
|
||||
|
||||
class DataCapabilityOverviewResponse(BaseModel):
|
||||
"""Read-only data capability and dataset quality overview."""
|
||||
|
||||
as_of: str
|
||||
providers: List[DataProviderCapability] = Field(default_factory=list)
|
||||
datasets: List[DataDatasetQuality] = Field(default_factory=list)
|
||||
priorities: List[DataPriorityView] = Field(default_factory=list)
|
||||
warnings: List[str] = Field(default_factory=list)
|
||||
@@ -31,7 +31,7 @@ from tenacity import (
|
||||
before_sleep_log,
|
||||
)
|
||||
|
||||
from .base import BaseFetcher, DataFetchError, STANDARD_COLUMNS, is_bse_code
|
||||
from .base import BaseFetcher, DataFetchError, STANDARD_COLUMNS, _is_hk_market, is_bse_code
|
||||
from .realtime_types import UnifiedRealtimeQuote, RealtimeSource
|
||||
from .us_index_mapping import get_us_index_yf_symbol, is_us_stock_code
|
||||
from .yfinance_fundamental_adapter import _safe_float
|
||||
@@ -793,7 +793,7 @@ class YfinanceFetcher(BaseFetcher):
|
||||
|
||||
def get_realtime_quote(self, stock_code: str) -> Optional[UnifiedRealtimeQuote]:
|
||||
"""
|
||||
获取美股/美股指数实时行情数据
|
||||
获取美股、港股、日韩台股票或美股指数实时行情数据
|
||||
|
||||
支持美股股票(AAPL、TSLA)和美股指数(SPX、DJI 等)。
|
||||
数据来源:yfinance Ticker.info
|
||||
@@ -815,13 +815,14 @@ class YfinanceFetcher(BaseFetcher):
|
||||
index_name=index_name,
|
||||
)
|
||||
|
||||
# 仅处理美股股票或 JP/KR/TW suffix-only 股票
|
||||
# 仅处理美股、港股或 JP/KR/TW suffix-only 股票
|
||||
if not (
|
||||
self._is_us_stock(stock_code)
|
||||
or _is_hk_market(stock_code)
|
||||
or self._is_jp_kr_suffix_stock(stock_code)
|
||||
or self._is_tw_suffix_stock(stock_code)
|
||||
):
|
||||
logger.debug(f"[Yfinance] {stock_code} 不是美股或日韩 suffix 代码,跳过")
|
||||
logger.debug(f"[Yfinance] {stock_code} 不是支持的美股、港股或日韩台代码,跳过")
|
||||
return None
|
||||
|
||||
try:
|
||||
@@ -911,7 +912,7 @@ class YfinanceFetcher(BaseFetcher):
|
||||
code=symbol,
|
||||
name=name,
|
||||
source=RealtimeSource.FALLBACK,
|
||||
market=suffix_market or ("us" if is_us_symbol else None),
|
||||
market=suffix_market or ("hk" if _is_hk_market(stock_code) else "us" if is_us_symbol else None),
|
||||
currency=str(ticker_info.get("currency") or "").upper() or None,
|
||||
data_quality="partial" if missing_fields else "ok",
|
||||
missing_fields=missing_fields or None,
|
||||
|
||||
@@ -14,6 +14,7 @@ and this project adheres to [Semantic Versioning](https://semver.org/).
|
||||
- [修复] 交易日过滤对市场未知的指数 code(如 `sh000016`/`csi930955`/`930955.CSI`)按 `market=cn` 参与 A 股休市过滤,避免休市日指数被 fail-open 保留;市场仍未知的非指数 code 继续保留。
|
||||
- [修复] 指数分析将实际命中的日线数据源归因保存到历史记录,并由 Dashboard/Brief aggregate 报告展示;来源无效时保持原有输出。
|
||||
- [文档] 在中英繁 README 顶部关联 DSA arXiv 论文,并新增 `CITATION.cff` 统一项目引用信息。
|
||||
- [新功能] 新增数据源能力与数据集质量只读契约,提供 `/api/v1/data/overview` 和 `/api/v1/data/capabilities`,为大盘看板、数据中心、个股详情和自选 2.0 统一暴露 provider capability、dataset quality 和 source priority。
|
||||
- [改进] PR CI 增加文档路径检测:仅修改普通文档、非治理 Markdown 或 LICENSE 时跳过后端测试分片、Docker、Web 与桌面打包,保留轻量治理和门禁汇总;契约文档、静态 API 规格与测试 fixture 仍执行后端回归。
|
||||
- [修复] Linux/Docker 分享图补齐 Noto CJK 字体与中韩文字体栈,避免 PNG 只显示数字和英文、中文或韩文内容消失。
|
||||
- [新功能] Web Chat 意图识别层新增分词模块:`web_intent_tokenizer` 六步管道(多股票全名实体扫描 → 标点/空白切分 → 代码形提取 → 市场关键词 → 无歧义关键词 → 残存 gap 多策略 DFS 匹配)把用户消息切分为携带语义标签的 Token 序列;配套 `web_intent_types` 数据字典(Token 结构、Market 枚举、21 个语义 tag、clean/extend 双词池与正则机器)。核心原则"宁可不做,不可做错":Step 1~5 只做精确匹配,Step 6 要求整段 TAG 全覆盖(交叉验证)才产出,未覆盖片段保持空 tag 交下游 LLM 兜底;代码形 token 辨认为 `stock_code`(附 code/name/market 三元组)/ `wrong_{market}_code` / `unknown_{market}_code` 三态,token 层代码拼写统一 canonical 归一(a=6 位裸数字、hk=HK+5 位、us=大写 ticker)。意图枚举与意图识别结果随后续 `web_intent_resolver` PR 引入。新增 183 个分词单元测试。
|
||||
|
||||
@@ -2672,6 +2672,70 @@
|
||||
}
|
||||
]
|
||||
}
|
||||
},
|
||||
"/api/v1/data/overview": {
|
||||
"get": {
|
||||
"tags": [
|
||||
"Data"
|
||||
],
|
||||
"summary": "Get data capability overview",
|
||||
"description": "Return provider capabilities, dataset quality, and source priority without exposing secrets.",
|
||||
"operationId": "get_data_overview_api_v1_data_overview_get",
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "Successful Response",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/DataCapabilityOverviewResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"500": {
|
||||
"description": "Internal Server Error",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"/api/v1/data/capabilities": {
|
||||
"get": {
|
||||
"tags": [
|
||||
"Data"
|
||||
],
|
||||
"summary": "Get data provider capabilities",
|
||||
"description": "Alias of /data/overview for clients that only need capability metadata.",
|
||||
"operationId": "get_data_capabilities_api_v1_data_capabilities_get",
|
||||
"responses": {
|
||||
"200": {
|
||||
"description": "Successful Response",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/DataCapabilityOverviewResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"500": {
|
||||
"description": "Internal Server Error",
|
||||
"content": {
|
||||
"application/json": {
|
||||
"schema": {
|
||||
"$ref": "#/components/schemas/ErrorResponse"
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
},
|
||||
"components": {
|
||||
@@ -7438,6 +7502,309 @@
|
||||
"metadata"
|
||||
],
|
||||
"title": "DecisionSignalPreview"
|
||||
},
|
||||
"DataCapabilityOverviewResponse": {
|
||||
"properties": {
|
||||
"as_of": {
|
||||
"type": "string",
|
||||
"title": "As Of"
|
||||
},
|
||||
"providers": {
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/DataProviderCapability"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Providers"
|
||||
},
|
||||
"datasets": {
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/DataDatasetQuality"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Datasets"
|
||||
},
|
||||
"priorities": {
|
||||
"items": {
|
||||
"$ref": "#/components/schemas/DataPriorityView"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Priorities"
|
||||
},
|
||||
"warnings": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Warnings"
|
||||
}
|
||||
},
|
||||
"type": "object",
|
||||
"required": [
|
||||
"as_of"
|
||||
],
|
||||
"title": "DataCapabilityOverviewResponse",
|
||||
"description": "Read-only data capability and dataset quality overview."
|
||||
},
|
||||
"DataDatasetQuality": {
|
||||
"properties": {
|
||||
"dataset": {
|
||||
"type": "string",
|
||||
"title": "Dataset",
|
||||
"description": "Stable dataset identifier"
|
||||
},
|
||||
"status": {
|
||||
"type": "string",
|
||||
"enum": [
|
||||
"ok",
|
||||
"degraded",
|
||||
"partial",
|
||||
"unconfigured",
|
||||
"unavailable",
|
||||
"unknown",
|
||||
"stale"
|
||||
],
|
||||
"title": "Status",
|
||||
"description": "Current quality status"
|
||||
},
|
||||
"source": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "string"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Source",
|
||||
"description": "Selected source token when known"
|
||||
},
|
||||
"stale": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "boolean"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Stale",
|
||||
"description": "Whether the data is known stale"
|
||||
},
|
||||
"last_success": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "string"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Last Success",
|
||||
"description": "Last successful load timestamp"
|
||||
},
|
||||
"last_error": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "string"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Last Error",
|
||||
"description": "Last non-sensitive error summary"
|
||||
},
|
||||
"fallback_from": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Fallback From",
|
||||
"description": "Earlier priority sources skipped before selected source"
|
||||
},
|
||||
"coverage": {
|
||||
"anyOf": [
|
||||
{
|
||||
"additionalProperties": true,
|
||||
"type": "object"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Coverage",
|
||||
"description": "Optional dataset coverage summary"
|
||||
},
|
||||
"warnings": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Warnings",
|
||||
"description": "Stable warning codes"
|
||||
}
|
||||
},
|
||||
"type": "object",
|
||||
"required": [
|
||||
"dataset",
|
||||
"status"
|
||||
],
|
||||
"title": "DataDatasetQuality",
|
||||
"description": "Dataset-level quality view consumed by dashboards and data center."
|
||||
},
|
||||
"DataPriorityView": {
|
||||
"properties": {
|
||||
"scenario": {
|
||||
"type": "string",
|
||||
"title": "Scenario",
|
||||
"description": "Stable scenario identifier"
|
||||
},
|
||||
"providers": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Providers",
|
||||
"description": "Configured provider/source tokens"
|
||||
},
|
||||
"source": {
|
||||
"type": "string",
|
||||
"title": "Source",
|
||||
"description": "Where this priority list comes from",
|
||||
"default": ""
|
||||
},
|
||||
"warnings": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Warnings",
|
||||
"description": "Stable warning codes"
|
||||
}
|
||||
},
|
||||
"type": "object",
|
||||
"required": [
|
||||
"scenario"
|
||||
],
|
||||
"title": "DataPriorityView",
|
||||
"description": "Configured provider order for one usage scenario."
|
||||
},
|
||||
"DataProviderCapability": {
|
||||
"properties": {
|
||||
"name": {
|
||||
"type": "string",
|
||||
"title": "Name",
|
||||
"description": "Stable provider token"
|
||||
},
|
||||
"label": {
|
||||
"type": "string",
|
||||
"title": "Label",
|
||||
"description": "Human-readable provider name",
|
||||
"default": ""
|
||||
},
|
||||
"enabled": {
|
||||
"type": "boolean",
|
||||
"title": "Enabled",
|
||||
"description": "Whether the provider is enabled for runtime routing"
|
||||
},
|
||||
"configured": {
|
||||
"type": "boolean",
|
||||
"title": "Configured",
|
||||
"description": "Whether required configuration is present"
|
||||
},
|
||||
"status": {
|
||||
"type": "string",
|
||||
"enum": [
|
||||
"ok",
|
||||
"partial",
|
||||
"unconfigured",
|
||||
"unavailable",
|
||||
"unknown"
|
||||
],
|
||||
"title": "Status",
|
||||
"description": "Configuration/runtime availability summary"
|
||||
},
|
||||
"priority": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "integer"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Priority",
|
||||
"description": "Runtime fetcher priority when available"
|
||||
},
|
||||
"markets": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Markets",
|
||||
"description": "Supported markets"
|
||||
},
|
||||
"datasets": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Datasets",
|
||||
"description": "Supported dataset identifiers"
|
||||
},
|
||||
"dataset_markets": {
|
||||
"additionalProperties": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array"
|
||||
},
|
||||
"type": "object",
|
||||
"title": "Dataset Markets",
|
||||
"description": "Exact supported markets for each dataset; consumers must not infer a markets × datasets cross-product"
|
||||
},
|
||||
"warnings": {
|
||||
"items": {
|
||||
"type": "string"
|
||||
},
|
||||
"type": "array",
|
||||
"title": "Warnings",
|
||||
"description": "Stable warning codes"
|
||||
},
|
||||
"last_error": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "string"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Last Error",
|
||||
"description": "Last known non-sensitive error summary"
|
||||
},
|
||||
"cooldown": {
|
||||
"anyOf": [
|
||||
{
|
||||
"type": "boolean"
|
||||
},
|
||||
{
|
||||
"type": "null"
|
||||
}
|
||||
],
|
||||
"title": "Cooldown",
|
||||
"description": "Whether the provider is currently in cooldown"
|
||||
}
|
||||
},
|
||||
"type": "object",
|
||||
"required": [
|
||||
"name",
|
||||
"enabled",
|
||||
"configured",
|
||||
"status"
|
||||
],
|
||||
"title": "DataProviderCapability",
|
||||
"description": "Visible capability metadata for one data provider."
|
||||
}
|
||||
},
|
||||
"securitySchemes": {
|
||||
|
||||
@@ -12,7 +12,7 @@
|
||||
- 已登记 A 股指数:固定按 Tencent → AkShare → TickFlow → YFinance 降级,不读取普通日 K 的 `*_PRIORITY` 配置。
|
||||
- A 股大盘复盘:配置 `TICKFLOW_API_KEY` 后,复盘聚合所需的指数和市场宽度会优先尝试 TickFlow,失败后回退现有免费源;这与单标的指数日线的 Tencent-first 固定链是不同入口。
|
||||
- 港股:配置 `FUTU_OPEND_HOST` 后,Futu 可作为港股实时与基本面主源;`FUTU_HK_REALTIME_SOURCE_PRIORITY` 控制港股实时行情顺序,Longbridge、AkShare、YFinance 保留为 fallback。
|
||||
- 美股:配置 `LONGBRIDGE_*` 后优先使用 Longbridge,YFinance、Finnhub、AlphaVantage 继续兜底。
|
||||
- 美股:配置 `LONGBRIDGE_*` 后可优先使用 Longbridge,YFinance 保留实时主链 fallback;Finnhub、AlphaVantage 只在已有主源报价后补充字段或参与日线链,不能在 Longbridge/YFinance 全部失败时独立返回实时报价。
|
||||
- 热点题材:选股的热点实现参考 AlphaSift,默认走 EastMoney provider,并使用本地 last-good cache 降低实时接口失败影响。
|
||||
|
||||
## 已接入数据源矩阵
|
||||
@@ -27,7 +27,20 @@
|
||||
| 选股日线补特征 | `DataFetcherManager` | 选股引擎优先复用现有日线与缓存链路 | 现有链路失败后才回到引擎自身的日线源 |
|
||||
| 选股热点题材 | EastMoney provider、参考 AlphaSift 的 hotspot 实现、last-good cache | 未指定 provider 时默认使用 EastMoney provider | 实时失败时回退热点缓存;无缓存时返回稳定空态和可读错误 |
|
||||
| 港股 | Futu、Longbridge、YFinance、AkShare、Tushare | 配置 `FUTU_OPEND_HOST` 后 Futu 作为实时与基本面主源,按 `FUTU_HK_REALTIME_SOURCE_PRIORITY` 顺序尝试 | Futu 失败时回退 Longbridge / AkShare / YFinance;Longbridge 冷却或失败时继续回退 YFinance / 其他可用源 |
|
||||
| 美股 | Longbridge、YFinance、AkShare、Tushare、Finnhub、AlphaVantage、Stooq | 配置 Longbridge 凭证后参与美股日线/实时兜底;YFinance 保持基础兜底 | Longbridge 冷却或失败时回退 YFinance / 其他可用源 |
|
||||
| 美股 | Longbridge、YFinance、AkShare、Tushare、Finnhub、AlphaVantage、Stooq | 配置 Longbridge 凭证后参与美股日线/实时主链;YFinance 保持实时基础 fallback;Finnhub/AlphaVantage 可补充已有报价字段并参与日线链 | Longbridge 冷却或失败时实时主链回退 YFinance;两者都失败时不把补充源伪装成可独立返回报价的 fallback |
|
||||
|
||||
## 统一数据能力只读契约
|
||||
|
||||
Web/API 提供 `GET /api/v1/data/overview` 和等价别名 `GET /api/v1/data/capabilities`,用于把数据源能力、数据集质量和场景优先级以同一个只读结构暴露给首页看板、数据中心、个股详情、自选和选股页面。
|
||||
|
||||
首版契约只读取配置和 `DataFetcherManager` 的 fetcher 快照,不触发外部行情请求,也不返回任何原始密钥。响应分三层:
|
||||
|
||||
- `providers`:每个 provider 的 `enabled`、`configured`、`status`、`markets`、`datasets`、精确的 `dataset_markets` 和非敏感 warning。`markets` 与 `datasets` 只是聚合索引,不能推断为笛卡尔积;例如 YFinance 可参与 A 股行情,但 A 股基本面固定走 AkShare,因此 `dataset_markets.financial.snapshot` 不包含 `cn`。Finnhub/AlphaVantage 的美股报价 handler 当前只在 Longbridge/YFinance 已返回主报价后补字段,因此不声明 `quote.realtime` 独立路由能力。
|
||||
- `datasets`:`quote.realtime`、`kline.daily`、`index.daily`、`market.overview`、`financial.snapshot`、`news.events`、`strategy.screening`、`alert.monitor`、`portfolio.account` 的 `status/source/fallback_from/stale/warnings`;其中 `quote.realtime`、`kline.daily`、`market.overview` 和 `financial.snapshot` 都聚合市场级 coverage,避免把单市场健康度误报成全局可用。A 股实时额外拆分 `cn` 股票优先级、`cn.index.exchange` 固定指数链和 `cn.index.csi` Efinance-only 链;美股实时拆分 `us` 个股动态优先级与 `us.index` 固定 YFinance-first 链。AkShare Tencent/Sina/EM、AkShare HK 双子源、Efinance 股票/指数子源级熔断及日线 breaker 会先于 provider-wide unknown 判定,使 overview 与运行时跳过行为一致。YFinance 的港股实时声明对应其实际 `.HK` 执行路径;未配置 OpenD 时 HK priority 与运行时一致地跳过 Futu;美股个股实时也按当前请求可用性跳过处于连接冷却的 Longbridge,使 YFinance 成为实际主源而不是伪 fallback。PyTDX 不声明未接入统一 route 的实时能力,YFinance 指数日线只声明实际可执行的 CN/US route。任何按运行时顺序尝试的数据集在首个优先源尚未探测且没有更具体的 open breaker 证据时保持 `unknown`,不会越过它宣称后续源已被选中。`alert.monitor` 只有在开关启用且当前 API 进程的 scheduler 已实际注册并运行 `agent_event_monitor` 时才为 `ok`;开关关闭返回 `agent_event_monitor_disabled`,scheduler 未运行则返回 `agent_event_monitor_not_running`。
|
||||
- `priorities`:`cn.realtime`、`hk.realtime`、`us.realtime`、`daily.generic`、`cn.index.daily`、`market.overview`、`screening.snapshot`、`news.events` 的当前 source order;`daily.generic` 与运行时一样先排除当前请求不可用或处于连接冷却的 fetcher,不把已跳过的数据源记录成 fallback。
|
||||
- `index.daily` 额外区分 `cn.exchange`、`cn.csi` 与 `us` coverage:沪深交易所指数使用 Tencent、AkShare、TickFlow、YFinance 固定多源链,CSI 指数按运行时契约只认 AkShare,美股指数只声明当前能够接收指数 symbol 的 YFinance;Finnhub 当前 fetcher 会拒绝指数 symbol,因此不作为伪 fallback。Tushare 不声明指数日线能力。CN/HK realtime 都只接受对应运行时实际有 handler 的 source token,Tushare realtime 能力仅声明 A 股。
|
||||
|
||||
其中 `unknown` 表示尚未执行运行态探测,`unconfigured` 表示缺少必要配置,`degraded` 表示前置优先源不可用但后续源仍可消费。选股 source health 在冷启动、成功数和失败数都为 0 时保持 `unknown`,不会把尚未调用的数据源提前标成 `ok`。后续 Data Center 可以在同一结构上追加 last_success、coverage、cooldown 和运行历史,不需要各页面各自维护数据质量口径。
|
||||
|
||||
## 总体链路图
|
||||
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
@@ -644,6 +644,23 @@ class RuntimeSchedulerService:
|
||||
}
|
||||
return {"accepted": True, "running": True}
|
||||
|
||||
def is_background_task_active(self, name: str) -> bool:
|
||||
"""Return whether a named task is registered on the live scheduler."""
|
||||
with self._lock:
|
||||
scheduler = self._scheduler
|
||||
thread = self._thread
|
||||
if (
|
||||
not self._enabled
|
||||
or scheduler is None
|
||||
or thread is None
|
||||
or not thread.is_alive()
|
||||
):
|
||||
return False
|
||||
return any(
|
||||
str(entry.get("name") or "") == name
|
||||
for entry in getattr(scheduler, "_background_tasks", [])
|
||||
)
|
||||
|
||||
def status(self) -> Dict[str, Any]:
|
||||
scheduler = self._scheduler
|
||||
jobs = scheduler.schedule.get_jobs() if scheduler is not None else []
|
||||
|
||||
@@ -61,6 +61,16 @@ P6_SIGNAL_LINKED_SCHEMAS = (
|
||||
"PortfolioDecisionSignalRiskItem",
|
||||
"PortfolioRiskResponse",
|
||||
)
|
||||
DATA_CAPABILITY_PATHS = (
|
||||
"/api/v1/data/overview",
|
||||
"/api/v1/data/capabilities",
|
||||
)
|
||||
DATA_CAPABILITY_SCHEMAS = (
|
||||
"DataCapabilityOverviewResponse",
|
||||
"DataDatasetQuality",
|
||||
"DataPriorityView",
|
||||
"DataProviderCapability",
|
||||
)
|
||||
|
||||
|
||||
def _collect_component_schema_refs(node: Any) -> set[str]:
|
||||
@@ -231,6 +241,10 @@ def test_decision_signal_static_api_spec_matches_runtime_paths() -> None:
|
||||
assert static_spec["paths"][path] == runtime_spec["paths"][path]
|
||||
for schema_name in P6_SIGNAL_LINKED_SCHEMAS:
|
||||
assert static_spec["components"]["schemas"][schema_name] == runtime_spec["components"]["schemas"][schema_name]
|
||||
for path in DATA_CAPABILITY_PATHS:
|
||||
assert static_spec["paths"][path] == runtime_spec["paths"][path]
|
||||
for schema_name in DATA_CAPABILITY_SCHEMAS:
|
||||
assert static_spec["components"]["schemas"][schema_name] == runtime_spec["components"]["schemas"][schema_name]
|
||||
schema_refs = _collect_component_schema_refs(static_spec)
|
||||
missing_schema_refs = sorted(schema_refs - set(static_spec["components"]["schemas"]))
|
||||
assert missing_schema_refs == []
|
||||
|
||||
@@ -0,0 +1,993 @@
|
||||
# -*- coding: utf-8 -*-
|
||||
"""Tests for the data capability overview service."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import json
|
||||
import tempfile
|
||||
from pathlib import Path
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import patch
|
||||
|
||||
from fastapi.testclient import TestClient
|
||||
|
||||
from api.app import create_app
|
||||
from src.services.data_capability_service import DataCapabilityService
|
||||
|
||||
|
||||
class _Fetcher:
|
||||
def __init__(
|
||||
self,
|
||||
name: str,
|
||||
priority: int,
|
||||
*,
|
||||
available=None,
|
||||
last_error: str = "",
|
||||
is_available=None,
|
||||
is_available_for_request=None,
|
||||
) -> None:
|
||||
self.name = name
|
||||
self.priority = priority
|
||||
if available is not None:
|
||||
self._available = available
|
||||
if last_error:
|
||||
self.last_error = last_error
|
||||
if is_available is not None:
|
||||
self.is_available = lambda: is_available
|
||||
if is_available_for_request is not None:
|
||||
self.is_available_for_request = lambda _capability="": is_available_for_request
|
||||
|
||||
|
||||
class _FetcherManager:
|
||||
def __init__(self, fetchers) -> None:
|
||||
self._fetchers = fetchers
|
||||
|
||||
def _get_fetchers_snapshot(self):
|
||||
return list(self._fetchers)
|
||||
|
||||
|
||||
class _RuntimeScheduler:
|
||||
def __init__(self, active: bool) -> None:
|
||||
self.active = active
|
||||
|
||||
def is_background_task_active(self, name: str) -> bool:
|
||||
return self.active and name == "agent_event_monitor"
|
||||
|
||||
|
||||
def _config(**overrides):
|
||||
values = {
|
||||
"tushare_token": None,
|
||||
"tickflow_api_key": None,
|
||||
"tickflow_priority": 2,
|
||||
"futu_opend_host": None,
|
||||
"longbridge_app_key": None,
|
||||
"longbridge_app_secret": None,
|
||||
"longbridge_access_token": None,
|
||||
"longbridge_oauth_client_id": None,
|
||||
"finnhub_api_key": None,
|
||||
"alphavantage_api_key": None,
|
||||
"enable_realtime_quote": True,
|
||||
"enable_fundamental_pipeline": True,
|
||||
"realtime_source_priority": "tencent,akshare_sina,efinance,akshare_em",
|
||||
"futu_hk_realtime_source_priority": "futu,longbridge,akshare,yfinance",
|
||||
"screening_enabled": False,
|
||||
"agent_event_monitor_enabled": False,
|
||||
}
|
||||
values.update(overrides)
|
||||
return SimpleNamespace(**values)
|
||||
|
||||
|
||||
def _dataset(overview, name: str):
|
||||
return next(item for item in overview["datasets"] if item["dataset"] == name)
|
||||
|
||||
|
||||
def _provider(overview, name: str):
|
||||
return next(item for item in overview["providers"] if item["name"] == name)
|
||||
|
||||
|
||||
def test_overview_marks_tickflow_priority_gap_without_leaking_secret() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
_Fetcher("EfinanceFetcher", 3, available=True),
|
||||
_Fetcher("TickFlowFetcher", 2),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
secret = "tickflow-secret-value"
|
||||
service = DataCapabilityService(
|
||||
config=_config(tickflow_api_key=secret, realtime_source_priority="tencent,efinance"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
|
||||
assert _provider(overview, "tickflow")["configured"] is True
|
||||
assert _provider(overview, "tickflow")["status"] == "unknown"
|
||||
assert _provider(overview, "tickflow")["warnings"] == ["runtime_probe_not_performed"]
|
||||
assert "tickflow_configured_but_not_in_realtime_priority" in overview["warnings"]
|
||||
assert secret not in json.dumps(overview, ensure_ascii=False)
|
||||
|
||||
|
||||
def test_realtime_dataset_degrades_when_first_priority_source_is_unconfigured() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True),
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(realtime_source_priority="tushare,efinance"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
quote_quality = _dataset(overview, "quote.realtime")
|
||||
|
||||
assert quote_quality["status"] == "degraded"
|
||||
assert quote_quality["source"] is None
|
||||
assert quote_quality["fallback_from"] == ["cn:tushare", "hk:longbridge"]
|
||||
assert quote_quality["coverage"]["markets"]["cn"]["status"] == "degraded"
|
||||
assert quote_quality["coverage"]["markets"]["cn"]["source"] == "efinance"
|
||||
assert "cn:source_status:tushare:unconfigured" in quote_quality["warnings"]
|
||||
|
||||
|
||||
def test_provider_runtime_probe_preserves_unknown_until_checked() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("TickFlowFetcher", 1),
|
||||
_Fetcher("TushareFetcher", 2, is_available=False),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(tickflow_api_key="secret", tushare_token="token"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
|
||||
assert _provider(overview, "tickflow")["status"] == "unknown"
|
||||
assert _provider(overview, "tushare")["status"] == "unavailable"
|
||||
|
||||
|
||||
def test_provider_runtime_probe_honors_request_time_unavailable_over_cached_available_flag() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("LongbridgeFetcher", 1, available=True, is_available_for_request=False),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(longbridge_app_key="key"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
|
||||
assert _provider(overview, "longbridge")["status"] == "unavailable"
|
||||
assert _provider(overview, "longbridge")["warnings"] == ["provider_marked_unavailable"]
|
||||
|
||||
|
||||
def test_us_realtime_priority_skips_longbridge_during_request_cooldown() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("LongbridgeFetcher", 1, available=True, is_available_for_request=False),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(longbridge_app_key="key"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
us_quality = _dataset(overview, "quote.realtime")["coverage"]["markets"]["us"]
|
||||
|
||||
assert priorities["us.realtime"]["providers"] == ["yfinance", "longbridge"]
|
||||
assert us_quality["status"] == "ok"
|
||||
assert us_quality["source"] == "yfinance"
|
||||
assert us_quality["fallback_from"] == []
|
||||
|
||||
|
||||
def test_provider_dataset_market_matrix_matches_fundamental_runtime_routes() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(
|
||||
tushare_token="token",
|
||||
longbridge_app_key="key",
|
||||
futu_opend_host="127.0.0.1",
|
||||
),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
|
||||
assert _provider(overview, "akshare")["dataset_markets"]["financial.snapshot"] == ["cn"]
|
||||
assert _provider(overview, "futu")["dataset_markets"]["financial.snapshot"] == ["hk"]
|
||||
assert _provider(overview, "yfinance")["dataset_markets"]["financial.snapshot"] == [
|
||||
"hk",
|
||||
"us",
|
||||
"jp",
|
||||
"kr",
|
||||
"tw",
|
||||
]
|
||||
assert "financial.snapshot" not in _provider(overview, "tushare")["datasets"]
|
||||
assert "financial.snapshot" not in _provider(overview, "longbridge")["datasets"]
|
||||
assert _provider(overview, "tushare")["dataset_markets"]["quote.realtime"] == ["cn"]
|
||||
assert "index.daily" not in _provider(overview, "tushare")["datasets"]
|
||||
assert _provider(overview, "yfinance")["dataset_markets"]["quote.realtime"] == [
|
||||
"hk",
|
||||
"us",
|
||||
"jp",
|
||||
"kr",
|
||||
"tw",
|
||||
]
|
||||
assert _provider(overview, "yfinance")["dataset_markets"]["index.daily"] == [
|
||||
"cn",
|
||||
"us",
|
||||
]
|
||||
assert "quote.realtime" not in _provider(overview, "pytdx")["datasets"]
|
||||
assert "quote.realtime" not in _provider(overview, "finnhub")["datasets"]
|
||||
assert "quote.realtime" not in _provider(overview, "alphavantage")["datasets"]
|
||||
assert "index.daily" not in _provider(overview, "finnhub")["datasets"]
|
||||
|
||||
|
||||
def test_hk_realtime_priority_skips_futu_when_opend_is_unconfigured() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(futu_opend_host=None),
|
||||
fetcher_manager=_FetcherManager([_Fetcher("YfinanceFetcher", 4, available=True)]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
hk_quality = _dataset(overview, "quote.realtime")["coverage"]["markets"]["hk"]
|
||||
|
||||
assert priorities["hk.realtime"]["providers"] == ["longbridge", "akshare", "yfinance"]
|
||||
assert "futu" not in hk_quality["fallback_from"]
|
||||
|
||||
|
||||
def test_realtime_dataset_quality_is_market_aware() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(
|
||||
realtime_source_priority="efinance",
|
||||
futu_hk_realtime_source_priority="futu,longbridge",
|
||||
),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
quote_quality = _dataset(overview, "quote.realtime")
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
|
||||
assert priorities["us.realtime"]["providers"] == ["yfinance", "longbridge"]
|
||||
assert quote_quality["status"] == "partial"
|
||||
assert quote_quality["coverage"]["markets"]["cn"]["status"] == "ok"
|
||||
assert quote_quality["coverage"]["markets"]["cn"]["source"] == "efinance"
|
||||
assert quote_quality["coverage"]["markets"]["hk"]["status"] == "unavailable"
|
||||
assert quote_quality["coverage"]["markets"]["us"]["status"] == "ok"
|
||||
for market in ("jp", "kr", "tw"):
|
||||
assert quote_quality["coverage"]["markets"][market] == {
|
||||
"status": "ok",
|
||||
"source": "yfinance",
|
||||
"fallback_from": [],
|
||||
"warnings": [],
|
||||
}
|
||||
|
||||
|
||||
def test_cn_realtime_quality_honors_akshare_subsource_circuit_breaker() -> None:
|
||||
from data_provider.realtime_types import get_realtime_circuit_breaker
|
||||
|
||||
breaker = get_realtime_circuit_breaker()
|
||||
breaker.reset("akshare_tencent")
|
||||
try:
|
||||
for _ in range(3):
|
||||
breaker.record_failure("akshare_tencent", "test failure")
|
||||
service = DataCapabilityService(
|
||||
config=_config(realtime_source_priority="tencent,efinance"),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("AkshareFetcher", 0),
|
||||
_Fetcher("EfinanceFetcher", 1, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
quote_quality = _dataset(service.get_overview(), "quote.realtime")
|
||||
cn_quality = quote_quality["coverage"]["markets"]["cn"]
|
||||
|
||||
assert cn_quality["status"] == "degraded"
|
||||
assert cn_quality["source"] == "efinance"
|
||||
assert cn_quality["fallback_from"] == ["tencent"]
|
||||
assert cn_quality["warnings"] == ["source_status:tencent:cooldown"]
|
||||
finally:
|
||||
breaker.reset("akshare_tencent")
|
||||
|
||||
|
||||
def test_daily_circuit_open_precedes_unknown_provider_probe() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0),
|
||||
_Fetcher("PytdxFetcher", 1, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
with patch(
|
||||
"data_provider.base.DataFetcherManager._is_daily_source_available",
|
||||
side_effect=lambda fetcher, _market: fetcher.name != "EfinanceFetcher",
|
||||
):
|
||||
cn_quality = _dataset(service.get_overview(), "kline.daily")["coverage"]["markets"]["cn"]
|
||||
|
||||
assert cn_quality["status"] == "degraded"
|
||||
assert cn_quality["source"] == "pytdx"
|
||||
assert cn_quality["fallback_from"] == ["efinance"]
|
||||
assert cn_quality["warnings"] == ["source_status:efinance:cooldown"]
|
||||
|
||||
|
||||
def test_cn_realtime_reports_fixed_index_route_separately_from_stock_priority() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(
|
||||
tushare_token="token",
|
||||
realtime_source_priority="tushare",
|
||||
),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("TushareFetcher", 0, available=True),
|
||||
_Fetcher("EfinanceFetcher", 1, available=False),
|
||||
]),
|
||||
)
|
||||
|
||||
quote_quality = _dataset(service.get_overview(), "quote.realtime")
|
||||
markets = quote_quality["coverage"]["markets"]
|
||||
|
||||
assert markets["cn"]["status"] == "ok"
|
||||
assert markets["cn"]["source"] == "tushare"
|
||||
assert markets["cn.index.exchange"]["status"] == "unavailable"
|
||||
assert markets["cn.index.csi"]["status"] == "unavailable"
|
||||
assert quote_quality["status"] == "partial"
|
||||
|
||||
|
||||
def test_efinance_realtime_breakers_apply_to_stock_and_index_routes() -> None:
|
||||
from data_provider.realtime_types import get_realtime_circuit_breaker
|
||||
|
||||
breaker = get_realtime_circuit_breaker()
|
||||
try:
|
||||
for key in ("efinance", "efinance_index"):
|
||||
breaker.reset(key)
|
||||
for _ in range(3):
|
||||
breaker.record_failure(key, "test")
|
||||
service = DataCapabilityService(
|
||||
config=_config(realtime_source_priority="efinance,tencent"),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True),
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
markets = _dataset(service.get_overview(), "quote.realtime")["coverage"]["markets"]
|
||||
|
||||
assert markets["cn"]["source"] == "tencent"
|
||||
assert markets["cn"]["fallback_from"] == ["efinance"]
|
||||
assert markets["cn.index.csi"]["status"] == "unavailable"
|
||||
assert markets["cn.index.csi"]["warnings"] == ["source_status:efinance:cooldown"]
|
||||
finally:
|
||||
breaker.reset("efinance")
|
||||
breaker.reset("efinance_index")
|
||||
|
||||
|
||||
def test_hk_realtime_falls_back_only_when_both_akshare_routes_are_open() -> None:
|
||||
from data_provider.realtime_types import get_realtime_circuit_breaker
|
||||
|
||||
breaker = get_realtime_circuit_breaker()
|
||||
try:
|
||||
for key in ("akshare_hk_em", "akshare_hk_sina"):
|
||||
breaker.reset(key)
|
||||
for _ in range(3):
|
||||
breaker.record_failure(key, "test")
|
||||
service = DataCapabilityService(
|
||||
config=_config(futu_hk_realtime_source_priority="akshare,yfinance"),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
hk_quality = _dataset(service.get_overview(), "quote.realtime")["coverage"]["markets"]["hk"]
|
||||
|
||||
assert hk_quality["status"] == "degraded"
|
||||
assert hk_quality["source"] == "yfinance"
|
||||
assert hk_quality["fallback_from"] == ["akshare"]
|
||||
finally:
|
||||
breaker.reset("akshare_hk_em")
|
||||
breaker.reset("akshare_hk_sina")
|
||||
|
||||
|
||||
def test_us_index_realtime_keeps_yfinance_ahead_of_healthy_longbridge() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(longbridge_app_key="key"),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("LongbridgeFetcher", 1, available=True, is_available_for_request=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
markets = _dataset(service.get_overview(), "quote.realtime")["coverage"]["markets"]
|
||||
|
||||
assert markets["us"]["source"] == "longbridge"
|
||||
assert markets["us.index"]["source"] == "yfinance"
|
||||
assert markets["us.index"]["fallback_from"] == []
|
||||
|
||||
|
||||
def test_cn_realtime_rejects_tokens_without_runtime_handlers() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(realtime_source_priority="yfinance,efinance"),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("YfinanceFetcher", 0, available=True),
|
||||
_Fetcher("EfinanceFetcher", 1, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
cn_quality = _dataset(overview, "quote.realtime")["coverage"]["markets"]["cn"]
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
|
||||
assert priorities["cn.realtime"]["warnings"] == ["unknown_source:yfinance"]
|
||||
assert cn_quality["status"] == "degraded"
|
||||
assert cn_quality["source"] == "efinance"
|
||||
assert cn_quality["fallback_from"] == ["yfinance"]
|
||||
assert "source_status:yfinance:unsupported" in cn_quality["warnings"]
|
||||
|
||||
|
||||
def test_hk_realtime_rejects_configured_provider_without_runtime_handler() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(
|
||||
tushare_token="token",
|
||||
futu_hk_realtime_source_priority="tushare,yfinance",
|
||||
),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("TushareFetcher", 0, available=True),
|
||||
_Fetcher("YfinanceFetcher", 1, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
hk_quality = _dataset(overview, "quote.realtime")["coverage"]["markets"]["hk"]
|
||||
|
||||
assert hk_quality["status"] == "degraded"
|
||||
assert hk_quality["source"] == "yfinance"
|
||||
assert hk_quality["fallback_from"] == ["tushare"]
|
||||
assert "source_status:tushare:unsupported" in hk_quality["warnings"]
|
||||
|
||||
|
||||
def test_runtime_ordered_daily_route_stops_at_unprobed_preferred_source() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0),
|
||||
_Fetcher("PytdxFetcher", 1, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
cn_quality = _dataset(service.get_overview(), "kline.daily")["coverage"]["markets"]["cn"]
|
||||
|
||||
assert cn_quality["status"] == "unknown"
|
||||
assert cn_quality["source"] is None
|
||||
assert cn_quality["fallback_from"] == []
|
||||
assert cn_quality["warnings"] == ["source_status:efinance:unknown"]
|
||||
|
||||
|
||||
def test_daily_dataset_quality_is_market_aware() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=False),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
daily_quality = _dataset(overview, "kline.daily")
|
||||
|
||||
assert daily_quality["status"] == "partial"
|
||||
assert daily_quality["coverage"]["markets"]["cn"]["status"] == "ok"
|
||||
assert daily_quality["coverage"]["markets"]["cn"]["source"] == "efinance"
|
||||
assert daily_quality["coverage"]["markets"]["hk"]["status"] == "unavailable"
|
||||
assert daily_quality["coverage"]["markets"]["us"]["status"] == "unavailable"
|
||||
for market in ("jp", "kr", "tw"):
|
||||
assert daily_quality["coverage"]["markets"][market]["status"] == "unavailable"
|
||||
assert "hk:request_available_priority_empty" in daily_quality["warnings"]
|
||||
assert "us:request_available_priority_empty" in daily_quality["warnings"]
|
||||
|
||||
|
||||
def test_daily_dataset_quality_prefers_longbridge_for_us_when_available() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("LongbridgeFetcher", 5, available=True, is_available_for_request=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(longbridge_app_key="key"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
daily_quality = _dataset(overview, "kline.daily")
|
||||
|
||||
assert daily_quality["coverage"]["markets"]["us"]["status"] == "ok"
|
||||
assert daily_quality["coverage"]["markets"]["us"]["source"] == "longbridge"
|
||||
assert "us:finnhub" not in daily_quality["fallback_from"]
|
||||
|
||||
|
||||
def test_us_daily_priority_omits_unregistered_configured_sources() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(config=_config(), fetcher_manager=manager)
|
||||
|
||||
overview = service.get_overview()
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
us_quality = _dataset(overview, "kline.daily")["coverage"]["markets"]["us"]
|
||||
|
||||
assert priorities["daily.generic"]["providers"] == ["yfinance"]
|
||||
assert us_quality["status"] == "ok"
|
||||
assert us_quality["source"] == "yfinance"
|
||||
assert us_quality["fallback_from"] == []
|
||||
|
||||
|
||||
def test_daily_priority_excludes_request_unavailable_fetchers() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True, is_available_for_request=False),
|
||||
_Fetcher("PytdxFetcher", 1, available=True, is_available_for_request=True),
|
||||
]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
cn_daily = _dataset(overview, "kline.daily")["coverage"]["markets"]["cn"]
|
||||
|
||||
assert priorities["daily.generic"]["providers"] == ["pytdx"]
|
||||
assert cn_daily["source"] == "pytdx"
|
||||
assert cn_daily["fallback_from"] == []
|
||||
|
||||
|
||||
def test_daily_dataset_quality_honors_market_specific_circuit_breakers() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
with patch(
|
||||
"data_provider.base.DataFetcherManager._is_daily_source_available",
|
||||
side_effect=lambda fetcher, market: not (
|
||||
getattr(fetcher, "name", "") == "YfinanceFetcher" and market == "hk"
|
||||
),
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
daily_quality = _dataset(overview, "kline.daily")
|
||||
|
||||
assert daily_quality["status"] == "partial"
|
||||
assert daily_quality["coverage"]["markets"]["cn"]["status"] == "ok"
|
||||
assert daily_quality["coverage"]["markets"]["hk"]["status"] == "unavailable"
|
||||
assert daily_quality["coverage"]["markets"]["hk"]["source"] is None
|
||||
assert daily_quality["coverage"]["markets"]["us"]["status"] == "ok"
|
||||
assert daily_quality["coverage"]["markets"]["us"]["source"] == "yfinance"
|
||||
assert "hk:source_status:yfinance:cooldown" in daily_quality["warnings"]
|
||||
|
||||
|
||||
def test_index_daily_quality_uses_only_the_executable_us_runtime_route() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(finnhub_api_key="key"),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("TencentFetcher", 0, available=True),
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
_Fetcher("YfinanceFetcher", 2, available=False),
|
||||
_Fetcher("FinnhubFetcher", 3, available=True),
|
||||
]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
index_quality = _dataset(overview, "index.daily")
|
||||
us_quality = index_quality["coverage"]["markets"]["us"]
|
||||
|
||||
assert index_quality["status"] == "partial"
|
||||
assert us_quality["status"] == "unavailable"
|
||||
assert us_quality["source"] is None
|
||||
assert us_quality["fallback_from"] == ["yfinance"]
|
||||
assert "us:source_status:yfinance:unavailable" in index_quality["warnings"]
|
||||
assert not any("finnhub" in warning for warning in index_quality["warnings"])
|
||||
assert "index.daily" not in _provider(overview, "finnhub")["datasets"]
|
||||
|
||||
|
||||
def test_market_overview_dataset_quality_is_market_aware() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("EfinanceFetcher", 0, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
market_overview = _dataset(overview, "market.overview")
|
||||
|
||||
assert market_overview["status"] == "ok"
|
||||
assert market_overview["source"] is None
|
||||
assert market_overview["coverage"]["markets"]["cn"]["source"] == "efinance"
|
||||
assert market_overview["coverage"]["markets"]["hk"]["source"] == "yfinance"
|
||||
assert market_overview["coverage"]["markets"]["us"]["source"] == "yfinance"
|
||||
assert market_overview["coverage"]["markets"]["jp"]["source"] == "yfinance"
|
||||
assert market_overview["coverage"]["markets"]["kr"]["source"] == "yfinance"
|
||||
assert market_overview["coverage"]["markets"]["tw"]["source"] == "yfinance"
|
||||
|
||||
|
||||
def test_index_daily_quality_honors_cn_index_circuit_breakers() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("TencentFetcher", 0, available=True),
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
_Fetcher("YfinanceFetcher", 2, available=True),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
with patch(
|
||||
"data_provider.base.DataFetcherManager._is_daily_source_available",
|
||||
side_effect=lambda fetcher, market: not (
|
||||
getattr(fetcher, "name", "") == "TencentFetcher" and market == "cn_index"
|
||||
),
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
index_quality = _dataset(overview, "index.daily")
|
||||
|
||||
assert index_quality["status"] == "degraded"
|
||||
assert index_quality["source"] is None
|
||||
assert index_quality["coverage"]["markets"]["cn.exchange"]["status"] == "degraded"
|
||||
assert index_quality["coverage"]["markets"]["cn.exchange"]["fallback_from"] == ["tencent"]
|
||||
assert index_quality["coverage"]["markets"]["cn.csi"]["status"] == "ok"
|
||||
assert index_quality["coverage"]["markets"]["us"]["source"] == "yfinance"
|
||||
assert index_quality["fallback_from"] == ["cn.exchange:tencent"]
|
||||
assert index_quality["warnings"] == ["cn.exchange:source_status:tencent:cooldown"]
|
||||
|
||||
|
||||
def test_index_daily_requires_akshare_for_csi_family() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(),
|
||||
fetcher_manager=_FetcherManager([
|
||||
_Fetcher("TencentFetcher", 0, available=True),
|
||||
_Fetcher("AkshareFetcher", 1, available=False),
|
||||
]),
|
||||
)
|
||||
|
||||
index_quality = _dataset(service.get_overview(), "index.daily")
|
||||
|
||||
assert index_quality["status"] == "partial"
|
||||
assert index_quality["coverage"]["markets"]["cn.exchange"]["source"] == "tencent"
|
||||
assert index_quality["coverage"]["markets"]["cn.csi"]["status"] == "unavailable"
|
||||
assert index_quality["coverage"]["markets"]["cn.csi"]["source"] is None
|
||||
|
||||
|
||||
def test_fundamental_dataset_quality_is_market_aware() -> None:
|
||||
manager = _FetcherManager([
|
||||
_Fetcher("AkshareFetcher", 1, available=True),
|
||||
_Fetcher("TushareFetcher", 2, available=True),
|
||||
_Fetcher("YfinanceFetcher", 4, available=False),
|
||||
])
|
||||
service = DataCapabilityService(
|
||||
config=_config(tushare_token="token"),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
fundamental_quality = _dataset(overview, "financial.snapshot")
|
||||
|
||||
assert fundamental_quality["status"] == "partial"
|
||||
assert fundamental_quality["coverage"]["markets"]["cn"]["status"] == "ok"
|
||||
assert fundamental_quality["coverage"]["markets"]["cn"]["source"] == "akshare"
|
||||
assert fundamental_quality["coverage"]["markets"]["hk"]["status"] == "unavailable"
|
||||
assert fundamental_quality["coverage"]["markets"]["hk"]["source"] is None
|
||||
assert fundamental_quality["coverage"]["markets"]["us"]["status"] == "unavailable"
|
||||
assert fundamental_quality["coverage"]["markets"]["us"]["source"] is None
|
||||
for market in ("jp", "kr", "tw"):
|
||||
assert fundamental_quality["coverage"]["markets"][market]["status"] == "unavailable"
|
||||
assert fundamental_quality["coverage"]["markets"][market]["source"] is None
|
||||
assert "hk:source_status:yfinance:unavailable" in fundamental_quality["warnings"]
|
||||
assert "us:source_status:yfinance:unavailable" in fundamental_quality["warnings"]
|
||||
|
||||
|
||||
def test_screening_snapshot_priority_preserves_explicit_env_override() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(screening_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {"SNAPSHOT_SOURCE_PRIORITY": "tushare,em_datacenter"}, clear=False):
|
||||
with patch(
|
||||
"src.services.screening_service._resolve_screening_snapshot_source_priority",
|
||||
side_effect=AssertionError("resolver should not run when override is set"),
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
priorities = {item["scenario"]: item for item in overview["priorities"]}
|
||||
screening = priorities["screening.snapshot"]
|
||||
|
||||
assert screening["providers"] == ["tushare", "em_datacenter"]
|
||||
|
||||
|
||||
def test_screening_dataset_reports_engine_unavailable() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(screening_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {"SNAPSHOT_SOURCE_PRIORITY": "tushare,em_datacenter"}, clear=False):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_status_snapshot",
|
||||
return_value=({}, False, {"error": "screening_unavailable"}),
|
||||
):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_source_health_snapshot",
|
||||
return_value={},
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
screening_quality = _dataset(overview, "strategy.screening")
|
||||
|
||||
assert screening_quality["status"] == "unavailable"
|
||||
assert screening_quality["source"] is None
|
||||
assert screening_quality["warnings"] == ["screening_engine_unavailable"]
|
||||
|
||||
|
||||
def test_screening_dataset_reports_cooldown_sources_as_unavailable() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(screening_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {"SNAPSHOT_SOURCE_PRIORITY": "tushare,em_datacenter"}, clear=False):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_status_snapshot",
|
||||
return_value=({"available": True}, True, None),
|
||||
):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_source_health_snapshot",
|
||||
return_value={
|
||||
"snapshot": {
|
||||
"tushare": {"disabled": True},
|
||||
"em_datacenter": {"disabled": True},
|
||||
}
|
||||
},
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
screening_quality = _dataset(overview, "strategy.screening")
|
||||
|
||||
assert screening_quality["status"] == "unavailable"
|
||||
assert screening_quality["source"] is None
|
||||
assert screening_quality["fallback_from"] == ["tushare", "em_datacenter"]
|
||||
assert screening_quality["warnings"] == [
|
||||
"source_status:tushare:cooldown",
|
||||
"source_status:em_datacenter:cooldown",
|
||||
]
|
||||
|
||||
|
||||
def test_screening_dataset_keeps_known_sources_unknown_before_runtime_probe() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(screening_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {"SNAPSHOT_SOURCE_PRIORITY": "tushare,em_datacenter"}, clear=False):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_status_snapshot",
|
||||
return_value=({"available": True}, True, None),
|
||||
):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_source_health_snapshot",
|
||||
return_value={
|
||||
"snapshot": {
|
||||
"tushare": {"successes": 0, "failures": 0, "disabled": False},
|
||||
"em_datacenter": {"successes": 0, "failures": 0, "disabled": False},
|
||||
}
|
||||
},
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
screening_quality = _dataset(overview, "strategy.screening")
|
||||
|
||||
assert screening_quality["status"] == "unknown"
|
||||
assert screening_quality["source"] is None
|
||||
assert screening_quality["fallback_from"] == []
|
||||
assert screening_quality["warnings"] == [
|
||||
"source_status:tushare:unknown",
|
||||
]
|
||||
|
||||
|
||||
def test_screening_dataset_does_not_skip_unprobed_preferred_source() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(screening_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {"SNAPSHOT_SOURCE_PRIORITY": "tushare,em_datacenter"}, clear=False):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_status_snapshot",
|
||||
return_value=({"available": True}, True, None),
|
||||
):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_source_health_snapshot",
|
||||
return_value={
|
||||
"snapshot": {
|
||||
"tushare": {"successes": 0, "failures": 0, "disabled": False},
|
||||
"em_datacenter": {"successes": 1},
|
||||
}
|
||||
},
|
||||
):
|
||||
screening_quality = _dataset(service.get_overview(), "strategy.screening")
|
||||
|
||||
assert screening_quality["status"] == "unknown"
|
||||
assert screening_quality["source"] is None
|
||||
assert screening_quality["fallback_from"] == []
|
||||
assert screening_quality["warnings"] == ["source_status:tushare:unknown"]
|
||||
|
||||
|
||||
def test_screening_dataset_does_not_select_unknown_snapshot_source() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(screening_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
with patch.dict("os.environ", {"SNAPSHOT_SOURCE_PRIORITY": "mystery_source,em_datacenter"}, clear=False):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_status_snapshot",
|
||||
return_value=({"available": True}, True, None),
|
||||
):
|
||||
with patch(
|
||||
"src.services.screening_service._get_screening_source_health_snapshot",
|
||||
return_value={"snapshot": {"em_datacenter": {"successes": 1}}},
|
||||
):
|
||||
overview = service.get_overview()
|
||||
|
||||
screening_quality = _dataset(overview, "strategy.screening")
|
||||
|
||||
assert screening_quality["status"] == "degraded"
|
||||
assert screening_quality["source"] == "em_datacenter"
|
||||
assert screening_quality["fallback_from"] == ["mystery_source"]
|
||||
assert screening_quality["warnings"] == [
|
||||
"unknown_source:mystery_source",
|
||||
"source_status:mystery_source:unsupported",
|
||||
]
|
||||
|
||||
|
||||
def test_daily_capability_contract_includes_market_specific_daily_fetchers() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(
|
||||
futu_opend_host="127.0.0.1",
|
||||
finnhub_api_key="key",
|
||||
alphavantage_api_key="key",
|
||||
),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
|
||||
assert "kline.daily" in _provider(overview, "futu")["datasets"]
|
||||
assert "kline.daily" in _provider(overview, "finnhub")["datasets"]
|
||||
assert "kline.daily" in _provider(overview, "alphavantage")["datasets"]
|
||||
|
||||
|
||||
def test_disabled_runtime_features_surface_dataset_quality_warnings() -> None:
|
||||
manager = _FetcherManager([_Fetcher("AkshareFetcher", 1), _Fetcher("YfinanceFetcher", 4)])
|
||||
service = DataCapabilityService(
|
||||
config=_config(
|
||||
enable_realtime_quote=False,
|
||||
enable_fundamental_pipeline=False,
|
||||
screening_enabled=False,
|
||||
),
|
||||
fetcher_manager=manager,
|
||||
)
|
||||
|
||||
overview = service.get_overview()
|
||||
|
||||
assert _dataset(overview, "quote.realtime")["status"] == "unavailable"
|
||||
assert _dataset(overview, "quote.realtime")["warnings"] == ["realtime_quote_disabled"]
|
||||
assert _dataset(overview, "financial.snapshot")["status"] == "unavailable"
|
||||
assert _dataset(overview, "financial.snapshot")["warnings"] == ["fundamental_pipeline_disabled"]
|
||||
assert _dataset(overview, "strategy.screening")["status"] == "unconfigured"
|
||||
assert _dataset(overview, "strategy.screening")["warnings"] == ["screening_disabled"]
|
||||
assert _dataset(overview, "alert.monitor")["status"] == "unavailable"
|
||||
assert _dataset(overview, "alert.monitor")["source"] is None
|
||||
assert _dataset(overview, "alert.monitor")["warnings"] == ["agent_event_monitor_disabled"]
|
||||
assert _dataset(overview, "news.events")["source"] is None
|
||||
|
||||
|
||||
def test_enabled_alert_monitor_requires_a_registered_live_scheduler_task() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(agent_event_monitor_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
runtime_scheduler=_RuntimeScheduler(active=False),
|
||||
)
|
||||
|
||||
quality = _dataset(service.get_overview(), "alert.monitor")
|
||||
|
||||
assert quality["status"] == "unavailable"
|
||||
assert quality["source"] is None
|
||||
assert quality["warnings"] == ["agent_event_monitor_not_running"]
|
||||
|
||||
|
||||
def test_running_alert_monitor_surfaces_local_dataset_as_available() -> None:
|
||||
service = DataCapabilityService(
|
||||
config=_config(agent_event_monitor_enabled=True),
|
||||
fetcher_manager=_FetcherManager([]),
|
||||
runtime_scheduler=_RuntimeScheduler(active=True),
|
||||
)
|
||||
|
||||
quality = _dataset(service.get_overview(), "alert.monitor")
|
||||
|
||||
assert quality["status"] == "ok"
|
||||
assert quality["source"] == "alerts"
|
||||
assert quality["warnings"] == []
|
||||
|
||||
|
||||
def test_data_capability_api_paths_return_valid_contract() -> None:
|
||||
overview_payload = {
|
||||
"as_of": "2026-08-26T15:05:00+08:00",
|
||||
"providers": [
|
||||
{
|
||||
"name": "efinance",
|
||||
"label": "Efinance",
|
||||
"enabled": True,
|
||||
"configured": True,
|
||||
"status": "ok",
|
||||
"priority": 0,
|
||||
"markets": ["cn"],
|
||||
"datasets": ["quote.realtime"],
|
||||
"dataset_markets": {"quote.realtime": ["cn"]},
|
||||
"warnings": [],
|
||||
"last_error": None,
|
||||
"cooldown": None,
|
||||
}
|
||||
],
|
||||
"datasets": [
|
||||
{
|
||||
"dataset": "quote.realtime",
|
||||
"status": "ok",
|
||||
"source": "efinance",
|
||||
"stale": False,
|
||||
"last_success": None,
|
||||
"last_error": None,
|
||||
"fallback_from": [],
|
||||
"coverage": None,
|
||||
"warnings": [],
|
||||
}
|
||||
],
|
||||
"priorities": [
|
||||
{
|
||||
"scenario": "cn.realtime",
|
||||
"providers": ["efinance"],
|
||||
"source": "test",
|
||||
"warnings": [],
|
||||
}
|
||||
],
|
||||
"warnings": [],
|
||||
}
|
||||
|
||||
class _Service:
|
||||
def __init__(self, *, config, runtime_scheduler=None) -> None:
|
||||
self.config = config
|
||||
self.runtime_scheduler = runtime_scheduler
|
||||
|
||||
def get_overview(self):
|
||||
return overview_payload
|
||||
|
||||
with tempfile.TemporaryDirectory() as temp_dir:
|
||||
client = TestClient(create_app(static_dir=Path(temp_dir)))
|
||||
with patch("api.v1.endpoints.data.DataCapabilityService", _Service):
|
||||
overview_response = client.get("/api/v1/data/overview")
|
||||
capabilities_response = client.get("/api/v1/data/capabilities")
|
||||
|
||||
assert overview_response.status_code == 200
|
||||
assert capabilities_response.status_code == 200
|
||||
assert overview_response.json() == overview_payload
|
||||
assert capabilities_response.json() == overview_payload
|
||||
@@ -642,6 +642,20 @@ class RuntimeSchedulerServiceTestCase(unittest.TestCase):
|
||||
|
||||
self.assertEqual(calls, ["run"])
|
||||
|
||||
def test_background_task_active_requires_live_registered_scheduler(self) -> None:
|
||||
service = RuntimeSchedulerService()
|
||||
service._enabled = True
|
||||
service._scheduler = SimpleNamespace(
|
||||
_background_tasks=[{"name": "agent_event_monitor"}],
|
||||
)
|
||||
service._thread = SimpleNamespace(is_alive=lambda: True)
|
||||
|
||||
self.assertTrue(service.is_background_task_active("agent_event_monitor"))
|
||||
self.assertFalse(service.is_background_task_active("missing"))
|
||||
|
||||
service._thread = SimpleNamespace(is_alive=lambda: False)
|
||||
self.assertFalse(service.is_background_task_active("agent_event_monitor"))
|
||||
|
||||
def test_start_registers_event_monitor_background_task(self) -> None:
|
||||
class _FakeScheduler:
|
||||
def __init__(self, **kwargs):
|
||||
|
||||
@@ -6,6 +6,9 @@ must route to ``.HK`` rather than fall through to the ``.SZ`` default,
|
||||
otherwise Yahoo Finance returns 404 and the daily-data chain breaks.
|
||||
"""
|
||||
|
||||
from types import SimpleNamespace
|
||||
from unittest.mock import MagicMock, patch
|
||||
|
||||
from data_provider.yfinance_fetcher import YfinanceFetcher
|
||||
|
||||
|
||||
@@ -22,6 +25,28 @@ class TestHKPrefixStillRoutesToHK:
|
||||
# 1-3 digit HK prefix (e.g. hk0001 -> 0001.HK)
|
||||
assert YfinanceFetcher()._convert_stock_code("hk0001") == "0001.HK"
|
||||
|
||||
@patch("yfinance.Ticker")
|
||||
def test_hk_code_executes_realtime_route(self, ticker_factory: MagicMock) -> None:
|
||||
ticker = ticker_factory.return_value
|
||||
ticker.fast_info = SimpleNamespace(
|
||||
lastPrice=401.2,
|
||||
previousClose=398.0,
|
||||
open=399.0,
|
||||
dayHigh=403.0,
|
||||
dayLow=397.0,
|
||||
lastVolume=12_345,
|
||||
marketCap=3_800_000_000_000,
|
||||
)
|
||||
ticker.info = {"shortName": "Tencent", "currency": "HKD"}
|
||||
|
||||
quote = YfinanceFetcher().get_realtime_quote("HK00700")
|
||||
|
||||
ticker_factory.assert_called_once_with("0700.HK")
|
||||
assert quote is not None
|
||||
assert quote.code == "0700.HK"
|
||||
assert quote.market == "hk"
|
||||
assert quote.currency == "HKD"
|
||||
|
||||
|
||||
class TestBareHKCodeRoutesToHK:
|
||||
"""Regression for issue #2091: bare 4-5 digit numeric codes -> ``.HK``.
|
||||
|
||||
Reference in New Issue
Block a user