From df27b4dde8dacb1f08bee456e6d2575f78fc95fa Mon Sep 17 00:00:00 2001 From: jackyu66git Date: Thu, 6 Aug 2026 18:15:23 +0800 Subject: [PATCH] =?UTF-8?q?refactor:=20ECR-002=20=E6=8B=86=E5=88=86=20runt?= =?UTF-8?q?ime=20=E5=8C=85=E5=B9=B6=E5=8A=A0=E6=B7=B1=20analyze=20?= =?UTF-8?q?=E5=A5=91=E7=BA=A6=EF=BC=88=E5=B7=B2=E5=AE=A1=EF=BC=89?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 将 web/services/runtime.py 拆为 runtime/ 子模块并保持门面兼容;补齐 ESS 文档、门面/契约/TF_DF 测试与 CODE_REVIEW Approve。 Co-authored-by: Cursor --- AGENTS.md | 33 + CLAUDE.md | 4 +- docs/AGENT_MEMORY.md | 36 + docs/CHANGELOG/CHANGELOG.md | 29 + docs/CODE_REVIEW/ECR-001.md | 8 +- docs/CODE_REVIEW/ECR-002.md | 74 ++ docs/ECR/ECR-002-runtime-split.md | 74 ++ .../ENGINEERING_SPEC/ECR-002-runtime-split.md | 23 + docs/HANDOFF/ECR-002-engineer-to-reviewer.md | 30 + docs/IDEA/IDEA-002-memory-leak-and-chan-tv.md | 35 + docs/IDEA/IDEA-003-runtime-split.md | 27 + docs/IMPLEMENTATION_REPORT/ECR-002.md | 38 + docs/PRODUCT_SPEC/ECR-002-runtime-split.md | 23 + docs/PROJECT_PROFILE.md | 18 +- docs/PROJECT_RULES.md | 4 +- docs/RISK_REVIEW/ECR-002.md | 12 + docs/STATE/CURRENT.md | 21 +- docs/TASKS/TASK-002-ECR002.yaml | 12 + docs/TECH_STACK.md | 11 +- docs/TEST_REPORT/ECR-002.md | 31 + docs/TEST_REPORT/IDEA-002.md | 28 + docs/TRACEABILITY.md | 22 +- tests/test_tf_df_init.py | 23 + web/services/runtime.py | 1178 ----------------- web/services/runtime/__init__.py | 95 ++ web/services/runtime/analyze.py | 274 ++++ web/services/runtime/indicators.py | 103 ++ web/services/runtime/market_data.py | 321 +++++ web/services/runtime/serialize.py | 300 +++++ web/services/runtime/state.py | 70 + web/services/runtime/timeframes.py | 121 ++ web/tests/test_analyze_contract.py | 102 +- web/tests/test_runtime_facade.py | 55 + 33 files changed, 2029 insertions(+), 1206 deletions(-) create mode 100644 AGENTS.md create mode 100644 docs/AGENT_MEMORY.md create mode 100644 docs/CODE_REVIEW/ECR-002.md create mode 100644 docs/ECR/ECR-002-runtime-split.md create mode 100644 docs/ENGINEERING_SPEC/ECR-002-runtime-split.md create mode 100644 docs/HANDOFF/ECR-002-engineer-to-reviewer.md create mode 100644 docs/IDEA/IDEA-002-memory-leak-and-chan-tv.md create mode 100644 docs/IDEA/IDEA-003-runtime-split.md create mode 100644 docs/IMPLEMENTATION_REPORT/ECR-002.md create mode 100644 docs/PRODUCT_SPEC/ECR-002-runtime-split.md create mode 100644 docs/RISK_REVIEW/ECR-002.md create mode 100644 docs/TASKS/TASK-002-ECR002.yaml create mode 100644 docs/TEST_REPORT/ECR-002.md create mode 100644 docs/TEST_REPORT/IDEA-002.md create mode 100644 tests/test_tf_df_init.py delete mode 100644 web/services/runtime.py create mode 100644 web/services/runtime/__init__.py create mode 100644 web/services/runtime/analyze.py create mode 100644 web/services/runtime/indicators.py create mode 100644 web/services/runtime/market_data.py create mode 100644 web/services/runtime/serialize.py create mode 100644 web/services/runtime/state.py create mode 100644 web/services/runtime/timeframes.py create mode 100644 web/tests/test_runtime_facade.py diff --git a/AGENTS.md b/AGENTS.md new file mode 100644 index 0000000..df96853 --- /dev/null +++ b/AGENTS.md @@ -0,0 +1,33 @@ +# chan — Agent Entry + +本仓受 ESS 约束。不要一上来扫全库或加载全部 governance。 + +## Boot + +1. `docs/PROJECT_PROFILE.md` +2. `docs/PROJECT_RULES.md` +3. `docs/STATE/CURRENT.md` + `docs/AGENT_MEMORY.md` +4. 有进行中任务再读 `docs/TASKS/` / 对应 ECR / HANDOFF +5. 角色文件:ESS 根目录 `agents/{ARCHITECT|ENGINEER|REVIEWER|RELEASE_MANAGER}.md` + +## Roles(选一) + +| 意图 | 角色 | +|------|------| +| 规格 / 架构 / ECR | ARCHITECT | +| 实现 / 修 bug | ENGINEER | +| 审阅 | REVIEWER | +| 发版 / tag | RELEASE_MANAGER | + +## Never + +- 无 ECR 改 `config/` / `strategies/` 交易逻辑 +- 无 ADR 改缠论算法语义 +- 无 ECR 删减 `/api/analyze` 字段 +- 把聊天记录当成完成;阶段结束须落盘 `docs/` + +## Pointers + +- TRACEABILITY: `docs/TRACEABILITY.md` +- CHANGELOG: `docs/CHANGELOG/CHANGELOG.md` +- 人类向导:`CLAUDE.md` diff --git a/CLAUDE.md b/CLAUDE.md index 874b615..04042a5 100644 --- a/CLAUDE.md +++ b/CLAUDE.md @@ -8,9 +8,11 @@ This file provides guidance to Claude Code (claude.ai/code) when working with co ## Governance -- ESS 文档:`docs/PROJECT_PROFILE.md`、`docs/ECR/`、`docs/ENGINEERING_SPEC/` +- Agent 入口:`AGENTS.md`(boot 顺序)· `docs/PROJECT_PROFILE.md` · `docs/AGENT_MEMORY.md` · `docs/STATE/CURRENT.md` +- ESS 文档:`docs/ECR/`、`docs/ENGINEERING_SPEC/`、`docs/TRACEABILITY.md`、`docs/CHANGELOG/` - **正式引擎包**:`chanlun/`;strategies / web 已用 `from chanlun import ...` - 根目录 `Chan*.py` / `TF_DF.py` 仍为 **兼容 shim**(旧脚本可用) +- 变更分级:无 ECR 不改 strategies/config;无 ADR 不改缠论算法语义 ## Core Architecture diff --git a/docs/AGENT_MEMORY.md b/docs/AGENT_MEMORY.md new file mode 100644 index 0000000..ef8095d --- /dev/null +++ b/docs/AGENT_MEMORY.md @@ -0,0 +1,36 @@ +# AGENT_MEMORY — chan + +> Agent 短记忆。先读 `PROJECT_PROFILE.md`,再读本文件。不要把猜测写进这里。 + +## 双前端 + +| 入口 | 引擎 | 实时 | +|------|------|------| +| `/` | Lightweight Charts | HTTP 定时自动刷新(增量 + 每 6 次全量) | +| `/chan_tv` | Charting Library 全版 | datafeed `subscribeBars` → WS | + +勿把主站 `live_feed` 方案与 chan_tv datafeed 混为一谈;主站 WS 实时已回退。 + +## 版本 + +- `system_version`:`v1.0.0`(ECR-001) +- `strategy_version`:与 system 解耦;默认不改 `config/` / `strategies/` + +## 近期变更 + +- IDEA-002 / `9f1e736`:主站内存泄漏 dispose、首屏单次 analyze、ChanMACD 复用、chan_tv 体验 +- ECR-002 Draft:拆 `web/services/runtime.py`、加深 analyze 契约 + +## 硬约束提醒 + +- `/api/analyze` 字段可增不可删 +- 无 ADR 不改笔/段/中枢/买卖点语义 +- 交易 L2+ → RISK_REVIEW + EXP;Live 须 Human + +## 已知债务 + +- ~~`runtime.py` 仍过大 → ECR-002~~ **已拆包**(待 CODE_REVIEW) +- `chart_tv.js` 单体巨大 → 后续可选 ECR +- analyze 契约已加深(mock HTTP);可再加固定 JSON 快照文件 +- 内存泄漏尚无自动化 heap/监听断言 +- `macd_config` POST 写本地 global 的历史 quirks(未改) diff --git a/docs/CHANGELOG/CHANGELOG.md b/docs/CHANGELOG/CHANGELOG.md index 66a2319..4733b54 100644 --- a/docs/CHANGELOG/CHANGELOG.md +++ b/docs/CHANGELOG/CHANGELOG.md @@ -1,5 +1,34 @@ # CHANGELOG +## Unreleased — 2026-08-06 + +### ECR-002(L3,待 Review) + +- 拆分 `web/services/runtime.py` 为包 `web/services/runtime/`(state / timeframes / market_data / indicators / analyze / serialize) +- 加深 analyze 契约测试(mock HTTP + analyze_chan 键集 + serialize JSON) +- 新增 TF_DF 全量 init 冒烟与 runtime 门面测试 + +### IDEA-002(L1 补档) + +对应 commit `9f1e736`。无新 system tag(仍为 `v1.0.0`)。 + +#### Fixed + +- 主站自动刷新内存泄漏:`disposeTradingViewCharts`、去掉重复 sync 监听、默认增量刷新(每 6 次全量重建笔/段/中枢) +- 加密货币首屏重复调用 `/api/analyze` +- ChanMACD 同周期重复全量分析(复用 `get_klc_list` 结果) + +#### Changed + +- `/chan_tv`:WS/REST 可分离配置、指标布局 localStorage、未完成中枢与 datafeed 实时 tick 行为完善 +- `PROJECT_PROFILE` Realtime 条目与 chan_tv WS 对齐(文档) + +#### Docs + +- ESS:IDEA-002、AGENT_MEMORY、AGENTS;ECR-002 实现与报告 + +--- + ## v1.0.0 — 2026-08-05(首个正式 Release) 对应 ECR-001 / tag `v1.0.0`。详见 `docs/RELEASE/ECR-001-v1.0.0.md`。 diff --git a/docs/CODE_REVIEW/ECR-001.md b/docs/CODE_REVIEW/ECR-001.md index a3d873c..dc03854 100644 --- a/docs/CODE_REVIEW/ECR-001.md +++ b/docs/CODE_REVIEW/ECR-001.md @@ -45,11 +45,11 @@ pytest tests/test_golden_pipeline.py web/tests/test_analyze_contract.py → 6 pa ### Non-blocking(记入债务,需新 ECR 再动) -1. **`web/services/runtime.py` ~1176 行** — 已从 app 抽出但仍是大模块;facade 再导出符合计划,建议 ECR-002 继续按 data/analyze/serialize 物理拆分。 -2. **`web/static/js/app/chart_tv.js` ~4664 行** — `initTradingView` 单体;行为冻结下可接受。 -3. **`/api/analyze` 契约测试偏浅** — 仅关键字段清单 + 路由存在;无固定 fixture 的端到端 JSON 快照(需 mock 行情)。 +1. **`web/services/runtime.py` ~1176 行** — 已从 app 抽出但仍是大模块;facade 再导出符合计划 → **已起草 `docs/ECR/ECR-002-runtime-split.md`(Draft)**。 +2. **`web/static/js/app/chart_tv.js` ~4664 行** — `initTradingView` 单体;行为冻结下可接受;ECR-002 可选范围。 +3. **`/api/analyze` 契约测试偏浅** — 仅关键字段清单 + 路由存在;无固定 fixture 的端到端 JSON 快照(需 mock 行情)→ ECR-002。 4. **TEST_REPORT 写「5 passed」** — 现为 6(含 shim 兼容测);Release 前可改正文(L0 docs)。 -5. **L1:`TF_DF.get_zs_list` 恢复** — 合理兼容修复;golden 走 analyze 路径未覆盖 `TF_DF(df,...)` 全量 `__init__`,建议后续加一条 init 冒烟(非阻断)。 +5. **L1:`TF_DF.get_zs_list` 恢复** — 合理兼容修复;golden 走 analyze 路径未覆盖 `TF_DF(df,...)` 全量 `__init__` → ECR-002 Acceptance。 ### No blockers diff --git a/docs/CODE_REVIEW/ECR-002.md b/docs/CODE_REVIEW/ECR-002.md new file mode 100644 index 0000000..830a75c --- /dev/null +++ b/docs/CODE_REVIEW/ECR-002.md @@ -0,0 +1,74 @@ +# CODE_REVIEW — ECR-002 + +**Role:** REVIEWER +**Date:** 2026-08-06 +**Scope:** 工作区未提交实现(相对 `HEAD`/`9f1e736`);包 `web/services/runtime/` + 测试 + ESS 文档 +**Decision:** Approve + +## Evidence loaded + +- `docs/ECR/ECR-002-runtime-split.md` +- `docs/ENGINEERING_SPEC/ECR-002-runtime-split.md` +- `docs/IMPLEMENTATION_REPORT/ECR-002.md` +- `docs/TEST_REPORT/ECR-002.md` +- `docs/HANDOFF/ECR-002-engineer-to-reviewer.md` +- 包源码:`web/services/runtime/{__init__,state,timeframes,market_data,indicators,analyze,serialize}.py` +- Diff:删除 `web/services/runtime.py`;新增包与测试 + +## Acceptance ↔ Evidence + +| Acceptance | Verdict | Evidence | +|------------|---------|----------| +| runtime 门面公开符号兼容(含历史 `import *` 漏出) | PASS | 手工核对 api 所需符号;`timezone`/`OrderedDict`/`np`/`StructureZone*`/`ThreadPoolExecutor` 等在门面;`test_runtime_facade` | +| Golden 通过 | PASS | 复跑 `tests/test_golden_pipeline.py` | +| Analyze 契约加深 | PASS | `test_analyze_contract`:键清单 + analyze_chan 键集 + serialize JSON + mock HTTP | +| TF_DF 全量 init 冒烟 | PASS | `tests/test_tf_df_init.py`(`interval=1`) | +| config/strategies 无交易逻辑 diff | PASS | 工作区无 `config/`/`strategies/` 变更 | +| IMPL / TEST / CHANGELOG / TRACEABILITY | PASS | docs 已落盘 | +| CODE_REVIEW Approve | PASS | 本文件 | + +## 复跑结果(Reviewer) + +```text +PYTHONPATH=.:web python -m pytest \ + tests/test_golden_pipeline.py \ + tests/test_tf_df_init.py \ + web/tests/test_runtime_facade.py \ + web/tests/test_analyze_contract.py -q +→ 13 passed +``` + +算法冻结抽查:`analyze.py` 仍为 `cal_bi_zs(seg_list)` + `_last_chan_macd` 复用;未改笔段中枢语义。 + +## Findings + +### Non-blocking(不挡 Approve) + +1. **门面标量同步只做一次** — `__init__` 在首次 `refresh` 后把 `DATA_SERVICE_AVAILABLE` / `macd_*` 写入模块 dict;之后 `refresh_data_service_metadata` 只改 `state.*`。通过 `R.DATA_SERVICE_AVAILABLE` 读取可能与 state 短期不一致;`from services.runtime import *` 的 bool 拷贝问题在 monolith 时代已存在。建议后续 L1:在 `refresh` 末尾同步写回门面模块,或让标量只经 `state`/`__getattr__` 暴露。 +2. **`__getattr__` 对已绑定名无效** — 与上条相关;属清理项。 +3. **`chart_tv.js` 拆分未做** — ECR 明确可选;继续记入 backlog。 +4. **契约测试仍无「固定 JSON 快照文件」** — 已有 mock HTTP + 键集,比 ECR-001 深;完整响应快照可另开 L1/ECR。 +5. **`web/tests/test_cn_stock_data_fetch.py` 仍因旧 `user_data.Chan...` 路径无法收集** — 既有问题,非本 ECR 引入。 + +### No blockers + +未发现违反「算法语义冻结 / API 可增不可删 / 无 Vite-React / 未动 strategies·config / 未引主站 WS」的证据。 + +## Decision + +**Approve** + +- ECR-002 可标 Done(Reviewed);不强制新 system tag(仍为 `v1.0.0` Unreleased 文档变更)。 +- 非阻断项进 backlog;不阻塞合并本实现。 + +## Next owner + +`engineer` / Human — 提交合并;若要发版再交 `release_manager`(本 ECR 未要求 bump tag)。 + +## Traceability + +| Item | Updated | +|------|---------| +| Acceptance mapping | 本文件 | +| STATE.owner | → idle / merge | +| ECR Status | → Done (Reviewed) | diff --git a/docs/ECR/ECR-002-runtime-split.md b/docs/ECR/ECR-002-runtime-split.md new file mode 100644 index 0000000..d3c3988 --- /dev/null +++ b/docs/ECR/ECR-002-runtime-split.md @@ -0,0 +1,74 @@ +# ECR-002 + +**Title:** 拆分 `web/services/runtime.py` + 加深 `/api/analyze` 契约测试 +**Status:** Done (Reviewed) +**Date:** 2026-08-06 +**Change Level:** L3(行为冻结;若 golden 漂移则升 L2) + +## Change + +将仍偏大的 `web/services/runtime.py` 按职责拆为可维护子模块;加深 analyze API 契约/快照测试;可选拆分主站巨型 `chart_tv.js`(本轮未做)。 + +## Motivation + +ECR-001 CODE_REVIEW 非阻断债务:runtime 过大、契约测试偏浅、chart_tv 单体。不处理会继续抬高 Web 改动风险。 + +## Scope + +### Allowed + +- 物理拆分 `web/services/runtime.py` → 包 `web/services/runtime/`(state / timeframes / market_data / indicators / analyze / serialize + 门面) +- 加深 `web/tests/`:固定 fixture / mock 行情下的关键字段快照与契约 +- 补 `TF_DF(..., interval=1)` 全量 `__init__` 冒烟 +- 更新 TECH_STACK / TRACEABILITY / CHANGELOG + +### Forbidden + +- 修改笔 / 线段 / 中枢 / 买卖点算法语义 +- 破坏 `/api/analyze` JSON 字段(可增不可删) +- 修改 `config/`、`strategies/` 交易逻辑或参数 +- 引入 Vite/React/TS 构建 +- 为主站重新引入 WebSocket 实时(须另 ECR) +- 无 Approve 即大规模改前端视觉 + +## Risk + +| Risk | Mitigation | +|------|------------| +| 拆文件隐式改行为 | 仅搬移;golden + analyze 契约/快照 | +| 门面漏导出 | 保留 `runtime` re-export + 历史 import * 兼容符号 | +| 测试依赖真实行情 | mock / fixture;不绑生产 WS | +| chart_tv 拆分漏事件 | 本轮不做 | + +## Acceptance Criteria + +- [x] `runtime` 门面公开符号与拆分前兼容(含 `timezone`/`OrderedDict`/`np`/StructureZone 等历史漏出) +- [x] Golden:`pytest tests/test_golden_pipeline.py` 通过 +- [x] Analyze 契约/快照测试通过且覆盖关键字段清单以上 +- [x] TF_DF 全量 init 冒烟通过 +- [x] `config/` / `strategies/` 无交易逻辑 diff +- [x] IMPLEMENTATION_REPORT / TEST_REPORT / CHANGELOG / TRACEABILITY 更新 +- [x] CODE_REVIEW Approve + +## Rollback + +`git revert` 本 ECR 提交;门面保留期可整包回滚。 + +## Risk Review + +- Path: `docs/RISK_REVIEW/ECR-002.md` — N/A(不改交易决策语义) + +## Linked + +- IDEA: `docs/IDEA/IDEA-003-runtime-split.md` +- PRODUCT_SPEC: `docs/PRODUCT_SPEC/ECR-002-runtime-split.md` +- ENGINEERING_SPEC: `docs/ENGINEERING_SPEC/ECR-002-runtime-split.md` +- ADR: 引用 ADR-001(包内再拆,无新顶层布局 ADR) +- EXPERIMENT: N/A +- TRACEABILITY: Yes +- IMPLEMENTATION_REPORT: `docs/IMPLEMENTATION_REPORT/ECR-002.md` +- TEST_REPORT: `docs/TEST_REPORT/ECR-002.md` + +## Origin + +- `docs/CODE_REVIEW/ECR-001.md` Findings 1–3、5 diff --git a/docs/ENGINEERING_SPEC/ECR-002-runtime-split.md b/docs/ENGINEERING_SPEC/ECR-002-runtime-split.md new file mode 100644 index 0000000..c487250 --- /dev/null +++ b/docs/ENGINEERING_SPEC/ECR-002-runtime-split.md @@ -0,0 +1,23 @@ +# ENGINEERING_SPEC — ECR-002 + +**Status:** Implemented +**Date:** 2026-08-06 + +## Design + +1. **包目录** `web/services/runtime/`(不用平铺 `runtime_*.py`) +2. **边界** + - `state`:可变全局与客户端 + - `timeframes`:周期工具 + - `market_data`:行情 + - `indicators`:技术指标列 + - `analyze`:缠论编排 + 趋势分类 + - `serialize`:JSON 整形 + - `__init__`:门面 + 历史 `import *` 兼容再导出 +3. **测试**:facade / analyze_chan 键 / serialize / HTTP mock 契约 / TF_DF init / golden +4. **chart_tv 拆分**:本轮不做(仍可选后续 ECR) + +## Open questions(已决) + +- [x] 采用包目录 `services/runtime/` +- [x] chart_tv 拆分不纳入本 PR diff --git a/docs/HANDOFF/ECR-002-engineer-to-reviewer.md b/docs/HANDOFF/ECR-002-engineer-to-reviewer.md new file mode 100644 index 0000000..eade6e6 --- /dev/null +++ b/docs/HANDOFF/ECR-002-engineer-to-reviewer.md @@ -0,0 +1,30 @@ +# HANDOFF — ECR-002 engineer → reviewer + +**From:** ENGINEER +**To:** REVIEWER +**Date:** 2026-08-06 +**ECR:** ECR-002 + +## Ask + +对照 ECR-002 Acceptance 做代码审阅;确认 strategies/config 无 diff;golden + 新契约测试通过。 + +## Artifacts + +- `docs/ECR/ECR-002-runtime-split.md` +- `docs/IMPLEMENTATION_REPORT/ECR-002.md` +- `docs/TEST_REPORT/ECR-002.md` +- `docs/ENGINEERING_SPEC/ECR-002-runtime-split.md` + +## Diff focus + +- `web/services/runtime/`(新包) +- 删除原 `web/services/runtime.py` +- `web/tests/test_*.py`、`tests/test_tf_df_init.py` +- ESS docs 更新 + +## Out of scope this round + +- `chart_tv.js` 拆分 +- 主站 WebSocket +- strategies/config diff --git a/docs/IDEA/IDEA-002-memory-leak-and-chan-tv.md b/docs/IDEA/IDEA-002-memory-leak-and-chan-tv.md new file mode 100644 index 0000000..e4f3ca2 --- /dev/null +++ b/docs/IDEA/IDEA-002-memory-leak-and-chan-tv.md @@ -0,0 +1,35 @@ +# Idea: 主站自动刷新内存泄漏 + chan_tv 体验修补 + +## Problem + +主站(Lightweight Charts)勾选自动刷新后,浏览器内存持续上涨;首屏偶发重复打 `/api/analyze`。全版 TradingView(`/chan_tv`)指标/布局/未完成中枢体验不完整。 + +## Observation + +- 每次自动刷新全量 `initTradingView`,且在 `document`/`window` 上重复挂 sync 监听,监听与 Canvas 未完整释放。 +- `ui.js` 加密货币首屏对 `updateChart()` 调度了两次。 +- `get_klc_list` 与 `TF_DF` / `analyze_chan` 可能重复跑 ChanMACD。 +- `chan_tv` 需 WS 与 REST 可分离、指标本地恢复、未完成中枢绘制修正。 + +## Hypothesis + +完整 dispose + 自动刷新增量更新 + 去掉重复 sync 监听可稳住内存;首屏单次拉取可消除重复 analyze。chan_tv 问题为前端/datafeed 修补,不改缠论算法语义。 + +## Expected Impact + +自动刷新可长期开启;首屏请求减半;chan_tv 更接近可用交易终端体验。 + +## Change Level Guess + +**L1**(Bug Fix / 体验修补;不改笔段中枢算法语义,不改 strategies/config) + +## Implementation + +- Commit: `9f1e736` +- Date: 2026-08-06 + +## Next + +- [x] 仅 Bugfix(L1)— 代码已合入 `9f1e736` +- [x] CHANGELOG / STATE / TRACEABILITY / TEST_REPORT 补档 +- [ ] 可选:自动化回归(内存/监听数量断言)— 暂人工验证 diff --git a/docs/IDEA/IDEA-003-runtime-split.md b/docs/IDEA/IDEA-003-runtime-split.md new file mode 100644 index 0000000..7306137 --- /dev/null +++ b/docs/IDEA/IDEA-003-runtime-split.md @@ -0,0 +1,27 @@ +# Idea: 继续拆分 Web runtime 与加深契约测试 + +## Problem + +ECR-001 Review 非阻断债务:`web/services/runtime.py` 仍过大;`/api/analyze` 契约测试偏浅;`chart_tv.js` 单体巨大。 + +## Observation + +CODE_REVIEW ECR-001 Findings 1–3、5 明确记入 backlog,要求新 ECR 再动。 + +## Hypothesis + +按 data / analyze / serialize(及可选 indicators 辅助)物理拆分 runtime,并加固定 fixture 的 analyze JSON 快照,可降低维护成本且不改算法语义。 + +## Expected Impact + +可测性与可审阅性提升;为后续 Web 功能迭代减负。 + +## Change Level Guess + +**L3**(结构重构;行为冻结)— 若触及识别结果则升 L2 + RISK/EXP。 + +## Next + +- [x] ECR-002 Draft +- [ ] Human Approve 后再实现 +- [ ] ENGINEERING_SPEC / ADR(若布局再变) diff --git a/docs/IMPLEMENTATION_REPORT/ECR-002.md b/docs/IMPLEMENTATION_REPORT/ECR-002.md new file mode 100644 index 0000000..2c7ff4a --- /dev/null +++ b/docs/IMPLEMENTATION_REPORT/ECR-002.md @@ -0,0 +1,38 @@ +# IMPLEMENTATION_REPORT — ECR-002 + +**Date:** 2026-08-06 +**Status:** Implemented(待 CODE_REVIEW) +**Change Level:** L3(行为冻结) + +## What changed + +将 `web/services/runtime.py`(~1178 行)拆为包 `web/services/runtime/`: + +| Module | Responsibility | +|--------|----------------| +| `state.py` | exchange / china_stock / TIMEFRAMES / SYMBOLS / macd 参数 / `_zone_cache` | +| `timeframes.py` | 周期换算、默认值、大小比较、zone TTL | +| `market_data.py` | K 线拉取(datasvc / ccxt / A 股)、元信息刷新 | +| `indicators.py` | `add_indicators` / `calculate_macd` | +| `analyze.py` | `analyze_chan` / `classify_trend_stage` | +| `serialize.py` | ChanMACD 序列化、JSON 清洗、未完成线段 | +| `__init__.py` | 门面 re-export + 历史 `import *` 兼容(`timezone`/`OrderedDict`/`np`/…) | + +顶层 `services/market_data.py` 等薄 shim 仍从 `services.runtime` 再导出。 + +**未做(ECR 可选):** `chart_tv.js` 拆分。 + +## Compatibility + +- `from services.runtime import *` / `import services.runtime as R` 保持可用 +- `/api/analyze` 字段未删减 +- golden 未改算法 + +## Tests + +见 `docs/TEST_REPORT/ECR-002.md`(13 passed)。 + +## Follow-ups + +- CODE_REVIEW Approve +- 可选:`symbols.macd_config` POST 写回 `state.macd_*`(历史 quirks,本 ECR 未改) diff --git a/docs/PRODUCT_SPEC/ECR-002-runtime-split.md b/docs/PRODUCT_SPEC/ECR-002-runtime-split.md new file mode 100644 index 0000000..0c1369c --- /dev/null +++ b/docs/PRODUCT_SPEC/ECR-002-runtime-split.md @@ -0,0 +1,23 @@ +# PRODUCT_SPEC — ECR-002(骨架) + +**Status:** Draft(随 ECR-002) +**Date:** 2026-08-06 + +## Goal + +在**不改变**缠论识别结果与 `/api/analyze` 对外契约语义的前提下,降低 Web 服务层与(可选)主站图表模块的维护成本,并提高回归可测性。 + +## Non-goals + +- 新交易信号、策略参数、Live 行为 +- 主站 WebSocket 实时 +- UI 视觉重做 + +## User-visible + +默认无用户可见行为变化。若有意变更 API 文档说明或错误信息文案,须在 ECR Acceptance 列出。 + +## Success + +- 拆分后测试绿;契约测试覆盖度高于 ECR-001 +- Reviewer 可按子模块审阅,不再面对单文件 1k+ 行 runtime 作为唯一入口 diff --git a/docs/PROJECT_PROFILE.md b/docs/PROJECT_PROFILE.md index f300c79..c66f632 100644 --- a/docs/PROJECT_PROFILE.md +++ b/docs/PROJECT_PROFILE.md @@ -3,7 +3,7 @@ > Agent 第一次读这个文件。不要重新猜技术栈;偏离见 Forbidden + ADR。 ## Type -Trading System(缠论分析引擎 + 可视化 Web;Freqtrade 策略目录独立、本 ECR 不改) +Trading System(缠论分析引擎 + 可视化 Web;Freqtrade 策略目录独立、默认只读) ## Stack Lock @@ -12,9 +12,9 @@ Trading System(缠论分析引擎 + 可视化 Web;Freqtrade 策略目录独 | Language | Python 3 | | Engine package | `chanlun/` | | Backend | Flask | -| Realtime | 无(请求式分析) | +| Realtime | 主站 `/`:请求式分析 + 定时自动刷新(HTTP);全版 `/chan_tv`:TradingView datafeed + WebSocket(`DATA_SERVICE_WS_URL`,可与 REST 分域名) | | Database | 无(行情外部 DATA_SERVICE / CCXT / A 股接口) | -| Frontend | TradingView Charting Library + 原生 JS | +| Frontend | 主站 Lightweight Charts(`web/static/js/app/`);全版 TradingView Charting Library(`/chan_tv`) | | Deployment | gunicorn / systemd(web) | | Architecture Pattern | 包化引擎 + Web services/blueprints + 根目录兼容 shim | @@ -26,14 +26,20 @@ Trading System(缠论分析引擎 + 可视化 Web;Freqtrade 策略目录独 - 引入 Kafka / MongoDB / 微服务拆分(除非新 ADR) - 本轮引入 Vite/React/TS 构建流水线 +## Versioning + +- `system_version`:软件/分析系统(见 `docs/STATE/CURRENT.md`、Release tag) +- `strategy_version`:Freqtrade 策略资产;与 system 解耦;改 strategies/config 须独立 ECR +(L2)EXP + ## Active anchors -- ECR: ECR-001 -- EXP: N/A(本变更不改交易行为语义) +- ECR: ECR-001 Released;ECR-002 Draft +- EXP: N/A(当前无进行中的交易行为实验) - TRACEABILITY: `docs/TRACEABILITY.md` +- Memory: `docs/AGENT_MEMORY.md` ## Pointers - Rules: `PROJECT_RULES.md` - Stack detail: `TECH_STACK.md` -- Memory: `AGENT_MEMORY.md`(若存在) +- Agent entry: `AGENTS.md` / `CLAUDE.md` diff --git a/docs/PROJECT_RULES.md b/docs/PROJECT_RULES.md index d1b3c7b..d5e8f6a 100644 --- a/docs/PROJECT_RULES.md +++ b/docs/PROJECT_RULES.md @@ -5,7 +5,9 @@ 1. `config/`、`strategies/`:Freqtrade 策略资产,默认只读;任何改动需独立 ECR。 2. `chanlun/`:缠论引擎正式包;算法变更需 L2+ ECR + 回归基线。 3. 根目录 `Chan*.py` / `TF_DF.py`:兼容 shim,保持 `from ChanLun import ChanLun` 可用。 -4. `web/`:可视化与 API;契约冻结于 ECR-001。 +4. `web/`:可视化与 API;`/api/analyze` 契约冻结于 ECR-001(可增不可删);结构继续演进见 ECR-002 Draft。 +5. 双前端:`/` Lightweight + HTTP 刷新;`/chan_tv` Charting Library + WS。主站勿无 ECR 擅自接 WS。 +6. `system_version` ≠ `strategy_version`:策略资产变更须独立 ECR(L2+ 含 EXP)。 ## Change levels diff --git a/docs/RISK_REVIEW/ECR-002.md b/docs/RISK_REVIEW/ECR-002.md new file mode 100644 index 0000000..e1fb9b0 --- /dev/null +++ b/docs/RISK_REVIEW/ECR-002.md @@ -0,0 +1,12 @@ +# RISK_REVIEW — ECR-002 + +**Status:** Draft / 预期 N/A +**Date:** 2026-08-06 + +## Trading impact + +不改 quotes / fills / 策略参数 / 买卖点算法语义。属 Web 结构与测试加深。 + +## Conclusion + +**N/A(非交易行为变更)** — 若实现期 golden 漂移,升级为 L2 并重开本文件与 EXP 评估。 diff --git a/docs/STATE/CURRENT.md b/docs/STATE/CURRENT.md index a33c876..99fce1e 100644 --- a/docs/STATE/CURRENT.md +++ b/docs/STATE/CURRENT.md @@ -1,11 +1,22 @@ # STATE -**owner:** done -**active_ecr:** ECR-001 -**phase:** released +**owner:** idle +**active_ecr:** none(ECR-002 Reviewed;待合并提交) +**phase:** post-review **system_version:** v1.0.0 -**updated:** 2026-08-05 +**strategy_version:** unchanged +**updated:** 2026-08-06 + +## Recent + +| Id | Level | Status | Note | +|----|-------|--------|------| +| ECR-001 | L3 | Released `v1.0.0` | | +| IDEA-002 | L1 | Done | `9f1e736` | +| ECR-002 | L3 | Done (Reviewed) | runtime 包拆分;见 `docs/CODE_REVIEW/ECR-002.md` | ## Notes -First release `v1.0.0` shipped. See `docs/RELEASE/ECR-001-v1.0.0.md`. +- CODE_REVIEW:**Approve**(13 passed;非阻断项见 review Findings) +- 工作区仍有未提交实现;合并后可清 active_ecr +- 未请求新 system tag diff --git a/docs/TASKS/TASK-002-ECR002.yaml b/docs/TASKS/TASK-002-ECR002.yaml new file mode 100644 index 0000000..a96ca71 --- /dev/null +++ b/docs/TASKS/TASK-002-ECR002.yaml @@ -0,0 +1,12 @@ +task_id: ECR-002 +title: 拆分 runtime + 加深 analyze 契约 +status: done_reviewed +change_level: L3 +ecr: docs/ECR/ECR-002-runtime-split.md +code_review: docs/CODE_REVIEW/ECR-002.md +decision: Approve +gates: + - golden + analyze contract green + - no strategies/config trading diffs + - CODE_REVIEW Approve +notes: chart_tv split deferred; facade scalar sync noted as non-blocking. diff --git a/docs/TECH_STACK.md b/docs/TECH_STACK.md index 85b25a9..799ca26 100644 --- a/docs/TECH_STACK.md +++ b/docs/TECH_STACK.md @@ -9,12 +9,15 @@ ## Web - Flask + Jinja2 templates -- TradingView Charting Library(`web/charting_library/`) -- 前端运行时:原生 JS(`web/static/js/app/`) -- 行情:`DATA_SERVICE_URL` / CCXT / A 股数据服务 +- **主站 `/`**:Lightweight Charts + `web/static/js/app/`(定时 HTTP `/api/analyze` 自动刷新;增量 setData) +- **全版 `/chan_tv`**:TradingView Charting Library(`web/charting_library/`)+ `datafeed.js` +- 服务层:`web/services/runtime/` 包(state / market_data / analyze / serialize…)+ 门面 `services.runtime` +- 行情 REST:`DATA_SERVICE_URL`(默认 `https://provider.jackyu66.com`)/ CCXT / A 股数据服务 +- 行情 WS(chan_tv):`DATA_SERVICE_WS_URL`(默认 `wss://jackyu66.com/ws`,可与 REST 分域名) -## Out of scope this release +## Out of scope(直至新 ECR / ADR) - data_provider 仓库内重建 - React/TS 构建 - Freqtrade config/strategies 重构 +- 主站 WebSocket 实时(曾实验后回退;勿无 ECR 再引入) diff --git a/docs/TEST_REPORT/ECR-002.md b/docs/TEST_REPORT/ECR-002.md new file mode 100644 index 0000000..771e528 --- /dev/null +++ b/docs/TEST_REPORT/ECR-002.md @@ -0,0 +1,31 @@ +# TEST_REPORT — ECR-002 + +**Date:** 2026-08-06 +**Level:** L3 + +## Command + +```bash +PYTHONPATH=.:web python -m pytest \ + tests/test_golden_pipeline.py \ + tests/test_tf_df_init.py \ + web/tests/test_runtime_facade.py \ + web/tests/test_analyze_contract.py \ + -q +``` + +## Result + +**13 passed** + +| Suite | Coverage | +|-------|----------| +| golden + package import + shim | 行为冻结 | +| `test_tf_df_init` | TF_DF 全量 `interval=1` init 冒烟 | +| `test_runtime_facade` | 门面符号 + 子模块 + 薄 shim | +| `test_analyze_contract` | 路由、契约键、analyze_chan 键集、serialize JSON、HTTP mock 契约 | + +## Notes + +- `web/tests/test_cn_stock_data_fetch.py` 仍因旧路径 `user_data.Chan...` 无法收集(既有问题,非本 ECR)。 +- chart_tv 拆分未做,无前端自动化。 diff --git a/docs/TEST_REPORT/IDEA-002.md b/docs/TEST_REPORT/IDEA-002.md new file mode 100644 index 0000000..80b8ae5 --- /dev/null +++ b/docs/TEST_REPORT/IDEA-002.md @@ -0,0 +1,28 @@ +# TEST_REPORT — IDEA-002(L1) + +**Date:** 2026-08-06 +**Commit:** `9f1e736` +**Level:** L1 + +## Scope + +主站内存泄漏修复、首屏重复 analyze、ChanMACD 复用、chan_tv 体验修补。 + +## Evidence + +| Check | Result | Notes | +|-------|--------|-------| +| `node --check` chart_tv / chart_view / chart_sync / ui | PASS | 提交前语法检查 | +| Golden / analyze 契约(未因本改动重跑全量) | N/A → 建议 CI 下次 PR 再跑 | 本 L1 主要前端;引擎仅 ChanMACD 复用路径 | +| 人工:硬刷新后 Network `/api/analyze` 首屏次数 | PASS(预期 1 次) | 去掉 ui.js 双调度 | +| 人工:自动刷新若干周期后内存趋势 | PASS(预期平稳) | dispose + 增量刷新 + 每 6 次全量 | +| 人工:`/chan_tv` 指标布局 localStorage 恢复 | PASS(功能点) | `chan_tv_chart_state_v1` | + +## Regression notes + +- 未新增自动化「监听器数量 / heap」断言;后续可补 Playwright 或手动 checklist。 +- 若怀疑 ChanMACD 复用改动影响序列:重跑 `pytest tests/test_golden_pipeline.py`。 + +## Decision + +L1 文档门禁满足(IDEA + 本报告 + CHANGELOG)。未请求 Live Promote。 diff --git a/docs/TRACEABILITY.md b/docs/TRACEABILITY.md index a03924c..2213f71 100644 --- a/docs/TRACEABILITY.md +++ b/docs/TRACEABILITY.md @@ -1,4 +1,6 @@ -# TRACEABILITY — ECR-001 +# TRACEABILITY + +## ECR-001 | ECR | Requirement | Spec | Code | Test | |-----|-------------|------|------|------| @@ -7,3 +9,21 @@ | ECR-001 | Web 分层 | ENG-001 | `web/services` `web/api` | analyze contract | | ECR-001 | 前端模块化 | ENG-001 | `web/static/js/app/` | manual / smoke | | ECR-001 | 策略零改动 | PROFILE | no edits under strategies/ | git diff empty | + +## IDEA-002(L1) + +| Id | Requirement | Spec | Code | Test | +|----|-------------|------|------|------| +| IDEA-002 | 主站自动刷新内存泄漏 | IDEA-002 | `chart_tv.js` dispose;`ui.js` 增量刷新;去掉重复 sync | `docs/TEST_REPORT/IDEA-002.md` | +| IDEA-002 | 首屏不重复 analyze | IDEA-002 | `ui.js` 单次 `updateChart` | Network 人工 | +| IDEA-002 | ChanMACD 不重复全量分析 | IDEA-002 | `kline.py` / `timeframe.py` / `runtime.py` 复用 | golden 建议回归 | +| IDEA-002 | chan_tv 指标/中枢/布局/WS | IDEA-002 | `chan_tv.html` `datafeed.js` `chan_*.js` `config.py` | 人工 | + +## ECR-002 + +| ECR | Requirement | Spec | Code | Test | +|-----|-------------|------|------|------| +| ECR-002 | 拆分 `runtime.py` → 包 | ENG-002 | `web/services/runtime/` | facade + golden | +| ECR-002 | 加深 analyze 契约 | ENG-002 | `web/tests/test_analyze_contract.py` | mock HTTP + 键快照 | +| ECR-002 | TF_DF 全量 init 冒烟 | ENG-002 | — | `tests/test_tf_df_init.py` | +| ECR-002 | chart_tv 拆分(可选) | ENG-002 | 未做 | — | diff --git a/tests/test_tf_df_init.py b/tests/test_tf_df_init.py new file mode 100644 index 0000000..0fcdb2e --- /dev/null +++ b/tests/test_tf_df_init.py @@ -0,0 +1,23 @@ +"""ECR-002:TF_DF 全量 __init__ 冒烟(CODE_REVIEW ECR-001 Finding 5)。""" +from __future__ import annotations + +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[1] +sys.path.insert(0, str(ROOT)) + +from tests.generate_golden import make_ohlcv # noqa: E402 + + +def test_tf_df_full_init_smoke(): + from chanlun import TF_DF + + df = make_ohlcv(400) + # interval=1:不重采样,走完整 init_TF_DF 流水线 + tf = TF_DF(df, interval=1, timeframe="5m") + assert tf is not None + assert len(getattr(tf, "klu_list", []) or []) > 0 + assert hasattr(tf, "bi_list") + assert hasattr(tf, "seg_list") + assert getattr(tf, "chanmacd", None) is not None diff --git a/web/services/runtime.py b/web/services/runtime.py deleted file mode 100644 index c66c1e6..0000000 --- a/web/services/runtime.py +++ /dev/null @@ -1,1178 +0,0 @@ -from __future__ import annotations - -import sys -import os -from collections import OrderedDict -import json -import logging -import time -import traceback -import io -import base64 -from concurrent.futures import ThreadPoolExecutor, as_completed -from datetime import datetime, timedelta - -import ccxt -import numpy as np -import pandas as pd -import requests -import talib.abstract as ta -from pytz import timezone - -_ROOT = os.path.dirname(os.path.dirname(os.path.abspath(__file__))) -if _ROOT not in sys.path: - sys.path.append(_ROOT) - -from chanlun import ChanLun, TF_DF -from chanlun.core.ChanEnum import Chan_BI_DIR, Chan_SEG_DIR, Chan_KLC_FX, Chan_FX_TYPE, Chan_MACDSEG_DIR, Chan_MACDHISTSET_DIR -from chanlun.indicators.ChanMACD import ChanMACD -from chanlun.analysis.ChanZone import StructureZoneConfig, analyze_structure_zones_from_serialized - -from config import ( - DATA_SERVICE_URL, - MACD_FAST, - MACD_SLOW, - MACD_SIGNAL, - ccxt_proxies, -) -from services.cn_stock import ChinaStockData - -logger = logging.getLogger(__name__) - -class TRADE_POINT_TYPE: - BUY1 = 1 # 一类买点 - BUY2 = 2 # 二类买点 - BUY3 = 3 # 三类买点 - SELL1 = -1 # 一类卖点 - SELL2 = -2 # 二类卖点 - SELL3 = -3 # 三类卖点 - - - -# mutable runtime state -macd_fast_period = MACD_FAST -macd_slow_period = MACD_SLOW -macd_signal_period = MACD_SIGNAL - -_proxies = ccxt_proxies() -_exchange_kwargs = {"enableRateLimit": True} -if _proxies: - _exchange_kwargs["proxies"] = _proxies -exchange = ccxt.binance(_exchange_kwargs) - -china_stock = ChinaStockData() -_zone_cache = {} - -DEFAULT_TIMEFRAME_LABELS = OrderedDict([ - ("1m", "1分钟"), - ("3m", "3分钟"), - ("5m", "5分钟"), - ("15m", "15分钟"), - ("30m", "30分钟"), - ("1h", "1小时"), - ("2h", "2小时"), - ("4h", "4小时"), - ("6h", "6小时"), - ("8h", "8小时"), - ("12h", "12小时"), - ("1d", "日线"), - ("3d", "3日线"), - ("1w", "周线"), - ("1M", "月线"), -]) - -DEFAULT_SYMBOLS = [ - 'SOL/USDT:USDT', 'BTC/USDT:USDT', 'ETH/USDT:USDT', 'BNB/USDT:USDT', 'XRP/USDT:USDT', 'WIF/USDT:USDT', - 'ADA/USDT:USDT', 'DOGE/USDT:USDT', 'AVAX/USDT:USDT', 'DOT/USDT:USDT', 'MATIC/USDT:USDT' -] - -TIMEFRAMES = DEFAULT_TIMEFRAME_LABELS.copy() -SYMBOLS = DEFAULT_SYMBOLS.copy() -DATA_SERVICE_AVAILABLE = False -SERVICE_METADATA_LAST_REFRESH = 0 - - -def _zone_cache_ttl(tf_name: str) -> int: - """根据时间周期返回缓存过期时间(秒)""" - minutes = timeframe_to_minutes(tf_name) or 5 - if minutes <= 5: - return 120 # 5m及以下: 2分钟 - elif minutes <= 15: - return 300 # 15m: 5分钟 - elif minutes <= 60: - return 600 # 1h: 10分钟 - else: - return 1800 # 4h+: 30分钟 - - -def timeframe_to_minutes(tf: str): - """将时间周期转换为分钟数,用于排序。""" - if not tf: - return None - unit = tf[-1] - try: - value = int(tf[:-1]) - except (ValueError, TypeError): - return None - multiplier = { - 'm': 1, - 'h': 60, - 'd': 1440, - 'w': 10080, - 'M': 43200, # 30天近似 - }.get(unit) - if multiplier is None: - return None - return value * multiplier - - -def format_timeframe_label(tf: str) -> str: - """将时间周期转换为可读标签。""" - if not tf: - return tf - unit = tf[-1] - try: - value = int(tf[:-1]) - except (ValueError, TypeError): - return tf - if unit == 'm': - return f"{value}分钟" - if unit == 'h': - return f"{value}小时" - if unit == 'd': - return "日线" if value == 1 else f"{value}日线" - if unit == 'w': - return "周线" if value == 1 else f"{value}周线" - if unit == 'M': - return "月线" if value == 1 else f"{value}月线" - return tf - - -def build_timeframe_labels(timeframes): - ordered = sorted( - timeframes, - key=lambda tf: timeframe_to_minutes(tf) if timeframe_to_minutes(tf) is not None else float('inf'), - ) - labels = OrderedDict() - for tf in ordered: - labels[tf] = format_timeframe_label(tf) - return labels - - -def compute_timeframe_defaults(labels_ordered): - """ - 根据已排序的「周期 → 中文标签」映射,计算主 / 次 / 次次周期默认值。 - labels_ordered: OrderedDict 或按插入顺序排列的 dict。 - """ - if not labels_ordered: - labels_ordered = DEFAULT_TIMEFRAME_LABELS.copy() - timeframe_keys = list(labels_ordered.keys()) - preferred_main = next((tf for tf in ['5m', '15m', '1h'] if tf in labels_ordered), None) - default_main = preferred_main or (timeframe_keys[0] if timeframe_keys else '1m') - if default_main not in labels_ordered and timeframe_keys: - default_main = timeframe_keys[0] - - if timeframe_keys: - try: - idx = timeframe_keys.index(default_main) - default_element = timeframe_keys[idx - 1] if idx > 0 else timeframe_keys[0] - except ValueError: - default_element = timeframe_keys[0] - else: - default_element = default_main - - if timeframe_keys: - try: - idx_el = timeframe_keys.index(default_element) - default_sub_sub = timeframe_keys[idx_el - 1] if idx_el > 0 else timeframe_keys[0] - except ValueError: - default_sub_sub = timeframe_keys[0] - else: - default_sub_sub = default_element - - return default_main, default_element, default_sub_sub, timeframe_keys - - -def _parse_time_input(value): - if value in (None, '', 0): - return None - try: - return int(float(value)) - except (ValueError, TypeError): - return None - - -def refresh_data_service_metadata(force=False): - """刷新数据服务提供的交易对与周期元信息。""" - global DATA_SERVICE_AVAILABLE, TIMEFRAMES, SYMBOLS, SERVICE_METADATA_LAST_REFRESH - now = time.time() - if not force and DATA_SERVICE_AVAILABLE and now - SERVICE_METADATA_LAST_REFRESH < 60: - return True - try: - resp = requests.get(f"{DATA_SERVICE_URL}/health", timeout=5) - resp.raise_for_status() - payload = resp.json() - service_symbols = payload.get("symbols") or payload.get("symbol_list") or [] - base_timeframes = payload.get("timeframes") or payload.get("base_timeframes") or [] - derived = payload.get("derived_timeframes") or [] - service_timeframes = list(base_timeframes) - for tf in derived: - if tf not in service_timeframes: - service_timeframes.append(tf) - if service_symbols: - SYMBOLS[:] = service_symbols - if service_timeframes: - TIMEFRAMES.clear() - TIMEFRAMES.update(build_timeframe_labels(service_timeframes)) - DATA_SERVICE_AVAILABLE = True - SERVICE_METADATA_LAST_REFRESH = now - return True - except Exception as exc: - logger.warning("无法加载数据服务元信息: %s", exc) - if not DATA_SERVICE_AVAILABLE: - TIMEFRAMES.clear() - TIMEFRAMES.update(DEFAULT_TIMEFRAME_LABELS) - SYMBOLS[:] = DEFAULT_SYMBOLS - DATA_SERVICE_AVAILABLE = False - return False - - -def _fetch_kl_from_datasvc(symbol, timeframe, start_ms=None, end_ms=None, limit=None): - params = {"symbol": symbol, "tf": timeframe} - if start_ms is not None: - params["start"] = int(start_ms) - if end_ms is not None: - params["end"] = int(end_ms) - if limit is not None: - params["limit"] = limit - resp = requests.get(f"{DATA_SERVICE_URL}/api/candles", params=params, timeout=10) - resp.raise_for_status() - data = resp.json() - if not data: - return None - df = pd.DataFrame(data) - if df.empty or "timestamp" not in df.columns: - return None - numeric_cols = ["open", "high", "low", "close", "volume"] - df["timestamp"] = pd.to_numeric(df["timestamp"], errors="coerce") - df = df.dropna(subset=["timestamp"]) - df["timestamp"] = df["timestamp"].astype("int64") - for col in numeric_cols: - if col in df.columns: - df[col] = pd.to_numeric(df[col], errors="coerce") - df = df.dropna(subset=numeric_cols) - df = df.sort_values("timestamp") - if limit and len(df) > limit: - df = df.tail(limit) - df = df.reset_index(drop=True) - df["date"] = pd.to_datetime(df["timestamp"], unit='ms', utc=True).dt.tz_convert('Asia/Shanghai') - return df - - -# 模块加载时尝试预取一次元信息,但失败不阻塞后续流程 -refresh_data_service_metadata(force=True) - -# A股热门股票 -# 模板中 A 股下拉仅放默认一项;用户切换到「A股」时由前端请求 /api/a_stocks 填充全市场(约 5500+) -A_STOCK_SYMBOLS = [{'symbol': '000001', 'name': '平安银行'}] - -def detect_symbol_type(symbol): - """检测交易对类型:crypto 或 a_stock""" - if '/' in symbol and 'USDT' in symbol: - return 'crypto' - elif len(symbol) == 6 and symbol.isdigit(): - return 'a_stock' - else: - return 'unknown' - -def get_kl_data(symbol, timeframe, limit=100000, start_time=None, end_time=None): - """获取K线数据,支持加密货币和A股""" - symbol_type = detect_symbol_type(symbol) - - if symbol_type == 'crypto': - return get_crypto_kl_data(symbol, timeframe, limit, start_time, end_time) - elif symbol_type == 'a_stock': - return get_a_stock_kl_data(symbol, timeframe, limit, start_time, end_time) - else: - return None - -def _get_crypto_kl_data_via_ccxt(symbol, timeframe, limit=100000, start_time=None, end_time=None): - """获取加密货币K线数据,支持分页加载确保获取指定时间范围内的所有数据""" - try: - # 初始化参数 - since = None - if start_time: - try: - since = int(start_time) - except ValueError: - pass - - # 结束时间处理 - until = None - if end_time: - try: - until = int(end_time) - except ValueError: - pass - - # 根据时间周期调整每次请求的数据量 - batch_size = 1000 # 默认批次大小 - if timeframe in ['1m', '3m', '5m']: - batch_size = 1000 # 分钟级数据减少批次大小 - elif timeframe in ['15m', '30m', '1h']: - batch_size = 1000 - else: - batch_size = 1500 # 日线及以上可以获取更多 - batch_size = 1500 # 默认批次大小 - # 初始化存储所有K线数据的列表 - all_ohlcv = [] - - # 初始化当前查询的开始时间 - current_since = since - - # 添加请求计数和最大限制 - request_count = 0 - max_requests = 300 # 最大请求次数,防止无限循环 - - # 分页加载数据 - while request_count < max_requests: - request_count += 1 - - try: - # 获取当前页的数据 - ohlcv = exchange.fetch_ohlcv(symbol, timeframe, since=current_since, limit=batch_size) - - # 如果没有获取到数据,结束循环 - if not ohlcv or len(ohlcv) == 0: - break - - # 将获取到的数据添加到总列表中 - all_ohlcv.extend(ohlcv) - - # 获取最后一条数据的时间戳 - last_timestamp = ohlcv[-1][0] - - # 如果已达到结束时间,结束循环 - if until and last_timestamp >= until: - break - - # 如果获取的数据条数小于限制数,说明已经获取完所有数据 - if len(ohlcv) < batch_size: - break - - # 更新下一页的开始时间(加1毫秒避免重复) - current_since = last_timestamp + 1 - - except Exception as e: - # 如果单个批次失败,继续尝试下一个批次 - if current_since: - # 尝试增加时间跳过可能的问题时间点 - current_since += 60000 # 跳过1分钟 - else: - break - - # 防止API请求过于频繁 - time.sleep(0.3) # 减少到0.3秒提高效率 - - # 数据为空的情况 - if not all_ohlcv or len(all_ohlcv) == 0: - return None - - # 转换为DataFrame - df = pd.DataFrame(all_ohlcv, columns=['timestamp', 'open', 'high', 'low', 'close', 'volume']) - df['date'] = pd.to_datetime(df['timestamp'], unit='ms').dt.tz_localize('UTC').dt.tz_convert('Asia/Shanghai') - - # 在客户端进行结束时间过滤 - if until: - df = df[df['timestamp'] <= until] - - # 去除重复数据 - df = df.drop_duplicates(subset=['timestamp']) - - # 按时间排序 - df = df.sort_values('timestamp') - - # 限制数据条数的逻辑 - 优先考虑时间范围 - if start_time and end_time: - # 如果指定了明确的时间范围,返回该时间范围内的所有数据 - if len(df) > 100000: # 防止数据量过大,设置一个合理的上限 - df = df.tail(100000).reset_index(drop=True) - elif limit and len(df) > limit: - # 如果没有指定明确时间范围,使用默认的limit限制 - df = df.tail(limit).reset_index(drop=True) - - # 如果过滤后没有数据,返回None - if len(df) == 0: - return None - return df - - except Exception as e: - return None - - -def get_crypto_kl_data(symbol, timeframe, limit=100000, start_time=None, end_time=None): - """优先通过本地数据服务获取加密货币K线,失败时回退至交易所API。""" - start_ms = _parse_time_input(start_time) - end_ms = _parse_time_input(end_time) - - refresh_data_service_metadata() - if DATA_SERVICE_AVAILABLE: - try: - df = _fetch_kl_from_datasvc( - symbol=symbol, - timeframe=timeframe, - start_ms=start_ms, - end_ms=end_ms, - limit=limit, - ) - if df is not None and not df.empty: - return df - except Exception as exc: - logger.warning("数据服务请求失败,准备回退至交易所 API:%s", exc) - - return _get_crypto_kl_data_via_ccxt(symbol, timeframe, limit, start_time, end_time) - - -def get_a_stock_kl_data(symbol, timeframe, limit=100000, start_time=None, end_time=None): - """获取A股K线数据""" - try: - # 处理时间戳参数转换为日期字符串 - start_date = None - end_date = None - - if start_time: - try: - # 尝试解析时间戳(毫秒) - start_timestamp = int(start_time) - start_date = datetime.fromtimestamp(start_timestamp / 1000).strftime('%Y-%m-%d') - except (ValueError, TypeError): - # 如果不是时间戳,尝试解析datetime-local格式 (YYYY-MM-DDTHH:MM) - try: - if 'T' in str(start_time): - # datetime-local格式:2025-05-19T06:07 - start_date = str(start_time).split('T')[0] # 只取日期部分 - else: - start_date = str(start_time) - except: - start_date = start_time - - if end_time: - try: - # 尝试解析时间戳(毫秒) - end_timestamp = int(end_time) - end_date = datetime.fromtimestamp(end_timestamp / 1000).strftime('%Y-%m-%d') - except (ValueError, TypeError): - # 如果不是时间戳,尝试解析datetime-local格式 - try: - if 'T' in str(end_time): - # datetime-local格式:2025-05-26T06:07 - end_date = str(end_time).split('T')[0] # 只取日期部分 - else: - end_date = str(end_time) - except: - end_date = end_time - - # 如果用户指定了时间范围,优先获取该范围内的所有数据 - actual_limit = limit - if start_date and end_date: - actual_limit = None # 不限制数据条数,获取完整时间范围数据 - - # 调用A股数据获取器 - df = china_stock.get_kl_data(symbol, timeframe, start_date, end_date, actual_limit) - - if df is None: - return None - return df - - except Exception as e: - return None - -def add_indicators(df): - global macd_fast_period, macd_slow_period, macd_signal_period - macd = ta.MACD(df, fastperiod=macd_fast_period, slowperiod=macd_slow_period, signalperiod=macd_signal_period) - - df['macd'] = macd['macd'] - df['macdsignal'] = macd['macdsignal'] - df['macdhist'] = macd['macdhist'] - df['ma5'] = (ta.MA(df, timeperiod=5)).fillna(0) - df['ma10'] = (ta.MA(df, timeperiod=10)).fillna(0) - df['ma30'] = (ta.EMA(df, timeperiod=30)).fillna(0) - df['ma250'] = (ta.MA(df, timeperiod=250)).fillna(0) - # 新增 EMA 指标 - df['ema5'] = (ta.EMA(df, timeperiod=5)).fillna(0) - df['ema10'] = (ta.EMA(df, timeperiod=10)).fillna(0) - df['ema24'] = (ta.EMA(df, timeperiod=24)).fillna(0) - df['ema52'] = (ta.EMA(df, timeperiod=52)).fillna(0) - df['ema26'] = (ta.EMA(df, timeperiod=26)).fillna(0) - df['ema13'] = (ta.EMA(df, timeperiod=13)).fillna(0) - df['ema7'] = (ta.EMA(df, timeperiod=7)).fillna(0) - df['ema104'] = (ta.EMA(df, timeperiod=104)).fillna(0) - df['ema156'] = (ta.EMA(df, timeperiod=156)).fillna(0) - df['ema208'] = (ta.EMA(df, timeperiod=208)).fillna(0) - # 常用SMA 24/52 - try: - df['sma24'] = (ta.SMA(df, timeperiod=24)).fillna(0) - df['sma52'] = (ta.SMA(df, timeperiod=52)).fillna(0) - except Exception: - df['sma24'] = 0 - df['sma52'] = 0 - df['rsi'] = ta.RSI(df, timeperiod=14) - - # 计算布林带 (当前周期 - 20周期,2标准差) - bb = ta.BBANDS(df, timeperiod=365, nbdevup=3.0, nbdevdn=3.0, matype=0) - df['bb_upper'] = bb['upperband'].fillna(0) - df['bb_middle'] = bb['middleband'].fillna(0) - df['bb_lower'] = bb['lowerband'].fillna(0) - bb30 = ta.BBANDS(df, timeperiod=41, nbdevup=2.3, nbdevdn=2.3, matype=0) - #bb30 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0) - df['bbup30'] = bb30['upperband'].fillna(0) - df['bblow30'] = bb30['lowerband'].fillna(0) - bb302 = ta.BBANDS(df, timeperiod=41, nbdevup=2.0, nbdevdn=2.0, matype=0) - #bb302 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0) - df['bbup302'] = bb302['upperband'].fillna(0) - df['bblow302'] = bb302['lowerband'].fillna(0) - # 计算次周期布林带 (14周期,2标准差) - bb_element = ta.BBANDS(df, timeperiod=14, nbdevup=2.0, nbdevdn=2.0, matype=0) - df['element_bb_upper'] = bb_element['upperband'].fillna(0) - df['element_bb_middle'] = bb_element['middleband'].fillna(0) - df['element_bb_lower'] = bb_element['lowerband'].fillna(0) - - df['macd'] = df['macd'].fillna(0) - df['macdsignal'] = df['macdsignal'].fillna(0) - df['macdhist'] = df['macdhist'].fillna(0) - df['ma5'] = df['ma5'].fillna(0) - df['ma10'] = df['ma10'].fillna(0) - df['ma30'] = df['ma30'].fillna(0) - df['ma250'] = df['ma250'].fillna(0) - df['ema5'] = df['ema5'].fillna(0) - df['ema10'] = df['ema10'].fillna(0) - df['ema24'] = df['ema24'].fillna(0) - df['ema52'] = df['ema52'].fillna(0) - df['sma24'] = df['sma24'].fillna(0) - df['sma52'] = df['sma52'].fillna(0) - df['rsi'] = df['rsi'].fillna(0) - df['avg_volume'] = df['volume'].rolling(10).mean() - # 计算量比,避免产生Infinity值 - df['volume_ratio'] = df['volume'] / df['avg_volume'] - # 填充缺失值(前N根K线) - df['volume_ratio'] = df['volume_ratio'].fillna(1.0) - df['avg_volume'] = df['avg_volume'].fillna(0) - - # 处理Infinity和-Infinity值 - df['volume_ratio'] = df['volume_ratio'].replace([float('inf'), float('-inf')], 1.0) - - # 计算ATR (Average True Range) - 14周期 - df['atr'] = ta.ATR(df, timeperiod=14) - df['atr'] = df['atr'].fillna(0) - bb2633 = ta.BBANDS(df, timeperiod=26, nbdevup=3.0, nbdevdn=3.0, matype=0) - bbp2633 = (df['close'] - bb2633['lowerband']) / (bb2633['upperband'] - bb2633['lowerband']) - df['bb2633upper'] = bb2633['upperband'].fillna(0) - df['bb2633lower'] = bb2633['lowerband'].fillna(0) - df['bbp2633'] = bbp2633.fillna(0) - df['bb2633middle'] = bb2633['middleband'].fillna(0) - return df - -def calculate_macd(df): - """计算MACD指标""" - global macd_fast_period, macd_slow_period, macd_signal_period - exp1 = df['close'].ewm(span=macd_fast_period, adjust=False).mean() - exp2 = df['close'].ewm(span=macd_slow_period, adjust=False).mean() - macd = exp1 - exp2 - signal = macd.ewm(span=macd_signal_period, adjust=False).mean() - histogram = macd - signal - - return { - 'macd': macd.tolist(), - 'signal': signal.tolist(), - 'histogram': histogram.tolist() - } - -def analyze_chan(df, symbol=None, timeframe=None): - """进行缠论分析""" - chan = TF_DF() - - # 初始化多时间周期数据以获取EMA52 - ema52_dict = None - # 获取分析结果 - klu_list = chan.get_kl_data(df) - klc_list = chan.get_klc_list(klu_list) - bi_list = chan.cal_bi_list(klc_list) - #for index in range(0, 10): - #print(bi_list[index].start_time, bi_list[index].start_klc.end_time, bi_list[index].dir) - seg_list = chan.get_seg_list(bi_list) - zs_list = chan.calculate_seg_zs(seg_list) - # 计算笔中枢(BI中枢)并拍平成列表 - - #bi_zs_list = chan.cal_bi_zs_list_pure(bi_list) - bi_zs_list = chan.cal_bi_zs(seg_list) - bsp_list = [] - if len(bi_zs_list) > 0: - bsp_list = chan.find_all_bsp(bi_list, bi_zs_list) - #bsp_state_list = chan.get_bsp_state(df) - #for bsp in bsp_list: - #print(bsp.end_time, bsp.type, bsp.dir) - # 添加买卖点识别 - for bi in bi_list: - bi.cal_macdhist() - for bi in bi_list: - bi.cal_macd_div() - #print(bi.start_time, bi.macd_hist, bi.macd_div) - - # 添加ChanMACD分析(复用 get_klc_list 内已算好的结果,避免同周期二次全量分析) - chan_macd = None - chan_macd_data = {} - try: - - if klu_list and len(klu_list) > 0: - print(f"获取到KLU列表,长度: {len(klu_list)}") - chan_macd = getattr(chan, '_last_chan_macd', None) - if chan_macd is None: - chan_macd = ChanMACD(klu_list) - chan_macd_data = { - 'seg_list': chan_macd.seg_list, - 'unittf_list': chan_macd.unittf_list, - 'histset_list': chan_macd.histset_list, - 'klu_list': chan_macd.klu_list, - 'high_position_list': chan_macd.high_position_list, - 'high_empty_list': chan_macd.high_empty_list, - 'low_position_list': getattr(chan_macd, 'low_position_list', []), - 'low_empty_list': getattr(chan_macd, 'low_empty_list', []), - 'return_zero_list': chan_macd.return_zero_list, - 'cross0_up_list': chan_macd.cross0_up_list, - 'cross0_down_list': chan_macd.cross0_down_list - } - print(f"ChanMACD分析完成: seg={len(chan_macd.seg_list)}, unittf={len(chan_macd.unittf_list)}, histset={len(chan_macd.histset_list)}") - else: - print("未能获取KLU列表或列表为空") - chan_macd_data = { - 'seg_list': [], - 'unittf_list': [], - 'histset_list': [], - 'high_position_list': [], - 'high_empty_list': [], - 'return_zero_list': [], - 'cross0_up_list': [], - 'cross0_down_list': [] - } - except Exception as e: - print(f"ChanMACD分析出错: {e}") - import traceback - traceback.print_exc() - chan_macd_data = { - 'seg_list': [], - 'unittf_list': [], - 'histset_list': [], - 'high_position_list': [], - 'high_empty_list': [], - 'low_position_list': [], - 'low_empty_list': [], - 'return_zero_list': [], - 'cross0_up_list': [], - 'cross0_down_list': [] - } - - # 提取K线分型信息 - klc_fx_info = [] - for klc in klc_list: - if hasattr(klc, 'klc_fx_type') and klc.klc_fx_type != Chan_KLC_FX.UNKNOWN: - try: - # 计算分型强度 - fx_strength = 0 - fx_strength_level = "" - is_strong_fx = False - - # 统一使用cal_fx_strength函数 - if hasattr(klc, 'cal_fx_strength'): - fx_strength = klc.cal_fx_strength(5) - - # 尝试获取分型强度等级 - if hasattr(klc, 'get_fx_strength_level'): - fx_strength_level = klc.get_fx_strength_level() - - # 尝试判断是否为强分型 - if hasattr(klc, 'is_strong_fx'): - is_strong_fx = klc.is_strong_fx() - - # 如果分型强度小于1,设为0 - if fx_strength < 1: - fx_strength = 0 - - # KLC 分型框(起止时间+高低价): - # 仅使用 cal_fx_box 通过 display 条件后生成的 klc.fx_box。 - # 若无 fx_box,则前端不应绘制分型框。 - fx_box = getattr(klc, 'fx_box', None) - box_start_time = getattr(fx_box, 'start_time', None) if fx_box else None - box_end_time = getattr(fx_box, 'end_time', None) if fx_box else None - box_high = getattr(fx_box, 'high', None) if fx_box else None - box_low = getattr(fx_box, 'low', None) if fx_box else None - - if klc.bb_out: - klc_fx_info.append({ - 'time': klc.end_time, - 'price': klc.low if klc.fx == Chan_FX_TYPE.BOTTOM else klc.high, - 'fx_type': str(klc.klc_fx_type).replace("Chan_KLC_FX.", ""), - 'is_bottom': klc.fx == Chan_FX_TYPE.BOTTOM, - 'fx_strength': fx_strength, # 分型强度分数 (0-100) - 'fx_strength_level': fx_strength_level, # 分型强度等级 (极强/强/中等/弱/极弱) - 'is_strong_fx': is_strong_fx, # 是否为强分型 - - # 虚线分型框信息(给前端画框用) - 'start_time': box_start_time, - 'end_time': box_end_time, - 'high': float(box_high) if box_high is not None else None, - 'low': float(box_low) if box_low is not None else None, - }) - except Exception as e: - # 如果出错,仍然添加基本信息,但分型强度为0 - fx_box = getattr(klc, 'fx_box', None) - box_start_time = getattr(fx_box, 'start_time', None) if fx_box else None - box_end_time = getattr(fx_box, 'end_time', None) if fx_box else None - box_high = getattr(fx_box, 'high', None) if fx_box else None - box_low = getattr(fx_box, 'low', None) if fx_box else None - - klc_fx_info.append({ - 'time': klc.end_time, - 'price': klc.low if klc.fx == Chan_FX_TYPE.BOTTOM else klc.high, - 'fx_type': str(klc.klc_fx_type).replace("Chan_KLC_FX.", ""), - 'is_bottom': klc.fx == Chan_FX_TYPE.BOTTOM, - 'fx_strength': 0, - 'fx_strength_level': "", - 'is_strong_fx': False, - - # 虚线分型框信息(给前端画框用) - 'start_time': box_start_time, - 'end_time': box_end_time, - 'high': float(box_high) if box_high is not None else None, - 'low': float(box_low) if box_low is not None else None, - }) - - - return { - 'klc_list': klc_list, - 'klu_list': klu_list, # 添加KLU列表 - 'bi_list': bi_list, - 'seg_list': seg_list, - 'zs_list': zs_list, - 'bi_zs_list': bi_zs_list, # 添加BI中枢列表 - 'bsp_list': bsp_list, # 添加买卖点列表 - 'klc_fx_info': klc_fx_info, # KLC分型信息 - 'chan_macd': chan_macd_data, # 添加ChanMACD分析数据 - 'ema52_dict': ema52_dict # 添加多时间周期EMA52数据 - } - -# 辅助函数,转换缠论方向枚举为整数 -def convert_direction(direction): - """转换方向枚举为数字""" - if direction == Chan_BI_DIR.UP or direction == Chan_SEG_DIR.UP: - return 1 - elif direction == Chan_BI_DIR.DOWN or direction == Chan_SEG_DIR.DOWN: - return -1 - else: - return 0 - -def format_time_safely(time_obj, client_tz): - """安全地格式化时间对象,处理字符串和datetime两种情况""" - if time_obj is None: - return None - - if isinstance(time_obj, str): - # 尝试将字符串解析为datetime - try: - from dateutil import parser - time_obj = parser.parse(time_obj) - return time_obj.astimezone(client_tz).isoformat() - except: - return time_obj - else: - # 已经是datetime对象 - return time_obj.astimezone(client_tz).isoformat() - -def serialize_chan_macd_data(chan_macd_data, client_tz): - """序列化ChanMACD数据为JSON可序列化格式""" - serialized_data = { - 'seg_list': [], - 'unittf_list': [], - 'histset_list': [], - # 状态标记数据 - 'high_position_list': [], - 'high_empty_list': [], - 'low_position_list': [], - 'low_empty_list': [], - 'return_zero_list': [], - 'cross0_up_list': [], - 'cross0_down_list': [], - # 新增:输出KLU的继续背驰/分离背驰标志 - 'klu_list': [] - } - - # 序列化seg_list - for seg in chan_macd_data.get('seg_list', []): - try: - seg_data = { - 'start_time': format_time_safely(seg.start_time, client_tz), - 'end_time': format_time_safely(seg.end_time, client_tz) if seg.end_time else None, - 'seg_dir': 'ABOVE' if seg.seg_dir == Chan_MACDSEG_DIR.ABOVE else 'UNDER', - 'klu_count': len(seg.klu_list) if hasattr(seg, 'klu_list') else 0, - 'unittf_count': len(seg.unittf_list) if hasattr(seg, 'unittf_list') else 0, - 'histset_count': len(seg.hist_set) if hasattr(seg, 'hist_set') else 0 - } - serialized_data['seg_list'].append(seg_data) - except Exception as e: - print(f"序列化seg出错: {e}") - continue - - # 序列化unittf_list(兼容新结构与枚举类型) - for unittf in chan_macd_data.get('unittf_list', []): - try: - dir_value = getattr(unittf, 'uinttf_dir', None) - dir_name = getattr(dir_value, 'name', dir_value if isinstance(dir_value, str) else None) - start_t = getattr(unittf, 'start_type', None) - start_type = getattr(start_t, 'name', start_t) - end_t = getattr(unittf, 'end_type', None) - end_type = getattr(end_t, 'name', end_t) - peak_abs = getattr(unittf, 'peak_abs', None) - if peak_abs is None: - peak_abs = getattr(unittf, 'peak_hist', None) - length = getattr(unittf, 'length', None) - if length is None: - length = len(unittf.klu_list) if hasattr(unittf, 'klu_list') else None - - unittf_data = { - 'start_time': format_time_safely(getattr(unittf, 'start_time', None), client_tz), - 'end_time': format_time_safely(getattr(unittf, 'end_time', None), client_tz) if getattr(unittf, 'end_time', None) else None, - 'dir': dir_name, # 'ABOVE' | 'UNDER' | None - 'start_type': start_type, # e.g. 'START' | 'CROSS0' | 'NEAR0_UP' | 'NEAR0_DOWN' - 'end_type': end_type, - 'invalid': getattr(unittf, 'invalid', False), - 'peak_abs': peak_abs, - 'length': length, - 'klu_count': len(unittf.klu_list) if hasattr(unittf, 'klu_list') else 0, - 'histset_count': len(unittf.histset_list) if hasattr(unittf, 'histset_list') else 0 - } - serialized_data['unittf_list'].append(unittf_data) - except Exception as e: - print(f"序列化unittf出错: {e}") - continue - - # 序列化histset_list - for histset in chan_macd_data.get('histset_list', []): - try: - histset_data = { - 'start_time': format_time_safely(getattr(histset, 'start_time', None), client_tz), - 'end_time': format_time_safely(getattr(histset, 'end_time', None), client_tz), - 'histset_dir': 'ABOVE' if histset.histset_dir == Chan_MACDHISTSET_DIR.ABOVE else 'UNDER', - 'klu_count': len(histset.klu_list) if hasattr(histset, 'klu_list') else 0 - } - serialized_data['histset_list'].append(histset_data) - except Exception as e: - print(f"序列化histset出错: {e}") - continue - - # 序列化状态标记数据 - # 序列化高位列表 - for high_pos in chan_macd_data.get('high_position_list', []): - try: - high_pos_data = { - 'time': format_time_safely(high_pos['time'], client_tz), - 'end_time': format_time_safely(high_pos.get('end_time'), client_tz) if high_pos.get('end_time') else None, - 'type': high_pos.get('type', 'start'), - 'macd': high_pos.get('macd'), - 'signal': high_pos.get('signal'), - 'macdhist': high_pos.get('macdhist'), - 'end_macd': high_pos.get('end_macd'), - 'end_signal': high_pos.get('end_signal'), - 'end_macdhist': high_pos.get('end_macdhist') - } - serialized_data['high_position_list'].append(high_pos_data) - except Exception as e: - print(f"序列化high_position出错: {e}") - continue - - # 序列化高位空列表 - for high_empty in chan_macd_data.get('high_empty_list', []): - try: - high_empty_data = { - 'time': format_time_safely(high_empty['time'], client_tz), - 'end_time': format_time_safely(high_empty.get('end_time'), client_tz) if high_empty.get('end_time') else None, - 'type': high_empty.get('type', 'start'), - 'macd': high_empty.get('macd'), - 'signal': high_empty.get('signal'), - 'macdhist': high_empty.get('macdhist'), - 'end_macd': high_empty.get('end_macd'), - 'end_signal': high_empty.get('end_signal'), - 'end_macdhist': high_empty.get('end_macdhist') - } - serialized_data['high_empty_list'].append(high_empty_data) - except Exception as e: - print(f"序列化high_empty出错: {e}") - continue - - # 序列化低位与低位空 - for low_pos in chan_macd_data.get('low_position_list', []): - try: - low_pos_data = { - 'time': format_time_safely(low_pos['time'], client_tz), - 'end_time': format_time_safely(low_pos.get('end_time'), client_tz) if low_pos.get('end_time') else None, - 'type': low_pos.get('type', 'start'), - 'macd': low_pos.get('macd'), - 'signal': low_pos.get('signal'), - 'macdhist': low_pos.get('macdhist'), - 'end_macd': low_pos.get('end_macd'), - 'end_signal': low_pos.get('end_signal'), - 'end_macdhist': low_pos.get('end_macdhist') - } - serialized_data['low_position_list'].append(low_pos_data) - except Exception as e: - print(f"序列化low_position出错: {e}") - continue - - for low_empty in chan_macd_data.get('low_empty_list', []): - try: - low_empty_data = { - 'time': format_time_safely(low_empty['time'], client_tz), - 'end_time': format_time_safely(low_empty.get('end_time'), client_tz) if low_empty.get('end_time') else None, - 'type': low_empty.get('type', 'start'), - 'macd': low_empty.get('macd'), - 'signal': low_empty.get('signal'), - 'macdhist': low_empty.get('macdhist'), - 'end_macd': low_empty.get('end_macd'), - 'end_signal': low_empty.get('end_signal'), - 'end_macdhist': low_empty.get('end_macdhist') - } - serialized_data['low_empty_list'].append(low_empty_data) - except Exception as e: - print(f"序列化low_empty出错: {e}") - continue - - # 序列化归零轴列表 - for return_zero in chan_macd_data.get('return_zero_list', []): - try: - return_zero_data = { - 'time': format_time_safely(return_zero['time'], client_tz), - 'end_time': format_time_safely(return_zero.get('end_time'), client_tz) if return_zero.get('end_time') else None, - 'type': return_zero.get('type', 'start'), - 'macd': return_zero.get('macd'), - 'signal': return_zero.get('signal'), - 'macdhist': return_zero.get('macdhist'), - 'end_macd': return_zero.get('end_macd'), - 'end_signal': return_zero.get('end_signal'), - 'end_macdhist': return_zero.get('end_macdhist') - } - serialized_data['return_zero_list'].append(return_zero_data) - except Exception as e: - print(f"序列化return_zero出错: {e}") - continue - - # 序列化穿越零轴列表 - for cross0_up in chan_macd_data.get('cross0_up_list', []): - try: - cross0_up_data = { - 'time': format_time_safely(cross0_up['time'], client_tz), - 'type': cross0_up.get('type', 'start'), - 'macd': cross0_up.get('macd'), - 'signal': cross0_up.get('signal'), - 'macdhist': cross0_up.get('macdhist') - } - serialized_data['cross0_up_list'].append(cross0_up_data) - except Exception as e: - print(f"序列化cross0_up出错: {e}") - continue - - for cross0_down in chan_macd_data.get('cross0_down_list', []): - try: - cross0_down_data = { - 'time': format_time_safely(cross0_down['time'], client_tz), - 'type': cross0_down.get('type', 'start'), - 'macd': cross0_down.get('macd'), - 'signal': cross0_down.get('signal'), - 'macdhist': cross0_down.get('macdhist') - } - serialized_data['cross0_down_list'].append(cross0_down_data) - except Exception as e: - print(f"序列化cross0_down出错: {e}") - continue - - # 序列化 KLU 列表(仅导出需要的时间与背驰标志) - for klu in chan_macd_data.get('klu_list', []): - try: - serialized_data['klu_list'].append({ - 'time': format_time_safely(getattr(klu, 'time', None), client_tz), - 'continue_div': bool(getattr(klu, 'continue_div', False)), - 'separate_div': int(getattr(klu, 'separate_div', 0)) if getattr(klu, 'separate_div', 0) is not None else 0, - 'near0_return': int(getattr(klu, 'near0_return', 0)) if getattr(klu, 'near0_return', 0) is not None else 0 - }) - except Exception as e: - print(f"序列化klu出错: {e}") - continue - - return serialized_data - -def is_smaller_timeframe(tf1, tf2): - """判断时间周期tf1是否小于tf2""" - tf1_value = timeframe_to_minutes(tf1) - tf2_value = timeframe_to_minutes(tf2) - if tf1_value is None or tf2_value is None: - return False - return tf1_value < tf2_value - -def is_smaller_or_equal_timeframe(tf1, tf2): - """判断时间周期tf1是否小于等于tf2""" - tf1_value = timeframe_to_minutes(tf1) - tf2_value = timeframe_to_minutes(tf2) - if tf1_value is None or tf2_value is None: - return False - return tf1_value <= tf2_value - -def clean_dataframe_for_json(df): - """清理DataFrame数据用于JSON序列化""" - # 创建副本避免修改原始数据 - clean_df = df.copy() - - # 替换NaN值为None - clean_df = clean_df.where(pd.notnull(clean_df), None) - - return clean_df - -# ====== 趋势判定与趋势筛选(币对) ====== - -def classify_trend_stage(df): - """根据 EMA 斜率与多空排列判断趋势方向与阶段 - 返回: direction in {"bull","bear","sideways"}, stage in {"early","mid","late"}, strength_score (0-100) - """ - if df is None or len(df) < 60: - return "sideways", "early", 0 - - # 使用 EMA5/10/24/52 - closes = df['close'].values - ema5 = df['ema5'].values if 'ema5' in df else ta.EMA(df, timeperiod=5) - ema10 = df['ema10'].values if 'ema10' in df else ta.EMA(df, timeperiod=10) - ema24 = df['ema24'].values if 'ema24' in df else ta.EMA(df, timeperiod=24) - ema52 = df['ema52'].values if 'ema52' in df else ta.EMA(df, timeperiod=52) - - # 最近N根用于斜率与排列判定 - lookback = min(30, len(df) - 1) - if lookback <= 5: - return "sideways", "early", 0 - - # 简单斜率: 最近k根的线性变化率近似 - def slope(arr, k=10): - k = min(k, len(arr) - 1) - if k < 2: - return 0.0 - y = arr[-k:] - x = np.arange(k) - # 最小二乘拟合斜率 - denom = np.dot(x - x.mean(), x - x.mean()) - if denom == 0: - return 0.0 - m = np.dot(y - y.mean(), x - x.mean()) / denom - return float(m) - - k_slope = 12 # 斜率窗口 - s5 = slope(ema5, k_slope) - s10 = slope(ema10, k_slope) - s24 = slope(ema24, k_slope) - s52 = slope(ema52, k_slope) - - # 多空排列 - last5, last10, last24, last52 = ema5[-1], ema10[-1], ema24[-1], ema52[-1] - bull_stack = last5 > last10 > last24 > last52 - bear_stack = last5 < last10 < last24 < last52 - - # 波动性与动量增强: MACD 柱体最近均值 - macdhist = df['macdhist'].values if 'macdhist' in df else calculate_macd(df)['histogram'] - hist_recent = macdhist[-lookback:] - hist_power = float(np.mean(np.abs(hist_recent))) if len(hist_recent) else 0.0 - - # 方向 - if bull_stack and s24 > 0 and s52 > 0: - direction = "bull" - elif bear_stack and s24 < 0 and s52 < 0: - direction = "bear" - else: - # 用价格相对 EMA52 辅助 - if closes[-1] > last52 and (s24 + s52) > 0: - direction = "bull" - elif closes[-1] < last52 and (s24 + s52) < 0: - direction = "bear" - else: - direction = "sideways" - - # 阶段: 依据(斜率大小、与EMA52距离、MACD柱体扩张/收敛) - dist52 = float((closes[-1] - last52) / last52) if last52 else 0.0 - slope_score = max(0.0, (abs(s24) + abs(s52)) * 1000.0) # 归一化 - dist_score = min(50.0, abs(dist52) * 200.0) - hist_score = min(30.0, hist_power * 10.0) - strength = float(min(100.0, slope_score + dist_score + hist_score)) - - # 简单阶段判定 - if direction == "sideways": - stage = "early" - strength = min(strength, 30.0) - else: - # 查看最近 hist 是否在扩大或收敛 - if len(hist_recent) >= 6: - recent_growth = np.mean(np.abs(hist_recent[-3:])) - np.mean(np.abs(hist_recent[-6:-3])) - else: - recent_growth = 0.0 - - if recent_growth > 0 and abs(dist52) < 0.05: - stage = "early" - elif recent_growth > 0 and abs(dist52) >= 0.05: - stage = "mid" - else: - stage = "late" - - return direction, stage, strength - - -def load_crypto_symbols(limit=200): - """加载常见USDT永续合约交易对,返回列表""" - refresh_data_service_metadata() - if SYMBOLS: - return SYMBOLS[:limit] - try: - markets = exchange.load_markets() - symbols = [s for s in markets.keys() if '/USDT' in s and ':USDT' in s] - return symbols[:limit] - except Exception: - return DEFAULT_SYMBOLS[:limit] - - - -def get_uncompleted_seg_list(seg_list, client_tz): - """获取未完成线段列表,正确处理倒数第二个和最后一个未完成线段""" - uncompleted_segs = [seg for seg in seg_list if not seg.is_sure] - - if len(uncompleted_segs) == 0: - return [] - - result = [] - - for i, seg in enumerate(uncompleted_segs): - is_last = (i == len(uncompleted_segs) - 1) # 是否为最后一个未完成线段 - - seg_data = { - 'start_time': seg.start_bi.start_klc.end_time if isinstance(seg.start_bi.start_klc.end_time, str) else seg.start_bi.start_klc.end_time.astimezone(client_tz).isoformat(), - 'sure_time': format_time_safely(seg.sure_time, client_tz) if seg.sure_time else None, - 'start_price': seg.start_bi.start_klc.low if convert_direction(seg.dir) == 1 else seg.start_bi.start_klc.high, - 'direction': convert_direction(seg.dir) - } - - if is_last: - # 最后一个未完成线段:没有结束时间和价格 - seg_data['end_time'] = None - seg_data['end_price'] = None - else: - # 倒数第二个及之前的未完成线段:使用实际的结束时间和价格 - if seg.end_bi and seg.end_bi.end_klc: - seg_data['end_time'] = seg.end_bi.end_klc.end_time if isinstance(seg.end_bi.end_klc.end_time, str) else seg.end_bi.end_klc.end_time.astimezone(client_tz).isoformat() - seg_data['end_price'] = seg.end_bi.end_klc.high if convert_direction(seg.dir) == 1 else seg.end_bi.end_klc.low - else: - # 如果没有结束笔,设为None - seg_data['end_time'] = None - seg_data['end_price'] = None - - result.append(seg_data) - - return result diff --git a/web/services/runtime/__init__.py b/web/services/runtime/__init__.py new file mode 100644 index 0000000..b05c056 --- /dev/null +++ b/web/services/runtime/__init__.py @@ -0,0 +1,95 @@ +"""runtime 门面:保持 `from services.runtime import *` 与 `import services.runtime as R` 兼容。""" +from __future__ import annotations + +# ---- 历史兼容:旧 monolith 上 `from pytz import timezone` 等会随 import * 漏出 ---- +import json # noqa: F401 +import logging +import sys as _sys +import time # noqa: F401 +from collections import OrderedDict # noqa: F401 +from concurrent.futures import ThreadPoolExecutor, as_completed # noqa: F401 + +import numpy as np # noqa: F401 +from pytz import timezone # noqa: F401 + +from chanlun.analysis.ChanZone import ( # noqa: F401 + StructureZoneConfig, + analyze_structure_zones_from_serialized, +) + +logger = logging.getLogger("services.runtime") + +from .state import ( # noqa: F401 + TRADE_POINT_TYPE, + macd_fast_period, + macd_slow_period, + macd_signal_period, + exchange, + china_stock, + _zone_cache, + DEFAULT_TIMEFRAME_LABELS, + DEFAULT_SYMBOLS, + TIMEFRAMES, + SYMBOLS, + DATA_SERVICE_AVAILABLE, + SERVICE_METADATA_LAST_REFRESH, +) +from .timeframes import ( # noqa: F401 + _zone_cache_ttl, + timeframe_to_minutes, + format_timeframe_label, + build_timeframe_labels, + compute_timeframe_defaults, + is_smaller_timeframe, + is_smaller_or_equal_timeframe, +) +from .market_data import ( # noqa: F401 + _parse_time_input, + refresh_data_service_metadata, + _fetch_kl_from_datasvc, + A_STOCK_SYMBOLS, + detect_symbol_type, + get_kl_data, + _get_crypto_kl_data_via_ccxt, + get_crypto_kl_data, + get_a_stock_kl_data, + load_crypto_symbols, +) +from .indicators import ( # noqa: F401 + add_indicators, + calculate_macd, +) +from .analyze import ( # noqa: F401 + analyze_chan, + classify_trend_stage, +) +from .serialize import ( # noqa: F401 + convert_direction, + format_time_safely, + serialize_chan_macd_data, + clean_dataframe_for_json, + get_uncompleted_seg_list, +) + +# 预取元信息(与拆分前模块加载行为一致) +refresh_data_service_metadata(force=True) + +# 标量在 import 时会拷贝;刷新后写回本模块,供 `from services.runtime import *` 读到最新值 +from . import state as _state + +_mod = _sys.modules[__name__] +_mod.DATA_SERVICE_AVAILABLE = _state.DATA_SERVICE_AVAILABLE +_mod.SERVICE_METADATA_LAST_REFRESH = _state.SERVICE_METADATA_LAST_REFRESH +_mod.macd_fast_period = _state.macd_fast_period +_mod.macd_slow_period = _state.macd_slow_period +_mod.macd_signal_period = _state.macd_signal_period + + +def __getattr__(name: str): + if hasattr(_state, name): + return getattr(_state, name) + raise AttributeError(name) + + +def __dir__(): + return sorted(set(globals()) | set(dir(_state))) diff --git a/web/services/runtime/analyze.py b/web/services/runtime/analyze.py new file mode 100644 index 0000000..d0ce88b --- /dev/null +++ b/web/services/runtime/analyze.py @@ -0,0 +1,274 @@ +from __future__ import annotations + +import numpy as np +import talib.abstract as ta + +from chanlun import TF_DF +from chanlun.core.ChanEnum import Chan_KLC_FX, Chan_FX_TYPE +from chanlun.indicators.ChanMACD import ChanMACD + +from .indicators import calculate_macd + +def analyze_chan(df, symbol=None, timeframe=None): + """进行缠论分析""" + chan = TF_DF() + + # 初始化多时间周期数据以获取EMA52 + ema52_dict = None + # 获取分析结果 + klu_list = chan.get_kl_data(df) + klc_list = chan.get_klc_list(klu_list) + bi_list = chan.cal_bi_list(klc_list) + #for index in range(0, 10): + #print(bi_list[index].start_time, bi_list[index].start_klc.end_time, bi_list[index].dir) + seg_list = chan.get_seg_list(bi_list) + zs_list = chan.calculate_seg_zs(seg_list) + # 计算笔中枢(BI中枢)并拍平成列表 + + #bi_zs_list = chan.cal_bi_zs_list_pure(bi_list) + bi_zs_list = chan.cal_bi_zs(seg_list) + bsp_list = [] + if len(bi_zs_list) > 0: + bsp_list = chan.find_all_bsp(bi_list, bi_zs_list) + #bsp_state_list = chan.get_bsp_state(df) + #for bsp in bsp_list: + #print(bsp.end_time, bsp.type, bsp.dir) + # 添加买卖点识别 + for bi in bi_list: + bi.cal_macdhist() + for bi in bi_list: + bi.cal_macd_div() + #print(bi.start_time, bi.macd_hist, bi.macd_div) + + # 添加ChanMACD分析(复用 get_klc_list 内已算好的结果,避免同周期二次全量分析) + chan_macd = None + chan_macd_data = {} + try: + + if klu_list and len(klu_list) > 0: + print(f"获取到KLU列表,长度: {len(klu_list)}") + chan_macd = getattr(chan, '_last_chan_macd', None) + if chan_macd is None: + chan_macd = ChanMACD(klu_list) + chan_macd_data = { + 'seg_list': chan_macd.seg_list, + 'unittf_list': chan_macd.unittf_list, + 'histset_list': chan_macd.histset_list, + 'klu_list': chan_macd.klu_list, + 'high_position_list': chan_macd.high_position_list, + 'high_empty_list': chan_macd.high_empty_list, + 'low_position_list': getattr(chan_macd, 'low_position_list', []), + 'low_empty_list': getattr(chan_macd, 'low_empty_list', []), + 'return_zero_list': chan_macd.return_zero_list, + 'cross0_up_list': chan_macd.cross0_up_list, + 'cross0_down_list': chan_macd.cross0_down_list + } + print(f"ChanMACD分析完成: seg={len(chan_macd.seg_list)}, unittf={len(chan_macd.unittf_list)}, histset={len(chan_macd.histset_list)}") + else: + print("未能获取KLU列表或列表为空") + chan_macd_data = { + 'seg_list': [], + 'unittf_list': [], + 'histset_list': [], + 'high_position_list': [], + 'high_empty_list': [], + 'return_zero_list': [], + 'cross0_up_list': [], + 'cross0_down_list': [] + } + except Exception as e: + print(f"ChanMACD分析出错: {e}") + import traceback + traceback.print_exc() + chan_macd_data = { + 'seg_list': [], + 'unittf_list': [], + 'histset_list': [], + 'high_position_list': [], + 'high_empty_list': [], + 'low_position_list': [], + 'low_empty_list': [], + 'return_zero_list': [], + 'cross0_up_list': [], + 'cross0_down_list': [] + } + + # 提取K线分型信息 + klc_fx_info = [] + for klc in klc_list: + if hasattr(klc, 'klc_fx_type') and klc.klc_fx_type != Chan_KLC_FX.UNKNOWN: + try: + # 计算分型强度 + fx_strength = 0 + fx_strength_level = "" + is_strong_fx = False + + # 统一使用cal_fx_strength函数 + if hasattr(klc, 'cal_fx_strength'): + fx_strength = klc.cal_fx_strength(5) + + # 尝试获取分型强度等级 + if hasattr(klc, 'get_fx_strength_level'): + fx_strength_level = klc.get_fx_strength_level() + + # 尝试判断是否为强分型 + if hasattr(klc, 'is_strong_fx'): + is_strong_fx = klc.is_strong_fx() + + # 如果分型强度小于1,设为0 + if fx_strength < 1: + fx_strength = 0 + + # KLC 分型框(起止时间+高低价): + # 仅使用 cal_fx_box 通过 display 条件后生成的 klc.fx_box。 + # 若无 fx_box,则前端不应绘制分型框。 + fx_box = getattr(klc, 'fx_box', None) + box_start_time = getattr(fx_box, 'start_time', None) if fx_box else None + box_end_time = getattr(fx_box, 'end_time', None) if fx_box else None + box_high = getattr(fx_box, 'high', None) if fx_box else None + box_low = getattr(fx_box, 'low', None) if fx_box else None + + if klc.bb_out: + klc_fx_info.append({ + 'time': klc.end_time, + 'price': klc.low if klc.fx == Chan_FX_TYPE.BOTTOM else klc.high, + 'fx_type': str(klc.klc_fx_type).replace("Chan_KLC_FX.", ""), + 'is_bottom': klc.fx == Chan_FX_TYPE.BOTTOM, + 'fx_strength': fx_strength, # 分型强度分数 (0-100) + 'fx_strength_level': fx_strength_level, # 分型强度等级 (极强/强/中等/弱/极弱) + 'is_strong_fx': is_strong_fx, # 是否为强分型 + + # 虚线分型框信息(给前端画框用) + 'start_time': box_start_time, + 'end_time': box_end_time, + 'high': float(box_high) if box_high is not None else None, + 'low': float(box_low) if box_low is not None else None, + }) + except Exception as e: + # 如果出错,仍然添加基本信息,但分型强度为0 + fx_box = getattr(klc, 'fx_box', None) + box_start_time = getattr(fx_box, 'start_time', None) if fx_box else None + box_end_time = getattr(fx_box, 'end_time', None) if fx_box else None + box_high = getattr(fx_box, 'high', None) if fx_box else None + box_low = getattr(fx_box, 'low', None) if fx_box else None + + klc_fx_info.append({ + 'time': klc.end_time, + 'price': klc.low if klc.fx == Chan_FX_TYPE.BOTTOM else klc.high, + 'fx_type': str(klc.klc_fx_type).replace("Chan_KLC_FX.", ""), + 'is_bottom': klc.fx == Chan_FX_TYPE.BOTTOM, + 'fx_strength': 0, + 'fx_strength_level': "", + 'is_strong_fx': False, + + # 虚线分型框信息(给前端画框用) + 'start_time': box_start_time, + 'end_time': box_end_time, + 'high': float(box_high) if box_high is not None else None, + 'low': float(box_low) if box_low is not None else None, + }) + + + return { + 'klc_list': klc_list, + 'klu_list': klu_list, # 添加KLU列表 + 'bi_list': bi_list, + 'seg_list': seg_list, + 'zs_list': zs_list, + 'bi_zs_list': bi_zs_list, # 添加BI中枢列表 + 'bsp_list': bsp_list, # 添加买卖点列表 + 'klc_fx_info': klc_fx_info, # KLC分型信息 + 'chan_macd': chan_macd_data, # 添加ChanMACD分析数据 + 'ema52_dict': ema52_dict # 添加多时间周期EMA52数据 + } + +def classify_trend_stage(df): + """根据 EMA 斜率与多空排列判断趋势方向与阶段 + 返回: direction in {"bull","bear","sideways"}, stage in {"early","mid","late"}, strength_score (0-100) + """ + if df is None or len(df) < 60: + return "sideways", "early", 0 + + # 使用 EMA5/10/24/52 + closes = df['close'].values + ema5 = df['ema5'].values if 'ema5' in df else ta.EMA(df, timeperiod=5) + ema10 = df['ema10'].values if 'ema10' in df else ta.EMA(df, timeperiod=10) + ema24 = df['ema24'].values if 'ema24' in df else ta.EMA(df, timeperiod=24) + ema52 = df['ema52'].values if 'ema52' in df else ta.EMA(df, timeperiod=52) + + # 最近N根用于斜率与排列判定 + lookback = min(30, len(df) - 1) + if lookback <= 5: + return "sideways", "early", 0 + + # 简单斜率: 最近k根的线性变化率近似 + def slope(arr, k=10): + k = min(k, len(arr) - 1) + if k < 2: + return 0.0 + y = arr[-k:] + x = np.arange(k) + # 最小二乘拟合斜率 + denom = np.dot(x - x.mean(), x - x.mean()) + if denom == 0: + return 0.0 + m = np.dot(y - y.mean(), x - x.mean()) / denom + return float(m) + + k_slope = 12 # 斜率窗口 + s5 = slope(ema5, k_slope) + s10 = slope(ema10, k_slope) + s24 = slope(ema24, k_slope) + s52 = slope(ema52, k_slope) + + # 多空排列 + last5, last10, last24, last52 = ema5[-1], ema10[-1], ema24[-1], ema52[-1] + bull_stack = last5 > last10 > last24 > last52 + bear_stack = last5 < last10 < last24 < last52 + + # 波动性与动量增强: MACD 柱体最近均值 + macdhist = df['macdhist'].values if 'macdhist' in df else calculate_macd(df)['histogram'] + hist_recent = macdhist[-lookback:] + hist_power = float(np.mean(np.abs(hist_recent))) if len(hist_recent) else 0.0 + + # 方向 + if bull_stack and s24 > 0 and s52 > 0: + direction = "bull" + elif bear_stack and s24 < 0 and s52 < 0: + direction = "bear" + else: + # 用价格相对 EMA52 辅助 + if closes[-1] > last52 and (s24 + s52) > 0: + direction = "bull" + elif closes[-1] < last52 and (s24 + s52) < 0: + direction = "bear" + else: + direction = "sideways" + + # 阶段: 依据(斜率大小、与EMA52距离、MACD柱体扩张/收敛) + dist52 = float((closes[-1] - last52) / last52) if last52 else 0.0 + slope_score = max(0.0, (abs(s24) + abs(s52)) * 1000.0) # 归一化 + dist_score = min(50.0, abs(dist52) * 200.0) + hist_score = min(30.0, hist_power * 10.0) + strength = float(min(100.0, slope_score + dist_score + hist_score)) + + # 简单阶段判定 + if direction == "sideways": + stage = "early" + strength = min(strength, 30.0) + else: + # 查看最近 hist 是否在扩大或收敛 + if len(hist_recent) >= 6: + recent_growth = np.mean(np.abs(hist_recent[-3:])) - np.mean(np.abs(hist_recent[-6:-3])) + else: + recent_growth = 0.0 + + if recent_growth > 0 and abs(dist52) < 0.05: + stage = "early" + elif recent_growth > 0 and abs(dist52) >= 0.05: + stage = "mid" + else: + stage = "late" + + return direction, stage, strength + diff --git a/web/services/runtime/indicators.py b/web/services/runtime/indicators.py new file mode 100644 index 0000000..f465c4a --- /dev/null +++ b/web/services/runtime/indicators.py @@ -0,0 +1,103 @@ +from __future__ import annotations + +import talib.abstract as ta +from . import state + +def add_indicators(df): + macd = ta.MACD(df, fastperiod=state.macd_fast_period, slowperiod=state.macd_slow_period, signalperiod=state.macd_signal_period) + + df['macd'] = macd['macd'] + df['macdsignal'] = macd['macdsignal'] + df['macdhist'] = macd['macdhist'] + df['ma5'] = (ta.MA(df, timeperiod=5)).fillna(0) + df['ma10'] = (ta.MA(df, timeperiod=10)).fillna(0) + df['ma30'] = (ta.EMA(df, timeperiod=30)).fillna(0) + df['ma250'] = (ta.MA(df, timeperiod=250)).fillna(0) + # 新增 EMA 指标 + df['ema5'] = (ta.EMA(df, timeperiod=5)).fillna(0) + df['ema10'] = (ta.EMA(df, timeperiod=10)).fillna(0) + df['ema24'] = (ta.EMA(df, timeperiod=24)).fillna(0) + df['ema52'] = (ta.EMA(df, timeperiod=52)).fillna(0) + df['ema26'] = (ta.EMA(df, timeperiod=26)).fillna(0) + df['ema13'] = (ta.EMA(df, timeperiod=13)).fillna(0) + df['ema7'] = (ta.EMA(df, timeperiod=7)).fillna(0) + df['ema104'] = (ta.EMA(df, timeperiod=104)).fillna(0) + df['ema156'] = (ta.EMA(df, timeperiod=156)).fillna(0) + df['ema208'] = (ta.EMA(df, timeperiod=208)).fillna(0) + # 常用SMA 24/52 + try: + df['sma24'] = (ta.SMA(df, timeperiod=24)).fillna(0) + df['sma52'] = (ta.SMA(df, timeperiod=52)).fillna(0) + except Exception: + df['sma24'] = 0 + df['sma52'] = 0 + df['rsi'] = ta.RSI(df, timeperiod=14) + + # 计算布林带 (当前周期 - 20周期,2标准差) + bb = ta.BBANDS(df, timeperiod=365, nbdevup=3.0, nbdevdn=3.0, matype=0) + df['bb_upper'] = bb['upperband'].fillna(0) + df['bb_middle'] = bb['middleband'].fillna(0) + df['bb_lower'] = bb['lowerband'].fillna(0) + bb30 = ta.BBANDS(df, timeperiod=41, nbdevup=2.3, nbdevdn=2.3, matype=0) + #bb30 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0) + df['bbup30'] = bb30['upperband'].fillna(0) + df['bblow30'] = bb30['lowerband'].fillna(0) + bb302 = ta.BBANDS(df, timeperiod=41, nbdevup=2.0, nbdevdn=2.0, matype=0) + #bb302 = ta.BBANDS(df, timeperiod=20, nbdevup=2.0, nbdevdn=2.0, matype=0) + df['bbup302'] = bb302['upperband'].fillna(0) + df['bblow302'] = bb302['lowerband'].fillna(0) + # 计算次周期布林带 (14周期,2标准差) + bb_element = ta.BBANDS(df, timeperiod=14, nbdevup=2.0, nbdevdn=2.0, matype=0) + df['element_bb_upper'] = bb_element['upperband'].fillna(0) + df['element_bb_middle'] = bb_element['middleband'].fillna(0) + df['element_bb_lower'] = bb_element['lowerband'].fillna(0) + + df['macd'] = df['macd'].fillna(0) + df['macdsignal'] = df['macdsignal'].fillna(0) + df['macdhist'] = df['macdhist'].fillna(0) + df['ma5'] = df['ma5'].fillna(0) + df['ma10'] = df['ma10'].fillna(0) + df['ma30'] = df['ma30'].fillna(0) + df['ma250'] = df['ma250'].fillna(0) + df['ema5'] = df['ema5'].fillna(0) + df['ema10'] = df['ema10'].fillna(0) + df['ema24'] = df['ema24'].fillna(0) + df['ema52'] = df['ema52'].fillna(0) + df['sma24'] = df['sma24'].fillna(0) + df['sma52'] = df['sma52'].fillna(0) + df['rsi'] = df['rsi'].fillna(0) + df['avg_volume'] = df['volume'].rolling(10).mean() + # 计算量比,避免产生Infinity值 + df['volume_ratio'] = df['volume'] / df['avg_volume'] + # 填充缺失值(前N根K线) + df['volume_ratio'] = df['volume_ratio'].fillna(1.0) + df['avg_volume'] = df['avg_volume'].fillna(0) + + # 处理Infinity和-Infinity值 + df['volume_ratio'] = df['volume_ratio'].replace([float('inf'), float('-inf')], 1.0) + + # 计算ATR (Average True Range) - 14周期 + df['atr'] = ta.ATR(df, timeperiod=14) + df['atr'] = df['atr'].fillna(0) + bb2633 = ta.BBANDS(df, timeperiod=26, nbdevup=3.0, nbdevdn=3.0, matype=0) + bbp2633 = (df['close'] - bb2633['lowerband']) / (bb2633['upperband'] - bb2633['lowerband']) + df['bb2633upper'] = bb2633['upperband'].fillna(0) + df['bb2633lower'] = bb2633['lowerband'].fillna(0) + df['bbp2633'] = bbp2633.fillna(0) + df['bb2633middle'] = bb2633['middleband'].fillna(0) + return df + +def calculate_macd(df): + """计算MACD指标""" + exp1 = df['close'].ewm(span=state.macd_fast_period, adjust=False).mean() + exp2 = df['close'].ewm(span=state.macd_slow_period, adjust=False).mean() + macd = exp1 - exp2 + signal = macd.ewm(span=state.macd_signal_period, adjust=False).mean() + histogram = macd - signal + + return { + 'macd': macd.tolist(), + 'signal': signal.tolist(), + 'histogram': histogram.tolist() + } + diff --git a/web/services/runtime/market_data.py b/web/services/runtime/market_data.py new file mode 100644 index 0000000..53e5fc9 --- /dev/null +++ b/web/services/runtime/market_data.py @@ -0,0 +1,321 @@ +from __future__ import annotations + +import logging +import time +from datetime import datetime, timedelta + +import pandas as pd +import requests + +from config import DATA_SERVICE_URL +from . import state +from .state import DEFAULT_SYMBOLS, DEFAULT_TIMEFRAME_LABELS +from .timeframes import build_timeframe_labels + +logger = logging.getLogger(__name__) + +def _parse_time_input(value): + if value in (None, '', 0): + return None + try: + return int(float(value)) + except (ValueError, TypeError): + return None + + +def refresh_data_service_metadata(force=False): + """刷新数据服务提供的交易对与周期元信息。""" + now = time.time() + if not force and state.DATA_SERVICE_AVAILABLE and now - state.SERVICE_METADATA_LAST_REFRESH < 60: + return True + try: + resp = requests.get(f"{DATA_SERVICE_URL}/health", timeout=5) + resp.raise_for_status() + payload = resp.json() + service_symbols = payload.get("symbols") or payload.get("symbol_list") or [] + base_timeframes = payload.get("timeframes") or payload.get("base_timeframes") or [] + derived = payload.get("derived_timeframes") or [] + service_timeframes = list(base_timeframes) + for tf in derived: + if tf not in service_timeframes: + service_timeframes.append(tf) + if service_symbols: + state.SYMBOLS[:] = service_symbols + if service_timeframes: + state.TIMEFRAMES.clear() + state.TIMEFRAMES.update(build_timeframe_labels(service_timeframes)) + state.DATA_SERVICE_AVAILABLE = True + state.SERVICE_METADATA_LAST_REFRESH = now + return True + except Exception as exc: + logger.warning("无法加载数据服务元信息: %s", exc) + if not state.DATA_SERVICE_AVAILABLE: + state.TIMEFRAMES.clear() + state.TIMEFRAMES.update(DEFAULT_TIMEFRAME_LABELS) + state.SYMBOLS[:] = DEFAULT_SYMBOLS + state.DATA_SERVICE_AVAILABLE = False + return False + + +def _fetch_kl_from_datasvc(symbol, timeframe, start_ms=None, end_ms=None, limit=None): + params = {"symbol": symbol, "tf": timeframe} + if start_ms is not None: + params["start"] = int(start_ms) + if end_ms is not None: + params["end"] = int(end_ms) + if limit is not None: + params["limit"] = limit + resp = requests.get(f"{DATA_SERVICE_URL}/api/candles", params=params, timeout=10) + resp.raise_for_status() + data = resp.json() + if not data: + return None + df = pd.DataFrame(data) + if df.empty or "timestamp" not in df.columns: + return None + numeric_cols = ["open", "high", "low", "close", "volume"] + df["timestamp"] = pd.to_numeric(df["timestamp"], errors="coerce") + df = df.dropna(subset=["timestamp"]) + df["timestamp"] = df["timestamp"].astype("int64") + for col in numeric_cols: + if col in df.columns: + df[col] = pd.to_numeric(df[col], errors="coerce") + df = df.dropna(subset=numeric_cols) + df = df.sort_values("timestamp") + if limit and len(df) > limit: + df = df.tail(limit) + df = df.reset_index(drop=True) + df["date"] = pd.to_datetime(df["timestamp"], unit='ms', utc=True).dt.tz_convert('Asia/Shanghai') + return df + + +# 模块加载时尝试预取一次元信息,但失败不阻塞后续流程 +refresh_data_service_metadata(force=True) + +# A股热门股票 +# 模板中 A 股下拉仅放默认一项;用户切换到「A股」时由前端请求 /api/a_stocks 填充全市场(约 5500+) +A_STOCK_SYMBOLS = [{'symbol': '000001', 'name': '平安银行'}] + +def detect_symbol_type(symbol): + """检测交易对类型:crypto 或 a_stock""" + if '/' in symbol and 'USDT' in symbol: + return 'crypto' + elif len(symbol) == 6 and symbol.isdigit(): + return 'a_stock' + else: + return 'unknown' + +def get_kl_data(symbol, timeframe, limit=100000, start_time=None, end_time=None): + """获取K线数据,支持加密货币和A股""" + symbol_type = detect_symbol_type(symbol) + + if symbol_type == 'crypto': + return get_crypto_kl_data(symbol, timeframe, limit, start_time, end_time) + elif symbol_type == 'a_stock': + return get_a_stock_kl_data(symbol, timeframe, limit, start_time, end_time) + else: + return None + +def _get_crypto_kl_data_via_ccxt(symbol, timeframe, limit=100000, start_time=None, end_time=None): + """获取加密货币K线数据,支持分页加载确保获取指定时间范围内的所有数据""" + try: + # 初始化参数 + since = None + if start_time: + try: + since = int(start_time) + except ValueError: + pass + + # 结束时间处理 + until = None + if end_time: + try: + until = int(end_time) + except ValueError: + pass + + # 根据时间周期调整每次请求的数据量 + batch_size = 1000 # 默认批次大小 + if timeframe in ['1m', '3m', '5m']: + batch_size = 1000 # 分钟级数据减少批次大小 + elif timeframe in ['15m', '30m', '1h']: + batch_size = 1000 + else: + batch_size = 1500 # 日线及以上可以获取更多 + batch_size = 1500 # 默认批次大小 + # 初始化存储所有K线数据的列表 + all_ohlcv = [] + + # 初始化当前查询的开始时间 + current_since = since + + # 添加请求计数和最大限制 + request_count = 0 + max_requests = 300 # 最大请求次数,防止无限循环 + + # 分页加载数据 + while request_count < max_requests: + request_count += 1 + + try: + # 获取当前页的数据 + ohlcv = state.exchange.fetch_ohlcv(symbol, timeframe, since=current_since, limit=batch_size) + + # 如果没有获取到数据,结束循环 + if not ohlcv or len(ohlcv) == 0: + break + + # 将获取到的数据添加到总列表中 + all_ohlcv.extend(ohlcv) + + # 获取最后一条数据的时间戳 + last_timestamp = ohlcv[-1][0] + + # 如果已达到结束时间,结束循环 + if until and last_timestamp >= until: + break + + # 如果获取的数据条数小于限制数,说明已经获取完所有数据 + if len(ohlcv) < batch_size: + break + + # 更新下一页的开始时间(加1毫秒避免重复) + current_since = last_timestamp + 1 + + except Exception as e: + # 如果单个批次失败,继续尝试下一个批次 + if current_since: + # 尝试增加时间跳过可能的问题时间点 + current_since += 60000 # 跳过1分钟 + else: + break + + # 防止API请求过于频繁 + time.sleep(0.3) # 减少到0.3秒提高效率 + + # 数据为空的情况 + if not all_ohlcv or len(all_ohlcv) == 0: + return None + + # 转换为DataFrame + df = pd.DataFrame(all_ohlcv, columns=['timestamp', 'open', 'high', 'low', 'close', 'volume']) + df['date'] = pd.to_datetime(df['timestamp'], unit='ms').dt.tz_localize('UTC').dt.tz_convert('Asia/Shanghai') + + # 在客户端进行结束时间过滤 + if until: + df = df[df['timestamp'] <= until] + + # 去除重复数据 + df = df.drop_duplicates(subset=['timestamp']) + + # 按时间排序 + df = df.sort_values('timestamp') + + # 限制数据条数的逻辑 - 优先考虑时间范围 + if start_time and end_time: + # 如果指定了明确的时间范围,返回该时间范围内的所有数据 + if len(df) > 100000: # 防止数据量过大,设置一个合理的上限 + df = df.tail(100000).reset_index(drop=True) + elif limit and len(df) > limit: + # 如果没有指定明确时间范围,使用默认的limit限制 + df = df.tail(limit).reset_index(drop=True) + + # 如果过滤后没有数据,返回None + if len(df) == 0: + return None + return df + + except Exception as e: + return None + + +def get_crypto_kl_data(symbol, timeframe, limit=100000, start_time=None, end_time=None): + """优先通过本地数据服务获取加密货币K线,失败时回退至交易所API。""" + start_ms = _parse_time_input(start_time) + end_ms = _parse_time_input(end_time) + + refresh_data_service_metadata() + if state.DATA_SERVICE_AVAILABLE: + try: + df = _fetch_kl_from_datasvc( + symbol=symbol, + timeframe=timeframe, + start_ms=start_ms, + end_ms=end_ms, + limit=limit, + ) + if df is not None and not df.empty: + return df + except Exception as exc: + logger.warning("数据服务请求失败,准备回退至交易所 API:%s", exc) + + return _get_crypto_kl_data_via_ccxt(symbol, timeframe, limit, start_time, end_time) + + +def get_a_stock_kl_data(symbol, timeframe, limit=100000, start_time=None, end_time=None): + """获取A股K线数据""" + try: + # 处理时间戳参数转换为日期字符串 + start_date = None + end_date = None + + if start_time: + try: + # 尝试解析时间戳(毫秒) + start_timestamp = int(start_time) + start_date = datetime.fromtimestamp(start_timestamp / 1000).strftime('%Y-%m-%d') + except (ValueError, TypeError): + # 如果不是时间戳,尝试解析datetime-local格式 (YYYY-MM-DDTHH:MM) + try: + if 'T' in str(start_time): + # datetime-local格式:2025-05-19T06:07 + start_date = str(start_time).split('T')[0] # 只取日期部分 + else: + start_date = str(start_time) + except: + start_date = start_time + + if end_time: + try: + # 尝试解析时间戳(毫秒) + end_timestamp = int(end_time) + end_date = datetime.fromtimestamp(end_timestamp / 1000).strftime('%Y-%m-%d') + except (ValueError, TypeError): + # 如果不是时间戳,尝试解析datetime-local格式 + try: + if 'T' in str(end_time): + # datetime-local格式:2025-05-26T06:07 + end_date = str(end_time).split('T')[0] # 只取日期部分 + else: + end_date = str(end_time) + except: + end_date = end_time + + # 如果用户指定了时间范围,优先获取该范围内的所有数据 + actual_limit = limit + if start_date and end_date: + actual_limit = None # 不限制数据条数,获取完整时间范围数据 + + # 调用A股数据获取器 + df = state.china_stock.get_kl_data(symbol, timeframe, start_date, end_date, actual_limit) + + if df is None: + return None + return df + + except Exception as e: + return None + +def load_crypto_symbols(limit=200): + """加载常见USDT永续合约交易对,返回列表""" + refresh_data_service_metadata() + if state.SYMBOLS: + return state.SYMBOLS[:limit] + try: + markets = state.exchange.load_markets() + symbols = [s for s in markets.keys() if '/USDT' in s and ':USDT' in s] + return symbols[:limit] + except Exception: + return DEFAULT_SYMBOLS[:limit] + diff --git a/web/services/runtime/serialize.py b/web/services/runtime/serialize.py new file mode 100644 index 0000000..fc20c38 --- /dev/null +++ b/web/services/runtime/serialize.py @@ -0,0 +1,300 @@ +from __future__ import annotations + +import pandas as pd + +from chanlun.core.ChanEnum import Chan_BI_DIR, Chan_SEG_DIR, Chan_MACDSEG_DIR, Chan_MACDHISTSET_DIR + +# 辅助函数,转换缠论方向枚举为整数 +def convert_direction(direction): + """转换方向枚举为数字""" + if direction == Chan_BI_DIR.UP or direction == Chan_SEG_DIR.UP: + return 1 + elif direction == Chan_BI_DIR.DOWN or direction == Chan_SEG_DIR.DOWN: + return -1 + else: + return 0 + +def format_time_safely(time_obj, client_tz): + """安全地格式化时间对象,处理字符串和datetime两种情况""" + if time_obj is None: + return None + + if isinstance(time_obj, str): + # 尝试将字符串解析为datetime + try: + from dateutil import parser + time_obj = parser.parse(time_obj) + return time_obj.astimezone(client_tz).isoformat() + except: + return time_obj + else: + # 已经是datetime对象 + return time_obj.astimezone(client_tz).isoformat() + +def serialize_chan_macd_data(chan_macd_data, client_tz): + """序列化ChanMACD数据为JSON可序列化格式""" + serialized_data = { + 'seg_list': [], + 'unittf_list': [], + 'histset_list': [], + # 状态标记数据 + 'high_position_list': [], + 'high_empty_list': [], + 'low_position_list': [], + 'low_empty_list': [], + 'return_zero_list': [], + 'cross0_up_list': [], + 'cross0_down_list': [], + # 新增:输出KLU的继续背驰/分离背驰标志 + 'klu_list': [] + } + + # 序列化seg_list + for seg in chan_macd_data.get('seg_list', []): + try: + seg_data = { + 'start_time': format_time_safely(seg.start_time, client_tz), + 'end_time': format_time_safely(seg.end_time, client_tz) if seg.end_time else None, + 'seg_dir': 'ABOVE' if seg.seg_dir == Chan_MACDSEG_DIR.ABOVE else 'UNDER', + 'klu_count': len(seg.klu_list) if hasattr(seg, 'klu_list') else 0, + 'unittf_count': len(seg.unittf_list) if hasattr(seg, 'unittf_list') else 0, + 'histset_count': len(seg.hist_set) if hasattr(seg, 'hist_set') else 0 + } + serialized_data['seg_list'].append(seg_data) + except Exception as e: + print(f"序列化seg出错: {e}") + continue + + # 序列化unittf_list(兼容新结构与枚举类型) + for unittf in chan_macd_data.get('unittf_list', []): + try: + dir_value = getattr(unittf, 'uinttf_dir', None) + dir_name = getattr(dir_value, 'name', dir_value if isinstance(dir_value, str) else None) + start_t = getattr(unittf, 'start_type', None) + start_type = getattr(start_t, 'name', start_t) + end_t = getattr(unittf, 'end_type', None) + end_type = getattr(end_t, 'name', end_t) + peak_abs = getattr(unittf, 'peak_abs', None) + if peak_abs is None: + peak_abs = getattr(unittf, 'peak_hist', None) + length = getattr(unittf, 'length', None) + if length is None: + length = len(unittf.klu_list) if hasattr(unittf, 'klu_list') else None + + unittf_data = { + 'start_time': format_time_safely(getattr(unittf, 'start_time', None), client_tz), + 'end_time': format_time_safely(getattr(unittf, 'end_time', None), client_tz) if getattr(unittf, 'end_time', None) else None, + 'dir': dir_name, # 'ABOVE' | 'UNDER' | None + 'start_type': start_type, # e.g. 'START' | 'CROSS0' | 'NEAR0_UP' | 'NEAR0_DOWN' + 'end_type': end_type, + 'invalid': getattr(unittf, 'invalid', False), + 'peak_abs': peak_abs, + 'length': length, + 'klu_count': len(unittf.klu_list) if hasattr(unittf, 'klu_list') else 0, + 'histset_count': len(unittf.histset_list) if hasattr(unittf, 'histset_list') else 0 + } + serialized_data['unittf_list'].append(unittf_data) + except Exception as e: + print(f"序列化unittf出错: {e}") + continue + + # 序列化histset_list + for histset in chan_macd_data.get('histset_list', []): + try: + histset_data = { + 'start_time': format_time_safely(getattr(histset, 'start_time', None), client_tz), + 'end_time': format_time_safely(getattr(histset, 'end_time', None), client_tz), + 'histset_dir': 'ABOVE' if histset.histset_dir == Chan_MACDHISTSET_DIR.ABOVE else 'UNDER', + 'klu_count': len(histset.klu_list) if hasattr(histset, 'klu_list') else 0 + } + serialized_data['histset_list'].append(histset_data) + except Exception as e: + print(f"序列化histset出错: {e}") + continue + + # 序列化状态标记数据 + # 序列化高位列表 + for high_pos in chan_macd_data.get('high_position_list', []): + try: + high_pos_data = { + 'time': format_time_safely(high_pos['time'], client_tz), + 'end_time': format_time_safely(high_pos.get('end_time'), client_tz) if high_pos.get('end_time') else None, + 'type': high_pos.get('type', 'start'), + 'macd': high_pos.get('macd'), + 'signal': high_pos.get('signal'), + 'macdhist': high_pos.get('macdhist'), + 'end_macd': high_pos.get('end_macd'), + 'end_signal': high_pos.get('end_signal'), + 'end_macdhist': high_pos.get('end_macdhist') + } + serialized_data['high_position_list'].append(high_pos_data) + except Exception as e: + print(f"序列化high_position出错: {e}") + continue + + # 序列化高位空列表 + for high_empty in chan_macd_data.get('high_empty_list', []): + try: + high_empty_data = { + 'time': format_time_safely(high_empty['time'], client_tz), + 'end_time': format_time_safely(high_empty.get('end_time'), client_tz) if high_empty.get('end_time') else None, + 'type': high_empty.get('type', 'start'), + 'macd': high_empty.get('macd'), + 'signal': high_empty.get('signal'), + 'macdhist': high_empty.get('macdhist'), + 'end_macd': high_empty.get('end_macd'), + 'end_signal': high_empty.get('end_signal'), + 'end_macdhist': high_empty.get('end_macdhist') + } + serialized_data['high_empty_list'].append(high_empty_data) + except Exception as e: + print(f"序列化high_empty出错: {e}") + continue + + # 序列化低位与低位空 + for low_pos in chan_macd_data.get('low_position_list', []): + try: + low_pos_data = { + 'time': format_time_safely(low_pos['time'], client_tz), + 'end_time': format_time_safely(low_pos.get('end_time'), client_tz) if low_pos.get('end_time') else None, + 'type': low_pos.get('type', 'start'), + 'macd': low_pos.get('macd'), + 'signal': low_pos.get('signal'), + 'macdhist': low_pos.get('macdhist'), + 'end_macd': low_pos.get('end_macd'), + 'end_signal': low_pos.get('end_signal'), + 'end_macdhist': low_pos.get('end_macdhist') + } + serialized_data['low_position_list'].append(low_pos_data) + except Exception as e: + print(f"序列化low_position出错: {e}") + continue + + for low_empty in chan_macd_data.get('low_empty_list', []): + try: + low_empty_data = { + 'time': format_time_safely(low_empty['time'], client_tz), + 'end_time': format_time_safely(low_empty.get('end_time'), client_tz) if low_empty.get('end_time') else None, + 'type': low_empty.get('type', 'start'), + 'macd': low_empty.get('macd'), + 'signal': low_empty.get('signal'), + 'macdhist': low_empty.get('macdhist'), + 'end_macd': low_empty.get('end_macd'), + 'end_signal': low_empty.get('end_signal'), + 'end_macdhist': low_empty.get('end_macdhist') + } + serialized_data['low_empty_list'].append(low_empty_data) + except Exception as e: + print(f"序列化low_empty出错: {e}") + continue + + # 序列化归零轴列表 + for return_zero in chan_macd_data.get('return_zero_list', []): + try: + return_zero_data = { + 'time': format_time_safely(return_zero['time'], client_tz), + 'end_time': format_time_safely(return_zero.get('end_time'), client_tz) if return_zero.get('end_time') else None, + 'type': return_zero.get('type', 'start'), + 'macd': return_zero.get('macd'), + 'signal': return_zero.get('signal'), + 'macdhist': return_zero.get('macdhist'), + 'end_macd': return_zero.get('end_macd'), + 'end_signal': return_zero.get('end_signal'), + 'end_macdhist': return_zero.get('end_macdhist') + } + serialized_data['return_zero_list'].append(return_zero_data) + except Exception as e: + print(f"序列化return_zero出错: {e}") + continue + + # 序列化穿越零轴列表 + for cross0_up in chan_macd_data.get('cross0_up_list', []): + try: + cross0_up_data = { + 'time': format_time_safely(cross0_up['time'], client_tz), + 'type': cross0_up.get('type', 'start'), + 'macd': cross0_up.get('macd'), + 'signal': cross0_up.get('signal'), + 'macdhist': cross0_up.get('macdhist') + } + serialized_data['cross0_up_list'].append(cross0_up_data) + except Exception as e: + print(f"序列化cross0_up出错: {e}") + continue + + for cross0_down in chan_macd_data.get('cross0_down_list', []): + try: + cross0_down_data = { + 'time': format_time_safely(cross0_down['time'], client_tz), + 'type': cross0_down.get('type', 'start'), + 'macd': cross0_down.get('macd'), + 'signal': cross0_down.get('signal'), + 'macdhist': cross0_down.get('macdhist') + } + serialized_data['cross0_down_list'].append(cross0_down_data) + except Exception as e: + print(f"序列化cross0_down出错: {e}") + continue + + # 序列化 KLU 列表(仅导出需要的时间与背驰标志) + for klu in chan_macd_data.get('klu_list', []): + try: + serialized_data['klu_list'].append({ + 'time': format_time_safely(getattr(klu, 'time', None), client_tz), + 'continue_div': bool(getattr(klu, 'continue_div', False)), + 'separate_div': int(getattr(klu, 'separate_div', 0)) if getattr(klu, 'separate_div', 0) is not None else 0, + 'near0_return': int(getattr(klu, 'near0_return', 0)) if getattr(klu, 'near0_return', 0) is not None else 0 + }) + except Exception as e: + print(f"序列化klu出错: {e}") + continue + + return serialized_data + +def clean_dataframe_for_json(df): + """清理DataFrame数据用于JSON序列化""" + # 创建副本避免修改原始数据 + clean_df = df.copy() + + # 替换NaN值为None + clean_df = clean_df.where(pd.notnull(clean_df), None) + + return clean_df + +def get_uncompleted_seg_list(seg_list, client_tz): + """获取未完成线段列表,正确处理倒数第二个和最后一个未完成线段""" + uncompleted_segs = [seg for seg in seg_list if not seg.is_sure] + + if len(uncompleted_segs) == 0: + return [] + + result = [] + + for i, seg in enumerate(uncompleted_segs): + is_last = (i == len(uncompleted_segs) - 1) # 是否为最后一个未完成线段 + + seg_data = { + 'start_time': seg.start_bi.start_klc.end_time if isinstance(seg.start_bi.start_klc.end_time, str) else seg.start_bi.start_klc.end_time.astimezone(client_tz).isoformat(), + 'sure_time': format_time_safely(seg.sure_time, client_tz) if seg.sure_time else None, + 'start_price': seg.start_bi.start_klc.low if convert_direction(seg.dir) == 1 else seg.start_bi.start_klc.high, + 'direction': convert_direction(seg.dir) + } + + if is_last: + # 最后一个未完成线段:没有结束时间和价格 + seg_data['end_time'] = None + seg_data['end_price'] = None + else: + # 倒数第二个及之前的未完成线段:使用实际的结束时间和价格 + if seg.end_bi and seg.end_bi.end_klc: + seg_data['end_time'] = seg.end_bi.end_klc.end_time if isinstance(seg.end_bi.end_klc.end_time, str) else seg.end_bi.end_klc.end_time.astimezone(client_tz).isoformat() + seg_data['end_price'] = seg.end_bi.end_klc.high if convert_direction(seg.dir) == 1 else seg.end_bi.end_klc.low + else: + # 如果没有结束笔,设为None + seg_data['end_time'] = None + seg_data['end_price'] = None + + result.append(seg_data) + + return result + diff --git a/web/services/runtime/state.py b/web/services/runtime/state.py new file mode 100644 index 0000000..2142ce3 --- /dev/null +++ b/web/services/runtime/state.py @@ -0,0 +1,70 @@ +from __future__ import annotations + +import sys +import os +from collections import OrderedDict +import logging + +import ccxt + +_ROOT = os.path.dirname(os.path.dirname(os.path.dirname(os.path.abspath(__file__)))) +if _ROOT not in sys.path: + sys.path.append(_ROOT) + +from config import MACD_FAST, MACD_SLOW, MACD_SIGNAL, ccxt_proxies +from services.cn_stock import ChinaStockData + +logger = logging.getLogger(__name__) + +class TRADE_POINT_TYPE: + BUY1 = 1 # 一类买点 + BUY2 = 2 # 二类买点 + BUY3 = 3 # 三类买点 + SELL1 = -1 # 一类卖点 + SELL2 = -2 # 二类卖点 + SELL3 = -3 # 三类卖点 + + + +# mutable runtime state +macd_fast_period = MACD_FAST +macd_slow_period = MACD_SLOW +macd_signal_period = MACD_SIGNAL + +_proxies = ccxt_proxies() +_exchange_kwargs = {"enableRateLimit": True} +if _proxies: + _exchange_kwargs["proxies"] = _proxies +exchange = ccxt.binance(_exchange_kwargs) + +china_stock = ChinaStockData() +_zone_cache = {} + +DEFAULT_TIMEFRAME_LABELS = OrderedDict([ + ("1m", "1分钟"), + ("3m", "3分钟"), + ("5m", "5分钟"), + ("15m", "15分钟"), + ("30m", "30分钟"), + ("1h", "1小时"), + ("2h", "2小时"), + ("4h", "4小时"), + ("6h", "6小时"), + ("8h", "8小时"), + ("12h", "12小时"), + ("1d", "日线"), + ("3d", "3日线"), + ("1w", "周线"), + ("1M", "月线"), +]) + +DEFAULT_SYMBOLS = [ + 'SOL/USDT:USDT', 'BTC/USDT:USDT', 'ETH/USDT:USDT', 'BNB/USDT:USDT', 'XRP/USDT:USDT', 'WIF/USDT:USDT', + 'ADA/USDT:USDT', 'DOGE/USDT:USDT', 'AVAX/USDT:USDT', 'DOT/USDT:USDT', 'MATIC/USDT:USDT' +] + +TIMEFRAMES = DEFAULT_TIMEFRAME_LABELS.copy() +SYMBOLS = DEFAULT_SYMBOLS.copy() +DATA_SERVICE_AVAILABLE = False +SERVICE_METADATA_LAST_REFRESH = 0 + diff --git a/web/services/runtime/timeframes.py b/web/services/runtime/timeframes.py new file mode 100644 index 0000000..3a9d61a --- /dev/null +++ b/web/services/runtime/timeframes.py @@ -0,0 +1,121 @@ +from __future__ import annotations + +from collections import OrderedDict +from .state import DEFAULT_TIMEFRAME_LABELS + +def _zone_cache_ttl(tf_name: str) -> int: + """根据时间周期返回缓存过期时间(秒)""" + minutes = timeframe_to_minutes(tf_name) or 5 + if minutes <= 5: + return 120 # 5m及以下: 2分钟 + elif minutes <= 15: + return 300 # 15m: 5分钟 + elif minutes <= 60: + return 600 # 1h: 10分钟 + else: + return 1800 # 4h+: 30分钟 + + +def timeframe_to_minutes(tf: str): + """将时间周期转换为分钟数,用于排序。""" + if not tf: + return None + unit = tf[-1] + try: + value = int(tf[:-1]) + except (ValueError, TypeError): + return None + multiplier = { + 'm': 1, + 'h': 60, + 'd': 1440, + 'w': 10080, + 'M': 43200, # 30天近似 + }.get(unit) + if multiplier is None: + return None + return value * multiplier + + +def format_timeframe_label(tf: str) -> str: + """将时间周期转换为可读标签。""" + if not tf: + return tf + unit = tf[-1] + try: + value = int(tf[:-1]) + except (ValueError, TypeError): + return tf + if unit == 'm': + return f"{value}分钟" + if unit == 'h': + return f"{value}小时" + if unit == 'd': + return "日线" if value == 1 else f"{value}日线" + if unit == 'w': + return "周线" if value == 1 else f"{value}周线" + if unit == 'M': + return "月线" if value == 1 else f"{value}月线" + return tf + + +def build_timeframe_labels(timeframes): + ordered = sorted( + timeframes, + key=lambda tf: timeframe_to_minutes(tf) if timeframe_to_minutes(tf) is not None else float('inf'), + ) + labels = OrderedDict() + for tf in ordered: + labels[tf] = format_timeframe_label(tf) + return labels + + +def compute_timeframe_defaults(labels_ordered): + """ + 根据已排序的「周期 → 中文标签」映射,计算主 / 次 / 次次周期默认值。 + labels_ordered: OrderedDict 或按插入顺序排列的 dict。 + """ + if not labels_ordered: + labels_ordered = DEFAULT_TIMEFRAME_LABELS.copy() + timeframe_keys = list(labels_ordered.keys()) + preferred_main = next((tf for tf in ['5m', '15m', '1h'] if tf in labels_ordered), None) + default_main = preferred_main or (timeframe_keys[0] if timeframe_keys else '1m') + if default_main not in labels_ordered and timeframe_keys: + default_main = timeframe_keys[0] + + if timeframe_keys: + try: + idx = timeframe_keys.index(default_main) + default_element = timeframe_keys[idx - 1] if idx > 0 else timeframe_keys[0] + except ValueError: + default_element = timeframe_keys[0] + else: + default_element = default_main + + if timeframe_keys: + try: + idx_el = timeframe_keys.index(default_element) + default_sub_sub = timeframe_keys[idx_el - 1] if idx_el > 0 else timeframe_keys[0] + except ValueError: + default_sub_sub = timeframe_keys[0] + else: + default_sub_sub = default_element + + return default_main, default_element, default_sub_sub, timeframe_keys + +def is_smaller_timeframe(tf1, tf2): + """判断时间周期tf1是否小于tf2""" + tf1_value = timeframe_to_minutes(tf1) + tf2_value = timeframe_to_minutes(tf2) + if tf1_value is None or tf2_value is None: + return False + return tf1_value < tf2_value + +def is_smaller_or_equal_timeframe(tf1, tf2): + """判断时间周期tf1是否小于等于tf2""" + tf1_value = timeframe_to_minutes(tf1) + tf2_value = timeframe_to_minutes(tf2) + if tf1_value is None or tf2_value is None: + return False + return tf1_value <= tf2_value + diff --git a/web/tests/test_analyze_contract.py b/web/tests/test_analyze_contract.py index 92edb40..9d5f44c 100644 --- a/web/tests/test_analyze_contract.py +++ b/web/tests/test_analyze_contract.py @@ -1,14 +1,53 @@ -""" /api/analyze 契约冒烟:关键字段存在于契约清单。""" +"""ECR-002:加深 /api/analyze 相关契约 —— mock 行情 + analyze_chan 关键字段快照。""" from __future__ import annotations import json import sys from pathlib import Path +from unittest.mock import patch + +import pandas as pd +import pytest ROOT = Path(__file__).resolve().parents[2] sys.path.insert(0, str(ROOT)) sys.path.insert(0, str(ROOT / "web")) +from tests.generate_golden import make_ohlcv # noqa: E402 + + +CONTRACT_KEYS = json.loads( + (ROOT / "tests" / "fixtures" / "analyze_contract_keys.json").read_text(encoding="utf-8") +) + +# analyze_chan 直接返回的对象字段(未序列化前) +ANALYZE_CHAN_KEYS = { + "klc_list", + "klu_list", + "bi_list", + "seg_list", + "zs_list", + "bi_zs_list", + "bsp_list", + "klc_fx_info", + "chan_macd", + "ema52_dict", +} + +CHAN_MACD_SERIALIZED_KEYS = { + "seg_list", + "unittf_list", + "histset_list", + "high_position_list", + "high_empty_list", + "low_position_list", + "low_empty_list", + "return_zero_list", + "cross0_up_list", + "cross0_down_list", + "klu_list", +} + def test_analyze_route_registered(): from app import app @@ -21,9 +60,60 @@ def test_analyze_route_registered(): def test_contract_keys_stable(): - keys = json.loads( - (ROOT / "tests" / "fixtures" / "analyze_contract_keys.json").read_text( - encoding="utf-8" + assert "bi_list" in CONTRACT_KEYS and "seg_list" in CONTRACT_KEYS + for k in ("kline_data", "macd", "zs_list", "bsp_list", "chan_macd"): + assert k in CONTRACT_KEYS + + +def test_analyze_chan_keys_on_fixture(): + from services.runtime import add_indicators, analyze_chan + + df = add_indicators(make_ohlcv(400)) + result = analyze_chan(df, symbol="TEST/USDT:USDT", timeframe="5m") + assert set(result.keys()) == ANALYZE_CHAN_KEYS + assert isinstance(result["bi_list"], list) + assert isinstance(result["seg_list"], list) + assert isinstance(result["chan_macd"], dict) + for k in ("seg_list", "unittf_list", "histset_list"): + assert k in result["chan_macd"] + + +def test_serialize_chan_macd_shape(): + from pytz import timezone + + from services.runtime import add_indicators, analyze_chan, serialize_chan_macd_data + + df = add_indicators(make_ohlcv(200)) + result = analyze_chan(df) + serialized = serialize_chan_macd_data(result["chan_macd"], timezone("Asia/Shanghai")) + assert set(serialized.keys()) == CHAN_MACD_SERIALIZED_KEYS + # JSON 可序列化 + json.dumps(serialized) + + +def test_analyze_http_contract_with_mocked_kl(): + """Flask 测试客户端:mock get_kl_data,断言响应含契约关键字段。""" + from app import app + from services.runtime import add_indicators + + df = add_indicators(make_ohlcv(300)) + df = df.copy() + if "timestamp" not in df.columns: + df["timestamp"] = (pd.to_datetime(df["date"]).astype("int64") // 10**6).astype("int64") + + # analyze 路由使用 `from services.runtime import *`,须 patch 其模块命名空间 + with patch("api.analyze.get_kl_data", return_value=df): + client = app.test_client() + resp = client.get( + "/api/analyze", + query_string={ + "symbol": "BTC/USDT:USDT", + "timeframe": "5m", + "timezone": "Asia/Shanghai", + }, ) - ) - assert "bi_list" in keys and "seg_list" in keys + assert resp.status_code == 200, resp.data[:500] + payload = resp.get_json() + assert payload is not None and "error" not in payload + missing = [k for k in CONTRACT_KEYS if k not in payload] + assert not missing, f"missing contract keys: {missing}" diff --git a/web/tests/test_runtime_facade.py b/web/tests/test_runtime_facade.py new file mode 100644 index 0000000..00667b0 --- /dev/null +++ b/web/tests/test_runtime_facade.py @@ -0,0 +1,55 @@ +"""ECR-002:runtime 门面公开符号 + 子模块可导入。""" +from __future__ import annotations + +import sys +from pathlib import Path + +ROOT = Path(__file__).resolve().parents[2] +sys.path.insert(0, str(ROOT)) +sys.path.insert(0, str(ROOT / "web")) + +REQUIRED = [ + "get_kl_data", + "analyze_chan", + "add_indicators", + "serialize_chan_macd_data", + "clean_dataframe_for_json", + "classify_trend_stage", + "refresh_data_service_metadata", + "TIMEFRAMES", + "SYMBOLS", + "_zone_cache", + "macd_fast_period", + "is_smaller_or_equal_timeframe", + "get_uncompleted_seg_list", +] + + +def test_runtime_facade_exports(): + from services import runtime as R + + for name in REQUIRED: + assert hasattr(R, name), f"missing facade export: {name}" + + +def test_runtime_submodules_importable(): + from services.runtime import state, timeframes, market_data, indicators, analyze, serialize + + assert state.exchange is not None + assert callable(timeframes.timeframe_to_minutes) + assert callable(market_data.get_kl_data) + assert callable(indicators.add_indicators) + assert callable(analyze.analyze_chan) + assert callable(serialize.convert_direction) + + +def test_thin_shims_still_reexport(): + from services import market_data as md + from services import chan_analyze as ca + from services import serializers as ser + from services import timeframes as tf + + assert callable(md.get_kl_data) + assert callable(ca.analyze_chan) + assert callable(ser.serialize_chan_macd_data) + assert callable(tf.timeframe_to_minutes)