diff --git a/.trellis/spec/backend/index.md b/.trellis/spec/backend/index.md index f8cd15f..1f5f89c 100644 --- a/.trellis/spec/backend/index.md +++ b/.trellis/spec/backend/index.md @@ -9,6 +9,7 @@ | [目录与模块边界](./directory-structure.md) | 包结构、bounded context 和导入边界 | | [配置与运行时](./configuration-and-runtime.md) | `Settings`、应用工厂和部署环境 | | [市场数据同步](./market-data-sync.md) | Tushare qfq、PostgreSQL、CSV 快照和一次性 Job 契约 | +| [Tushare 当前上市股票范围](./tushare-listed-stock-universe.md) | 所有股票型功能统一只使用构建时 `stock_basic(list_status=L)` 母集 | | [历史选股](./selection.md) | selection bounded context、目标交易日、qfq 读取和信号结果契约 | | [HTTP 契约](./http-api-contracts.md) | 路由组合、响应模型和同源 API 路径 | | [错误处理](./error-handling.md) | 当前 FastAPI 错误行为及跨层错误传递 | @@ -17,6 +18,7 @@ ## 开发前检查 - 先阅读 `docs/adr/0001-bounded-context-first-modular-monolith.md`,确认新业务是否有清晰的语言和所有权边界。 +- 涉及 Tushare 个股数据时先阅读 `tushare-listed-stock-universe.md`,所有新功能都必须从当前 `L` 股票母集继续缩小范围,禁止重新引入 `D/P/G/UN`。 - 先阅读目标上下文的 `modules//README.md`(如已存在),再决定 domain、application、infrastructure、presentation 的位置。 - 变更 HTTP 字段时同时检查 `zhixing-server/tests/`、前端 feature API 类型以及 `docs/adr/0002-use-a-same-origin-browser-api.md`。 - 不要为了“未来可能需要”创建空的数据库、服务或日志层;当前仓库没有这些实现。 diff --git a/.trellis/spec/backend/tushare-listed-stock-universe.md b/.trellis/spec/backend/tushare-listed-stock-universe.md new file mode 100644 index 0000000..20f327a --- /dev/null +++ b/.trellis/spec/backend/tushare-listed-stock-universe.md @@ -0,0 +1,103 @@ +# Tushare 当前上市股票范围 + +## Scenario: 所有股票型功能统一使用当前 `L` 股票池 + +### 1. Scope / Trigger + +- 触发:新增或修改任何通过 Tushare 获取个股基础资料、行情、资金流、板块成员、财务或估值数据的后端功能。 +- 目标:所有功能统一以构建时 `stock_basic(list_status="L")` 返回的当前上市股票为证券母集,禁止为了历史回溯获取 `D/P/G/UN`。 +- 历史语义:功能上线日视为最早业务历史日期;以后重跑旧日期仍使用重跑当时的当前 `L` 股票池,不保证还原目标日的退市证券。 +- 边界:指数、基金、期货、宏观等非个股数据不适用本股票状态契约;若未来产品必须恢复历史时点证券生命周期,必须先显式修改本规格及对应任务设计,不能在单个 adapter 内局部绕过。 + +### 2. Signatures + +所有直接读取股票基础档案的 Tushare adapter 必须显式传入 `list_status="L"`: + +```python +client.query( + "stock_basic", + exchange="", + list_status="L", + fields="ts_code,symbol,name,market,exchange,list_status,list_date,delist_date", +) +``` + +应用层不得通过 adapter 隐式缓存推断股票范围;需要候选股票的端口必须显式接收已经排序、去重并与当前 `L` 股票池相交的代码集合,例如: + +```python +def fetch_moneyflow_dc( + trade_date: date, + candidate_codes: Sequence[str], +) -> SourceResult[MoneyflowDcRow]: ... +``` + +### 3. Contracts + +- `stock_basic` 请求必须显式设置 `list_status="L"`,不能依赖供应商默认值,也不能循环请求 `D/P/G/UN`。 +- 当前股票母集至少以 `ts_code` 唯一;返回的非 `L` 记录不得进入业务目标集合。严格 source adapter 应将与请求分区不符的状态视为来源契约错误,已有宽松同步边界至少必须在领域过滤时排除。 +- 股票型功能可以继续执行自身既有的市场边界,例如沪深 A 股、B 股、北交所、ST 或风险警示过滤;这些过滤只能缩小 `L` 母集,不能重新引入其他上市状态。 +- `daily`、`moneyflow_dc`、`dc_member` 等不支持 `list_status` 的接口可以按其最有效的方式获取原始响应,但进入计算、排名、覆盖率、缺口补拉或持久化业务事实前,候选代码必须与当前 `L` 母集取交集。 +- 为审计保存的全市场原始 snapshot 可以包含非 `L` 行;非 `L` 行不得进入规范化事实、策略计算或“应覆盖股票数”。 +- 当前 `L` 股票池必须带有构建时来源快照或等价审计信息。重试若复用旧下游 snapshot,必须确认它仍覆盖本轮候选集合;候选扩大时应在同一次重试中刷新相应下游来源。 +- 本契约不要求各 bounded context 共享数据库表、缓存或 Tushare client;共享的是证券范围语义,而不是运行时耦合。 + +### 4. Validation & Error Matrix + +| 条件 | 必须行为 | +| --- | --- | +| `stock_basic` 请求未显式传 `list_status="L"` | 测试失败;不得发布该功能 | +| `L` 分区返回 `D/P/G/UN` | 严格 adapter 抛来源契约错误,或在既有宽松边界明确排除;非 `L` 不得进入业务集合 | +| 板块成员包含非当前 `L` 股票 | 保留原始成员审计,计算候选与当前 `L` 集合取交集 | +| 目标日期早于当前 `L` 股票的 `list_date` | 从该目标日候选集合排除 | +| 行情或资金流全市场响应包含非 `L` 股票 | 原始 snapshot 可保留,规范化事实和覆盖率忽略这些股票 | +| 缺失补拉收到不在请求候选集合内的代码 | 按来源契约错误 fail closed,禁止合并 | +| 重试时成员恢复导致当前候选集合扩大 | 检查旧下游 snapshot 覆盖;不足时同轮刷新,不能先发布一次可预见的 `partial` | +| 新需求要求历史退市股票或历史时点生命周期 | 先修改本规格并完成独立设计评审,禁止直接请求 `D/P/G/UN` | + +### 5. Good/Base/Bad Cases + +- Good:资金雷达只请求一次 `stock_basic(list_status="L")`,将有效板块成员与当前沪深 A 股交集传给资金流 source;全市场原始资金流即使含额外股票,也只补拉和计算交集内代码。 +- Base:普通行情同步从 `L` 股票池再排除 ST、北交所或不属于目标市场的证券;这是允许的模块级缩小,不改变全局母集。 +- Good:重试刷新成员后发现新增两个当前 `L` 候选,旧资金流 checkpoint 少两只,于是同一次重试只刷新资金流来源组并恢复成功。 +- Bad:为了回填旧日期,将 `stock_basic` 改为循环获取 `L/D/P/G/UN`,或者直接把 `dc_member` 的全部代码作为资金流覆盖分母。 +- Bad:看到全市场原始 snapshot 含非 `L` 股票便将它们写入策略事实,造成候选数量、覆盖率或排名口径漂移。 + +### 6. Tests Required + +- Tushare adapter 测试必须断言 `stock_basic` 的调用参数包含且只包含 `list_status="L"`,并断言非 `L` 返回记录不会进入结果。 +- 应用编排测试必须构造板块成员、未来上市记录和当前 `L` 记录,断言传给下游 source 的候选集合是稳定排序后的交集。 +- 规范化或策略测试必须断言非 `L`、目标日尚未上市、B 股或模块已排除市场不会贡献金额、覆盖率或排名。 +- 重试测试必须覆盖“成员刷新后候选扩大但旧下游 checkpoint 不完整”,断言同一次重试刷新必要来源组。 +- 新增股票型 bounded context 时,至少有一个边界测试证明它没有请求或引入 `D/P/G/UN`。 + +### 7. Wrong vs Correct + +#### Wrong + +```python +# 禁止:为历史回填循环获取全部生命周期状态。 +rows = tuple( + client.query("stock_basic", list_status=status) + for status in ("L", "D", "P", "G", "UN") +) +candidate_codes = tuple(member.stock_code for member in memberships) +``` + +#### Correct + +```python +# 正确:当前 L 是唯一母集,模块规则只能继续缩小它。 +listed = client.query("stock_basic", list_status="L") +listed_codes = { + row.ts_code + for row in listed + if row.list_status == "L" and is_module_eligible(row, target_trade_date) +} +candidate_codes = tuple( + sorted( + member.stock_code + for member in memberships + if member.stock_code in listed_codes + ) +) +``` diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/check.jsonl b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/check.jsonl new file mode 100644 index 0000000..49a051b --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/check.jsonl @@ -0,0 +1,7 @@ +{"file":".trellis/spec/backend/index.md","reason":"核验后端模块边界与开发规范"} +{"file":".trellis/spec/backend/market-data-sync.md","reason":"核验共享 coordinator 未改变 market-data 语义"} +{"file":".trellis/spec/backend/tushare-listed-stock-universe.md","reason":"核验所有业务候选与当前 L 股票母集相交"} +{"file":".trellis/spec/backend/quality-guidelines.md","reason":"执行完整后端质量门禁"} +{"file":".trellis/spec/guides/code-reuse-thinking-guide.md","reason":"核验共享能力没有越界或重复实现"} +{"file":".trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/radar-build-tushare-call-chain.md","reason":"核验 source group、候选集与 retry/checkpoint 契约"} +{"file":".trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/tushare-global-start-interval.md","reason":"核验两路 worker、共享启动间隔及确定性测试覆盖"} diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/design.md b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/design.md new file mode 100644 index 0000000..423f446 --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/design.md @@ -0,0 +1,29 @@ +# 设计:当前上市股票池与 `moneyflow_dc` 缺口补拉 + +## 边界与契约 + +应用层在完成板块成员与 `stock_basic(L)` 采集后,先规范化成员并计算“有效成员代码与当前 L 股票代码的交集”,再调用扩展后的 `SectorRadarSource.fetch_moneyflow_dc(trade_date, candidate_codes)`。候选集合通过端口显式传递;adapter 不缓存先前 `fetch_stock_basics` 的响应,因此 publication replay 和 retry 仍是无隐式状态的。 + +`stock_basic` source 由五分区请求收敛为单一 `L` 分区。领域层继续用现有代码、市场和 `list_date` 规则处理当前上市候选,不新增 ST 过滤或历史退市语义。 + +## 资金流数据流 + +adapter 首先保存按 `trade_date` 获取的全市场 snapshot。它对 rows 执行 typed parsing、目标日期和 `ts_code` 唯一性校验;初始 snapshot 达到 6000 行只表示全市场可能截断,不再单独构成失败。随后计算 `candidate_codes - returned_codes`。 + +缺失集合为空时返回首批 snapshot 与 rows。存在缺失时,按排序后的 `ts_code` 使用固定两路 executor 请求 `moneyflow_dc(trade_date=..., ts_code=...)`。每个成功分片形成独立、带 `partition_key=ts_code` 的 snapshot;分片只允许为空或返回所请求股票在目标日期的唯一记录。空分片以及重试耗尽的普通 provider 异常不产生伪造 snapshot/row,并保留为覆盖缺口;来源 schema、日期、代码、唯一键或 row-limit 契约错误立即上浮。 + +主线程按输入代码顺序汇总 future,保证 `source_order=0` 始终是全市场 snapshot,后续分片按 `ts_code` 稳定排列。合并 rows 后再次验证 `(trade_date, ts_code)` 唯一,防止首批与分片重叠。所有 snapshot 继续归入 `PublicationSourceGroup.MONEYFLOW_DC`,现有数据库模型无需迁移。 + +## 并发与限流 + +共享 `RequestCoordinator` 增加默认值为 0 的 `request_interval_seconds` 和受现有 `threading.Condition` 保护的下次启动时刻。每次 attempt 在调用 provider 前原子等待 cooldown 并预约请求启动槽,预约完成后释放锁,再执行真实请求。资金雷达默认 coordinator 接收现有的 0.2 秒配置,移除 adapter 请求完成后的独立 sleep;因此两个 worker可重叠网络等待,但同一 adapter 中任意两次请求的启动时间仍至少相隔 0.2 秒。 + +当前锁定的 Tushare 1.4.29 `DataApi.query` 只读取 client 的 token、URL 和 timeout,在局部变量中构造参数并调用模块级 `requests.post`,未维护单次请求可变状态。两路 worker 共享该 client 的风险可接受,并由并发单元测试约束;该结论不扩展为 Tushare SDK 的通用线程安全保证。 + +## 兼容性、失败与回滚 + +`RequestCoordinator` 的新参数默认关闭,market-data bounded context 行为不变。端口签名变化同步更新 fake source 和 CLI/build 测试。普通补拉调用在 coordinator 的有限 retry 后仍失败时记录安全日志并留下覆盖缺口;契约错误保持 hard failure。日志不得包含 token 或完整 payload。 + +publication retry 只有在已保存的 `MONEYFLOW_DC` rows 仍覆盖本轮候选代码时才重放该来源组。若 `MEMBERS` 刷新后候选集合扩大,旧资金流 checkpoint 不足以覆盖新增候选,则在同一次 retry 中刷新 `MONEYFLOW_DC`,避免先发布一次可预见的 partial 再要求第二次重试。 + +回滚只需恢复 adapter 的单次 `moneyflow_dc` 请求、旧端口签名和协调器调用方式,不涉及 schema 或数据迁移。已生成的分片 snapshots 使用现有通用存储格式,旧版本即使不能主动生成,也仍可按 source group replay。 diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/implement.jsonl b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/implement.jsonl new file mode 100644 index 0000000..02d8e17 --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/implement.jsonl @@ -0,0 +1,9 @@ +{"file":".trellis/spec/backend/index.md","reason":"后端模块边界、开发前检查和质量入口"} +{"file":".trellis/spec/backend/directory-structure.md","reason":"共享协调器与 sector_radar bounded context 的所有权边界"} +{"file":".trellis/spec/backend/configuration-and-runtime.md","reason":"复用现有请求间隔配置并避免新增环境读取"} +{"file":".trellis/spec/backend/market-data-sync.md","reason":"现有 Tushare coordinator、并发与 checkpoint 相邻契约"} +{"file":".trellis/spec/backend/tushare-listed-stock-universe.md","reason":"所有股票型功能只使用构建时当前 L 股票母集"} +{"file":".trellis/spec/backend/quality-guidelines.md","reason":"后端 Ruff、Pyright 和 pytest 门禁"} +{"file":".trellis/spec/guides/code-reuse-thinking-guide.md","reason":"评估 RequestCoordinator 共享原语扩展"} +{"file":".trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/radar-build-tushare-call-chain.md","reason":"资金雷达调用顺序、候选集与 checkpoint 证据"} +{"file":".trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/tushare-global-start-interval.md","reason":"两路 worker 与全局启动间隔的适配点和测试模式"} diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/implement.md b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/implement.md new file mode 100644 index 0000000..94a31c1 --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/implement.md @@ -0,0 +1,25 @@ +# 实施计划 + +1. 创建 `codex/sector-radar-listed-moneyflow-recovery` 分支并读取目标 backend/shared、sector radar 代码及相关规格。 +2. 扩展 `RequestCoordinator`,实现线程安全的共享请求启动间隔,保留默认关闭与现有 cooldown/retry 行为;补充确定性单元测试。 +3. 将 `stock_basic` 收敛为单一 `L` 分区,并调整 source 测试与返回状态契约。 +4. 前移候选股票形成步骤,扩展 `SectorRadarSource.fetch_moneyflow_dc` 端口,向 adapter 显式传递稳定的当前 L 候选代码。 +5. 在 Tushare adapter 内实现全市场首拉、候选覆盖检查、两路缺失代码补拉、分片契约校验、稳定 snapshot/row 汇总和安全错误日志。 +6. 更新 FakeRadarSource、build/retry 测试和 CLI 组合测试,覆盖成功、6000 行、空分片、瞬时失败、错误日期/代码、重复键、分片触顶、两路 worker 与 checkpoint 重试。 +7. 更新必要的运维说明,明确当前 L 股票池、候选覆盖语义、同一 adapter 的 0.2 秒共享间隔以及不同定时任务不得重叠。 +8. 依次运行定向 pytest、Ruff format/lint、Pyright、完整 pytest,并由独立 Trellis check 代理核验规格和实现;修复所有本任务引入的问题后提交本地分支。 + +## 风险点与回滚检查 + +- `RequestCoordinator` 是共享模块,必须证明默认参数不改变 market-data 并发。 +- worker 完成顺序不能进入 publication source order 或 input hash。 +- 不能把普通 provider 异常与来源契约错误混为一类,也不能用空行伪造成功分片。 +- 端口签名变化必须同步所有 fake/replay 路径,完整测试前不得仅凭 source 单测判定完成。 +- 无数据库迁移;若验证失败,可按步骤分别回滚协调器启动槽和资金流分片逻辑。 + +## 验证结果 + +- `uv lock --check`、Ruff format/check、Pyright strict 全部通过。 +- 后端完整测试 `159 passed, 3 skipped`;跳过项均要求显式设置 `ZHIXING_TEST_DATABASE_URL`。 +- 根目录 `./dev.sh check` 与 `./dev.sh test` 通过;前端 `64 passed`。 +- 未执行真实 Tushare 账号并发调用与真实 PostgreSQL 集成测试,留待部署后的 capability/生产批次验证。 diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/prd.md b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/prd.md new file mode 100644 index 0000000..f133f22 --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/prd.md @@ -0,0 +1,40 @@ +# 资金雷达当前上市股票池与资金流缺口补拉 + +## Goal + +让板块资金雷达以“构建时当前上市股票”为唯一证券范围,并在 Tushare `moneyflow_dc` 单日响应触及 6000 行上限时,仍能安全验证和补齐雷达候选股票,而不是直接失败或接受可能截断的数据。 + +## Background + +- 生产构建目标日 `2026-08-28` 已在 `moneyflow_dc` 来源组因响应达到供应商 6000 行上限而失败。 +- Tushare 官方接口说明确认 `moneyflow_dc` 单次最多返回 6000 条,并支持按日期或股票代码循环提取。 +- 用户明确不要求历史时点证券生命周期还原;功能上线日视为最早历史日期,当前及未来均只研究构建时 `stock_basic(list_status=L)` 返回的股票。 +- Tushare 账号频率限制为 500 次/分钟;用户批准资金流缺口补拉使用 2 个 worker,但两个 worker必须共享同一请求启动限流器。 + +## Requirements + +1. `stock_basic` 只请求 `list_status=L`,不再请求 `D/P/G/UN`。保留资金雷达既有的沪深 A 股、B 股/北交所排除和上市日期校验,不额外引入 `market-data-sync` 的 ST 过滤语义。 +2. 资金流完整性只针对有效板块成员与当前 `L` 股票的交集。应用层必须在请求 `moneyflow_dc` 前形成稳定、去重的候选代码集合,并通过显式端口参数传给 source adapter,禁止依赖 adapter 内部调用顺序或缓存状态。 +3. `moneyflow_dc` 首次仍按目标交易日请求全市场。首次响应即使达到 6000 行,也必须先校验目标日期和业务唯一键,再检查候选股票覆盖率,不能直接接受或直接报截断。 +4. 首次响应缺少候选股票时,只按稳定排序后的缺失 `ts_code` 补拉。每个分片必须同时传入 `trade_date` 和 `ts_code`,并校验返回日期、返回代码、唯一键以及供应商是否忽略了分片参数。 +5. 缺口补拉固定使用 2 个 worker。同一 adapter 的所有首次请求、补拉请求及 retry 共享请求启动间隔,默认相邻请求启动至少间隔 0.2 秒;普通 provider 调用允许重叠,不得把整个请求放在协调器锁内。 +6. 空分片或重试耗尽的瞬时请求失败保留为真实缺口,不补零;构建继续走现有覆盖率逻辑并可发布 `partial`。日期错误、返回错误股票代码、重复业务键或分片再次触及供应商上限属于来源契约错误,必须 fail closed。 +7. 全市场首批 snapshot 与每个成功分片 snapshot 都属于现有 `MONEYFLOW_DC` source group,并按确定性顺序保存。重试继续复用已完成来源组,只刷新资金流来源组,不新增数据库表或 publication group。 +8. `RequestCoordinator` 的请求启动间隔默认关闭,只有资金雷达通过现有 `sector_radar_request_interval_seconds` 启用,不能改变 `market-data-sync` 当前八路并发语义。 + +## Acceptance Criteria + +- [x] 资金雷达构建只发出一次 `stock_basic(list_status=L)` 请求,并拒绝该分区返回非 `L` 状态。 +- [x] `moneyflow_dc` 首批低于或等于 6000 行且覆盖全部候选股票时均可成功解析;达到 6000 行本身不再导致 `SourceTruncatedError`。 +- [x] 首批未覆盖候选股票时,仅补拉缺失代码,调用总数为 `1 + 缺失代码数`,最终 snapshot 顺序与 worker 完成顺序无关。 +- [x] 两个补拉请求可以处于并发等待状态,但共享协调器记录的请求启动时间间隔不小于配置值;默认配置下理论总速率不超过约 300 次/分钟。 +- [x] 空补拉和瞬时请求失败不会被补零或伪装成完整覆盖;错误日期、错误代码、重复键和分片触顶会阻止发布错误结果。 +- [x] failed/partial publication 重试仍复用既有 source checkpoints,并只刷新需要重取的 `MONEYFLOW_DC` group。 +- [x] 共享协调器、sector radar source/build/CLI 相关单元测试、Ruff、Pyright 和完整后端 pytest 通过;需要真实 PostgreSQL 的测试若未配置,必须明确报告跳过状态。 + +## Out of Scope + +- 不保证历史日期按当时上市状态精确重建,也不保留已退市股票进入未来重跑结果。 +- 不复用 `market-data-sync` 的数据库股票池、行情或 Tushare client。 +- 不实现跨进程或跨定时任务的分布式限流;运维上仍要求 `market-data-sync` 与 `sector-radar-build` 不重叠运行。 +- 不改变板块评分公式、前端展示、数据库 schema 或其他 Tushare 来源组的请求策略。 diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/radar-build-tushare-call-chain.md b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/radar-build-tushare-call-chain.md new file mode 100644 index 0000000..8df6a8f --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/radar-build-tushare-call-chain.md @@ -0,0 +1,85 @@ +# Research: 板块资金雷达 build 到 Tushare source 调用链 + +- Query: 定位板块资金雷达从 application build 到 Tushare source 的完整调用链,解释 `stock_basic` 为什么请求 `L/D/P/G/UN`、`moneyflow_dc` 在哪里按 6000 行拒绝、候选股票集合何时形成,以及重试时 publication/source checkpoint 如何复用。 +- Scope: internal +- Date: 2026-08-31 + +## Findings + +### 1. 完整调用链 + +生产入口由 `zhixing-server/pyproject.toml:36-38` 将 `sector-radar-build` 绑定到 `presentation.cli:main`。CLI 在 `zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py:42-47` 构造 `BuildSectorRadarCommand`,在同文件 `:62-77` 用 token 创建 `TushareSectorRadarAdapter`、创建 PostgreSQL repository,并调用 `BuildSectorRadar(...).execute(command)`。 + +application 层从 `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:199-222` 的 `BuildSectorRadar.execute` 开始。普通单日/区间模式先由 `_resolve_targets` 调 `source.fetch_trade_calendar` 解析目标交易日(`:224-247`);retry 模式则直接读取原 publication 并复用其目标交易日(`:225-231`)。随后 `_build_target` 获取按交易日的 advisory lock(`:257-271`),`_build_locked` 恢复遗留 running publication、创建新的 running publication、加载可复用来源组,再进入 `_collect`(`:285-310`)。 + +`_collect` 以固定顺序调用 `_fetch_group`:calendar、concept indices、industry indices、members、stock basics、suspensions、daily、moneyflow_dc,见 `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:498-563`。application 依赖的 source port 定义在 `zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py:24-47`;生产实现是 `TushareSectorRadarAdapter`。每个 adapter 方法最终进入 `TushareSectorRadarAdapter._fetch_snapshot`,它组装显式 fields 后优先调用 `client.query(api_name, fields=..., **params)`,见 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:368-405`。因此资金流主链是 `BuildSectorRadar.execute -> _build_target -> _build_locked -> _collect -> _fetch_group(MONEYFLOW_DC) -> SectorRadarSource.fetch_moneyflow_dc -> TushareSectorRadarAdapter.fetch_moneyflow_dc -> _fetch_snapshot -> client.query("moneyflow_dc", trade_date=..., fields=...)`。 + +`_fetch_group` 不只是调用 source:新拉或重放成功后,它立即保存 content-addressed raw snapshot,并以 `(publication_id, source_group, source_order)` 建立 publication checkpoint 链接,见 `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:457-496`。这意味着失败发生前已完成的每一组都已经具备可恢复检查点。 + +### 2. `stock_basic` 为什么请求 `L/D/P/G/UN` + +adapter 明确说明不能依赖 Tushare 默认只返回 `L`,并逐一请求五个文档化生命周期分区,见 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:267-285`。每个响应还校验返回 `list_status` 必须与请求分区一致(`:280-282`),最后跨分区校验 `ts_code` 唯一(`:284`)。现有回归测试固定了五次请求顺序与五种状态均被汇总,见 `zhixing-server/tests/unit/sector_radar/test_tushare_source.py:308-333`。 + +业务原因是 radar 需要按目标交易日判断 point-in-time 生命周期,而不是只看“当前仍上市”的默认集合。`StockBasicRow` 保存 `list_date`/`delist_date`,见 `zhixing-server/src/zhixing_server/modules/sector_radar/domain/source.py:338-364`;真正的生命周期判断在 `zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py:200-210`,要求沪深 A 股、非 B 股/北交所、`list_date <= target` 且目标日不晚于 `delist_date`。因此完整状态分区主要用于避免历史目标日漏掉目前已退市/暂停等股票,并使未上市/过会等记录由日期规则明确排除。`list_status` 本身目前不直接决定资格,资格由代码、市场及上市/退市日期决定。 + +### 3. `moneyflow_dc` 的 6000 行拒绝点 + +`ROW_LIMITS["moneyflow_dc"]` 在 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:73-81` 固定为 `6_000`。`_fetch_snapshot` 把该上限传给 `build_source_snapshot`(同文件 `:397-405`);snapshot builder 在 `zhixing-server/src/zhixing_server/modules/sector_radar/domain/source.py:147-160` 计算 `row_count`,并以 `row_count >= row_limit` 标记 `limit_reached=True`,所以恰好返回 6000 行也视为可能截断。 + +具体拒绝发生在 `TushareSectorRadarAdapter.fetch_moneyflow_dc`:取到 snapshot 后立刻调用 `_reject_limit`,见 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:318-330`;`_reject_limit` 在同文件 `:465-468` 抛出 `SourceTruncatedError("moneyflow_dc reached its provider row limit")`。这里没有像 `dc_member` 那样的分区补拉逻辑;错误经 `_fetch_group` 和 `_build_locked` 上浮,最终 publication 被记为 failed(`application/build.py:391-421`)。 + +### 4. 候选股票集合形成时点 + +候选集不是在请求 `moneyflow_dc` 之前形成。`_collect` 先完成全部八个来源组,包括 full-market `daily` 与 `moneyflow_dc`(`zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:550-563`),然后才调用 `normalize_memberships`,从状态为 `AVAILABLE` 的板块成员记录中取非空 `stock_code`、去重并排序为 `candidate_codes`(`:565-574`)。随后 `candidate_codes` 才传入 `normalize_stock_facts`(`:575-582`)。 + +这个集合此时只是“当日概念/行业成员股票并集”,尚未完成生命周期过滤。`normalize_stock_facts` 在遍历候选代码时才逐只调用生命周期规则;不合法者被保留为 `LIFECYCLE_INVALID` fact,见 `zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py:159-170`。所以当前调用顺序无法用候选集缩小或分片本轮 `moneyflow_dc` 请求,这是本次“当前上市股票池与资金流缺口补拉”设计需要显式调整的结构性边界。 + +### 5. publication/source checkpoint 的重试复用 + +retry 命令只接受 `partial` 或 `failed` publication,并复用旧 publication 的目标交易日,见 `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:224-231`。实际重试不会继续写旧 publication,而是在持有日期锁后新建一个 running publication,再从旧 publication 加载 source checkpoints(`:285-310`)。 + +checkpoint 的领域模型是八个稳定的 `PublicationSourceGroup` 和有序 `PublicationSourceRecord`,后者持有 raw `SourceSnapshot` 与 `refresh_on_retry` 标志,见 `zhixing-server/src/zhixing_server/modules/sector_radar/domain/persistence.py:131-160`。数据库表以 `(publication_id, source_group, source_order)` 为主键,raw snapshot 外键采用 `RESTRICT`,见 `zhixing-server/migrations/versions/0005_radar_daily_aggregate.py:17-55`;repository 加载时 join raw snapshot 并按 group/order 返回,见 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py:192-222`。 + +`_reusable_sources` 会忽略 `refresh_on_retry=True` 的记录;其余记录按 `source_order` 排序,并要求编号从 0 连续,否则拒绝重放,见 `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:433-455`。`_fetch_group` 命中 reusable group 时不访问 Tushare,而是从保存的 snapshot rows 重新执行 typed parser;未命中时才调用 source。无论重放还是新拉,snapshot 都会再次链接到新的 running publication,见同文件 `:457-496`。 + +两类失败的复用语义不同。对于完整采集后因覆盖率不足形成的 partial,`_retry_source_groups` 根据 `membership_unknown`、缺失/空 `daily`、缺失/空 `moneyflow` 精确选择需刷新的来源组,见 `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:683-704`;`finalize_publication` 在同一事务中把这些组标记为 `refresh_on_retry=TRUE` 后再结束 publication,见 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py:562-571`。对于采集中途 hard failure,成功组已经由 `_fetch_group` 即时 checkpoint,失败组及其后的组没有记录;因此 retry 自动重放所有已完成组并从首个未完成组继续。测试证明 partial 资金缺口只再次调用 `moneyflow_dc`(`zhixing-server/tests/unit/sector_radar/test_build.py:389-406`),而 daily hard failure 后会复用此前六组,只再次调用 `daily` 和尚未执行的 `moneyflow_dc`(`:409-426`)。 + +收集完成后,所有 snapshot id 与指标版本参与 `input_hash`(`zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:706-716`);若已存在同目标日、同 input hash 的 publication,新 running publication 会被丢弃并返回 existing publication,见同文件 `:310-327`。这是 publication 级幂等复用,与 source-group 级断点重放互补。 + +## Files Found + +- `zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py`:生产 CLI 组合根,创建 Tushare adapter、PostgreSQL repository 与 application use case。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py`:目标日期解析、来源组采集顺序、候选集生成、publication 生命周期及重试复用核心。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py`:application 到 source adapter 的端口契约。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py`:七类 Tushare 接口、状态分区、行数上限与 `client.query` 边界。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/domain/source.py`:raw snapshot 的上限标记、typed row 解析与 source contract errors。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py`:成员候选并集之后的生命周期、停牌、行情与资金流事实归一化。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/domain/persistence.py`:source checkpoint 分组和值对象契约。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py`:checkpoint 的保存、加载与 partial 刷新标记事务。 +- `zhixing-server/migrations/versions/0005_radar_daily_aggregate.py`:publication-source checkpoint 表结构及完整性约束。 +- `zhixing-server/tests/unit/sector_radar/test_build.py`:partial/failed 重试只刷新未完成来源组的可执行证据。 +- `zhixing-server/tests/unit/sector_radar/test_tushare_source.py`:五种 `stock_basic` 状态分区请求的回归证据。 + +## Code Patterns + +- Port/adapter:application 只依赖 `SectorRadarSource`,CLI 注入 `TushareSectorRadarAdapter`(`domain/ports.py:24-47`;`presentation/cli.py:62-77`)。 +- Point-in-time master data:显式拉取所有生命周期状态,再按目标日 `list_date`/`delist_date` 判定(`infrastructure/tushare.py:267-285`;`domain/normalize.py:200-210`)。 +- Fail closed on provider limit:snapshot 以 `>=` 标记触顶,不能将潜在截断当成功(`domain/source.py:158-160`;`infrastructure/tushare.py:465-468`)。 +- Immediate source checkpoint:每个来源组一成功就保存 raw snapshot 及 publication link,而不是等待整个 publication 完成(`application/build.py:486-496`)。 +- Selective retry:partial 显式标记缺口组;failed 依赖已完成组存在、未完成组缺席来恢复(`application/build.py:433-475,683-704`)。 + +## External References + +- 本次为内部调用链研究,未新增外部资料检索。既有已归档研究 `.trellis/tasks/archive/2026-08/08-28-sector-capital-radar/research/tushare-radar-contract.md:7-15` 记录了原实现采用的 Tushare 接口边界:`stock_basic` 默认只返回 `L`,`moneyflow_dc` 单次上限 6000;上线前仍应以目标账号 capability probe 和当时官方文档为准。 + +## Related Specs + +- `.trellis/spec/backend/market-data-sync.md`:一次性 Tushare Job、可恢复 snapshot、失败保留旧发布的相邻上下文规范;sector radar 有独立 bounded context,不能直接套用选股/ST 股票池规则。 +- `.trellis/spec/backend/selection.md`:selection 的当前沪深非 ST 股票池契约不等于 radar 的 point-in-time 板块成员 universe。 +- `.trellis/tasks/archive/2026-08/08-28-sector-capital-radar/design.md:36-44`:原 radar 设计要求全部上市状态、目标日生命周期、沪深 A 股过滤及行数触顶时不得接受截断响应。 + +## Caveats / Not Found + +- 当前 `moneyflow_dc` 没有按候选股票或代码分片的实现;达到 6000 行只会 hard fail。`dc_member` 有按板块分区补拉,可作为模式参考,但不能直接证明 Tushare `moneyflow_dc` 支持同样的参数或批量行为。 +- 当前候选集形成得晚于 `moneyflow_dc` 请求,并且成员并集与“生命周期有效股票池”是两个阶段;讨论修复时必须明确要前移哪一个集合,避免误把所有板块成员都视为当前上市股票。 +- 代码中的 `ROW_LIMITS` 是本地契约常量,不是运行时从供应商元数据发现;若要改变请求策略,需要重新核对当前 Tushare `moneyflow_dc` 的可用过滤参数、单次限制及积分权限。 diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/tushare-global-start-interval.md b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/tushare-global-start-interval.md new file mode 100644 index 0000000..95bff71 --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/research/tushare-global-start-interval.md @@ -0,0 +1,92 @@ +# Research: Tushare 两路补拉的全局请求启动间隔 + +- Query: 检索仓库现有 Tushare 限流、并发 worker、线程安全、测试夹具与 `sector_radar` source 测试模式,定位实现“2 个 worker 共享全局 0.2 秒请求启动间隔”的最小适配点。 +- Scope: internal +- Date: 2026-08-31 + +## Findings + +### 结论与最小适配面 + +最小且边界清晰的实现是扩展共享的 `RequestCoordinator`,让它可选地协调“请求启动槽”,然后仅让 `TushareSectorRadarAdapter` 启用现有的 `request_interval_seconds=0.2`。两路 `moneyflow_dc` 补拉 worker 共享同一个 adapter,而该 adapter 已经只持有一个 `_coordinator`,因此不需要新建进程级 singleton,也不需要新增环境变量。 + +具体适配点如下。 + +1. 在 `zhixing-server/src/zhixing_server/shared/request_coordinator.py:35-66` 的 `RequestCoordinator` 增加默认关闭的启动间隔参数及 `_next_request_at` 状态;继续复用现有 `threading.Condition`,在同一临界区内读取单调时钟、计算 `max(_cooldown_until, _next_request_at)`、等待并预约下一启动时刻。只有“预约”需要持锁,真实 provider 请求必须在锁外执行,才能保持两个 worker 的请求重叠能力。 +2. 在 `zhixing-server/src/zhixing_server/shared/request_coordinator.py:75-82` 的每次 attempt 开始前,把当前只等待 cooldown 的 `_wait_for_cooldown` 收敛成“等待 cooldown 并原子预约启动槽”。预约完成时令 `_next_request_at = actual_start + interval`。这样初次请求和 retry 都服从同一个启动间隔;当 403/429 创建 cooldown 后,等待中的 worker 还会在醒来时重新检查 cooldown。 +3. 新参数必须默认 `0.0`。`RequestCoordinator` 还被 market-data 使用,而且它的现有契约明确是“普通请求不串行,只共享命中限流后的 cooldown”(`zhixing-server/src/zhixing_server/shared/request_coordinator.py:35-41`;`docs/market-data-sync.md:65-69`)。默认关闭可避免顺带改变八路行情同步语义。 +4. 在 `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:92-116` 构造默认 coordinator 时,把已经存在的 `request_interval_seconds` 传入协调器;删除或停用 `_fetch_snapshot` 成功返回后的逐线程休眠(`zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:368-389`)。当前休眠发生在请求完成后,两个线程可以同时启动,不能表达“全局请求启动间隔”。 +5. `request_interval_seconds` 的配置链已经完整:`Settings.sector_radar_request_interval_seconds` 默认 0.2(`zhixing-server/src/zhixing_server/bootstrap/config.py:27-31`),CLI 将它传给 `TushareSectorRadarAdapter.from_token`(`zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py:62-67`)。因此不需要改 `.env`、Compose 或配置模型。 + +这里的“全局”只能可靠地解释为“同一 adapter/coordinator 实例覆盖的两个 worker”。现有 coordinator 不是模块 singleton,也不能跨进程协调;market-data job、sector-radar job 或两个独立进程各自创建 coordinator。若需求是全系统或跨进程的 5 requests/s,则本方案不满足,需要外部/分布式限流器,这会明显扩大范围。 + +### 两个 worker 的落点 + +`moneyflow_dc` 的当前全市场入口完全串行,只请求一次并校验结果(`zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:318-330`)。补拉属于 Tushare 响应分区细节,最小落点是该 infrastructure adapter 内部:先保留全市场快照,再对缺失的当前上市股票代码使用 `ThreadPoolExecutor(max_workers=2)` 发起按 `ts_code` 分区请求。worker 共享 `self._coordinator`,所以每一个 `_fetch_snapshot` 最终都经过同一个启动槽。 + +仓库已有 worker 写法可复用:`SyncMarketData` 在 `zhixing-server/src/zhixing_server/modules/market_data/application/sync.py:342-361` 使用具名的 `ThreadPoolExecutor` 和 future-to-business-key 映射;并发测试用带 `threading.Lock` 的 fake 统计 active/max-active(`zhixing-server/tests/unit/market_data/test_sync_concurrency.py:22-59`),并断言两路上限(`zhixing-server/tests/unit/market_data/test_sync_concurrency.py:156-180`)。sector-radar 不宜照搬其数据库副作用模型,只应复用“有界 executor + 主线程汇总”的形状。 + +如果补拉需要由当前上市股票池驱动,应用层已经先得到 `stock_basics`、后取 `moneyflow`(`zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:536-563`)。最小跨层契约是从 `stock_basics.rows` 中取 `list_status == "L"` 的代码并传给 `fetch_moneyflow_dc`;这会同步影响 `SectorRadarSource.fetch_moneyflow_dc`(`zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py:24-47`)和测试 fake。不要让 adapter 缓存上一次 `fetch_stock_basics` 的结果,否则 retry/replay 和调用顺序会形成隐式状态。 + +worker 完成顺序不得直接决定 snapshot 顺序。publication checkpoint 会按 `result.snapshots` 的枚举顺序持久化(`zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:486-495`),重放又要求 `source_order` 从 0 连续并按序恢复(`zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py:433-455`)。因此 future 结果应按输入股票代码或明确排序后汇总;虽然单个 snapshot 内部的 hash 已对 rows 做顺序稳定化(`zhixing-server/src/zhixing_server/modules/sector_radar/domain/source.py:134-163`),snapshot 元组自身仍需稳定。 + +固定“两路”不要求新增 `Settings`。最小做法是在 adapter 内使用命名常量或默认值为 2 的构造参数;只有产品要求运行时可调时,才需要扩展 config、CLI、`.env.example` 和 Compose。仓库现有可调 worker 的完整链路可参考 `Settings.market_data_max_workers`(`zhixing-server/src/zhixing_server/bootstrap/config.py:20-25`)和 `SyncMarketData(max_workers=...)`(`zhixing-server/src/zhixing_server/modules/market_data/application/sync.py:131-143`)。 + +### 现有限流与线程安全证据 + +共享协调器已经用 `threading.Condition` 保护 `_cooldown_until` 和 `_rate_limit_count`(`zhixing-server/src/zhixing_server/shared/request_coordinator.py:43-66`),读取、创建 cooldown 和成功后清理也都在该条件锁内(同文件 `:68-73`、`:127-152`)。403、429 和稳定中英文提示的分类位于同文件 `:13-28`、`:154-163`,cooldown 阶梯为 60/120/180 秒(`:13`)。这正是承载全 worker 启动槽的现有线程安全原语。 + +当前 `call` 明确允许普通请求并发(`zhixing-server/src/zhixing_server/shared/request_coordinator.py:35-41`),而 `_wait_for_cooldown` 在锁外调用 `wait_fn`(`:127-138`),不会把 provider 调用包在全局锁中。新增启动间隔也应保持这一点;如果把整次 `client.query` 放进锁中,虽然间隔成立,但会把两个 worker 退化为串行请求。 + +`TushareSectorRadarAdapter` 对同一个 SDK client 调用 `client.query` 或接口方法(`zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py:378-385`)。仓库内没有 Tushare SDK client 线程安全保证,也没有为 `_client` 加锁。锁定版本是 Tushare 1.4.29(`zhixing-server/uv.lock:731-739`)。因此两路并发是否可共享同一 SDK client 是实现前仍需确认的风险;启动间隔只保护频率状态,不等于保证 SDK client 内部线程安全。若无法确认,选择独立 client 会需要 token/factory 生命周期改造,选择锁住整个 client 请求则无法获得网络调用并发收益。 + +### 测试模式与建议入口 + +仓库没有 `tests/**/conftest.py` 或 sector-radar pytest fixture。`zhixing-server/tests/unit/sector_radar/test_tushare_source.py:23-44` 采用文件内 `QueryClient` 和 `make_adapter`:响应按 `(api_name, ts_code/list_status)` 分区,adapter 关闭 retry 和真实 sleep,并固定 `now_fn`。`dc_member` 达上限后按分区补拉的测试(同文件 `:216-274`)是 moneyflow 缺失补拉最接近的现有测试模板;当前上市状态查询测试在 `:308-333`。 + +启动间隔的直接先例是 `zhixing-server/tests/unit/market_data/test_tushare.py:53-87`:用可注入 fake monotonic clock 和 wait 函数验证一个请求触发的 cooldown 会阻塞后续请求。新增测试应延续 fake clock,而不是用真实 `sleep(0.2)` 和宽松 wall-clock 断言,以避免并发测试抖动。 + +建议最少覆盖两层行为: + +- 在共享 coordinator 的单元测试中让两个线程共享一个 coordinator,用 `threading.Event` 保持首个 request 未完成,fake clock/wait 将第二个启动推进到 0.2;记录两个真实 request callback 的开始时刻并断言差值为 0.2。该形状能证明“请求可重叠,但启动槽不重叠”,也能避免单纯顺序调用掩盖线程竞态。 +- 在 `test_tushare_source.py` 增加 moneyflow 缺失回补测试:全市场响应遗漏若干 `list_status=L` 代码,按 `ts_code` 的 fake 分区返回补拉结果,断言 executor 最大 active 不超过 2、最终 rows 和 snapshots 顺序稳定、非上市状态不补拉。并发 fake 的 `calls`、`active` 和 `max_active` 必须用 `threading.Lock`;现有 `QueryClient.calls.append`(`:23-34`)只适合串行测试。 + +现有测试入口为: + +```bash +cd zhixing-server +uv run pytest tests/unit/market_data/test_tushare.py +uv run pytest tests/unit/sector_radar/test_tushare_source.py +uv run pytest tests/unit/sector_radar/test_build.py +uv run pytest tests/unit/sector_radar/test_cli.py +``` + +共享 coordinator 改动至少应运行前两个入口;若 `fetch_moneyflow_dc` 端口增加上市代码参数,还必须运行后两个入口以覆盖 `FakeRadarSource`、应用编排和 CLI 组合。完整后端门禁由 `.trellis/spec/backend/quality-guidelines.md:3-20` 和 `zhixing-server/pyproject.toml:40-49` 定义,包括 Ruff format/lint、Pyright strict 和完整 pytest。 + +### Files found + +- `zhixing-server/src/zhixing_server/shared/request_coordinator.py`:跨 bounded context 的 retry、限流识别和共享 cooldown 协调器,是全局启动槽的最小所有权位置。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py`:sector-radar Tushare adapter、source 分区与当前逐请求休眠位置。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py`:stock basics 到 moneyflow 的调用顺序、source checkpoint 稳定顺序契约。 +- `zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py`:`fetch_moneyflow_dc` 的应用端口签名。 +- `zhixing-server/src/zhixing_server/modules/market_data/application/sync.py`:仓库现有有界 `ThreadPoolExecutor` 模式。 +- `zhixing-server/tests/unit/market_data/test_tushare.py`:fake clock/wait 的 coordinator 测试模式。 +- `zhixing-server/tests/unit/market_data/test_sync_concurrency.py`:两路 worker 上限与加锁 fake 的测试模式。 +- `zhixing-server/tests/unit/sector_radar/test_tushare_source.py`:source fake、分区响应、禁用真实 sleep 及 schema/limit 测试入口。 +- `zhixing-server/tests/unit/sector_radar/test_build.py`:应用端口 fake 和 source-group replay/retry 覆盖。 +- `zhixing-server/src/zhixing_server/bootstrap/config.py`、`zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py`:现有 0.2 秒配置传递链。 + +### Related specs + +- `.trellis/spec/backend/directory-structure.md`:无业务归属的小型跨上下文能力应放在 `shared/`;Tushare 请求启动协调符合这一边界。 +- `.trellis/spec/backend/market-data-sync.md`:Tushare client、共享限流和后端测试门禁的既有契约。 +- `.trellis/spec/backend/configuration-and-runtime.md`:运行时配置只能通过 `Settings` 注入;本最小方案复用既有配置,不新增环境读取。 +- `.trellis/spec/backend/quality-guidelines.md`:Pyright strict、pytest 严格模式和后端质量命令。 +- `.trellis/spec/guides/code-reuse-thinking-guide.md`:跨上下文且无业务所有权的原语才进入 `shared/`,并要求复用前先核对生命周期与错误语义。 + +## Caveats / Not Found + +- 未在仓库中找到 Tushare 1.4.29 对 `pro_api` client 的线程安全声明;不能仅凭 Python 对 `list.append` 或对象读取的实现细节宣称 SDK client 可安全并发。 +- 未找到 sector-radar 专用 `conftest.py`、pytest fixture 或现成的 request-start 间隔测试;需要沿用文件内 fake 和 coordinator fake clock 模式。 +- 当前 PRD 仍为 TBD,未定义“全局”是否跨 adapter/进程,也未定义单只股票补拉失败是整组失败还是保留部分回补。以上结论按“一个 sector-radar adapter 内两路 worker、任一补拉失败则 source group 失败”的最小解释给出。 +- 本次只读研究未运行 pytest;研究代理只写入本文件,未修改产品代码或测试代码。 diff --git a/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/task.json b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/task.json new file mode 100644 index 0000000..a507d54 --- /dev/null +++ b/.trellis/tasks/08-31-sector-radar-listed-moneyflow-recovery/task.json @@ -0,0 +1,26 @@ +{ + "id": "sector-radar-listed-moneyflow-recovery", + "name": "sector-radar-listed-moneyflow-recovery", + "title": "资金雷达当前上市股票池与资金流缺口补拉", + "description": "", + "status": "in_progress", + "dev_type": null, + "scope": null, + "package": null, + "priority": "P1", + "creator": "yuxuanhui", + "assignee": "yuxuanhui", + "createdAt": "2026-08-31", + "completedAt": null, + "branch": null, + "base_branch": "develop", + "worktree_path": null, + "commit": null, + "pr_url": null, + "subtasks": [], + "children": [], + "parent": null, + "relatedFiles": [], + "notes": "", + "meta": {} +} \ No newline at end of file diff --git a/docs/market-data-sync.md b/docs/market-data-sync.md index 6567d9c..22ad10c 100644 --- a/docs/market-data-sync.md +++ b/docs/market-data-sync.md @@ -98,7 +98,7 @@ docker compose -f docker-compose.prod.yml --profile jobs config ## 板块资金雷达 Job -`sector-radar-build` 同样是外部调度器触发的一次性任务,FastAPI 不会在进程内启动定时器。它只读取 Tushare 的 `trade_cal`、`dc_index`、`dc_member`、`stock_basic`、`suspend_d`、`daily` 和 `moneyflow_dc`,保存 point-in-time 原始快照与规范化事实,再生成明确标注为“知行独立实现”的版本化指标。生产运行时不请求 OneChartLab。 +`sector-radar-build` 同样是外部调度器触发的一次性任务,FastAPI 不会在进程内启动定时器。它只读取 Tushare 的 `trade_cal`、`dc_index`、`dc_member`、`stock_basic`、`suspend_d`、`daily` 和 `moneyflow_dc`,保存 point-in-time 原始快照与规范化事实,再生成明确标注为“知行独立实现”的版本化指标。生产运行时不请求 OneChartLab。雷达股票范围固定为构建时 `stock_basic(list_status=L)` 返回的沪深 A 股与有效板块成员的交集;历史回填也采用构建时当前上市股票池,不还原目标日当时已经退市的证券。 开发环境没有 token 时可以检查命令契约,但不能执行真实构建: @@ -120,4 +120,6 @@ docker compose -f docker-compose.prod.yml --profile jobs run --rm sector-radar-b --retry-publication-id ``` +`moneyflow_dc` 先按交易日拉取全市场快照;即使首批达到 6000 行,也会根据上述候选股票检查实际覆盖,并用固定两路 worker 逐只补拉缺失代码。全市场请求、补拉和 retry 在同一 adapter 内共享 `ZHIXING_SECTOR_RADAR_REQUEST_INTERVAL_SECONDS`(默认 0.2 秒)的请求启动间隔;空分片或普通请求重试耗尽会保留为覆盖缺口并形成 `partial`,来源返回错误日期、错误代码、重复键或分片再次触顶则整次构建失败。该限流只在单进程 adapter 内生效,生产调度仍不得让 `market-data-sync` 与 `sector-radar-build` 重叠运行。 + 重复输入通过内容 hash 复用已有成功发布,不产生无意义修订;同一目标日由 PostgreSQL advisory lock 阻止并发构建。`success` 或 `unchanged` 返回 0,覆盖率不足的 `partial` 返回 2,输入、上游、锁或基础设施失败返回 1。`partial`/`failed` 会保留审计,但读取端只选择 `success` 作为 last-good。当前版本只提供手工和外部调度入口,不新增生产 Cron;待真实账号 capability、到达时点和首轮回填验证完成后再单独启用调度。 diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py b/zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py index d784704..a53a1ef 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/application/build.py @@ -33,7 +33,11 @@ from ..domain.models import ( StockDailyFact, StockFactStatus, ) -from ..domain.normalize import normalize_memberships, normalize_stock_facts +from ..domain.normalize import ( + is_current_listed_stock, + normalize_memberships, + normalize_stock_facts, +) from ..domain.persistence import ( DailyAggregateRecord, MembershipRecord, @@ -461,18 +465,20 @@ class BuildSectorRadar: reusable: Mapping[PublicationSourceGroup, tuple[SourceSnapshot, ...]], fetch: Callable[[], SourceResult[T]], parser: Callable[[Mapping[str, SourceScalar]], T], + reuse_if: Callable[[SourceResult[T]], bool] | None = None, ) -> SourceResult[T]: - """Replay a completed group or fetch and checkpoint it immediately.""" + """Replay a compatible completed group or fetch and checkpoint it immediately.""" try: snapshots = reusable.get(source_group) if snapshots is None: result = fetch() else: - result = SourceResult( + replayed = SourceResult( snapshots=snapshots, rows=tuple(parser(row) for snapshot in snapshots for row in snapshot.rows), ) + result = replayed if reuse_if is None or reuse_if(replayed) else fetch() except SourceContractError as exc: if exc.claim_diagnostic(): logger.error( @@ -540,6 +546,22 @@ class BuildSectorRadar: self.source.fetch_stock_basics, StockBasicRow.from_mapping, ) + memberships = normalize_memberships(indices, members) + member_codes = tuple( + sorted( + { + item.stock_code + for item in memberships + if item.status is MembershipStatus.AVAILABLE and item.stock_code is not None + } + ) + ) + current_listed_codes = { + row.ts_code for row in stock_basics.rows if is_current_listed_stock(row, target) + } + moneyflow_candidate_codes = tuple( + code for code in member_codes if code in current_listed_codes + ) suspensions = self._fetch_group( publication_id, PublicationSourceGroup.SUSPENSIONS, @@ -558,23 +580,16 @@ class BuildSectorRadar: publication_id, PublicationSourceGroup.MONEYFLOW_DC, reusable, - lambda: self.source.fetch_moneyflow_dc(target), + lambda: self.source.fetch_moneyflow_dc(target, moneyflow_candidate_codes), MoneyflowDcRow.from_mapping, + reuse_if=lambda result: set(moneyflow_candidate_codes).issubset( + {row.ts_code for row in result.rows} + ), ) - memberships = normalize_memberships(indices, members) - candidate_codes = tuple( - sorted( - { - item.stock_code - for item in memberships - if item.status is MembershipStatus.AVAILABLE and item.stock_code is not None - } - ) - ) stock_facts = normalize_stock_facts( target_trade_date=target, - candidate_codes=candidate_codes, + candidate_codes=member_codes, stock_basics=stock_basics, suspensions=suspensions, daily=daily, diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py index 0c4dda7..69a16d7 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/normalize.py @@ -117,7 +117,7 @@ def normalize_stock_facts( Args: target_trade_date: Date whose point-in-time lifecycle is evaluated. candidate_codes: Union of stocks in that date's sector memberships. - stock_basics: All explicit Tushare listing-status partitions. + stock_basics: Current ``list_status=L`` Tushare listings. suspensions: Same-date suspend/resume events. daily: Same-date stock turnover rows in source units. moneyflow: Same-date DC main-moneyflow rows in source units. @@ -165,7 +165,7 @@ def normalize_stock_facts( turnover_yuan = None net_amount_yuan = None - if basic is None or not _is_lifecycle_candidate(basic, target_trade_date): + if basic is None or not is_current_listed_stock(basic, target_trade_date): status = StockFactStatus.LIFECYCLE_INVALID elif ts_code in suspended_codes and daily_row is None: status = StockFactStatus.SUSPENDED @@ -197,7 +197,23 @@ def normalize_stock_facts( return tuple(records) -def _is_lifecycle_candidate(stock: StockBasicRow, target: date) -> bool: +def is_current_listed_stock(stock: StockBasicRow, target: date) -> bool: + """Return whether one current ``L`` row is an eligible radar security. + + The radar intentionally uses the listings observed at build time rather than + reconstructing historical delistings. Code, market, and list-date checks keep + the existing Shanghai/Shenzhen A-share boundary intact. + + Args: + stock: One validated ``stock_basic`` row. + target: Radar date whose list date must already have arrived. + + Returns: + Whether the security belongs to the build-time radar universe. + """ + + if stock.list_status != "L": + return False if not stock.ts_code.endswith((".SH", ".SZ")): return False if stock.symbol.startswith(("200", "900")): @@ -205,9 +221,7 @@ def _is_lifecycle_candidate(stock: StockBasicRow, target: date) -> bool: market = stock.market or "" if "北交" in market or "B股" in market.upper(): return False - if stock.list_date is None or stock.list_date > target: - return False - return stock.delist_date is None or target <= stock.delist_date + return stock.list_date is not None and stock.list_date <= target def _is_suspend_event(value: str) -> bool: diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py index 06f8ce6..9206963 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ports.py @@ -42,7 +42,11 @@ class SectorRadarSource(Protocol): def fetch_daily(self, trade_date: date) -> SourceResult[DailyRow]: ... - def fetch_moneyflow_dc(self, trade_date: date) -> SourceResult[MoneyflowDcRow]: ... + def fetch_moneyflow_dc( + self, + trade_date: date, + candidate_codes: Sequence[str], + ) -> SourceResult[MoneyflowDcRow]: ... def probe(self, trade_date: date) -> CapabilityProbeResult: ... diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py b/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py index e9dca18..c2a8dc3 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/tushare.py @@ -5,12 +5,14 @@ from __future__ import annotations import logging import time from collections.abc import Callable, Iterable, Mapping, Sequence +from concurrent.futures import ThreadPoolExecutor from datetime import UTC, date, datetime from typing import TypeVar, cast from zhixing_server.shared.request_coordinator import ( DEFAULT_RATE_LIMIT_COOLDOWNS, RequestCoordinator, + TushareSourceError, ) from ..domain.models import SectorType @@ -84,6 +86,7 @@ _SECTOR_TYPE_PARAM = { SectorType.CONCEPT: "概念板块", SectorType.INDUSTRY: "行业板块", } +_MONEYFLOW_WORKERS = 2 class TushareSectorRadarAdapter: @@ -104,12 +107,11 @@ class TushareSectorRadarAdapter: """Create an adapter around one already-authenticated SDK client.""" self._client = client - self._sleep_fn = sleep_fn - self._request_interval_seconds = max(0.0, request_interval_seconds) self._now_fn = now_fn self._coordinator = request_coordinator or RequestCoordinator( max_retries=max_retries, backoff_seconds=backoff_seconds, + request_interval_seconds=request_interval_seconds, cooldown_seconds=cooldown_seconds, wait_fn=sleep_fn, sleep_fn=sleep_fn, @@ -265,24 +267,19 @@ class TushareSectorRadarAdapter: ) def fetch_stock_basics(self) -> SourceResult[StockBasicRow]: - """Fetch every documented listing status instead of relying on the L default.""" + """Fetch the build-time current ``L`` listings in one explicit partition.""" - snapshots: list[SourceSnapshot] = [] - rows: list[StockBasicRow] = [] - for status in ("L", "D", "P", "G", "UN"): - snapshot = self._fetch_snapshot( - "stock_basic", - {"exchange": "", "list_status": status}, - target_trade_date=None, - partition_key=status, - ) - snapshots.append(snapshot) - parsed = tuple(StockBasicRow.from_mapping(row) for row in snapshot.rows) - if any(row.list_status != status for row in parsed): - raise SourceContractError("stock_basic returned an unexpected list_status") - rows.extend(parsed) + snapshot = self._fetch_snapshot( + "stock_basic", + {"exchange": "", "list_status": "L"}, + target_trade_date=None, + partition_key="L", + ) + rows = tuple(StockBasicRow.from_mapping(row) for row in snapshot.rows) + if any(row.list_status != "L" for row in rows): + raise SourceContractError("stock_basic returned an unexpected list_status") self._require_unique(rows, key=lambda row: row.ts_code, api_name="stock_basic") - return SourceResult(tuple(snapshots), tuple(sorted(rows, key=lambda row: row.ts_code))) + return SourceResult((snapshot,), tuple(sorted(rows, key=lambda row: row.ts_code))) def fetch_suspensions(self, trade_date: date) -> SourceResult[SuspendRow]: """Fetch explicit suspend/resume events for one date.""" @@ -315,19 +312,121 @@ class TushareSectorRadarAdapter: self._require_unique(rows, key=lambda row: row.ts_code, api_name="daily") return SourceResult((snapshot,), tuple(sorted(rows, key=lambda row: row.ts_code))) - def fetch_moneyflow_dc(self, trade_date: date) -> SourceResult[MoneyflowDcRow]: - """Fetch a full-market DC moneyflow snapshot in its documented source unit.""" + def fetch_moneyflow_dc( + self, + trade_date: date, + candidate_codes: Sequence[str], + ) -> SourceResult[MoneyflowDcRow]: + """Fetch full-market moneyflow and refill uncovered current candidates.""" - snapshot = self._fetch_snapshot( + expected_codes = tuple(sorted(set(candidate_codes))) + if tuple(candidate_codes) != expected_codes or any( + not code.strip() for code in expected_codes + ): + raise ValueError("candidate_codes must be sorted unique non-empty values") + initial = self._fetch_snapshot( "moneyflow_dc", {"trade_date": trade_date.strftime("%Y%m%d")}, target_trade_date=trade_date, + partition_key="all", ) - self._reject_limit(snapshot) - rows = tuple(MoneyflowDcRow.from_mapping(row) for row in snapshot.rows) - self._require_target_date(rows, trade_date, "moneyflow_dc") - self._require_unique(rows, key=lambda row: row.ts_code, api_name="moneyflow_dc") - return SourceResult((snapshot,), tuple(sorted(rows, key=lambda row: row.ts_code))) + try: + initial_rows = tuple(MoneyflowDcRow.from_mapping(row) for row in initial.rows) + self._require_target_date(initial_rows, trade_date, "moneyflow_dc") + self._require_unique( + initial_rows, + key=lambda row: (row.trade_date, row.ts_code), + api_name="moneyflow_dc", + ) + except SourceContractError as exc: + self._log_contract_failure("moneyflow_dc", "all", exc) + raise + + returned_codes = {row.ts_code for row in initial_rows} + missing_codes = tuple(code for code in expected_codes if code not in returned_codes) + if not missing_codes: + return SourceResult( + (initial,), + tuple(sorted(initial_rows, key=lambda row: row.ts_code)), + ) + + with ThreadPoolExecutor( + max_workers=_MONEYFLOW_WORKERS, + thread_name_prefix="sector-radar-moneyflow", + ) as executor: + futures = { + code: executor.submit(self._fetch_moneyflow_partition, trade_date, code) + for code in missing_codes + } + partition_results = tuple(futures[code].result() for code in missing_codes) + + snapshots = [initial] + merged_rows = list(initial_rows) + for result in partition_results: + if result is None: + continue + snapshot, rows = result + snapshots.append(snapshot) + merged_rows.extend(rows) + try: + self._require_unique( + merged_rows, + key=lambda row: (row.trade_date, row.ts_code), + api_name="moneyflow_dc", + ) + except SourceContractError as exc: + self._log_contract_failure("moneyflow_dc", "merged", exc) + raise + return SourceResult( + tuple(snapshots), + tuple(sorted(merged_rows, key=lambda row: row.ts_code)), + ) + + def _fetch_moneyflow_partition( + self, + trade_date: date, + ts_code: str, + ) -> tuple[SourceSnapshot, tuple[MoneyflowDcRow, ...]] | None: + """Return one validated refill partition or preserve an ordinary gap.""" + + try: + snapshot = self._fetch_snapshot( + "moneyflow_dc", + { + "trade_date": trade_date.strftime("%Y%m%d"), + "ts_code": ts_code, + }, + target_trade_date=trade_date, + partition_key=ts_code, + ) + except TushareSourceError: + logger.warning( + "sector_radar_moneyflow_partition_failed partition_key=%s error_type=%s", + self._safe_partition_key(ts_code), + TushareSourceError.__name__, + ) + return None + if not snapshot.rows: + logger.warning( + "sector_radar_moneyflow_partition_empty partition_key=%s", + self._safe_partition_key(ts_code), + ) + return None + try: + self._reject_limit(snapshot) + rows = tuple(MoneyflowDcRow.from_mapping(row) for row in snapshot.rows) + self._require_target_date(rows, trade_date, "moneyflow_dc") + self._require_unique( + rows, + key=lambda row: (row.trade_date, row.ts_code), + api_name="moneyflow_dc", + ) + if any(row.ts_code != ts_code for row in rows): + raise SourceContractError("moneyflow_dc partition returned a different ts_code") + except SourceContractError as exc: + self._log_contract_failure("moneyflow_dc", ts_code, exc) + raise + return snapshot, rows def probe(self, trade_date: date) -> CapabilityProbeResult: """Probe required interfaces while returning only safe classifications.""" @@ -360,7 +459,7 @@ class TushareSectorRadarAdapter: ("stock_basic", self.fetch_stock_basics), ("suspend_d", lambda: self.fetch_suspensions(trade_date)), ("daily", lambda: self.fetch_daily(trade_date)), - ("moneyflow_dc", lambda: self.fetch_moneyflow_dc(trade_date)), + ("moneyflow_dc", lambda: self.fetch_moneyflow_dc(trade_date, ())), ): results.append(self._probe_call(api_name, operation)[0]) return CapabilityProbeResult(observed_at=self._now_fn(), interfaces=tuple(results)) @@ -385,7 +484,6 @@ class TushareSectorRadarAdapter: return method(fields=fields, **params) result = self._coordinator.call(api_name, request) - self._sleep_fn(self._request_interval_seconds) try: columns = getattr(result, "columns", None) returned_fields = ( diff --git a/zhixing-server/src/zhixing_server/shared/request_coordinator.py b/zhixing-server/src/zhixing_server/shared/request_coordinator.py index 5491d23..3e5f28c 100644 --- a/zhixing-server/src/zhixing_server/shared/request_coordinator.py +++ b/zhixing-server/src/zhixing_server/shared/request_coordinator.py @@ -33,11 +33,12 @@ class TushareSourceError(RuntimeError): class RequestCoordinator: - """Coordinate retries and shared rate-limit cooling for one provider client. + """Coordinate retries, rate-limit cooling, and optional request start spacing. - Normal requests are not serialized. Only a classified provider limit creates - a shared cooldown. Injectable time functions keep long cooldowns deterministic - in tests without coupling the coordinator to any business bounded context. + Provider calls execute outside the coordinator lock and may overlap. When a + positive request interval is configured, only their start times are serialized. + Injectable time functions keep waits deterministic in tests without coupling + the coordinator to any business bounded context. """ def __init__( @@ -45,6 +46,7 @@ class RequestCoordinator: *, max_retries: int = 3, backoff_seconds: float = 1.0, + request_interval_seconds: float = 0.0, cooldown_seconds: Sequence[float] = DEFAULT_RATE_LIMIT_COOLDOWNS, random_fn: Callable[[], float] = random.random, clock: Callable[[], float] = time.monotonic, @@ -56,6 +58,7 @@ class RequestCoordinator: raise ValueError("cooldown_seconds must contain non-negative values") self.max_retries = max(0, max_retries) self.backoff_seconds = max(0.0, backoff_seconds) + self.request_interval_seconds = max(0.0, request_interval_seconds) self.cooldown_seconds = cooldowns self.random_fn = random_fn self.clock = clock @@ -63,6 +66,7 @@ class RequestCoordinator: self.sleep_fn = sleep_fn or wait_fn self._condition = threading.Condition() self._cooldown_until = 0.0 + self._next_request_start = 0.0 self._rate_limit_count = 0 @property @@ -77,7 +81,7 @@ class RequestCoordinator: last_error: BaseException | None = None for attempt in range(self.max_retries + 1): - self._wait_for_cooldown(method_name) + self._wait_for_request_start(method_name) try: result = request() except Exception as exc: @@ -124,17 +128,29 @@ class RequestCoordinator: return self.call(method_name, operation) - def _wait_for_cooldown(self, method_name: str) -> None: + def _wait_for_request_start(self, method_name: str) -> None: + """Reserve one start slot after both shared wait deadlines have elapsed.""" + while True: with self._condition: - delay = self._cooldown_until - self.clock() - if delay <= 0: - return - logger.info( - "provider_rate_limit_wait method=%s wait_seconds=%.1f", - method_name, - delay, - ) + now = self.clock() + start_at = max(self._cooldown_until, self._next_request_start) + delay = start_at - now + if delay <= 0: + self._next_request_start = now + self.request_interval_seconds + return + if start_at == self._cooldown_until: + logger.info( + "provider_rate_limit_wait method=%s wait_seconds=%.1f", + method_name, + delay, + ) + else: + logger.debug( + "provider_request_interval_wait method=%s wait_seconds=%.3f", + method_name, + delay, + ) self.wait_fn(delay) def _set_rate_limit_cooldown(self) -> float: diff --git a/zhixing-server/tests/unit/market_data/test_tushare.py b/zhixing-server/tests/unit/market_data/test_tushare.py index 04d74fd..a6ea2bc 100644 --- a/zhixing-server/tests/unit/market_data/test_tushare.py +++ b/zhixing-server/tests/unit/market_data/test_tushare.py @@ -1,3 +1,4 @@ +import threading from datetime import date import pytest @@ -87,6 +88,109 @@ def test_rate_limit_cooldown_is_shared_by_following_requests() -> None: assert waits == [60] +def test_request_start_interval_allows_overlapping_provider_calls() -> None: + current = [0.0] + state_lock = threading.Lock() + first_started = threading.Event() + release_first = threading.Event() + waits: list[float] = [] + starts: list[tuple[str, float]] = [] + errors: list[BaseException] = [] + + def clock() -> float: + with state_lock: + return current[0] + + def wait(seconds: float) -> None: + with state_lock: + waits.append(seconds) + current[0] += seconds + + coordinator = RequestCoordinator( + max_retries=0, + request_interval_seconds=0.2, + clock=clock, + wait_fn=wait, + sleep_fn=wait, + ) + + def first_request() -> object: + starts.append(("first", clock())) + first_started.set() + if not release_first.wait(timeout=2): + raise AssertionError("first provider call was not released") + return "first" + + def run_first() -> None: + try: + coordinator.call("first", first_request) + except BaseException as exc: # pragma: no cover - surfaced by the assertion below + errors.append(exc) + + first_thread = threading.Thread(target=run_first) + first_thread.start() + assert first_started.wait(timeout=2) + + second = coordinator.call( + "second", + lambda: starts.append(("second", clock())) or "second", + ) + + assert second == "second" + assert first_thread.is_alive() + release_first.set() + first_thread.join(timeout=2) + assert not first_thread.is_alive() + assert errors == [] + assert starts == [("first", 0.0), ("second", 0.2)] + assert waits == [0.2] + + +def test_request_start_interval_is_disabled_by_default() -> None: + waits: list[float] = [] + starts: list[str] = [] + coordinator = RequestCoordinator( + max_retries=0, + clock=lambda: 0.0, + wait_fn=waits.append, + ) + + coordinator.call("first", lambda: starts.append("first")) + coordinator.call("second", lambda: starts.append("second")) + + assert starts == ["first", "second"] + assert waits == [] + + +def test_request_start_interval_applies_to_retry_attempts() -> None: + current = [0.0] + waits: list[float] = [] + starts: list[float] = [] + + def wait(seconds: float) -> None: + waits.append(seconds) + current[0] += seconds + + coordinator = RequestCoordinator( + max_retries=1, + backoff_seconds=0, + request_interval_seconds=0.2, + clock=lambda: current[0], + wait_fn=wait, + sleep_fn=wait, + ) + + def request() -> object: + starts.append(current[0]) + if len(starts) == 1: + raise RuntimeError("transient provider failure") + return "ok" + + assert coordinator.call("daily", request) == "ok" + assert starts == [0.0, 0.2] + assert waits == [0.0, 0.2] + + def test_pro_bar_qfq_calls_are_bound_to_the_shared_coordinator( monkeypatch: pytest.MonkeyPatch, ) -> None: diff --git a/zhixing-server/tests/unit/sector_radar/test_build.py b/zhixing-server/tests/unit/sector_radar/test_build.py index b8b9511..c1b981d 100644 --- a/zhixing-server/tests/unit/sector_radar/test_build.py +++ b/zhixing-server/tests/unit/sector_radar/test_build.py @@ -1,5 +1,6 @@ import logging from collections.abc import Sequence +from dataclasses import replace from datetime import UTC, date, datetime, timedelta from decimal import Decimal @@ -51,6 +52,7 @@ class FakeRadarSource: self.net_scale = net_scale self.fail_daily = False self.calls: list[str] = [] + self.moneyflow_candidate_codes: list[tuple[str, ...]] = [] def _result[T]( self, api_name: str, target: date | None, rows: tuple[T, ...] @@ -239,8 +241,13 @@ class FakeRadarSource: ) return self._result("daily", trade_date, rows) - def fetch_moneyflow_dc(self, trade_date: date) -> SourceResult[MoneyflowDcRow]: + def fetch_moneyflow_dc( + self, + trade_date: date, + candidate_codes: Sequence[str], + ) -> SourceResult[MoneyflowDcRow]: self.calls.append("moneyflow_dc") + self.moneyflow_candidate_codes.append(tuple(candidate_codes)) count = 4 if self.missing_moneyflow else 5 rows = tuple( MoneyflowDcRow( @@ -288,6 +295,28 @@ def test_successful_build_is_idempotent_and_failed_retry_preserves_last_good() - assert any(item.status is PublicationStatus.FAILED for item in repository.publications.values()) +def test_build_passes_stable_current_listing_member_intersection_to_moneyflow() -> None: + class FutureListingSource(FakeRadarSource): + def fetch_stock_basics(self) -> SourceResult[StockBasicRow]: + result = super().fetch_stock_basics() + rows = result.rows[:-1] + (replace(result.rows[-1], list_date=date(2027, 1, 1)),) + return self._result("stock_basic", None, rows) + + source = FutureListingSource() + + summary = BuildSectorRadar( + source, + InMemorySectorRadarRepository(), + today=TARGET_DATE, + now_fn=lambda: NOW, + ).execute(BuildSectorRadarCommand(trade_date=TARGET_DATE)) + + assert summary.status == "success" + assert source.moneyflow_candidate_codes == [ + ("000001.SZ", "000002.SZ", "000003.SZ", "000004.SZ") + ] + + def test_source_contract_failure_is_logged_with_safe_build_context( caplog: pytest.LogCaptureFixture, ) -> None: @@ -386,6 +415,85 @@ def test_unknown_membership_is_persisted_as_partial_and_retried_independently() assert source.calls == ["members"] +def test_membership_retry_refreshes_moneyflow_when_replay_misses_new_candidates() -> None: + class ExpandingMembershipSource(FakeRadarSource): + def fetch_sector_members( + self, + trade_date: date, + sector_codes: Sequence[str], + ) -> SourceResult[SectorMemberRow]: + self.calls.append("members") + rows: list[SectorMemberRow] = [] + snapshots: list[SourceSnapshot] = [] + for index, sector_code in enumerate(sector_codes, start=1): + sector_rows = ( + () + if self.missing_membership and index == len(sector_codes) + else ( + SectorMemberRow( + trade_date, + sector_code, + f"00000{index}.SZ", + f"股票{index}", + ), + ) + ) + snapshots.append( + build_source_snapshot( + api_name="dc_member", + params={ + "trade_date": trade_date.isoformat(), + "ts_code": sector_code, + }, + rows=tuple(self._raw_row(row) for row in sector_rows), + target_trade_date=trade_date, + partition_key=sector_code, + observed_at=NOW, + ) + ) + rows.extend(sector_rows) + return SourceResult(tuple(snapshots), tuple(rows)) + + def fetch_moneyflow_dc( + self, + trade_date: date, + candidate_codes: Sequence[str], + ) -> SourceResult[MoneyflowDcRow]: + self.calls.append("moneyflow_dc") + self.moneyflow_candidate_codes.append(tuple(candidate_codes)) + rows = tuple( + MoneyflowDcRow( + trade_date, + code, + code, + Decimal(1), + Decimal(0), + Decimal(0), + Decimal(10), + ) + for code in candidate_codes + ) + return self._result("moneyflow_dc", trade_date, rows) + + repository = InMemorySectorRadarRepository() + source = ExpandingMembershipSource(missing_membership=True) + use_case = BuildSectorRadar(source, repository, today=TARGET_DATE, now_fn=lambda: NOW) + + partial = use_case.execute(BuildSectorRadarCommand(trade_date=TARGET_DATE)) + partial_id = partial.outcomes[0].publication_id + assert partial.status == "partial" + assert partial_id is not None + assert source.moneyflow_candidate_codes == [("000001.SZ",)] + + source.missing_membership = False + source.calls.clear() + retried = use_case.execute(BuildSectorRadarCommand(retry_publication_id=partial_id)) + + assert retried.status == "success" + assert source.calls == ["members", "moneyflow_dc"] + assert source.moneyflow_candidate_codes[-1] == ("000001.SZ", "000002.SZ") + + def test_range_builds_dates_in_order_and_retry_uses_old_target() -> None: repository = InMemorySectorRadarRepository() source = FakeRadarSource(missing_moneyflow=True) diff --git a/zhixing-server/tests/unit/sector_radar/test_cli.py b/zhixing-server/tests/unit/sector_radar/test_cli.py index 9820748..e9101b9 100644 --- a/zhixing-server/tests/unit/sector_radar/test_cli.py +++ b/zhixing-server/tests/unit/sector_radar/test_cli.py @@ -76,7 +76,9 @@ def test_cli_main_returns_summary_exit_code_and_json( @staticmethod def from_token(token: str, **kwargs: object) -> object: assert token == "secret-token" - assert kwargs + assert kwargs["max_retries"] == 3 + assert kwargs["backoff_seconds"] == 1.0 + assert kwargs["request_interval_seconds"] == 0.2 return object() class FakeBuild: diff --git a/zhixing-server/tests/unit/sector_radar/test_tushare_source.py b/zhixing-server/tests/unit/sector_radar/test_tushare_source.py index 5901ebd..b889a91 100644 --- a/zhixing-server/tests/unit/sector_radar/test_tushare_source.py +++ b/zhixing-server/tests/unit/sector_radar/test_tushare_source.py @@ -1,4 +1,5 @@ import logging +import threading from collections.abc import Mapping from datetime import UTC, date, datetime from decimal import Decimal @@ -24,11 +25,13 @@ class QueryClient: def __init__(self, responses: Mapping[tuple[str, str], object]) -> None: self.responses = dict(responses) self.calls: list[tuple[str, dict[str, object]]] = [] + self._lock = threading.Lock() def query(self, api_name: str, **kwargs: object) -> object: - self.calls.append((api_name, kwargs)) partition = str(kwargs.get("ts_code") or kwargs.get("list_status") or "") - response = self.responses.get((api_name, partition), ()) + with self._lock: + self.calls.append((api_name, kwargs)) + response = self.responses.get((api_name, partition), ()) if isinstance(response, BaseException): raise response return response @@ -44,6 +47,22 @@ def make_adapter(client: object) -> TushareSectorRadarAdapter: ) +def moneyflow_record( + ts_code: str, + *, + trade_date: str = "20260828", +) -> dict[str, object]: + return { + "trade_date": trade_date, + "ts_code": ts_code, + "name": ts_code, + "net_amount": "1", + "net_amount_rate": "0.1", + "pct_change": "1", + "close": "10", + } + + def test_daily_and_moneyflow_keep_source_units_and_distinguish_missing_from_zero() -> None: client = QueryClient( { @@ -98,7 +117,7 @@ def test_daily_and_moneyflow_keep_source_units_and_distinguish_missing_from_zero adapter = make_adapter(client) daily = adapter.fetch_daily(TARGET_DATE) - moneyflow = adapter.fetch_moneyflow_dc(TARGET_DATE) + moneyflow = adapter.fetch_moneyflow_dc(TARGET_DATE, ("000001.SZ", "000002.SZ")) assert daily.rows[0].amount_thousand_yuan == Decimal("12.5") assert daily.rows[0].turnover_yuan == Decimal("12500.0") @@ -109,6 +128,168 @@ def test_daily_and_moneyflow_keep_source_units_and_distinguish_missing_from_zero assert client.calls[0][1]["fields"] == ",".join(source_module.FIELDS["daily"]) +def test_moneyflow_accepts_a_full_initial_snapshot_at_the_provider_limit( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setitem(source_module.ROW_LIMITS, "moneyflow_dc", 2) + client = QueryClient( + { + ("moneyflow_dc", ""): ( + moneyflow_record("000001.SZ"), + moneyflow_record("000002.SZ"), + ) + } + ) + + result = make_adapter(client).fetch_moneyflow_dc( + TARGET_DATE, + ("000001.SZ", "000002.SZ"), + ) + + assert result.snapshots[0].limit_reached is True + assert [row.ts_code for row in result.rows] == ["000001.SZ", "000002.SZ"] + assert len(client.calls) == 1 + + +@pytest.mark.parametrize( + ("initial_rows", "message"), + ( + ((moneyflow_record("000001.SZ", trade_date="20260827"),), "trade_date"), + ( + (moneyflow_record("000001.SZ"), moneyflow_record("000001.SZ")), + "duplicate business keys", + ), + ), +) +def test_moneyflow_initial_contract_errors_fail_closed( + initial_rows: tuple[dict[str, object], ...], + message: str, +) -> None: + client = QueryClient({("moneyflow_dc", ""): initial_rows}) + + with pytest.raises(SourceContractError, match=message): + make_adapter(client).fetch_moneyflow_dc(TARGET_DATE, ()) + + +def test_moneyflow_refills_only_missing_codes_in_stable_snapshot_order( + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setitem(source_module.ROW_LIMITS, "moneyflow_dc", 3) + third_finished = threading.Event() + completion_order: list[str] = [] + completion_lock = threading.Lock() + + class ReverseCompletionClient(QueryClient): + def query(self, api_name: str, **kwargs: object) -> object: + response = super().query(api_name, **kwargs) + ts_code = str(kwargs.get("ts_code") or "") + if ts_code == "000004.SZ": + if not third_finished.wait(timeout=2): + raise AssertionError("second moneyflow worker did not start") + elif ts_code == "000005.SZ": + third_finished.set() + if ts_code: + with completion_lock: + completion_order.append(ts_code) + return response + + client = ReverseCompletionClient( + { + ("moneyflow_dc", ""): tuple( + moneyflow_record(f"00000{index}.SZ") for index in range(1, 4) + ), + ("moneyflow_dc", "000004.SZ"): (moneyflow_record("000004.SZ"),), + ("moneyflow_dc", "000005.SZ"): (moneyflow_record("000005.SZ"),), + } + ) + + result = make_adapter(client).fetch_moneyflow_dc( + TARGET_DATE, + tuple(f"00000{index}.SZ" for index in range(1, 6)), + ) + + assert completion_order == ["000005.SZ", "000004.SZ"] + assert [snapshot.partition_key for snapshot in result.snapshots] == [ + "all", + "000004.SZ", + "000005.SZ", + ] + assert [row.ts_code for row in result.rows] == [ + "000001.SZ", + "000002.SZ", + "000003.SZ", + "000004.SZ", + "000005.SZ", + ] + assert len(client.calls) == 3 + + +def test_moneyflow_empty_and_exhausted_refills_remain_real_gaps( + caplog: pytest.LogCaptureFixture, +) -> None: + client = QueryClient( + { + ("moneyflow_dc", ""): (moneyflow_record("000001.SZ"),), + ("moneyflow_dc", "000002.SZ"): (), + ("moneyflow_dc", "000003.SZ"): RuntimeError("private provider payload"), + } + ) + caplog.set_level( + logging.WARNING, + logger="zhixing_server.modules.sector_radar.infrastructure.tushare", + ) + + result = make_adapter(client).fetch_moneyflow_dc( + TARGET_DATE, + ("000001.SZ", "000002.SZ", "000003.SZ"), + ) + + assert [row.ts_code for row in result.rows] == ["000001.SZ"] + assert [snapshot.partition_key for snapshot in result.snapshots] == ["all"] + messages = "\n".join(record.getMessage() for record in caplog.records) + assert "partition_empty partition_key=000002.SZ" in messages + assert "partition_failed partition_key=000003.SZ" in messages + assert "private provider payload" not in messages + + +@pytest.mark.parametrize( + ("partition_rows", "row_limit", "message"), + ( + ((moneyflow_record("000002.SZ", trade_date="20260827"),), 6_000, "trade_date"), + ((moneyflow_record("000099.SZ"),), 6_000, "different ts_code"), + ( + (moneyflow_record("000002.SZ"), moneyflow_record("000002.SZ")), + 6_000, + "duplicate business keys", + ), + ( + (moneyflow_record("000002.SZ"), moneyflow_record("000002.SZ")), + 2, + "provider row limit", + ), + ), +) +def test_moneyflow_partition_contract_errors_fail_closed( + monkeypatch: pytest.MonkeyPatch, + partition_rows: tuple[dict[str, object], ...], + row_limit: int, + message: str, +) -> None: + monkeypatch.setitem(source_module.ROW_LIMITS, "moneyflow_dc", row_limit) + client = QueryClient( + { + ("moneyflow_dc", ""): (moneyflow_record("000001.SZ"),), + ("moneyflow_dc", "000002.SZ"): partition_rows, + } + ) + + with pytest.raises(SourceContractError, match=message): + make_adapter(client).fetch_moneyflow_dc( + TARGET_DATE, + ("000001.SZ", "000002.SZ"), + ) + + def test_non_finite_source_values_are_rejected() -> None: client = QueryClient( { @@ -305,32 +486,54 @@ def test_dc_member_preserves_an_explicit_empty_partition() -> None: assert result.snapshots[1].row_count == 0 -def test_stock_basic_explicitly_requests_all_lifecycle_statuses() -> None: - responses = { - ( - "stock_basic", - status, - ): ( - { - "ts_code": f"00000{index}.SZ", - "symbol": f"00000{index}", - "name": status, - "market": None if status == "D" else "主板", - "exchange": "SZSE", - "list_status": status, - "list_date": "20200101", - "delist_date": None, - }, - ) - for index, status in enumerate(("L", "D", "P", "G", "UN"), start=1) - } - client = QueryClient(responses) +def test_stock_basic_requests_only_current_listings() -> None: + client = QueryClient( + { + ( + "stock_basic", + "L", + ): ( + { + "ts_code": "000001.SZ", + "symbol": "000001", + "name": "L", + "market": "主板", + "exchange": "SZSE", + "list_status": "L", + "list_date": "20200101", + "delist_date": None, + }, + ) + } + ) result = make_adapter(client).fetch_stock_basics() - assert {row.list_status for row in result.rows} == {"L", "D", "P", "G", "UN"} - assert next(row for row in result.rows if row.list_status == "D").market is None - assert [call[1]["list_status"] for call in client.calls] == ["L", "D", "P", "G", "UN"] + assert {row.list_status for row in result.rows} == {"L"} + assert [snapshot.partition_key for snapshot in result.snapshots] == ["L"] + assert [call[1]["list_status"] for call in client.calls] == ["L"] + + +def test_stock_basic_rejects_a_non_listed_row_from_the_l_partition() -> None: + client = QueryClient( + { + ("stock_basic", "L"): ( + { + "ts_code": "000001.SZ", + "symbol": "000001", + "name": "unexpected", + "market": "主板", + "exchange": "SZSE", + "list_status": "D", + "list_date": "20200101", + "delist_date": "20260828", + }, + ) + } + ) + + with pytest.raises(SourceContractError, match="unexpected list_status"): + make_adapter(client).fetch_stock_basics() def test_suspend_timing_may_be_missing_while_suspend_type_remains_required() -> None: