From d6c2bd944c8f330a9606cada17ff6a379124b9a1 Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 30 Jul 2026 11:11:42 +0800 Subject: [PATCH 1/4] =?UTF-8?q?feat(data):=20=E7=BB=9F=E4=B8=80=E5=A4=9A?= =?UTF-8?q?=E5=B8=82=E5=9C=BA=E8=B4=A2=E7=BB=8F=E6=96=B0=E9=97=BB=E4=B8=8E?= =?UTF-8?q?=E6=8A=AB=E9=9C=B2?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- services/data/pyproject.toml | 2 + services/data/src/inalpha_data/api/market.py | 29 ++- services/data/src/inalpha_data/api/news.py | 61 +++--- services/data/src/inalpha_data/config.py | 11 ++ .../src/inalpha_data/connectors/baostock.py | 8 +- .../inalpha_data/connectors/news/__init__.py | 5 + .../src/inalpha_data/connectors/news/base.py | 33 ++++ .../inalpha_data/connectors/news/dedupe.py | 73 ++++++++ .../connectors/news/feed_models.py | 61 ++++++ .../src/inalpha_data/connectors/news/hkex.py | 101 ++++++++++ .../connectors/news/hkex_parser.py | 66 +++++++ .../inalpha_data/connectors/news/legacy.py | 111 +++++++++++ .../connectors/news/market_proxy.py | 82 +++++++++ .../inalpha_data/connectors/news/router.py | 109 +++++++++++ .../src/inalpha_data/connectors/news/rss.py | 77 ++++++++ .../src/inalpha_data/connectors/news/sec.py | 85 +++++++++ .../connectors/news/sec_parser.py | 63 +++++++ .../inalpha_data/connectors/yfinance_conn.py | 6 +- services/data/src/inalpha_data/main.py | 20 ++ services/data/src/inalpha_data/news_models.py | 96 ++++++++++ services/data/src/inalpha_data/schemas.py | 38 ---- services/data/tests/test_market.py | 116 +++++++++++- services/data/tests/test_news.py | 174 +++++++++++++++++- services/data/tests/test_news_providers.py | 94 ++++++++++ services/data/uv.lock | 20 ++ 25 files changed, 1452 insertions(+), 89 deletions(-) create mode 100644 services/data/src/inalpha_data/connectors/news/__init__.py create mode 100644 services/data/src/inalpha_data/connectors/news/base.py create mode 100644 services/data/src/inalpha_data/connectors/news/dedupe.py create mode 100644 services/data/src/inalpha_data/connectors/news/feed_models.py create mode 100644 services/data/src/inalpha_data/connectors/news/hkex.py create mode 100644 services/data/src/inalpha_data/connectors/news/hkex_parser.py create mode 100644 services/data/src/inalpha_data/connectors/news/legacy.py create mode 100644 services/data/src/inalpha_data/connectors/news/market_proxy.py create mode 100644 services/data/src/inalpha_data/connectors/news/router.py create mode 100644 services/data/src/inalpha_data/connectors/news/rss.py create mode 100644 services/data/src/inalpha_data/connectors/news/sec.py create mode 100644 services/data/src/inalpha_data/connectors/news/sec_parser.py create mode 100644 services/data/src/inalpha_data/news_models.py create mode 100644 services/data/tests/test_news_providers.py diff --git a/services/data/pyproject.toml b/services/data/pyproject.toml index a7c6546e..386bb490 100644 --- a/services/data/pyproject.toml +++ b/services/data/pyproject.toml @@ -28,6 +28,8 @@ dependencies = [ # FRED 全球宏观(圣路易斯联储),免费 key 注册;80 万+ 经济序列 "fredapi>=0.5.2", "trafilatura>=2.1.0", + # RSS/Atom 仅做 bytes 解析;HTTP、超时和条件请求由 async httpx 管理 + "feedparser>=6.0.12", ] [dependency-groups] diff --git a/services/data/src/inalpha_data/api/market.py b/services/data/src/inalpha_data/api/market.py index afc482bc..4de763f5 100644 --- a/services/data/src/inalpha_data/api/market.py +++ b/services/data/src/inalpha_data/api/market.py @@ -25,6 +25,8 @@ from inalpha_shared.errors import InalphaError, ValidationError from ..connectors.cn_market import CnMarketConnector, CnMarketError, get_connector +from ..connectors.news import get_router +from ..news_models import NewsQuery, NewsResponse from ..schemas import ( MarketNewsItem, MarketNewsResponse, @@ -37,7 +39,10 @@ _logger = get_logger(__name__) router = APIRouter(tags=["market"]) -_SUPPORTED_MARKETS = ("cn",) +_SUPPORTED_MARKETS = ( + "cn", "us", "hk", "jp", "kr", "au", "in", "uk", "de", "fr", "ca", "br", "global", "crypto" +) +_CN_ONLY_MARKETS = ("cn",) class MarketDataUnavailableError(InalphaError): @@ -52,17 +57,29 @@ def _resolve(market: str) -> CnMarketConnector: raise ValidationError( f"market {market!r} not supported", code="MARKET_NOT_SUPPORTED", - details={"market": market, "supported": list(_SUPPORTED_MARKETS)}, + details={"market": market, "supported": list(_CN_ONLY_MARKETS)}, ) -@router.get("/market/news", response_model=MarketNewsResponse) +@router.get("/market/news", response_model=NewsResponse | MarketNewsResponse) async def market_news( _user: Annotated[User, Depends(get_current_user)], - market: Annotated[str, Query(description="市场,当前支持 cn")] = "cn", + market: Annotated[ + str, Query(description="市场代码;新闻支持全部已声明市场") + ] = "cn", limit: Annotated[int, Query(ge=1, le=50)] = 20, -) -> MarketNewsResponse: - """全市场财经快讯(东财 7×24),无需 symbol。""" +) -> NewsResponse | MarketNewsResponse: + """全市场财经快讯;非 A 股转调统一新闻路由。""" + if market not in _SUPPORTED_MARKETS: + raise ValidationError( + f"market {market!r} not supported", + code="MARKET_NOT_SUPPORTED", + details={"market": market, "supported": list(_SUPPORTED_MARKETS)}, + ) + if market != "cn": + return await get_router().fetch( + NewsQuery(market=market, limit=limit, kinds=["market_news", "media"]) + ) conn = _resolve(market) try: items = await conn.fetch_market_news(limit=limit) diff --git a/services/data/src/inalpha_data/api/news.py b/services/data/src/inalpha_data/api/news.py index 876dc5b0..41902240 100644 --- a/services/data/src/inalpha_data/api/news.py +++ b/services/data/src/inalpha_data/api/news.py @@ -1,8 +1,4 @@ -"""``GET /news`` —— 拉新闻头条(D-9,零 key,给 research analyst 喂真数据用)。 - -当前支持 venue=yfinance(全球)和 venue=baostock(A股)。 -""" - +"""``GET /news`` —— 统一多市场财经新闻与官方披露。""" from __future__ import annotations from typing import Annotated @@ -11,12 +7,15 @@ from inalpha_shared.auth import User, get_current_user from inalpha_shared.errors import ValidationError -from ..connectors import yfinance_conn -from ..connectors._base import get_connector_for_venue -from ..schemas import NewsItem, NewsQuery, NewsResponse +from ..connectors.news import get_router +from ..news_models import NewsQuery, NewsResponse from ..venues import canonicalize_market_identity router = APIRouter(tags=["news"]) +_SUPPORTED_MARKETS = { + "au", "br", "ca", "cn", "crypto", "de", "fr", "global", "hk", "in", "jp", "kr", "uk", "us" +} +_SUPPORTED_LEGACY_VENUES = {"yfinance", "baostock", "akshare"} @router.get("/news", response_model=NewsResponse) @@ -24,32 +23,28 @@ async def get_news( _user: Annotated[User, Depends(get_current_user)], query: Annotated[NewsQuery, Query()], ) -> NewsResponse: - """拉指定 ticker 的最新新闻。 - - venue 支持 yfinance / baostock(其它返 422);不支持 ticker 返空 list 而非错。 - """ - 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(effective_symbol, limit=query.limit) - except Exception: - raw = [] - elif effective_venue == "baostock": - conn = get_connector_for_venue("baostock") - if not hasattr(conn, "fetch_news"): - raise ValidationError( - f"news fetch not available for venue {query.venue!r}", - code="NEWS_FETCH_NOT_SUPPORTED", - details={"venue": query.venue}, - ) - raw = await conn.fetch_news(effective_symbol, limit=query.limit) # type: ignore[union-attr] - else: + """按市场和标的聚合新闻;旧 ``venue + symbol`` 请求保持兼容。""" + if query.market and query.market not in _SUPPORTED_MARKETS: + raise ValidationError( + f"news market {query.market!r} not supported", + code="NEWS_MARKET_NOT_SUPPORTED", + details={"market": query.market, "supported": sorted(_SUPPORTED_MARKETS)}, + ) + requested_venue = query.venue + requested_symbol = query.symbol + if not query.market and query.symbol and not query.venue: + query = query.model_copy(update={"venue": "yfinance"}) + requested_venue = "yfinance" + if query.venue and query.symbol: + venue, symbol = canonicalize_market_identity(query.venue, query.symbol) + query = query.model_copy(update={"venue": venue, "symbol": symbol}) + if not query.market and query.venue not in _SUPPORTED_LEGACY_VENUES: raise ValidationError( f"news venue {query.venue!r} not supported", code="NEWS_VENUE_NOT_SUPPORTED", - details={"venue": query.venue, "supported": ["yfinance", "baostock"]}, + details={"venue": query.venue, "supported": sorted(_SUPPORTED_LEGACY_VENUES)}, ) - - items = [NewsItem(**r) for r in raw] - return NewsResponse(venue=query.venue, symbol=query.symbol, items=items) + response = await get_router().fetch(query) + return response.model_copy( + update={"venue": requested_venue, "symbol": requested_symbol} + ) diff --git a/services/data/src/inalpha_data/config.py b/services/data/src/inalpha_data/config.py index 969666ac..489ffc61 100644 --- a/services/data/src/inalpha_data/config.py +++ b/services/data/src/inalpha_data/config.py @@ -75,6 +75,17 @@ class DataSettings(BaseSettings): """进程内缓存 TTL(秒)。快讯/板块榜分钟级更新,60s 挡住 analyst fan-out 同一轮重复打源站;响应带 fetched_at,fresh 语义不破。""" + news_timeout_s: float = Field(default=15.0, alias="NEWS_TIMEOUT_S") + """SEC、HKEX 与 RSS provider 的单请求超时。""" + + sec_user_agent: str = Field( + default="Inalpha/0.2 contact@inalpha.dev", alias="SEC_USER_AGENT" + ) + """SEC 要求可识别应用和联系方式;生产可覆盖为维护者邮箱。""" + + sec_min_interval_s: float = Field(default=0.11, alias="SEC_MIN_INTERVAL_S") + """SEC host 请求最小间隔;默认略低于每秒 10 次上限。""" + constituent_snapshot_indices: str = Field(default="", alias="CONSTITUENT_SNAPSHOT_INDICES") """每日成分快照追踪的指数代码,逗号分隔(如 ``000300,000905``)。空=禁用调度 (ADR-0053 阶段 C 向前累积:免费源只回当前成分,唯一 PIT 路径是从启用日起每日落库)。 diff --git a/services/data/src/inalpha_data/connectors/baostock.py b/services/data/src/inalpha_data/connectors/baostock.py index accbd18a..932139ab 100644 --- a/services/data/src/inalpha_data/connectors/baostock.py +++ b/services/data/src/inalpha_data/connectors/baostock.py @@ -538,7 +538,9 @@ async def fetch_news(self, symbol: str, limit: int = 20) -> list[dict[str, Any]] raw = await asyncio.to_thread(_fetch_news_sync, symbol=code) except Exception as exc: _logger.warning("baostock_news_fetch_failed", symbol=symbol, error=str(exc)) - return [] + raise RuntimeError( + f"eastmoney news for {symbol} unavailable: {exc}" + ) from exc if not raw or not isinstance(raw, list): return [] @@ -561,10 +563,12 @@ async def fetch_news(self, symbol: str, limit: int = 20) -> list[dict[str, Any]] published_at = dt_dt_dt.fromtimestamp(int(ts_raw), tz=UTC).isoformat() else: from datetime import datetime as dt_dt_dt + from zoneinfo import ZoneInfo published_at = ( dt_dt_dt.strptime(str(ts_raw)[:19], "%Y-%m-%d %H:%M:%S") - .replace(tzinfo=UTC) + .replace(tzinfo=ZoneInfo("Asia/Shanghai")) + .astimezone(UTC) .isoformat() ) except (ValueError, OSError): diff --git a/services/data/src/inalpha_data/connectors/news/__init__.py b/services/data/src/inalpha_data/connectors/news/__init__.py new file mode 100644 index 00000000..09d1c01e --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/__init__.py @@ -0,0 +1,5 @@ +"""统一多市场财经新闻 connector。""" + +from .router import NewsRouter, close_router, get_router, init_router + +__all__ = ["NewsRouter", "close_router", "get_router", "init_router"] diff --git a/services/data/src/inalpha_data/connectors/news/base.py b/services/data/src/inalpha_data/connectors/news/base.py new file mode 100644 index 00000000..d3f412a3 --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/base.py @@ -0,0 +1,33 @@ +"""财经新闻 provider 的共享类型与错误分类。""" +from __future__ import annotations + +from dataclasses import dataclass, field +from datetime import UTC, datetime +from typing import Protocol + +from ...news_models import NewsItem, NewsQuery, ProviderStatusCode + + +@dataclass(slots=True) +class ProviderResult: + """单个 provider 的标准化返回。""" + + provider: str + status: ProviderStatusCode + fetched_at: datetime = field(default_factory=lambda: datetime.now(UTC)) + items: list[NewsItem] = field(default_factory=list) + error: str | None = None + + +class NewsProvider(Protocol): + """新闻 provider 的最小接口。""" + + name: str + + async def fetch(self, query: NewsQuery) -> ProviderResult: + """拉取并标准化当前 provider 能覆盖的事件。""" + ... + + async def close(self) -> None: + """释放底层资源。""" + ... diff --git a/services/data/src/inalpha_data/connectors/news/dedupe.py b/services/data/src/inalpha_data/connectors/news/dedupe.py new file mode 100644 index 00000000..04864495 --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/dedupe.py @@ -0,0 +1,73 @@ +"""统一新闻去重与 point-in-time 过滤。""" +from __future__ import annotations + +import re +from datetime import UTC, datetime +from urllib.parse import parse_qsl, urlencode, urlsplit, urlunsplit + +from ...news_models import NewsItem, NewsQuery + +_TRACKING_KEYS = {"gclid", "fbclid", "ref", "source"} +_TIER_WEIGHT = {"official": 3, "professional_media": 2, "aggregator": 1} + + +def filter_and_dedupe(items: list[NewsItem], query: NewsQuery) -> list[NewsItem]: + """按时间窗、类型过滤并跨 provider 去重。""" + filtered = [item for item in items if _visible(item, query)] + winners: dict[str, NewsItem] = {} + for item in filtered: + key = _event_key(item) + current = winners.get(key) + if current is None: + winners[key] = item + continue + if _TIER_WEIGHT[item.source_tier] > _TIER_WEIGHT[current.source_tier]: + item.alternative_sources = _sources(current, item) + winners[key] = item + else: + current.alternative_sources = _sources(current, item) + epoch = datetime.min.replace(tzinfo=UTC) + return sorted(winners.values(), key=lambda item: item.published_at or epoch, reverse=True)[ + : query.limit + ] + + +def canonical_url(url: str) -> str: + """移除常见追踪参数并稳定化 URL。""" + if not url: + return "" + parts = urlsplit(url) + query = [ + (key, value) + for key, value in parse_qsl(parts.query, keep_blank_values=True) + if not key.lower().startswith("utm_") and key.lower() not in _TRACKING_KEYS + ] + return urlunsplit((parts.scheme.lower(), parts.netloc.lower(), parts.path, urlencode(query), "")) + + +def _visible(item: NewsItem, query: NewsQuery) -> bool: + if query.kinds and item.kind not in query.kinds: + return False + if query.as_of: + if item.published_at is None or item.published_at > query.as_of: + return False + if query.since and (item.published_at is None or item.published_at < query.since): + return False + return True + + +def _event_key(item: NewsItem) -> str: + if item.source_id: + return f"id:{item.source_name}:{item.source_id}" + url = canonical_url(item.link) + if url: + return f"url:{url}" + title = re.sub(r"\W+", "", item.title.casefold()) + bucket = item.published_at.strftime("%Y%m%d%H") if item.published_at else "unknown" + return f"title:{title}:{bucket}" + + +def _sources(left: NewsItem, right: NewsItem) -> list[str]: + values = [*left.alternative_sources, *right.alternative_sources] + values.extend(value for value in (left.source_name, right.source_name) if value) + return sorted(set(values)) diff --git a/services/data/src/inalpha_data/connectors/news/feed_models.py b/services/data/src/inalpha_data/connectors/news/feed_models.py new file mode 100644 index 00000000..9e634a38 --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/feed_models.py @@ -0,0 +1,61 @@ +"""RSS/Atom feed 定义与条目转换。""" +from __future__ import annotations + +import calendar +from dataclasses import dataclass +from datetime import UTC, datetime +from typing import Any + +from ...news_models import NewsItem, NewsQuery, SourceTier + + +@dataclass(frozen=True, slots=True) +class FeedDefinition: + """显式声明 feed 的身份与覆盖边界。""" + + id: str + name: str + url: str + tier: SourceTier + language: str + + +DEFAULT_CRYPTO_FEEDS = ( + FeedDefinition( + "coindesk", "CoinDesk", "https://www.coindesk.com/arc/outboundfeeds/rss/", + "professional_media", "en" + ), + FeedDefinition( + "kraken_blog", "Kraken Blog", "https://blog.kraken.com/feed", "official", "en" + ), +) + + +def parse_entry( + value: dict[str, Any], definition: FeedDefinition, query: NewsQuery, fetched_at: datetime +) -> NewsItem: + """把 feedparser entry 转为统一新闻条目。""" + return NewsItem( + title=str(value.get("title") or ""), + publisher=definition.name, + link=str(value.get("link") or ""), + published_at=_entry_time(value), + summary=str(value.get("summary") or "")[:500], + kind="media", + source_id=str(value.get("id") or value.get("guid") or ""), + source_name=definition.id, + source_tier=definition.tier, + fetched_at=fetched_at, + market="crypto", + language=definition.language, + ) + + +def _entry_time(value: dict[str, Any]) -> datetime | None: + parsed = value.get("published_parsed") or value.get("updated_parsed") + if parsed: + try: + return datetime.fromtimestamp(calendar.timegm(parsed), tz=UTC) + except (TypeError, ValueError, OverflowError): + return None + return None diff --git a/services/data/src/inalpha_data/connectors/news/hkex.py b/services/data/src/inalpha_data/connectors/news/hkex.py new file mode 100644 index 00000000..da2cc66b --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/hkex.py @@ -0,0 +1,101 @@ +"""HKEXnews 公告 provider。""" +from __future__ import annotations + +import asyncio +import json +import re +from datetime import UTC, datetime, timedelta +from typing import Any +from urllib.parse import urlencode + +import httpx + +from ...news_models import NewsQuery +from .base import ProviderResult +from .hkex_parser import parse_rows + +_BASE_URL = "https://www1.hkexnews.hk" + + +class HkexNewsProvider: + name = "hk" + + def __init__(self, *, timeout_s: float) -> None: + self._client = httpx.AsyncClient( + timeout=timeout_s, + trust_env=False, + headers={"User-Agent": "Mozilla/5.0 Inalpha"}, + ) + + async def fetch(self, query: NewsQuery) -> ProviderResult: + fetched_at = datetime.now(UTC) + if query.market != "hk" or not query.symbol: + return ProviderResult(self.name, "unsupported", fetched_at=fetched_at) + try: + symbol = query.symbol.split(".", 1)[0].lstrip("0") or "0" + stock_id = await self._resolve_stock_id(symbol) + if stock_id is None: + return ProviderResult(self.name, "no_results", fetched_at=fetched_at) + english, chinese = await asyncio.gather( + self._search("en", stock_id, query), self._search("zh", stock_id, query) + ) + items = parse_rows([*chinese, *english], query, fetched_at) + return ProviderResult( + self.name, "ok" if items else "no_results", fetched_at=fetched_at, items=items + ) + except httpx.TimeoutException as exc: + return ProviderResult(self.name, "timeout", fetched_at=fetched_at, error=str(exc)) + except httpx.HTTPStatusError as exc: + status = "rate_limited" if exc.response.status_code == 429 else "upstream_error" + return ProviderResult(self.name, status, fetched_at=fetched_at, error=str(exc)) + except Exception as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + + async def close(self) -> None: + """关闭 client。""" + await self._client.aclose() + + async def _resolve_stock_id(self, symbol: str) -> str | None: + response = await self._client.get( + f"{_BASE_URL}/search/prefix.do", + params={ + "callback": "callback", "lang": "EN", "type": "A", + "name": symbol, "market": "SEHK", + }, + ) + response.raise_for_status() + match = re.search(r"callback\((.*)\);?\s*$", response.text, re.S) + if not match: + raise ValueError("HKEX issuer lookup returned invalid JSONP") + payload = json.loads(match.group(1)) + for candidate in payload.get("stockInfo", []): + code = str(candidate.get("code", "")).lstrip("0") or "0" + if code == symbol: + return str(candidate.get("stockId")) + return None + + async def _search( + self, language: str, stock_id: str, query: NewsQuery + ) -> list[dict[str, Any]]: + end = query.as_of or datetime.now(UTC) + start = query.since or end - timedelta(days=365) + params = { + "lang": language, "sortDir": "0", "sortByOptions": "DateTime", + "category": "0", "market": "SEHK", "stockId": stock_id, + "documentType": "-1", "fromDate": start.strftime("%Y%m%d"), + "toDate": end.strftime("%Y%m%d"), "title": "", + } + response = await self._client.get( + f"{_BASE_URL}/search/titleSearchServlet.do?{urlencode(params)}" + ) + response.raise_for_status() + envelope = response.json() + result = envelope.get("result") if isinstance(envelope, dict) else None + if not isinstance(result, str): + raise ValueError("HKEX title search missing result") + rows = json.loads(result) + for row in rows: + row["_language"] = "zh-HK" if language == "zh" else "en-HK" + return rows diff --git a/services/data/src/inalpha_data/connectors/news/hkex_parser.py b/services/data/src/inalpha_data/connectors/news/hkex_parser.py new file mode 100644 index 00000000..5a4b0eaa --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/hkex_parser.py @@ -0,0 +1,66 @@ +"""HKEX 公告结果解析。""" +from __future__ import annotations + +import html +import re +from datetime import UTC, datetime +from typing import Any +from zoneinfo import ZoneInfo + +import httpx + +from ...news_models import NewsItem, NewsQuery + +_BASE_URL = "https://www1.hkexnews.hk" + + +def parse_rows( + rows: list[dict[str, Any]], query: NewsQuery, fetched_at: datetime +) -> list[NewsItem]: + """按 NEWS_ID 去重并转换公告条目。""" + seen: set[str] = set() + items: list[NewsItem] = [] + for row in rows: + news_id = str(row.get("NEWS_ID") or "").strip() + if not news_id or news_id in seen: + continue + seen.add(news_id) + published = parse_hk_time(row.get("DATE_TIME")) + file_link = str(row.get("FILE_LINK") or "").strip() + if not published or not file_link: + continue + items.append( + NewsItem( + title=html.unescape( + str(row.get("TITLE") or row.get("LONG_TEXT") or "HKEX announcement") + ), + publisher="Hong Kong Exchanges and Clearing", + link=str(httpx.URL(_BASE_URL).join(file_link)), + published_at=published, + kind="disclosure", + source_id=news_id, + source_name="hkexnews", + source_tier="official", + fetched_at=fetched_at, + market="hk", + symbols=[query.symbol] if query.symbol else [], + language=str(row.get("_language") or ""), + ) + ) + return items + + +def parse_hk_time(value: Any) -> datetime | None: + """把 HKEX ``DD/MM/YYYY HH:MM`` 香港时间转成 UTC。""" + match = re.match(r"^(\d{2})/(\d{2})/(\d{4})\s+(\d{2}):(\d{2})$", str(value or "")) + if not match: + return None + source = datetime( + int(match[3]), + int(match[2]), + int(match[1]), + int(match[4]), + int(match[5]), + tzinfo=ZoneInfo("Asia/Hong_Kong"), + ) + return source.astimezone(UTC) diff --git a/services/data/src/inalpha_data/connectors/news/legacy.py b/services/data/src/inalpha_data/connectors/news/legacy.py new file mode 100644 index 00000000..44318f1b --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/legacy.py @@ -0,0 +1,111 @@ +"""现有东财与 Yahoo 新闻能力的 provider 适配。""" +from __future__ import annotations + +from datetime import UTC, datetime + +from ...news_models import NewsItem, NewsQuery +from .. import yfinance_conn +from .._base import get_connector_for_venue +from ..cn_market import CnMarketError +from ..cn_market import get_connector as get_cn_market +from .base import ProviderResult + + +class CnNewsProvider: + """A 股市场快讯与个股东财新闻。""" + + name = "eastmoney" + + async def fetch(self, query: NewsQuery) -> ProviderResult: + fetched_at = datetime.now(UTC) + if query.market != "cn" and query.venue not in {"baostock", "akshare"}: + return ProviderResult(self.name, "unsupported", fetched_at=fetched_at) + try: + if query.symbol: + connector = get_connector_for_venue("baostock") + raw = await connector.fetch_news(query.symbol, limit=query.limit) # type: ignore[attr-defined] + items = _items(raw, query, fetched_at, "media", "eastmoney", "professional_media") + else: + raw = await get_cn_market().fetch_market_news(limit=query.limit) + items = _items( + raw, query, fetched_at, "market_news", "eastmoney", "professional_media" + ) + except CnMarketError as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + except Exception as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + return ProviderResult( + self.name, "ok" if items else "no_results", fetched_at=fetched_at, items=items + ) + + async def close(self) -> None: + """底层 connector 由既有生命周期管理。""" + + +class YahooNewsProvider: + """Yahoo Finance 全球 ticker 新闻聚合兜底。""" + + name = "yfinance" + + async def fetch(self, query: NewsQuery) -> ProviderResult: + fetched_at = datetime.now(UTC) + if not query.symbol or query.market == "crypto": + return ProviderResult(self.name, "unsupported", fetched_at=fetched_at) + try: + raw = await yfinance_conn.get_connector().fetch_news(query.symbol, limit=query.limit) + except Exception as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + items = _items(raw, query, fetched_at, "media", "yfinance", "aggregator") + return ProviderResult( + self.name, "ok" if items else "no_results", fetched_at=fetched_at, items=items + ) + + async def close(self) -> None: + """底层 connector 由既有生命周期管理。""" + + +def _items( + raw: list[dict[str, object]], + query: NewsQuery, + fetched_at: datetime, + kind: str, + source_name: str, + source_tier: str, +) -> list[NewsItem]: + """批量适配既有 connector 字段。""" + return [ + _item(value, query, fetched_at, kind, source_name, source_tier) for value in raw + ] + + +def _item( + value: dict[str, object], + query: NewsQuery, + fetched_at: datetime, + kind: str, + source_name: str, + source_tier: str, +) -> NewsItem: + """把既有 connector 字段转为统一模型。""" + published = value.get("published_at") + return NewsItem( + title=str(value.get("title") or ""), + publisher=str(value.get("publisher") or ""), + link=str(value.get("link") or ""), + published_at=published, # type: ignore[arg-type] + summary=str(value.get("summary") or ""), + kind=kind, # type: ignore[arg-type] + source_id=str(value.get("source_id") or ""), + source_name=source_name, + source_tier=source_tier, # type: ignore[arg-type] + fetched_at=fetched_at, + market=query.market, + symbols=[query.symbol] if query.symbol else [], + language=query.language, + ) diff --git a/services/data/src/inalpha_data/connectors/news/market_proxy.py b/services/data/src/inalpha_data/connectors/news/market_proxy.py new file mode 100644 index 00000000..6508a50b --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/market_proxy.py @@ -0,0 +1,82 @@ +"""Yahoo 代表性指数/ETF 的市场级新闻代理。""" +from __future__ import annotations + +from datetime import UTC, datetime + +from ...news_models import NewsItem, NewsQuery +from .. import yfinance_conn +from .base import ProviderResult + +_MARKET_PROXIES = { + "us": ("SPY", "S&P 500 market proxy"), + "hk": ("^HSI", "Hang Seng Index market proxy"), + "jp": ("^N225", "Nikkei 225 market proxy"), + "kr": ("^KS11", "KOSPI market proxy"), + "au": ("^AXJO", "S&P/ASX 200 market proxy"), + "in": ("^NSEI", "Nifty 50 market proxy"), + "uk": ("^FTSE", "FTSE 100 market proxy"), + "de": ("^GDAXI", "DAX market proxy"), + "fr": ("^FCHI", "CAC 40 market proxy"), + "ca": ("^GSPTSE", "S&P/TSX Composite market proxy"), + "br": ("^BVSP", "Bovespa market proxy"), + "global": ("ACWI", "MSCI ACWI ETF market proxy"), +} + + +class YahooMarketNewsProvider: + """用代表性市场载体聚合无 symbol 市场新闻。""" + + name = "yfinance_market_proxy" + + async def fetch(self, query: NewsQuery) -> ProviderResult: + """拉市场代理 ticker 新闻,并明确标注其代理性质。""" + fetched_at = datetime.now(UTC) + proxy = _MARKET_PROXIES.get(query.market or "") + if query.symbol or proxy is None: + return ProviderResult(self.name, "unsupported", fetched_at=fetched_at) + ticker, label = proxy + try: + raw = await yfinance_conn.get_connector().fetch_news(ticker, limit=query.limit) + except Exception as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + items = [_item(value, query, fetched_at, ticker, label) for value in raw] + return ProviderResult( + self.name, "ok" if items else "no_results", fetched_at=fetched_at, items=items + ) + + async def close(self) -> None: + """底层 Yahoo connector 由既有生命周期管理。""" + + +def _item( + value: dict[str, object], + query: NewsQuery, + fetched_at: datetime, + ticker: str, + label: str, +) -> NewsItem: + """把 Yahoo 结果转换为明确标记的市场代理新闻。""" + summary = str(value.get("summary") or "") + note = f"Market-level proxy via {ticker} ({label}); not a complete market newswire." + return NewsItem( + title=str(value.get("title") or ""), + publisher=str(value.get("publisher") or ""), + link=str(value.get("link") or ""), + published_at=value.get("published_at"), # type: ignore[arg-type] + summary=f"{note} {summary}".strip(), + kind="market_news", + source_id=str(value.get("source_id") or ""), + source_name=self_source(query.market), + source_tier="aggregator", + fetched_at=fetched_at, + market=query.market, + symbols=[ticker], + language=query.language, + ) + + +def self_source(market: str | None) -> str: + """给 provider 状态之外的条目生成稳定来源 ID。""" + return f"yfinance_{market}_market_proxy" diff --git a/services/data/src/inalpha_data/connectors/news/router.py b/services/data/src/inalpha_data/connectors/news/router.py new file mode 100644 index 00000000..023da483 --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/router.py @@ -0,0 +1,109 @@ +"""多市场财经新闻聚合路由。""" +from __future__ import annotations + +import asyncio +from datetime import UTC, datetime + +from ...news_models import NewsProviderStatus, NewsQuery, NewsResponse +from .base import NewsProvider, ProviderResult +from .dedupe import filter_and_dedupe +from .legacy import CnNewsProvider, YahooNewsProvider +from .market_proxy import YahooMarketNewsProvider + + +class NewsRouter: + """选择 provider,并保留部分失败与覆盖状态。""" + + def __init__(self, providers: list[NewsProvider]) -> None: + self._providers = providers + + async def fetch(self, query: NewsQuery) -> NewsResponse: + """并发调用 provider,聚合、PIT 过滤并去重。""" + fetched_at = datetime.now(UTC) + selected = self._select(query) + results = await asyncio.gather(*(provider.fetch(query) for provider in selected)) + items = filter_and_dedupe( + [item for result in results for item in result.items], query + ) + statuses = [_status(result) for result in results] + failures = {"timeout", "rate_limited", "upstream_error"} + return NewsResponse( + venue=query.venue, + market=query.market, + symbol=query.symbol, + as_of=query.as_of, + since=query.since, + fetched_at=fetched_at, + items=items, + providers=statuses, + is_partial=any(status.status in failures for status in statuses), + ) + + def _select(self, query: NewsQuery) -> list[NewsProvider]: + """按市场选择来源;未知市场仍允许 Yahoo ticker 兜底。""" + names: set[str] + if query.market == "cn" or query.venue in {"baostock", "akshare"}: + names = {"eastmoney"} + elif query.market in {"us", "hk"} and query.symbol: + wants_disclosures = not query.kinds or "disclosure" in query.kinds + names = {"yfinance", *({query.market} if wants_disclosures else set())} + elif query.symbol and query.market in { + "jp", "kr", "au", "in", "uk", "de", "fr", "ca", "br", "global" + }: + names = {"yfinance"} + elif query.market in {"us", "hk", "jp", "kr", "au", "in", "uk", "de", "fr", "ca", "br", "global"}: + names = {"yfinance_market_proxy"} + elif query.market == "crypto": + names = {"rss"} + elif query.venue == "yfinance" or query.symbol: + names = {"yfinance"} + else: + names = set() + return [ + provider + for provider in self._providers + if provider.name in names or ("rss" in names and provider.name.startswith("rss:")) + ] + + async def close(self) -> None: + """关闭自有 provider 资源。""" + await asyncio.gather(*(provider.close() for provider in self._providers)) + + +_router: NewsRouter | None = None + + +def init_router(extra_providers: list[NewsProvider] | None = None) -> None: + """初始化模块级 router。""" + global _router + if _router is not None: + raise RuntimeError("news router already initialized") + _router = NewsRouter( + [CnNewsProvider(), YahooNewsProvider(), YahooMarketNewsProvider(), *(extra_providers or [])] + ) + + +async def close_router() -> None: + """关闭并清除模块级 router。""" + global _router + if _router is not None: + await _router.close() + _router = None + + +def get_router() -> NewsRouter: + """返回已初始化 router。""" + if _router is None: + raise RuntimeError("news router not initialized") + return _router + + +def _status(result: ProviderResult) -> NewsProviderStatus: + """转换 provider 状态模型。""" + return NewsProviderStatus( + provider=result.provider, + status=result.status, + error=result.error, + fetched_at=result.fetched_at, + item_count=len(result.items), + ) diff --git a/services/data/src/inalpha_data/connectors/news/rss.py b/services/data/src/inalpha_data/connectors/news/rss.py new file mode 100644 index 00000000..2477d29b --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/rss.py @@ -0,0 +1,77 @@ +"""RSS/Atom 财经新闻 provider。""" +from __future__ import annotations + +import asyncio +from datetime import UTC, datetime + +import feedparser +import httpx + +from ...news_models import NewsItem, NewsQuery +from .base import ProviderResult +from .feed_models import FeedDefinition, parse_entry + + +class RssFeedProvider: + """单一 RSS/Atom feed;provider 状态精确到来源。""" + + def __init__(self, definition: FeedDefinition, *, timeout_s: float) -> None: + self.definition = definition + self.name = f"rss:{definition.id}" + self._client = httpx.AsyncClient( + timeout=timeout_s, + trust_env=False, + follow_redirects=True, + headers={"User-Agent": "Inalpha/0.2 financial-news"}, + ) + self._etag: str | None = None + self._last_modified: str | None = None + self._cached: list[NewsItem] = [] + + async def fetch(self, query: NewsQuery) -> ProviderResult: + """条件请求 feed,并标准化条目时间和来源。""" + fetched_at = datetime.now(UTC) + if query.market != "crypto" or query.symbol: + return ProviderResult(self.name, "unsupported", fetched_at=fetched_at) + headers = {} + if self._etag: + headers["If-None-Match"] = self._etag + if self._last_modified: + headers["If-Modified-Since"] = self._last_modified + try: + response = await self._client.get(self.definition.url, headers=headers) + if response.status_code == 304: + return ProviderResult( + self.name, + "ok" if self._cached else "no_results", + fetched_at=fetched_at, + items=self._cached, + ) + response.raise_for_status() + parsed = await asyncio.to_thread(feedparser.parse, response.content) + if parsed.bozo and not parsed.entries: + raise ValueError(f"malformed feed: {parsed.bozo_exception}") + self._etag = response.headers.get("etag") + self._last_modified = response.headers.get("last-modified") + self._cached = [ + parse_entry(item, self.definition, query, fetched_at) for item in parsed.entries + ] + return ProviderResult( + self.name, + "ok" if self._cached else "no_results", + fetched_at=fetched_at, + items=self._cached, + ) + except httpx.TimeoutException as exc: + return ProviderResult(self.name, "timeout", fetched_at=fetched_at, error=str(exc)) + except httpx.HTTPStatusError as exc: + status = "rate_limited" if exc.response.status_code == 429 else "upstream_error" + return ProviderResult(self.name, status, fetched_at=fetched_at, error=str(exc)) + except Exception as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + + async def close(self) -> None: + """关闭 feed HTTP client。""" + await self._client.aclose() diff --git a/services/data/src/inalpha_data/connectors/news/sec.py b/services/data/src/inalpha_data/connectors/news/sec.py new file mode 100644 index 00000000..72b81757 --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/sec.py @@ -0,0 +1,85 @@ +"""SEC EDGAR 官方披露 provider。""" +from __future__ import annotations + +import asyncio +from datetime import UTC, datetime +from time import monotonic +from typing import Any + +import httpx + +from ...news_models import NewsQuery +from .base import ProviderResult +from .sec_parser import parse_submissions + +_TICKERS_URL = "https://www.sec.gov/files/company_tickers.json" +_SUBMISSIONS_URL = "https://data.sec.gov/submissions/CIK{cik}.json" + + +class SecNewsProvider: + """通过 SEC 官方 JSON 获取美国上市公司披露。""" + + name = "us" + + def __init__(self, *, user_agent: str, timeout_s: float, min_interval_s: float) -> None: + self._client = httpx.AsyncClient( + timeout=timeout_s, + trust_env=False, + headers={"User-Agent": user_agent, "Accept-Encoding": "gzip, deflate"}, + ) + self._min_interval_s = min_interval_s + self._request_lock = asyncio.Lock() + self._last_request_at = 0.0 + self._ticker_map: dict[str, str] | None = None + + async def fetch(self, query: NewsQuery) -> ProviderResult: + """返回 ticker 截至查询时点的近期 SEC filings。""" + fetched_at = datetime.now(UTC) + if query.market != "us" or not query.symbol: + return ProviderResult(self.name, "unsupported", fetched_at=fetched_at) + try: + cik = await self._resolve_cik(query.symbol) + if cik is None: + return ProviderResult(self.name, "no_results", fetched_at=fetched_at) + payload = await self._get_json(_SUBMISSIONS_URL.format(cik=cik)) + items = parse_submissions(payload, query, fetched_at, cik) + return ProviderResult( + self.name, "ok" if items else "no_results", fetched_at=fetched_at, items=items + ) + except httpx.TimeoutException as exc: + return ProviderResult(self.name, "timeout", fetched_at=fetched_at, error=str(exc)) + except httpx.HTTPStatusError as exc: + status = "rate_limited" if exc.response.status_code == 429 else "upstream_error" + return ProviderResult(self.name, status, fetched_at=fetched_at, error=str(exc)) + except Exception as exc: + return ProviderResult( + self.name, "upstream_error", fetched_at=fetched_at, error=str(exc) + ) + + async def close(self) -> None: + """关闭 SEC HTTP client。""" + await self._client.aclose() + + async def _resolve_cik(self, symbol: str) -> str | None: + if self._ticker_map is None: + payload = await self._get_json(_TICKERS_URL) + self._ticker_map = { + str(value.get("ticker", "")).upper(): str(value.get("cik_str", "")).zfill(10) + for value in payload.values() + if isinstance(value, dict) + } + ticker = symbol.split(".", 1)[0].upper() + return self._ticker_map.get(ticker) + + async def _get_json(self, url: str) -> dict[str, Any]: + async with self._request_lock: + wait_s = self._min_interval_s - (monotonic() - self._last_request_at) + if wait_s > 0: + await asyncio.sleep(wait_s) + response = await self._client.get(url) + self._last_request_at = monotonic() + response.raise_for_status() + payload = response.json() + if not isinstance(payload, dict): + raise ValueError("SEC returned a non-object payload") + return payload diff --git a/services/data/src/inalpha_data/connectors/news/sec_parser.py b/services/data/src/inalpha_data/connectors/news/sec_parser.py new file mode 100644 index 00000000..91ae1de4 --- /dev/null +++ b/services/data/src/inalpha_data/connectors/news/sec_parser.py @@ -0,0 +1,63 @@ +"""SEC submissions JSON 转统一披露事件。""" +from __future__ import annotations + +from datetime import UTC, datetime +from typing import Any + +from ...news_models import NewsItem, NewsQuery + +_ARCHIVES_URL = "https://www.sec.gov/Archives/edgar/data" + + +def parse_submissions( + payload: dict[str, Any], query: NewsQuery, fetched_at: datetime, cik: str +) -> list[NewsItem]: + """把 submissions.recent 的并行数组转成披露事件。""" + recent = payload.get("filings", {}).get("recent", {}) + if not isinstance(recent, dict): + return [] + accessions = recent.get("accessionNumber", []) + items: list[NewsItem] = [] + for index, accession in enumerate(accessions): + if not accession: + continue + accepted = _parse_datetime(_at(recent, "acceptanceDateTime", index)) + published = accepted or _parse_datetime(_at(recent, "filingDate", index)) + primary = str(_at(recent, "primaryDocument", index) or "") + form = str(_at(recent, "form", index) or "Filing") + accession_path = str(accession).replace("-", "") + link = f"{_ARCHIVES_URL}/{int(cik)}/{accession_path}/{primary}" if primary else "" + items.append( + NewsItem( + title=f"SEC {form}: {payload.get('name') or query.symbol}", + publisher="U.S. Securities and Exchange Commission", + link=link, + published_at=published, + accepted_at=accepted, + summary=f"Official SEC filing {form}; primary document: {primary}", + kind="disclosure", + source_id=str(accession), + source_name="sec_edgar", + source_tier="official", + fetched_at=fetched_at, + market="us", + symbols=[query.symbol] if query.symbol else [], + language="en", + ) + ) + return items + + +def _at(data: dict[str, Any], key: str, index: int) -> Any: + values = data.get(key) + return values[index] if isinstance(values, list) and index < len(values) else None + + +def _parse_datetime(value: Any) -> datetime | None: + if not value: + return None + try: + parsed = datetime.fromisoformat(str(value).replace("Z", "+00:00")) + except ValueError: + return None + return parsed.replace(tzinfo=UTC) if parsed.tzinfo is None else parsed.astimezone(UTC) diff --git a/services/data/src/inalpha_data/connectors/yfinance_conn.py b/services/data/src/inalpha_data/connectors/yfinance_conn.py index fe2d7177..db89090d 100644 --- a/services/data/src/inalpha_data/connectors/yfinance_conn.py +++ b/services/data/src/inalpha_data/connectors/yfinance_conn.py @@ -343,9 +343,9 @@ def _fetch_news_sync(symbol: str, limit: int) -> list[dict[str, Any]]: ticker = yf.Ticker(symbol) try: raw_news = ticker.news or [] - except Exception: - _logger.warning("yfinance_news_fetch_failed", symbol=symbol) - return [] + except Exception as exc: + _logger.warning("yfinance_news_fetch_failed", symbol=symbol, error=str(exc)) + raise RuntimeError(f"yfinance news for {symbol} unavailable: {exc}") from exc out: list[dict[str, Any]] = [] for item in raw_news[:limit]: diff --git a/services/data/src/inalpha_data/main.py b/services/data/src/inalpha_data/main.py index de072104..d6d88200 100644 --- a/services/data/src/inalpha_data/main.py +++ b/services/data/src/inalpha_data/main.py @@ -39,10 +39,15 @@ from .connectors import binance as binance_conn from .connectors import cn_market as cn_market_conn from .connectors import fred as fred_conn +from .connectors import news as news_conn from .connectors import symbol_search as symbol_search_conn from .connectors import web_fetch as web_fetch_conn from .connectors import web_search as web_search_conn from .connectors import yfinance_conn +from .connectors.news.feed_models import DEFAULT_CRYPTO_FEEDS +from .connectors.news.hkex import HkexNewsProvider +from .connectors.news.rss import RssFeedProvider +from .connectors.news.sec import SecNewsProvider from .scheduler import ConstituentSnapshotScheduler, parse_indices _settings = get_data_settings() @@ -70,6 +75,20 @@ async def lifespan(_app: FastAPI) -> AsyncIterator[None]: fred_conn.init_connector(api_key=_settings.fred_api_key) web_search_conn.init_connector() cn_market_conn.init_connector() + news_conn.init_router( + [ + SecNewsProvider( + user_agent=_settings.sec_user_agent, + timeout_s=_settings.news_timeout_s, + min_interval_s=_settings.sec_min_interval_s, + ), + HkexNewsProvider(timeout_s=_settings.news_timeout_s), + *[ + RssFeedProvider(feed, timeout_s=_settings.news_timeout_s) + for feed in DEFAULT_CRYPTO_FEEDS + ], + ] + ) web_fetch_conn.init_connector() symbol_search_conn.init_connector() # 成分快照每日调度(ADR-0053 阶段 C 向前累积)——无追踪指数则自动禁用 @@ -83,6 +102,7 @@ async def lifespan(_app: FastAPI) -> AsyncIterator[None]: finally: await snapshot_scheduler.stop() await symbol_search_conn.close_connector() + await news_conn.close_router() await web_fetch_conn.close_connector() await fred_conn.close_connector() await yfinance_conn.close_connector() diff --git a/services/data/src/inalpha_data/news_models.py b/services/data/src/inalpha_data/news_models.py new file mode 100644 index 00000000..a33d1ef0 --- /dev/null +++ b/services/data/src/inalpha_data/news_models.py @@ -0,0 +1,96 @@ +"""统一财经新闻的数据契约。""" +from __future__ import annotations + +from datetime import UTC, datetime +from typing import Literal + +from pydantic import BaseModel, Field, field_validator, model_validator + +NewsKind = Literal["market_news", "media", "disclosure"] +SourceTier = Literal["official", "professional_media", "aggregator"] +ProviderStatusCode = Literal[ + "ok", "no_results", "timeout", "rate_limited", "upstream_error", "unsupported" +] + + +class NewsQuery(BaseModel): + """``GET /news`` 查询参数;兼容旧 ``venue + symbol`` 调用。""" + + venue: str | None = Field(default=None) + market: str | None = Field(default=None) + symbol: str | None = Field(default=None) + as_of: datetime | None = Field(default=None) + since: datetime | None = Field(default=None) + kinds: list[NewsKind] | None = Field(default=None) + language: str | None = Field(default=None, max_length=35) + limit: int = Field(default=10, ge=1, le=50) + + @field_validator("kinds", mode="before") + @classmethod + def split_kinds(cls, value: object) -> object: + """兼容 HTTP client 发送的逗号分隔 kinds。""" + values = [value] if isinstance(value, str) else value + if not isinstance(values, list): + return values + return [part.strip() for item in values for part in str(item).split(",") if part.strip()] + + @field_validator("as_of", "since", mode="after") + @classmethod + def assume_utc_if_naive(cls, value: datetime | None) -> datetime | None: + """与 bars 契约一致:无时区输入按 UTC 解释。""" + if value is None: + return None + return value.replace(tzinfo=UTC) if value.tzinfo is None else value.astimezone(UTC) + + @model_validator(mode="after") + def validate_scope(self) -> NewsQuery: + """要求 market 或 symbol 至少存在一个,并校验时间窗口。""" + if not self.market and not self.symbol: + raise ValueError("market or symbol is required") + if self.since and self.as_of and self.since > self.as_of: + raise ValueError("since must not be later than as_of") + return self + + +class NewsItem(BaseModel): + """标准化新闻或披露事件。""" + + title: str + publisher: str = "" + link: str = "" + published_at: datetime | None = None + summary: str = "" + kind: NewsKind = "media" + source_id: str = "" + source_name: str = "" + source_tier: SourceTier = "aggregator" + fetched_at: datetime | None = None + accepted_at: datetime | None = None + market: str | None = None + symbols: list[str] = Field(default_factory=list) + language: str | None = None + alternative_sources: list[str] = Field(default_factory=list) + + +class NewsProviderStatus(BaseModel): + """单个 provider 的可观察结果。""" + + provider: str + status: ProviderStatusCode + error: str | None = None + fetched_at: datetime + item_count: int = 0 + + +class NewsResponse(BaseModel): + """统一新闻响应,保留旧 venue/symbol 字段。""" + + venue: str | None = None + market: str | None = None + symbol: str | None = None + as_of: datetime | None = None + since: datetime | None = None + fetched_at: datetime + items: list[NewsItem] = Field(default_factory=list) + providers: list[NewsProviderStatus] = Field(default_factory=list) + is_partial: bool = False diff --git a/services/data/src/inalpha_data/schemas.py b/services/data/src/inalpha_data/schemas.py index 6105c2ea..4e747b35 100644 --- a/services/data/src/inalpha_data/schemas.py +++ b/services/data/src/inalpha_data/schemas.py @@ -82,44 +82,6 @@ class HealthResponse(BaseModel): # ──────────────────────────────────────────────────────────────────── -# ──────────────────────────────────────────────────────────────────── -# News(D-9 加:给 research macro/sentiment analyst 喂真新闻) -# ──────────────────────────────────────────────────────────────────── - - -class NewsQuery(BaseModel): - """``GET /news`` 的 query 参数。""" - - venue: str = Field( - default="yfinance", - description="新闻数据源 venue。支持 yfinance(全球零 key)和 baostock(A股)。", - ) - symbol: str = Field( - ..., - examples=["AAPL", "^GSPC", "005930.KS", "sh.600519"], - description="ticker 标识:yfinance 用 Yahoo ticker,baostock 用 sh./sz. 前缀。", - ) - limit: int = Field(default=10, ge=1, le=30, description="最多返回多少条") - - -class NewsItem(BaseModel): - """单条新闻头条。""" - - title: str - publisher: str = "" - link: str = "" - published_at: datetime | None = None - summary: str = "" - - -class NewsResponse(BaseModel): - """``GET /news`` 响应:按发布时间倒序(最新在 items[0])。""" - - venue: str - symbol: str - items: list[NewsItem] - - class TickerQuery(BaseModel): """``GET /ticker`` 的 query 参数。""" diff --git a/services/data/tests/test_market.py b/services/data/tests/test_market.py index c37f01af..5dbdcbc1 100644 --- a/services/data/tests/test_market.py +++ b/services/data/tests/test_market.py @@ -31,15 +31,119 @@ def test_market_sectors_requires_auth(client: TestClient) -> None: assert r.status_code == 401 -def test_market_unsupported_market_rejected( +def test_market_us_news_uses_market_proxy( client: TestClient, auth_headers: dict[str, str] ) -> None: - """未实装的 market 返 400 MARKET_NOT_SUPPORTED(不要硬调后静默空)。""" - r = client.get("/market/news", headers=auth_headers, params={"market": "us"}) - assert r.status_code == 400 + """无 symbol 美股快讯应通过 SPY 市场代理,而不是 ticker provider unsupported。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + seen: list[str] = [] + + async def mock_news(symbol: str, limit: int = 20) -> list[dict[str, Any]]: + seen.append(symbol) + return [{ + "title": "US market update", + "publisher": "Reuters", + "link": "https://example.com/us", + "published_at": "2026-07-29T05:00:00Z", + "summary": "Stocks moved after macro news.", + }] + + yf._connector.fetch_news = mock_news + try: + r = client.get("/market/news", headers=auth_headers, params={"market": "us"}) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + body = r.json() + assert seen == ["SPY"] + assert body["providers"][0]["provider"] == "yfinance_market_proxy" + assert body["providers"][0]["status"] == "ok" + assert body["items"][0]["kind"] == "market_news" + assert body["items"][0]["symbols"] == ["SPY"] + assert "not a complete market newswire" in body["items"][0]["summary"] + + + +def test_market_hk_news_uses_hsi_proxy( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """无 symbol 港股快讯使用恒指代理,并保留代理 ticker。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + seen: list[str] = [] + + async def mock_news(symbol: str, limit: int = 20) -> list[dict[str, Any]]: + seen.append(symbol) + return [{"title": "HK market", "published_at": "2026-07-29T05:00:00Z"}] + + yf._connector.fetch_news = mock_news + try: + r = client.get("/market/news", headers=auth_headers, params={"market": "hk"}) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + assert seen == ["^HSI"] + assert r.json()["items"][0]["symbols"] == ["^HSI"] + + + +@pytest.mark.parametrize( + ("market", "ticker"), + [ + ("jp", "^N225"), ("kr", "^KS11"), ("au", "^AXJO"), ("in", "^NSEI"), + ("uk", "^FTSE"), ("de", "^GDAXI"), ("fr", "^FCHI"), + ("ca", "^GSPTSE"), ("br", "^BVSP"), ("global", "ACWI"), + ], +) +def test_global_market_news_proxy_routes_expected_ticker( + client: TestClient, auth_headers: dict[str, str], market: str, ticker: str +) -> None: + """已声明的全球股票市场都应映射到可审计的代表性载体。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + seen: list[str] = [] + + async def mock_news(symbol: str, limit: int = 20) -> list[dict[str, Any]]: + seen.append(symbol) + return [{"title": f"{market} market", "published_at": "2026-07-29T05:00:00Z"}] + + yf._connector.fetch_news = mock_news + try: + r = client.get("/market/news", headers=auth_headers, params={"market": market}) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + assert seen == [ticker] + assert r.json()["items"][0]["symbols"] == [ticker] + + + +def test_market_proxy_preserves_yahoo_failure( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """Yahoo 故障不能被包装成 no_results 或假成功。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + + async def failed_news(symbol: str, limit: int = 20) -> list[dict[str, Any]]: + raise RuntimeError("Yahoo rate limited") + + yf._connector.fetch_news = failed_news + try: + r = client.get("/market/news", headers=auth_headers, params={"market": "us"}) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 body = r.json() - assert body["code"] == "MARKET_NOT_SUPPORTED" - assert body["details"]["supported"] == ["cn"] + assert body["items"] == [] + assert body["is_partial"] is True + assert body["providers"][0]["status"] == "upstream_error" + assert "rate limited" in body["providers"][0]["error"] def test_market_news_mock_returns_items( diff --git a/services/data/tests/test_news.py b/services/data/tests/test_news.py index a87a8a98..9b8b9108 100644 --- a/services/data/tests/test_news.py +++ b/services/data/tests/test_news.py @@ -2,9 +2,14 @@ from __future__ import annotations +import asyncio + import pytest from fastapi.testclient import TestClient +from inalpha_data.connectors.news.legacy import CnNewsProvider +from inalpha_data.news_models import NewsQuery + pytestmark = pytest.mark.anyio @@ -105,7 +110,7 @@ async def mock_news(symbol, limit=20): def test_news_unsupported_venue(client: TestClient, auth_headers: dict[str, str]) -> None: - """Venue=binance should return 422.""" + """Venue=binance should return 400.""" r = client.get( "/news", headers=auth_headers, @@ -113,3 +118,170 @@ def test_news_unsupported_venue(client: TestClient, auth_headers: dict[str, str] ) assert r.status_code == 400 assert "NEWS" in r.json()["code"] + + + +def test_news_normalizes_legacy_venue_before_validation( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """旧客户端的 venue 大小写和空白继续兼容,响应仍回显原请求。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + + async def mock_news(symbol: str, limit: int = 20): + return [] + + yf._connector.fetch_news = mock_news + try: + r = client.get( + "/news", + headers=auth_headers, + params={"venue": " YFINANCE ", "symbol": " AAPL "}, + ) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + assert r.json()["venue"] == " YFINANCE " + assert r.json()["symbol"] == " AAPL " + + +def test_news_naive_as_of_is_assumed_utc( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """无 offset 的查询时间按 UTC 解释,避免 PIT 比较触发 500。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + + async def mock_news(symbol: str, limit: int = 20): + return [{"title": "Before cutoff", "published_at": "2026-07-29T11:00:00Z"}] + + yf._connector.fetch_news = mock_news + try: + r = client.get( + "/news", + headers=auth_headers, + params={ + "market": "us", + "symbol": "AAPL", + "kinds": "media", + "as_of": "2026-07-29T12:00:00", + }, + ) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + assert r.json()["as_of"] == "2026-07-29T12:00:00Z" + assert len(r.json()["items"]) == 1 + + +def test_market_only_news_does_not_claim_yfinance_venue( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """市场级 scope 与实际 provider 分离,响应不回填错误 venue。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + + async def mock_news(symbol: str, limit: int = 20): + return [] + + yf._connector.fetch_news = mock_news + try: + r = client.get( + "/news", + headers=auth_headers, + params={"market": "us", "limit": 5}, + ) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + assert r.json()["venue"] is None + + +def test_crypto_symbol_news_reports_feed_unsupported() -> None: + """市场级 RSS 不得冒充单标的新闻流。""" + from inalpha_data.connectors.news.feed_models import DEFAULT_CRYPTO_FEEDS + from inalpha_data.connectors.news.rss import RssFeedProvider + from inalpha_data.news_models import NewsQuery + + provider = RssFeedProvider(DEFAULT_CRYPTO_FEEDS[0], timeout_s=1) + + async def run(): + try: + return await provider.fetch( + NewsQuery(market="crypto", symbol="BTC/USDT", limit=5) + ) + finally: + await provider.close() + + result = asyncio.run(run()) + assert result.items == [] + assert result.status == "unsupported" + + +def test_global_stock_news_uses_symbol_not_market_proxy( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """全球单股新闻必须查标的 ticker,不能被市场代理覆盖。""" + from inalpha_data.connectors import yfinance_conn as yf + + original = yf._connector.fetch_news + seen: list[str] = [] + + async def mock_news(symbol: str, limit: int = 20) -> list[dict[str, object]]: + seen.append(symbol) + return [{"title": "Sony news", "published_at": "2026-07-29T05:00:00Z"}] + + yf._connector.fetch_news = mock_news + try: + r = client.get( + "/news", + headers=auth_headers, + params={"market": "jp", "symbol": "6758.T", "kinds": "media"}, + ) + finally: + yf._connector.fetch_news = original + assert r.status_code == 200 + assert seen == ["6758.T"] + assert r.json()["providers"][0]["provider"] == "yfinance" + + +async def test_cn_provider_preserves_failure_and_provenance( + monkeypatch: pytest.MonkeyPatch, +) -> None: + """东财故障不能伪装空结果,成功条目也不能误标 Yahoo 或官方来源。""" + from inalpha_data.connectors.news import legacy + + class FakeConnector: + async def fetch_news(self, symbol: str, limit: int = 20): + raise RuntimeError("eastmoney unavailable") + + monkeypatch.setattr(legacy, "get_connector_for_venue", lambda venue: FakeConnector()) + provider = CnNewsProvider() + failed = await provider.fetch(NewsQuery(market="cn", symbol="sh.600519")) + assert failed.status == "upstream_error" + + async def mock_news(self, symbol: str, limit: int = 20): + return [{"title": "公告摘要", "published_at": "2026-07-29T08:00:00Z"}] + + monkeypatch.setattr(FakeConnector, "fetch_news", mock_news) + succeeded = await provider.fetch(NewsQuery(market="cn", symbol="sh.600519")) + assert succeeded.items[0].source_name == "eastmoney" + assert succeeded.items[0].source_tier == "professional_media" + + +def test_news_accepts_comma_separated_kinds( + client: TestClient, auth_headers: dict[str, str] +) -> None: + """TS client 的逗号分隔 kinds 应在 FastAPI list 包装后继续展开。""" + r = client.get( + "/news", + headers=auth_headers, + params={"market": "us", "kinds": "market_news,media", "limit": 5}, + ) + assert r.status_code == 200 + body = r.json() + assert body["market"] == "us" + assert body["providers"][0]["provider"] == "yfinance_market_proxy" diff --git a/services/data/tests/test_news_providers.py b/services/data/tests/test_news_providers.py new file mode 100644 index 00000000..2dfffd1a --- /dev/null +++ b/services/data/tests/test_news_providers.py @@ -0,0 +1,94 @@ +"""统一财经新闻层测试。""" +from __future__ import annotations + +from datetime import UTC, datetime + +import httpx +import pytest + +from inalpha_data.connectors.news.dedupe import filter_and_dedupe +from inalpha_data.connectors.news.feed_models import FeedDefinition +from inalpha_data.connectors.news.hkex_parser import parse_rows +from inalpha_data.connectors.news.rss import RssFeedProvider +from inalpha_data.connectors.news.sec_parser import parse_submissions +from inalpha_data.news_models import NewsItem, NewsQuery + +pytestmark = pytest.mark.anyio + + +def test_sec_submissions_are_official_disclosures() -> None: + query = NewsQuery(market="us", symbol="AAPL", as_of="2026-07-01T00:00:00Z") + payload = { + "name": "Apple Inc.", + "filings": {"recent": { + "accessionNumber": ["0000320193-26-000001", "0000320193-26-000002"], + "acceptanceDateTime": ["2026-06-30T18:01:02Z", "2026-07-02T18:01:02Z"], + "filingDate": ["2026-06-30", "2026-07-02"], + "form": ["8-K", "10-Q"], + "primaryDocument": ["aapl-8k.htm", "aapl-10q.htm"], + }}, + } + items = filter_and_dedupe( + parse_submissions(payload, query, datetime.now(UTC), "0000320193"), query + ) + assert len(items) == 1 + assert items[0].kind == "disclosure" + assert items[0].source_tier == "official" + assert "/320193/000032019326000001/aapl-8k.htm" in items[0].link + + +def test_hkex_rows_dedupe_languages_and_convert_timezone() -> None: + query = NewsQuery(market="hk", symbol="0700.HK") + rows = [ + {"NEWS_ID": "1", "TITLE": "Annual Results", "DATE_TIME": "28/07/2026 18:30", + "FILE_LINK": "/listedco/listconews/sehk/2026/1.pdf", "_language": "en-HK"}, + {"NEWS_ID": "1", "TITLE": "全年業績", "DATE_TIME": "28/07/2026 18:30", + "FILE_LINK": "/listedco/listconews/sehk/2026/1c.pdf", "_language": "zh-HK"}, + ] + items = parse_rows(rows, query, datetime.now(UTC)) + assert len(items) == 1 + assert items[0].published_at == datetime(2026, 7, 28, 10, 30, tzinfo=UTC) + + +async def test_rss_provider_reuses_cached_items_on_304() -> None: + feed = FeedDefinition("test", "Test Feed", "https://feed.test/rss", "professional_media", "en") + provider = RssFeedProvider(feed, timeout_s=1) + calls = 0 + + async def handler(request: httpx.Request) -> httpx.Response: + nonlocal calls + calls += 1 + if calls == 1: + return httpx.Response( + 200, + headers={"etag": '"v1"'}, + content=b"xT" + b"https://x.testTue, 28 Jul 2026 10:00:00 GMT" + b"", + ) + assert request.headers["if-none-match"] == '"v1"' + return httpx.Response(304) + + await provider._client.aclose() + provider._client = httpx.AsyncClient(transport=httpx.MockTransport(handler)) + try: + first = await provider.fetch(NewsQuery(market="crypto", limit=5)) + second = await provider.fetch(NewsQuery(market="crypto", limit=5)) + finally: + await provider.close() + assert first.status == second.status == "ok" + assert first.items[0].source_name == "test" + assert second.items[0].title == "T" + + +def test_dedupe_prefers_official_source() -> None: + query = NewsQuery(market="us", symbol="AAPL") + ts = datetime(2026, 7, 28, tzinfo=UTC) + media = NewsItem(title="Event", link="https://x.test/a?utm_source=z", published_at=ts, + source_name="wire", source_tier="professional_media") + official = NewsItem(title="Event", link="https://x.test/a", published_at=ts, + source_name="sec", source_tier="official", kind="disclosure") + result = filter_and_dedupe([media, official], query) + assert len(result) == 1 + assert result[0].source_name == "sec" + assert "wire" in result[0].alternative_sources diff --git a/services/data/uv.lock b/services/data/uv.lock index 1e93f3b4..ecb6be0b 100644 --- a/services/data/uv.lock +++ b/services/data/uv.lock @@ -828,6 +828,18 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/5a/ff/2e4eca3ade2c22fe1dea7043b8ee9dabe47753349eb1b56a202de8af6349/fastapi-0.136.1-py3-none-any.whl", hash = "sha256:a6e9d7eeada96c93a4d69cb03836b44fa34e2854accb7244a1ece36cd4781c3f", size = 117683, upload-time = "2026-04-23T16:49:42.437Z" }, ] +[[package]] +name = "feedparser" +version = "6.0.12" +source = { registry = "https://pypi.org/simple" } +dependencies = [ + { name = "sgmllib3k" }, +] +sdist = { url = "https://files.pythonhosted.org/packages/dc/79/db7edb5e77d6dfbc54d7d9df72828be4318275b2e580549ff45a962f6461/feedparser-6.0.12.tar.gz", hash = "sha256:64f76ce90ae3e8ef5d1ede0f8d3b50ce26bcce71dd8ae5e82b1cd2d4a5f94228", size = 286579, upload-time = "2025-09-10T13:33:59.486Z" } +wheels = [ + { url = "https://files.pythonhosted.org/packages/4e/eb/c96d64137e29ae17d83ad2552470bafe3a7a915e85434d9942077d7fd011/feedparser-6.0.12-py3-none-any.whl", hash = "sha256:6bbff10f5a52662c00a2e3f86a38928c37c48f77b3c511aedcd51de933549324", size = 81480, upload-time = "2025-09-10T13:33:58.022Z" }, +] + [[package]] name = "fredapi" version = "0.5.2" @@ -1096,6 +1108,7 @@ dependencies = [ { name = "ccxt" }, { name = "curl-cffi" }, { name = "ddgs" }, + { name = "feedparser" }, { name = "fredapi" }, { name = "inalpha-shared" }, { name = "trafilatura" }, @@ -1121,6 +1134,7 @@ requires-dist = [ { name = "ccxt", specifier = ">=4.5.0" }, { name = "curl-cffi", specifier = ">=0.7" }, { name = "ddgs", specifier = ">=9.0.0" }, + { name = "feedparser", specifier = ">=6.0.12" }, { name = "fredapi", specifier = ">=0.5.2" }, { name = "inalpha-shared", editable = "../_shared" }, { name = "trafilatura", specifier = ">=2.1.0" }, @@ -2453,6 +2467,12 @@ wheels = [ { url = "https://files.pythonhosted.org/packages/9d/76/f789f7a86709c6b087c5a2f52f911838cad707cc613162401badc665acfe/setuptools-82.0.1-py3-none-any.whl", hash = "sha256:a59e362652f08dcd477c78bb6e7bd9d80a7995bc73ce773050228a348ce2e5bb", size = 1006223, upload-time = "2026-03-09T12:47:15.026Z" }, ] +[[package]] +name = "sgmllib3k" +version = "1.0.0" +source = { registry = "https://pypi.org/simple" } +sdist = { url = "https://files.pythonhosted.org/packages/9e/bd/3704a8c3e0942d711c1299ebf7b9091930adae6675d7c8f476a7ce48653c/sgmllib3k-1.0.0.tar.gz", hash = "sha256:7868fb1c8bfa764c1ac563d3cf369c381d1325d36124933a726f29fcdaa812e9", size = 5750, upload-time = "2010-08-24T14:33:52.445Z" } + [[package]] name = "six" version = "1.17.0" From 65567438af0094657f5928f314bf1368ecf07a5f Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 30 Jul 2026 11:12:09 +0800 Subject: [PATCH 2/4] =?UTF-8?q?feat(orchestration):=20=E6=9A=B4=E9=9C=B2?= =?UTF-8?q?=E5=85=A8=E7=90=83=E5=B8=82=E5=9C=BA=E6=96=B0=E9=97=BB=E5=B7=A5?= =?UTF-8?q?=E5=85=B7?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- packages/orchestration/src/clients/data.ts | 55 +++++++++++++++++-- .../src/mastra/agents/instructions/market.ts | 7 +++ .../agents/instructions/tool-catalog.ts | 14 +++-- packages/orchestration/src/tools/data.ts | 43 +++++++++++++++ packages/orchestration/src/tools/index.ts | 3 + packages/orchestration/src/tools/market.ts | 32 +++++++---- .../orchestration/tests/market-tools.test.ts | 49 +++++++++++++++++ 7 files changed, 181 insertions(+), 22 deletions(-) diff --git a/packages/orchestration/src/clients/data.ts b/packages/orchestration/src/clients/data.ts index 4ab4ca34..033dccfc 100644 --- a/packages/orchestration/src/clients/data.ts +++ b/packages/orchestration/src/clients/data.ts @@ -35,6 +35,10 @@ export type Ticker = { stale_seconds: number; }; +export type NewsMarket = + | "cn" | "us" | "hk" | "jp" | "kr" | "au" | "in" | "uk" + | "de" | "fr" | "ca" | "br" | "global" | "crypto"; + export class DataClient { private readonly http: HttpClient; @@ -111,20 +115,61 @@ export class DataClient { } } + async getNews(params: { + market?: NewsMarket; + venue?: string; + symbol?: string; + asOf?: string; + since?: string; + kinds?: Array<"market_news" | "media" | "disclosure">; + language?: string; + limit?: number; + }): Promise> { + try { + return await this.http.get>("/news", { + ...(params.market ? { market: params.market } : {}), + ...(params.venue ? { venue: params.venue } : {}), + ...(params.symbol ? { symbol: params.symbol } : {}), + ...(params.asOf ? { as_of: params.asOf } : {}), + ...(params.since ? { since: params.since } : {}), + ...(params.kinds?.length ? { kinds: params.kinds.join(",") } : {}), + ...(params.language ? { language: params.language } : {}), + limit: String(params.limit ?? 10), + }); + } catch (err) { + return { + market: params.market, + symbol: params.symbol, + items: [], + providers: [], + is_partial: true, + error: err instanceof HttpClientError ? `upstream ${err.status}: ${err.message}` : String(err), + }; + } + } + async getMarketNews(params: { market?: string; limit?: number; }): Promise> { + if ((params.market ?? "cn") !== "cn") { + return await this.getNews({ + market: params.market as NewsMarket, + limit: params.limit ?? 20, + kinds: ["market_news", "media"], + }); + } try { return await this.http.get>("/market/news", { - market: params.market ?? "cn", + market: "cn", limit: String(params.limit ?? 20), }); } catch (err) { - if (err instanceof HttpClientError) { - return { market: params.market ?? "cn", items: [], error: `upstream ${err.status}: ${err.message}` }; - } - return { market: params.market ?? "cn", items: [], error: String(err) }; + return { + market: "cn", + items: [], + error: err instanceof HttpClientError ? `upstream ${err.status}: ${err.message}` : String(err), + }; } } diff --git a/packages/orchestration/src/mastra/agents/instructions/market.ts b/packages/orchestration/src/mastra/agents/instructions/market.ts index d1c859ac..fd389ba8 100644 --- a/packages/orchestration/src/mastra/agents/instructions/market.ts +++ b/packages/orchestration/src/mastra/agents/instructions/market.ts @@ -110,6 +110,13 @@ Inalpha 是金融 agent —— "数据 stale 几天" 等于"建议过时"。任 市场级工具传对应 market;该市场没有市场级工具时,用 "web.search_news + 该市场代表性指数的 get_bars"组合替代维度 1-4。 +**市场快讯 market code**(供 data.get_market_news 使用): +- A股 \`cn\`、美股 \`us\`、港股 \`hk\`、日股 \`jp\`、韩股 \`kr\`、澳股 \`au\` +- 印股 \`in\`、英股 \`uk\`、德股 \`de\`、法股 \`fr\`、加拿大 \`ca\`、巴西 \`br\` +- 全球市场 \`global\`、加密市场 \`crypto\` +- 除 A股专业快讯和 Crypto 专业 feed 外,股票市场返回代表性指数/ETF 新闻代理; + 必须根据响应 \`symbols\` 标注代理口径,不得称为完整市场新闻线 + **结论纪律**: - 每个维度的结论必须指得回工具返回的数据;某维度拿不到数据 → **显式声明该维度缺失 并跳过**,继续完成其余维度——不要因为一个维度空就放弃整个归因、只讲技术面 diff --git a/packages/orchestration/src/mastra/agents/instructions/tool-catalog.ts b/packages/orchestration/src/mastra/agents/instructions/tool-catalog.ts index e7ecb992..79a80150 100644 --- a/packages/orchestration/src/mastra/agents/instructions/tool-catalog.ts +++ b/packages/orchestration/src/mastra/agents/instructions/tool-catalog.ts @@ -23,6 +23,9 @@ export const TOOL_CATALOG = ` 其他市场返 yahoo 格式(venue 字段标明配哪个数据源)。 **你已判断出市场分类时显式传 venue**(美股/港股/全球 → yfinance,A股 → baostock); query 的语言 ≠ 市场——中文名问美股公司极常见,别把市场判断丢给 auto 兜底。 +- **data.get_news —— 多市场单标的新闻与官方披露**。A股走东财、美股 SEC+Yahoo、 + 港股 HKEX+Yahoo、Crypto 走已注册财经 feed;需要历史截止点时传 asOf。响应的 + providers/status 与 is_partial 区分“无结果”和“源站故障”,后者不能解释为没有消息。 **Web 搜索**(D-10 新 · 零 key,ddgs 聚合多引擎): - web.search —— 搜索互联网。query 用自然语言;backend 默认 auto,中文自动走 bing。 @@ -44,9 +47,10 @@ export const TOOL_CATALOG = ` 对 A股用 venue=baostock(支持 1d/1wk/1mo + 分钟级 5m/15m/30m/1h),美股/港股用 venue=yfinance **市场级行情(D-12+ 新 · 行情归因专用,无需 symbol)**: -- data.get_market_news —— 市场级财经快讯流。用户问"某市场 / 大盘今天有什么消息 / - 为什么涨跌"时**优先于 web.search_news**(专业财经快讯源,免搜索引擎噪声)。 - 不用于单标的新闻深挖(标的级仍走 web.search_news + web.fetch) +- data.get_market_news —— 多市场财经快讯流。用户问“某市场 / 大盘今天有什么消息 / + 为什么涨跌”时优先于 web.search_news;支持全球市场分类。A股用东财、Crypto 用 + 专业 feed,其余股票市场使用代表性指数/ETF 新闻代理,引用时必须带代理口径。 + 单标的消息或官方披露改用 data.get_news。 - data.get_market_sectors —— 行业板块涨跌幅榜(涨跌两端 + 领涨股)。 判断"普涨还是结构性、哪些板块领涨领跌";归因个股时先看它所属板块在榜单的位置 - data.get_market_moneyflow —— 跨境资金流(A股=沪深港通)。资金面维度。 @@ -55,8 +59,8 @@ export const TOOL_CATALOG = ` - data.get_market_movers —— 当日强势股 + 人工题材标签。归因"什么主线在涨"的 最直接证据(对 tags 聚类看热点)。坑:标签是媒体归纳**非因果实锤**, 措辞用"市场归因于 / 题材标签显示" -- 四个工具按 market 参数路由(同"全球市场覆盖"分类);当前仅实装 cn(A股)。 - **未实装的市场不要硬调**(会返 400),降级走 web.search_news + 该市场代表性指数 get_bars +- sectors / moneyflow / movers 当前仅实装 cn(A股);其它市场只可调用 get_market_news, + 其余盘面维度降级走 web.search_news + 代表性指数 get_bars。 **有效因子择时(接现成因子库 pandas-ta / Alpha101 / qlib)**: - factor.timing —— 给一个标的/周期,返回**当前最有效的因子**(按时序 Rank IC 排序)+ 读数 + 方向 + 强度。 diff --git a/packages/orchestration/src/tools/data.ts b/packages/orchestration/src/tools/data.ts index 09f2dc67..637485cb 100644 --- a/packages/orchestration/src/tools/data.ts +++ b/packages/orchestration/src/tools/data.ts @@ -397,10 +397,53 @@ export const dataSearchSymbolTool = createTool({ }, }); +// ──────────────────────────────────────────────────────────────────── +// data.get_news +// ──────────────────────────────────────────────────────────────────── + +export const dataGetNewsTool = createTool({ + id: "data.get_news", + description: ` + 统一财经新闻与官方披露。按 market 路由:A股东财、美股 SEC+Yahoo、港股 + HKEX+Yahoo、Crypto 专业媒体/官方 feed;返回逐来源 status、内容时点和部分失败标志。 + + 何时用:单标的新闻、公司公告、市场消息面,以及需要 asOf 截止时间的研究证据。 + 何时不用:价格/K线用 data.get_bars;要读全文时对返回 URL 再用 web.fetch。 + 坑:disclosure 不等于媒体新闻;is_partial=true 表示至少一个专业源故障,不能把 + 空结果解释成“没有消息”。历史查询会排除无可靠发布时间和晚于 asOf 的条目。 + `.trim(), + inputSchema: z.object({ + market: z.enum([ + "cn", "us", "hk", "jp", "kr", "au", "in", "uk", + "de", "fr", "ca", "br", "global", "crypto", + ]), + symbol: SymbolSchema.optional(), + asOf: z.string().datetime({ offset: true }).optional(), + since: z.string().datetime({ offset: true }).optional(), + kinds: z.array(z.enum(["market_news", "media", "disclosure"])).optional(), + language: z.string().min(1).max(35).optional(), + limit: z.number().int().min(1).max(50).default(10), + }), + execute: async (inputData, ctx) => { + const tc = ctx?.requestContext as ToolRequestContext | undefined; + const client = await getClient(tc); + return await client.getNews({ + market: inputData.market, + symbol: inputData.symbol, + asOf: inputData.asOf, + since: inputData.since, + kinds: inputData.kinds, + language: inputData.language, + limit: inputData.limit ?? 10, + }); + }, +}); + export const dataTools = [ dataGetBarsTool, dataBackfillBarsTool, dataGetTickerTool, dataGetFundamentalsTool, dataSearchSymbolTool, + dataGetNewsTool, ] as const; diff --git a/packages/orchestration/src/tools/index.ts b/packages/orchestration/src/tools/index.ts index 7f90147b..e0e2c380 100644 --- a/packages/orchestration/src/tools/index.ts +++ b/packages/orchestration/src/tools/index.ts @@ -12,6 +12,7 @@ import { dataBackfillBarsTool, dataGetBarsTool, dataGetFundamentalsTool, + dataGetNewsTool, dataGetTickerTool, dataSearchSymbolTool, dataTools, @@ -115,6 +116,7 @@ export { dataGetMarketMoversTool, dataGetMarketNewsTool, dataGetMarketSectorsTool, + dataGetNewsTool, dataGetTickerTool, dataSearchSymbolTool, divinationCastHexagramTool, @@ -244,6 +246,7 @@ export const orchestratorToolList = [ dataBackfillBarsTool, dataGetTickerTool, dataGetFundamentalsTool, + dataGetNewsTool, // 公司名 → ticker 解析(候选池构建,禁训练记忆猜代码) dataSearchSymbolTool, // D-10:web 搜索;web.fetch 读原文补证据链最后一公里 diff --git a/packages/orchestration/src/tools/market.ts b/packages/orchestration/src/tools/market.ts index e5fa9f3c..d5268534 100644 --- a/packages/orchestration/src/tools/market.ts +++ b/packages/orchestration/src/tools/market.ts @@ -2,8 +2,8 @@ * 市场级行情归因 tool(D-12+)—— services/data 的 /market/* 端点包装。 * * 行情归因("今天为什么涨/跌")的四个数据维度,全部无需 symbol。 - * venue 按 market 参数路由:当前实装 cn(A股,直连东财/同花顺,配方源自 - * a-stock-data);未实装的市场后端返 400,不要硬调。 + * venue 按 market 参数路由:新闻支持全部已声明市场;板块、资金与强势股当前仅支持 + * cn(A股,直连东财/同花顺,配方源自 a-stock-data)。 * * Tool 设计遵循 docs/05-tool-skill-discipline.md:description 四要素。 */ @@ -22,8 +22,14 @@ async function getClient(ctx?: ToolRequestContext): Promise { return new DataClient({ baseUrl: settings.dataServiceUrl, token }); } -const MarketSchema = z.enum(["cn"]).default("cn") - .describe("市场;当前仅实装 cn(A股),其它市场归因用 web.search_news + 指数 get_bars 替代"); +const NewsMarketSchema = z.enum([ + "cn", "us", "hk", "jp", "kr", "au", "in", "uk", + "de", "fr", "ca", "br", "global", "crypto", +]).default("cn") + .describe("市场;A股走东财,Crypto 走财经 feed,其余股票市场走明确标注的代表指数/ETF 新闻代理"); + +const CnMarketSchema = z.enum(["cn"]).default("cn") + .describe("市场;当前仅支持 cn(A股)"); // ──────────────────────────────────────────────────────────────────── // data.get_market_news @@ -32,8 +38,9 @@ const MarketSchema = z.enum(["cn"]).default("cn") export const dataGetMarketNewsTool = createTool({ id: "data.get_market_news", description: ` - 全市场财经快讯流(A股=东财 7×24 全球资讯),**无需 symbol**。返回标题 / 摘要 / - UTC 时间戳 / 关联代码。 + 多市场财经快讯流,**无需 symbol**。A股走东财,Crypto 走专业财经 feed;其他 + 股票市场走代表性指数/ETF 的 Yahoo 新闻代理,条目会保留代理 ticker 与口径声明。 + 响应含逐 provider status、UTC 时间戳与部分失败标志。 何时用: - 行情归因:"某市场 / 大盘今天为什么涨跌、有什么消息"——**优先于 web.search_news** @@ -47,11 +54,12 @@ export const dataGetMarketNewsTool = createTool({ 坑: - published_at 已转 UTC;引用时按 §3.1 标注数据时点 - 60s 进程内缓存;快讯≠结论,结论级引用先 web.fetch 读原文 - - 源站故障返回 error 字段(502 MARKET_DATA_UNAVAILABLE)——此时降级 - web.search 并显式说明快讯源不可用 + - A股源站故障返回 error 字段(502 MARKET_DATA_UNAVAILABLE);其它市场通过 + providers[].status / is_partial 表示部分或全部来源故障。此时降级 web.search_news + 并显式说明结构化快讯源不可用 `.trim(), inputSchema: z.object({ - market: MarketSchema, + market: NewsMarketSchema, limit: z.number().int().min(1).max(50).default(20), }), execute: async (inputData, ctx) => { @@ -88,7 +96,7 @@ export const dataGetMarketSectorsTool = createTool({ - 60s 缓存;不要循环逐板块调用 `.trim(), inputSchema: z.object({ - market: MarketSchema, + market: CnMarketSchema, topN: z.number().int().min(1).max(50).default(10), }), execute: async (inputData, ctx) => { @@ -125,7 +133,7 @@ export const dataGetMarketMoneyflowTool = createTool({ - as_of_time 是北京时间 HH:MM;非交易时段拿的是上一交易日尾盘值 `.trim(), inputSchema: z.object({ - market: MarketSchema, + market: CnMarketSchema, }), execute: async (inputData, ctx) => { const tc = ctx?.requestContext as ToolRequestContext | undefined; @@ -158,7 +166,7 @@ export const dataGetMarketMoversTool = createTool({ - 含 ST / 涨停板个股,标签可能滞后;60s 缓存 `.trim(), inputSchema: z.object({ - market: MarketSchema, + market: CnMarketSchema, limit: z.number().int().min(1).max(50).default(30), }), execute: async (inputData, ctx) => { diff --git a/packages/orchestration/tests/market-tools.test.ts b/packages/orchestration/tests/market-tools.test.ts index 2686fa8b..d6a6b574 100644 --- a/packages/orchestration/tests/market-tools.test.ts +++ b/packages/orchestration/tests/market-tools.test.ts @@ -9,6 +9,7 @@ import { dataGetMarketMoversTool, dataGetMarketNewsTool, dataGetMarketSectorsTool, + dataGetNewsTool, } from "../src/tools/index.js"; const TEST_TOKEN = "test-token-doesnt-need-to-be-real"; @@ -35,6 +36,35 @@ function mockFetch(impl: (url: string, init?: RequestInit) => Promise) const ctx = (authToken: string | undefined = TEST_TOKEN): never => ({ requestContext: { authToken } }) as never; +describe("data.get_news", () => { + it("forwards PIT window, kinds, symbol, and token", async () => { + let capturedUrl = ""; + mockFetch(async (url) => { + capturedUrl = String(url); + return new Response( + JSON.stringify({ market: "us", symbol: "AAPL", items: [], providers: [] }), + { status: 200, headers: { "content-type": "application/json" } }, + ); + }); + await dataGetNewsTool.execute!( + { + market: "us", + symbol: "AAPL", + asOf: "2026-07-28T10:00:00Z", + since: "2026-07-01T00:00:00Z", + kinds: ["disclosure", "media"], + limit: 12, + }, + ctx(), + ); + expect(capturedUrl).toContain("/news?"); + expect(capturedUrl).toContain("market=us"); + expect(capturedUrl).toContain("symbol=AAPL"); + expect(capturedUrl).toContain("as_of=2026-07-28T10%3A00%3A00Z"); + expect(capturedUrl).toContain("kinds=disclosure%2Cmedia"); + }); +}); + describe("data.get_market_news", () => { it("calls /market/news with market+limit and forwards token", async () => { let capturedUrl = ""; @@ -58,6 +88,24 @@ describe("data.get_market_news", () => { expect((out.items as unknown[]).length).toBe(1); }); + it.each(["us", "hk", "jp", "kr", "au", "in", "uk", "de", "fr", "ca", "br", "global", "crypto"] as const)( + "routes %s market news through unified endpoint", + async (market) => { + let capturedUrl = ""; + mockFetch(async (url) => { + capturedUrl = String(url); + return new Response(JSON.stringify({ market, items: [], providers: [] }), { + status: 200, + headers: { "content-type": "application/json" }, + }); + }); + await dataGetMarketNewsTool.execute!({ market, limit: 7 }, ctx()); + expect(capturedUrl).toContain("/news?"); + expect(capturedUrl).toContain(`market=${market}`); + expect(capturedUrl).toContain("kinds=market_news%2Cmedia"); + }, + ); + it("upstream 502 returns error field instead of throwing", async () => { mockFetch(async () => new Response( @@ -133,6 +181,7 @@ describe("data.get_market_movers", () => { describe("tool ids", () => { it("market tool ids follow data. convention", () => { + expect(dataGetNewsTool.id).toBe("data.get_news"); expect(dataGetMarketNewsTool.id).toBe("data.get_market_news"); expect(dataGetMarketSectorsTool.id).toBe("data.get_market_sectors"); expect(dataGetMarketMoneyflowTool.id).toBe("data.get_market_moneyflow"); From 828892ccfc2c64530bd62b51d2722b271e11e103 Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 30 Jul 2026 11:16:29 +0800 Subject: [PATCH 3/4] =?UTF-8?q?feat(research):=20=E6=8E=A5=E5=85=A5?= =?UTF-8?q?=E5=A4=9A=E5=B8=82=E5=9C=BA=E6=96=B0=E9=97=BB=E8=AF=81=E6=8D=AE?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- .../src/inalpha_research/analysts/macro.py | 11 +++- .../inalpha_research/analysts/sentiment.py | 47 +++++++++++++---- .../src/inalpha_research/data_client.py | 52 ++++++++++++------- 3 files changed, 80 insertions(+), 30 deletions(-) diff --git a/services/research/src/inalpha_research/analysts/macro.py b/services/research/src/inalpha_research/analysts/macro.py index 1bf9fa3f..947320ba 100644 --- a/services/research/src/inalpha_research/analysts/macro.py +++ b/services/research/src/inalpha_research/analysts/macro.py @@ -216,10 +216,17 @@ async def build_user_prompt( # D-9 L3:拉 SPY 当宏观 proxy(美国宏观环境主导全球风险偏好)。 # 拉不到时返空 list,prompt 里清晰标注,LLM 走纯 calendar + 训练知识。 # D-12:FRED 读数与新闻并发拉(互不依赖,各自独立降级)。 - macro_news, readings = await asyncio.gather( - self._data.get_news(symbol="SPY", limit=8), + macro_news_result, readings = await asyncio.gather( + self._data.get_news( + market="us", + symbol="SPY", + as_of=as_of, + kinds=["media"], + limit=8, + ), _fetch_macro_readings(self._data, as_of=as_of), ) + macro_news = macro_news_result.get("items", []) # 双档 cap(run() 里代码级 clamp):有任一 live 数据 0.7,全无 0.5 self._confidence_cap = 0.7 if (readings or macro_news) else 0.5 return _format_user_prompt( diff --git a/services/research/src/inalpha_research/analysts/sentiment.py b/services/research/src/inalpha_research/analysts/sentiment.py index d6026ddb..3515b8e6 100644 --- a/services/research/src/inalpha_research/analysts/sentiment.py +++ b/services/research/src/inalpha_research/analysts/sentiment.py @@ -18,7 +18,7 @@ """ from __future__ import annotations -from datetime import datetime +from datetime import datetime, timedelta from typing import Any import httpx @@ -98,9 +98,23 @@ async def build_user_prompt( lookback_days: int, ) -> str: market_type = infer_asset_type(venue=venue, symbol=symbol) - - # crypto → 拉 FNG;非 crypto → 拉 yfinance news 喂 LLM + news_market = { + "crypto": "crypto", + "us_stock": "us", + "cn_stock": "cn", + "hk_stock": "hk", + }.get(market_type) + + # crypto → FNG + 专业新闻;两类证据并列,不互相替代。 if market_type == "crypto": + news_result = await self._data.get_news( + market="crypto", + as_of=as_of, + since=as_of - timedelta(days=lookback_days), + kinds=["media"], + limit=8, + ) + crypto_news = news_result.get("items", []) try: entries = await _fetch_fng(limit=_DEFAULT_LIMIT) except Exception as exc: @@ -121,7 +135,7 @@ async def build_user_prompt( as_of=as_of, market_type=market_type, fng_note="(Fear & Greed API unavailable — using web search)", - news=[], + news=crypto_news, web_results=web_results, ) latest = entries[0] @@ -134,22 +148,32 @@ async def build_user_prompt( latest=latest, recent_values=recent_values, trend=trend, + news=crypto_news, ) - # 非 crypto:拉 yfinance ticker news + web search,真新闻锚定 sentiment - # symbol 不一定能直接给 yfinance(akshare 的 sh.600519 等格式不通); - # 这里直接用原 symbol 试一次;data-service 拉不到会返空 list,自然降级 LLM-only。 - news = await self._data.get_news(symbol=symbol, limit=8) + # 非 crypto:结构化本地新闻 + 通用搜索兜底 + news_result = await self._data.get_news( + market=news_market, + venue=venue, + symbol=symbol, + as_of=as_of, + since=as_of - timedelta(days=lookback_days), + limit=8, + ) + news = news_result.get("items", []) + provider_status = news_result.get("providers", []) web_news = await self._data.get_web_search( f"{symbol} stock news sentiment analysis", max_results=5 ) + partial_note = " structured sources partial" if news_result.get("is_partial") else "" return _format_user_prompt_llm_only( symbol=symbol, as_of=as_of, market_type=market_type, fng_note=( - f"(non-crypto market — no Fear & Greed; {len(news)} news headlines, {len(web_news)} web results)" - if news or web_news + f"(non-crypto market — no Fear & Greed; {len(news)} structured headlines, " + f"{len(web_news)} web results; providers={provider_status}{partial_note})" + if news or web_news or provider_status else "(non-crypto market — no Fear & Greed; all sources returned empty)" ), news=news, @@ -195,7 +219,9 @@ def _format_user_prompt_with_fng( latest: dict[str, Any], recent_values: list[int], trend: dict[str, Any], + news: list[dict[str, Any]] | None = None, ) -> str: + news_block = _render_news_block(news or []) return ( f"asset: {symbol}\n" f"market_type: {market_type}\n" @@ -206,6 +232,7 @@ def _format_user_prompt_with_fng( f" timestamp: {latest.get('timestamp')}\n" f" trend_snapshot: {trend}\n" f" recent_30d_values (newest first): {recent_values}\n\n" + f"{news_block}\n" f"Output the required JSON only." ) diff --git a/services/research/src/inalpha_research/data_client.py b/services/research/src/inalpha_research/data_client.py index 1156a4a7..0f6fdcdf 100644 --- a/services/research/src/inalpha_research/data_client.py +++ b/services/research/src/inalpha_research/data_client.py @@ -153,30 +153,46 @@ async def _best_effort_backfill( async def get_news( self, *, - venue: str = "yfinance", - symbol: str, + venue: str | None = None, + market: str | None = None, + symbol: str | None = None, + as_of: datetime | None = None, + since: datetime | None = None, + kinds: list[str] | None = None, limit: int = 10, - ) -> list[dict[str, Any]]: - """``GET /news`` —— ticker-specific news 头条(按时间倒序)。 - - 失败(venue 不支持 / 网络 / 测试 mock 未注册)时返**空 list**,不抛—— - 让 analyst 兜底走 LLM-only 而不是整条链路 500。 - """ + ) -> dict[str, Any]: + """``GET /news`` —— 保留条目与 provider 故障状态。""" + params: dict[str, Any] = {"limit": limit} + if venue: + params["venue"] = venue + if market: + params["market"] = market + if symbol: + params["symbol"] = symbol + if as_of: + params["as_of"] = as_of.isoformat() + if since: + params["since"] = since.isoformat() + if kinds: + params["kinds"] = kinds try: - r = await self._client.get( - "/news", - params={"venue": venue, "symbol": symbol, "limit": limit}, - ) - except Exception: - return [] + r = await self._client.get("/news", params=params) + except Exception as exc: + return {"items": [], "providers": [], "is_partial": True, "error": str(exc)} if r.status_code >= 400: - return [] + return { + "items": [], + "providers": [], + "is_partial": True, + "error": f"upstream {r.status_code}", + } try: payload = r.json() except Exception: - return [] - items = payload.get("items") if isinstance(payload, dict) else None - return items if isinstance(items, list) else [] + return {"items": [], "providers": [], "is_partial": True, "error": "invalid json"} + return payload if isinstance(payload, dict) else { + "items": [], "providers": [], "is_partial": True, "error": "invalid payload" + } async def get_fundamentals( self, venue: str, symbol: str, as_of: datetime | None = None From 0c1e435f33433f38c0d63e5a673ede7259945360 Mon Sep 17 00:00:00 2001 From: Miro Date: Thu, 30 Jul 2026 11:34:43 +0800 Subject: [PATCH 4/4] =?UTF-8?q?docs(data):=20=E8=AE=B0=E5=BD=95=E5=A4=9A?= =?UTF-8?q?=E5=B8=82=E5=9C=BA=E6=96=B0=E9=97=BB=E8=83=BD=E5=8A=9B?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Co-Authored-By: Claude Opus 4.8 --- README.md | 1 + README.zh-CN.md | 1 + docs/04-current-state.md | 6 ++++++ 3 files changed, 8 insertions(+) diff --git a/README.md b/README.md index 86ecf64f..20d65336 100644 --- a/README.md +++ b/README.md @@ -235,6 +235,7 @@ Where each capability stands today. Live module inventory and the end-to-end dec | ✅ 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 · `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 | Unified multi-market financial news | D-12 | `GET /news` + `data.get_news` · Eastmoney / SEC / HKEX / Crypto feeds · PIT `as_of` filtering · observable provider failures · explicit index/ETF proxy labels for market-level stock news | | ✅ 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` | diff --git a/README.zh-CN.md b/README.zh-CN.md index 78a687b8..c62f4cda 100644 --- a/README.zh-CN.md +++ b/README.zh-CN.md @@ -233,6 +233,7 @@ Inalpha 把*调度*和*算力*分开:agent runtime 负责扇出网格、聚合 | ✅ 已上线 | Scheduler / cron agent 模式 | D-9 | `scheduler_jobs` + advisory lock + `/api/scheduler/*` 管理面 | | ✅ 已上线 | RiskGuard 账户级隔离 | D-9.1a | `RiskGuardFactory` 去除跨账户状态串联 | | ✅ 已上线 | 多市场数据源 — web 搜索 + 财报基本面 | D-10 | DDGS 零 key web search · `baostock` 是 A 股逻辑 venue:腾讯 HTTPS 行情/最新价 + Baostock 基本面/日历/成分;yfinance 覆盖全球(含港股)· analyst 接入 + 兜底 · lookbackDays 按市场分化 | +| ✅ 已上线 | 统一多市场财经新闻 | D-12 | `GET /news` + `data.get_news` · 东财 / SEC / HKEX / Crypto feed · PIT `as_of` 过滤 · provider 故障可观察 · 股票市场级新闻明确标注指数/ETF 代理口径 | | ✅ 已上线 | 风控引擎 — 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` 硬编码 | diff --git a/docs/04-current-state.md b/docs/04-current-state.md index 09a4ec4b..5c0dbd11 100644 --- a/docs/04-current-state.md +++ b/docs/04-current-state.md @@ -199,6 +199,12 @@ D-9 把决策护栏(Plan/Exec + 风控 + 沙盒)做扎实后,D-10 在**数 `lookbackDays`(crypto 1h/4h ≈30d;A股/港股/日股 akshare 1d ≈180d;美股/全球指数 yfinance 1d ≈90d),避免"30 天 K 线<20 根交易日"无统计意义;research client 超时放宽到 300s 给 A股慢路径留余量。 +- **统一多市场新闻**:`GET /news` 与 `data.get_news` 统一 `market_news` / `media` / + `disclosure` 契约,首批接入东财、SEC、HKEX 与 Crypto 专业 feed;股票市场级快讯 + 用代表指数/ETF 的 Yahoo 新闻代理并在条目中保留代理 ticker,不冒充完整新闻线。 + 请求支持 `as_of` / `since` 的 PIT 过滤,响应保留逐 provider 状态、`is_partial` 与 + 来源等级;旧 `venue + symbol` 请求继续兼容。`SentimentAnalyst` 与 `MacroAnalyst` + 已消费结构化新闻,来源不可用时再显式降级 web 搜索。 - **金融时效**:web/news 与基本面均无 API key 依赖;analyst 数据源不可用时降级到 LLM-only 并标低 confidence,不静默用过时数据。