From 7aaeee36553417c9c3e23ed76101f55f1da9335c Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 23 Jul 2026 15:34:44 +0800 Subject: [PATCH 1/4] =?UTF-8?q?fix(data):=20=E5=88=87=E6=8D=A2=20A=20?= =?UTF-8?q?=E8=82=A1=E8=A1=8C=E6=83=85=E5=B9=B6=E7=BB=9F=E4=B8=80=E5=B8=82?= =?UTF-8?q?=E5=9C=BA=E8=BA=AB=E4=BB=BD?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- services/data/scripts/probe_baostock_batch.py | 112 ++++++++ .../data/src/inalpha_data/api/backfill.py | 48 ++-- services/data/src/inalpha_data/api/bars.py | 7 +- .../data/src/inalpha_data/api/constituents.py | 11 +- .../data/src/inalpha_data/api/fundamentals.py | 8 +- services/data/src/inalpha_data/api/news.py | 8 +- services/data/src/inalpha_data/api/ticker.py | 16 +- services/data/src/inalpha_data/config.py | 15 +- .../src/inalpha_data/connectors/__init__.py | 5 +- .../src/inalpha_data/connectors/baostock.py | 248 +++++++++++++----- services/data/src/inalpha_data/venues.py | 45 +++- services/data/tests/test_backfill_router.py | 80 +++++- services/data/tests/test_baostock_batch.py | 84 ------ services/data/tests/test_connectors.py | 220 +++++++++++++++- services/data/tests/test_fundamentals.py | 24 +- services/data/tests/test_news.py | 35 ++- services/data/tests/test_symbol_search.py | 29 +- services/data/tests/test_ticker_router.py | 95 +++++-- services/data/tests/test_venues.py | 50 ++++ 19 files changed, 869 insertions(+), 271 deletions(-) create mode 100644 services/data/scripts/probe_baostock_batch.py delete mode 100644 services/data/tests/test_baostock_batch.py create mode 100644 services/data/tests/test_venues.py diff --git a/services/data/scripts/probe_baostock_batch.py b/services/data/scripts/probe_baostock_batch.py new file mode 100644 index 00000000..898f55f3 --- /dev/null +++ b/services/data/scripts/probe_baostock_batch.py @@ -0,0 +1,112 @@ +"""手工探测 baostock 是否支持批量查询。 + +本文件不是 pytest 测试。需要访问真实 Baostock 网络时显式运行: + +``uv run python scripts/probe_baostock_batch.py`` +""" +from __future__ import annotations + +from datetime import datetime, timedelta +from typing import Any + +import baostock as bs + + +def main() -> int: + """运行单标的、多标的和分钟线探测,返回进程退出码。""" + login = bs.login() + if login.error_code != "0": + print(f"登录失败: {login.error_msg}") + return 1 + + try: + print("登录成功") + _probe_single_symbol() + count, stocks = _probe_multiple_symbols() + minute_count, minute_stocks = _probe_minute_symbols() + finally: + bs.logout() + print("\n登出成功") + + print("\n=== 结论 ===") + if count > 0 and len(stocks) > 1 and minute_count > 0 and len(minute_stocks) > 1: + print("✅ baostock 支持批量查询(多股票逗号分隔)") + print(" 可以在 1 次请求中获取多只股票数据,节省配额") + return 0 + + print("❌ baostock 不支持批量查询,需要单独查询每只股票") + return 1 + + +def _date_range(days: int) -> tuple[str, str]: + """返回 Baostock 所需的起止日期字符串。""" + now = datetime.now() + return (now - timedelta(days=days)).strftime("%Y-%m-%d"), now.strftime("%Y-%m-%d") + + +def _probe_single_symbol() -> None: + """探测单标的日线查询。""" + print("\n=== 测试单股票查询 ===") + start_date, end_date = _date_range(7) + result = bs.query_history_k_data_plus( + "sh.600519", + "date,code,open,high,low,close,volume", + start_date=start_date, + end_date=end_date, + frequency="d", + ) + print(f"单股票查询结果: error_code={result.error_code}") + count = 0 + while result.error_code == "0" and result.next(): + count += 1 + if count <= 3: + print(f" {result.get_row_data()}") + print(f" 共 {count} 条数据") + + +def _probe_multiple_symbols() -> tuple[int, set[str]]: + """探测多标的日线查询。""" + print("\n=== 测试多股票查询(逗号分隔)===") + start_date, end_date = _date_range(7) + result = bs.query_history_k_data_plus( + "sh.600519,sh.600036,sh.601318", + "date,code,open,high,low,close,volume", + start_date=start_date, + end_date=end_date, + frequency="d", + ) + print(f"多股票查询结果: error_code={result.error_code}") + return _consume_rows(result, code_index=1) + + +def _probe_minute_symbols() -> tuple[int, set[str]]: + """探测多标的五分钟线查询。""" + print("\n=== 测试分钟 K 线批量查询 ===") + start_date, end_date = _date_range(1) + result = bs.query_history_k_data_plus( + "sh.600519,sh.600036", + "date,time,code,open,high,low,close,volume", + start_date=start_date, + end_date=end_date, + frequency="5", + ) + print(f"分钟 K 线查询结果: error_code={result.error_code}") + return _consume_rows(result, code_index=2) + + +def _consume_rows(result: Any, *, code_index: int) -> tuple[int, set[str]]: + """消费 Baostock ResultData,打印样本并返回行数和标的集合。""" + count = 0 + stocks: set[str] = set() + while result.error_code == "0" and result.next(): + count += 1 + row = result.get_row_data() + stocks.add(row[code_index]) + if count <= 5: + print(f" {row}") + print(f" 共 {count} 条数据,涉及 {len(stocks)} 只股票: {stocks}") + return count, stocks + + +if __name__ == "__main__": + raise SystemExit(main()) diff --git a/services/data/src/inalpha_data/api/backfill.py b/services/data/src/inalpha_data/api/backfill.py index 5ed26ac2..d30a6a0e 100644 --- a/services/data/src/inalpha_data/api/backfill.py +++ b/services/data/src/inalpha_data/api/backfill.py @@ -12,7 +12,7 @@ from inalpha_shared import get_logger from inalpha_shared.auth import User, get_current_user from inalpha_shared.db import DBConn -from inalpha_shared.errors import ValidationError +from inalpha_shared.errors import InalphaError, ValidationError from ..connectors import Connector, get_connector_for_venue, list_registered_venues from ..connectors.alpaca import TIMEFRAME_SECONDS as ALPACA_TIMEFRAME_SECONDS @@ -21,7 +21,7 @@ from ..connectors.yfinance_conn import TIMEFRAME_SECONDS as YFINANCE_TIMEFRAME_SECONDS from ..schemas import BackfillRequest, BackfillResponse from ..storage.bars import insert_bars, latest_bar_ts -from ..venues import canonicalize_venue, is_legacy_a_share_venue +from ..venues import canonicalize_market_identity, is_legacy_a_share_venue router = APIRouter(tags=["backfill"]) _logger = get_logger(__name__) @@ -33,8 +33,8 @@ # D-8b' review 高风险 #6:长跨度同步 backfill 会卡死请求线程 _MAX_BARS_PER_REQUEST = 50_000 -# 分钟级查询限制(baostock 配额优化) -# baostock 分钟 K 每条调用一次 API,长跨度会快速消耗 5 万日配额 +# 分钟级查询限制(腾讯行情请求量与同步延迟保护) +# 腾讯分钟 K 单请求有返回条数上限;长跨度既无法完整返回,也会占用串行抓取通道。 _MINUTE_LOOKBACK_LIMITS = { "5m": 7, # 7 天 = 336 条 "15m": 14, # 14 天 = 672 条 @@ -43,6 +43,13 @@ } +class BarsUpstreamUnavailableError(InalphaError): + """外部 K 线源不可用,避免把网络故障伪装成成功的空结果。""" + + code = "BARS_UPSTREAM_UNAVAILABLE" + status_code = 502 + + # venue → 该 venue 支持的 timeframe → 秒数 # baostock(baostock venue)支持日级 + 分钟级(5/15/30/60 分钟) _VENUE_TIMEFRAME_SECONDS: dict[str, dict[str, int]] = { @@ -80,14 +87,13 @@ async def backfill_bars( # ─── venue 路由 ────────────────────────────────────────────── # **向后兼容**:venue="akshare" 且 symbol 带 sh./sz. 前缀 → 自动路由到 baostock - effective_venue = canonicalize_venue(req.venue, req.symbol) + effective_venue, effective_symbol = canonicalize_market_identity(req.venue, req.symbol) if is_legacy_a_share_venue(req.venue, req.symbol): _logger.warning( "venue_akshare_deprecated", symbol=req.symbol, reason="venue 'akshare' is deprecated for A-share; use 'baostock' instead", ) - effective_venue = "baostock" try: connector: Connector = get_connector_for_venue(effective_venue) @@ -107,8 +113,8 @@ async def backfill_bars( }, ) - # ─── 分钟级强制限制(baostock 配额优化)────────────────────── - # baostock 分钟 K 每条调用一次 API,长跨度会快速消耗 5 万日配额 + # ─── 分钟级强制限制(腾讯行情请求量与同步延迟保护)───────────────── + # canonical baostock venue 的分钟线由腾讯 HTTPS 提供,限制长跨度无效拉取。 # 仅对 baostock venue 生效(binance/alpaca/yfinance 不受此限制) effective_from_ts = req.from_ts if effective_venue == "baostock" and req.timeframe in _MINUTE_LOOKBACK_LIMITS: @@ -121,12 +127,12 @@ async def backfill_bars( _logger.warning( "backfill_minute_lookback_capped", venue=req.venue, - symbol=req.symbol, + symbol=effective_symbol, timeframe=req.timeframe, original_from_ts=req.from_ts.isoformat(), capped_from_ts=effective_from_ts.isoformat(), max_lookback_days=max_lookback_days, - reason="baostock quota optimization", + reason="minute kline request bound", ) span_seconds = (req.to_ts - effective_from_ts).total_seconds() @@ -153,7 +159,7 @@ async def backfill_bars( # 但不早于请求的 from_ts;空缓存则从 from_ts 全量。 # 注:仅按 max(ts) 续拉,中间空洞(非连续缓存,罕见)不会回补;需要时显式重拉窗口。 cached_latest = await latest_bar_ts( - db, effective_venue, req.symbol, req.timeframe, upto=req.to_ts + db, effective_venue, effective_symbol, req.timeframe, upto=req.to_ts ) if cached_latest is not None and cached_latest > effective_from_ts: cursor = cached_latest @@ -173,7 +179,7 @@ async def backfill_bars( while cursor < req.to_ts: try: bars = await connector.fetch_bars( - symbol=req.symbol, + symbol=effective_symbol, timeframe=req.timeframe, since=cursor, limit=_BATCH_LIMIT, @@ -182,18 +188,26 @@ async def backfill_bars( _logger.warning( "backfill_connector_failed", venue=req.venue, - symbol=req.symbol, + symbol=effective_symbol, error=str(exc), cursor=cursor.isoformat(), ) - break + raise BarsUpstreamUnavailableError( + f"bars source unavailable for {effective_symbol}@{effective_venue}", + details={ + "venue": req.venue, + "symbol": req.symbol, + "timeframe": req.timeframe, + "reason": str(exc), + }, + ) from exc if not bars: if fetched_total == 0: # 首批就空——上游根本没返数据,比增量结束严重 _logger.warning( "backfill_no_more_bars_first_batch_empty", venue=req.venue, - symbol=req.symbol, + symbol=effective_symbol, timeframe=req.timeframe, cursor=cursor.isoformat(), ) @@ -201,7 +215,7 @@ async def backfill_bars( _logger.info( "backfill_no_more_bars", venue=req.venue, - symbol=req.symbol, + symbol=effective_symbol, cursor=cursor.isoformat(), ) break @@ -211,7 +225,7 @@ async def backfill_bars( if not bars: break - n = await insert_bars(db, effective_venue, req.symbol, req.timeframe, bars) + n = await insert_bars(db, effective_venue, effective_symbol, req.timeframe, bars) fetched_total += len(bars) inserted_total += n diff --git a/services/data/src/inalpha_data/api/bars.py b/services/data/src/inalpha_data/api/bars.py index c213b341..4dc99deb 100644 --- a/services/data/src/inalpha_data/api/bars.py +++ b/services/data/src/inalpha_data/api/bars.py @@ -12,7 +12,7 @@ from ..schemas import BarResponse, BarsQuery from ..storage.bars import query_bars -from ..venues import canonicalize_venue, is_legacy_a_share_venue +from ..venues import canonicalize_market_identity, is_legacy_a_share_venue router = APIRouter(tags=["bars"]) _logger = get_logger(__name__) @@ -38,19 +38,18 @@ async def list_bars( ) # 向后兼容:venue="akshare" + sh./sz. 前缀 → 自动用 baostock 查 DB - effective_venue = canonicalize_venue(query.venue, query.symbol) + effective_venue, effective_symbol = canonicalize_market_identity(query.venue, query.symbol) if is_legacy_a_share_venue(query.venue, query.symbol): _logger.warning( "venue_akshare_deprecated", symbol=query.symbol, reason="venue 'akshare' is deprecated for A-share; use 'baostock' instead", ) - effective_venue = "baostock" rows = await query_bars( db, venue=effective_venue, - symbol=query.symbol, + symbol=effective_symbol, timeframe=query.timeframe, from_ts=query.from_ts, to_ts=query.to_ts, diff --git a/services/data/src/inalpha_data/api/constituents.py b/services/data/src/inalpha_data/api/constituents.py index b965fe78..87a5b92a 100644 --- a/services/data/src/inalpha_data/api/constituents.py +++ b/services/data/src/inalpha_data/api/constituents.py @@ -7,6 +7,7 @@ 横截面选股/轮动回测的存活者偏差前提:每期取 as_of 那刻的真实成分,而非"今天还在的票"回看。 """ + from __future__ import annotations from datetime import UTC, datetime @@ -37,14 +38,12 @@ async def snapshot_constituents( ) -> SnapshotConstituentsResponse: """拉 ``index_code`` 当前成分(akshare)落库,``as_of_date=今天``。 - 每日(或按需)调用一次即向前累积一份 PIT 快照。akshare 只回当前成分,故本接口是 + 每日(或按需)调用一次即向前累积一份 PIT 快照。免费源只回当前成分,故本接口是 PIT 史的**唯一来源**;源站失败 → 502,不静默写空(§3.1)。与每日调度器共用 :func:`record_snapshot`,手动触发与自动累积行为一致。 """ snap_date, n = await record_snapshot(db, index_code=req.index_code) - return SnapshotConstituentsResponse( - index_code=req.index_code, as_of_date=snap_date, count=n - ) + return SnapshotConstituentsResponse(index_code=req.index_code, as_of_date=snap_date, count=n) @router.get("/constituents", response_model=ConstituentsResponse) @@ -72,9 +71,7 @@ async def get_constituents_pit( else: as_of_date = datetime.now(UTC).date() - snap_date, members = await store.get_constituents( - db, index_code=index_code, as_of=as_of_date - ) + snap_date, members = await store.get_constituents(db, index_code=index_code, as_of=as_of_date) return ConstituentsResponse( index_code=index_code, as_of=as_of_date.isoformat(), diff --git a/services/data/src/inalpha_data/api/fundamentals.py b/services/data/src/inalpha_data/api/fundamentals.py index 98fef75b..299bc398 100644 --- a/services/data/src/inalpha_data/api/fundamentals.py +++ b/services/data/src/inalpha_data/api/fundamentals.py @@ -14,7 +14,7 @@ from ..connectors import yfinance_conn from ..connectors._base import get_connector_for_venue from ..schemas import FinancialsResponse -from ..venues import canonicalize_venue +from ..venues import canonicalize_market_identity router = APIRouter(tags=["fundamentals"]) @@ -36,11 +36,11 @@ async def get_fundamentals( venue 支持 baostock / yfinance;其它 venue 返 422。 """ - effective_venue = canonicalize_venue(venue, symbol) + effective_venue, effective_symbol = canonicalize_market_identity(venue, symbol) if effective_venue == "yfinance": try: conn = yfinance_conn.get_connector() - data = await conn.fetch_financials(symbol, as_of=as_of) + data = await conn.fetch_financials(effective_symbol, as_of=as_of) except Exception as exc: return FinancialsResponse( venue=venue, @@ -58,7 +58,7 @@ async def get_fundamentals( code="FUNDAMENTALS_NOT_SUPPORTED", details={"venue": venue}, ) - data = await conn.fetch_financials(symbol, as_of=as_of) # type: ignore[union-attr] + data = await conn.fetch_financials(effective_symbol, as_of=as_of) # type: ignore[union-attr] return FinancialsResponse(**data) raise ValidationError( diff --git a/services/data/src/inalpha_data/api/news.py b/services/data/src/inalpha_data/api/news.py index 52887cdc..876dc5b0 100644 --- a/services/data/src/inalpha_data/api/news.py +++ b/services/data/src/inalpha_data/api/news.py @@ -14,7 +14,7 @@ from ..connectors import yfinance_conn from ..connectors._base import get_connector_for_venue from ..schemas import NewsItem, NewsQuery, NewsResponse -from ..venues import canonicalize_venue +from ..venues import canonicalize_market_identity router = APIRouter(tags=["news"]) @@ -28,11 +28,11 @@ async def get_news( venue 支持 yfinance / baostock(其它返 422);不支持 ticker 返空 list 而非错。 """ - effective_venue = canonicalize_venue(query.venue, query.symbol) + effective_venue, effective_symbol = canonicalize_market_identity(query.venue, query.symbol) if effective_venue == "yfinance": try: conn = yfinance_conn.get_connector() - raw = await conn.fetch_news(query.symbol, limit=query.limit) + raw = await conn.fetch_news(effective_symbol, limit=query.limit) except Exception: raw = [] elif effective_venue == "baostock": @@ -43,7 +43,7 @@ async def get_news( code="NEWS_FETCH_NOT_SUPPORTED", details={"venue": query.venue}, ) - raw = await conn.fetch_news(query.symbol, limit=query.limit) # type: ignore[union-attr] + raw = await conn.fetch_news(effective_symbol, limit=query.limit) # type: ignore[union-attr] else: raise ValidationError( f"news venue {query.venue!r} not supported", diff --git a/services/data/src/inalpha_data/api/ticker.py b/services/data/src/inalpha_data/api/ticker.py index cc7bdc7f..09040e28 100644 --- a/services/data/src/inalpha_data/api/ticker.py +++ b/services/data/src/inalpha_data/api/ticker.py @@ -15,13 +15,14 @@ D-9 ``fresh=true`` 路由从硬编码 binance 改成**按 connector capability 鸭子分发**: - venue 在注册表 + connector 实现了 ``TickerCapable`` Protocol(含 ``fetch_ticker``) - → 调该 venue 的 fetch_ticker(binance / yfinance / alpaca 当前实现) -- venue 在注册表但 connector 无 fetch_ticker(akshare / fred)→ 返 + → 调该 venue 的 fetch_ticker(binance / yfinance / alpaca / baostock 当前实现) +- venue 在注册表但 connector 无 fetch_ticker(fred)→ 返 ``FRESH_NOT_SUPPORTED_FOR_VENUE``,引导 caller 切 fresh=false 走 DB cache - venue 未注册 → 422 unsupported venue(与 ``/backfill/bars`` 一致的错误形态) 返回 ``is_stale=true`` 时 caller 决定是否信任(paper 当前都信,后续可加阈值拒绝)。 """ + from __future__ import annotations from datetime import UTC, datetime, timedelta @@ -39,7 +40,7 @@ ) from ..schemas import TickerQuery, TickerResponse from ..storage.bars import query_bars -from ..venues import canonicalize_venue +from ..venues import canonicalize_market_identity router = APIRouter(tags=["ticker"]) @@ -62,6 +63,7 @@ class FreshNotSupportedError(InalphaError): class TickerUnavailableError(InalphaError): """外部数据源实时报价不可用(限流 / 超时 / 网络分区),非代码 bug。""" + code = "TICKER_UNAVAILABLE" status_code = 502 @@ -76,12 +78,12 @@ async def get_ticker( - ``fresh=false``(默认):DB 优先 1m → 1h;都没有抛 ``NO_PRICE_AVAILABLE``。 - ``fresh=true``:调 venue 实时 ticker;venue connector 须实现 ``TickerCapable`` - (当前 binance / yfinance / alpaca)。akshare / fred 不实现则返 + (当前 binance / yfinance / alpaca / baostock)。fred 不实现则返 ``FRESH_NOT_SUPPORTED_FOR_VENUE`` 引导走 fresh=false。 网络抖动失败时**不**自动 fallback 到 DB(让 caller 看到原因),由 caller 决定重试。 """ now = datetime.now(UTC) - effective_venue = canonicalize_venue(query.venue, query.symbol) + effective_venue, effective_symbol = canonicalize_market_identity(query.venue, query.symbol) if query.fresh: try: @@ -100,7 +102,7 @@ async def get_ticker( }, ) try: - ts, price = await connector.fetch_ticker(query.symbol) + ts, price = await connector.fetch_ticker(effective_symbol) except RuntimeError as exc: raise TickerUnavailableError( str(exc), @@ -128,7 +130,7 @@ async def get_ticker( rows = await query_bars( db, venue=effective_venue, - symbol=query.symbol, + symbol=effective_symbol, timeframe=timeframe, from_ts=lookback_start, to_ts=now, diff --git a/services/data/src/inalpha_data/config.py b/services/data/src/inalpha_data/config.py index 9ed6191a..969666ac 100644 --- a/services/data/src/inalpha_data/config.py +++ b/services/data/src/inalpha_data/config.py @@ -2,6 +2,7 @@ 继承 ``inalpha_shared.Settings``,加 Binance 凭证字段(公开接口可以为空)。 """ + from __future__ import annotations from functools import lru_cache @@ -38,16 +39,12 @@ class DataSettings(BaseSettings): web_search_timeout_s: int = Field(default=8, alias="WEB_SEARCH_TIMEOUT_S") """ddgs 单引擎 HTTP 超时(秒)。原默认 15s,叠多引擎可到 30s+,收紧到 8s 砍长尾。""" - web_search_overall_timeout_s: int = Field( - default=20, alias="WEB_SEARCH_OVERALL_TIMEOUT_S" - ) + web_search_overall_timeout_s: int = Field(default=20, alias="WEB_SEARCH_OVERALL_TIMEOUT_S") """单次搜索整体超时(秒),避免 backend="auto" 顺序试 8 个引擎把调用方拖死。 12s 对 bing 中文查询偏紧(本地网络实测常恰好掐死);auto 失败会换引擎再兜一次, 最坏耗时 = 2 × 本值。""" - web_search_max_concurrency: int = Field( - default=4, alias="WEB_SEARCH_MAX_CONCURRENCY" - ) + web_search_max_concurrency: int = Field(default=4, alias="WEB_SEARCH_MAX_CONCURRENCY") """同时在飞的搜索数上限。analyst 常 ~10 个并行查询,限并发避免线程池 + GIL 把 async 事件循环饿死。""" web_search_cache_ttl_s: int = Field(default=600, alias="WEB_SEARCH_CACHE_TTL_S") @@ -78,11 +75,9 @@ class DataSettings(BaseSettings): """进程内缓存 TTL(秒)。快讯/板块榜分钟级更新,60s 挡住 analyst fan-out 同一轮重复打源站;响应带 fetched_at,fresh 语义不破。""" - constituent_snapshot_indices: str = Field( - default="", alias="CONSTITUENT_SNAPSHOT_INDICES" - ) + constituent_snapshot_indices: str = Field(default="", alias="CONSTITUENT_SNAPSHOT_INDICES") """每日成分快照追踪的指数代码,逗号分隔(如 ``000300,000905``)。空=禁用调度 - (ADR-0053 阶段 C 向前累积:akshare 只回当前成分,唯一 PIT 路径是从启用日起每日落库)。 + (ADR-0053 阶段 C 向前累积:免费源只回当前成分,唯一 PIT 路径是从启用日起每日落库)。 手动 ``POST /constituents/snapshot`` 不受本项影响。""" constituent_snapshot_interval_h: float = Field( diff --git a/services/data/src/inalpha_data/connectors/__init__.py b/services/data/src/inalpha_data/connectors/__init__.py index 07578aae..397c7ff5 100644 --- a/services/data/src/inalpha_data/connectors/__init__.py +++ b/services/data/src/inalpha_data/connectors/__init__.py @@ -1,14 +1,15 @@ """外部市场 / 经纪商接入。 注册表模式(``_base.Connector`` Protocol + venue → connector dict)让 ``api`` 层 -按 venue 找 connector,不必关心后端是 CCXT / alpaca-py / akshare。 +按 venue 找 connector,不必关心后端是 CCXT / alpaca-py / 腾讯财经 / Baostock。 D-9 起: - ``binance`` :CCXT spot OHLCV(crypto) - ``alpaca`` :alpaca-py IEX free feed(美股 OHLCV) -- ``akshare`` :akshare 公开页(A股 sh./sz. + 港股 hk.) +- ``baostock`` :A 股逻辑 venue;腾讯 HTTPS 行情 + Baostock 基本面/日历/成分 """ + from ._base import ( Connector, TickerCapable, diff --git a/services/data/src/inalpha_data/connectors/baostock.py b/services/data/src/inalpha_data/connectors/baostock.py index 479d5615..accbd18a 100644 --- a/services/data/src/inalpha_data/connectors/baostock.py +++ b/services/data/src/inalpha_data/connectors/baostock.py @@ -1,13 +1,11 @@ -"""baostock connector —— A 股全栈数据源(证券宝,免费零 key)。 +"""baostock venue connector —— A 股腾讯行情 + 证券宝基本面(免费零 key)。 -2026-07 起,原 akshare venue 的 A 股能力全部迁移至此独立 venue。 +2026-07 起,原 akshare venue 的 A 股能力迁移至此独立 venue。 -baostock(证券宝)功能覆盖: +功能覆盖: -- **K 线**:日/周/月/分钟(5m/15m/30m/1h),OHLCV 齐全(含真实成交量) -- **财报**:利润表 + 负债表 + 成长指标 + 现金流 + 运营能力 + 杜邦分析 + 分红记录 -- **交易日历**:A 股交易/非交易日查询 -- **指数成分**:沪深300 / 上证50 / 中证500 当前成分股 +- **K 线 / 最新价**:腾讯财经 HTTPS,日/周/月/分钟(5m/15m/30m/1h) +- **财报 / 交易日历 / 核心指数成分**:baostock(证券宝) - **全部免费零 key**,无需注册 **symbol 格式约定**(venue=``"baostock"``): @@ -17,9 +15,9 @@ **timeframe 支持**: -- ``"1d"`` / ``"1wk"`` / ``"1mo"`` → baostock ``frequency="d"/"w"/"m"`` -- ``"5m"`` / ``"15m"`` / ``"30m"`` / ``"1h"`` → baostock ``frequency="5"/"15"/"30"/"60"`` -- 不支持 1 分钟(baostock 限制) +- ``"1d"`` / ``"1wk"`` / ``"1mo"`` → 腾讯 ``day/week/month`` +- ``"5m"`` / ``"15m"`` / ``"30m"`` / ``"1h"`` → 腾讯 ``m5/m15/m30/m60`` +- 不支持 1 分钟 """ from __future__ import annotations @@ -29,6 +27,7 @@ import time from datetime import UTC, datetime, timedelta from typing import Any +from zoneinfo import ZoneInfo from inalpha_shared import get_logger @@ -38,31 +37,39 @@ VENUE = "baostock" +#: 腾讯财经 HTTPS K 线。A 股行情使用此源,避免 Baostock 专用 TCP 端口在海外生产机被阻断。 +_TENCENT_KLINE_URL = "https://web.ifzq.gtimg.cn/appstock/app/fqkline/get" +_TENCENT_MINUTE_KLINE_URL = "https://ifzq.gtimg.cn/appstock/app/kline/mkline" +_TENCENT_QUOTE_URL = "https://qt.gtimg.cn/q=" +_TENCENT_HEADERS = { + "User-Agent": "Mozilla/5.0", + "Referer": "https://gu.qq.com/", +} +_TZ_SHANGHAI = ZoneInfo("Asia/Shanghai") + #: fundamentals 进程内缓存 TTL(秒)。A股一次 fundamentals 要打 1 次财报摘要 + 3 次 #: Baidu 估值(串行 ~4-5s),research 多 analyst fan-out 会重复问同一标的;60s 缓存挡掉 #: 重复打源站,兼顾防封与延迟。基本面日级更新,60s 内复用不损时效。 _FIN_CACHE_TTL_S = 60.0 -# akshare 的 ``period`` 字符串映射(日级 + 分钟级) -# baostock frequency 参数:日级 "d"/"w"/"m",分钟级 "5"/"15"/"30"/"60" +# connector 的统一 period 映射;行情实现再转换为腾讯 day/week/month 或分钟名称。 _PERIOD_MAP: dict[str, str] = { "1d": "daily", "1wk": "weekly", "1mo": "monthly", - # 分钟级(baostock 支持 5/15/30/60 分钟) + # 分钟级(腾讯支持 5/15/30/60 分钟) "5m": "5", "15m": "15", "30m": "30", "1h": "60", } -#: 串行锁——akshare 走公开页聚合(东财/同花顺/中证等),并发突发会触发反爬; -#: 进程级串行 + 最小间隔把突发摊成节流串行,避免 429/空返(yfinance 同类模式)。 +#: 串行锁——腾讯 / AkShare 公开接口并发突发会触发反爬;进程级串行 + 最小间隔 +#: 把突发摊成节流串行,避免 429/空返(yfinance 同类模式)。 _FETCH_LOCK = asyncio.Lock() -#: 最小拉取间隔(秒)。公开页比 Yahoo API 更脆弱,≥1s;A 股数据走东财/同花顺 -#: 直连配方,memory ``a_stock_data_recipes`` 已记录"防封串行≥1s"。 +#: 最小拉取间隔(秒)。公开页比 Yahoo API 更脆弱,≥1s;A 股按直连配方限速。 _MIN_FETCH_INTERVAL_S = 1.0 -#: 锁内单次拉取超时上限——TCP 挂起时快速放锁,不把整个 panel 拖死。 +#: 锁内单次拉取超时上限——网络挂起时快速放锁,不把整个 panel 拖死。 _FETCH_TIMEOUT_S = 30.0 _last_fetch_mono: float = 0.0 @@ -130,7 +137,7 @@ def _parse_symbol(symbol: str) -> tuple[str, str]: # connector 内部做格式适配。只匹配纯数字代码 + 已知市场后缀。 if "." in symbol and symbol.rsplit(".", 1)[1].lower() in {"sh", "sz"}: code, suffix = symbol.rsplit(".", 1) - if code.replace(".", "", 1).isdigit(): + if code.isdigit() and len(code) == 6: symbol = f"{suffix.lower()}.{code}" if "." not in symbol: raise ValueError( @@ -140,8 +147,8 @@ def _parse_symbol(symbol: str) -> tuple[str, str]: prefix = prefix.lower() if prefix not in _ALLOWED_PREFIXES: raise ValueError(f"baostock unknown prefix {prefix!r},allow: {sorted(_ALLOWED_PREFIXES)}") - if not code: - raise ValueError(f"baostock code is empty: {symbol!r}") + if not code.isdigit() or len(code) != 6: + raise ValueError(f"baostock code must be exactly 6 digits: {symbol!r}") return prefix, code @@ -187,13 +194,13 @@ async def fetch_bars( since: datetime, limit: int = 1000, ) -> list[tuple[datetime, float, float, float, float, float]]: - """从 akshare 拉 OHLCV。 + """从腾讯财经 HTTPS 拉 OHLCV。 Args: symbol: ``"sh.600519"`` / ``"sz.000001"`` timeframe: 支持 ``"1d"`` / ``"1wk"`` / ``"1mo"`` 及分钟级 ``"5m"`` / ``"15m"`` / ``"30m"`` / ``"1h"`` - since: UTC datetime;akshare 接 ``YYYYMMDD`` 字符串 - limit: 不直接生效(akshare 不接 limit,整段拉回;上层切片) + since: UTC datetime;转换为腾讯接口的起始日期 + limit: 单次最多请求的 K 线数量,上层仍会再次截断尾部 Returns: list of ``(ts, open, high, low, close, volume)``,UTC aware。 @@ -206,8 +213,7 @@ async def fetch_bars( prefix, code = _parse_symbol(symbol) period = _PERIOD_MAP[timeframe] start_str = since.strftime("%Y%m%d") - # end 给 today 让 baostock 一口气拉全 - end_str = datetime.now(UTC).strftime("%Y%m%d") + end_str = _tencent_window_end(since, period, limit).strftime("%Y%m%d") _logger.debug( "baostock_fetch_bars", @@ -225,6 +231,7 @@ async def fetch_bars( period=period, start_str=start_str, end_str=end_str, + limit=limit, ) except Exception as exc: _logger.warning( @@ -235,13 +242,17 @@ async def fetch_bars( end_str=end_str, error=str(exc), ) - return [] + raise RuntimeError(f"A-share bars unavailable for {symbol}: {exc}") from exc - # akshare 返的是 DataFrame;列名中文 / 英文都见过,做防御性归一化 + # 腾讯响应已归一为 dict;同时兼容历史中文字段,避免 fallback 数据形态漂移。 out: list[tuple[datetime, float, float, float, float, float]] = [] for r in rows: - # 日级:用 date 字段;分钟级:用 time 字段(格式 YYYYMMDDHHMMSS) - ts_raw = r.get("日期") or r.get("date") or r.get("Date") or r.get("time") + # 日级:用 date 字段;分钟级:优先用 time 字段(格式 YYYYMMDDHHMMSS) + ts_raw = ( + r.get("time") + if timeframe in {"5m", "15m", "30m", "1h"} + else r.get("日期") or r.get("date") or r.get("Date") + ) o = _to_float(r.get("开盘") or r.get("open")) h = _to_float(r.get("最高") or r.get("high")) low = _to_float(r.get("最低") or r.get("low")) @@ -254,20 +265,36 @@ async def fetch_bars( out.append((ts, o or 0.0, h or 0.0, low or 0.0, c, v or 0.0)) if not out and rows: - # 上游返了行但全部被列名解析跳过 → 列名漂移告警 + # 上游返了行但全部被列名解析跳过,属于响应契约漂移而不是正常空结果。 + sample_keys = list(rows[0].keys()) if rows else [] _logger.warning( "baostock_fetch_bars_all_rows_skipped", symbol=symbol, timeframe=timeframe, row_count=len(rows), - sample_keys=list(rows[0].keys()) if rows else [], + sample_keys=sample_keys, + ) + raise RuntimeError( + f"A-share bars payload invalid for {symbol}: " + f"all {len(rows)} rows were skipped (keys={sample_keys})" ) - # 按 limit 截断尾部(akshare 不接 limit,整段返) + # 防御性截断尾部;上游可能忽略或扩大 limit。 if limit and len(out) > limit: out = out[-limit:] return out + async def fetch_ticker(self, symbol: str) -> tuple[datetime, float]: + """从腾讯财经实时行情接口返回最新报价时间和价格。""" + prefix, code = _parse_symbol(symbol) + try: + return await asyncio.wait_for( + asyncio.to_thread(_fetch_tencent_ticker_sync, symbol=f"{prefix}.{code}"), + timeout=_FETCH_TIMEOUT_S, + ) + except Exception as exc: + raise RuntimeError(f"A-share ticker unavailable for {symbol}: {exc}") from exc + @staticmethod async def _throttled_fetch_sync( *, @@ -276,10 +303,11 @@ async def _throttled_fetch_sync( period: str, start_str: str, end_str: str, + limit: int, ) -> list[dict[str, Any]]: - """串行 + 最小间隔跑 akshare fetch,防并发突发触发反爬。 + """串行 + 最小间隔跑 A 股公开行情请求,防并发突发触发反爬。 - 进程级 ``_FETCH_LOCK`` 保证同一时刻只有一个 akshare 请求在飞;锁内再补足 + 进程级 ``_FETCH_LOCK`` 保证同一时刻只有一个请求在飞;锁内再补足 ``_MIN_FETCH_INTERVAL_S`` 的最小间隔。锁内每请求超时 ``_FETCH_TIMEOUT_S``, TCP 挂起时快速放锁让队列继续。 @@ -299,6 +327,7 @@ async def _throttled_fetch_sync( period=period, start_str=start_str, end_str=end_str, + limit=limit, ), timeout=_FETCH_TIMEOUT_S, ) @@ -1010,6 +1039,114 @@ def _fetch_news_sync(symbol: str) -> list[dict[str, Any]]: return [] +def _fetch_tencent_ticker_sync(*, symbol: str) -> tuple[datetime, float]: + """从腾讯财经 GBK 行情接口拉一个 A 股标的的实时价格。""" + import httpx + + compact_symbol = symbol.replace(".", "") + response = httpx.get( + f"{_TENCENT_QUOTE_URL}{compact_symbol}", + headers=_TENCENT_HEADERS, + timeout=_FETCH_TIMEOUT_S, + trust_env=False, + ) + response.raise_for_status() + text = response.content.decode("gbk") + if '="' not in text: + raise RuntimeError(f"tencent quote response invalid: {text[:120]}") + fields = text.split('="', 1)[1].split('"', 1)[0].split("~") + if len(fields) <= 30 or not fields[3] or not fields[30]: + raise RuntimeError(f"tencent quote fields missing for {symbol}") + quote_time = datetime.strptime(fields[30], "%Y%m%d%H%M%S").replace(tzinfo=_TZ_SHANGHAI) + return quote_time.astimezone(UTC), float(fields[3]) + + +def _tencent_window_end(since: datetime, period: str, limit: int) -> datetime: + """按腾讯单次条数上限计算历史请求窗口终点,供 backfill 分页推进。""" + days_per_bar = { + "daily": 2, + "weekly": 8, + "monthly": 32, + }.get(period) + if days_per_bar is None: + return datetime.now(UTC) + return min( + datetime.now(UTC), + since + timedelta(days=max(limit, 1) * days_per_bar), + ) + + +def _fetch_tencent_bars_sync( + *, + symbol: str, + period: str, + start_date: str, + end_date: str, + limit: int, +) -> list[dict[str, Any]]: + """从腾讯财经 HTTPS 接口拉 A 股 K 线。""" + import httpx + + compact_symbol = symbol.replace(".", "") + if period in {"daily", "weekly", "monthly"}: + period_name = {"daily": "day", "weekly": "week", "monthly": "month"}[period] + start_fmt = f"{start_date[:4]}-{start_date[4:6]}-{start_date[6:8]}" + end_fmt = f"{end_date[:4]}-{end_date[4:6]}-{end_date[6:8]}" + params = {"param": (f"{compact_symbol},{period_name},{start_fmt},{end_fmt},{limit},")} + response = httpx.get( + _TENCENT_KLINE_URL, + params=params, + headers=_TENCENT_HEADERS, + timeout=_FETCH_TIMEOUT_S, + trust_env=False, + ) + response.raise_for_status() + data = response.json() + payload = (data.get("data") or {}).get(compact_symbol) or {} + raw_rows = payload.get(period_name) or [] + if data.get("code") != 0 or not isinstance(raw_rows, list): + raise RuntimeError(f"tencent kline response invalid: {str(data)[:200]}") + if not raw_rows: + raise RuntimeError(f"tencent kline returned no rows for {symbol}") + return [_tencent_bar_row(row, intraday=False) for row in raw_rows if len(row) >= 6] + + minute_name = f"m{period}" + response = httpx.get( + _TENCENT_MINUTE_KLINE_URL, + params={"param": f"{compact_symbol},{minute_name},,{limit}"}, + headers=_TENCENT_HEADERS, + timeout=_FETCH_TIMEOUT_S, + trust_env=False, + ) + response.raise_for_status() + data = response.json() + payload = (data.get("data") or {}).get(compact_symbol) or {} + raw_rows = payload.get(minute_name) or [] + if data.get("code") != 0 or not isinstance(raw_rows, list): + raise RuntimeError(f"tencent minute kline response invalid: {str(data)[:200]}") + if not raw_rows: + raise RuntimeError(f"tencent minute kline returned no rows for {symbol}") + return [_tencent_bar_row(row, intraday=True) for row in raw_rows if len(row) >= 6] + + +def _tencent_bar_row(row: list[Any], *, intraday: bool) -> dict[str, Any]: + """把腾讯 K 线数组转换为 connector 通用字段。""" + time_value, open_s, close_s, high_s, low_s, volume_s = row[:6] + volume = _to_float(volume_s) + out: dict[str, Any] = { + "date": str(time_value)[:8] if intraday else str(time_value), + "open": open_s, + "close": close_s, + "high": high_s, + "low": low_s, + # 腾讯 K 线以“手”返回成交量;bars 的统一单位是“股”。 + "volume": volume * 100 if volume is not None else None, + } + if intraday: + out["time"] = str(time_value) + return out + + def _fetch_baostock_sync( *, symbol: str, @@ -1107,6 +1244,7 @@ def _fetch_sync( period: str, start_str: str, end_str: str, + limit: int, ) -> list[dict[str, Any]]: """同步调数据源 —— 按市场前缀路由。 @@ -1130,31 +1268,14 @@ def _fetch_sync( ) if prefix in ("sh", "sz"): - # A股走 baostock(证券宝)。东财 push2his 2026-07 起失效,腾讯源仅日线且无 volume; - # baostock 免费零 key、日/周/月/分钟全支持、有真实成交量。 - # symbol 格式 "sh.600519" / "sz.000001"(与 _parse_symbol 产物一致,零转换)。 - # period 已由上层 fetch_bars 经 _PERIOD_MAP 转为 "daily"/"weekly"/"monthly"/"5"/"15"/"30"/"60" - _baostock_freq = { - "daily": "d", - "weekly": "w", - "monthly": "m", - # 分钟级直接透传 - "5": "5", - "15": "15", - "30": "30", - "60": "60", - } - baostock_freq = _baostock_freq.get(period) - if baostock_freq is None: - raise NotImplementedError( - f"baostock does not support period {period!r}; " - f"supported: {sorted(_baostock_freq.keys())}" - ) - return _fetch_baostock_sync( + # 行情走腾讯 HTTPS。生产机无法访问 Baostock 专用 TCP 10030,但腾讯接口可直连, + # 同时覆盖指数、个股和日/周/月/分钟 K 线。Baostock 继续承载财报、日历与成分股。 + return _fetch_tencent_bars_sync( symbol=f"{prefix}.{code}", - frequency=baostock_freq, + period=period, start_date=start_str, end_date=end_str, + limit=limit, ) elif prefix == "hk": # 港股走东财 stock_hk_hist(push2his API,当前不可用;orchestrator 已将 @@ -1318,18 +1439,19 @@ def _parse_date(v: Any) -> datetime: return datetime(v.year, v.month, v.day, tzinfo=UTC) # 字符串解析 s = str(v) - # 分钟级:YYYYMMDDHHMMSS(14 位数字) - if len(s) == 14 and s.isdigit(): - # 20260709093500000 → 2026-07-09 09:35:00 - return datetime( + # 分钟级:腾讯可能返回 YYYYMMDDHHMM 或 YYYYMMDDHHMMSS。时间是北京时间,转 UTC 后落库。 + if len(s) in {12, 14} and s.isdigit(): + second = int(s[12:14]) if len(s) == 14 else 0 + local_dt = datetime( int(s[:4]), int(s[4:6]), int(s[6:8]), int(s[8:10]), int(s[10:12]), - int(s[12:14]), - tzinfo=UTC, + second, + tzinfo=_TZ_SHANGHAI, ) + return local_dt.astimezone(UTC) # 日级:YYYY-MM-DD 或 YYYYMMDD return datetime.fromisoformat(s).replace(tzinfo=UTC) diff --git a/services/data/src/inalpha_data/venues.py b/services/data/src/inalpha_data/venues.py index 7e7e02cf..8bc18693 100644 --- a/services/data/src/inalpha_data/venues.py +++ b/services/data/src/inalpha_data/venues.py @@ -1,25 +1,50 @@ -"""数据服务 venue 规范化。 +"""数据服务的 venue 与 symbol 规范化。 ``akshare`` 曾同时承载多个市场。A 股行情现已迁移到 ``baostock``;过渡期仍接受 -旧客户端发送的 ``akshare`` + ``sh./sz.``,但所有存储和 connector 路由统一使用新 venue。 +旧客户端发送的 ``akshare``,但 connector 路由和持久化统一使用 canonical market identity。 """ + from __future__ import annotations LEGACY_A_SHARE_VENUE = "akshare" A_SHARE_VENUE = "baostock" -_A_SHARE_PREFIXES = ("sh.", "sz.") +_A_SHARE_PREFIXES = frozenset({"sh", "sz"}) + + +def canonicalize_market_identity(venue: str, symbol: str) -> tuple[str, str]: + """返回用于 connector 和持久化的 canonical ``(venue, symbol)``。""" + normalized_venue = venue.strip().lower() + normalized_symbol = _canonicalize_a_share_symbol(symbol) + if normalized_symbol is not None and normalized_venue in { + LEGACY_A_SHARE_VENUE, + A_SHARE_VENUE, + }: + return A_SHARE_VENUE, normalized_symbol + return normalized_venue, symbol.strip() def canonicalize_venue(venue: str, symbol: str) -> str: """返回用于 connector 与持久化的 canonical venue。""" - normalized = venue.strip().lower() - if normalized == LEGACY_A_SHARE_VENUE and symbol.strip().lower().startswith(_A_SHARE_PREFIXES): - return A_SHARE_VENUE - return normalized + return canonicalize_market_identity(venue, symbol)[0] def is_legacy_a_share_venue(venue: str, symbol: str) -> bool: """旧 A 股 venue 是否被映射到 ``baostock``。""" - return venue.strip().lower() == LEGACY_A_SHARE_VENUE and canonicalize_venue( - venue, symbol - ) == A_SHARE_VENUE + normalized_venue, _ = canonicalize_market_identity(venue, symbol) + return venue.strip().lower() == LEGACY_A_SHARE_VENUE and normalized_venue == A_SHARE_VENUE + + +def _canonicalize_a_share_symbol(symbol: str) -> str | None: + """把 ``SH.600519`` / ``600519.SH`` 归一为 ``sh.600519``。""" + normalized = symbol.strip() + if "." not in normalized: + return None + + prefix, code = normalized.split(".", 1) + if prefix.lower() in _A_SHARE_PREFIXES and code.isdigit() and len(code) == 6: + return f"{prefix.lower()}.{code}" + + code, suffix = normalized.rsplit(".", 1) + if suffix.lower() in _A_SHARE_PREFIXES and code.isdigit() and len(code) == 6: + return f"{suffix.lower()}.{code}" + return None diff --git a/services/data/tests/test_backfill_router.py b/services/data/tests/test_backfill_router.py index 9711d781..dfef9ffd 100644 --- a/services/data/tests/test_backfill_router.py +++ b/services/data/tests/test_backfill_router.py @@ -11,6 +11,7 @@ **不**走真实 yfinance / baostock / alpaca / fred 网络 —— 这些是 connector 层职责(已在 ``test_connectors.py`` 覆盖);本文件只验 router/registry 装配是否正确。 """ + from __future__ import annotations from collections.abc import AsyncIterator @@ -42,7 +43,11 @@ async def fetch_bars( if start.tzinfo is None: start = start.replace(tzinfo=UTC) # 时间步长按 timeframe 粗略给一个,避免下次循环 cursor 越界产生死循环 - step = timedelta(days=1) if "d" in timeframe or "w" in timeframe or "mo" in timeframe else timedelta(hours=1) + step = ( + timedelta(days=1) + if "d" in timeframe or "w" in timeframe or "mo" in timeframe + else timedelta(hours=1) + ) return [ (start + step * i, 100.0 + i, 101.0 + i, 99.0 + i, 100.5 + i, 1000.0 + i) for i in range(3) @@ -75,6 +80,50 @@ def all_venues_client(app_with_all_venues_mocked: Any) -> TestClient: return TestClient(app_with_all_venues_mocked) +class _FailingConnector: + """模拟外部行情源网络故障。""" + + async def fetch_bars( + self, + symbol: str, + timeframe: str, + since: datetime, + limit: int = 1000, + ) -> list[tuple[datetime, float, float, float, float, float]]: + raise RuntimeError("upstream timeout") + + +@pytest.fixture +async def app_with_failing_baostock() -> AsyncIterator[Any]: + from inalpha_data.connectors import _base as _connectors_base + from inalpha_data.main import app + + async with app.router.lifespan_context(app): + _connectors_base._REGISTRY["baostock"] = _FailingConnector() + yield app + app.dependency_overrides.clear() + + +def test_backfill_upstream_failure_returns_502( + app_with_failing_baostock: Any, auth_headers: dict[str, str] +) -> None: + """上游网络失败必须显式 502,不能 200 + bars_fetched=0。""" + response = TestClient(app_with_failing_baostock).post( + "/backfill/bars", + headers=auth_headers, + json={ + "venue": "baostock", + "symbol": "sh.600518", + "timeframe": "1d", + "from_ts": "2026-04-01T00:00:00Z", + "to_ts": "2026-04-04T00:00:00Z", + }, + ) + + assert response.status_code == 502 + assert response.json()["code"] == "BARS_UPSTREAM_UNAVAILABLE" + + # ─── 5 venue 路由 round-trip ────────────────────────────────────────── @@ -100,7 +149,11 @@ def test_backfill_routes_each_venue_to_its_connector( 回归用户报"TSLA 报错只支持 Binance"的 root cause —— TS schema 解锁后,TSLA + yfinance + 1d 必须能通到底层 connector。 """ - unique_symbol = f"{symbol}-{uuid4().hex[:8]}" + unique_symbol = ( + f"sh.{uuid4().int % 1_000_000:06d}" + if venue == "baostock" + else f"{symbol}-{uuid4().hex[:8]}" + ) # 1d / 1h 跨度都用 3 天 → 不会触发 50k bars 上限 from_ts = "2026-04-01T00:00:00Z" to_ts = "2026-04-04T00:00:00Z" @@ -192,16 +245,17 @@ def test_backfill_rejects_timeframe_unsupported_by_venue( assert "1d" in supported_tfs # 两个 venue 至少都支持 1d -def test_legacy_akshare_alias_applies_baostock_minute_cap( +def test_legacy_akshare_alias_normalizes_symbol_and_applies_minute_cap( all_venues_client: TestClient, auth_headers: dict[str, str] ) -> None: - """旧 A 股 venue 别名必须与 baostock 走同一配额保护。""" + """旧 venue + Yahoo 后缀 symbol 必须统一写入 canonical Baostock identity。""" + symbol = "600518.SH" r = all_venues_client.post( "/backfill/bars", headers=auth_headers, json={ "venue": "akshare", - "symbol": f"sh.600519-{uuid4().hex[:8]}", + "symbol": symbol, "timeframe": "5m", "from_ts": "2026-01-01T00:00:00Z", "to_ts": "2026-04-01T00:00:00Z", @@ -210,6 +264,22 @@ def test_legacy_akshare_alias_applies_baostock_minute_cap( assert r.status_code == 200, r.json() assert r.json()["from_ts"] == "2026-03-25T00:00:00Z" + from inalpha_shared.db import get_conn + + from inalpha_data.storage.bars import count_bars + + async def _count() -> tuple[int, int]: + async with get_conn() as conn: + canonical = await count_bars(conn, "baostock", "sh.600518", "5m") + legacy = await count_bars(conn, "akshare", symbol, "5m") + return canonical, legacy + + import asyncio + + canonical_count, legacy_count = asyncio.run(_count()) + assert canonical_count > 0 + assert legacy_count == 0 + # ─── 增量 backfill:已缓存则从 max(ts) 续拉,不从 from_ts 全量重拉 ────────── diff --git a/services/data/tests/test_baostock_batch.py b/services/data/tests/test_baostock_batch.py deleted file mode 100644 index 3498b651..00000000 --- a/services/data/tests/test_baostock_batch.py +++ /dev/null @@ -1,84 +0,0 @@ -"""测试 baostock 是否支持批量查询(多股票逗号分隔)""" -from datetime import datetime, timedelta - -import baostock as bs - -# 登录 -lg = bs.login() -if lg.error_code != "0": - print(f"登录失败: {lg.error_msg}") - exit(1) - -print("登录成功") - -# 测试单股票查询 -print("\n=== 测试单股票查询 ===") -rs = bs.query_history_k_data_plus( - "sh.600519", # 贵州茅台 - "date,code,open,high,low,close,volume", - start_date=(datetime.now() - timedelta(days=7)).strftime("%Y-%m-%d"), - end_date=datetime.now().strftime("%Y-%m-%d"), - frequency="d", -) - -print(f"单股票查询结果: error_code={rs.error_code}") -count = 0 -while (rs.error_code == "0") & rs.next(): - count += 1 - if count <= 3: - print(f" {rs.get_row_data()}") -print(f" 共 {count} 条数据") - -# 测试多股票查询(逗号分隔) -print("\n=== 测试多股票查询(逗号分隔)===") -rs = bs.query_history_k_data_plus( - "sh.600519,sh.600036,sh.601318", # 茅台 + 招行 + 平安 - "date,code,open,high,low,close,volume", - start_date=(datetime.now() - timedelta(days=7)).strftime("%Y-%m-%d"), - end_date=datetime.now().strftime("%Y-%m-%d"), - frequency="d", -) - -print(f"多股票查询结果: error_code={rs.error_code}") -count = 0 -stocks = set() -while (rs.error_code == "0") & rs.next(): - count += 1 - row = rs.get_row_data() - stocks.add(row[1]) # code 字段 - if count <= 5: - print(f" {row}") -print(f" 共 {count} 条数据,涉及 {len(stocks)} 只股票: {stocks}") - -# 测试分钟 K 线批量查询 -print("\n=== 测试分钟 K 线批量查询 ===") -rs = bs.query_history_k_data_plus( - "sh.600519,sh.600036", # 茅台 + 招行 - "date,time,code,open,high,low,close,volume", - start_date=(datetime.now() - timedelta(days=1)).strftime("%Y-%m-%d"), - end_date=datetime.now().strftime("%Y-%m-%d"), - frequency="5", # 5 分钟 -) - -print(f"分钟 K 线查询结果: error_code={rs.error_code}") -count = 0 -stocks = set() -while (rs.error_code == "0") & rs.next(): - count += 1 - row = rs.get_row_data() - stocks.add(row[2]) # code 字段(分钟 K 线是第 3 列) - if count <= 5: - print(f" {row}") -print(f" 共 {count} 条数据,涉及 {len(stocks)} 只股票: {stocks}") - -# 登出 -bs.logout() -print("\n登出成功") - -# 结论 -print("\n=== 结论 ===") -if count > 0 and len(stocks) > 1: - print("✅ baostock 支持批量查询(多股票逗号分隔)") - print(" 可以在 1 次请求中获取多只股票数据,节省配额") -else: - print("❌ baostock 不支持批量查询,需要单独查询每只股票") diff --git a/services/data/tests/test_connectors.py b/services/data/tests/test_connectors.py index e859f812..ab3e31db 100644 --- a/services/data/tests/test_connectors.py +++ b/services/data/tests/test_connectors.py @@ -1,4 +1,5 @@ """multi-venue connector 路由 + baostock symbol 解析单测(不打网络)。""" + from __future__ import annotations import asyncio @@ -15,7 +16,7 @@ register_connector, unregister_connector, ) -from inalpha_data.connectors.baostock import _parse_symbol +from inalpha_data.connectors.baostock import _fetch_sync, _parse_date, _parse_symbol # ──────────────────────────────────────────────────────────────────── # Registry @@ -127,10 +128,225 @@ def test_parse_symbol_unknown_prefix_raises() -> None: def test_parse_symbol_empty_code_raises() -> None: - with pytest.raises(ValueError, match="code is empty"): + with pytest.raises(ValueError, match="exactly 6 digits"): _parse_symbol("sh.") +@pytest.mark.parametrize("raw", ["sh.600.519", "sh.60051", "600.519.SH"]) +def test_parse_symbol_rejects_noncanonical_a_share_code(raw: str) -> None: + """A 股代码必须严格为 6 位数字,不能让点号被腾讯 compact symbol 吞掉。""" + with pytest.raises(ValueError, match=r"exactly 6 digits|prefix"): + _parse_symbol(raw) + + +def test_baostock_intraday_timestamp_converts_shanghai_time_to_utc() -> None: + """腾讯分钟 K 线时间是北京时间,12/14 位格式都必须转换为 UTC。""" + expected = datetime(2026, 7, 22, 1, 35, tzinfo=UTC) + assert _parse_date("202607220935") == expected + assert _parse_date("20260722093500") == expected + + +def test_baostock_connector_uses_intraday_time_field(monkeypatch) -> None: # type: ignore[no-untyped-def] + """分钟 K 线不能误用仅含日期的 date 字段把整天压成同一主键。""" + from inalpha_data.connectors.baostock import BaostockConnector + + async def _fake_fetch(_self: object, **_kwargs: object) -> list[dict[str, str]]: + return [ + { + "date": "20260722", + "time": "20260722093500", + "open": "10.00", + "high": "10.10", + "low": "9.90", + "close": "10.05", + "volume": "100", + } + ] + + monkeypatch.setattr(BaostockConnector, "_throttled_fetch_sync", _fake_fetch) + rows = asyncio.run( + BaostockConnector().fetch_bars( + "sh.600519", + "5m", + datetime(2026, 7, 22, tzinfo=UTC), + limit=5, + ) + ) + + assert rows[0][0] == datetime(2026, 7, 22, 1, 35, tzinfo=UTC) + + +def test_baostock_fetch_ticker_parses_tencent_quote(monkeypatch) -> None: # type: ignore[no-untyped-def] + """A 股 fresh ticker 返回腾讯报价的真实时间与价格。""" + from inalpha_data.connectors.baostock import BaostockConnector + + fields = [""] * 31 + fields[1] = "上证指数" + fields[3] = "3867.03" + fields[30] = "20260722150000" + + class _Response: + content = f'v_sh000001="{"~".join(fields)}";'.encode("gbk") + + def raise_for_status(self) -> None: + return None + + def _fake_get(url: str, **kwargs: object) -> _Response: + assert url == "https://qt.gtimg.cn/q=sh000001" + assert kwargs["trust_env"] is False + return _Response() + + monkeypatch.setattr("httpx.get", _fake_get) + ts, price = asyncio.run(BaostockConnector().fetch_ticker("sh.000001")) + assert ts == datetime(2026, 7, 22, 7, 0, tzinfo=UTC) + assert price == 3867.03 + + +def test_baostock_fetch_ticker_rejects_missing_quote_fields(monkeypatch) -> None: # type: ignore[no-untyped-def] + """腾讯实时行情缺字段时必须上抛,不得伪造价格或时间。""" + from inalpha_data.connectors.baostock import BaostockConnector + + class _Response: + content = b'v_sh000001="1~index";' + + def raise_for_status(self) -> None: + return None + + monkeypatch.setattr("httpx.get", lambda *_args, **_kwargs: _Response()) + with pytest.raises(RuntimeError, match="A-share ticker unavailable"): + asyncio.run(BaostockConnector().fetch_ticker("sh.000001")) + + +@pytest.mark.parametrize( + ("timeframe", "expected_url", "expected_period", "expected_end"), + [ + ("1wk", "https://web.ifzq.gtimg.cn/appstock/app/fqkline/get", "week", "2026-07-23"), + ("1mo", "https://web.ifzq.gtimg.cn/appstock/app/fqkline/get", "month", "2026-07-23"), + ("5m", "https://ifzq.gtimg.cn/appstock/app/kline/mkline", "m5", None), + ("15m", "https://ifzq.gtimg.cn/appstock/app/kline/mkline", "m15", None), + ("30m", "https://ifzq.gtimg.cn/appstock/app/kline/mkline", "m30", None), + ("1h", "https://ifzq.gtimg.cn/appstock/app/kline/mkline", "m60", None), + ], +) +def test_baostock_bars_map_all_tencent_periods( + monkeypatch, + timeframe: str, + expected_url: str, + expected_period: str, + expected_end: str | None, +) -> None: # type: ignore[no-untyped-def] + """周/月/分钟周期必须映射到腾讯对应参数。""" + from inalpha_data.connectors.baostock import BaostockConnector + + captured: dict[str, object] = {} + intraday = timeframe not in {"1wk", "1mo"} + raw_time = "202607220935" if intraday else "2026-07-22" + row = [raw_time, "10.00", "10.05", "10.10", "9.90", "100"] + + class _Response: + def raise_for_status(self) -> None: + return None + + def json(self) -> dict[str, object]: + return {"code": 0, "data": {"sh600519": {expected_period: [row]}}} + + def _fake_get(url: str, **kwargs: object) -> _Response: + captured["url"] = url + captured["params"] = kwargs["params"] + return _Response() + + monkeypatch.setattr("httpx.get", _fake_get) + bars = asyncio.run( + BaostockConnector().fetch_bars( + "sh.600519", timeframe, datetime(2026, 7, 1, tzinfo=UTC), limit=5 + ) + ) + assert captured["url"] == expected_url + assert expected_period in str(captured["params"]) + if expected_end is not None: + assert "2026-07-01" in str(captured["params"]) + assert expected_end in str(captured["params"]) + assert str(captured["params"]).endswith(",5,'}") + assert len(bars) == 1 + assert bars[0][-1] == 10_000.0 + + +def test_baostock_bars_reject_invalid_tencent_payload(monkeypatch) -> None: # type: ignore[no-untyped-def] + """腾讯 K 线错误响应必须转换为显式上游异常。""" + from inalpha_data.connectors.baostock import BaostockConnector + + class _Response: + def raise_for_status(self) -> None: + return None + + def json(self) -> dict[str, object]: + return {"code": 1, "data": {}} + + monkeypatch.setattr("httpx.get", lambda *_args, **_kwargs: _Response()) + with pytest.raises(RuntimeError, match="A-share bars unavailable"): + asyncio.run( + BaostockConnector().fetch_bars( + "sh.600519", "1d", datetime(2026, 7, 1, tzinfo=UTC), limit=5 + ) + ) + + +def test_baostock_bars_use_https_feed_when_binary_service_is_unreachable(monkeypatch) -> None: # type: ignore[no-untyped-def] + """A 股 K 线不能依赖生产环境不可达的 Baostock TCP 10030 服务。""" + from inalpha_data.connectors import baostock + + captured: dict[str, object] = {} + + class _Response: + def raise_for_status(self) -> None: + return None + + def json(self) -> dict[str, object]: + return { + "code": 0, + "data": { + "sh000001": { + "day": [ + ["2026-07-21", "3835.68", "3864.37", "3864.60", "3834.72", "119330136"], + ["2026-07-22", "3839.67", "3861.65", "3884.44", "3839.67", "589368135"], + ] + } + }, + } + + def _fake_get(url: str, **kwargs: object) -> _Response: + captured["url"] = url + captured.update(kwargs) + return _Response() + + def _binary_feed_must_not_run(**_kwargs: object) -> list[dict[str, object]]: + raise AssertionError("Baostock TCP feed must not be used for bars") + + monkeypatch.setattr("httpx.get", _fake_get) + monkeypatch.setattr(baostock, "_fetch_baostock_sync", _binary_feed_must_not_run) + + rows = _fetch_sync( + prefix="sh", + code="000001", + period="daily", + start_str="20260715", + end_str="20260722", + limit=1000, + ) + + assert captured["url"] == "https://web.ifzq.gtimg.cn/appstock/app/fqkline/get" + assert captured["trust_env"] is False + assert captured["params"] == {"param": "sh000001,day,2026-07-15,2026-07-22,1000,"} + assert rows[-1] == { + "date": "2026-07-22", + "open": "3839.67", + "close": "3861.65", + "high": "3884.44", + "low": "3839.67", + "volume": 58936813500.0, + } + + # ──────────────────────────────────────────────────────────────────── # alpaca connector skipped when keys missing # ──────────────────────────────────────────────────────────────────── diff --git a/services/data/tests/test_fundamentals.py b/services/data/tests/test_fundamentals.py index fd7e5315..f7c3d772 100644 --- a/services/data/tests/test_fundamentals.py +++ b/services/data/tests/test_fundamentals.py @@ -1,4 +1,5 @@ """Tests for GET /fundamentals endpoint.""" + from __future__ import annotations import pytest @@ -13,9 +14,7 @@ def test_fundamentals_requires_auth(client: TestClient) -> None: assert r.status_code == 401 -def test_fundamentals_baostock_venue( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_fundamentals_baostock_venue(client: TestClient, auth_headers: dict[str, str]) -> None: """Mocked baostock connector returns financial data.""" from inalpha_data.connectors import baostock as bs @@ -55,12 +54,14 @@ async def mock_fin(symbol, as_of=None): def test_fundamentals_legacy_akshare_alias( client: TestClient, auth_headers: dict[str, str] ) -> None: - """旧 A 股 venue 在迁移窗口内仍路由到 baostock。""" + """旧 venue 和 Yahoo 后缀 symbol 在迁移窗口内统一路由到 baostock。""" from inalpha_data.connectors import baostock as bs original = bs._connector.fetch_financials + seen: list[str] = [] async def mock_fin(symbol, as_of=None): + seen.append(symbol) return {"venue": "baostock", "symbol": symbol, "available": False, "reason": "no data"} bs._connector.fetch_financials = mock_fin @@ -68,17 +69,16 @@ async def mock_fin(symbol, as_of=None): r = client.get( "/fundamentals", headers=auth_headers, - params={"venue": "akshare", "symbol": "sh.600519"}, + params={"venue": "akshare", "symbol": "600519.SH"}, ) assert r.status_code == 200 assert r.json()["venue"] == "baostock" + assert seen == ["sh.600519"] finally: bs._connector.fetch_financials = original -def test_fundamentals_yfinance_venue( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_fundamentals_yfinance_venue(client: TestClient, auth_headers: dict[str, str]) -> None: """GET /fundamentals with venue=yfinance works.""" from inalpha_data.connectors import yfinance_conn as yf @@ -109,9 +109,7 @@ async def mock_fin(symbol, as_of=None): yf._connector.fetch_financials = original -def test_fundamentals_unsupported_venue( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_fundamentals_unsupported_venue(client: TestClient, auth_headers: dict[str, str]) -> None: """Venue=binance should return 422.""" r = client.get( "/fundamentals", @@ -122,9 +120,7 @@ def test_fundamentals_unsupported_venue( assert "FUNDAMENTALS" in r.json()["code"] -def test_fundamentals_unavailable_data( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_fundamentals_unavailable_data(client: TestClient, auth_headers: dict[str, str]) -> None: """Connector returns available=False → endpoint returns 200 with available=False.""" from inalpha_data.connectors import baostock as bs diff --git a/services/data/tests/test_news.py b/services/data/tests/test_news.py index 210aeba3..a87a8a98 100644 --- a/services/data/tests/test_news.py +++ b/services/data/tests/test_news.py @@ -1,4 +1,5 @@ """Tests for GET /news endpoint — extended to support baostock venue.""" + from __future__ import annotations import pytest @@ -13,9 +14,7 @@ def test_news_requires_auth(client: TestClient) -> None: assert r.status_code == 401 -def test_news_yfinance_venue( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_news_yfinance_venue(client: TestClient, auth_headers: dict[str, str]) -> None: """GET /news with venue=yfinance returns results.""" from inalpha_data.connectors import yfinance_conn as yf @@ -47,9 +46,7 @@ async def mock_news(symbol, limit=20): yf._connector.fetch_news = original -def test_news_baostock_venue( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_news_baostock_venue(client: TestClient, auth_headers: dict[str, str]) -> None: """GET /news with venue=baostock returns A-share news.""" from inalpha_data.connectors import baostock as bs @@ -81,9 +78,33 @@ async def mock_news(symbol, limit=20): bs._connector.fetch_news = original -def test_news_unsupported_venue( +def test_news_legacy_akshare_alias_normalizes_symbol( client: TestClient, auth_headers: dict[str, str] ) -> None: + """旧 venue 和大小写前缀应统一传给 Baostock connector。""" + from inalpha_data.connectors import baostock as bs + + original = bs._connector.fetch_news + seen: list[str] = [] + + async def mock_news(symbol, limit=20): + seen.append(symbol) + return [] + + bs._connector.fetch_news = mock_news + try: + r = client.get( + "/news", + headers=auth_headers, + params={"venue": "akshare", "symbol": "SH.600519"}, + ) + assert r.status_code == 200 + assert seen == ["sh.600519"] + finally: + bs._connector.fetch_news = original + + +def test_news_unsupported_venue(client: TestClient, auth_headers: dict[str, str]) -> None: """Venue=binance should return 422.""" r = client.get( "/news", diff --git a/services/data/tests/test_symbol_search.py b/services/data/tests/test_symbol_search.py index a9f77148..b6bcfc6c 100644 --- a/services/data/tests/test_symbol_search.py +++ b/services/data/tests/test_symbol_search.py @@ -1,4 +1,5 @@ """symbol_search connector + GET /symbols/search 端点测试(不打外网)。""" + from __future__ import annotations import pytest @@ -39,7 +40,7 @@ async def test_cjk_query_hits_a_share_table(monkeypatch: pytest.MonkeyPatch) -> "symbol": "sh.600519", "name": "贵州茅台", "exchange": "XSHG", - "venue": "akshare", + "venue": "baostock", "quote_type": "EQUITY", } ] @@ -65,7 +66,7 @@ def fake_yahoo(query: str, max_results: int): # Yahoo 以完整 max_results 被查(不是 A股填剩的余额) assert yahoo_calls == [("平安", 4)] # 轮替合并:a1, y1, a2(无), y2 → A股命中 1 条时 Yahoo 两条都进结果 - assert [r["venue"] for r in out] == ["akshare", "yfinance", "yfinance"] + assert [r["venue"] for r in out] == ["baostock", "yfinance", "yfinance"] assert out[0]["symbol"] == "sz.000001" assert {r["symbol"] for r in out if r["venue"] == "yfinance"} == {"AAPL", "0700.HK"} @@ -93,11 +94,11 @@ def fake_load(): async def test_code_prefix_match_and_bj_skipped(monkeypatch: pytest.MonkeyPatch) -> None: monkeypatch.setattr(ss, "_load_a_share_table_sync", lambda: _FAKE_A_SHARE) conn = SymbolSearchConnector() - out = await conn.search("430047", venue="akshare") + out = await conn.search("430047", venue="baostock") assert out == [] # bj 代码被跳过而不是返回错误格式 - out2 = await conn.search("0005", venue="akshare") + out2 = await conn.search("0005", venue="baostock") assert [r["symbol"] for r in out2] == [] # 前缀不匹配(000001 不以 0005 开头) - out3 = await conn.search("00000", venue="akshare") + out3 = await conn.search("00000", venue="baostock") assert [r["symbol"] for r in out3] == ["sz.000001"] @@ -125,18 +126,18 @@ def fake_load(): monkeypatch.setattr(ss, "_load_a_share_table_sync", fake_load) conn = SymbolSearchConnector() - await conn.search("茅台", venue="akshare") - await conn.search("平安", venue="akshare") + await conn.search("茅台", venue="baostock") + await conn.search("平安", venue="baostock") assert len(loads) == 1 # 第二次走缓存 async def test_table_load_failure_returns_empty(monkeypatch: pytest.MonkeyPatch) -> None: def boom(): - raise RuntimeError("akshare down") + raise RuntimeError("baostock down") monkeypatch.setattr(ss, "_load_a_share_table_sync", boom) conn = SymbolSearchConnector() - out = await conn.search("茅台", venue="akshare") + out = await conn.search("茅台", venue="baostock") assert out == [] # fail-open @@ -145,9 +146,7 @@ def test_symbols_search_requires_auth(client: TestClient) -> None: assert r.status_code == 401 -def test_symbols_search_endpoint_shape( - client: TestClient, auth_headers: dict[str, str] -) -> None: +def test_symbols_search_endpoint_shape(client: TestClient, auth_headers: dict[str, str]) -> None: original = ss._connector.search async def mock_search(query: str, venue: str = "auto", max_results: int = 10): @@ -156,16 +155,14 @@ async def mock_search(query: str, venue: str = "auto", max_results: int = 10): "symbol": "sh.600519", "name": "贵州茅台", "exchange": "XSHG", - "venue": "akshare", + "venue": "baostock", "quote_type": "EQUITY", } ] ss._connector.search = mock_search try: - r = client.get( - "/symbols/search", headers=auth_headers, params={"query": "茅台"} - ) + r = client.get("/symbols/search", headers=auth_headers, params={"query": "茅台"}) assert r.status_code == 200 body = r.json() assert body["results"][0]["symbol"] == "sh.600519" diff --git a/services/data/tests/test_ticker_router.py b/services/data/tests/test_ticker_router.py index 621d5c76..8e47f268 100644 --- a/services/data/tests/test_ticker_router.py +++ b/services/data/tests/test_ticker_router.py @@ -2,12 +2,13 @@ D-9 ``ticker.py`` 从硬编码 binance 改成按 ``TickerCapable`` Protocol 鸭子分发。本测试覆盖: -- ``fetch_ticker`` 已实现的 venue(binance / yfinance / alpaca)走通 fresh=true 路径 -- 已注册但未实现 ``fetch_ticker`` 的 venue(akshare / fred)→ 422 +- ``fetch_ticker`` 已实现的 venue(binance / yfinance / alpaca / baostock)走通 fresh=true 路径 +- 已注册但未实现 ``fetch_ticker`` 的 venue(fred)→ 422 FRESH_NOT_SUPPORTED_FOR_VENUE + hint 提示切 fresh=false - venue 未注册 → 422 + supported 列表(与 /backfill/bars 错误形态一致) - fresh=false 路径仍然支持任意 venue(走 DB cache,由 conftest 的 binance mock 覆盖原路径) """ + from __future__ import annotations from collections.abc import AsyncIterator @@ -22,7 +23,10 @@ class _TickerCapableFake: - """fake connector,实现 ``TickerCapable`` —— fetch_ticker 返固定 (now, 123.45)。""" + """fake connector,实现 ``TickerCapable`` 并记录 canonical symbol。""" + + def __init__(self) -> None: + self.seen_symbols: list[str] = [] async def fetch_bars( self, @@ -34,6 +38,7 @@ async def fetch_bars( return [] async def fetch_ticker(self, symbol: str) -> tuple[datetime, float]: + self.seen_symbols.append(symbol) return datetime.now(UTC), 123.45 async def close(self) -> None: @@ -43,7 +48,7 @@ async def close(self) -> None: class _NoTickerFake: """fake connector,**只**实现 fetch_bars,不实现 fetch_ticker。 - 模拟 akshare / fred 这类不支持实时 ticker 的 venue。 + 模拟 fred 这类不支持实时 ticker 的 venue。 """ async def fetch_bars( @@ -61,15 +66,14 @@ async def close(self) -> None: @pytest.fixture async def app_with_ticker_capabilities() -> AsyncIterator[Any]: - """覆盖 registry:binance/yfinance/alpaca 用 _TickerCapableFake;akshare/fred 用 _NoTickerFake。""" + """覆盖 registry:binance/yfinance/alpaca/baostock 可取 ticker;fred 不可。""" from inalpha_data.connectors import _base as _connectors_base from inalpha_data.main import app async with app.router.lifespan_context(app): - for v in ("binance", "yfinance", "alpaca"): + for v in ("binance", "yfinance", "alpaca", "baostock"): _connectors_base._REGISTRY[v] = _TickerCapableFake() - for v in ("akshare", "fred"): - _connectors_base._REGISTRY[v] = _NoTickerFake() + _connectors_base._REGISTRY["fred"] = _NoTickerFake() yield app app.dependency_overrides.clear() @@ -89,6 +93,8 @@ def ticker_client(app_with_ticker_capabilities: Any) -> TestClient: ("binance", "BTC/USDT"), ("yfinance", "TSLA"), ("alpaca", "AAPL"), + ("baostock", "sh.600519"), + ("akshare", "600519.SH"), ], ) def test_ticker_fresh_true_routes_via_capability( @@ -97,10 +103,7 @@ def test_ticker_fresh_true_routes_via_capability( venue: str, symbol: str, ) -> None: - """3 个 TickerCapable venue × 各自 symbol → fresh=true 路径返 200 + 假价 + source 含 venue 名。 - - 回归用户报"特斯拉现价拿不到":yfinance 走 TickerCapable 分支,不再被硬编码挡。 - """ + """TickerCapable venue 和旧 A 股 alias 均返回外部报价。""" r = ticker_client.get( "/ticker", headers=auth_headers, @@ -112,6 +115,12 @@ def test_ticker_fresh_true_routes_via_capability( assert body["symbol"] == symbol assert body["price"] == 123.45 assert body["source"] == f"{venue}_ticker" + if venue == "akshare": + from inalpha_data.connectors import _base as _connectors_base + + connector = _connectors_base._REGISTRY["baostock"] + assert isinstance(connector, _TickerCapableFake) + assert connector.seen_symbols[-1] == "sh.600519" # fake 返 now() → stale_seconds 应接近 0 assert body["stale_seconds"] < 5 assert body["is_stale"] is False @@ -147,14 +156,38 @@ def test_ticker_fresh_true_marks_stale_when_quote_time_old( assert body["stale_seconds"] > 3 * 3600 +class _FailingTickerFake(_TickerCapableFake): + """模拟外部实时报价源故障。""" + + async def fetch_ticker(self, symbol: str) -> tuple[datetime, float]: + raise RuntimeError("upstream timeout") + + +def test_ticker_upstream_failure_returns_502( + app_with_ticker_capabilities: Any, auth_headers: dict[str, str] +) -> None: + """外部实时报价失败应返回标准 TICKER_UNAVAILABLE。""" + from inalpha_data.connectors import _base as _connectors_base + + _connectors_base._REGISTRY["baostock"] = _FailingTickerFake() + response = TestClient(app_with_ticker_capabilities).get( + "/ticker", + headers=auth_headers, + params={"venue": "baostock", "symbol": "sh.000001", "fresh": "true"}, + ) + + assert response.status_code == 502 + assert response.json()["code"] == "TICKER_UNAVAILABLE" + + # ─── 已注册但无 fetch_ticker 的 venue ─────────────────────────────── -@pytest.mark.parametrize("venue", ["akshare", "fred"]) +@pytest.mark.parametrize("venue", ["fred"]) def test_ticker_fresh_true_returns_friendly_error_for_non_ticker_venue( ticker_client: TestClient, auth_headers: dict[str, str], venue: str ) -> None: - """akshare / fred connector 没 fetch_ticker → 422 FRESH_NOT_SUPPORTED_FOR_VENUE + hint。""" + """fred connector 没 fetch_ticker → 422 FRESH_NOT_SUPPORTED_FOR_VENUE + hint。""" r = ticker_client.get( "/ticker", headers=auth_headers, @@ -184,13 +217,45 @@ def test_ticker_fresh_true_unknown_venue_lists_supported( assert body["code"] == "VALIDATION_ERROR" assert "bitfinex" in body["message"] supported = body["details"]["supported"] - for v in ("binance", "yfinance", "alpaca", "akshare", "fred"): + for v in ("binance", "yfinance", "alpaca", "baostock", "fred"): assert v in supported # ─── fresh=false 仍走 DB cache,不动 connector ──────────────────── +def test_ticker_legacy_alias_reads_canonical_db_symbol( + ticker_client: TestClient, auth_headers: dict[str, str] +) -> None: + """旧 venue + Yahoo 后缀应从 canonical Baostock namespace 读取。""" + import asyncio + + from inalpha_shared.db import get_conn + + from inalpha_data.storage.bars import insert_bars + + now = datetime.now(UTC).replace(microsecond=0, second=0) + + async def _insert() -> None: + async with get_conn() as conn: + await insert_bars( + conn, + "baostock", + "sh.600518", + "1h", + [(now - timedelta(minutes=1), 10.0, 11.0, 9.0, 10.5, 100.0)], + ) + + asyncio.run(_insert()) + r = ticker_client.get( + "/ticker", + headers=auth_headers, + params={"venue": "akshare", "symbol": "600518.SH", "fresh": "false"}, + ) + assert r.status_code == 200, r.json() + assert r.json()["price"] == 10.5 + + def test_ticker_fresh_false_still_returns_404_when_no_db_data( ticker_client: TestClient, auth_headers: dict[str, str] ) -> None: diff --git a/services/data/tests/test_venues.py b/services/data/tests/test_venues.py new file mode 100644 index 00000000..69d0bbfb --- /dev/null +++ b/services/data/tests/test_venues.py @@ -0,0 +1,50 @@ +"""venue 与 symbol 规范化单元测试。""" + +from inalpha_data.venues import ( + canonicalize_market_identity, + canonicalize_venue, + is_legacy_a_share_venue, +) + + +def test_canonicalize_legacy_a_share_prefix() -> None: + assert canonicalize_market_identity(" AKShare ", " SH.600519 ") == ( + "baostock", + "sh.600519", + ) + + +def test_canonicalize_legacy_a_share_yahoo_suffix() -> None: + assert canonicalize_market_identity("akshare", "600519.SH") == ( + "baostock", + "sh.600519", + ) + assert canonicalize_market_identity("akshare", "000001.sz") == ( + "baostock", + "sz.000001", + ) + + +def test_canonicalize_current_a_share_symbol() -> None: + assert canonicalize_market_identity("BAOSTOCK", "600519.SH") == ( + "baostock", + "sh.600519", + ) + + +def test_non_a_share_identity_only_normalizes_venue_whitespace() -> None: + assert canonicalize_market_identity(" YFINANCE ", " 0700.HK ") == ( + "yfinance", + "0700.HK", + ) + + +def test_malformed_a_share_symbols_are_not_canonicalized() -> None: + for symbol in ("sh.12345", "sh.1234567", "sh.ABCDEF", "sh.123.456", "12345.SH"): + assert canonicalize_market_identity("baostock", symbol) == ("baostock", symbol) + + +def test_legacy_alias_requires_a_share_symbol() -> None: + assert canonicalize_venue("akshare", "hk.00700") == "akshare" + assert is_legacy_a_share_venue("akshare", "sh.600519") is True + assert is_legacy_a_share_venue("baostock", "sh.600519") is False From 6f1ba88104c09fd21c5ecf9c220b3f89512d6fd1 Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 23 Jul 2026 15:34:44 +0800 Subject: [PATCH 2/4] =?UTF-8?q?fix(db):=20=E8=A7=84=E8=8C=83=E5=8C=96=20A?= =?UTF-8?q?=20=E8=82=A1=20bars=20symbol?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- .../0036_bars_a_share_symbol_canonical.py | 93 +++++++++++++++++++ services/data/tests/test_migration_0036.py | 85 +++++++++++++++++ 2 files changed, 178 insertions(+) create mode 100644 infra/migrations/versions/0036_bars_a_share_symbol_canonical.py create mode 100644 services/data/tests/test_migration_0036.py diff --git a/infra/migrations/versions/0036_bars_a_share_symbol_canonical.py b/infra/migrations/versions/0036_bars_a_share_symbol_canonical.py new file mode 100644 index 00000000..45d3a1eb --- /dev/null +++ b/infra/migrations/versions/0036_bars_a_share_symbol_canonical.py @@ -0,0 +1,93 @@ +"""规范化 A 股 bars 的历史 symbol 格式。 + +兼容期曾接受 ``SH.600519``、``600519.SH`` 等形式。venue 已统一为 ``baostock`` 后, +这些 symbol 仍会形成不同的持久化 namespace;本迁移将其合并为 ``sh.600519`` / ``sz.000001``。 +""" + +from __future__ import annotations + +from alembic import op + +revision: str = "0036" +down_revision: str | None = "0035" +branch_labels: str | tuple[str, ...] | None = None +depends_on: str | tuple[str, ...] | None = None + +_PREFIX_SYMBOL_PATTERN = "^(sh|sz)[.][0-9][0-9][0-9][0-9][0-9][0-9]$" +_SUFFIX_SYMBOL_PATTERN = "^[0-9][0-9][0-9][0-9][0-9][0-9][.](sh|sz)$" + +_CANONICAL_SYMBOL_SQL = f""" +CASE + WHEN btrim(symbol) ~* '{_PREFIX_SYMBOL_PATTERN}' THEN + lower(split_part(btrim(symbol), '.', 1)) || '.' || + split_part(btrim(symbol), '.', 2) + WHEN btrim(symbol) ~* '{_SUFFIX_SYMBOL_PATTERN}' THEN + lower(substring(btrim(symbol) from '\\.([^.]+)$')) || '.' || + regexp_replace(btrim(symbol), '\\.(sh|sz)$', '', 'i') +END +""" + + +def upgrade() -> None: + """把旧 venue / symbol 组合合并到 canonical Baostock identity。""" + op.execute( + f""" + INSERT INTO bars (ts, venue, symbol, timeframe, open, high, low, close, volume) + SELECT DISTINCT ON (ts, canonical_symbol, timeframe) + ts, + 'baostock', + canonical_symbol, + timeframe, + open, + high, + low, + close, + volume + FROM ( + SELECT + ts, + venue, + symbol, + timeframe, + open, + high, + low, + close, + volume, + {_CANONICAL_SYMBOL_SQL} AS canonical_symbol + FROM bars + WHERE venue IN ('akshare', 'baostock') + AND ( + btrim(symbol) ~* '{_PREFIX_SYMBOL_PATTERN}' + OR btrim(symbol) ~* '{_SUFFIX_SYMBOL_PATTERN}' + ) + ) AS candidates + WHERE canonical_symbol IS NOT NULL + AND NOT (venue = 'baostock' AND symbol = canonical_symbol) + ORDER BY + ts, + canonical_symbol, + timeframe, + CASE WHEN venue = 'baostock' THEN 0 ELSE 1 END, + symbol + ON CONFLICT (ts, venue, symbol, timeframe) DO NOTHING + """ + ) + op.execute( + f""" + DELETE FROM bars + WHERE venue IN ('akshare', 'baostock') + AND ( + btrim(symbol) ~* '{_PREFIX_SYMBOL_PATTERN}' + OR btrim(symbol) ~* '{_SUFFIX_SYMBOL_PATTERN}' + ) + AND NOT ( + venue = 'baostock' + AND symbol = ({_CANONICAL_SYMBOL_SQL}) + ) + """ + ) + + +def downgrade() -> None: + """symbol 原始大小写/前后缀不可无损恢复,降级保持 canonical 数据。""" diff --git a/services/data/tests/test_migration_0036.py b/services/data/tests/test_migration_0036.py new file mode 100644 index 00000000..7444f945 --- /dev/null +++ b/services/data/tests/test_migration_0036.py @@ -0,0 +1,85 @@ +"""0036 A 股 bars identity 规范化迁移集成测试。""" + +from __future__ import annotations + +import importlib.util +import sys +from datetime import UTC, datetime +from pathlib import Path +from types import SimpleNamespace +from typing import Any + +import psycopg +import pytest + +pytestmark = pytest.mark.integration + +_MIGRATION_PATH = ( + Path(__file__).resolve().parents[3] + / "infra" + / "migrations" + / "versions" + / "0036_bars_a_share_symbol_canonical.py" +) +_DB_URL = "postgresql://quant:devpass@localhost:5433/inalpha_test" + + +def _load_migration() -> Any: + spec = importlib.util.spec_from_file_location("migration_0036", _MIGRATION_PATH) + assert spec is not None and spec.loader is not None + module = importlib.util.module_from_spec(spec) + spec.loader.exec_module(module) + return module + + +def test_0036_merges_legacy_symbols_without_deleting_unrelated_rows(monkeypatch) -> None: # type: ignore[no-untyped-def] + """canonical 行优先;legacy 变体合并,非 A 股与畸形代码保留。""" + ts = datetime(2026, 7, 22, tzinfo=UTC) + rows = [ + ("baostock", "sh.600519", 10.0), + ("akshare", "SH.600519", 20.0), + ("baostock", "600519.SH", 30.0), + ("akshare", "000001.SZ", 40.0), + ("yfinance", "0700.HK", 50.0), + ("baostock", "sh.BAD", 60.0), + ] + + with psycopg.connect(_DB_URL) as conn: + with conn.cursor() as cursor: + cursor.execute( + "DELETE FROM bars WHERE ts = %s AND symbol IN " + "('sh.600519', 'SH.600519', '600519.SH', '000001.SZ', '0700.HK', 'sh.BAD')", + (ts,), + ) + cursor.executemany( + "INSERT INTO bars " + "(ts, venue, symbol, timeframe, open, high, low, close, volume) " + "VALUES (%s, %s, %s, '1d', %s, %s, %s, %s, 1)", + [(ts, venue, symbol, close, close, close, close) for venue, symbol, close in rows], + ) + conn.commit() + + def _execute(sql: str) -> None: + with conn.cursor() as cursor: + cursor.execute(str(sql)) + + monkeypatch.setitem( + sys.modules, "alembic", SimpleNamespace(op=SimpleNamespace(execute=_execute)) + ) + _load_migration().upgrade() + conn.commit() + with conn.cursor() as cursor: + cursor.execute( + "SELECT venue, symbol, close FROM bars WHERE ts = %s " + "AND symbol IN ('sh.600519', 'SH.600519', '600519.SH', 'sz.000001', " + "'000001.SZ', '0700.HK', 'sh.BAD') ORDER BY venue, symbol", + (ts,), + ) + result = cursor.fetchall() + + assert result == [ + ("baostock", "sh.600519", 10.0), + ("baostock", "sh.BAD", 60.0), + ("baostock", "sz.000001", 40.0), + ("yfinance", "0700.HK", 50.0), + ] From 0ed867da6a9e5279d252e8a5b95970ecbecc7539 Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 23 Jul 2026 15:34:45 +0800 Subject: [PATCH 3/4] =?UTF-8?q?test(services):=20=E8=A1=A5=E9=BD=90=20A=20?= =?UTF-8?q?=E8=82=A1=E8=B7=A8=E6=9C=8D=E5=8A=A1=E5=85=BC=E5=AE=B9=E5=A5=91?= =?UTF-8?q?=E7=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- .../paper/tests/test_currency_resolver.py | 6 ++++-- .../paper/tests/test_exchange_resolver.py | 6 ++++-- services/research/tests/test_debate.py | 4 +++- .../research/tests/test_fundamentals_pit.py | 4 ++-- services/research/tests/test_venue_routing.py | 21 +++++++++++++++++++ 5 files changed, 34 insertions(+), 7 deletions(-) create mode 100644 services/research/tests/test_venue_routing.py diff --git a/services/paper/tests/test_currency_resolver.py b/services/paper/tests/test_currency_resolver.py index 645acdac..66881e79 100644 --- a/services/paper/tests/test_currency_resolver.py +++ b/services/paper/tests/test_currency_resolver.py @@ -17,9 +17,11 @@ # 美股 / yfinance 无后缀 → USD ("yfinance", "AAPL", "USD"), ("alpaca", "TSLA", "USD"), - # A股 / 港股(akshare 前缀) + # A股 canonical venue + 迁移前旧值 + ("baostock", "sh.600519", "CNY"), + ("baostock", "sz.000001", "CNY"), ("akshare", "sh.600519", "CNY"), - ("akshare", "sz.000001", "CNY"), + # 旧 akshare 全球前缀只为已持久化记录保留 ("akshare", "hk.00700", "HKD"), # 全球单股(yfinance 后缀) ("yfinance", "005930.KS", "KRW"), diff --git a/services/paper/tests/test_exchange_resolver.py b/services/paper/tests/test_exchange_resolver.py index 6242ffe5..78f89aa2 100644 --- a/services/paper/tests/test_exchange_resolver.py +++ b/services/paper/tests/test_exchange_resolver.py @@ -9,9 +9,11 @@ @pytest.mark.parametrize( ("venue", "symbol", "expected"), [ - # akshare 前缀(sz 复用 XSHG) + # A股 canonical venue(sz 复用 XSHG) + ("baostock", "sh.600519", "XSHG"), + ("baostock", "sz.000001", "XSHG"), + # 迁移前持久化的 akshare 记录继续兼容 ("akshare", "sh.600519", "XSHG"), - ("akshare", "sz.000001", "XSHG"), ("akshare", "hk.00700", "XHKG"), ("akshare", "jp.6758", "XTKS"), ("akshare", "uk.VOD", "XLON"), diff --git a/services/research/tests/test_debate.py b/services/research/tests/test_debate.py index 28858d06..6c4802ea 100644 --- a/services/research/tests/test_debate.py +++ b/services/research/tests/test_debate.py @@ -591,8 +591,10 @@ def test_infer_asset_type_classifies_all_venues() -> None: assert infer_asset_type(venue="yfinance", symbol="AAPL") == "us_stock" assert infer_asset_type(venue="yfinance", symbol="SPY") == "us_stock" # A 股 + assert infer_asset_type(venue="baostock", symbol="sh.600519") == "cn_stock" + assert infer_asset_type(venue="baostock", symbol="sz.000001") == "cn_stock" + # 兼容迁移前持久化的旧 venue assert infer_asset_type(venue="akshare", symbol="sh.600519") == "cn_stock" - assert infer_asset_type(venue="akshare", symbol="sz.000001") == "cn_stock" # 港股 assert infer_asset_type(venue="akshare", symbol="hk.00700") == "hk_stock" # 日 / 英 / 德(akshare 归 global_stock) diff --git a/services/research/tests/test_fundamentals_pit.py b/services/research/tests/test_fundamentals_pit.py index 76fd9f49..24954fcf 100644 --- a/services/research/tests/test_fundamentals_pit.py +++ b/services/research/tests/test_fundamentals_pit.py @@ -20,7 +20,7 @@ async def test_get_fundamentals_threads_as_of() -> None: ) async with DataClient("http://data-mock.test", "t") as data: await data.get_fundamentals( - venue="akshare", + venue="baostock", symbol="sh.600519", as_of=datetime(2020, 1, 1, tzinfo=UTC), ) @@ -34,5 +34,5 @@ async def test_get_fundamentals_omits_as_of_when_none() -> None: return_value=Response(200, json={"available": True}) ) async with DataClient("http://data-mock.test", "t") as data: - await data.get_fundamentals(venue="akshare", symbol="sh.600519") + await data.get_fundamentals(venue="baostock", symbol="sh.600519") assert "as_of" not in route.calls.last.request.url.params diff --git a/services/research/tests/test_venue_routing.py b/services/research/tests/test_venue_routing.py new file mode 100644 index 00000000..8f517c22 --- /dev/null +++ b/services/research/tests/test_venue_routing.py @@ -0,0 +1,21 @@ +"""research 的市场分类与 fundamentals 路由契约测试。""" +from inalpha_research.analysts.utils import fundamentals_route +from inalpha_research.researchers.base import infer_asset_type + + +def test_baostock_a_share_routes_to_cn_fundamentals() -> None: + market_type = infer_asset_type(venue="baostock", symbol="sh.600519") + assert market_type == "cn_stock" + assert fundamentals_route(venue="baostock", market_type=market_type) == "baostock" + + +def test_legacy_akshare_a_share_routes_to_baostock_fundamentals() -> None: + market_type = infer_asset_type(venue="akshare", symbol="SH.600519") + assert market_type == "cn_stock" + assert fundamentals_route(venue="akshare", market_type=market_type) == "baostock" + + +def test_hk_fundamentals_use_yfinance() -> None: + market_type = infer_asset_type(venue="yfinance", symbol="0700.HK") + assert market_type == "hk_stock" + assert fundamentals_route(venue="yfinance", market_type=market_type) == "yfinance" From 929b981c19a5560a732f5384b52d9377551db28b Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 23 Jul 2026 15:34:45 +0800 Subject: [PATCH 4/4] =?UTF-8?q?docs(data):=20=E6=9B=B4=E6=96=B0=20A=20?= =?UTF-8?q?=E8=82=A1=E6=B7=B7=E5=90=88=E6=95=B0=E6=8D=AE=E6=BA=90=E5=A5=91?= =?UTF-8?q?=E7=BA=A6?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- .env.example | 7 ++++--- README.md | 6 +++--- README.zh-CN.md | 6 +++--- docs/04-current-state.md | 4 ++-- packages/orchestration/skills/cn-equity-research/SKILL.md | 2 +- packages/orchestration/src/tools/data.ts | 6 +++--- 6 files changed, 16 insertions(+), 15 deletions(-) diff --git a/.env.example b/.env.example index e48efde4..12cf52b0 100644 --- a/.env.example +++ b/.env.example @@ -103,12 +103,13 @@ FRED_API_KEY= # 给调度器一份「每天去给哪些指数拍成分快照」的点名册:逗号分隔指数代码,留空=禁用调度 # (后台循环不起,constituent_snapshot 表保持空,panel_score 传 index_code 一律走 non-PIT 降级)。 # -# 为什么要它:akshare 只回**当前**成分、免费历史拿不到 → 唯一 PIT 路径是从启用日起每日把 +# 为什么要它:免费源只回**当前**成分、历史成分拿不到 → 唯一 PIT 路径是从启用日起每日把 # 当前成分落库、向前累积。**覆盖从配置那天起向前增长,不回填历史,越早配可用于无偏回测的历史越长。** # -# 目前 connector 走 akshare 中证接口,只支持 **A股指数代码**(美股/港股成分快照需另接数据源,未做): +# 目前 connector 对三大核心指数优先走 Baostock,其余 A股指数回退 AkShare 中证接口; +# 美股/港股成分快照需另接数据源,未做: # 000300 沪深300 / 000905 中证500 / 000016 上证50 / 000852 中证1000 / 399006 创业板指 -# 不预设市场:按你实际要回测的指数填,别贪多(akshare 是爬虫接口,抓越多越易被封)。 +# 不预设市场:按你实际要回测的指数填,别贪多(公开接口抓越多越容易限流)。 # 手动 POST /constituents/snapshot 不受本项影响(仍可随时按需落一份)。 CONSTITUENT_SNAPSHOT_INDICES= # 调度检查间隔(小时,默认 12)。幂等:每轮只补「今天还没快照」的指数,<24h 不会重复打源站; diff --git a/README.md b/README.md index 718f421a..71d222a1 100644 --- a/README.md +++ b/README.md @@ -186,7 +186,7 @@ Inalpha splits *scheduling* from *compute*. The agent runtime fans out the grid A deep dive doesn't hand you one "correct answer." Beyond the usual technical, fundamental, and sentiment analysts, you can convene a panel of master personas — Buffett (value / moats), Lynch (GARP growth), Wood (disruptive innovation), Burry (contrarian / bubbles), Druckenmiller (macro trends), Marks (cycles / risk): each argues in their own style, naturally forming opposing views that feed a synthesized judgment. - **Opt-in, cost-controlled.** A plain deep dive costs the same as before; you only pay for the masters you actually convene. -- **Views grounded in data.** Each persona reads technicals / fundamentals / web intel with `as_of` pinned to *now* — no passing a stale forecast off as the present. Fundamentals are read point-in-time, filtered by report period and release lag so a backtest never sees a number before it was public (akshare today; yfinance v1 not yet PIT, flagged in place). +- **Views grounded in data.** Each persona reads technicals / fundamentals / web intel with `as_of` pinned to *now* — no passing a stale forecast off as the present. Fundamentals are read point-in-time, filtered by report period and release lag so a backtest never sees a number before it was public (Baostock today; yfinance v1 not yet PIT, flagged in place). - **Structured bull / bear / risk debate — now live.** Beyond the parallel analysts, opposing-stance bull and bear researchers argue across rounds while a risk researcher stress-tests both — triggered only when the analysts genuinely disagree, with a soft early-stop when arguments stop changing and the full decision chain (why it debated, why it stopped, how it was synthesized) persisted for replay. ### 6. Skills — absorb outside research playbooks @@ -234,7 +234,7 @@ Where each capability stands today. Live module inventory and the end-to-end dec | ✅ Shipped | Bull / bear researcher debate | D-9 | opposing-stance researchers under `services/research` | | ✅ Shipped | Scheduler / cron agent mode | D-9 | `scheduler_jobs` + advisory lock + `/api/scheduler/*` management plane | | ✅ Shipped | RiskGuard per-account isolation | D-9.1a | `RiskGuardFactory` removes cross-account state bleed | -| ✅ Shipped | Multi-market data sources — web search + financial fundamentals | D-10 | zero-key DDGS web search · akshare (A-shares/HK) + yfinance (global) fundamentals · analyst integration + fallback · per-market lookbackDays | +| ✅ Shipped | Multi-market data sources — web search + financial fundamentals | D-10 | zero-key DDGS web search · `baostock` is the logical A-share venue: Tencent HTTPS bars/ticker + Baostock fundamentals/calendar/constituents; yfinance covers global markets including HK · analyst integration + fallback · per-market lookbackDays | | ✅ Shipped | Risk engine — all 5 rules live in HTTP path | D-9 closed | `closed_trades` writes from HTTP order flow; `RoutingCalendar` for US equity + crypto; all trade-based rules trigger on real data | | ✅ Shipped | `askUserChoice` — `ask` permission path | D-11 (issue #2) | pending-permission flow resolves the `ask` state (no longer a workaround) | | ✅ Shipped | `permissions.yaml` configuration | D-11 (issue #4) | `config/permissions.default.yaml` + `yaml_loader.ts` replace the hard-coded `defaults.ts` | @@ -249,7 +249,7 @@ Where each capability stands today. Live module inventory and the end-to-end dec | ✅ Shipped | Factor discovery — L1 | D-12 | restricted qlib-style DSL (zero eval/exec) · `factor_candidates` pool · multiple-testing correction + null-IC benchmark · discovery workflow → propose; register is dashboard-only | | ✅ Shipped | Cross-sectional factor scoring | D-12 | `factor.panel_score` · `POST /panel/score` · cross-sectional Rank IC (rank the pool each period vs forward cross-sectional return) · native Alpha101 a1/a3 · orthogonal to single-name timing | | ✅ Shipped | Time-series cross-validation — anti-overfitting | D-12 | WalkForward / PurgedKFold / Combinatorial Purged CV + Deflated Sharpe · `POST /backtest/cv` · test fold always includes the latest bar · auto-fallback to walk-forward when samples are short | -| ✅ Shipped | Point-in-time fundamentals | D-12 | akshare financials filtered by report-period + release lag · `GET /fundamentals?as_of=` · prevents look-ahead (yfinance v1 not yet PIT, explicitly flagged) | +| ✅ Shipped | Point-in-time fundamentals | D-12 | Baostock financials filtered by actual publication date · `GET /fundamentals?as_of=` · prevents look-ahead (yfinance v1 not yet PIT, explicitly flagged) | | 🗓️ Planned | Strategy evolution — E2 | E2 | multi-generation loop + MAP-Elites + Island Model + `unified-diff` mutations (E1 single-generation closed loop already shipped in D-9) | | 🗓️ Planned | Factor discovery — L2 / L3 | L2 / L3 | multi-agent factor crew (L2) + weekly automated scans (L3), on top of the L1 DSL pipeline already shipped | | 🗓️ Planned | Automated decay handling | TBD | reflection-driven backtest + auto-trim of decaying factors — today the decay patrol only alerts, never moves the book | diff --git a/README.zh-CN.md b/README.zh-CN.md index 85b43c10..8f15748c 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -184,7 +184,7 @@ Inalpha 把*调度*和*算力*分开:agent runtime 负责扇出网格、聚合 一次深度研究不只给你一个"标准答案"。除了技术面、基本面、情绪面的常规 analyst,你还能叫上一支"大师团"——巴菲特(价值 / 护城河)、林奇(GARP 成长)、伍德(颠覆创新)、伯里(逆向 / 泡沫)、德鲁肯米勒(宏观趋势)、马克斯(周期 / 风险):各按自己的风格给观点,天然形成对立视角,再汇进综合判断。 - **按需启用,成本可控。** 不点名时普通研究成本不变;点了哪几位大师,才多跑那几次。 -- **观点要落到数据。** 大师视角接技术面 / 基本面 / web 情报,`as_of` 是真现在,不许拿过时预测当现在。财报按 point-in-time 读取,依报告期 + 发布滞后过滤,回测看不到尚未公开的数字(akshare 已生效;yfinance v1 尚未 PIT,已就地标注)。 +- **观点要落到数据。** 大师视角接技术面 / 基本面 / web 情报,`as_of` 是真现在,不许拿过时预测当现在。财报按 point-in-time 读取,依实际披露日期过滤,回测看不到尚未公开的数字(Baostock 已生效;yfinance v1 尚未 PIT,已就地标注)。 - **bull / bear / risk 结构化辩论——已上线。** 在并行 analyst 之外,立场对抗的多头与空头研究员多轮交锋,风险官同时压测双方——只在 analyst 真出现分歧时才触发,论点不再变化时软早停,完整决策链路(为什么辩、为什么停、如何综合)全程落盘可复盘。 ### 6. Skills · 吸收外部投研方法论 @@ -232,7 +232,7 @@ Inalpha 把*调度*和*算力*分开:agent runtime 负责扇出网格、聚合 | ✅ 已上线 | Bull / Bear 研究员辩论 | D-9 | `services/research` 立场对抗研究员 | | ✅ 已上线 | Scheduler / cron agent 模式 | D-9 | `scheduler_jobs` + advisory lock + `/api/scheduler/*` 管理面 | | ✅ 已上线 | RiskGuard 账户级隔离 | D-9.1a | `RiskGuardFactory` 去除跨账户状态串联 | -| ✅ 已上线 | 多市场数据源 — web 搜索 + 财报基本面 | D-10 | DDGS 零 key web search · akshare(A股/港股)+ yfinance(全球)基本面 · analyst 接入 + 兜底 · lookbackDays 按市场分化 | +| ✅ 已上线 | 多市场数据源 — web 搜索 + 财报基本面 | D-10 | DDGS 零 key web search · `baostock` 是 A 股逻辑 venue:腾讯 HTTPS 行情/最新价 + Baostock 基本面/日历/成分;yfinance 覆盖全球(含港股)· analyst 接入 + 兜底 · lookbackDays 按市场分化 | | ✅ 已上线 | 风控引擎 — 5 条规则全在 HTTP 路径激活 | D-9 收口 | `closed_trades` 由 HTTP 订单流写入;`RoutingCalendar` 覆盖美股 + crypto;trade-based 规则全在真实数据上触发 | | ✅ 已上线 | `askUserChoice` — `ask` 权限路径 | D-11(issue #2) | pending-permission 流程把 `ask` 状态从 workaround 救回(已收口) | | ✅ 已上线 | `permissions.yaml` 配置化 | D-11(issue #4) | `config/permissions.default.yaml` + `yaml_loader.ts` 替代 `defaults.ts` 硬编码 | @@ -247,7 +247,7 @@ Inalpha 把*调度*和*算力*分开:agent runtime 负责扇出网格、聚合 | ✅ 已上线 | 因子发现 — L1 | D-12 | 受限 qlib 风格 DSL(零 eval/exec)· `factor_candidates` 候选池 · 多重检验校正 + null IC 基准 · discovery workflow → propose;register 仅 dashboard 人工 | | ✅ 已上线 | 横截面因子打分 | D-12 | `factor.panel_score` · `POST /panel/score` · 横截面 Rank IC(每期对全池按因子排序 vs 跨标的前瞻收益)· Alpha101 a1/a3 原生 · 与单标的择时口径正交 | | ✅ 已上线 | 时序交叉验证 — 防过拟合 | D-12 | WalkForward / PurgedKFold / Combinatorial Purged CV + Deflated Sharpe · `POST /backtest/cv` · test 段始终含最新 bar · 样本不足自动回落 walk-forward | -| ✅ 已上线 | 财报 point-in-time | D-12 | akshare 财报按报告期 + 发布滞后过滤 · `GET /fundamentals?as_of=` · 防前视(yfinance v1 尚未 PIT,已显式标注) | +| ✅ 已上线 | 财报 point-in-time | D-12 | Baostock 财报按实际披露日期过滤 · `GET /fundamentals?as_of=` · 防前视(yfinance v1 尚未 PIT,已显式标注) | | 🗓️ 已规划 | 策略进化 — E2 | E2 | 多代演化 + MAP-Elites + Island Model + `unified-diff` 变异(E1 单代闭环已在 D-9 上线) | | 🗓️ 已规划 | 因子发现 — L2 / L3 | L2 / L3 | 多 agent 因子小组(L2)+ 每周自动扫描(L3),建立在已上线的 L1 DSL pipeline 之上 | | 🗓️ 已规划 | 因子衰减自动处置 | 待定 | 反思驱动回测 + 衰减因子自动剔除——当前衰减巡检只告警、绝不替你动仓 | diff --git a/docs/04-current-state.md b/docs/04-current-state.md index 8f24b8c3..09a4ec4b 100644 --- a/docs/04-current-state.md +++ b/docs/04-current-state.md @@ -430,8 +430,8 @@ manager 综合已覆盖)、debate 轮内并行化、reasoning 进活动流前 Sharpe 分布 + DSR;`POST /backtest/cv`(cpcv 数据不足自动回落 walk_forward,整体甩 ProcessPool 不阻塞事件循环);`paper.cv_backtest` tool。**单策略报 DSR、不报 PBO**(CSCV-PBO 需多配置作输入,留 grid 层)。 -- **财报 / bar Point-in-Time(阶段 A)**(ADR-0053 · #100):akshare 财报按"报告期末 + 120 天 - 发布滞后 ≤ as_of"过滤未披露报告期(防未来函数);`GET /fundamentals?as_of=`;paper +- **财报 / bar Point-in-Time(阶段 A)**(ADR-0053 · #100):Baostock 财报按实际公告日 + `pubDate ≤ as_of` 过滤未披露报告期(防未来函数);`GET /fundamentals?as_of=`;paper `DataClient.get_bars_pit` 把 bar 截断到 as_of。yfinance v1 不做 PIT 但响应显式标注;与 FRED 宏观 PIT(ADR-0044)口径对齐。 diff --git a/packages/orchestration/skills/cn-equity-research/SKILL.md b/packages/orchestration/skills/cn-equity-research/SKILL.md index fa65edcb..8ecdd4cf 100644 --- a/packages/orchestration/skills/cn-equity-research/SKILL.md +++ b/packages/orchestration/skills/cn-equity-research/SKILL.md @@ -24,7 +24,7 @@ fail 就 stop,不跨层补救**。同一结论要多源交叉验证,不信 | 公司名 → 可交易代码(一切个股研究的第一步) | `data.search_symbol`(A股返 sh./sz. 格式) | | 实时价格 / 当日涨跌 | `data.get_ticker` | | 历史 K 线(趋势 / 区间定位) | `data.get_bars`(保持默认 fresh) | -| PE_TTM / PB / ROE / 营收与利润增速 | `data.get_fundamentals`(venue=akshare) | +| PE_TTM / PB / ROE / 营收与利润增速 | `data.get_fundamentals`(venue=baostock) | | 市场级财经快讯(盘面消息) | `data.get_market_news` | | 行业板块涨跌幅榜(顺风 / 逆风) | `data.get_market_sectors` | | 当日强势股 + 题材标签(主线确认) | `data.get_market_movers` | diff --git a/packages/orchestration/src/tools/data.ts b/packages/orchestration/src/tools/data.ts index a5563820..09f2dc67 100644 --- a/packages/orchestration/src/tools/data.ts +++ b/packages/orchestration/src/tools/data.ts @@ -266,7 +266,7 @@ export const dataGetTickerTool = createTool({ 取 venue/symbol 的"最新价"(单值,不是 K 线)。 - fresh=true(**默认**):直接调外部市场实时报价,绕过 DB 缓存。 - **支持 venue**:binance / yfinance / alpaca;baostock / fred 不支持(返 + **支持 venue**:binance / yfinance / alpaca / baostock;fred 不支持(返 FRESH_NOT_SUPPORTED_FOR_VENUE,hint 提示切 fresh=false)。 网络抖动 ~200-800ms;不要高频循环调(rate-limit)。 - fresh=false:从 DB 拿最新 1m → fallback 1h,任意 venue 都支持。 @@ -281,8 +281,8 @@ export const dataGetTickerTool = createTool({ 何时不用: - 要历史走势 / 做技术分析 → 用 data.get_bars - 要 N 根 K 线 → 用 data.get_bars - - A 股 / FRED 想要"实时"价 → baostock/fred 没这能力,改 fresh=false 走 DB - cache(先 backfill_bars 灌一遍最新数据再调);港股走 yfinance 有实时 ticker + - A 股现价可用 baostock + fresh=true;FRED 没有盘中报价能力,改用 data.get_bars + 或显式 fresh=false 读 DB cache。港股走 yfinance 坑: - 默认 fresh=true!想要 DB cache 必须显式传 fresh:false