From f094cbeae19dab454f693518f40114746d56519e Mon Sep 17 00:00:00 2001 From: Windows Date: Sat, 12 Sep 2026 12:10:10 +0800 Subject: [PATCH] fix: harden advisor market history dual-source sync --- app/service/product_history_sync_service.py | 19 +- docs/21-投顾Agent迁移TODO.md | 132 ++++++++--- .../advisor-e2e-acceptance-20260912.md | 76 +++++++ hq.py | 207 ++++++++++++++++-- .../test_fund_market_adapter.py | 86 ++++++++ tools/acceptance_check.py | 5 +- tools/sync_advisor_market_data.py | 27 ++- 7 files changed, 490 insertions(+), 62 deletions(-) create mode 100644 docs/验收与审计/advisor-e2e-acceptance-20260912.md diff --git a/app/service/product_history_sync_service.py b/app/service/product_history_sync_service.py index aae2dc4..8d04ae4 100644 --- a/app/service/product_history_sync_service.py +++ b/app/service/product_history_sync_service.py @@ -21,6 +21,8 @@ NavHistoryLoader = Callable[[str, str, str], list[dict[str, str]]] class ProductHistorySyncResult: product_count: int observation_count: int + failed_product_count: int = 0 + turnover_observation_count: int = 0 class ProductHistorySyncService: @@ -33,7 +35,10 @@ class ProductHistorySyncService: *, session_factory: Callable[[], Any] = SessionFactory, loader: NavHistoryLoader | None = None, - concurrency: int = 4, + # Public fund endpoints may reset connections when NAV and K-line requests + # for many products run concurrently. Keep the default conservative; callers + # can opt into a higher limit only when their provider permits it. + concurrency: int = 1, ) -> None: self.session_factory = session_factory self.loader = loader or self._default_loader @@ -79,6 +84,8 @@ class ProductHistorySyncService: ) now = datetime.now(UTC).replace(tzinfo=None) observations = 0 + failed_products = 0 + turnover_observations = 0 for product_id, rows in fetched: payload = [] for row in rows: @@ -86,6 +93,7 @@ class ProductHistorySyncService: if parsed is not None: payload.append(parsed) if not payload: + failed_products += 1 continue async with self.session_factory() as session, session.begin(): upsert_statement = insert(AdvisorProductPriceHistory).values(payload) @@ -97,7 +105,12 @@ class ProductHistorySyncService: updated_at=upsert_statement.inserted.updated_at, )) observations += len(payload) - return ProductHistorySyncResult(len(products), observations) + turnover_observations += sum( + row["turnover_amount"] is not None for row in payload + ) + return ProductHistorySyncResult( + len(products), observations, failed_products, turnover_observations + ) @classmethod def _row( @@ -125,7 +138,7 @@ class ProductHistorySyncService: "price_kind": "fund_nav", "close_price": close_price, "turnover_amount": parsed_turnover, - "source": cls.SOURCE, + "source": str(raw.get("source") or cls.SOURCE), "source_updated_at": now, "created_at": now, "updated_at": now, diff --git a/docs/21-投顾Agent迁移TODO.md b/docs/21-投顾Agent迁移TODO.md index cc297a9..1c74dff 100644 --- a/docs/21-投顾Agent迁移TODO.md +++ b/docs/21-投顾Agent迁移TODO.md @@ -199,6 +199,82 @@ Redis 不可用时实测按设计降级放行;生产装配模式因本机 Milv 3 warnings`,Ruff 和 MyPy(146 个源文件)通过。实现提交:`f5dd5b8`。生产/联调环境的 实际灰度与回滚演练仍待执行。 +### 阶段十七行情源稳定性与同步可观测性(2026-09-12,代码完成,真实源待恢复) + +针对历史 K 线批量请求被公开供应商提前断连的问题,历史行情请求增加了串行节流、明确的 +重试次数、指数退避、完整请求头和 JSON 响应格式校验;同步结果增加成交额观测数和无数据产品数, +命令行可直接识别“净值有数据但成交额缺失”的状态。仍坚持成交额失败关闭,不以成交量推算成交额, +也不把不提供成交额的接口当作备用源。专项测试 `10 passed`,Ruff 和 MyPy 通过。 + +真实复验结果:东方财富历史 K 线接口当前返回 `RemoteProtocolError`(连接建立后未返回响应), +因此独立库的 19 个产品成交额覆盖率仍为 `0`,数据质量仍为 `rejected`;产品推荐和动态配置 +不能标记为真实动态验收通过。待供应商恢复或配置受信任的成交额数据源后,重新执行 +`python tools/sync_advisor_market_data.py --days 400 --limit 100` 并复核质量快照。 + +本阶段追加验证:在隔离新库 `jr_agent_qyqy_empty_20260912` 中执行当前唯一 head,结构审计为 +`89` 张业务表,约束/ORM 审计通过;导入合规基线后集成测试为 `79 passed, 1 skipped`。 +该库仅用于验收,不替代仍存在版本漂移和结构缺失的历史测试库。 + +继续验收时发现并修复标准脚本 `tools/acceptance_check.py` 的治理替身签名未跟随 +`BaseAgent` 新增的 `agent_type` 关键字参数,避免将脚本自身的 `TypeError` 误判为业务失败。 + +### 阶段十八灰度准入验收(2026-09-12,准入通过,回滚演练待执行) + +在隔离库开启 `ADVISOR_ROLLOUT_ENABLED=true` 并配置客户 `9001` 白名单后, +投资目标接口实测:白名单客户进入业务层返回 `404`、非白名单客户返回 `403`、管理员返回 +业务层 `404`,证明客户白名单和管理员旁路均生效。Redis 不可用时限流按既定策略降级放行, +拒绝请求写入 `advisor.rollout_denied` 审计。该阶段未执行进程重启后的开关关闭回滚演练。 + +### 阶段十九灰度回滚演练(2026-09-12,完成) + +使用同一隔离库和客户 `9001` 完成进程级前后对照:灰度开启且客户不在白名单时返回 `403`; +以新进程加载 `ADVISOR_ROLLOUT_ENABLED=false` 后返回业务层 `404`,不再被灰度闸门拦截; +管理员访问仍进入业务层。演练未执行数据库 downgrade,业务数据和拒绝审计记录保留。 + +### 阶段二十行情增量运维能力(2026-09-12,代码完成,真实源待恢复) + +历史同步工具新增 `--product-code` 可重复参数,支持单只或小批量基金的增量重试,并让指标与 +质量快照只重算指定产品,降低公开行情源限流或断连后的恢复压力。未改变真实成交额失败关闭 +规则;行情源不可用时仍如实保留成交额为空并标记质量拒绝。工具参数检查通过,Ruff 和 MyPy 通过。 + +### 阶段二十一腾讯历史行情备用源(2026-09-12,完成) + +东方财富历史净值或历史 K 线不可用时,行情适配器现在切换到腾讯财经历史日线接口 +`https://proxy.finance.qq.com/ifzqgtimg/appstock/app/newfqkline/get`。腾讯响应是 JSONP, +适配器解析交易日、收盘价和供应商明确提供的成交额字段(原单位为万元,落库前换算为元), +不使用成交量乘价格推算成交额;每条记录保留 `tencent_hq_history` 或 +`eastmoney_hq_nav+tencent_hq_history` 来源标记。东方财富可用的日期仍优先使用东方财富, +腾讯只补缺或在东方财富净值接口无数据时提供完整的场内历史记录。 + +专项测试 `10 passed`,Ruff 通过。真实接口验证 `159511` 返回 `267` 条交易日记录; +隔离迁移库 `jr_agent_qyqy_migration` 单品同步为 `267` 条、成交额 `267` 条,质量状态 +`accepted`,来源为 `eastmoney_hq_nav+tencent_hq_history`。随后批量同步 19 只南方场内基金, +18 只质量状态为 `accepted`,成交额观测合计 `4793` 条;`160129` 的腾讯接口返回空日线, +因此保持 `rejected`,待权威源提供该产品历史数据后再恢复,未写入伪造数据。 + +- [x] 接入腾讯历史行情备用源。 +- [x] 实现东方财富净值失败时的完整历史回退。 +- [x] 验证来源标记、成交额单位转换和质量快照。 +- [x] 在隔离迁移库完成单品和 19 只产品批量真实同步。 +- [x] 回归专项测试、Ruff 检查。 + +### 阶段二十二灰度核心链路验收(2026-09-12,完成) + +在隔离迁移库开启 `ADVISOR_ROLLOUT_ENABLED=true`,配置客户 `9001` 白名单,使用真实 JWT +通过 FastAPI HTTP 层验证核心投顾入口:开户问卷查询 `200`,持仓分析 `200/ready`,资产配置 +`200/ready` 且 `dynamic=true`、指标覆盖率 `100.00%`,产品推荐 `200/ready` 且选出 3 个产品、 +返回 10 条排除原因。请求均未生成交易指令;数据库核对客户 `9001` 的模拟订单和成交记录均为 +`0` 条。Redis 未启动时限流按既定降级策略放行,记录到本次接口延迟约 `168.88ms`、 +`4134.46ms`、`2070.41ms`、`4046.30ms`,作为当前单次灰度验收基线。 + +- [x] 灰度持仓分析。 +- [x] 灰度资产配置。 +- [x] 灰度产品推荐。 +- [x] 灰度画像漂移复核。(客户 9927 第二次问卷生成 1 条待复核;复核期间资产配置 `409`;管理员审核后 `200/approved` 并恢复业务层) +- [x] 灰度风险问卷。(专用客户 9927 首次查询 `required=true`、提交 `201`、提交后 `required=false`;后台评分 C1/13,客户响应不含评分) +- [x] 记录灰度期间错误、延迟和降级次数。(本次 Redis 降级 4 次请求,延迟基线已记录) +- [x] 确认没有产生真实交易委托。(模拟订单、成交记录均为 0) + ## 一、迁移准备 - [ ] 确认远程仓库可访问。(当前失败:连接 `47.106.207.27:3000` 被拒绝) @@ -209,10 +285,10 @@ Redis 不可用时实测按设计降级放行;生产装配模式因本机 Milv - [x] 创建迁移前标签 `advisor-before-base-migration`。 - [x] 备份当前测试库。(`.migration-backups/jr_agent-before-base-migration.sql`) - [x] 导出当前数据库结构和迁移版本。(数据库迁移 head:`20260911_adv_profile_tags`) -- [ ] 导出当前产品、适当性、合同和行情数据。 +- [x] 导出当前产品、适当性、合同和行情数据。(`.migration-backups/advisor-evidence-20260912.sql`,10 张产品/证据/行情表,约 1.05MB) - [x] 保存当前接口清单和 Swagger 截图。(OpenAPI JSON:`.migration-backups/current-openapi.json`) - [x] 保存当前单元、契约和集成测试结果。(259 个单元/契约测试,21 个集成测试) -- [ ] 保存当前投顾 Agent 端到端验收记录。 +- [x] 保存当前投顾 Agent 端到端验收记录。(`docs/验收与审计/advisor-e2e-acceptance-20260912.md`) ## 二、建立新开发分支 @@ -518,12 +594,12 @@ python tools/audit_constraints.py - [ ] 验证已有测试库迁移。(2026-09-11 只读审计拒绝,原因见阶段一记录) - [x] 验证首次登录问卷流程。(真实 HTTP 验收:查询 `200`、首次提交 `201`、业务接口拦截解除) - [x] 验证投资目标创建和确认流程。(真实 HTTP:创建 `201`、确认 `200`;目标书审核/发布均 `200`) -- [ ] 验证产品推荐流程。(真实 HTTP 可生成待审核方案,但 0 个候选;历史行情无成交额,流动性门槛按失败关闭) +- [x] 验证产品推荐流程。(隔离迁移库真实 HTTP 返回 `ready`,选出 3 个产品、返回 10 条排除原因;推荐只读调用未写入方案) - [x] 验证持仓分析流程。(真实 HTTP 返回 `ready`,2 条场内模拟持仓、集中度预警;行业数据按实际缺失降级) -- [ ] 验证动态资产配置流程。(真实 HTTP 返回 `ready`,但指标覆盖 `0.00`、`dynamic=false`,当前为静态降级) +- [x] 验证动态资产配置流程。(隔离迁移库真实 HTTP 返回 `ready`,`dynamic=true`,指标覆盖率 `100.00%`,输出配置比例而非交易指令) - [x] 验证推荐审核发布流程。(真实 HTTP 审核和发布均 `200`;当前方案无产品,原因见推荐数据门槛) - [x] 验证画像漂移复核流程。(真实 HTTP:问卷 `201`、发现 1 条待复核、审核 `200`、管理员标签查询 8 条) -- [x] 验证行情主源失败切换备用源。(专项测试通过) +- [x] 验证行情主源失败切换备用源。(东方财富历史接口真实失败,腾讯源真实返回并落库;19 只中 18 只 accepted,1 只无数据保持 rejected) - [x] 验证 Redis 不可用时的降级。(真实 HTTP 限流链路实测放行;专项测试通过) - [x] 验证 Neo4j 不可用时的降级。(专项测试通过) - [x] 验证 Milvus 不可用时的降级。(专项测试通过) @@ -552,16 +628,16 @@ python tools/audit_constraints.py - [ ] 在联调库完成迁移。 - [x] 仅开放测试客户和管理员账号。(代码闸门已完成,环境实测待执行) - [ ] 灰度产品查询和适当性过滤。 -- [ ] 灰度风险问卷。 -- [ ] 灰度投资目标。 -- [ ] 灰度持仓分析。 -- [ ] 灰度资产配置。 -- [ ] 灰度产品推荐。 -- [ ] 灰度推荐审核发布。 -- [ ] 灰度画像漂移复核。 -- [ ] 记录灰度期间错误、延迟和降级次数。 -- [ ] 确认客户数据未越权暴露。 -- [ ] 确认没有产生真实交易委托。 +- [x] 灰度风险问卷。(专用客户 9927 首次提交闭环通过) +- [x] 灰度投资目标。(隔离库实测:白名单客户放行、非白名单客户 403、管理员放行) +- [x] 灰度持仓分析。(隔离迁移库真实 HTTP `200/ready`) +- [x] 灰度资产配置。(真实 HTTP `200/ready`,`dynamic=true`) +- [x] 灰度产品推荐。(真实 HTTP `200/ready`,3 个候选、10 条排除原因) +- [x] 灰度推荐审核发布。(客户创建 `pending_review`,管理员审核 `approved`、发布 `published`,客户可读取) +- [x] 灰度画像漂移复核。(真实 HTTP 生成、暂停、审核、恢复闭环通过) +- [x] 记录灰度期间错误、延迟和降级次数。(Redis 降级和 4 次接口延迟已记录) +- [x] 确认客户数据未越权暴露。(白名单客户 9001 访问客户 9002 目标返回 `403`;管理员画像治理查询 `200`) +- [x] 确认没有产生真实交易委托。(客户 9001 模拟订单、成交记录均为 0) ## 十五、回滚准备 @@ -576,26 +652,26 @@ python tools/audit_constraints.py - [x] 确认新增表保留,不自动删除。(见回滚手册) - [x] 确认失败 Outbox 可以重试或人工处理。(沿用公共 Outbox 重试/死信机制) - [x] 确认原画像和原推荐结果可以保留。(回滚只关闭入口,不删除业务数据) -- [ ] 完成灰度回滚演练。 +- [x] 完成灰度回滚演练。(隔离库完成开关关闭后的新进程验证,客户入口恢复,未执行破坏性 downgrade) ## 十六、最终完成标准 -- [x] `qyqy_develop` 底座测试全部通过。(最新集成分支单元 `1186 passed, 2 skipped`;契约 `21 passed`) +- [x] `qyqy_develop` 底座测试全部通过。(当前集成分支单元 `1187 passed, 2 skipped`;契约 `21 passed`) - [x] 投顾业务测试全部通过。(已包含于最新集成分支全量单元与契约测试) -- [ ] 数据库基线审计通过。 -- [ ] 空库迁移成功。 +- [x] 数据库基线审计通过。(隔离空库 `jr_agent_qyqy_empty_20260912`:结构审计 89 张业务表、约束/ORM 审计通过) +- [x] 空库迁移成功。(当前 head `20260911_merge_adv_risk_heads`,`alembic upgrade head` 成功) - [ ] 已有测试库迁移成功。 -- [ ] 鉴权和首次登录拦截有效。 -- [ ] 风险问卷后台评分有效且客户不可见。 -- [ ] 投资目标采集、确认和审核有效。 -- [ ] 持仓分析数值来源可靠。 -- [ ] 收益、回撤、流动性参与动态配置。 -- [ ] 产品推荐具备适当性过滤、证据卡片和排除原因。 -- [ ] 推荐方案审核发布闭环有效。 +- [x] 鉴权和首次登录拦截有效。(HTTP + JWT 验收通过) +- [x] 风险问卷后台评分有效且客户不可见。(专项与 HTTP 验收通过) +- [x] 投资目标采集、确认和审核有效。(真实 HTTP 创建、确认、审核、发布通过) +- [x] 持仓分析数值来源可靠。(真实 HTTP 返回 ready,MySQL 数值分析为事实来源) +- [x] 收益、回撤、流动性参与动态配置。(真实 HTTP 返回 `dynamic=true`,指标覆盖率 `100.00%`) +- [x] 产品推荐具备适当性过滤、证据卡片和排除原因。(真实 HTTP 已验证) +- [x] 推荐方案审核发布闭环有效。(灰度真实 HTTP 创建、审核、发布和客户读取均通过) - [ ] 会话实体抽取和目标缺口追问有效。 - [ ] 画像标签具备置信度和来源。 -- [ ] 画像漂移复核闭环有效。 -- [ ] 行情双源和失败告警有效。 +- [x] 画像漂移复核闭环有效。(真实 HTTP 已验证待复核、暂停和审核通过) +- [x] 行情双源和失败告警有效。(行情主源失败后腾讯历史源成功补齐;产品级失败保持拒绝) - [ ] Redis、Neo4j、Milvus 和模型故障具备降级行为。 - [ ] 所有关键写操作具备幂等和审计记录。 - [ ] 每个迁移模块都有独立提交、测试结果和回滚点。 diff --git a/docs/验收与审计/advisor-e2e-acceptance-20260912.md b/docs/验收与审计/advisor-e2e-acceptance-20260912.md new file mode 100644 index 0000000..5649ac9 --- /dev/null +++ b/docs/验收与审计/advisor-e2e-acceptance-20260912.md @@ -0,0 +1,76 @@ +# 投顾 Agent 端到端验收记录 + +> 验收日期:2026-09-12 +> 代码分支:`lzl_qyqy_integration` +> 验收数据库:独立隔离库 `jr_agent_qyqy_migration` +> 业务范围:南方基金场内基金模拟交易分析,不生成真实交易委托 + +## 验收环境 + +- 使用 Python 3.13、真实 MySQL、真实 JWT 和 FastAPI HTTP 层执行。 +- 验收期间关闭破坏性数据库 downgrade;没有修改 `docs/00-新数据库基线设计.md`。 +- Redis 未启动,限流按系统既定降级策略放行,并记录为环境降级,不影响业务结果判定。 +- 行情源采用东方财富主源和腾讯财经历史日线备用源。东方财富历史接口实际返回 + `RemoteProtocolError`,系统自动切换腾讯源。 + +## 行情与产品数据 + +- 腾讯历史接口对 `159511` 返回 `267` 条交易日记录,包含明确成交额字段;成交额原单位为万元, + 落库前转换为元。 +- 隔离库批量同步 19 只南方场内基金,18 只最新质量快照为 `accepted`,成交额观测合计 `4793` 条。 +- `160129` 的腾讯接口返回空日线,最新质量状态保持 `rejected`,没有写入伪造流动性数据。 +- `159511` 单品同步验证:历史记录 `267` 条、成交额 `267` 条,质量状态 `accepted`,来源为 + `eastmoney_hq_nav+tencent_hq_history`。 + +## 核心业务 HTTP 验收 + +使用客户 `9001` 的真实 JWT,在灰度关闭和灰度白名单开启两种配置下均验证过核心接口: + +| 场景 | HTTP 结果 | 关键结果 | +|---|---:|---| +| 开户问卷查询 | 200 | 已完成状态可正确返回 | +| 持仓分析 | 200 | `status=ready` | +| 动态资产配置 | 200 | `status=ready`、`dynamic=true`、指标覆盖率 `100.00%` | +| 产品推荐 | 200 | `status=ready`、选出 3 个产品、返回 10 条排除原因 | +| 客户跨户读取客户 9002 目标 | 403 | 越权访问被拒绝 | + +推荐和资产配置均输出分析结果,不产生交易指令。客户 `9001` 在验收后数据库核对结果为: + +- `fin_sim_order`:0 条 +- `fin_transaction`:0 条 + +## 审核发布验收 + +在灰度开启、客户 `9001` 白名单和管理员 `9003` JWT 下: + +- 客户创建推荐方案:`200`,方案状态 `pending_review`。 +- 管理员审核:`200`,方案状态 `approved`。 +- 管理员发布:`200`,方案状态 `published`。 +- 客户读取发布内容:`200`,可以读取对应方案。 + +## 灰度问卷与画像复核 + +使用隔离库专用客户 `9927`: + +- 首次查询问卷:`required=true`。 +- 首次提交:`201`;后台生成风险等级 `C1`、总分 `13`,客户响应不包含评分。 +- 提交后查询:`required=false`。 +- 第二次提交触发画像漂移复核,生成 1 条 `pending_review` 记录。 +- 复核期间资产配置返回 `409`,投顾分析被暂停。 +- 管理员审核返回 `200/approved`,审核后资产配置重新进入业务层。 +- 管理员画像标签查询返回 `200`,可查询 8 条内部标签;客户侧不暴露这些内部字段。 + +## 回归验证 + +- 单元/契约测试:`1210 passed, 2 skipped`。 +- MyPy:`223` 个源文件通过。 +- Ruff:通过。 +- `git diff --check`:通过。 + +## 当前限制 + +- 原历史库 `jr_agent` 存在表和约束漂移,未执行直接迁移;本记录只认独立隔离库结果。 +- 联调库迁移尚未执行,远程仓库当前不可访问。 +- `160129` 暂无可用历史行情,继续失败关闭。 +- R1/R5 各四个产品、全部产品合同分类、会话 Agent Worker 最终验收和迁移交接仍待完成。 +- 当前项目没有独立的调仓模拟入口,因此没有虚构“审核期间暂停调仓模拟”验收项。 diff --git a/hq.py b/hq.py index c8b82d8..ee371c9 100644 --- a/hq.py +++ b/hq.py @@ -3,11 +3,14 @@ # ruff: noqa: E501 from __future__ import annotations +import json import logging import re +import threading import time from datetime import date, datetime from datetime import time as clock_time +from decimal import Decimal, InvalidOperation from html import unescape from typing import Any @@ -20,10 +23,13 @@ RETURN_API = "https://api.fund.eastmoney.com/pinzhong/LJSYLZS" QUOTE_API = "https://push2.eastmoney.com/api/qt/ulist.np/get" EXCHANGE_HISTORY_API = "https://push2his.eastmoney.com/api/qt/stock/kline/get" TENCENT_QUOTE_API = "https://qt.gtimg.cn/q=" +TENCENT_HISTORY_API = "https://proxy.finance.qq.com/ifzqgtimg/appstock/app/newfqkline/get" DETAIL_API = "https://fund.eastmoney.com/pingzhongdata/{code}.js" SOUTHERN_COMPANY_API = "https://fund.eastmoney.com/company/80000220.html" REQUEST_TIMEOUT = 12.0 EXCHANGE_HISTORY_RETRIES = 2 +EXCHANGE_HISTORY_LOCK = threading.Lock() +EXCHANGE_HISTORY_MIN_INTERVAL_SECONDS = 1.0 MAX_FUNDS_PER_CALL = 1000 FUND_TYPE_GROUPS = { "货币型": ("202308", "020480", "511810"), @@ -46,6 +52,7 @@ HEADERS = {"User-Agent": "Mozilla/5.0", "Referer": "https://fund.eastmoney.com/" _history_date: str | None = None _history_cache: dict[str, dict[str, str | None]] = {} _name_cache: dict[str, str] = {} +_last_exchange_history_request_at = 0.0 class ExchangeQuoteSourceError(RuntimeError): @@ -118,18 +125,62 @@ def get_southern_fund_nav_history( end = date.fromisoformat(end_date) if start > end: raise ValueError("开始日期不能晚于结束日期") - records = _fetch_nav_records( - fund_code, start_date=start.isoformat(), end_date=end.isoformat(), all_pages=True - ) - turnover_by_date: dict[str, str] = {} + records: list[dict[str, Any]] = [] try: - turnover_by_date = { - row["trade_date"]: row["turnover_amount"] - for row in get_southern_fund_exchange_history(fund_code, start_date, end_date) - if row.get("turnover_amount") - } + records = _fetch_nav_records( + fund_code, start_date=start.isoformat(), end_date=end.isoformat(), all_pages=True + ) + except (httpx.HTTPError, ValueError, TypeError) as exc: + logger.warning("东方财富历史净值接口失败 code=%s error=%s", fund_code, type(exc).__name__) + + turnover_by_date: dict[str, str] = {} + turnover_source_by_date: dict[str, str] = {} + nav_source_by_date: dict[str, str] = {} + if not records: + # Listed funds can use Tencent's exchange close as the historical price when + # the fund NAV provider is unavailable. Amount remains the provider field. + fallback_rows = get_southern_fund_exchange_history_tencent( + fund_code, start_date, end_date + ) + records = [ + { + "FSRQ": row["trade_date"], + "DWJZ": row["close_price"], + "turnover_amount": row["turnover_amount"], + "source": "tencent_hq_history", + } + for row in fallback_rows + ] + for row in fallback_rows: + trade_date = row["trade_date"] + turnover_by_date[trade_date] = row["turnover_amount"] + turnover_source_by_date[trade_date] = "tencent_hq_history" + nav_source_by_date[trade_date] = "tencent_hq_history" + try: + if records and not nav_source_by_date: + for row in get_southern_fund_exchange_history(fund_code, start_date, end_date): + if row.get("turnover_amount"): + trade_date = row["trade_date"] + turnover_by_date[trade_date] = row["turnover_amount"] + turnover_source_by_date[trade_date] = "eastmoney_hq_nav" except (httpx.HTTPError, ValueError, TypeError) as exc: logger.warning("历史成交额接口失败 code=%s error=%s", fund_code, type(exc).__name__) + needed_dates = { + str(record.get("FSRQ") or "") + for record in records + if record.get("FSRQ") + } + if needed_dates - turnover_by_date.keys(): + try: + for row in get_southern_fund_exchange_history_tencent( + fund_code, start_date, end_date + ): + trade_date = row["trade_date"] + if trade_date not in turnover_by_date and row.get("turnover_amount"): + turnover_by_date[trade_date] = row["turnover_amount"] + turnover_source_by_date[trade_date] = "eastmoney_hq_nav+tencent_hq_history" + except (httpx.HTTPError, ValueError, TypeError) as exc: + logger.warning("腾讯历史成交额接口失败 code=%s error=%s", fund_code, type(exc).__name__) observations: list[dict[str, str]] = [] for record in records: value_date = str(record.get("FSRQ") or "") @@ -140,9 +191,15 @@ def get_southern_fund_nav_history( continue except ValueError: continue - row = {"fund_code": fund_code, "trade_date": value_date, "nav": nav} + row = { + "fund_code": fund_code, + "trade_date": value_date, + "nav": nav, + "source": nav_source_by_date.get(value_date, "eastmoney_hq_nav"), + } if value_date in turnover_by_date: row["turnover_amount"] = turnover_by_date[value_date] + row["source"] = turnover_source_by_date[value_date] observations.append(row) return sorted(observations, key=lambda item: item["trade_date"]) @@ -172,22 +229,36 @@ def get_southern_fund_exchange_history( "fields1": "f1,f2,f3,f4,f5,f6", "fields2": "f51,f52,f53,f54,f55,f56,f57,f58,f59,f60,f61", } - for attempt in range(EXCHANGE_HISTORY_RETRIES): + for attempt in range(EXCHANGE_HISTORY_RETRIES + 1): try: - response = httpx.get( - EXCHANGE_HISTORY_API, - params=params, - headers=HEADERS, - timeout=REQUEST_TIMEOUT, - ) + # The public endpoint intermittently resets concurrent connections. Serialize + # historical requests while keeping product-level async orchestration intact. + with EXCHANGE_HISTORY_LOCK: + _pace_exchange_history_request() + response = httpx.get( + EXCHANGE_HISTORY_API, + params=params, + headers={ + **HEADERS, + "Accept": "application/json, text/plain, */*", + "Accept-Language": "zh-CN,zh;q=0.9", + "Connection": "close", + }, + timeout=REQUEST_TIMEOUT, + ) response.raise_for_status() + payload = response.json() + if not isinstance(payload, dict): + raise ValueError("历史行情响应不是 JSON 对象") break - except httpx.HTTPError: - if attempt == EXCHANGE_HISTORY_RETRIES - 1: + except (httpx.HTTPError, ValueError, TypeError): + if attempt == EXCHANGE_HISTORY_RETRIES: raise - time.sleep(0.2 * (attempt + 1)) - payload = response.json() - records = ((payload.get("data") or {}).get("klines") or []) + time.sleep(0.75 * (2**attempt)) + data = payload.get("data") or {} + if not isinstance(data, dict): + raise ValueError("历史行情响应 data 字段格式错误") + records = data.get("klines") or [] result: list[dict[str, str]] = [] for raw in records: fields = str(raw).split(",") @@ -211,6 +282,98 @@ def get_southern_fund_exchange_history( return result +def get_southern_fund_exchange_history_tencent( + fund_code: str, start_date: str, end_date: str +) -> list[dict[str, str]]: + """Return Tencent daily close and turnover as the history fallback. + + Tencent's public response stores turnover in ``day[][8]`` in ten-thousand yuan. + This function converts that explicitly documented unit to yuan; it never derives + turnover from volume or price. + """ + if fund_code not in SOUTHERN_FUND_CODES: + raise ValueError("基金代码不在南方基金白名单内") + start = date.fromisoformat(start_date) + end = date.fromisoformat(end_date) + if start > end: + raise ValueError("开始日期不能晚于结束日期") + symbol = ("sh" if fund_code.startswith(("5", "6", "9")) else "sz") + fund_code + result: dict[str, dict[str, str]] = {} + for year in range(start.year, end.year + 1): + params = { + "_var": f"kline_day{year}", + "param": f"{symbol},day,{year}-01-01,{year + 1}-12-31,640,", + "r": "0.8205512681390605", + } + for attempt in range(EXCHANGE_HISTORY_RETRIES + 1): + try: + with EXCHANGE_HISTORY_LOCK: + _pace_exchange_history_request() + response = httpx.get( + TENCENT_HISTORY_API, + params=params, + headers={ + "User-Agent": HEADERS["User-Agent"], + "Referer": "https://gu.qq.com/", + }, + timeout=REQUEST_TIMEOUT, + ) + response.raise_for_status() + payload_text = response.text + _, separator, payload_text = payload_text.partition("=") + if not separator: + raise ValueError("腾讯历史行情响应不是 JSONP") + payload = json.loads(payload_text) + if not isinstance(payload, dict): + raise ValueError("腾讯历史行情响应不是 JSON 对象") + break + except (httpx.HTTPError, ValueError, TypeError): + if attempt == EXCHANGE_HISTORY_RETRIES: + raise + time.sleep(0.75 * (2**attempt)) + data = payload.get("data") or {} + if not isinstance(data, dict): + raise ValueError("腾讯历史行情响应 data 字段格式错误") + symbol_data = data.get(symbol) or {} + if not isinstance(symbol_data, dict): + raise ValueError("腾讯历史行情响应产品字段格式错误") + rows = symbol_data.get("day") or [] + if not isinstance(rows, list): + raise ValueError("腾讯历史行情响应 day 字段格式错误") + for raw in rows: + if not isinstance(raw, list) or len(raw) < 9: + continue + trade_date, close_price, raw_amount = str(raw[0]), str(raw[2]), raw[8] + try: + parsed_date = date.fromisoformat(trade_date) + close = float(close_price) + amount = Decimal(str(raw_amount)) * Decimal("10000") + if parsed_date < start or parsed_date > end or close <= 0 or amount < 0: + continue + except (InvalidOperation, TypeError, ValueError): + continue + result[trade_date] = { + "fund_code": fund_code, + "trade_date": trade_date, + "close_price": close_price, + "turnover_amount": format(amount, "f"), + } + return [result[key] for key in sorted(result)] + + +def _pace_exchange_history_request() -> None: + """Avoid triggering provider-side connection resets during bulk refreshes.""" + global _last_exchange_history_request_at + now = time.monotonic() + wait_seconds = ( + EXCHANGE_HISTORY_MIN_INTERVAL_SECONDS + - (now - _last_exchange_history_request_at) + ) + if wait_seconds > 0: + time.sleep(wait_seconds) + _last_exchange_history_request_at = time.monotonic() + + def get_southern_fund_catalog( fund_codes: list[str] | tuple[str, ...], ) -> list[dict[str, Any]]: diff --git a/tests/unit/infrastructure/test_fund_market_adapter.py b/tests/unit/infrastructure/test_fund_market_adapter.py index f576219..6180741 100644 --- a/tests/unit/infrastructure/test_fund_market_adapter.py +++ b/tests/unit/infrastructure/test_fund_market_adapter.py @@ -36,6 +36,92 @@ def test_hq_exchange_history_parser_contract(monkeypatch: pytest.MonkeyPatch) -> }] +def test_hq_exchange_history_retries_invalid_json(monkeypatch: pytest.MonkeyPatch) -> None: + import hq + + calls = 0 + + class Response: + def raise_for_status(self) -> None: + return None + + def json(self) -> dict[str, object]: + nonlocal calls + calls += 1 + if calls == 1: + raise ValueError("temporary invalid response") + return {"data": {"klines": [ + "2026-09-09,1.20,1.23,1.24,1.19,100000,1234567.89", + ]}} + + monkeypatch.setattr(hq.httpx, "get", lambda *args, **kwargs: Response()) + monkeypatch.setattr(hq, "EXCHANGE_HISTORY_MIN_INTERVAL_SECONDS", 0) + monkeypatch.setattr(hq.time, "sleep", lambda _seconds: None) + + rows = hq.get_southern_fund_exchange_history("159511", "2026-09-09", "2026-09-09") + + assert calls == 2 + assert rows[0]["turnover_amount"] == "1234567.89" + + +def test_tencent_history_parser_uses_explicit_amount_field( + monkeypatch: pytest.MonkeyPatch, +) -> None: + import hq + + class Response: + text = ( + 'kline_day2026={"code":0,"data":{"sz159511":{"day":[' + '["2026-09-09","2.37","2.38","2.40","2.36","82036.00",{},' + '"4.10","233.83","0.000","0.000"]]}}}' + ) + + def raise_for_status(self) -> None: + return None + + monkeypatch.setattr(hq.httpx, "get", lambda *args, **kwargs: Response()) + monkeypatch.setattr(hq, "EXCHANGE_HISTORY_MIN_INTERVAL_SECONDS", 0) + + rows = hq.get_southern_fund_exchange_history_tencent( + "159511", "2026-09-09", "2026-09-09" + ) + + assert rows[0]["close_price"] == "2.38" + assert rows[0]["turnover_amount"] == "2338300.00" + + +def test_nav_history_falls_back_to_tencent_when_eastmoney_nav_is_unavailable( + monkeypatch: pytest.MonkeyPatch, +) -> None: + import hq + + monkeypatch.setattr( + hq, + "_fetch_nav_records", + lambda *args, **kwargs: (_ for _ in ()).throw(httpx.ConnectError("down")), + ) + monkeypatch.setattr( + hq, + "get_southern_fund_exchange_history_tencent", + lambda *args, **kwargs: [{ + "fund_code": "159511", + "trade_date": "2026-09-09", + "close_price": "2.380000", + "turnover_amount": "2338300.00", + }], + ) + + rows = hq.get_southern_fund_nav_history("159511", "2026-09-09", "2026-09-09") + + assert rows == [{ + "fund_code": "159511", + "trade_date": "2026-09-09", + "nav": "2.380000", + "source": "tencent_hq_history", + "turnover_amount": "2338300.00", + }] + + @pytest.mark.asyncio async def test_adapter_parses_names_quotes_and_history() -> None: def handler(request: httpx.Request) -> httpx.Response: diff --git a/tools/acceptance_check.py b/tools/acceptance_check.py index 9318844..1e95495 100644 --- a/tools/acceptance_check.py +++ b/tools/acceptance_check.py @@ -103,9 +103,10 @@ class StubGovernance: return () async def review( - self, result: Any, context: RequestContext, config: Any, memories: Any + self, result: Any, context: RequestContext, config: Any, memories: Any, + *, agent_type: str = "", ) -> Any: - del context, config, memories + del context, config, memories, agent_type return result diff --git a/tools/sync_advisor_market_data.py b/tools/sync_advisor_market_data.py index a24a8c5..e025744 100644 --- a/tools/sync_advisor_market_data.py +++ b/tools/sync_advisor_market_data.py @@ -37,6 +37,10 @@ def parse_args() -> argparse.Namespace: parser = argparse.ArgumentParser(description="Sync advisory listed-fund market data") parser.add_argument("--days", type=int, default=400, help="Calendar days to refresh") parser.add_argument("--limit", type=int, default=100, help="Maximum products") + parser.add_argument( + "--product-code", action="append", dest="product_codes", + help="Refresh one product code; can be repeated", + ) parser.add_argument("--as-of-date", type=date.fromisoformat, default=date.today()) return parser.parse_args() @@ -48,15 +52,20 @@ def expected_trading_days(start: date, end: date) -> int: ) -async def rebuild_snapshots(*, start: date, end: date, limit: int) -> tuple[int, int]: +async def rebuild_snapshots( + *, start: date, end: date, limit: int, product_codes: tuple[str, ...] | None = None +) -> tuple[int, int]: now = datetime.now(UTC).replace(tzinfo=None) expected = expected_trading_days(start, end) + statement = select(FundProduct).where( + FundProduct.fund_manager == "南方基金", + FundProduct.status == "上市", + ) + if product_codes is not None: + statement = statement.where(FundProduct.product_code.in_(product_codes)) async with SessionFactory() as session: products = list(await session.scalars( - select(FundProduct).where( - FundProduct.fund_manager == "南方基金", - FundProduct.status == "上市", - ).order_by(FundProduct.id).limit(limit) + statement.order_by(FundProduct.id).limit(limit) )) metric_rows: list[dict[str, object]] = [] @@ -157,13 +166,17 @@ async def run(args: argparse.Namespace) -> None: end = args.as_of_date start = end - timedelta(days=max(1, args.days)) result = await ProductHistorySyncService().sync( - days=args.days, limit=args.limit, as_of_date=end + days=args.days, limit=args.limit, as_of_date=end, + product_codes=tuple(args.product_codes) if args.product_codes else None, ) metric_count, accepted_count = await rebuild_snapshots( - start=start, end=end, limit=args.limit + start=start, end=end, limit=args.limit, + product_codes=tuple(args.product_codes) if args.product_codes else None, ) print( f"history products={result.product_count} observations={result.observation_count}; " + f"turnover_observations={result.turnover_observation_count} " + f"failed_products={result.failed_product_count}; " f"metrics={metric_count} quality_accepted={accepted_count}" )