diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/check.jsonl b/.trellis/tasks/09-21-radar-weighted-score-rank-change/check.jsonl new file mode 100644 index 0000000..72258ad --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/check.jsonl @@ -0,0 +1,9 @@ +{"file": ".trellis/spec/backend/index.md", "reason": "后端包边界与必需检查"} +{"file": ".trellis/spec/backend/tushare-listed-stock-universe.md", "reason": "继续遵守当前上市证券母集,不能为对齐网站擅自改变范围"} +{"file": ".trellis/spec/backend/http-api-contracts.md", "reason": "评分与三指标变化字段的端到端契约"} +{"file": ".trellis/spec/frontend/component-guidelines.md", "reason": "镜像榜单、控件与可访问性"} +{"file": ".trellis/spec/frontend/type-safety.md", "reason": "API nullable 字段与严格解析"} +{"file": ".trellis/tasks/09-21-radar-weighted-score-rank-change/research/formula-evidence.md", "reason": "已验证的评分结构及数值证据边界"} +{"file": ".trellis/tasks/09-21-radar-weighted-score-rank-change/research/code-evidence.md", "reason": "现有实现定位与需保持的版本/持久化行为"} +{"file": ".trellis/spec/backend/quality-guidelines.md", "reason": "后端 Ruff、Pyright、pytest 和持久化检查"} +{"file": ".trellis/spec/frontend/quality-guidelines.md", "reason": "前端格式、lint、类型、行为与构建检查"} diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/design.md b/.trellis/tasks/09-21-radar-weighted-score-rank-change/design.md new file mode 100644 index 0000000..f1d4747 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/design.md @@ -0,0 +1,48 @@ +# 设计:加权评分与排名变化 + +## 计算与领域边界 + +评分留在 sector_radar 领域/发布阶段,不在 React 中根据当前页数据计算百分位。当前每板块 MetricStrategy 只能得到该板块历史,因此新增全池评分步骤,接收同一目标日所有板块的原始特征,按日期、类型分池计算。 + +新版本拟为 `zhixing_ratio_weighted_v2` 和 `zhixing_swing_weighted_v2`。设 A 为现有聚合成交额元、F 为聚合主力净额元;按前置研究的可复算输入约定取 `r=F/(A+100)`,A<=0 或输入未知时不制造比率。W=`log10(1+MA5(A))/10`;P 为原始指标升序平均名次/N。 + +```text +单日原值 = r +单日评分 = 1000 * P(r) * W +波段原值 = .5*MA3(r) + .5*MA10(r) +波段评分 = 1000*(.5*P(MA3(r)) + .5*P(MA10(r)))*W +``` + +使用 Decimal 计算金额、比率和对数;仅在输出展示时保留 1 位小数。P 的并列值采用平均名次,最终分数并列沿用 sector_code 升序作为稳定破同分规则。原始值与 weighted_score 分离;窗口不足时保留可确定的原始值,评分和对应名次为空,不以原始值替代缺失评分。质量状态和覆盖率继续传播。 + +新版本使用保存的交易日历确定窗口与对比日期,缺少某个应有交易日输入时保留未知;不得悄悄以更早成功发布替代缺失交易日。第一个可复算日期受库内真实历史覆盖限制。 + +## 排名、历史与版本 + +`rank_metric_observations` 对金额继续使用原始金额,对新版本单日/波段使用 weighted_score。上下普通榜筛选使用最终百分位;底榜排序也须与评分键一致。排名变化在同版本最终名次上计算 `past_rank-current_rank`;取 ceil(有效排名池大小×10%),历史不可比者不进入变化候选。 + +一行排名变化响应提供所选基准兼容字段 rank_change,以及 amount/ratio/swing 三个变化值,供中心主列和两侧辅助列复用。查询选出的前后榜由所选基准确定;点击辅助列只改变该侧当前候选排列。字段均允许 null。 + +旧发布按其实际 metric_versions 解析定义与历史,不用新版本常量把旧数据过滤成空,也不比较 v1 与 v2 名次。榜单、详情和历史折线共享版本解析与排名事实。首次切换期间旧发布评分显示“—”;历史重算完成后,同日最新成功发布提供新评分。 + +## 持久化与历史重算 + +新增 Alembic migration,为 `sector_radar_ranking` 增加 nullable NUMERIC weighted_score,保持 metric_value 原义。同步所有 INSERT、SELECT、序列化和反序列化;内存仓库保持同等语义。 + +发布构建使用统一评分服务,版本参与 input_hash。增加离线重算入口,仅使用已落库聚合事实和原发布的来源关联,不初始化 Tushare 客户端。按交易日先后创建新派生发布,保留旧发布与原始快照;固定源 publication IDs 避免重算期间新发布改变本轮输入。复用现有按日锁、事务和 last-good 规则,重算幂等键包含源发布身份/输入指纹和新策略版本。 + +重算时必须复制/关联原聚合的 pct_change、leading_code 及来源组,使新发布的详情仍可读取;不能只写 ranking 而产生空详情。适配器需要提供 publication 精确的 aggregate records 与来源读取,避免现有 history 方法丢失这些元信息。 + +生产数据库的迁移和重算是独立运行步骤;在本地代码和验证准备完成前不执行外部写入。回退可恢复旧代码版本并重新选择保留的旧发布;nullable 新列本身保持兼容。 + +## HTTP 与前端 + +- 在榜单行和必要详情摘要中扩展 weighted_score;排名变化补充三指标变化映射与对比日期信息,同步 Pydantic、TypeScript、解析器和 query 测试。 +- 单日/波段左右增加评分列,1 位小数、缺失“—”;其余原始百分比保留独立展示。 +- rank_change 使用独立次行,左侧三枚排序基准按钮,右侧近 1–5 日下拉。默认波段率和 1 日,仅在 URL 没有显式值时采用默认。 +- 标题为排名飙升榜/排名暴跌榜,中央为“{基准全名}排名变化”;主列显示所选指标变化,辅助列显示另外两项,增加涨跌幅。 +- 延续现有主题、表格镜像、板块详情入口和响应式横向滚动;不添加未要求的全局导出、搜索改版或主题重制。 + +## 主要风险 + +新算法会改变单日/波段榜及历史名次;必须依靠新版本与离线重算切换,不能给旧名次套新评分。知行的 1000 个板块及按日成员与 OneChart 的 791 板块输入不同,算法结构一致不保证逐值一致。原站并列、缺失与极低成交额策略未完全公开,本系统明确采用上述确定规则。 diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/implement.jsonl b/.trellis/tasks/09-21-radar-weighted-score-rank-change/implement.jsonl new file mode 100644 index 0000000..bcc64e9 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/implement.jsonl @@ -0,0 +1,7 @@ +{"file": ".trellis/spec/backend/index.md", "reason": "后端包边界与必需检查"} +{"file": ".trellis/spec/backend/tushare-listed-stock-universe.md", "reason": "继续遵守当前上市证券母集,不能为对齐网站擅自改变范围"} +{"file": ".trellis/spec/backend/http-api-contracts.md", "reason": "评分与三指标变化字段的端到端契约"} +{"file": ".trellis/spec/frontend/component-guidelines.md", "reason": "镜像榜单、控件与可访问性"} +{"file": ".trellis/spec/frontend/type-safety.md", "reason": "API nullable 字段与严格解析"} +{"file": ".trellis/tasks/09-21-radar-weighted-score-rank-change/research/formula-evidence.md", "reason": "已验证的评分结构及数值证据边界"} +{"file": ".trellis/tasks/09-21-radar-weighted-score-rank-change/research/code-evidence.md", "reason": "现有实现定位与需保持的版本/持久化行为"} diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/implement.md b/.trellis/tasks/09-21-radar-weighted-score-rank-change/implement.md new file mode 100644 index 0000000..502c4c9 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/implement.md @@ -0,0 +1,39 @@ +# 执行计划 + +状态:本地实现与范围内验收完成;全仓既有失败已在实施前版本复现。详细结果及生产应用命令见 `validation.md`。 + +1. 完整阅读要修改的文件;确认已保存的评分研究、接口契约、版本及缺失语义。核对需要使用的 Alembic/Pydantic 等当前版本与官方文档。 +2. 实现版本化评分特征和全池计算,分离 raw value/score;更换单日/波段排序键,补独立数值样例、并列、未来数据排除及缺失窗口测试。 +3. 新增 nullable score 列迁移;同步 PostgreSQL/内存仓库所有读写与 historical ranking 查询。补充真实数据库往返、旧行 NULL 和事务失败测试。 +4. 接入发布构建与 input_hash;实现不访问 Tushare 的历史重算入口、新旧版本解析及严格对比日期。验证幂等、旧发布保留、失败 last-good 与详情来源完整。 +5. 扩展榜单/详情 HTTP 输出和前端类型解析,支持评分、三指标变化值及对比日期;覆盖不匹配版本与缺失历史。 +6. 单日/波段表格添加评分列与侧内排序;按截图调整排名变化次级控制行、双榜标题、中央标题、涨跌幅和两项辅助变化列,保持 URL/查询联动。 +7. 执行领域、发布、读取、HTTP、仓库、前端 API 和页面行为测试;用本地测试数据库验证迁移和离线重算,避免使用用户生产库进行测试写入。 +8. 完成后端 Ruff/Pyright/pytest 与前端格式、lint、typecheck、Vitest/build。浏览器实际检查三基准、1–5 日、空历史、评分列和窄屏;保存必要截图和验证结果。 +9. 检查 diff 范围、已有用户改动和未验证事项;交付本地成果及明确的迁移/重算命令,不自行提交、部署或写入生产库。 + +## 验证命令 + +```bash +cd zhixing-server +uv run ruff format --check . +uv run ruff check . +uv run pyright +uv run pytest +``` + +```bash +cd zhixing-web +pnpm check +pnpm build +``` + +实现期间先执行相关测试,最后执行项目规定完整检查;通过后仅在新改动或未解决问题需要时重复。 + +## 高风险核验点 + +- P 必须来自完整同类池,不能来自当前页面 TOP10% 子集。 +- Swing 分数是两次排名后合成,不是合成 raw 后再排名。 +- 历史不足时不能伪造评分/0变化;不能跨失败交易日跳位或跨版本相减。 +- migration 列顺序涉及所有 ranking SELECT 与 fake row,不得只改写入。 +- 离线重算的新发布必须可继续打开详情,且不得重新请求上游或改变旧快照。 diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/prd.md b/.trellis/tasks/09-21-radar-weighted-score-rank-change/prd.md new file mode 100644 index 0000000..c8d4ea2 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/prd.md @@ -0,0 +1,46 @@ +# 雷达加权评分与排名变化复刻 + +## Goal + +根据已验证的 OneChart 评分公式,为单日和波段榜增加加权评分,并复刻排名变化的排序基准、1至5日窗口和双榜交互。 + +## Requirements + +- R1:单日流入率与波段流入率双榜各增加“加权评分”列,左右镜像排列,保留 1 位小数,支持表头排序;原始流入率仍单独展示。 +- R2:采用本轮前置研究验证的评分结构。单日和波段榜按各自加权评分形成最终排名;波段原始流入率采用 3 日、10 日单日流入率均值各 50%。单日净额榜按净额形成排名。 +- R3:排名变化视图按参考截图提供独立次级控制行:左侧“排序基准”含波段率、单日率、单日额,右侧“统计天数”含近 1–5 日;首次进入默认波段率、近 1 日,并保留 URL 中显式指定的选择。 +- R4:排名变化视图展示“排名飙升榜 TOP 10% / 排名暴跌榜 BOTTOM 10%”,中央标题随基准切换;每行显示涨跌幅、所选基准的名次变化以及其余两个指标的名次变化。上下榜由所选基准决定,表头排序只重排该榜现有候选。 +- R5:名次变化为过去名次减当前名次,统计基于交易日,概念与行业分池。历史缺失、算法版本不兼容时显示“—”,不得伪造 0 或混用新旧算法名次。 +- R6:评分、最终名次、排名变化和详情历史使用相同算法版本;提供基于已有聚合事实重算历史的能力,旧发布保持可追溯。已有旧版本数据升级前仍能安全读取。 +- R7:继续使用知行当前数据库、板块池、按日成员快照与当前上市股票范围。复刻评分结构与交互,不以抓取 OneChart 结果替代业务计算,也不承诺与不同输入口径的网站逐值相等。 + +## 范围边界 + +- 本任务包含必要的领域计算、数据库派生字段迁移、HTTP 契约、前端及测试,以及历史重算命令。 +- 不包含调整证券母集、将网站当前成员回填历史、额外的榜单导出或搜索功能重写。 +- 本地实现与非破坏性验证按任务执行;生产迁移、生产历史重算、部署及提交不在本轮已获授权的执行范围内。 + +## Acceptance Criteria + +- [x] AC1/R1:单日/波段榜显示左右“加权评分”列,有限值为 1 位小数,缺失为“—”,排序与表头状态正确。 +- [x] AC2/R2:独立样例验证百分位、5 日成交额权重、3/10 日各 50% 及“先分别排名再合成”;高原始比率但低评分的反例能按评分正确排名。 +- [x] AC3/R3–R4:三基准 × 五窗口切换正确更新双榜、中央标题和辅助两列,URL 刷新可恢复;正负方向、镜像布局和窄屏横向滚动正确。 +- [x] AC4/R5:交易日跨周末、缺发布、缺板块、零变化、并列及新旧版本不匹配均有测试,榜单选取为对应池 ceil(N×10%) 个有效变化候选的上限。 +- [x] AC5/R6:旧行 weighted_score 为空时可读取;历史重算创建可追溯新发布且幂等,不请求外部数据;失败不影响最近成功发布,详情与榜单的当前名次一致。 +- [ ] AC6/R7:沿用当前股票/板块数据范围;单元、HTTP、前端、持久化及浏览器验证覆盖变更,项目规定的质量检查通过。 + +## 已确认的现状与依据 + +- 前置研究:`research/formula-evidence.md`。公开输入 9,492 条样本的评分与最终名次可精确重现;数据库输入存在成员/板块池差异。 +- 当前单日/波段只有原始 metric_value;旧波段策略实际为 3 至 10 日八个净额/成交额窗口等权,并非新研究中的两个窗口,见 `domain/metrics.py:140` 起。 +- 当前排名变化已有基础 API 与 1–5 日变化值,但 UI 控件、标题、辅助列不匹配截图,且历史名次基于旧原始指标,定位见 `research/code-evidence.md`。 + +## 决策状态 + +用户已明确回复“同意,开始实施”,批准本版方案。执行现有数据口径下的评分和交互升级;按个人约定由主代理编码与最终验证,子代理仅承担只读探索和独立核验。 + +## 本轮验收状态 + +AC1–AC5 已完成。AC6 的变更范围检查和浏览器验收通过;全仓检查仍有在实施前 HEAD 上独立复现的既有失败,因此保留未完全通过状态。详见 `validation.md`。本地实现已交付;生产迁移、重算、部署与 Git 提交均未执行。 + +用户已于本轮授权提交并推送。本任务的实现和范围内验证完成;AC6 的全仓基线失败仍明确保留为已知限制,不扩展修复选股/行情模块。线上重算说明已补入 `docs/market-data-sync.md`。 diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/code-evidence.md b/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/code-evidence.md new file mode 100644 index 0000000..15d1bc5 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/code-evidence.md @@ -0,0 +1,36 @@ +# 现有实现与改动定位 + +只读研究已完成:`/root/radar_storage_research`、`/root/radar_ui_research`,主会话另行阅读领域算法、models、read、build 和 HTTP 相关契约。 + +## 领域与应用 + +- `domain/metrics.py:140` 起,旧 SwingEqualThreeToTenStrategy 取 3–10 日八个窗口的累计净额/累计成交额,再等权平均;需采用新版本,不能把它误认作前置研究的两个窗口。 +- `domain/ranking.py:42` 的排序键按 observation.value 排序;`with_rank_changes` 同版本过去名次减当前名次;select_rank_change_side 已实现 ceil(N×10%),不需另写重复算法。 +- `application/build.py:694` 的 _rank 读取 9 个历史成功日期再拼当前日;历史排名取前 5 次成功发布,并非严格前 5 交易日。 +- `application/build.py:753` 的 input_hash 包含策略版本;重算可沿用版本隔离与幂等思路。 +- `application/read.py` 的 _METRIC_DEFINITIONS 和 `application/details.py:25` 的 _METRIC_VERSIONS 都硬编码当前策略;升级必须让旧发布保持可读。 +- `application/read.py:289` 的 extras 只服务 amount/ratio,swing/rank_change 的涨跌幅、净额和辅助字段需要补齐。 +- `application/details.py:173` 的 history_data 已按发布来源中的交易日历补出缺口,适合统一严格交易日语义;必须避免每行单独查库。 + +## 存储 + +- `domain/persistence.py:294–355` 定义 aggregate、ranking、history 读写契约。 +- `infrastructure/postgres.py:532–710` 为事务发布,`:741–778` 为 ranking 批写,`:1145–1173` 为序列化,`:1239–1267` 为固定列反序列化。 +- ranking SELECT 位于 `postgres.py:871–964`、`:996–1033`,当前只含 metric_value,没有 score。 +- `postgres.py:966–994` 的 aggregate history 只返回原始金额/coverage,不包含 publication_id、pct_change、leading_code,离线重算需要精确来源记录。 +- 最新迁移为 `0009_radar_sector_detail`。新增 score 应采用独立 migration,旧值可空。 +- 单测 `tests/unit/sector_radar/test_postgres_repository.py:43–65` 使用 17 列 fake row,必须随查询同步。现有集成测试尚未覆盖 ranking/aggregate 往返。 + +## 前端 + +- `pages/sector-radar-page.tsx:209–229` 已有两个下拉,分别为变化指标/对比区间;改为截图的按钮组与统计天数次行。 +- `pages/sector-radar-page.tsx:406–428` 顶部双榜和中央标题仍是普通资金榜文案。 +- `pages/sector-radar-page.tsx:430–463`、`:643–708` 是单日/波段镜像列,无 score。 +- `pages/sector-radar-page.tsx:465–500` 的变化榜目前只显示样本、排名百分位与一个变化值;应展示涨跌幅及另外两项指标变化。 +- `api/sector-radar.types.ts:61–114` 与 `api/sector-radar.api.ts:277–377` 仅承载单个 rank_change,需扩展。 +- 路由与 query key 已保存基准/天数,保留这一结构;URL 显式值优先。 +- 页面测试 `sector-radar-page.test.tsx:654–683` 已覆盖正变化/缺历史,需增加评分、三基准五窗口、镜像辅助列、侧内排序与截图文案。 + +## 参考站核对 + +本会话下载的 `/tmp/onechart-public-evidence/index.html` 中:`:1197` 默认 Swing/1 日,`:1380` 附近 mirrorMetricLabel 随基准变,`:1397–1427` 的 mirrorAuxColumns 在变化榜显示另外两个基准,在单日/波段榜显示 score/净额/在榜。前置研究完整来源和数值证据位于已归档的评分研究任务。 diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/formula-evidence.md b/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/formula-evidence.md new file mode 100644 index 0000000..feb1491 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/formula-evidence.md @@ -0,0 +1,19 @@ +# 加权公式与验证边界 + +前置研究在 2026-09-04 至 2026-09-21 的 9,492 条公开样本上复现单日、波段评分与最终排名;分数最大误差 3.41e-13,属于浮点运算误差。这是观测数据还原,不是原站源码证据。 + +对本项目已有聚合数据,F 为主力净额元、A 为成交额元: + +```text +r = F / (A + 100) +W = log10(1 + MA5(A)) / 10 +单日评分 = 1000 × P(r) × W +波段评分 = 500 × (P(MA3(r)) + P(MA10(r))) × W +波段原值 = (MA3(r) + MA10(r)) / 2 +``` + +P 为同日、同类型全池升序平均名次 / N。波段必须先分别排名,再合成。最终评分并列按板块代码稳定排序。权重可能大于 1,分数不是固定上限 1000 的百分制。 + +公开字段可以确认的等价权重是 `log10(MA5(F/r)-99)/10`;`A=F/r-100` 是本任务采用的可复算口径,不能证明原站源码中具体常量的用途。原站边界与并列规则未完全公开,本地明确使用 design.md 中的完整交易日窗口、平均并列名次及缺失语义。 + +公开验证池为 791 个板块,本地数据库池和按日成员不同。本项目保留当前上市股票与已保存历史输入,因此不承诺与网站逐值相同。原始研究数据和复算快照保留在本地前置研究任务中,不是运行本功能的依赖。 diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/seed_local_preview.py b/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/seed_local_preview.py new file mode 100644 index 0000000..b177a9b --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/research/seed_local_preview.py @@ -0,0 +1,166 @@ +"""Synthetic browser fixtures; restricted to the disposable local radar test DB.""" + +import hashlib +import math +import os +from dataclasses import replace +from datetime import UTC, date, datetime, timedelta +from decimal import Decimal +from urllib.parse import urlsplit + +import psycopg + +from zhixing_server.modules.sector_radar.application.scoring import calculate_rankings +from zhixing_server.modules.sector_radar.domain.metrics import ( + AmountNetStrategy, + RatioTurnoverStrategy, + SwingEqualThreeToTenStrategy, +) +from zhixing_server.modules.sector_radar.domain.models import ( + PublicationStatus, + RadarPublication, + SectorDailyAggregate, + SectorType, +) +from zhixing_server.modules.sector_radar.domain.persistence import ( + DailyAggregateRecord, + PublicationSourceGroup, + PublicationSourceRecord, + RankingRecord, +) +from zhixing_server.modules.sector_radar.domain.source import build_source_snapshot +from zhixing_server.modules.sector_radar.infrastructure.postgres import ( + PostgresSectorRadarRepository, +) + +url = os.environ["ZHIXING_TEST_DATABASE_URL"] +location = urlsplit(url) +if (location.hostname, location.port, location.path) != ("127.0.0.1", 55439, "/radar_test"): + raise SystemExit("Preview seeding is restricted to the disposable local radar test database") +with psycopg.connect(url) as connection: + connection.execute("TRUNCATE sector_radar_publication CASCADE") + +start = date(2026, 8, 10) +days = tuple( + start + timedelta(days=n) for n in range(43) if (start + timedelta(days=n)).weekday() < 5 +) +strategies = (AmountNetStrategy(), RatioTurnoverStrategy(), SwingEqualThreeToTenStrategy()) +names = ( + "人工智能", + "机器人", + "商业航天", + "低空经济", + "半导体", + "创新药", + "新能源", + "电力设备", + "数字经济", + "消费电子", +) +now = datetime.now(UTC) +calendar = build_source_snapshot( + api_name="trade_cal", + params={"fixture": "weighted-radar-preview"}, + rows=[{"exchange": "SSE", "cal_date": day, "is_open": 1} for day in days], + target_trade_date=days[-1], + observed_at=now, +) +repository = PostgresSectorRadarRepository(url) +history = [] +try: + repository.save_source_snapshots((calendar,)) + for index, day in enumerate(days): + aggregates = [] + for sector_type, size in ((SectorType.CONCEPT, 300), (SectorType.INDUSTRY, 100)): + for n in range(size): + turnover = Decimal( + str(10 ** (8 + n % 4) * (1 + 0.2 * math.cos(n + index))) + ).quantize(Decimal(".01")) + ratio = Decimal( + str( + 0.09 * math.sin(n * 1.37 + index * 0.43) + 0.02 * math.cos(n * 0.71 + index) + ) + ) + aggregates.append( + SectorDailyAggregate( + day, + sector_type, + f"TEST-{sector_type.value[0]}-{n:04}", + f"{names[n % len(names)]} {n + 1}", + 10, + 10, + (turnover * ratio).quantize(Decimal(".01")), + turnover, + Decimal(1), + Decimal(1), + ) + ) + publication = RadarPublication( + f"preview-{day}", + day, + PublicationStatus.RUNNING, + "local-preview-v1", + "synthetic-preview-v1", + tuple(strategy.metric_version for strategy in strategies), + None, + Decimal(1), + now, + ) + repository.create_publication(publication) + repository.save_publication_sources( + ( + PublicationSourceRecord( + publication.publication_id, PublicationSourceGroup.CALENDAR, 0, calendar + ), + ) + ) + for sector_type, group in ( + (SectorType.CONCEPT, PublicationSourceGroup.CONCEPT_INDICES), + (SectorType.INDUSTRY, PublicationSourceGroup.INDUSTRY_INDICES), + ): + snapshot = build_source_snapshot( + api_name="dc_index", + params={"fixture": "preview", "type": sector_type.value}, + rows=[ + { + "trade_date": day, + "ts_code": row.sector_code, + "name": row.sector_name, + "pct_change": Decimal(str(3 * math.sin(n + index))).quantize( + Decimal(".01") + ), + } + for n, row in enumerate(aggregates) + if row.sector_type is sector_type + ], + target_trade_date=day, + observed_at=now, + ) + repository.save_source_snapshots((snapshot,)) + repository.save_publication_sources( + (PublicationSourceRecord(publication.publication_id, group, 0, snapshot),) + ) + ranks = calculate_rankings(day, aggregates, history, days, strategies) + repository.finalize_publication( + replace( + publication, + status=PublicationStatus.SUCCESS, + input_hash=hashlib.sha256(str(day).encode()).hexdigest(), + finished_at=now, + ), + memberships=(), + stock_facts=(), + daily_aggregates=( + DailyAggregateRecord( + publication.publication_id, + row, + Decimal(str(3 * math.sin(n + index))).quantize(Decimal(".01")), + ) + for n, row in enumerate(aggregates) + ), + rankings=(RankingRecord(publication.publication_id, row) for row in ranks), + ) + history.extend(aggregates) + print(f"Seeded {len(days)} days × 400 synthetic sectors; no production data or provider calls.") +finally: + repository.close() diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/task.json b/.trellis/tasks/09-21-radar-weighted-score-rank-change/task.json new file mode 100644 index 0000000..5f8a469 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/task.json @@ -0,0 +1,40 @@ +{ + "id": "radar-weighted-score-rank-change", + "name": "radar-weighted-score-rank-change", + "title": "雷达加权评分与排名变化复刻", + "description": "根据已验证的 OneChart 评分公式,为单日和波段榜增加加权评分,并复刻排名变化的排序基准、1至5日窗口和双榜交互。", + "status": "in_progress", + "dev_type": "fullstack", + "scope": "sector-radar", + "package": null, + "priority": "P2", + "creator": "yuxuanhui", + "assignee": "yuxuanhui", + "createdAt": "2026-09-21", + "completedAt": null, + "branch": "develop", + "base_branch": "main", + "worktree_path": null, + "commit": null, + "pr_url": null, + "subtasks": [], + "children": [], + "parent": null, + "relatedFiles": [ + "zhixing-server/src/zhixing_server/modules/sector_radar/domain/weighted.py", + "zhixing-server/src/zhixing_server/modules/sector_radar/application/recompute.py", + "zhixing-server/migrations/versions/0010_radar_weighted_score.py", + "zhixing-web/src/features/sector-radar/pages/sector-radar-page.tsx", + "zhixing-web/src/features/sector-radar/components/radar-table-sort.ts" + ], + "notes": "本地实现与范围内验收完成;用户已授权提交并推送。后端雷达 110 项、前端雷达 74 项及浏览器 30 组合通过;全仓既有失败见 validation.md。生产迁移、重算、部署未执行。", + "meta": { + "planning_status": "approved", + "planning_summary_version": 1, + "production_writes_authorized": false, + "implementation_status": "complete", + "validation_status": "scoped_pass_with_verified_baseline_failures", + "commit_push_authorized": true, + "spec_sync": "No shared spec promotion; reviewed formula and limitations retained in task artifacts." + } +} diff --git a/.trellis/tasks/09-21-radar-weighted-score-rank-change/validation.md b/.trellis/tasks/09-21-radar-weighted-score-rank-change/validation.md new file mode 100644 index 0000000..b620564 --- /dev/null +++ b/.trellis/tasks/09-21-radar-weighted-score-rank-change/validation.md @@ -0,0 +1,71 @@ +# 实施与验证记录 + +2026-09-21,本地实施完成。实现轮结束时尚未提交;本轮用户已授权提交并推送。生产数据库未迁移、未重算,未部署。 + +## 已实现 + +- 单日/波段流入率增加镜像加权评分列,保留 1 位小数,原始率独立展示。左右列头分别排序,加载该侧完整候选榜后应用排序,NULL 始终置后。 +- 新算法版本为 `zhixing_ratio_weighted_v2`、`zhixing_swing_weighted_v2`。全池平均名次百分位与 5 日成交额权重在后端计算,波段使用 MA3/MA10 分别排名后各占 50%。 +- 排名变化提供波段率/单日率/单日额、近 1–5 天、涨跌幅、另外两项指标的变化、镜像双榜和随基准更新的标题。URL 无显式值时默认波段率、近 1 天。 +- 对比使用保存的交易日历和同版本名次。缺发布、缺板块、历史不足或版本不匹配时保持未知,不跳过缺失交易日、不伪造零。 +- 迁移 `0010_radar_weighted_score` 增加 nullable NUMERIC(28,12) 和有限值约束。全部排名读写路径、HTTP、前端解析与详情历史同步扩展。 +- `sector-radar-build --recompute` 只读取已有聚合事实与源快照,按日期生成新发布,保留原发布。当前名次、历史和详情使用相同版本。 + +## 实际验证 + +| 检查 | 结果 | +| --- | --- | +| 后端雷达领域、构建、读取、HTTP、PostgreSQL 集成 | 110 passed | +| 前端雷达 API、query、详情、页面 | 74 passed | +| 后端 Ruff format/check | 通过 | +| 本次后端模块、测试和迁移 Pyright | 0 errors | +| 前端格式、ESLint、TypeScript 和生产 build | 通过;Vite 仍提示既有大 chunk | +| migration 0010 → 0009 → 0010 | 288 条原始排名保留,新增评分为 NULL | +| 真实 PostgreSQL 重算 | 保留旧发布、评分在各读取路径一致、详情来源完整、失败不替换 last-good | +| CLI 完整重算 | 31 日 × 400 个合成板块;首次 success,重复 31 日均 unchanged | +| 浏览器 2 类型 × 3 基准 × 5 天数 | 30 组合通过;首行变化与 API 一致,刷新恢复 URL 选择 | +| 浏览器评分排序 | 左榜完整 31 个、右榜完整 30 个候选;升降序正确且两侧独立 | +| 浏览器窄屏与缺历史 | 390px 页面无整页横向溢出;960px 表格可横向滚动;空历史文案完整可见 | +| Git diff/context manifests | diff --check 与 7/9 条上下文清单校验通过 | + +数据库验证均在单独创建的本地 PostgreSQL 17 容器执行,未连接生产数据库。浏览器预览为合成测试数据;`research/seed_local_preview.py` 限制只能写入本地测试库。控制台仅观察到既有 favicon.ico 404,没有本次页面异常。 + +## 全仓基线问题 + +完整命令已执行,但不能报告全仓检查通过。已用 `git archive HEAD`(实施前 HEAD `669e89d`)及相同依赖独立复现: + +- 后端全量 Pyright:14 个相同错误,位于 selection 的 `chart.py`、`gold_brick.py`、`test_run.py`。 +- 后端全量 pytest:本次 230 passed / 6 failed,修改前 221 passed / 同样 6 failed。既有 market_data 集成测试调用不存在的 Connection.executemany;旧迁移测试日志配置导致后续 5 个日志断言失败。雷达相关 110 项在独立运行时全部通过。 +- 前端 `pnpm check` 的格式、lint 和类型均通过;全量 Vitest 为 148 passed / 4 failed。修改前 selection 页面同样 4 个执行状态弹窗测试失败(其余 20 个通过)。 + +这些既有选股/行情问题未在本任务扩展修复。前端 `.prettierignore` 补充已被 Git 忽略的 `.playwright-cli`,避免旧浏览器快照影响格式检查。 + +## 生产应用步骤(尚未执行) + +生产 Compose 的完整构建、迁移、重算顺序见 `docs/market-data-sync.md`。以下为直接使用 Python 环境时的等价命令。 + +在目标环境确认 `ZHIXING_DATABASE_URL` 指向正确数据库,先执行迁移,再上线代码并重算现有日期区间: + +```bash +cd zhixing-server +uv run alembic upgrade head +uv run sector-radar-build --recompute --start-date 2026-08-28 --end-date 2026-09-21 +``` + +日期应覆盖需要展示的已有历史;命令只处理区间内已有成功发布。单日评分需要连续 5 个交易日成交额,波段评分需要连续 10 个交易日输入;要比较近 5 日波段名次,还需要相应更早窗口。缺失日期保持未知,重算不补造原始行情或成员数据。 + +新列允许旧行为 NULL,旧发布保留追溯。回退需切回保留的旧版本发布;不能在仍运行新代码时直接删列。当前代码库没有自动生产回退动作,本任务没有执行任何生产状态变更。 + +## 预览证据 + +- `output/playwright/radar-rank-change-desktop.png` +- `output/playwright/radar-swing-desktop.png` +- `output/playwright/radar-ratio-desktop.png` +- `output/playwright/radar-rank-change-mobile.png` +- `output/playwright/radar-missing-history-mobile.png` + +本地预览:`http://127.0.0.1:15173/sector-radar?view=rank_change&tradeDate=2026-09-21`。前端端口 15173,后端端口 18081,数据库容器 `zhixing-radar-score-test`(本地端口 55439);均用于隔离验收。 + +## 提交前复核 + +本轮独立只读检查由 `/root/radar_commit_review` 完成,未发现新增正确性或部署阻塞;范围内 Ruff/Pyright 和 diff 检查通过。Compose Job 的 entrypoint 与新版 CLI 参数已核对,并使用无生产凭据的配置通过 `docker compose config --quiet`。线上操作说明已补入 `docs/market-data-sync.md`;本轮只执行 Git 提交与推送,不执行线上迁移、重算或部署。 diff --git a/docs/market-data-sync.md b/docs/market-data-sync.md index 22ad10c..71dc247 100644 --- a/docs/market-data-sync.md +++ b/docs/market-data-sync.md @@ -123,3 +123,36 @@ docker compose -f docker-compose.prod.yml --profile jobs run --rm sector-radar-b `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、到达时点和首轮回填验证完成后再单独启用调度。 + +### 升级加权评分并重算历史 + +`--recompute` 只使用数据库中已有的板块聚合事实、交易日历和源快照,生成新的评分与排名发布;不调用 Tushare,不需要 token,不修改旧发布或历史成员。它与重新拉取上游的普通构建不同,不能和 `--retry-publication-id` 同用。 + +在生产项目目录中,将代码更新到包含加权评分的提交,沿用生产 `.env`,依次执行以下命令;每步成功后再执行下一步。Compose 为不同服务使用不同镜像名,因此除服务端和前端外,也要构建迁移与重算 Job 的新镜像。 + +```bash +# 构建新代码,暂不替换运行中的服务 +docker compose -f docker-compose.prod.yml --profile jobs build \ + server web migrate sector-radar-build + +# 先增加 nullable weighted_score 列(0010),再启动新服务 +docker compose -f docker-compose.prod.yml --profile jobs run --rm --no-deps migrate +docker compose -f docker-compose.prod.yml up -d server web + +# 使用新 Job 镜像,按交易日顺序重算已有历史 +docker compose -f docker-compose.prod.yml --profile jobs run --rm --no-deps sector-radar-build \ + --recompute --start-date 2026-08-28 --end-date 2026-09-21 +``` + +这里的日期是示例,应覆盖需要展示的已有历史,区间包含起止日。单日评分需要连续 5 个交易日成交额,波段评分需要连续 10 个交易日流入率;近 5 日波段排名变化还要求更早的对比日也有足够历史并已采用同一算法版本。首次升级建议从最早已有成功发布开始重算整个展示区间。缺少交易日、板块或可比版本时显示“—”,重算不会补造缺失数据。 + +输出的 `outcomes` 列出逐日结果;整个命令 `status` 为 `success` 或 `unchanged` 且退出码为 0 表示完成。相同输入重复运行返回 `unchanged`;失败会保留原 last-good,可修复原因后重跑同一区间。重算期间避免同时启动覆盖相同日期的普通构建或另一个重算任务。 + +以后只重算某一天,可复用已构建的新 Job 镜像: + +```bash +docker compose -f docker-compose.prod.yml --profile jobs run --rm --no-deps sector-radar-build \ + --recompute --trade-date 2026-09-21 +``` + +`--no-deps` 用于上述已手动完成迁移的流程,避免再次启动依赖服务。`sector-radar-build` 服务配置了同名 entrypoint,后面直接传 `--recompute` 等参数即可。 diff --git a/output/playwright/radar-missing-history-mobile.png b/output/playwright/radar-missing-history-mobile.png new file mode 100644 index 0000000..b5e1fec Binary files /dev/null and b/output/playwright/radar-missing-history-mobile.png differ diff --git a/output/playwright/radar-rank-change-desktop.png b/output/playwright/radar-rank-change-desktop.png new file mode 100644 index 0000000..dbbbae5 Binary files /dev/null and b/output/playwright/radar-rank-change-desktop.png differ diff --git a/output/playwright/radar-rank-change-mobile.png b/output/playwright/radar-rank-change-mobile.png new file mode 100644 index 0000000..501bfe7 Binary files /dev/null and b/output/playwright/radar-rank-change-mobile.png differ diff --git a/output/playwright/radar-ratio-desktop.png b/output/playwright/radar-ratio-desktop.png new file mode 100644 index 0000000..6325df4 Binary files /dev/null and b/output/playwright/radar-ratio-desktop.png differ diff --git a/output/playwright/radar-swing-desktop.png b/output/playwright/radar-swing-desktop.png new file mode 100644 index 0000000..008b6fd Binary files /dev/null and b/output/playwright/radar-swing-desktop.png differ diff --git a/zhixing-server/migrations/versions/0010_radar_weighted_score.py b/zhixing-server/migrations/versions/0010_radar_weighted_score.py new file mode 100644 index 0000000..d7b802e --- /dev/null +++ b/zhixing-server/migrations/versions/0010_radar_weighted_score.py @@ -0,0 +1,27 @@ +"""Keep weighted ranking scores separate from raw flow values.""" + +import sqlalchemy as sa +from alembic import op + +revision = "0010_radar_weighted_score" +down_revision = "0009_radar_sector_detail" +branch_labels = None +depends_on = None + + +def upgrade() -> None: + """Preserve existing publication rows with an unknown (NULL) score.""" + op.add_column( + "sector_radar_ranking", sa.Column("weighted_score", sa.Numeric(28, 12), nullable=True) + ) + op.create_check_constraint( + "ck_radar_weighted_score_finite", + "sector_radar_ranking", + "weighted_score IS NULL OR weighted_score NOT IN " + "('NaN'::numeric, 'Infinity'::numeric, '-Infinity'::numeric)", + ) + + +def downgrade() -> None: + """Remove the additional score projection without changing raw values.""" + op.drop_column("sector_radar_ranking", "weighted_score") 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 984d05e..e535e12 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 @@ -18,14 +18,10 @@ from zhixing_server.shared.request_coordinator import TushareSourceError from ..domain.facts import aggregate_sector_snapshot from ..domain.metrics import ( - AmountNetStrategy, MetricStrategy, - RatioTurnoverStrategy, - SwingEqualThreeToTenStrategy, ) from ..domain.models import ( MembershipStatus, - MetricObservation, PublicationStatus, RadarPublication, RankedMetric, @@ -50,7 +46,6 @@ from ..domain.persistence import ( StockFactRecord, ) from ..domain.ports import ActiveMoneyflowSource, SectorRadarSource -from ..domain.ranking import rank_metric_observations, with_rank_changes from ..domain.source import ( DailyRow, MoneyflowDcRow, @@ -65,6 +60,7 @@ from ..domain.source import ( SuspendRow, TradeCalendarRow, ) +from .scoring import calculate_rankings, calendar_rank_changes, default_strategies BuildOutcomeStatus = Literal["success", "partial", "failed", "locked", "unchanged"] SHANGHAI = ZoneInfo("Asia/Shanghai") @@ -194,14 +190,7 @@ class BuildSectorRadar: self.repository = repository self.coverage_threshold = coverage_threshold self.today = today or self.now_fn().astimezone(SHANGHAI).date() - self.strategies = tuple( - strategies - or ( - AmountNetStrategy(), - RatioTurnoverStrategy(), - SwingEqualThreeToTenStrategy(), - ) - ) + self.strategies = tuple(strategies) if strategies is not None else default_strategies() def execute(self, command: BuildSectorRadarCommand | None = None) -> BuildSummary: """Build each selected trade date sequentially for deterministic history.""" @@ -314,7 +303,13 @@ class BuildSectorRadar: publication_created = True reusable = self._reusable_sources(target.retry_publication_id) collected = self._collect(target.trade_date, publication_id, reusable) - input_hash = self._input_hash(collected.snapshots) + previous_publications = tuple( + item + for item in self.repository.load_history_publications(target.trade_date) + if item.target_trade_date < target.trade_date + and item.source_version == self.source_version + )[:9] + input_hash = self._input_hash(collected.snapshots, previous_publications) existing = self.repository.find_reusable_publication(target.trade_date, input_hash) if existing is not None: self.repository.discard_running_publication(publication_id) @@ -339,7 +334,9 @@ class BuildSectorRadar: input_hash=input_hash, ) aggregates = self._aggregate(collected) - rankings = self._rank(target.trade_date, aggregates) + rankings = self._rank( + target.trade_date, aggregates, collected.trading_dates, previous_publications + ) coverage = self._coverage(collected.stock_facts) membership_complete = all( item.status is MembershipStatus.AVAILABLE for item in collected.memberships @@ -638,6 +635,7 @@ class BuildSectorRadar: ) return _CollectedInputs( target_trade_date=target, + trading_dates=tuple(sorted({row.cal_date for row in calendar.rows if row.is_open})), indices=concepts.rows + industries.rows, snapshots=snapshots, membership_snapshots=members.snapshots, @@ -692,24 +690,29 @@ class BuildSectorRadar: return aggregates def _rank( - self, target: date, aggregates: Sequence[SectorDailyAggregate] + self, + target: date, + aggregates: Sequence[SectorDailyAggregate], + trading_dates: Sequence[date], + previous_publications: Sequence[RadarPublication], ) -> tuple[RankedMetric, ...]: - history = tuple(self.repository.load_daily_aggregate_history(target, limit_dates=9)) - observations: list[MetricObservation] = [] - for current in aggregates: - sector_history = tuple( - item - for item in history - if (item.sector_type, item.sector_code) - == (current.sector_type, current.sector_code) - ) + (current,) - observations.extend( - strategy.evaluate(sector_history, target) for strategy in self.strategies + history = tuple( + record.aggregate + for record in self.repository.load_publication_aggregates( + tuple(item.publication_id for item in previous_publications) ) - current_rankings = rank_metric_observations(observations) - previous = self.repository.load_previous_rankings(target, limit_dates=5) - history_by_days = {days: rankings for days, (_, rankings) in enumerate(previous, start=1)} - return with_rank_changes(current_rankings, history_by_days) + ) + current_rankings = calculate_rankings( + target, aggregates, history, trading_dates, self.strategies + ) + dates_by_id = { + item.publication_id: item.target_trade_date for item in previous_publications[:5] + } + previous = { + dates_by_id[key]: rows + for key, rows in self.repository.load_publication_rankings(tuple(dates_by_id)) + } + return calendar_rank_changes(current_rankings, previous, trading_dates, target) @staticmethod def _coverage(stock_facts: Sequence[StockFactRecord]) -> Decimal: @@ -750,12 +753,15 @@ class BuildSectorRadar: groups.append(PublicationSourceGroup.MONEYFLOW_DC) return tuple(groups) - def _input_hash(self, snapshots: Sequence[SourceSnapshot]) -> str: + def _input_hash( + self, snapshots: Sequence[SourceSnapshot], previous: Sequence[RadarPublication] = () + ) -> str: payload = json.dumps( { "snapshot_ids": sorted(snapshot.snapshot_id for snapshot in snapshots), "metric_versions": sorted(strategy.metric_version for strategy in self.strategies), "normalizer": "zhixing_stock_fact_v2", + "history": [(item.publication_id, item.input_hash) for item in previous], }, sort_keys=True, separators=(",", ":"), @@ -785,6 +791,7 @@ class BuildSectorRadar: @dataclass(frozen=True, slots=True) class _CollectedInputs: target_trade_date: date + trading_dates: tuple[date, ...] indices: tuple[SectorIndexRow, ...] snapshots: tuple[SourceSnapshot, ...] membership_snapshots: tuple[SourceSnapshot, ...] diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/application/details.py b/zhixing-server/src/zhixing_server/modules/sector_radar/application/details.py index f8c9e7f..b4d9ba7 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/application/details.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/application/details.py @@ -7,7 +7,6 @@ from dataclasses import dataclass from datetime import date from decimal import Decimal -from ..domain.metrics import AmountNetStrategy, RatioTurnoverStrategy, SwingEqualThreeToTenStrategy from ..domain.models import MetricKind, RadarPublication, RankedMetric, RankSide, SectorType from ..domain.normalize import is_current_listed_stock from ..domain.persistence import HistoricalRanking, PublicationSourceGroup, SectorRadarRepository @@ -20,12 +19,7 @@ from ..domain.source import ( StockBasicRow, TradeCalendarRow, ) - -_METRIC_VERSIONS = { - MetricKind.AMOUNT: AmountNetStrategy.metric_version, - MetricKind.RATIO: RatioTurnoverStrategy.metric_version, - MetricKind.SWING: SwingEqualThreeToTenStrategy.metric_version, -} +from ..domain.weighted import resolve_metric_version @dataclass(frozen=True, slots=True) @@ -50,6 +44,7 @@ class HistoryMetric: missing: bool = True in_top: bool = False in_bottom: bool = False + weighted_score: Decimal | None = None @dataclass(frozen=True, slots=True) @@ -134,7 +129,8 @@ class SectorDetail: class _RankHistory: """Index one day's selected ranks once, retaining complete versioned pool counts.""" - def __init__(self, records: Sequence[HistoricalRanking]) -> None: + def __init__(self, records: Sequence[HistoricalRanking], versions: Sequence[str] = ()) -> None: + self.versions = tuple(versions) or tuple(record.metric_version for record in records) self.pools: dict[tuple[SectorType, MetricKind, str], int] = {} self.rankings: dict[tuple[SectorType, str, MetricKind, str], RankedMetric] = {} self.names: dict[tuple[SectorType, str], str] = {} @@ -158,7 +154,9 @@ class _RankHistory: ) def metric(self, sector_type: SectorType, sector_code: str, kind: MetricKind) -> HistoryMetric: - version = _METRIC_VERSIONS[kind] + version = resolve_metric_version(kind, self.versions) + if version is None: + return HistoryMetric() size = self.pools.get((sector_type, kind, version), 0) row = self.rankings.get((sector_type, sector_code, kind, version)) return _metric_from_ranking(row, size) @@ -187,9 +185,9 @@ class ReadRadarDetails: ): by_id.setdefault(record.publication_id, []).append(record) rows = { - day: _RankHistory(by_id.get(item.publication_id, ())) + day: _RankHistory(by_id.get(item.publication_id, ()), publication.metric_versions) if item.source_version == publication.source_version - else _RankHistory(()) + else _RankHistory((), publication.metric_versions) for day, item in publications.items() } sources = self.repository.load_publication_rows( @@ -209,7 +207,11 @@ class ReadRadarDetails: | set(publications) )[-30:] ) - return publications, {day: rows.get(day, _RankHistory(())) for day in dates}, dates + return ( + publications, + {day: rows.get(day, _RankHistory((), publication.metric_versions)) for day in dates}, + dates, + ) def ranking_extras( self, publication: RadarPublication, rankings: Sequence[RankedMetric], side: RankSide @@ -438,12 +440,13 @@ def metric_at( rows: Sequence[RankedMetric], sector_type: SectorType, sector_code: str, kind: MetricKind ) -> HistoryMetric: """Select compatible rank values; pool thresholds use each day's actual percentile.""" + version = resolve_metric_version(kind, (row.observation.metric_version for row in rows)) pool = [ row for row in rows if row.observation.sector_type is sector_type and row.observation.metric_kind is kind - and row.observation.metric_version == _METRIC_VERSIONS[kind] + and row.observation.metric_version == version ] size = sum(row.rank_position is not None for row in pool) row = next((row for row in pool if row.observation.sector_code == sector_code), None) @@ -463,6 +466,7 @@ def _metric_from_ranking(row: RankedMetric | None, size: int) -> HistoryMetric: row.rank_position is None, percentile is not None and percentile >= 90, percentile is not None and percentile <= 10, + row.observation.weighted_score, ) diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/application/read.py b/zhixing-server/src/zhixing_server/modules/sector_radar/application/read.py index 8622f44..7c81399 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/application/read.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/application/read.py @@ -27,7 +27,13 @@ from ..domain.persistence import ( StockMembershipEntry, ) from ..domain.ranking import select_percentile_side, select_rank_change_side +from ..domain.weighted import ( + RatioWeightedStrategy, + SwingWeightedStrategy, + resolve_metric_version, +) from .details import RankingExtras, ReadRadarDetails +from .scoring import calendar_rank_changes, comparison_dates, publication_trading_dates ReadStatus = Literal["success", "no_data"] @@ -62,7 +68,7 @@ class RadarQuery: trade_date: date | None = None sector_type: SectorType = SectorType.CONCEPT view: RadarView = RadarView.AMOUNT - rank_change_metric: MetricKind = MetricKind.AMOUNT + rank_change_metric: MetricKind = MetricKind.SWING rank_change_days: int = 1 side: RankSide = RankSide.ALL search: str | None = None @@ -108,6 +114,10 @@ class RankingPage: rows: tuple[RankedMetric, ...] total: int extras: dict[str, RankingExtras] = field(default_factory=lambda: dict[str, RankingExtras]()) + rank_change_values: dict[str, dict[MetricKind, int | None]] = field( + default_factory=lambda: dict[str, dict[MetricKind, int | None]]() + ) + comparison_trade_date: date | None = None @dataclass(frozen=True, slots=True) @@ -175,27 +185,54 @@ class SectorMembersSnapshot: _METRIC_DEFINITIONS = { - MetricKind.AMOUNT: RadarMetricDefinition( + AmountNetStrategy.metric_version: RadarMetricDefinition( metric_kind=MetricKind.AMOUNT, metric_version=AmountNetStrategy.metric_version, label="主力净流入(知行独立实现)", unit=MetricUnit.CNY_100M, ), - MetricKind.RATIO: RadarMetricDefinition( + RatioTurnoverStrategy.metric_version: RadarMetricDefinition( metric_kind=MetricKind.RATIO, metric_version=RatioTurnoverStrategy.metric_version, label="主力净流入/成交额(知行独立实现)", unit=MetricUnit.RATIO, ), - MetricKind.SWING: RadarMetricDefinition( + SwingEqualThreeToTenStrategy.metric_version: RadarMetricDefinition( metric_kind=MetricKind.SWING, metric_version=SwingEqualThreeToTenStrategy.metric_version, label="3—10 日等权资金率(知行独立实现)", unit=MetricUnit.RATIO, ), + RatioWeightedStrategy.metric_version: RadarMetricDefinition( + metric_kind=MetricKind.RATIO, + metric_version=RatioWeightedStrategy.metric_version, + label="单日流入率", + unit=MetricUnit.RATIO, + disclaimer="基于知行数据的资金排名与成交额加权评分", + ), + SwingWeightedStrategy.metric_version: RadarMetricDefinition( + metric_kind=MetricKind.SWING, + metric_version=SwingWeightedStrategy.metric_version, + label="波段流入率", + unit=MetricUnit.RATIO, + disclaimer="3 日与 10 日资金率排名各占 50%,再按成交额加权", + ), } +def metric_definition( + kind: MetricKind, publication: RadarPublication | None +) -> RadarMetricDefinition: + """Resolve a publication's actual algorithm instead of hiding legacy rows.""" + current = { + MetricKind.AMOUNT: AmountNetStrategy.metric_version, + MetricKind.RATIO: RatioWeightedStrategy.metric_version, + MetricKind.SWING: SwingWeightedStrategy.metric_version, + }[kind] + version = resolve_metric_version(kind, publication.metric_versions) if publication else None + return _METRIC_DEFINITIONS[version or current] + + class ReadSectorRadar: """Hide last-good selection, ranking filters, search, and pagination.""" @@ -219,18 +256,45 @@ class ReadSectorRadar: if query.view is RadarView.RANK_CHANGE else MetricKind(query.view.value) ) - definition = _METRIC_DEFINITIONS[metric_kind] publication = ( self.repository.get_successful_publication(query.trade_date) if query.trade_date is not None else self.repository.get_last_good_publication() ) + definition = metric_definition(metric_kind, publication) if publication is None: return RankingPage("no_data", query, None, definition, (), 0) + all_rows = tuple(self.repository.load_rankings(publication.publication_id)) + compared_date = None + if query.view is RadarView.RANK_CHANGE: + calendar = publication_trading_dates( + self.repository, publication.publication_id, publication.target_trade_date + ) + previous_publications = tuple( + item + for item in self.repository.load_history_publications(publication.target_trade_date) + if item.target_trade_date < publication.target_trade_date + and item.source_version == publication.source_version + )[:5] + dates_by_id = { + item.publication_id: item.target_trade_date for item in previous_publications + } + previous = { + dates_by_id[publication_id]: rows + for publication_id, rows in self.repository.load_publication_rankings( + tuple(dates_by_id) + ) + } + all_rows = calendar_rank_changes( + all_rows, previous, calendar, publication.target_trade_date + ) + compared_date = comparison_dates(calendar, publication.target_trade_date)[ + query.rank_change_days + ] metric_rows = tuple( row - for row in self.repository.load_rankings(publication.publication_id) + for row in all_rows if row.observation.sector_type is query.sector_type and row.observation.metric_kind is metric_kind and row.observation.metric_version == definition.metric_version @@ -276,18 +340,32 @@ class ReadSectorRadar: or search in row.observation.sector_name.casefold() ) start = (query.page - 1) * query.page_size + page_rows = searched[start : start + query.page_size] + selected_codes = {row.observation.sector_code for row in page_rows} + changes: dict[str, dict[MetricKind, int | None]] = {} + for row in all_rows: + observation = row.observation + if ( + observation.sector_type is query.sector_type + and observation.sector_code in selected_codes + and observation.metric_version + == resolve_metric_version(observation.metric_kind, publication.metric_versions) + ): + changes.setdefault(observation.sector_code, {})[observation.metric_kind] = ( + row.rank_change(query.rank_change_days) + ) return RankingPage( status="success", query=query, publication=publication, definition=definition, - rows=searched[start : start + query.page_size], + rows=page_rows, total=len(searched), extras=ReadRadarDetails(self.repository).ranking_extras( - publication, searched[start : start + query.page_size], query.side - ) - if query.view in {RadarView.AMOUNT, RadarView.RATIO} - else {}, + publication, page_rows, query.side + ), + rank_change_values=changes, + comparison_trade_date=compared_date, ) def stock_membership(self, query: StockSectorQuery) -> StockSectorMembership: diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/application/recompute.py b/zhixing-server/src/zhixing_server/modules/sector_radar/application/recompute.py new file mode 100644 index 0000000..e086bf6 --- /dev/null +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/application/recompute.py @@ -0,0 +1,268 @@ +"""Rescore immutable local inputs without refetching or altering source history.""" + +from __future__ import annotations + +import hashlib +import json +from collections.abc import Callable, Sequence +from dataclasses import replace +from datetime import UTC, date, datetime +from decimal import Decimal +from uuid import uuid4 + +from ..domain.models import PublicationStatus, RadarPublication, RankedMetric +from ..domain.persistence import ( + DailyAggregateRecord, + PublicationSourceRecord, + RankingRecord, + SectorRadarRepository, +) +from .build import BuildDateOutcome, BuildSectorRadarCommand, BuildSummary +from .scoring import ( + calculate_rankings, + calendar_rank_changes, + default_strategies, + publication_trading_dates, +) + + +class RecomputeSectorRadar: + """Publish a new algorithm revision from pinned, already-normalized facts. + + No source adapter is accepted: current stock membership cannot accidentally + replace a historical snapshot. Failed rebuilds leave the old last-good intact. + """ + + def __init__( + self, + repository: SectorRadarRepository, + now_fn: Callable[[], datetime] = lambda: datetime.now(UTC), + ) -> None: + self.repository = repository + self.now_fn = now_fn + self.versions = tuple(strategy.metric_version for strategy in default_strategies()) + + def execute(self, command: BuildSectorRadarCommand) -> BuildSummary: + """Rescore successful dates oldest first, stopping if one date fails. + + The requested range is inclusive. A single date or no date selects one + existing publication. All aggregate inputs are pinned before writes. + """ + if command.retry_publication_id is not None: + raise ValueError("offline recompute does not accept a source retry") + available = sorted(self.repository.list_successful_dates()) + if command.trade_date is not None: + targets = [command.trade_date] if command.trade_date in available else [] + elif command.start_date is not None and command.end_date is not None: + targets = [day for day in available if command.start_date <= day <= command.end_date] + else: + targets = available[-1:] + if not targets: + return BuildSummary( + ( + BuildDateOutcome( + command.trade_date or command.start_date or self.now_fn().date(), + "failed", + None, + Decimal(0), + 0, + 0, + "no_local_publication", + "no successful local input publication", + ), + ) + ) + needed = [day for day in available if day < targets[0]][-9:] + targets + pinned: dict[date, RadarPublication] = {} + for day in needed: + publication = self.repository.get_successful_publication(day) + if publication is None: + raise ValueError("selected local publication disappeared") + pinned[day] = publication + inputs = self.repository.load_publication_aggregates( + tuple(item.publication_id for item in pinned.values()) + ) + by_publication: dict[str, list[DailyAggregateRecord]] = {} + for record in inputs: + by_publication.setdefault(record.publication_id, []).append(record) + existing_rankings = dict( + self.repository.load_publication_rankings( + tuple(item.publication_id for item in pinned.values()) + ) + ) + ranks_by_day = { + day: existing_rankings.get(item.publication_id, ()) for day, item in pinned.items() + } + outcomes: list[BuildDateOutcome] = [] + for target in targets: + base = pinned[target] + current = by_publication.get(base.publication_id, []) + prior_dates = [day for day in needed if day < target][-9:] + history = [ + record + for day in prior_dates + if pinned[day].source_version == base.source_version + for record in by_publication.get(pinned[day].publication_id, []) + ] + previous = { + day: ranks_by_day[day] + for day in prior_dates + if pinned[day].source_version == base.source_version + } + outcome, rankings = self._rescore(base, current, history, previous) + outcomes.append(outcome) + if outcome.status in {"failed", "locked"}: + break + ranks_by_day[target] = rankings + return BuildSummary(tuple(outcomes)) + + def _rescore( + self, + base: RadarPublication, + current: Sequence[DailyAggregateRecord], + history: Sequence[DailyAggregateRecord], + previous: dict[date, Sequence[RankedMetric]], + ) -> tuple[BuildDateOutcome, Sequence[RankedMetric]]: + target = base.target_trade_date + pending: RadarPublication | None = None + try: + with self.repository.advisory_lock(target) as acquired: + if not acquired: + return BuildDateOutcome(target, "locked", None, Decimal(0), 0, 0), () + if not current: + raise ValueError("local aggregate input is missing") + calendar = publication_trading_dates(self.repository, base.publication_id, target) + rankings = calendar_rank_changes( + calculate_rankings( + target, + [item.aggregate for item in current], + [item.aggregate for item in history], + calendar, + ), + previous, + calendar, + target, + ) + sources = tuple(self.repository.load_publication_sources(base.publication_id)) + digest = self._fingerprint(base, current, history, sources, rankings) + reusable = self.repository.find_reusable_publication(target, digest) + if reusable is not None and reusable.status is PublicationStatus.SUCCESS: + return BuildDateOutcome( + target, + "unchanged", + reusable.publication_id, + reusable.coverage, + len(current), + len(rankings), + ), rankings + started = self.now_fn() + self.repository.recover_running_publications(target, finished_at=started) + pending = replace( + base, + publication_id=f"radar-{target:%Y%m%d}-rescore-{uuid4().hex[:20]}", + status=PublicationStatus.RUNNING, + metric_versions=self.versions, + input_hash=None, + started_at=started, + finished_at=None, + error_summary=None, + ) + self.repository.create_publication(pending) + self.repository.save_publication_sources( + replace(item, publication_id=pending.publication_id, refresh_on_retry=False) + for item in sources + ) + finished = replace( + pending, + status=PublicationStatus.SUCCESS, + input_hash=digest, + finished_at=self.now_fn(), + ) + self.repository.finalize_publication( + finished, + memberships=(), + stock_facts=(), + daily_aggregates=( + replace(item, publication_id=pending.publication_id) for item in current + ), + rankings=(RankingRecord(pending.publication_id, row) for row in rankings), + ) + return BuildDateOutcome( + target, + "success", + pending.publication_id, + base.coverage, + len(current), + len(rankings), + ), rankings + except (ValueError, RuntimeError, ArithmeticError) as exc: + if pending is not None: + saved = self.repository.get_publication(pending.publication_id) + if saved is not None and saved.status is PublicationStatus.RUNNING: + self.repository.finish_publication( + replace( + pending, + status=PublicationStatus.FAILED, + finished_at=self.now_fn(), + error_summary="offline_recompute_failed", + ) + ) + return BuildDateOutcome( + target, + "failed", + pending.publication_id if pending else None, + Decimal(0), + 0, + 0, + type(exc).__name__, + "local radar rescoring failed", + ), () + + def _fingerprint( + self, + base: RadarPublication, + current: Sequence[DailyAggregateRecord], + history: Sequence[DailyAggregateRecord], + sources: Sequence[PublicationSourceRecord], + rankings: Sequence[RankedMetric], + ) -> str: + """Hash source content, not derived publication IDs, so reruns are idempotent.""" + payload = { + "source_version": base.source_version, + "universe_version": base.universe_version, + "metric_versions": self.versions, + "snapshots": sorted(item.snapshot.snapshot_id for item in sources), + "inputs": [ + ( + item.aggregate.trade_date.isoformat(), + item.aggregate.sector_type.value, + item.aggregate.sector_code, + str(item.aggregate.net_amount_yuan), + str(item.aggregate.turnover_yuan), + item.aggregate.member_count, + item.aggregate.valid_sample_count, + str(item.aggregate.membership_coverage), + str(item.aggregate.moneyflow_coverage), + str(item.pct_change), + item.leading_code, + ) + for item in sorted( + (*history, *current), + key=lambda item: ( + item.aggregate.trade_date, + item.aggregate.sector_type, + item.aggregate.sector_code, + ), + ) + ], + "rank_changes": [ + ( + row.observation.sector_type.value, + row.observation.sector_code, + row.observation.metric_version, + [(change.days, change.value) for change in row.rank_changes], + ) + for row in rankings + ], + } + return hashlib.sha256(json.dumps(payload, sort_keys=True).encode()).hexdigest() diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/application/scoring.py b/zhixing-server/src/zhixing_server/modules/sector_radar/application/scoring.py new file mode 100644 index 0000000..31b28be --- /dev/null +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/application/scoring.py @@ -0,0 +1,98 @@ +"""Shared publication scoring and strict trading-session comparisons.""" + +from __future__ import annotations + +from collections import defaultdict +from collections.abc import Mapping, Sequence +from datetime import date + +from ..domain.metrics import AmountNetStrategy, MetricStrategy +from ..domain.models import MetricObservation, RankedMetric, SectorDailyAggregate, SectorType +from ..domain.persistence import PublicationSourceGroup, SectorRadarRepository +from ..domain.ranking import rank_metric_observations, with_rank_changes +from ..domain.source import TradeCalendarRow +from ..domain.weighted import ( + FlowFeatures, + RatioWeightedStrategy, + SwingWeightedStrategy, + attach_weighted_scores, + calendar_sector_history, + flow_features, +) + + +def default_strategies() -> tuple[MetricStrategy, ...]: + """Return versioned strategies shared by online builds and offline rescoring.""" + return (AmountNetStrategy(), RatioWeightedStrategy(), SwingWeightedStrategy()) + + +def publication_trading_dates( + repository: SectorRadarRepository, + publication_id: str, + target: date, +) -> tuple[date, ...]: + """Read the exact publication's calendar without querying an external provider.""" + sources = repository.load_publication_rows(publication_id, (PublicationSourceGroup.CALENDAR,)) + return tuple( + sorted( + { + parsed.cal_date + for row in sources.get(PublicationSourceGroup.CALENDAR, ()) + for parsed in (TradeCalendarRow.from_mapping(row),) + if parsed.is_open and parsed.cal_date <= target + } + ) + ) + + +def calculate_rankings( + target: date, + aggregates: Sequence[SectorDailyAggregate], + history: Sequence[SectorDailyAggregate], + trading_dates: Sequence[date], + strategies: Sequence[MetricStrategy] | None = None, +) -> tuple[RankedMetric, ...]: + """Evaluate raw features, score full pools, then generate authoritative ranks.""" + histories: defaultdict[tuple[SectorType, str], dict[date, SectorDailyAggregate]] = defaultdict( + dict + ) + for row in history: + if row.trade_date >= target: + continue + days = histories[(row.sector_type, row.sector_code)] + if row.trade_date in days: + raise ValueError("aggregate history must contain unique sector dates") + days[row.trade_date] = row + observations: list[MetricObservation] = [] + features: dict[tuple[SectorType, str], FlowFeatures] = {} + selected = default_strategies() if strategies is None else strategies + for current in aggregates: + if current.trade_date != target: + raise ValueError("current aggregates must match the target date") + key = (current.sector_type, current.sector_code) + sector_history = calendar_sector_history(current, histories[key], trading_dates) + features[key] = flow_features(sector_history, target) + observations.extend(strategy.evaluate(sector_history, target) for strategy in selected) + return rank_metric_observations(attach_weighted_scores(observations, features)) + + +def comparison_dates(trading_dates: Sequence[date], target: date) -> dict[int, date | None]: + """Resolve actual prior trading sessions; missing calendar history stays unknown.""" + previous = sorted({day for day in trading_dates if day < target}, reverse=True) + return {days: previous[days - 1] if len(previous) >= days else None for days in range(1, 6)} + + +def calendar_rank_changes( + current: Sequence[RankedMetric], + previous: Mapping[date, Sequence[RankedMetric]], + trading_dates: Sequence[date], + target: date, +) -> tuple[RankedMetric, ...]: + """Compare matching versions on exact calendar dates, without skipping failed days.""" + return with_rank_changes( + current, + { + days: previous.get(day, ()) if day is not None else () + for days, day in comparison_dates(trading_dates, target).items() + }, + ) diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/metrics.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/metrics.py index c665322..32266ff 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/metrics.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/metrics.py @@ -33,22 +33,24 @@ class MetricStrategy(Protocol): ... -def _target_aggregate( +def target_aggregate( history: Iterable[SectorDailyAggregate], target_trade_date: date ) -> SectorDailyAggregate: + """Require one target row; missing or duplicate targets are invalid inputs.""" matches = tuple(row for row in history if row.trade_date == target_trade_date) if len(matches) != 1: raise ValueError("history must contain exactly one target-date aggregate") return matches[0] -def _quality(row: SectorDailyAggregate) -> MetricQuality: +def aggregate_quality(row: SectorDailyAggregate) -> MetricQuality: + """Retain limited-sample status for incomplete membership or flow coverage.""" if row.valid_sample_count < 5 or row.membership_coverage < 1 or row.moneyflow_coverage < 1: return MetricQuality.AVAILABLE_LIMITED_SAMPLE return MetricQuality.AVAILABLE -def _observation( +def make_metric_observation( row: SectorDailyAggregate, *, metric_kind: MetricKind, @@ -57,6 +59,7 @@ def _observation( value: Decimal | None, quality: MetricQuality | None = None, ) -> MetricObservation: + """Carry source coverage into a metric and mark unknown raw values unavailable.""" return MetricObservation( trade_date=row.trade_date, sector_type=row.sector_type, @@ -72,7 +75,7 @@ def _observation( if value is None else quality if quality is not None - else _quality(row) + else aggregate_quality(row) ), member_count=row.member_count, valid_sample_count=row.valid_sample_count, @@ -95,9 +98,9 @@ class AmountNetStrategy: ) -> MetricObservation: """Return the target net amount; missing moneyflow remains unavailable.""" - row = _target_aggregate(history, target_trade_date) + row = target_aggregate(history, target_trade_date) value = None if row.net_amount_yuan is None else row.net_amount_yuan / Decimal("100000000") - return _observation( + return make_metric_observation( row, metric_kind=self.metric_kind, metric_version=self.metric_version, @@ -120,7 +123,7 @@ class RatioTurnoverStrategy: ) -> MetricObservation: """Return a ratio only when numerator and positive denominator exist.""" - row = _target_aggregate(history, target_trade_date) + row = target_aggregate(history, target_trade_date) value = None if ( row.net_amount_yuan is not None @@ -128,7 +131,7 @@ class RatioTurnoverStrategy: and row.turnover_yuan > 0 ): value = row.net_amount_yuan / row.turnover_yuan - return _observation( + return make_metric_observation( row, metric_kind=self.metric_kind, metric_version=self.metric_version, @@ -156,7 +159,7 @@ class SwingEqualThreeToTenStrategy: """Calculate eight complete trading-day windows ending at the target.""" rows = tuple(sorted(history, key=lambda row: row.trade_date)) - target = _target_aggregate(rows, target_trade_date) + target = target_aggregate(rows, target_trade_date) eligible = tuple(row for row in rows if row.trade_date <= target_trade_date) if any( (row.sector_type, row.sector_code) != (target.sector_type, target.sector_code) @@ -190,11 +193,11 @@ class SwingEqualThreeToTenStrategy: value = sum(window_ratios, start=Decimal(0)) / Decimal(8) quality = ( MetricQuality.AVAILABLE_LIMITED_SAMPLE - if any(_quality(row) is not MetricQuality.AVAILABLE for row in latest) + if any(aggregate_quality(row) is not MetricQuality.AVAILABLE for row in latest) else MetricQuality.AVAILABLE ) - return _observation( + return make_metric_observation( target, metric_kind=self.metric_kind, metric_version=self.metric_version, diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/models.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/models.py index 2944604..4cff2f3 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/models.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/models.py @@ -260,11 +260,15 @@ class MetricObservation: valid_sample_count: int membership_coverage: Decimal moneyflow_coverage: Decimal + weighted_score: Decimal | None = None def __post_init__(self) -> None: """Keep unavailable and finite-value states internally consistent.""" _validate_finite_decimal(self.value, "value") + _validate_finite_decimal(self.weighted_score, "weighted_score") + if self.weighted_score is not None and self.value is None: + raise ValueError("a weighted score requires an observed raw value") if self.value is None and self.quality is not MetricQuality.UNAVAILABLE: raise ValueError("a missing metric value must be unavailable") if self.value is not None and self.quality is MetricQuality.UNAVAILABLE: diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/persistence.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/persistence.py index cde8b2c..306b70e 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/persistence.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/persistence.py @@ -293,6 +293,12 @@ class SectorRadarRepository(Protocol): def save_daily_aggregates(self, records: Iterable[DailyAggregateRecord]) -> WriteCounts: ... + def load_publication_aggregates( + self, publication_ids: Sequence[str] + ) -> Sequence[DailyAggregateRecord]: + """Read exact revisions, retaining source details for offline rescoring.""" + ... + def create_publication(self, publication: RadarPublication) -> WriteCounts: ... def finish_publication(self, publication: RadarPublication) -> None: ... diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ranking.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ranking.py index d39bdc6..7010025 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ranking.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/ranking.py @@ -16,6 +16,7 @@ from .models import ( RankSide, SectorType, ) +from .weighted import ranking_value PoolKey = tuple[date, SectorType, MetricKind, str] SectorMetricKey = tuple[SectorType, str, MetricKind, str] @@ -40,9 +41,10 @@ def _sector_metric_key(observation: MetricObservation) -> SectorMetricKey: def _available_sort_key(observation: MetricObservation) -> tuple[Decimal, str]: - if observation.value is None: + value = ranking_value(observation) + if value is None: raise ValueError("unavailable observations cannot use the ranking sort key") - return (-observation.value, observation.sector_code) + return (-value, observation.sector_code) def rank_metric_observations( @@ -70,7 +72,7 @@ def rank_metric_observations( raise ValueError("a ranking pool must not contain duplicate sector codes") available = sorted( - (observation for observation in pool if observation.value is not None), + (observation for observation in pool if ranking_value(observation) is not None), key=_available_sort_key, ) pool_size = len(available) @@ -93,7 +95,7 @@ def rank_metric_observations( rank_percentile=None, ) for observation in sorted( - (observation for observation in pool if observation.value is None), + (observation for observation in pool if ranking_value(observation) is None), key=lambda observation: observation.sector_code, ) ) @@ -134,7 +136,7 @@ def select_percentile_side( sorted( pools[pool_key], key=lambda row: ( - row.observation.value if row.observation.value is not None else Decimal(0), + ranking_value(row.observation) or Decimal(0), row.observation.sector_code, ), ) diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/domain/weighted.py b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/weighted.py new file mode 100644 index 0000000..a42609e --- /dev/null +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/domain/weighted.py @@ -0,0 +1,246 @@ +"""Versioned flow scores: separate observable ratios from cross-sectional scores.""" + +from __future__ import annotations + +from collections import defaultdict +from collections.abc import Iterable, Mapping, Sequence +from dataclasses import dataclass, replace +from datetime import date +from decimal import Decimal + +from .metrics import ( + AmountNetStrategy, + RatioTurnoverStrategy, + SwingEqualThreeToTenStrategy, + aggregate_quality, + make_metric_observation, + target_aggregate, +) +from .models import ( + MetricKind, + MetricObservation, + MetricQuality, + MetricUnit, + SectorDailyAggregate, + SectorType, +) + +RATIO_WEIGHTED_VERSION = "zhixing_ratio_weighted_v2" +SWING_WEIGHTED_VERSION = "zhixing_swing_weighted_v2" +WEIGHTED_VERSIONS = frozenset((RATIO_WEIGHTED_VERSION, SWING_WEIGHTED_VERSION)) +SectorKey = tuple[SectorType, str] + + +def resolve_metric_version(kind: MetricKind, versions: Iterable[str]) -> str | None: + """Select a supported version actually present in an immutable publication.""" + available = set(versions) + candidates = { + MetricKind.AMOUNT: (AmountNetStrategy.metric_version,), + MetricKind.RATIO: (RATIO_WEIGHTED_VERSION, RatioTurnoverStrategy.metric_version), + MetricKind.SWING: (SWING_WEIGHTED_VERSION, SwingEqualThreeToTenStrategy.metric_version), + } + return next((version for version in candidates[kind] if version in available), None) + + +def ranking_value(observation: MetricObservation) -> Decimal | None: + """Never substitute a raw ratio for a missing score during v2 warm-up.""" + return ( + observation.weighted_score + if observation.metric_version in WEIGHTED_VERSIONS + else observation.value + ) + + +def daily_flow_ratio(row: SectorDailyAggregate) -> Decimal | None: + """Use yuan consistently, with the reviewed 100-yuan denominator offset.""" + if row.net_amount_yuan is None or row.turnover_yuan is None or row.turnover_yuan <= 0: + return None + return row.net_amount_yuan / (row.turnover_yuan + Decimal(100)) + + +@dataclass(frozen=True, slots=True) +class FlowFeatures: + """Past-only inputs before any ranking or page filtering is applied.""" + + daily_ratio: Decimal | None + short_ratio: Decimal | None + long_ratio: Decimal | None + liquidity_weight: Decimal | None + quality: MetricQuality + + +def flow_features(history: Iterable[SectorDailyAggregate], target: date) -> FlowFeatures: + """Extract complete 3/5/10-session windows; unknown inputs stay unknown. + + The caller supplies calendar-aligned history, including explicit missing + aggregates for calendar holes. Future observations never enter a window. + """ + rows = tuple( + sorted((row for row in history if row.trade_date <= target), key=lambda row: row.trade_date) + ) + current = target_aggregate(rows, target) + if len({row.trade_date for row in rows}) != len(rows): + raise ValueError("history must not contain duplicate trade dates") + if any( + (row.sector_type, row.sector_code) != (current.sector_type, current.sector_code) + for row in rows + ): + raise ValueError("history must contain exactly one sector identity") + + def mean_ratio(window: int) -> Decimal | None: + ratios = tuple(daily_flow_ratio(row) for row in rows[-window:]) + if len(ratios) < window or any(value is None for value in ratios): + return None + return sum((value for value in ratios if value is not None), Decimal(0)) / window + + turnovers = tuple(row.turnover_yuan for row in rows[-5:]) + weight = None + if len(turnovers) == 5 and all(value is not None and value > 0 for value in turnovers): + average = sum((value for value in turnovers if value is not None), Decimal(0)) / 5 + weight = (average + 1).log10() / 10 + quality = ( + MetricQuality.AVAILABLE_LIMITED_SAMPLE + if any(aggregate_quality(row) is not MetricQuality.AVAILABLE for row in rows[-10:]) + else MetricQuality.AVAILABLE + ) + return FlowFeatures(daily_flow_ratio(current), mean_ratio(3), mean_ratio(10), weight, quality) + + +class RatioWeightedStrategy: + """Keep the raw daily ratio; attach its score only after observing the full pool.""" + + metric_kind = MetricKind.RATIO + metric_version = RATIO_WEIGHTED_VERSION + unit = MetricUnit.RATIO + + def evaluate( + self, history: Iterable[SectorDailyAggregate], target_trade_date: date + ) -> MetricObservation: + """Return a daily observation without pretending that a score is a ratio.""" + rows = tuple(history) + current = target_aggregate(rows, target_trade_date) + features = flow_features(rows, target_trade_date) + return make_metric_observation( + current, + metric_kind=self.metric_kind, + metric_version=self.metric_version, + unit=self.unit, + value=features.daily_ratio, + quality=features.quality, + ) + + +class SwingWeightedStrategy: + """Combine 3- and 10-session simple means, equally weighted.""" + + metric_kind = MetricKind.SWING + metric_version = SWING_WEIGHTED_VERSION + unit = MetricUnit.RATIO + + def evaluate( + self, history: Iterable[SectorDailyAggregate], target_trade_date: date + ) -> MetricObservation: + """Expose the mean ratio independently of the two percentile ranks.""" + rows = tuple(history) + current = target_aggregate(rows, target_trade_date) + features = flow_features(rows, target_trade_date) + value = ( + (features.short_ratio + features.long_ratio) / 2 + if features.short_ratio is not None and features.long_ratio is not None + else None + ) + return make_metric_observation( + current, + metric_kind=self.metric_kind, + metric_version=self.metric_version, + unit=self.unit, + value=value, + quality=features.quality, + ) + + +def calendar_sector_history( + current: SectorDailyAggregate, + previous: Mapping[date, SectorDailyAggregate], + trading_dates: Sequence[date], +) -> tuple[SectorDailyAggregate, ...]: + """Represent missing trading sessions explicitly rather than skipping them.""" + dates = sorted({day for day in trading_dates if day <= current.trade_date})[-10:] + if not dates or dates[-1] != current.trade_date: + raise ValueError("the target must belong to the observed trading calendar") + return tuple( + current + if day == current.trade_date + else previous.get(day) + or replace( + current, + trade_date=day, + member_count=0, + valid_sample_count=0, + net_amount_yuan=None, + turnover_yuan=None, + membership_coverage=Decimal(0), + moneyflow_coverage=Decimal(0), + ) + for day in dates + ) + + +def _percentiles(values: Mapping[str, Decimal | None]) -> dict[str, Decimal]: + """Compute ascending average-rank percentiles, preserving ties deterministically.""" + ordered = sorted((value, code) for code, value in values.items() if value is not None) + result: dict[str, Decimal] = {} + start = 0 + while start < len(ordered): + end = start + 1 + while end < len(ordered) and ordered[end][0] == ordered[start][0]: + end += 1 + percentile = Decimal(start + 1 + end) / (2 * len(ordered)) + result.update((code, percentile) for _, code in ordered[start:end]) + start = end + return result + + +def attach_weighted_scores( + observations: Sequence[MetricObservation], + features: Mapping[SectorKey, FlowFeatures], +) -> tuple[MetricObservation, ...]: + """Score complete date/type/version pools, never a paginated subset. + + Missing history can leave a raw value available without a final rank. + The two swing percentiles are calculated separately before combining. + """ + pools: defaultdict[tuple[date, SectorType, str], list[MetricObservation]] = defaultdict(list) + for row in observations: + if row.metric_version in WEIGHTED_VERSIONS: + pools[(row.trade_date, row.sector_type, row.metric_version)].append(row) + scores: dict[tuple[date, SectorType, str, str], Decimal | None] = {} + for pool in pools.values(): + feature = {row.sector_code: features[(row.sector_type, row.sector_code)] for row in pool} + daily = _percentiles({code: item.daily_ratio for code, item in feature.items()}) + short = _percentiles({code: item.short_ratio for code, item in feature.items()}) + long = _percentiles({code: item.long_ratio for code, item in feature.items()}) + for row in pool: + code = row.sector_code + weight = feature[code].liquidity_weight + percentile = ( + daily.get(code) + if row.metric_kind is MetricKind.RATIO + else ((short[code] + long[code]) / 2 if code in short and code in long else None) + ) + scores[(row.trade_date, row.sector_type, code, row.metric_version)] = ( + Decimal(1000) * percentile * weight + if percentile is not None and weight is not None and row.value is not None + else None + ) + return tuple( + replace( + row, + weighted_score=scores[ + (row.trade_date, row.sector_type, row.sector_code, row.metric_version) + ], + ) + if row.metric_version in WEIGHTED_VERSIONS + else row + for row in observations + ) diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/memory.py b/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/memory.py index 4d7278e..2ddbbdf 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/memory.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/memory.py @@ -254,6 +254,26 @@ class InMemorySectorRadarRepository: ), ) + def load_publication_aggregates( + self, publication_ids: Sequence[str] + ) -> Sequence[DailyAggregateRecord]: + """Load exact immutable aggregate revisions, including detail fields.""" + wanted = set(publication_ids) + return tuple( + sorted( + ( + record + for record in self.daily_aggregates.values() + if record.publication_id in wanted + ), + key=lambda record: ( + record.aggregate.trade_date, + record.aggregate.sector_type, + record.aggregate.sector_code, + ), + ) + ) + def create_publication(self, publication: RadarPublication) -> WriteCounts: """Create one running publication without replacing an existing identity.""" diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py b/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py index 8525d5b..24f6a40 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/infrastructure/postgres.py @@ -493,6 +493,34 @@ class PostgresSectorRadarRepository: rows, ) + def load_publication_aggregates( + self, publication_ids: Sequence[str] + ) -> Sequence[DailyAggregateRecord]: + """Batch-read exact input revisions without dropping detail provenance.""" + if not publication_ids: + return () + with self._connection() as connection: + rows = connection.execute( + """ + SELECT publication_id, trade_date, sector_type, sector_code, sector_name, + member_count, valid_sample_count, net_amount_yuan, turnover_yuan, + membership_coverage, moneyflow_coverage, pct_change, leading_code + FROM sector_radar_daily_aggregate + WHERE publication_id = ANY(%s) + ORDER BY trade_date, sector_type, sector_code + """, + (list(publication_ids),), + ).fetchall() + return tuple( + DailyAggregateRecord( + str(row[0]), + self._aggregate_from_row(row[1:11]), + None if row[11] is None else Decimal(str(row[11])), + None if row[12] is None else str(row[12]), + ) + for row in rows + ) + def create_publication(self, publication: RadarPublication) -> WriteCounts: """Insert a new running publication identity idempotently.""" @@ -694,6 +722,7 @@ class PostgresSectorRadarRepository: "rank_position", "rank_percentile", "rank_changes", + "weighted_score", ), ("publication_id", "sector_type", "sector_code", "metric_version"), self._ranking_rows(ranking_items), @@ -772,6 +801,7 @@ class PostgresSectorRadarRepository: "rank_position", "rank_percentile", "rank_changes", + "weighted_score", ), ("publication_id", "sector_type", "sector_code", "metric_version"), self._ranking_rows(items), @@ -877,7 +907,8 @@ class PostgresSectorRadarRepository: SELECT trade_date, sector_type, sector_code, sector_name, metric_kind, metric_version, implementation_kind, unit, metric_value, quality, member_count, valid_sample_count, membership_coverage, - moneyflow_coverage, rank_position, rank_percentile, rank_changes + moneyflow_coverage, rank_position, rank_percentile, rank_changes, + weighted_score FROM sector_radar_ranking WHERE publication_id = %s ORDER BY sector_type, metric_version, rank_position NULLS LAST, sector_code @@ -898,7 +929,8 @@ class PostgresSectorRadarRepository: SELECT publication_id, trade_date, sector_type, sector_code, sector_name, metric_kind, metric_version, implementation_kind, unit, metric_value, quality, member_count, valid_sample_count, membership_coverage, - moneyflow_coverage, rank_position, rank_percentile, rank_changes + moneyflow_coverage, rank_position, rank_percentile, rank_changes, + weighted_score FROM sector_radar_ranking WHERE publication_id = ANY(%s) ORDER BY publication_id, sector_type, metric_kind, rank_position NULLS LAST, @@ -938,7 +970,8 @@ class PostgresSectorRadarRepository: ranking.implementation_kind, ranking.unit, ranking.metric_value, ranking.quality, ranking.member_count, ranking.valid_sample_count, ranking.membership_coverage, ranking.moneyflow_coverage, - ranking.rank_position, ranking.rank_percentile, ranking.rank_changes + ranking.rank_position, ranking.rank_percentile, ranking.rank_changes, + ranking.weighted_score FROM pools AS pool LEFT JOIN sector_radar_ranking AS ranking ON ranking.publication_id = pool.publication_id @@ -1016,7 +1049,7 @@ class PostgresSectorRadarRepository: ranking.metric_value, ranking.quality, ranking.member_count, ranking.valid_sample_count, ranking.membership_coverage, ranking.moneyflow_coverage, ranking.rank_position, - ranking.rank_percentile, ranking.rank_changes + ranking.rank_percentile, ranking.rank_changes, ranking.weighted_score FROM sector_radar_ranking AS ranking JOIN selected ON selected.id = ranking.publication_id ORDER BY selected.target_trade_date DESC, ranking.sector_type, @@ -1168,6 +1201,7 @@ class PostgresSectorRadarRepository: ranking.rank_position, ranking.rank_percentile, Jsonb({str(change.days): change.value for change in ranking.rank_changes}), + observation.weighted_score, ) ) return tuple(rows) @@ -1258,6 +1292,7 @@ class PostgresSectorRadarRepository: valid_sample_count=int(row[11]), membership_coverage=Decimal(str(row[12])), moneyflow_coverage=Decimal(str(row[13])), + weighted_score=None if row[17] is None else Decimal(str(row[17])), ) return RankedMetric( observation=observation, diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py b/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py index 4a8c3a0..7b3638a 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/cli.py @@ -10,6 +10,7 @@ from datetime import date from ....bootstrap.config import get_settings from ..application.build import BuildSectorRadar, BuildSectorRadarCommand +from ..application.recompute import RecomputeSectorRadar from ..infrastructure.postgres import PostgresSectorRadarRepository from ..infrastructure.tushare import TushareSectorRadarAdapter @@ -30,6 +31,11 @@ def build_parser() -> argparse.ArgumentParser: ) mode.add_argument("--start-date", type=_parse_date, help="inclusive backfill start date") parser.add_argument("--end-date", type=_parse_date, help="inclusive backfill end date") + parser.add_argument( + "--recompute", + action="store_true", + help="rescore saved aggregates without fetching external sources", + ) return parser @@ -39,6 +45,8 @@ def main(argv: Sequence[str] | None = None) -> int: args = build_parser().parse_args(argv) if (args.start_date is None) != (args.end_date is None): raise SystemExit("--start-date and --end-date must be provided together") + if args.recompute and args.retry_publication_id: + raise SystemExit("--recompute cannot be combined with --retry-publication-id") command = BuildSectorRadarCommand( trade_date=args.trade_date, start_date=args.start_date, @@ -59,22 +67,25 @@ def main(argv: Sequence[str] | None = None) -> int: command.end_date or "none", bool(command.retry_publication_id), ) - source = TushareSectorRadarAdapter.from_token( - settings.tushare_token, - max_retries=settings.sector_radar_max_retries, - backoff_seconds=settings.sector_radar_retry_backoff_seconds, - request_interval_seconds=settings.sector_radar_request_interval_seconds, - ) repository = PostgresSectorRadarRepository( settings.database_url, advisory_lock_key=settings.sector_radar_advisory_lock_key, ) try: - summary = BuildSectorRadar( - source, - repository, - coverage_threshold=settings.sector_radar_coverage_threshold, - ).execute(command) + if args.recompute: + summary = RecomputeSectorRadar(repository).execute(command) + else: + source = TushareSectorRadarAdapter.from_token( + settings.tushare_token, + max_retries=settings.sector_radar_max_retries, + backoff_seconds=settings.sector_radar_retry_backoff_seconds, + request_interval_seconds=settings.sector_radar_request_interval_seconds, + ) + summary = BuildSectorRadar( + source, + repository, + coverage_threshold=settings.sector_radar_coverage_threshold, + ).execute(command) finally: repository.close() except Exception as exc: # noqa: BLE001 - CLI boundary returns a redacted scheduler result diff --git a/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/http.py b/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/http.py index faf7440..9b6eb06 100644 --- a/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/http.py +++ b/zhixing-server/src/zhixing_server/modules/sector_radar/presentation/http.py @@ -98,6 +98,7 @@ class RadarRankingRowResponse(BaseModel): implementation_kind: Literal["independent"] unit: MetricUnit metric_value: Decimal | None + weighted_score: Decimal | None = None quality: MetricQuality member_count: int = Field(ge=0) valid_sample_count: int = Field(ge=0) @@ -107,6 +108,9 @@ class RadarRankingRowResponse(BaseModel): rank_percentile: Decimal | None = Field(default=None, gt=0, le=100) rank_change_days: int = Field(ge=1, le=5) rank_change: int | None + rank_change_values: dict[MetricKind, int | None] = Field( + default_factory=lambda: dict[MetricKind, int | None]() + ) pct_change: Decimal | None = None daily_net_amount_yuan: Decimal | None = None daily_ratio: Decimal | None = None @@ -127,6 +131,7 @@ class RadarHistoryMetricResponse(RadarDetailModel): rank_percentile: Decimal | None pool_size: int metric_value: Decimal | None + weighted_score: Decimal | None = None missing: bool in_top: bool in_bottom: bool @@ -222,6 +227,7 @@ class RadarRankingsResponse(BaseModel): view: RadarView rank_change_metric: MetricKind rank_change_days: int = Field(ge=1, le=5) + comparison_trade_date: date | None = None side: RankSide search: str | None publication: RadarPublicationResponse | None @@ -303,7 +309,7 @@ def get_sector_radar_rankings( trade_date: date | None = None, sector_type: SectorType = SectorType.CONCEPT, view: RadarView = RadarView.AMOUNT, - rank_change_metric: MetricKind = MetricKind.AMOUNT, + rank_change_metric: MetricKind = MetricKind.SWING, rank_change_days: Annotated[int, Query(ge=1, le=5)] = 1, side: RankSide = RankSide.ALL, search: Annotated[str | None, Query(max_length=100)] = None, @@ -458,6 +464,7 @@ def _rankings_response(page: RankingPage) -> RadarRankingsResponse: view=query.view, rank_change_metric=query.rank_change_metric, rank_change_days=query.rank_change_days, + comparison_trade_date=page.comparison_trade_date, side=query.side, search=query.search, publication=( @@ -469,7 +476,10 @@ def _rankings_response(page: RankingPage) -> RadarRankingsResponse: total=page.total, rows=[ _ranking_response( - row, query.rank_change_days, page.extras.get(row.observation.sector_code) + row, + query.rank_change_days, + page.extras.get(row.observation.sector_code), + page.rank_change_values.get(row.observation.sector_code), ) for row in page.rows ], @@ -525,7 +535,10 @@ def _definition_response( def _ranking_response( - row: RankedMetric, rank_change_days: int, extras: RankingExtras | None = None + row: RankedMetric, + rank_change_days: int, + extras: RankingExtras | None = None, + changes: dict[MetricKind, int | None] | None = None, ) -> RadarRankingRowResponse: observation = row.observation extras = extras or RankingExtras() @@ -539,6 +552,7 @@ def _ranking_response( implementation_kind=observation.implementation_kind, unit=observation.unit, metric_value=observation.value, + weighted_score=observation.weighted_score, quality=observation.quality, member_count=observation.member_count, valid_sample_count=observation.valid_sample_count, @@ -548,6 +562,7 @@ def _ranking_response( rank_percentile=row.rank_percentile, rank_change_days=rank_change_days, rank_change=row.rank_change(rank_change_days), + rank_change_values={kind: (changes or {}).get(kind) for kind in MetricKind}, pct_change=extras.pct_change, daily_net_amount_yuan=extras.daily_net_amount_yuan, daily_ratio=extras.daily_ratio, diff --git a/zhixing-server/tests/test_sector_radar_http.py b/zhixing-server/tests/test_sector_radar_http.py index 491078c..41a9a8a 100644 --- a/zhixing-server/tests/test_sector_radar_http.py +++ b/zhixing-server/tests/test_sector_radar_http.py @@ -225,6 +225,35 @@ def test_no_data_is_a_stable_200_response() -> None: assert rankings.json()["rows"] == [] +def test_weighted_score_and_all_rank_changes_cross_the_http_boundary() -> None: + reader = FakeReader() + ranking = _ranking() + reader.page = replace( + reader.page, + rows=( + replace( + ranking, + observation=replace( + ranking.observation, weighted_score=Decimal("712.345678901234") + ), + ), + ), + comparison_trade_date=date(2026, 8, 21), + rank_change_values={ + "BK0001.DC": {MetricKind.AMOUNT: 3, MetricKind.RATIO: 0, MetricKind.SWING: None} + }, + ) + response = _client(reader).get("/api/v1/sector-radar/rankings", params={"view": "rank_change"}) + assert response.status_code == 200 + assert ( + reader.last_query is not None and reader.last_query.rank_change_metric is MetricKind.SWING + ) + payload = response.json() + assert payload["comparison_trade_date"] == "2026-08-21" + assert payload["rows"][0]["weighted_score"] == "712.345678901234" + assert payload["rows"][0]["rank_change_values"] == {"amount": 3, "ratio": 0, "swing": None} + + def test_http_contract_rejects_zero_rank_percentile() -> None: payload = _client(FakeReader()).get("/api/v1/sector-radar/rankings").json()["rows"][0] payload["rank_percentile"] = "0" diff --git a/zhixing-server/tests/unit/sector_radar/test_build.py b/zhixing-server/tests/unit/sector_radar/test_build.py index 375b5fb..b14fbb2 100644 --- a/zhixing-server/tests/unit/sector_radar/test_build.py +++ b/zhixing-server/tests/unit/sector_radar/test_build.py @@ -599,10 +599,14 @@ def test_tenth_trading_day_publishes_swing_and_five_rank_changes() -> None: swing = tuple( ranking for ranking in current - if ranking.observation.metric_version == "zhixing_swing_equal_3_10_v1" + if ranking.observation.metric_version == "zhixing_swing_weighted_v2" ) assert len(swing) == 2 - assert all(ranking.observation.value == Decimal("0.03") for ranking in swing) + assert all( + ranking.observation.value == pytest.approx(Decimal(150000) / Decimal(5000100)) + for ranking in swing + ) + assert all(ranking.observation.weighted_score is not None for ranking in swing) assert all( tuple(change.value for change in ranking.rank_changes) == (None, None, None, None, None) for ranking in swing @@ -786,7 +790,7 @@ def test_detail_history_and_ranking_extras_http_use_the_same_publication() -> No row = ranking.json()["rows"][0] assert row["pct_change"] == "1" assert Decimal(row["daily_net_amount_yuan"]) == 150000 - assert Decimal(row["daily_ratio"]) == Decimal("0.03") + assert Decimal(row["daily_ratio"]) == Decimal(150000) / Decimal(5000100) assert row["on_list_count"] == row["history_available_days"] == 1 assert client.get(base + "/detail").status_code == 422 absent = client.get(base + "/detail", params={"trade_date": "2020-01-01"}) @@ -805,6 +809,13 @@ def test_postgres_detail_migration_and_build_roundtrip(monkeypatch: pytest.Monke from zhixing_server.bootstrap.config import sqlalchemy_database_url from zhixing_server.modules.sector_radar.application.details import ReadRadarDetails + from zhixing_server.modules.sector_radar.application.recompute import RecomputeSectorRadar + from zhixing_server.modules.sector_radar.domain.metrics import ( + AmountNetStrategy, + RatioTurnoverStrategy, + SwingEqualThreeToTenStrategy, + ) + from zhixing_server.modules.sector_radar.domain.models import MetricKind from zhixing_server.modules.sector_radar.infrastructure.postgres import ( PostgresSectorRadarRepository, ) @@ -843,7 +854,7 @@ def test_postgres_detail_migration_and_build_roundtrip(monkeypatch: pytest.Monke assert detail.summary["amount"].metric_value == Decimal("0.0015") with psycopg.connect(database_url) as connection: assert connection.execute("SELECT version_num FROM alembic_version").fetchone() == ( - "0009_radar_sector_detail", + "0010_radar_weighted_score", ) row = connection.execute( "SELECT pct_change, leading_code FROM sector_radar_daily_aggregate " @@ -863,5 +874,140 @@ def test_postgres_detail_migration_and_build_roundtrip(monkeypatch: pytest.Monke "SET active_buy_net_amount_yuan = 'NaN'::numeric WHERE trade_date = %s", (target,), ) + # Exercise every ranking projection against persisted v2 scores, while + # keeping the original source publication and its NULL score readable. + start, end = target + timedelta(days=31), target + timedelta(days=46) + interval = BuildSectorRadarCommand(start_date=start, end_date=end) + source = FakeRadarSource() + legacy = BuildSectorRadar( + source, + repository, + now_fn=lambda: NOW, + strategies=( + AmountNetStrategy(), + RatioTurnoverStrategy(), + SwingEqualThreeToTenStrategy(), + ), + ).execute(interval) + assert legacy.status in ("success", "unchanged") + old_id = legacy.outcomes[-1].publication_id + assert old_id is not None + old_rows = tuple(repository.load_rankings(old_id)) + assert all(row.observation.weighted_score is None for row in old_rows) + source.calls.clear() + rescore = RecomputeSectorRadar(repository, now_fn=lambda: NOW + timedelta(hours=1)) + updated = rescore.execute(interval) + assert updated.status in ("success", "unchanged") + updated_id = updated.outcomes[-1].publication_id + assert updated_id is not None and updated_id != old_id + assert source.calls == [] + assert tuple(repository.load_rankings(old_id)) == old_rows + current_rows = tuple(repository.load_rankings(updated_id)) + scored = next( + row + for row in current_rows + if row.observation.metric_kind is MetricKind.SWING + and row.observation.sector_type is SectorType.CONCEPT + ) + expected = ((Decimal(5000000) + 1).log10() * 100).quantize(Decimal("0.000000000001")) + assert scored.observation.weighted_score == expected + assert scored.rank_change(5) == 0 + assert dict(repository.load_publication_rankings((updated_id,)))[updated_id] == current_rows + previous = dict(repository.load_previous_rankings(end + timedelta(days=1), limit_dates=1)) + assert previous[end] == current_rows + historical = repository.load_ranked_history((updated_id,), ("BK0001.DC",)) + assert any(row.ranking == scored for row in historical) + detail = ReadRadarDetails(repository).detail(end, SectorType.CONCEPT, "BK0001.DC") + assert detail.summary["swing"].weighted_score == expected + assert detail.pct_change == 1 and len(detail.members) == 5 + assert rescore.execute(interval).status == "unchanged" + with ( + psycopg.connect(database_url) as connection, + pytest.raises(psycopg.errors.CheckViolation), + connection.transaction(), + ): + connection.execute( + "UPDATE sector_radar_ranking SET weighted_score = 'NaN'::numeric " + "WHERE publication_id = %s", + (updated_id,), + ) finally: repository.close() + + +def test_offline_rescore_is_idempotent_preserves_old_publications_and_details() -> None: + from zhixing_server.modules.sector_radar.application.details import ReadRadarDetails + from zhixing_server.modules.sector_radar.application.read import ( + RadarQuery, + RadarView, + ReadSectorRadar, + ) + from zhixing_server.modules.sector_radar.application.recompute import RecomputeSectorRadar + from zhixing_server.modules.sector_radar.domain.metrics import ( + AmountNetStrategy, + RatioTurnoverStrategy, + SwingEqualThreeToTenStrategy, + ) + + source = FakeRadarSource() + repository = InMemorySectorRadarRepository() + end = TARGET_DATE + timedelta(days=15) + command = BuildSectorRadarCommand(start_date=TARGET_DATE, end_date=end) + initial = BuildSectorRadar( + source, + repository, + now_fn=lambda: NOW, + strategies=(AmountNetStrategy(), RatioTurnoverStrategy(), SwingEqualThreeToTenStrategy()), + ).execute(command) + assert initial.status == "success" + old_publications = repository.publications.copy() + old_rankings = repository.rankings.copy() + snapshot_count = len(repository.source_snapshots) + source.calls.clear() + source.fail_daily = True + service = RecomputeSectorRadar(repository, now_fn=lambda: NOW + timedelta(hours=1)) + result = service.execute(command) + assert result.status == "success" + assert all(repository.publications[key] == value for key, value in old_publications.items()) + assert all(repository.rankings[key] == value for key, value in old_rankings.items()) + assert len(repository.source_snapshots) == snapshot_count + assert source.calls == [] + latest = repository.get_last_good_publication(end) + assert latest is not None + assert "zhixing_swing_weighted_v2" in latest.metric_versions + count = len(repository.publications) + assert service.execute(command).status == "unchanged" + assert len(repository.publications) == count + reader = ReadSectorRadar(repository) + page = reader.query(RadarQuery(trade_date=end, view=RadarView.SWING)) + assert len(page.rows) == 1 + assert page.rows[0].observation.weighted_score is not None + assert page.rows[0].rank_change(5) == 0 + detail = ReadRadarDetails(repository).detail(end, SectorType.CONCEPT, "BK0001.DC") + assert len(detail.members) == 5 + assert detail.pct_change == 1 + assert detail.summary["swing"].weighted_score == page.rows[0].observation.weighted_score + assert detail.summary["swing"].rank_position == page.rows[0].rank_position + + +def test_offline_rescore_failure_does_not_replace_last_good( + monkeypatch: pytest.MonkeyPatch, +) -> None: + from unittest.mock import Mock + + from zhixing_server.modules.sector_radar.application.recompute import RecomputeSectorRadar + + repository = InMemorySectorRadarRepository() + command = BuildSectorRadarCommand(trade_date=TARGET_DATE) + BuildSectorRadar(FakeRadarSource(), repository, now_fn=lambda: NOW).execute(command) + before = repository.get_last_good_publication() + monkeypatch.setattr( + repository, "finalize_publication", Mock(side_effect=RuntimeError("injected failure")) + ) + result = RecomputeSectorRadar(repository, now_fn=lambda: NOW + timedelta(hours=1)).execute( + command + ) + assert result.status == "failed" + assert repository.get_last_good_publication() == before + latest = repository.get_latest_publication() + assert latest is not None and latest.status is PublicationStatus.FAILED diff --git a/zhixing-server/tests/unit/sector_radar/test_postgres_repository.py b/zhixing-server/tests/unit/sector_radar/test_postgres_repository.py index f8629ad..6d52482 100644 --- a/zhixing-server/tests/unit/sector_radar/test_postgres_repository.py +++ b/zhixing-server/tests/unit/sector_radar/test_postgres_repository.py @@ -61,6 +61,7 @@ class FakeConnection: 1, Decimal(100), {"1": 3, "2": None}, + None, ), ) ) @@ -161,6 +162,7 @@ def test_load_rankings_reconstructs_values_and_rank_changes() -> None: ranking = rankings[0] assert ranking.observation.metric_version == "zhixing_amount_net_bn_v1" assert ranking.observation.value == Decimal("12.5") + assert ranking.observation.weighted_score is None assert ranking.rank_position == 1 assert ranking.rank_change(1) == 3 assert ranking.rank_change(2) is None diff --git a/zhixing-server/tests/unit/sector_radar/test_read.py b/zhixing-server/tests/unit/sector_radar/test_read.py index b00c309..84e24a0 100644 --- a/zhixing-server/tests/unit/sector_radar/test_read.py +++ b/zhixing-server/tests/unit/sector_radar/test_read.py @@ -158,7 +158,45 @@ def test_percentile_side_is_selected_before_search_and_pagination() -> None: def test_rank_change_uses_selected_metric_days_and_pool_sides() -> None: - reader = ReadSectorRadar(_published_repository()) + from zhixing_server.modules.sector_radar.domain.persistence import ( + PublicationSourceGroup, + PublicationSourceRecord, + ) + from zhixing_server.modules.sector_radar.domain.source import build_source_snapshot + + repository = _published_repository() + dates = [date(2026, 8, day) for day in [21, 24, 25, 26, 27, 28]] + calendar = build_source_snapshot( + api_name="trade_cal", + params={}, + observed_at=NOW, + target_trade_date=TARGET_DATE, + rows=tuple({"exchange": "SSE", "cal_date": day.isoformat(), "is_open": 1} for day in dates), + ) + repository.save_source_snapshots((calendar,)) + repository.save_publication_sources( + ( + PublicationSourceRecord( + "publication-success", PublicationSourceGroup.CALENDAR, 0, calendar + ), + ) + ) + old = _running("past-publication", dates[0]) + repository.create_publication(old) + past = tuple( + replace( + row.observation, + trade_date=dates[0], + value=Decimal(index), + sector_code="BK9999.DC" if index == 5 else row.observation.sector_code, + ) + for index, row in enumerate(_amount_rankings(), 1) + ) + repository.save_rankings( + RankingRecord(old.publication_id, row) for row in rank_metric_observations(past) + ) + repository.finish_publication(_finish(old, PublicationStatus.SUCCESS)) + reader = ReadSectorRadar(repository) query = RadarQuery( view=RadarView.RANK_CHANGE, rank_change_metric=MetricKind.AMOUNT, @@ -170,12 +208,14 @@ def test_rank_change_uses_selected_metric_days_and_pool_sides() -> None: all_rows = reader.query(query) assert top.total == 1 - assert top.rows[0].rank_change(5) == 5 + assert top.rows[0].rank_change(5) == 9 assert bottom.total == 1 - assert bottom.rows[0].rank_change(5) == -4 + assert bottom.rows[0].rank_change(5) == -9 assert all_rows.total == 10 assert all_rows.rows[-1].observation.sector_code == "BK0005.DC" assert all_rows.rows[-1].rank_change(5) is None + assert top.comparison_trade_date == dates[0] + assert top.rank_change_values["BK0001.DC"][MetricKind.AMOUNT] == 9 def test_latest_partial_attempt_is_visible_but_does_not_replace_last_good() -> None: diff --git a/zhixing-server/tests/unit/sector_radar/test_weighted.py b/zhixing-server/tests/unit/sector_radar/test_weighted.py new file mode 100644 index 0000000..effa634 --- /dev/null +++ b/zhixing-server/tests/unit/sector_radar/test_weighted.py @@ -0,0 +1,160 @@ +from dataclasses import replace +from datetime import date, timedelta +from decimal import Decimal + +import pytest + +from zhixing_server.modules.sector_radar.application.scoring import ( + calculate_rankings, + calendar_rank_changes, + comparison_dates, +) +from zhixing_server.modules.sector_radar.domain.models import ( + MetricKind, + RankedMetric, + RankSide, + SectorDailyAggregate, + SectorType, +) +from zhixing_server.modules.sector_radar.domain.ranking import select_percentile_side +from zhixing_server.modules.sector_radar.domain.weighted import ( + RATIO_WEIGHTED_VERSION, + SWING_WEIGHTED_VERSION, + resolve_metric_version, +) + +DAYS = tuple( + date(2026, 9, 1) + timedelta(days=n) + for n in range(21) + if (date(2026, 9, 1) + timedelta(days=n)).weekday() < 5 +) + + +def aggregate( + code: str, + day: date, + ratio: str, + turnover: str = "9999999999", + kind: SectorType = SectorType.CONCEPT, +) -> SectorDailyAggregate: + amount = Decimal(turnover) + return SectorDailyAggregate( + day, + kind, + code, + code, + 10, + 10, + Decimal(ratio) * (amount + 100), + amount, + Decimal(1), + Decimal(1), + ) + + +def scores( + rows: list[SectorDailyAggregate], target: date = DAYS[9] +) -> dict[tuple[SectorType, str, MetricKind], RankedMetric]: + return { + (row.observation.sector_type, row.observation.sector_code, row.observation.metric_kind): row + for row in calculate_rankings( + target, + [item for item in rows if item.trade_date == target], + [item for item in rows if item.trade_date != target], + DAYS, + ) + } + + +def test_liquidity_weight_can_reverse_raw_ratio_order_and_separates_pools() -> None: + rows = [ + aggregate(code, day, ratio, turnover, kind) + for day in DAYS[:10] + for code, ratio, turnover, kind in [ + ("small", ".3", "999999", SectorType.CONCEPT), + ("large", ".2", "999999999999", SectorType.CONCEPT), + ("medium", ".1", "9999999999", SectorType.CONCEPT), + ("industry", "-.5", "9999999999", SectorType.INDUSTRY), + ] + ] + ranked = scores(rows) + small = ranked[(SectorType.CONCEPT, "small", MetricKind.RATIO)] + large = ranked[(SectorType.CONCEPT, "large", MetricKind.RATIO)] + assert small.observation.value == Decimal(".3") + assert small.observation.weighted_score == 600 + assert large.observation.weighted_score == 800 + assert large.rank_position == 1 and small.rank_position == 2 + assert ( + ranked[(SectorType.INDUSTRY, "industry", MetricKind.RATIO)].observation.weighted_score + == 1000 + ) + + +def test_swing_ranks_each_window_before_combining_not_the_averaged_raw_ratio() -> None: + rows = [ + aggregate(code, day, ratio) + for index, day in enumerate(DAYS[:10]) + for code, ratio in [("A", ".1" if index >= 7 else "-.1"), ("B", ".03"), ("C", ".05")] + ] + ranked = scores(rows) + a = ranked[(SectorType.CONCEPT, "A", MetricKind.SWING)] + b = ranked[(SectorType.CONCEPT, "B", MetricKind.SWING)] + assert a.observation.value == b.observation.value == Decimal(".03") + assert a.observation.weighted_score == pytest.approx(Decimal("666.6666666666667")) + assert b.observation.weighted_score == 500 + assert a.rank_position == 2 and b.rank_position == 3 + + +def test_tied_features_use_average_percentiles_and_code_breaks_final_ties() -> None: + ranked = scores([aggregate(code, day, "0") for day in DAYS[:10] for code in ["B", "A"]]) + assert ranked[(SectorType.CONCEPT, "A", MetricKind.RATIO)].observation.weighted_score == 750 + assert ranked[(SectorType.CONCEPT, "B", MetricKind.RATIO)].observation.weighted_score == 750 + assert ranked[(SectorType.CONCEPT, "A", MetricKind.RATIO)].rank_position == 1 + + +def test_missing_calendar_session_is_not_replaced_by_an_older_success() -> None: + rows = [aggregate("A", day, ".2") for day in DAYS[:10] if day != DAYS[7]] + ranked = scores(rows) + daily = ranked[(SectorType.CONCEPT, "A", MetricKind.RATIO)] + assert daily.observation.value == Decimal(".2") + assert daily.observation.weighted_score is None and daily.rank_position is None + assert ranked[(SectorType.CONCEPT, "A", MetricKind.SWING)].observation.value is None + + +def test_future_inputs_do_not_affect_scores_and_bottom_uses_score_order() -> None: + rows = [ + aggregate(f"C{n}", day, str(n), str(10 ** (6 + n) - 1)) + for day in DAYS[:10] + for n in range(1, 11) + ] + original = scores(rows) + future = scores(rows + [aggregate("C1", DAYS[10], "1000000")]) + assert future == original + ratio_rows = [ + row for row in original.values() if row.observation.metric_kind is MetricKind.RATIO + ] + bottom = select_percentile_side(ratio_rows, RankSide.BOTTOM) + assert [row.observation.sector_code for row in bottom] == ["C1"] + + +def test_rank_changes_resolve_trading_dates_and_do_not_mix_versions() -> None: + rows = [aggregate("A", day, ".1") for day in DAYS[:11]] + before = tuple(scores(rows, DAYS[9]).values()) + current = tuple(scores(rows, DAYS[10]).values()) + assert comparison_dates(DAYS, DAYS[10])[1] == DAYS[9] + missing = calendar_rank_changes(current, {DAYS[8]: before}, DAYS, DAYS[10]) + assert all(row.rank_change(1) is None for row in missing) + older = tuple( + replace(row, observation=replace(row.observation, metric_version="legacy")) + for row in before + ) + incompatible = calendar_rank_changes(current, {DAYS[9]: older}, DAYS, DAYS[10]) + assert all(row.rank_change(1) is None for row in incompatible) + comparable = calendar_rank_changes(current, {DAYS[9]: before}, DAYS, DAYS[10]) + assert all(row.rank_change(1) == 0 for row in comparable) + assert ( + resolve_metric_version(MetricKind.RATIO, [RATIO_WEIGHTED_VERSION]) == RATIO_WEIGHTED_VERSION + ) + assert ( + resolve_metric_version(MetricKind.SWING, [SWING_WEIGHTED_VERSION]) == SWING_WEIGHTED_VERSION + ) diff --git a/zhixing-web/.prettierignore b/zhixing-web/.prettierignore index a46f1e4..ac98bba 100644 --- a/zhixing-web/.prettierignore +++ b/zhixing-web/.prettierignore @@ -3,3 +3,4 @@ dist node_modules pnpm-lock.yaml .pnpm-store +.playwright-cli diff --git a/zhixing-web/src/features/sector-radar/api/sector-radar.api.test.ts b/zhixing-web/src/features/sector-radar/api/sector-radar.api.test.ts index b451c9f..1356283 100644 --- a/zhixing-web/src/features/sector-radar/api/sector-radar.api.test.ts +++ b/zhixing-web/src/features/sector-radar/api/sector-radar.api.test.ts @@ -8,6 +8,7 @@ import { getSectorRadarHistory, getSectorRadarDetail, parseRadarDetailResponse, + parseRadarRankingsResponse, getSectorRadarDates, getSectorRadarRankings, getStockSectorMembership, @@ -133,6 +134,72 @@ describe("sector radar API adapters", () => { requestJson.mockReset() }) + it("preserves score precision, signed changes, zero and unknown historical values", () => { + const response = parseRadarRankingsResponse({ + ...rankingPayload, + comparison_trade_date: "2026-08-21", + rows: [ + { + ...rankingPayload.rows[0], + weighted_score: "712.345678901234", + rank_change_values: { amount: 3, ratio: 0, swing: null }, + }, + ], + }) + expect(response.comparison_trade_date).toBe("2026-08-21") + expect(response.rows[0]?.weighted_score).toBe(712.345678901234) + expect(response.rows[0]?.rank_change_values).toEqual({ + amount: 3, + ratio: 0, + swing: null, + }) + const legacy = parseRadarRankingsResponse(rankingPayload) + expect(legacy.rows[0]?.weighted_score).toBeNull() + expect(legacy.rows[0]?.rank_change_values).toEqual({ + amount: null, + ratio: null, + swing: null, + }) + expect(legacy.comparison_trade_date).toBeNull() + const detail = parseRadarDetailResponse({ + ...detailPayload, + summary: { + ...detailPayload.summary, + ratio: { ...historyMetric, weighted_score: "712.345" }, + }, + }) + expect(detail.summary.ratio.weighted_score).toBe(712.345) + }) + + it.each(["NaN", "Infinity", "0x10", "", true])( + "rejects malformed weighted scores %s", + (weighted_score) => { + expect(() => + parseRadarRankingsResponse({ + ...rankingPayload, + rows: [{ ...rankingPayload.rows[0], weighted_score }], + }), + ).toThrow("weighted_score") + }, + ) + + it("rejects fractional changes and invalid comparison dates", () => { + expect(() => + parseRadarRankingsResponse({ + ...rankingPayload, + rows: [ + { ...rankingPayload.rows[0], rank_change_values: { amount: 0.5 } }, + ], + }), + ).toThrow("rank_change_values.amount") + expect(() => + parseRadarRankingsResponse({ + ...rankingPayload, + comparison_trade_date: "yesterday", + }), + ).toThrow("comparison_trade_date") + }) + it("reads persisted sector resources and preserves nullable independent metrics", async () => { const signal = new AbortController().signal const query = { diff --git a/zhixing-web/src/features/sector-radar/api/sector-radar.api.ts b/zhixing-web/src/features/sector-radar/api/sector-radar.api.ts index 269ff41..267d33c 100644 --- a/zhixing-web/src/features/sector-radar/api/sector-radar.api.ts +++ b/zhixing-web/src/features/sector-radar/api/sector-radar.api.ts @@ -121,6 +121,10 @@ export function parseRadarRankingsResponse( 5, "rankings.rank_change_days", ), + comparison_trade_date: readNullableDate( + record.comparison_trade_date ?? null, + "rankings.comparison_trade_date", + ), side: readEnum(record.side, radarRankSides, "rankings.side"), search: readNullableString(record.search, "rankings.search"), publication: readNullablePublication( @@ -305,6 +309,14 @@ function readRankingRow(value: unknown, index: number): RadarRankingRow { record.metric_value, `${path}.metric_value`, ), + weighted_score: readNullableFiniteNumber( + record.weighted_score ?? null, + `${path}.weighted_score`, + ), + rank_change_values: readRankChangeValues( + record.rank_change_values, + `${path}.rank_change_values`, + ), pct_change: readNullableFiniteNumber( record.pct_change ?? null, `${path}.pct_change`, @@ -377,6 +389,30 @@ function readRankingRow(value: unknown, index: number): RadarRankingRow { } } +function readRankChangeValues( + value: unknown, + path: string, +): RadarRankingRow["rank_change_values"] { + const record = value == null ? {} : readRecord(value, path) + return { + amount: readNullableInteger( + record.amount ?? null, + Number.MIN_SAFE_INTEGER, + `${path}.amount`, + ), + ratio: readNullableInteger( + record.ratio ?? null, + Number.MIN_SAFE_INTEGER, + `${path}.ratio`, + ), + swing: readNullableInteger( + record.swing ?? null, + Number.MIN_SAFE_INTEGER, + `${path}.swing`, + ), + } +} + function readRecord(value: unknown, path: string): JsonRecord { if (typeof value !== "object" || value === null || Array.isArray(value)) { throw contractError(path, "must be an object") @@ -627,6 +663,10 @@ function readHistoryMetric(value: unknown, path: string): RadarHistoryMetric { record.metric_value, `${path}.metric_value`, ), + weighted_score: readNullableFiniteNumber( + record.weighted_score ?? null, + `${path}.weighted_score`, + ), missing: readBoolean(record.missing, `${path}.missing`), in_top: readBoolean(record.in_top, `${path}.in_top`), in_bottom: readBoolean(record.in_bottom, `${path}.in_bottom`), diff --git a/zhixing-web/src/features/sector-radar/api/sector-radar.types.ts b/zhixing-web/src/features/sector-radar/api/sector-radar.types.ts index e78d5bc..1cf49fe 100644 --- a/zhixing-web/src/features/sector-radar/api/sector-radar.types.ts +++ b/zhixing-web/src/features/sector-radar/api/sector-radar.types.ts @@ -68,6 +68,7 @@ export interface RadarRankingRow { implementation_kind: "independent" unit: RadarMetricUnit metric_value: number | null + weighted_score: number | null quality: RadarMetricQuality member_count: number valid_sample_count: number @@ -77,6 +78,7 @@ export interface RadarRankingRow { rank_percentile: number | null rank_change_days: number rank_change: number | null + rank_change_values: Record pct_change: number | null daily_net_amount_yuan: number | null daily_ratio: number | null @@ -91,6 +93,7 @@ export interface RadarRankingsResponse { view: RadarView rank_change_metric: RadarMetricKind rank_change_days: number + comparison_trade_date: string | null side: RadarRankSide search: string | null publication: RadarPublication | null @@ -141,6 +144,7 @@ export interface RadarHistoryMetric { rank_percentile: number | null pool_size: number metric_value: number | null + weighted_score: number | null missing: boolean in_top: boolean in_bottom: boolean diff --git a/zhixing-web/src/features/sector-radar/components/radar-detail.test.tsx b/zhixing-web/src/features/sector-radar/components/radar-detail.test.tsx index 5cd2b3b..bacec97 100644 --- a/zhixing-web/src/features/sector-radar/components/radar-detail.test.tsx +++ b/zhixing-web/src/features/sector-radar/components/radar-detail.test.tsx @@ -33,6 +33,7 @@ const metric = { rank_percentile: 90, pool_size: 20, metric_value: 0.02, + weighted_score: null, missing: false, in_top: true, in_bottom: false, @@ -111,6 +112,7 @@ const row: RadarRankingRow = { implementation_kind: "independent", unit: "ratio", metric_value: 0.02, + weighted_score: null, quality: "available", member_count: 1, valid_sample_count: 1, @@ -120,6 +122,7 @@ const row: RadarRankingRow = { rank_percentile: 90, rank_change_days: 1, rank_change: null, + rank_change_values: { amount: null, ratio: null, swing: null }, pct_change: -2, daily_net_amount_yuan: 2e8, daily_ratio: 0.02, diff --git a/zhixing-web/src/features/sector-radar/components/radar-table-sort.ts b/zhixing-web/src/features/sector-radar/components/radar-table-sort.ts new file mode 100644 index 0000000..4028dda --- /dev/null +++ b/zhixing-web/src/features/sector-radar/components/radar-table-sort.ts @@ -0,0 +1,107 @@ +import type { + RadarMetricKind, + RadarRankingRow, + RadarView, +} from "../api/sector-radar.types" + +export const radarMetricLabels: Record = { + swing: "波段流入率", + ratio: "单日流入率", + amount: "单日净额", +} + +export type RadarSortKey = + | "pct_change" + | "weighted_score" + | "daily_net_amount_yuan" + | "daily_ratio" + | "on_list_count" + | "metric_value" + | "change_amount" + | "change_ratio" + | "change_swing" + +export type RadarSort = { + key: RadarSortKey + direction: "ascending" | "descending" +} + +export type RadarColumn = { label: string; key?: RadarSortKey } + +export function rankChangeAuxiliaryMetrics(metric: RadarMetricKind) { + return (["swing", "ratio", "amount"] as const).filter( + (candidate) => candidate !== metric, + ) +} + +/** Keep the two table halves in the same field order, mirrored at the center. */ +export function radarColumns( + view: RadarView, + metric: RadarMetricKind, + side: "top" | "bottom", +): RadarColumn[] { + const columns: RadarColumn[] = [ + { label: "板块/排名" }, + { label: "涨跌幅", key: "pct_change" }, + ] + if (view === "rank_change") { + columns.push( + ...rankChangeAuxiliaryMetrics(metric).map((kind): RadarColumn => ({ + label: radarMetricLabels[kind], + key: `change_${kind}`, + })), + { + label: side === "top" ? "上升" : "下降", + key: `change_${metric}`, + }, + ) + } else if (view === "amount") { + columns.push( + { label: "单日流入率", key: "daily_ratio" }, + { label: "净额", key: "metric_value" }, + ) + } else { + columns.push( + { label: "加权评分", key: "weighted_score" }, + { label: "净额", key: "daily_net_amount_yuan" }, + { label: "在榜", key: "on_list_count" }, + { + label: view === "swing" ? (side === "top" ? "流入" : "流出") : "流入率", + key: "metric_value", + }, + ) + } + return side === "top" ? columns : columns.reverse() +} + +function sortValue(row: RadarRankingRow, key: RadarSortKey) { + switch (key) { + case "change_amount": + return row.rank_change_values.amount + case "change_ratio": + return row.rank_change_values.ratio + case "change_swing": + return row.rank_change_values.swing + default: + return row[key] + } +} + +/** Reorder an entire side's candidate pool; unknown values remain last in both directions. */ +export function sortRadarRows(rows: RadarRankingRow[], sort: RadarSort | null) { + if (!sort) return rows + return [...rows].sort((a, b) => { + const left = sortValue(a, sort.key) + const right = sortValue(b, sort.key) + if (left == null && right != null) return 1 + if (right == null && left != null) return -1 + if (left != null && right != null && left !== right) { + return (left - right) * (sort.direction === "ascending" ? 1 : -1) + } + return ( + (a.rank_position ?? Number.MAX_SAFE_INTEGER) - + (b.rank_position ?? Number.MAX_SAFE_INTEGER) || + a.sector_code.localeCompare(b.sector_code) + ) + }) +} diff --git a/zhixing-web/src/features/sector-radar/pages/sector-radar-page.test.tsx b/zhixing-web/src/features/sector-radar/pages/sector-radar-page.test.tsx index 0318f0c..b906204 100644 --- a/zhixing-web/src/features/sector-radar/pages/sector-radar-page.test.tsx +++ b/zhixing-web/src/features/sector-radar/pages/sector-radar-page.test.tsx @@ -71,6 +71,7 @@ const rankingsResponse: RadarRankingsResponse = { view: "amount", rank_change_metric: "amount", rank_change_days: 1, + comparison_trade_date: "2026-08-27", side: "all", search: null, publication: successPublication, @@ -96,6 +97,8 @@ const rankingsResponse: RadarRankingsResponse = { implementation_kind: "independent", unit: "CNY_100M", metric_value: 12.5, + weighted_score: 712.345, + rank_change_values: { amount: 3, ratio: 0, swing: -2 }, quality: "available", member_count: 20, valid_sample_count: 19, @@ -121,6 +124,8 @@ const rankingsResponse: RadarRankingsResponse = { implementation_kind: "independent", unit: "CNY_100M", metric_value: null, + weighted_score: null, + rank_change_values: { amount: null, ratio: null, swing: null }, quality: "available_limited_sample", member_count: 4, valid_sample_count: 3, @@ -272,22 +277,28 @@ describe("SectorRadarPage", () => { expect( within(screen.getAllByRole("row")[1]!) .getAllByRole("columnheader") - .map((cell) => cell.textContent), + .map((cell) => cell.textContent?.replace(/[↕↑↓]/g, "")), ).toEqual([ "板块/排名", "涨跌幅", + "加权评分", "净额", "在榜", "流入", "流出", "在榜", "净额", + "加权评分", "涨跌幅", "板块/排名", ]) expect(screen.getAllByText("-1.25%")).toHaveLength(2) expect(screen.getAllByText("+12.5 亿元")).toHaveLength(2) } + expect(screen.getAllByText("712.3")).toHaveLength(2) + expect( + screen.getAllByRole("columnheader", { name: "加权评分" }), + ).toHaveLength(2) expect(screen.getAllByText("-1%")).toHaveLength(4) const trigger = screen.getAllByRole("button", { name: "机器人展开在榜历史", @@ -537,7 +548,7 @@ describe("SectorRadarPage", () => { }) }) - it("loads the next page once when repeated scroll events reach the bottom", () => { + it("loads each side once when repeated scroll events reach the bottom", () => { fetchNextPage.mockReturnValue(new Promise(() => undefined)) useSectorRadarRankings.mockReturnValue( rankingQueryResult({ hasNextPage: true }), @@ -555,7 +566,7 @@ describe("SectorRadarPage", () => { fireEvent.scroll(scrollRegion) fireEvent.scroll(scrollRegion) - expect(fetchNextPage).toHaveBeenCalledOnce() + expect(fetchNextPage).toHaveBeenCalledTimes(2) }) it("announces that the next page is loading without hiding loaded rows", () => { @@ -676,10 +687,143 @@ describe("SectorRadarPage", () => { render() expect( - screen.getByRole("columnheader", { name: "排名上升" }), + screen.getByRole("columnheader", { name: "上升" }), ).toBeInTheDocument() expect(screen.getAllByText("+3 名")[0]).toBeInTheDocument() - expect(screen.getAllByText("暂无可比历史")[0]).toBeInTheDocument() + expect(screen.getAllByTitle("暂无可比历史")[0]).toHaveTextContent("—") + }) + + it.each( + (["swing", "ratio", "amount"] as const).flatMap((metric) => + [1, 2, 3, 4, 5].map((days) => ({ metric, days })), + ), + )( + "restores the $metric / $days-session view and changes URL filters", + async ({ metric, days }) => { + routeSearch = { + ...routeSearch, + view: "rank_change", + rankChangeMetric: metric, + rankChangeDays: days, + } + const labels = { + swing: "波段流入率", + ratio: "单日流入率", + amount: "单日净额", + } + useSectorRadarRankings.mockReturnValue( + rankingQueryResult({ + data: { + pageParams: [1], + pages: [ + { + ...rankingsResponse, + view: "rank_change", + rank_change_metric: metric, + rank_change_days: days, + }, + ], + }, + }), + ) + render() + expect( + screen.getByRole("columnheader", { name: `${labels[metric]}排名变化` }), + ).toBeInTheDocument() + expect(screen.getByText("排名飙升榜 TOP 10%")).toBeInTheDocument() + expect(screen.getByText("BOTTOM 10% 排名暴跌榜")).toBeInTheDocument() + expect( + screen.getByRole("combobox", { name: "统计天数" }), + ).toHaveTextContent(`近 ${days} 天`) + expect(useSectorRadarRankings).toHaveBeenCalledWith( + expect.objectContaining({ + rankChangeMetric: metric, + rankChangeDays: days, + }), + ) + const tabs = within(screen.getByRole("tablist", { name: "排序基准" })) + expect( + tabs.getByRole("tab", { + name: { swing: "波段率", ratio: "单日率", amount: "单日额" }[metric], + }), + ).toHaveAttribute("aria-selected", "true") + fireEvent.click(tabs.getByRole("tab", { name: "单日率" })) + expect(navigate.mock.calls.at(-1)?.[0].search(routeSearch)).toMatchObject( + { rankChangeMetric: "ratio", page: 1 }, + ) + fireEvent.click(screen.getByRole("combobox", { name: "统计天数" })) + const option = await screen.findByRole("option", { name: "近 5 天" }) + fireEvent.pointerDown(option, { pointerType: "mouse" }) + fireEvent.click(option) + expect(navigate.mock.calls.at(-1)?.[0].search(routeSearch)).toMatchObject( + { rankChangeDays: 5, page: 1 }, + ) + expect(screen.getAllByRole("columnheader")).toHaveLength(13) + }, + ) + + it("loads all side candidates before sorting, keeps missing scores last, and leaves the other side intact", () => { + routeSearch = { ...routeSearch, view: "ratio" } + const first = { ...rankingsResponse, view: "ratio" as const, total: 3 } + const last = { + ...first, + page: 2, + rows: [ + { + ...first.rows[0]!, + sector_code: "BK03", + sector_name: "后页高分", + weighted_score: 900, + }, + ], + } + const nextTop = vi.fn().mockResolvedValue(undefined) + const nextBottom = vi.fn().mockResolvedValue(undefined) + let complete = false + useSectorRadarRankings.mockImplementation((query: { side: string }) => + rankingQueryResult({ + data: { + pageParams: complete ? [1, 2] : [1], + pages: + query.side === "top" + ? complete + ? [first, last] + : [first] + : [{ ...bottomRankingsResponse, view: "ratio" }], + }, + hasNextPage: query.side === "top" && !complete, + fetchNextPage: query.side === "top" ? nextTop : nextBottom, + }), + ) + const { container, rerender } = render() + const order = () => + Array.from( + container.querySelectorAll("tbody tr"), + (row) => row.firstElementChild?.textContent, + ) + fireEvent.click(screen.getByRole("button", { name: "左榜按加权评分排序" })) + expect(nextTop).toHaveBeenCalledOnce() + expect(nextBottom).not.toHaveBeenCalled() + expect(order()[0]).toContain("机器人") + complete = true + rerender() + expect(order()).toEqual([ + expect.stringContaining("后页高分"), + expect.stringContaining("机器人"), + expect.stringContaining("低空经济"), + ]) + expect( + container.querySelector("tbody tr")?.lastElementChild, + ).toHaveTextContent("新能源") + const score = screen.getByRole("button", { name: "左榜按加权评分排序" }) + expect(score.closest("th")).toHaveAttribute("aria-sort", "descending") + fireEvent.click(score) + expect(order()).toEqual([ + expect.stringContaining("机器人"), + expect.stringContaining("后页高分"), + expect.stringContaining("低空经济"), + ]) + expect(score.closest("th")).toHaveAttribute("aria-sort", "ascending") }) it("warns when a partial attempt has not replaced last-good", () => { diff --git a/zhixing-web/src/features/sector-radar/pages/sector-radar-page.tsx b/zhixing-web/src/features/sector-radar/pages/sector-radar-page.tsx index b47a129..f8c129b 100644 --- a/zhixing-web/src/features/sector-radar/pages/sector-radar-page.tsx +++ b/zhixing-web/src/features/sector-radar/pages/sector-radar-page.tsx @@ -1,5 +1,5 @@ import { AlertTriangle, Database, RefreshCw } from "lucide-react" -import { useRef, useState } from "react" +import { useCallback, useEffect, useRef, useState } from "react" import { useNavigate, useSearch } from "@tanstack/react-router" import { PageLayout } from "@/app/layout/page-layout" @@ -39,23 +39,32 @@ import { RadarDetailDialog } from "../components/radar-detail-dialog" import { RadarHistory } from "../components/radar-history" import { formatRadarValue, radarValueTone } from "../components/radar-format" +import { + radarColumns, + radarMetricLabels, + rankChangeAuxiliaryMetrics, + sortRadarRows, + type RadarColumn, + type RadarSort, +} from "../components/radar-table-sort" + const sectorTypeOptions = [ { label: "概念", value: "concept" }, { label: "行业", value: "industry" }, ] as const const viewOptions = [ - { label: "波段资金率", value: "swing" }, + { label: "波段流入率", value: "swing" }, { label: "单日流入率", value: "ratio" }, { label: "单日净额", value: "amount" }, { label: "排名变化", value: "rank_change" }, ] as const const metricOptions = [ - { label: "主力净额", value: "amount" }, - { label: "单日流入率", value: "ratio" }, - { label: "波段资金率", value: "swing" }, + { label: "波段率", value: "swing" }, + { label: "单日率", value: "ratio" }, + { label: "单日额", value: "amount" }, ] as const const rankChangeDayOptions = [1, 2, 3, 4, 5].map((value) => ({ - label: `${value} 日变化`, + label: value === 1 ? "近 1 天(相比昨日)" : `近 ${value} 天`, value: String(value), })) @@ -149,6 +158,7 @@ export function SectorRadarPage() { /> ) : null} onChange({ page: 1, view: value as RadarView })} /> + {search.view === "rank_change" ? ( - <> - + + onChange={(value) => onChange({ page: 1, rankChangeMetric: value as RadarMetricKind }) } /> onChange({ page: 1, rankChangeDays: Number(value) }) } /> - + ) : null} - ) } @@ -239,30 +250,36 @@ function RadarFilters({ function TabList({ ariaLabel, className, + inline = false, onChange, options, value, }: { ariaLabel: string className?: string + inline?: boolean onChange: (value: string) => void options: readonly { label: string; value: string }[] value: string }) { return (
{ariaLabel}
{options.map((option) => (