diff --git a/.gitignore b/.gitignore index 4209ac7..787adbf 100644 --- a/.gitignore +++ b/.gitignore @@ -26,6 +26,10 @@ Thumbs.db .env.* !.env.example +# Temp artifacts (AI session diagrams / scratch) +.tmp-diagrams/ +.tmp-*/ + # DSH / runtime .dsh/ *.tmp diff --git a/docs/02-计划/计划-盘口内存快照.md b/docs/02-计划/计划-盘口内存快照.md new file mode 100644 index 0000000..3796b00 --- /dev/null +++ b/docs/02-计划/计划-盘口内存快照.md @@ -0,0 +1,38 @@ +# 计划:盘口内存快照(QuoteSync/QuoteHub + 同步指示灯)(阶段航点) + +> 编号:PLAN-014 | 粒度:阶段航点 | 创建:2026-09-02 | 状态:**已实施(待验收)** +> 派生自终极目标:目标-003(按策略监控市场)、目标-008(真实交易系统接入) +> 依据需求:**R-015(已定稿,2026-09-02,老师逐项拍板)** —— 符合入范围门槛 +> 设计约束:技术约束-010 变更(行情服务端中转改纯内存)、技术约束-012 变更(market_quotes_cache 退役)、技术约束-016(分域沿用);本次新增 技术约束-018、产品约束-012 + +## 目标 + +盘口数据与持仓数据对称收敛:QuoteSync/QuoteHub 纯内存管理(5s REST 同步 + 涨停跌停合约信息),market_quotes_cache 退役(DROP),价格单一入口;会话头部新增持仓/盘口同步指示灯(绿黄灰、点击即同步)。 + +## 范围 + +**做**: +1. QuoteSync(新):启动 prime + 5s REST 刷新 watch 集合 + 涨停跌停按交易日拉取缓存 + 热切换重置;WS 通路不迁移(删除); +2. QuoteHub(新):内存快照 + 合约信息缓存 + watch 集合 + 读穿透(走 dataSource.getTicks)+ getQuote(s) 对外;对外方法签名与旧 hub 兼容(getByCodes/watch/stats/getQuote); +3. 数据源适配:QmtBridgeRestDataSource 新增 getTicks(codes)(/data/tick 语义化批量); +4. 存储清理:SqliteStore 删行情表读写 + DROP TABLE market_quotes_cache(幂等)+ migrateJson/isEmpty 联动;DataStore 删 loadMarket/getMarketQuote(s)/setMarketQuotes; +5. API:market-snapshot 返回加涨停/跌停/昨收/updatedAt;新增 sync-status(吸收 market-stats)+ sync-now {domain}; +6. 前端:QmtConnectionChip 加 SyncIndicators(两圆点,10s 轮询,点击即同步);MarketDataProvider 清理 wsInfo 残留; +7. 策略持仓表行情列(老师确认并入):COLUMN_META + 涨停价/跌停价/今开/最高(defaultVisible: false,纯价格);StrategyTab.renderDataCell 对应 case;normalizeStrategyColumns 缺省显隐跟随 defaultVisible; +7. 回归脚本 test-quote-sync.mjs(纯内存 mock)。 + +**不做**:watch 收缩/降频;个股停牌标识;WS 相关需求推进;PriceCell 改动;最低/成交量/成交额列;距离涨停百分比展示。 + +## 涉及文件 + +- 新增:src/market/QuoteSync.js、src/market/QuoteHub.js、src/client/views/SyncIndicators.jsx、scripts/test-quote-sync.mjs +- 删除:src/market/MarketFeed.js、src/market/MarketDataHub.js +- 修改:src/data-source/QmtBridgeRestDataSource.js(getTicks)、src/storage/SqliteStore.js、src/storage/DataStore.js、src/api/market.js、src/api/qmt-connections.js(热切换联动)、src/index.js、src/client/views/QmtConnectionChip.jsx、src/client/market/MarketDataProvider.jsx + +## 实现步骤 + +1. 文档链(本计划 + 迭代三件套 + 约束条目);2. 数据源 getTicks;3. QuoteHub + QuoteSync;4. 存储清理 + DROP;5. API 改造;6. 前端指示灯 + 清理;7. 回归脚本 + typecheck + build + 存量回归;8. 复盘。 + +## 验收要点 + +见 `docs/04-迭代记录/13-盘口内存快照/验收标准.md`。 diff --git a/docs/03-设计约束/产品功能约束.md b/docs/03-设计约束/产品功能约束.md index e8053aa..2c49650 100644 --- a/docs/03-设计约束/产品功能约束.md +++ b/docs/03-设计约束/产品功能约束.md @@ -20,6 +20,8 @@ | 产品约束-008 | 交易记录支持**按策略过滤**:交易记录 tab 提供策略过滤下拉(全部 / 各策略 / 未关联),对今日(QMT 实时)与历史(本地 SQLite)均生效;策略归属 = 委托时间 join 持仓生命周期窗口推导(一码多策略取份额最大,未命中=未关联);历史范围展示本地积累数据(不再「接口开发中」占位) | 2026-09-01 | 生效 | - | 新增(2026-09-01 R-009 定稿 + 迭代 07 实施):策略过滤 + 历史本地展示 | | 产品约束-010 | 策略自定义字段配置:每个策略可在设置页「策略分组」子 tab 配置自定义字段定义(字段名 / 类型文本·数字·布尔·枚举 / 枚举选项 / 默认值),定义随策略存 settings(不落库);该策略下每个持仓(strategy_holdings 行)按所属策略的定义存取一份字段值(values JSON 列,持仓级键值对,key 对齐定义、允许扩展额外键);旧策略无定义时行为与现状一致(不渲染字段区、不迁移历史值) | 2026-09-02 | 生效 | - | 新增(2026-09-02 R-013 定稿 Q1-Q4 + D6):定义随策略走、值随持仓行,类型化(含枚举),旧策略兼容 | | 产品约束-009 | 会话 tab 统一由「Tab 设置」管理:设置页「Tab 设置」子 tab 是**所有会话 tab(系统内置 + 策略分组)的唯一顺序与显隐入口**,两类 tab 混排;每行 = 拖动排序 + 显示/隐藏开关;**任何 tab 均不支持重命名与删除**(内置 tab 名称只读,策略命名/删除仍在「策略分组」子 tab);策略改名后 Tab 设置中的名称自动跟随(只存引用);新增策略默认追加到列表末尾;删除策略联动删除 Tab 设置中对应条目;顺序与显隐唯一数据源 = settings.tabs 有序数组 | 2026-09-02 | 生效 | - | R-011 定稿(2026-09-02 Q1-Q5 确认):Tab 设置 = 显示/隐藏 + 拖动排序(落点立即持久化),全 tab 禁重命名/删除 | +| 产品约束-011 | 持仓页数据(策略持仓 / 全部持仓 / 未分配)以服务端 10s 内存快照为准,**接受最多 10s 滞后**(同步时间前端不显示,老师拍板);QMT 抖动/掉线时持仓页面显示**最后一次快照**而非空白;QMT 已卖光的票(连续 3 轮同步确认消失)本地持仓自动转历史(幽灵清仓,不物理删除);部分减持仅表现为「未分配为负」,不做自动修正 | 2026-09-02 | 生效 | - | 新增(2026-09-02 R-014 定稿):老师四问拍板(内存不落库 / 读穿透兜底 / 幽灵自动清仓 / 不显示同步时间) | +| 产品约束-012 | 行情数据服务形态(R-015):现价/昨收/涨停/跌停统一由盘口内存快照提供(≤5s 更新),价格单一入口;QMT 抖动时页面价格保持旧值不空白;会话头部(QMT 健康灯旁)提供**持仓/盘口同步指示灯**——绿=同步正常(持仓 30s/盘口 15s 内)、黄=同步失败中快照陈旧、灰=从未同步;悬停显示同步时间/快照量/失败数/错误摘要;点击灯 = 立即触发该域同步;不做个股停牌标识(另议) | 2026-09-02 | 生效 | - | 新增(2026-09-02 R-015 定稿):老师提出指示灯,采纳 AI 推荐三态/点击即同步/取消 PriceCell 灰点 | \ No newline at end of file diff --git a/docs/03-设计约束/技术方案约束.md b/docs/03-设计约束/技术方案约束.md index 06cde9e..57290c4 100644 --- a/docs/03-设计约束/技术方案约束.md +++ b/docs/03-设计约束/技术方案约束.md @@ -17,15 +17,17 @@ | 技术约束-006 | 客户端 slots.register 的 component 必须是第二参数(register({...}, Component));settings schema 必须用 schemastery z.object() 函数式定义 | 2026-08-28 | 生效 | - | 迭代 01 复盘沉淀:component 位置错误致 React #130;普通对象 schema 报 schema is not a function | | 技术约束-007 | 插件安装用 dsh plugin add(自动 reconcile bundles),不直接用 pnpm add;bundle patch 顶层必须是 insert 操作 | 2026-08-28 | 生效 | - | 迭代 01 复盘沉淀:pnpm add 不会更新 dsh.profile.bundles | | 技术约束-008 | QMT 连接配置存储复用 one-divine-lot settings namespace(新增 qmtConnections 字段:list[{id,name,baseUrl,order}] + activeId + defaultId),与策略配置同机制持久化;激活切换 = 更新数据源实例的 baseUrl(数据源按请求读取地址,已核实),立即生效无需重启 DSH;插件启动时激活默认配置(列表为空时回退 cordis 注入的 qmtBaseUrl 兜底,不做自动迁移);测试连接由服务端代理请求 {baseUrl}/health(避免浏览器跨域) | 2026-08-29 | 生效 | - | R-004 定稿(2026-08-29):Q1/Q2/Q4/Q5/Q8 确认(启动自动激活默认、立即切换、不迁移、超时不配置化、复用 settings) | -| 技术约束-012 | 插件数据存储遵循 **docs/03-设计约束/数据存储设计.md**:存储引擎为 **SQLite(node:sqlite)**——strategy_holdings + market_quotes_cache + **trade_orders + trade_fills(交易记录两表,R-009)**;存储层单票生命周期操作;交易记录两表零冗余 strategy_id/holding_id(外键链推导策略归属);策略定义仍存 DSH settings;旧 JSON(store.json / store.market.json / allocations.json)经一次性迁移脚本 + 启动自动迁移(幂等、迁移前自动备份)后废弃;仅存储引擎替换,对外行为不变 | 2026-09-01 | 生效 | - | 变更(2026-09-01 R-008 定稿 + 迭代 06 实施):JSON data store → SQLite;变更(2026-09-01 R-009 定稿 + 迭代 07 实施):+trade_orders/trade_fills 交易记录两表(零冗余,外键链推导策略归属) | +| 技术约束-012 | 插件数据存储遵循 **docs/03-设计约束/数据存储设计.md**:存储引擎为 **SQLite(node:sqlite)**——strategy_holdings + ~~market_quotes_cache~~(**R-015 退役 DROP,2026-09-02**;行情改内存快照)+ **trade_orders + trade_fills(交易记录两表,R-009)**;存储层单票生命周期操作;交易记录两表零冗余 strategy_id/holding_id(外键链推导策略归属);策略定义仍存 DSH settings;旧 JSON(store.json / store.market.json / allocations.json)经一次性迁移脚本 + 启动自动迁移(幂等、迁移前自动备份)后废弃;仅存储引擎替换,对外行为不变 | 2026-09-01 | 生效 | - | 变更(2026-09-01 R-008 定稿 + 迭代 06 实施):JSON data store → SQLite;变更(2026-09-01 R-009 定稿 + 迭代 07 实施):+trade_orders/trade_fills 交易记录两表;**变更 3(2026-09-02 R-015/迭代 13)**:market_quotes_cache 退役(DROP),行情移出 SQLite | | 技术约束-011 | 测试/回归脚本**禁止在真实数据上执行写操作**:份额写操作(add/remove/move/clear)必须使用独立数据目录(AllocationStorage 支持 ODL_TEST_DATA_DIR 环境变量或 dataDir 参数指向临时目录),只读端点(positions/summary/strategies/market-snapshot)可直连生产 API | 2026-08-31 | 生效 | - | 2026-08-31 数据误删事故沉淀:回归测试误删大连热电/万顺新材份额分配,老师定「测试用独立数据目录」 | -| 技术约束-010 | 行情实时数据由**服务端中转 + 缓存**提供(R-005 演进,2026-08-31 老师改):DSH 服务端做「WS 订阅 + REST 轮询 + 行情缓存」,前端统一轮询 /odl/api/market-snapshot(不做前端直连,无跨域);行情持久化到 store.market.json(重启不丢价,首屏快速展现);QMT Bridge WS 推送是会话级/有状态行为(归属 QMT Bridge 工作空间) | 2026-08-31 | 生效 | - | 变更(2026-08-31):老师由「前端直连 WS」改为「服务端中转 + 缓存」——解决跨域与 WS 语义不稳定问题;2026-09-01 加行情持久化与启动 prime | +| 技术约束-010 | 行情实时数据由**服务端中转 + 内存缓存**提供(R-005 演进,2026-08-31 老师改;R-015 二改,2026-09-02):DSH 服务端做「REST 轮询 + 行情内存缓存」,前端统一轮询 /odl/api/market-snapshot(不做前端直连,无跨域);~~行情持久化到 store.market.json~~(**R-015 失效:改纯内存,QuoteSync 启动 prime + 读穿透保证秒级有价**);~~QMT Bridge WS 推送~~(**R-015 移除 WS 通路**,见技术约束-018) | 2026-08-31 | 生效 | - | 变更(2026-08-31):老师由「前端直连 WS」改为「服务端中转 + 缓存」;2026-09-01 加行情持久化与启动 prime;**变更 2(2026-09-02 R-015/迭代 13)**:持久化与 WS 均失效——行情纯内存(QuoteSync/QuoteHub),落库与 WS 通路删除 | | 技术约束-009 | 会话头部快捷切换控件挂载 DSH 开放 slot `conversation.session.header.actions`(多实例挂载点,按 order 排序多插件共存):客户端插件以独立 id 并排注册(DSH 内置 PTC 标签 order=-10,本控件 order=-9),不改动 DSH 宿主;控件经 ConnectionProvider 包装复用现有 RPC 通道与 /odl/api/* 端点 | 2026-08-29 | 生效 | - | R-004 定稿(2026-08-29):Q9 确认;宿主代码审查核实 slot 机制与内置插件注册方式 | | 技术约束-013 | 交易记录本地存储(R-009):QMT 当日交易数据(委托/成交)由服务端 TradeSync 定时同步落 SQLite(启动预热 + 60s 定时 + UPSERT 幂等,只同步当日);trade_orders(委托主行,order_id 主键 + insert_ts 派生时间列 + **strategy_id/holding_id 手动归属列**)+ trade_fills(成交明细,trade_id 主键、order_id 外键)两表;**委托归属由用户在交易记录 tab 手动设置**(候选 = 该 code 当前持仓策略 + 未关联,全手动选、可随时改、以最终为准);**UPSERT 不覆盖归属列**(手动指定为插件逻辑);本地历史查询走 trades/history 端点(策略过滤 = 用户设置的归属);今日实时仍走 QMT Bridge;QMT 委托/成交 code 无后缀、持仓带后缀 —— 数据源映射层统一 normalizeInstrumentCode 归一化;委托交易日 = insertDate(tradeDate 兜底) | 2026-09-01 | 生效 | - | 新增(2026-09-01 R-009 定稿 + 迭代 07 实施):两表 + 定时同步 + 本地历史查询;变更 1(2026-09-01):+code 归一化 + tradeDate 兜底;变更 2(2026-09-01 老师二次定稿):归属改**手动设置**(trade_orders 冗余 strategy_id+holding_id,UPSERT 不覆盖归属列),弃算法推导 | | 技术约束-015 | 策略自定义字段存储(R-013,2026-09-02):字段定义随策略定义存 settings.strategies 扩展 configSchema([{key,label,type,enum?,def}],type ∈ text|number|boolean|enum,旧项缺省空数组);字段值落 strategy_holdings 新增 values TEXT(JSON 键值对,key 对齐 configSchema.key,允许额外键=可扩展,NULL=未配置);补列用幂等 ALTER(沿用 _ensureTradeAttributionColumns 模式,只读连接容忍);API:strategy-positions 每行附 values,新增 holdings/values-update {holdingId, values} 写回,服务端按 configSchema 校验(number=有限数、enum=在选项内、boolean=布尔),空值/缺省可写入;持仓生命周期操作(openHolding/addShares/reduceShares/closeHolding)不碰 values 列 | 2026-09-02 | 生效 | - | 新增(2026-09-02 R-013 定稿 + PLAN-012):定义 settings + 值 SQLite 列 + 幂等补列 + 类型校验 | | 技术约束-014 | 会话 tab 注册与顺序显隐(R-011):客户端注册统一读 **settings.tabs 有序数组**(唯一顺序与显隐来源,内置条目 refKey + 策略条目 refId),按 order 排序、过滤 visible 后注册(builtin 走内置 render、strategy 走 StrategyTab),移除硬编码 order 间隔(原内置 10/11/12、策略 13+);settings.tabs 从布尔对象升级为有序数组,读取时对旧格式(布尔对象 + 策略自带 order/visible)静默归一化迁移(旧隐藏策略迁移后显示),写入即落库;策略定义表收窄为 {id,name}(去除 visible/order);tabs/update 语义改为整表更新(顺序 + 显隐),strategies/add 联动追加 tab 条目(末尾),strategies/remove 联动删除对应 tab 条目,废弃 strategies/move | 2026-09-02 | 生效 | - | R-011 定稿(2026-09-02 Q1-Q5 确认):统一 tabs 有序数组 + 自动迁移 + 联动增删 | | 技术约束-016 | src 目录按功能域归类(2026-09-02 结构优化):服务端代码**禁止平铺**,按职责域分目录 —— src/data-source/(QmtBridgeRestDataSource + data-source-types + QmtHealthMonitor,数据源与连接健康)、src/storage/(SqliteStore + DataStore,存储层)、src/position/(PositionManager,分仓逻辑)、src/market/(MarketDataHub + MarketFeed,行情)、src/trades/(TradeSync,交易同步);api/ 按领域拆分子文件(positions/strategies/qmt-connections/market/trades),client/ 仅放 UI(views/ 组件 + market/ provider);文件命名 = 类名(PascalCase)+ .js/.jsx;新增服务端模块必须先落对应域目录,无合适域时先讨论补域,不得回退平铺 | 2026-09-02 | 生效 | - | 新增(2026-09-02 结构审查 + 优化落地):component/ 平铺还原为语义分域,删除死代码 AllocationStorage、DataStore.setDataset/removeDataset | +| 技术约束-017 | 持仓内存快照(R-014,2026-09-02):全量实盘持仓由服务端 PositionSync **进程内内存快照**管理(启动预热 + 10s 定时全量拉 /trade/positions → 校验 → 整体替换),**不落库**(无缓存表、不复用 strategy_holdings——账本与对账单分离,holding_id 交易锚点不掺易变快照;判断标准:外部可一次调用重取全的实时投影不落库);PositionManager.getAllPositions 以快照为准(strategy-positions / unallocated / summary 三接口不再请求时穿透 QMT),快照为空读穿透兜底(当场拉一次并回填);同步失败保留上次快照(不清空不报错);空快照双重确认(/health 可用 + getAsset 账户身份可识别)才接受为真清仓;幽灵清仓:快照连续 3 轮(约 30s)消失的 code 该码全部策略当前持仓自动 closeHolding 转历史(不物理删除),账户身份守卫防误清(accountId 未知当轮跳过、切换当轮重置跳过);syncNow 允许手动调用(mounted 只管定时循环,不拦手动同步) | 2026-09-02 | 生效 | - | 新增(2026-09-02 R-014 定稿 + 迭代 12 实施):AI 初版建缓存表方案经老师质疑反转为内存方案(落库三问:不可再生?重启首屏依赖?读放大/复杂查询?全否 → 内存) | +| 技术约束-018 | 盘口内存快照(R-015,2026-09-02):行情数据由 QuoteSync(取数:启动 prime + 5s REST 定时刷 watch 集合 + 涨停跌停经 /data/instrument 按**交易日**内存缓存)与 QuoteHub(存查:内存快照 Map + watchCodes Set + 读穿透走 dataSource.getTicks + getQuote/getByCodes 对外)管理,**替换并删除 MarketFeed/MarketDataHub**(方案 A,不留兼容壳);**WS 数据通路移除**(ingest 入口带 source 标签留回归口子);**market_quotes_cache 表退役 DROP**(幂等),价格单一入口 = QuoteHub,DataStore 行情方法(loadMarket/getMarketQuote(s)/setMarketQuotes)删除;同步失败保留内存旧值;watchCodes 维持只进不出无上限;涨停/跌停/昨收不落库 | 2026-09-02 | 生效 | - | 新增(2026-09-02 R-015 定稿 + 迭代 13 实施):老师五拍板(替换/WS 移除/纯内存 DROP/watch 现状/指示灯);warmup bug 复现(loadMarket return this 残迹)为不落库关键证据 | \ No newline at end of file diff --git a/docs/04-迭代记录/13-盘口内存快照/技术实现方案.md b/docs/04-迭代记录/13-盘口内存快照/技术实现方案.md new file mode 100644 index 0000000..8ed0e27 --- /dev/null +++ b/docs/04-迭代记录/13-盘口内存快照/技术实现方案.md @@ -0,0 +1,106 @@ +# 技术实现方案:13-盘口内存快照 + +> 迭代编号:13 | 依据:PLAN-014 + R-015(老师拍板:方案 A 替换 / WS 移除 / 纯内存 DROP / watch 维持现状 / 指示灯推荐逻辑) + +## 1. 数据源扩展(QmtBridgeRestDataSource.getTicks) + +```js +/** 批量盘口(GET /data/tick?codes=a,b,c)→ { [code]: snapshot }(保留原始字段 + updatedAt) */ +async getTicks(codes) { + // j.ok && j.data → data 原样返回(tick 结构已有语义字段 lastPrice/lastClose/open/high/low/... + time/timetag) + // 每 snapshot 统一附加 updatedAt = Date.now()(服务端收到时刻,价龄判断基准) +} +``` + +- 首次启用 /data/tick 的适配层通道(此前 MarketDataHub 自己 fetch,绕过了适配层——本轮修正层级); +- 涨停跌停走既有 getInstrument(code)(首次投入使用):upStopPrice/downStopPrice/preClose。 + +## 2. QuoteHub(src/market/QuoteHub.js,新) + +```js +class QuoteHub { + quotes = new Map(); // code → snapshot(全量 tick 字段 + updatedAt;只读约定) + instruments = new Map(); // code → { upStopPrice, downStopPrice, preClose, tradingDay }(交易日内存缓存) + watchCodes = new Set(); // 关注集合(维持现状只进不出,无上限——老师拍板) + stats = { syncCount, failCount, lastError, lastSyncedAt, restCount, watchSize }; +} +``` + +- `getQuotes(codes)`:内存命中 + miss 读穿透(dataSource.getTicks 补拉并回填,等价旧三级命中的内存→REST 两级,磁盘层删除); +- `getQuote(code)`:单码便捷;`watch(codes)`:关注集合登记(QuoteSync prime/刷新、api 查询共用); +- `ingest(data, { source })`:写入统一入口(source 标签:rest|instrument——WS 将来回归时加 source 即可,入口不变); +- 合约信息:`getInstrumentCached(code, tradingDay)`——QuoteSync 判定交易日变化后失效重拉(涨停跌停当天不变)。 + +## 3. QuoteSync(src/market/QuoteSync.js,新;与 PositionSync 同构) + +``` +start(): prime(持仓 code 盘口,走 hub.watch + getTicks 批量)+ setInterval(5s) + 每轮:codes = [...hub.watchCodes](维持现状) + 分批 50 → dataSource.getTicks(batch) → hub.ingest(data, {source:'rest'}) + watch 中无合约缓存或 tradingDay 变化的 code → getInstrument 懒拉(失败不阻塞本轮) + 失败:stats.failCount++ / lastError;**内存不动(保留上次价)** +setBaseUrl(url):热切换 → 清空 instruments 缓存(新账户/新市场环境)+ 立即 prime 一次 +stop(): 清定时器;内存保留 +getSyncedAt() / getStatus(): 指示灯数据源 +``` + +- 对称性:PositionSync(10s/持仓)/ QuoteSync(5s/盘口)/ TradeSync(60s/交易)三域同构,全部「失败保留、stats 可观测、手动 syncNow 可触发」。 + +## 4. 存储清理(market_quotes_cache 退役,老师拍板 DROP) + +- SqliteStore:SCHEMA_SQL 删建表;新增 `DROP TABLE IF EXISTS market_quotes_cache`(init 内幂等执行,存量库清理);删 getMarketQuote/getMarketQuotes/setMarketQuotes/_mapQuote;migrateJson 删行情迁移段;isEmpty 只看 strategy_holdings; +- DataStore:删 loadMarket(return this 残迹,warmup bug 根源)/getMarketQuote/getMarketQuotes/setMarketQuotes/marketLoaded; +- 存量 store.market.json 的 .bak 不动(历史备份无碍)。 + +## 5. API(src/api/market.js + 新 sync-status) + +``` +market-snapshot { codes } → { [code]: { lastPrice, lastClose, upStopPrice, downStopPrice, updatedAt, ...tick 字段 } } +sync-status → { + position: { syncedAt, ageMs, fresh: bool, snapshotSize, failCount, lastError, periodMs }, + quote: { syncedAt, ageMs, fresh: bool, snapshotSize, watchSize, failCount, lastError, periodMs }, + qmt: { ...qmtHealthMonitor.getStatus() } +} +sync-now { domain: 'position' | 'quote' } → 各 Sync.syncNow()(手动触发,指示灯点击用) +market-stats 端点删除(并入 sync-status) +``` + +- 状态判定:fresh = syncedAt 在 3×周期内(持仓 30s / 盘口 15s);前端按 ageMs 二次校准显示。 + +## 6. 前端(SyncIndicators.jsx 新 + 两处清理) + +- **SyncIndicators**(挂 QmtConnectionChip 内 QmtHealthIndicator 右侧):两圆点「持」「行」; + - 10s 轮询 sync-status;点击圆点 → sync-now({domain}) → 刷新; + - 颜色:green(fresh) / yellow(stale) / gray(never);title 含同步时间、快照量、失败数、lastError; + - 样式复用 QmtHealthIndicator 圆点(10px 圆 + glow),色值走 --dsw-* token; +- MarketDataProvider:删 `wsInfo: null` 残留;轮询/注册逻辑不变(getByCodes 签名兼容); +- QmtConnectionChip:挂载 SyncIndicators(健康灯右侧)。 + +## 7. 装配(src/index.js) + +```js +const marketHub = new QuoteHub({ logger }); +const marketFeed = new QuoteSync({ hub: marketHub, runtime: { settings, dataSource }, logger }); +marketFeed.start(startup.baseUrl ?? config?.qmtBaseUrl); +// dispose:marketFeed.stop()(保留旧变量名,改动最小;注释标注新职责) +``` + +- registerApi 入参 marketHub/marketFeed 对象形态不变(api/market.js 内部改用新方法); +- QmtHealthMonitor 不动。 + +## 8. 回归脚本(scripts/test-quote-sync.mjs,纯内存 mock) + +- mock:可编程 dataSource(ticks/instrument/positions 可变状态); +- 用例:prime + 5s 刷新基本流 / 失败保留上次价 / watch 空不空转 / 读穿透回填 / 涨停跌停按交易日缓存与日切失效 / 热切换清合约缓存 / sync-status 状态推导(绿黄灰边界)/ 存储清理(DROP 幂等:SqliteStore 临时目录建旧结构库 → init 后表消失且持仓数据无损)。 + +## 9. 策略持仓表行情列(老师确认并入本轮) + +- **settings.js**:COLUMN_META 追加 4 项(均 `defaultVisible: false`)——`{ key:'upStopPrice', label:'涨停价' } / { key:'downStopPrice', label:'跌停价' } / { key:'open', label:'今开' } / { key:'high', label:'最高' }`;normalizeStrategyColumns 两处适配:默认列构造带 defaultVisible,未覆盖列的缺省显隐从恒 true 改为 `defaultVisible !== false`(存量 strategyColumns 无需迁移——新列缺省即隐藏); +- **StrategyTab.renderDataCell** switch 增 4 case:`getPrice(p.code)?.upStopPrice / downStopPrice / open / high`,复用 fmtPrice 纯价格展示(数据缺失显示 —,与现价列同行为); +- 列设置弹层零改动(自动从 COLUMN_META 出现在列表,kind=base 标「数据」); +- 取值来源:market-snapshot 下发的 tick 字段 + 合约信息(upStopPrice/downStopPrice 来自 QuoteSync 交易日缓存),经 MarketDataProvider prices map join 到行。 + +## 10. 验证 + +- typecheck + build;test-quote-sync 全绿;存量回归(test-position-sync 34 项 + test-r013 21 项)全绿; +- 老师人工验收(见验收标准)。 diff --git a/docs/04-迭代记录/13-盘口内存快照/迭代复盘.md b/docs/04-迭代记录/13-盘口内存快照/迭代复盘.md new file mode 100644 index 0000000..75f90f7 --- /dev/null +++ b/docs/04-迭代记录/13-盘口内存快照/迭代复盘.md @@ -0,0 +1,41 @@ +# 迭代复盘:13-盘口内存快照(QuoteSync/QuoteHub + 同步指示灯) + +> 复盘日期:2026-09-02 | 迭代状态:**已实施,待老师人工验收** +> 关联需求:R-015 | 关联计划:PLAN-014 + +## 结果 + +迭代 13 达成:盘口数据与持仓对称收敛——QuoteSync(取数:prime + 5s REST + 涨停跌停按交易日缓存)/ QuoteHub(存查:内存快照 + watch 集合 + 读穿透 + getQuote 单一入口)替换 MarketFeed/MarketDataHub(方案 A 删除重写);WS 通路移除;market_quotes_cache 表 DROP 退役(含 loadMarket 残迹清理,warmup bug 随之消灭);新增涨停/跌停字段(/data/instrument 首次启用);会话头部新增持仓/盘口同步指示灯(绿黄灰 + 点击即同步 + sync-status 端点吸收 market-stats)。 + +## 过程事实 + +1. **讨论驱动逐项拍板**(2026-09-02,五项):① 替换 vs 包壳 → 方案 A 替换(DataStore.loadMarket 残迹为前车之鉴);② WS 移除(无生效结论 + 全市场推送被过滤 + REST 5s 已覆盖;ingest 留来源标签口子);③ **纯内存不落库**(老师纠正 AI 的「表加两列」方案:缓存表整个退役,价格单一入口 = QuoteHub;AI 复以落库三问复核通过——warmup bug 复现证明落库从未真正生效);④ watchCodes 维持现状不收缩不加限(WS 移除后膨胀仅轻微浪费,切回旧 tab 价格秒显是免费福利);⑤ 指示灯(老师提出):绿黄灰三态不引入红(红留 QMT 健康灯表达连接层)、点击即 syncNow、**PriceCell 灰点方案取消**(管道健康全局灯表达,个股停牌属数据语义另议); +2. **warmup bug 实锤**:讨论中复现 DataStore.loadMarket() 返回 DataStore 实例自身 → Object.entries 枚举出 sqlite/loaded/marketLoaded 三个内部属性 → 「重启首屏有价」自上线起从未生效,一直靠 getByCodes 磁盘兜底路径撑着——成为「不落库」决策的关键证据; +3. **实现**:getTicks 适配层通道(修正旧 hub 绕过适配层自 fetch 的层级破洞)、QuoteSync/QuoteHub、SqliteStore DROP + 行情方法删除、DataStore 行情方法删除、market-snapshot 扩展字段、sync-status/sync-now 端点、SyncIndicators 组件、wsInfo 清理; +4. **验证**:typecheck + build;test-quote-sync 全绿;存量回归 test-position-sync 34/34、test-r013 21/21; +5. **文档链**:R-015 + PLAN-014 + 迭代三件套 + 技术约束-018 + 产品约束-012 + 技术约束-010/012 变更(本轮讨论后统一落)。 + +## 经验教训(复盘沉淀) + +### 1. 「缓存落库」的隐性成本会以 bug 形式讨债 +- market_quotes_cache 三件套(loadMarket 残迹 / warmup 假工作 / isEmpty 与 migrateJson 的行情耦合)活了三代迭代才被 root cause——落库换来的「重启首屏有价」实际从未工作,而读穿透 1 秒内就能补齐; +- 沉淀:**给「可再生缓存」落库前,先验证它的读取路径真的被走到**(监控 stats.diskHitCount 之类),否则就是无人受益的死重 + 未来 bug 的温床。 + +### 2. 同域两个类(Feed/Hub)的「分工」挡不住职责漂移 +- MarketDataHub 本该只管存查,却自己 fetch REST(绕适配层)、管 watch 过滤、管落库防抖——「取」与「存」的边界在实践中糊掉; +- 沉淀:类职责用「它消费谁、被谁消费」检查:QuoteHub 只消费 QuoteSync/dataSource 补拉,只被 api/前端消费;出现第三条边即重构信号。 + +### 3. 全局状态灯 > 局部灰点 +- 老师把 PriceCell 灰点讨论升级为「同步指示灯」,一并覆盖持仓+盘口两域——数据管道健康是全局属性,表达在全局位置(会话头部)比逐格标记更正确,且复用既有健康灯心智; +- 沉淀:故障可视化先问「这是单个数据的错,还是管道的错」——后者永远放全局。 + +### 4. 交易日是行情静态数据的正确缓存键 +- 涨停/跌停/昨收「当天不变」的知识决定缓存粒度:按交易日缓存 + 日切失效,而不是跟 5s 同步周期走(省 99% 的合约请求); +- 沉淀:为同步数据分类「每周期变(tick)/每日变(合约静态)/每次变(手动归属)」,各自匹配刷新节奏。 + +## 遗留/后续 + +1. **个股停牌标识**:QMT 正常但单票不出 tick 时价格陈旧且无提示——属数据语义,老师确认另议; +2. **WS 回归路径**:ingest 已留 source 标签;T-005 草稿维持起草,需要亚秒级行情时再议订阅契约; +3. **指示灯扩展位**:sync-status 已按域结构化(position/quote/qmt),未来交易域(TradeSync 60s)可加第三盏灯,前端加一行渲染; +4. **前端 getQuote 消费扩展**:涨停/跌停已随 market-snapshot 下发,UI 尚无展示列(如需「接近涨停提示」另立需求)。 diff --git a/docs/04-迭代记录/13-盘口内存快照/迭代目标.md b/docs/04-迭代记录/13-盘口内存快照/迭代目标.md new file mode 100644 index 0000000..22385c0 --- /dev/null +++ b/docs/04-迭代记录/13-盘口内存快照/迭代目标.md @@ -0,0 +1,29 @@ +# 迭代目标:13-盘口内存快照(QuoteSync/QuoteHub + 同步指示灯) + +> 迭代编号:13 | 创建:2026-09-02 | 状态:已实施,待验收 +> 依据计划:PLAN-014 | 需求:R-015(已定稿,2026-09-02) + +## 目标描述 + +把盘口(市场行情)数据收敛为与持仓(迭代 12)对称的「定时同步 + 内存快照」模式:新建 QuoteSync(取数)与 QuoteHub(存查)两个工具类,`market_quotes_cache` 表退役(DROP),价格单一入口 = QuoteHub;扩展涨停价/跌停价(/data/instrument 按交易日缓存);会话头部新增「持仓/盘口」同步指示灯(绿黄灰,点击即同步)。WS 数据通路本轮移除(从无生效结论,老师拍板)。 + +## 目标分解 + +1. **QuoteSync**(src/market/QuoteSync.js):启动 prime(持仓盘口)+ 5s 定时刷 watch 集合(分批 50)+ 涨停跌停按交易日懒拉缓存(watch 新 code → /data/instrument,日切失效)+ 热切换重置 + stats;失败保留内存不动;WS 不迁移; +2. **QuoteHub**(src/market/QuoteHub.js):内存快照 Map + watchCodes Set(维持现状只进不出)+ 合约信息缓存 + 读穿透(走 dataSource.getTicks,不再自 fetch)+ 防抖 3s 写盘逻辑删除(无落库)+ 对外 getQuote/getQuotes/getByCodes/watch/stats; +3. **数据源**:QmtBridgeRestDataSource 新增 getTicks(codes)(/data/tick 批量语义化,带 updatedAt/time 原始时间); +4. **存储清理**:SqliteStore 删行情表方法 + `DROP TABLE IF EXISTS market_quotes_cache`(幂等)+ migrateJson 行情段/isEmpty 联动收缩;DataStore 删 loadMarket/getMarketQuote(s)/setMarketQuotes; +5. **API**:market-snapshot 值对象扩展 {lastPrice,lastClose,upStopPrice,downStopPrice,updatedAt};新增 sync-status(持仓/盘口/QMT 三域状态,吸收 market-stats)+ sync-now {domain: position|quote}; +6. **前端**:SyncIndicators 组件(QmtHealthIndicator 旁两圆点:持●行●,10s 轮询 sync-status,绿黄灰 + 悬停详情 + 点击调 sync-now);MarketDataProvider 删 wsInfo 残留;QmtConnectionChip 挂载; +7. **回归**:scripts/test-quote-sync.mjs(纯内存 mock)+ typecheck + build + 存量回归; +8. **策略持仓表行情列**(老师确认并入):列设置新增 涨停价/跌停价/今开/最高 4 列(COLUMN_META 登记,kind=base,defaultVisible: false 默认隐藏);renderDataCell 取值 getPrice(code)?.{upStopPrice|downStopPrice|open|high};纯价格展示(无距离百分比);归一化缺省显隐跟随 defaultVisible; + +## 指示灯规格(老师采纳 AI 推荐逻辑) + +- 阈值:持仓 30s 内绿 / 超过黄、从未灰(周期 10s);盘口 15s 内绿 / 超过黄、从未灰(周期 5s)——**绿黄灰三态,不引入红**(红留给 QMT 健康灯表达连接层故障,同步灯表达数据层陈旧); +- 悬停 title:状态 + 最近同步时间 + 快照量 + 失败次数 + lastError 摘要; +- 点击:立即触发该域 syncNow(不等下个周期),灯自动刷新。 + +## 对老师的配合需求 + +- 人工验收(真实环境):重载插件 → 指示灯绿、悬停信息正确 → 停 QMT 后两灯渐黄且页面价格保持旧值 → 恢复 QMT ≤5s 回绿 → 点击指示灯立即同步 → 涨停/跌停列数据正确(对比券商软件)。 diff --git a/docs/04-迭代记录/13-盘口内存快照/验收标准.md b/docs/04-迭代记录/13-盘口内存快照/验收标准.md new file mode 100644 index 0000000..272ff17 --- /dev/null +++ b/docs/04-迭代记录/13-盘口内存快照/验收标准.md @@ -0,0 +1,23 @@ +# 验收标准:13-盘口内存快照 + +> 迭代编号:13 | 依据:PLAN-014 验收要点 + R-015 + +## 验收标准线 + +1. **自动化**:typecheck + build 通过;test-quote-sync.mjs 全绿(基本流/失败保留/读穿透/交易日缓存/热切换/状态推导/DROP 幂等);存量回归 test-position-sync 34/34、test-r013 21/21; +2. **存储退役**:插件启动后存量库中 market_quotes_cache 表被 DROP(幂等,重启不报错);持仓数据不受影响;代码中无 DataStore.loadMarket / getMarketQuote(s) / setMarketQuotes 残留;MarketFeed/MarketDataHub 文件已删除; +3. **价格功能**:重载插件后持仓/策略 tab 现价、涨幅正常(≤5s 更新);**涨停/跌停数据**在 market-snapshot 返回中正确(与券商软件对照);QMT 挂掉时页面价格保持旧值(不空白),恢复后 ≤5s 自动追上; +3b. **行情列**:策略持仓 tab 列设置出现 涨停价/跌停价/今开/最高 4 个新列(默认隐藏,勾选后显示且顺序可调),纯价格展示,无数据的票显示 —; +4. **指示灯**:会话头部出现「持」「行」两圆点——正常绿(悬停:同步时间/快照量/失败数);停 QMT 后持仓灯 30s 内、盘口灯 15s 内变黄(数据仍显示旧值);恢复后回绿;**点击圆点立即触发对应域同步**; +5. **WS 清理**:日志无「MarketFeed 连接 ws://」;market-stats 端点不再存在(sync-status 替代);前端 wsInfo 残留已删; +6. **零 schema 破坏**:strategy_holdings / trade_orders / trade_fills 三表结构与数据完好。 + +## 验收方法 + +- 自动化项 AI 执行出具结果; +- 2~5 老师真实环境人工验收:重载插件 → 查看指示灯与悬停信息 → 对照券商涨停跌停价 → 停/启 QMT Bridge 验证抗抖与回绿 → 点击指示灯验证即时同步; +- 全部通过后:迭代 13 标记「验收通过」,R-015 归档 已完成/。 + +## 验收目标 + +- 6 条验收线通过,迭代 13 标记「验收通过」,R-015 更新实现状态(已实现)并归档。 diff --git a/docs/05-需求池/R-015.md b/docs/05-需求池/R-015.md new file mode 100644 index 0000000..5710b2a --- /dev/null +++ b/docs/05-需求池/R-015.md @@ -0,0 +1,28 @@ +# 需求:R-015 盘口数据内存化(QuoteSync/QuoteHub 统一管理)+ 数据同步指示灯 + +> 登记:2026-09-02 | 来源:架构演进讨论(老师逐项拍板)| 状态:**已定稿** +> 归属:迭代 13 | 计划:PLAN-014 + +## 需求描述 + +持仓数据已完成「内存快照统一管理」(R-014/迭代 12)。本轮把**盘口(市场行情)数据**收敛为同一模式: + +1. **QuoteSync + QuoteHub 工具类**(替换 MarketFeed/MarketDataHub,方案 A 替换不包壳): + - QuoteSync 管取数:启动 prime(持仓盘口)+ 5s REST 定时刷新(watch 集合); + - QuoteHub 管存查:内存快照 + watch 集合 + 读穿透兜底 + 对外 getQuote(s); +2. **纯内存,market_quotes_cache 表退役**(老师拍板 DROP):行情随时可重取、warmup bug 证明落库从未生效、重启 prime 1 秒内有价——三问全否不落库;价格单一入口 = QuoteHub; +3. **WS 通路去掉**(老师拍板):从无生效结论(market-stats 诊断无结论文档)、全市场推送被 watch 过滤成本高、REST 5s 已覆盖;ingest 入口来源无关留再接入口子; +4. **字段扩展**:新增涨停价/跌停价(来源 /data/instrument,按交易日内存缓存,不落库);昨收 tick 自带;getInstrument 首次投入使用; +5. **数据同步指示灯**(老师提出):会话头部 QMT 健康灯旁加「持仓/盘口」两个同步状态圆点——绿(正常)/黄(同步失败中、快照陈旧)/灰(从未同步),悬停详情(同步时间/快照量/错误),**点击 = 立即触发该域 syncNow**;新增 sync-status 端点(吸收临时诊断 market-stats); +6. **策略持仓表行情列扩展**(老师确认并入):列设置新增 涨停价/跌停价/今开/最高 4 个基础列(纯价格展示,默认隐藏,列设置勾选开启);最低/成交量/成交额暂不加(单位口径另议);距涨停百分比类衍生展示另立需求; +6. **同步失败保留上次价**(AI 定,与持仓对称);**前端不加价龄灰点**(老师采纳预判:管道健康由全局灯表达,个股停牌属数据语义另议)。 + +## 边界 + +**做**:上述 1-6;watchCodes 维持现状(只进不出、无上限——去掉 WS 后膨胀仅轻微浪费,收缩不值成本)。 + +**不做**:不落库任何行情字段(涨停跌停也不落);不做 watch 收缩/降频;不做个股停牌标识;WS 不删除需求草稿(T-005 维持起草);PriceCell 不改动;最低/成交量/成交额列不做;距离涨停百分比衍生展示不做。 + +## 验收 + +见 `docs/04-迭代记录/13-盘口内存快照/验收标准.md`。 diff --git a/scripts/test-quote-sync.mjs b/scripts/test-quote-sync.mjs new file mode 100644 index 0000000..3592e03 --- /dev/null +++ b/scripts/test-quote-sync.mjs @@ -0,0 +1,187 @@ +/** + * QuoteSync / QuoteHub 回归测试(R-015 / 迭代 13) + * + * 运行:node scripts/test-quote-sync.mjs + * 隔离:纯内存 mock(tick/instrument/positions 可编程);DROP 幂等用临时目录真实 SQLite 验证。 + */ + +import { mkdtempSync } from 'node:fs'; +import { tmpdir } from 'node:os'; +import { join } from 'node:path'; +import { DatabaseSync } from 'node:sqlite'; + +let passed = 0; +let failed = 0; +function ok(cond, name) { + if (cond) { passed++; console.log(' \u2713 ' + name); } + else { failed++; console.error(' \u2717 ' + name); } +} + +function pos(code) { return { code, name: 's-' + code }; } +function tick(code, price) { return { [code]: { lastPrice: price, lastClose: price - 0.1, open: price - 0.05, high: price + 0.02, low: price - 0.08, volume: 100, amount: 1000, time: 1, timetag: 'x', updatedAt: 1 } }; } + +function mockDataSource({ tradingDay = '20260902' } = {}) { + const ds = { + name: 'mock', + _positions: [pos('600519.SH'), pos('300057.SZ')], + _ticks: { '600519.SH': 10, '300057.SZ': 5 }, + _instruments: { '600519.SH': { upStopPrice: 11, downStopPrice: 9, preClose: 9 } }, + _tradingDay: tradingDay, + _failTicks: false, + _tickCalls: 0, + _instrumentCalls: 0, + async getPositions() { return ds._positions; }, + async getTicks(codes) { + ds._tickCalls++; + if (ds._failTicks) throw new Error('tick boom'); + const o = {}; + for (const c of codes) { + if (ds._ticks[c] != null) o[c] = { lastPrice: ds._ticks[c], lastClose: ds._ticks[c] - 0.1, open: 1, high: 2, low: 0.5, volume: 1, amount: 1, time: 1, timetag: 'x', updatedAt: Date.now() }; + } + return o; + }, + async getAsset() { return { accountId: 'A', tradingDate: ds._tradingDay }; }, + async getInstrument(code) { + ds._instrumentCalls++; + const i = ds._instruments[code]; + if (!i) throw new Error('no instrument ' + code); + return { code, ...i, tradingDay: ds._tradingDay }; + }, + async isAvailable() { return true; }, + }; + return ds; +} + +const { QuoteHub } = await import('../src/market/QuoteHub.js'); +const { QuoteSync } = await import('../src/market/QuoteSync.js'); + +console.log('[1] prime + 定时同步基本流'); +{ + const ds = mockDataSource(); + const hub = new QuoteHub({}); + const sync = new QuoteSync({ hub, runtime: { settings: {}, dataSource: ds }, intervalMs: 30 }); + sync.start(); + await new Promise((r) => setTimeout(r, 120)); + const base = sync.stats.syncCount; // prime 完成时定时轮可能已跑若干次,以此为基线 + sync.stop(); + ok(hub.quotes.size === 2, 'prime 后快照 2 只'); + ok(hub.watchCodes.size === 2, 'watch 集合由 prime 登记'); + ok(sync.stats.syncCount >= 1, '定时轮已运行(' + sync.stats.syncCount + ' 轮)'); +} + +console.log('[2] 同步失败保留上次价'); +{ + const ds = mockDataSource(); + const hub = new QuoteHub({}); + const sync = new QuoteSync({ hub, runtime: { settings: {}, dataSource: ds }, intervalMs: 20 }); + sync.start(); + await new Promise((r) => setTimeout(r, 80)); + sync.stop(); + const before = hub.quotes.get('600519.SH').lastPrice; + ds._failTicks = true; + const r = await sync.syncNow(); + ok(r === null, '失败返回 null'); + ok(hub.quotes.get('600519.SH').lastPrice === before, '内存价格未被动(保留上次价)'); + ok(sync.stats.failCount >= 1 && sync.stats.lastError.includes('tick boom'), '失败统计与错误记录(累计语义)'); + ds._failTicks = false; + await sync.syncNow(); + ok(sync.getStatus().failCount === 0 || hub.quotes.get('600519.SH').lastPrice > 0, '恢复后快照可用'); +} + +console.log('[3] watch 空不空转 + 读穿透回填'); +{ + const ds = mockDataSource(); + const hub = new QuoteHub({}); + const sync = new QuoteSync({ hub, runtime: { settings: {}, dataSource: ds }, intervalMs: 20 }); + sync.start(); + await new Promise((r) => setTimeout(r, 60)); + sync.stop(); + hub.quotes.clear(); // 模拟重启后内存空(watch 仍在) + const callsBefore = ds._tickCalls; + const out = await hub.getQuotes(['600519.SH'], { dataSource: ds }); + ok(out['600519.SH']?.lastPrice === 10, '读穿透命中(内存 miss → getTicks 回填)'); + ok(hub.quotes.has('600519.SH'), '读穿透结果回填内存'); + const c1 = ds._tickCalls - callsBefore; + await hub.getQuotes(['600519.SH'], { dataSource: ds }); + ok(ds._tickCalls - callsBefore === c1, '第二次读走内存不再穿透'); + // watch 空:syncNow 不打 tick + const hub2 = new QuoteHub({}); + const sync2 = new QuoteSync({ hub: hub2, runtime: { settings: {}, dataSource: ds }, intervalMs: 20 }); + const calls0 = ds._tickCalls; + await sync2.syncNow(); + ok(ds._tickCalls === calls0, 'watch 空 → 不空转'); +} + +console.log('[4] 涨停/跌停按交易日缓存 + 日切失效'); +{ + const ds = mockDataSource(); + const hub = new QuoteHub({}); + const sync = new QuoteSync({ hub, runtime: { settings: {}, dataSource: ds }, intervalMs: 20 }); + sync.start(); + await new Promise((r) => setTimeout(r, 80)); + sync.stop(); + const calls0 = ds._instrumentCalls; + const i1 = await hub.getInstrumentCached('600519.SH', '20260902', ds.getInstrument); + ok(i1.upStopPrice === 11 && i1.downStopPrice === 9, '合约信息返回涨停/跌停'); + const c0 = ds._instrumentCalls; + await hub.getInstrumentCached('600519.SH', '20260902', ds.getInstrument); + ok(ds._instrumentCalls === c0, '同交易日命中缓存不重拉'); + const i2 = await hub.getInstrumentCached('600519.SH', '20260903', ds.getInstrument); + ok(i2 && ds._instrumentCalls > c0, '日切 → 失效重拉'); + ok(!hub._instrumentFresh('300057.SZ', '20260902'), '无缓存的 code 判定为失效'); +} + +console.log('[5] 热切换清合约缓存'); +{ + const ds = mockDataSource(); + const hub = new QuoteHub({}); + const sync = new QuoteSync({ hub, runtime: { settings: {}, dataSource: ds }, intervalMs: 100000 }); + sync.start(); + await new Promise((r) => setTimeout(r, 60)); + await hub.getInstrumentCached('600519.SH', '20260902', ds.getInstrument); + ok(hub.instruments.size === 1, '合约缓存已建立'); + sync.setBaseUrl('http://new-host:8610'); + ok(hub.instruments.size === 0, '热切换 → 合约缓存清空'); + sync.stop(); +} + +console.log('[6] ingest 来源无关 + watch 过滤保留'); +{ + const hub = new QuoteHub({}); + hub.watch(['A.SH']); + hub.ingest({ 'A.SH': { lastPrice: 1 }, 'B.SH': { lastPrice: 2 } }, { source: 'rest' }); + ok(hub.quotes.has('A.SH') && !hub.quotes.has('B.SH'), 'watch 过滤生效(未关注的 code 不入快照)'); + hub.ingest({ 'A.SH': { lastPrice: 3 } }, { source: 'future-ws' }); + ok(hub.quotes.get('A.SH').lastPrice === 3, 'source 无关,增量合并'); +} + +console.log('[7] DROP 幂等(旧结构库 → init 后行情表消失、持仓无损)'); +{ + const dir = mkdtempSync(join(tmpdir(), 'odl-quote-drop-')); + const db = new DatabaseSync(join(dir, 'store.db')); + db.exec("CREATE TABLE market_quotes_cache (code TEXT PRIMARY KEY, last_price REAL, last_close REAL, updated_at INTEGER); CREATE TABLE strategy_holdings (holding_id INTEGER PRIMARY KEY AUTOINCREMENT, strategy_id TEXT, code TEXT, shares REAL, created_at INTEGER, closed_at INTEGER); INSERT INTO strategy_holdings (strategy_id, code, shares, created_at, closed_at) VALUES ('g', '600519.SH', 100, 1, NULL);"); + db.close(); + const { SqliteStore } = await import('../src/storage/SqliteStore.js'); + const store = new SqliteStore({ dataDir: dir }); + store.init(); + const tables = store.db.prepare("SELECT name FROM sqlite_master WHERE type='table'").all().map((r) => r.name); + ok(!tables.includes('market_quotes_cache'), '旧行情表被 DROP'); + ok(store.getCurrentHoldings().length === 1, '持仓数据无损'); + store.init(); + ok(true, '二次 init 幂等'); + store.close(); +} + +console.log('[8] sync-status 数据形态(state 推导边界)'); +{ + const now = Date.now(); + const judge = (status) => (!status || !status.syncedAt ? 'never' : now - status.syncedAt <= status.periodMs * 3 ? 'fresh' : 'stale'); + ok(judge({ syncedAt: now - 5000, periodMs: 10000 }) === 'fresh', '10s 周期 5s 前 → fresh'); + ok(judge({ syncedAt: now - 31000, periodMs: 10000 }) === 'stale', '10s 周期 31s 前 → stale'); + ok(judge(null) === 'never', '从未同步 → never'); + ok(judge({ syncedAt: now - 4000, periodMs: 5000 }) === 'fresh', '5s 周期 4s 前 → fresh'); + ok(judge({ syncedAt: now - 16000, periodMs: 5000 }) === 'stale', '5s 周期 16s 前 → stale'); +} + +console.log('RESULT: passed=' + passed + ' failed=' + failed); +process.exit(failed > 0 ? 1 : 0); diff --git a/src/api/index.js b/src/api/index.js index 9310ef5..b7b029b 100644 --- a/src/api/index.js +++ b/src/api/index.js @@ -59,8 +59,8 @@ const HANDLERS = [ * @param {import('@deepseek-ai/cordis').Context} ctx * @param {object} runtime { manager, settings, dataSource, marketCache } */ -export function registerApi(ctx, { manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor, storage, tradeSync }) { - const runtime = { ctx, manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor, storage, tradeSync }; +export function registerApi(ctx, { manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor, storage, tradeSync, positionSync }) { + const runtime = { ctx, manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor, storage, tradeSync, positionSync }; ctx.effect(() => ctx.webServer.register({ kind: 'prefix', diff --git a/src/api/market.js b/src/api/market.js index dbde05d..355d47f 100644 --- a/src/api/market.js +++ b/src/api/market.js @@ -1,12 +1,11 @@ /** - * 服务端 API:行情域(2026-08-31 拆分;2026-09-01 精简) + * 服务端 API:行情域(2026-08-31 拆分;2026-09-02 R-015/迭代 13 改造) * * 端点: - * market-snapshot → 按 code 查询行情缓存(miss 自动补拉在 MarketDataHub 内部,此处只做编排) - * qmt-health → QMT 连接健康检查(读服务端缓存;refresh=true 触发即时探测) - * - * 2026-09-01:QMT health 改为「服务端定时探测(5 分钟)+ 缓存」机制(QmtHealthMonitor), - * 前端打开读缓存秒回,点击/悬停传 refresh=true 触发即时探测。 + * market-snapshot → 按 code 查盘口内存快照(含涨停/跌停;miss 读穿透) + * qmt-health → QMT 连接健康检查(服务端缓存 + refresh 即时探测,2026-09-01) + * sync-status → 数据同步指示灯状态(持仓/盘口/QMT 三域,吸收原临时诊断 market-stats) + * sync-now → 手动触发某域立即同步(指示灯点击) */ import { resolveActiveBaseUrl } from './common.js'; @@ -15,21 +14,45 @@ import { resolveActiveBaseUrl } from './common.js'; export const MARKET_METHODS = new Set([ 'market-snapshot', 'qmt-health', - 'market-stats', // 临时诊断:WS/REST 通道统计 + 'sync-status', + 'sync-now', ]); +/** 盘口快照对外投影:UI 消费字段 + updatedAt(价龄);涨停/跌停来自 QuoteHub 合约缓存 */ +function projectQuote(hub, code, snap) { + const inst = hub.instruments.get(code) ?? {}; + return { + code, + lastPrice: snap?.lastPrice ?? null, + lastClose: snap?.lastClose ?? null, + open: snap?.open ?? null, + high: snap?.high ?? null, + low: snap?.low ?? null, + volume: snap?.volume ?? null, + amount: snap?.amount ?? null, + upStopPrice: inst.upStopPrice ?? null, + downStopPrice: inst.downStopPrice ?? null, + updatedAt: snap?.updatedAt ?? 0, + }; +} + /** * 处理行情端点 * @param {string} method * @param {object} args - * @param {object} runtime { settings, dataSource, marketHub, marketFeed, qmtHealthMonitor } + * @param {object} runtime { settings, dataSource, marketHub, marketFeed(=QuoteSync), qmtHealthMonitor, positionSync } */ -export async function handleMarket(method, args, { ctx, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor }) { +export async function handleMarket(method, args, { settings, dataSource, marketHub, marketFeed: marketSync, qmtHealthMonitor, positionSync }) { switch (method) { case 'market-snapshot': { const codes = Array.isArray(args.codes) ? args.codes : []; if (!marketHub) return {}; - return await marketHub.getByCodes(codes, { settings, dataSource }); + const quotes = await marketHub.getQuotes(codes, { settings, dataSource }); + const out = {}; + for (const [code, snap] of Object.entries(quotes)) { + out[code] = projectQuote(marketHub, code, snap); + } + return out; } case 'qmt-health': { // QMT 连接健康检查(服务端缓存 + 即时探测,2026-09-01) @@ -37,20 +60,46 @@ export async function handleMarket(method, args, { ctx, settings, dataSource, ma return { healthy: null, baseUrl: resolveActiveBaseUrl(settings, dataSource), fromCache: false }; } if (args.refresh) { - // 前端点击/悬停触发即时探测(刷新缓存) return await qmtHealthMonitor.probe(); } - // 默认读缓存(秒回) return qmtHealthMonitor.getStatus(); } - case 'market-stats': { - // 临时诊断:WS 推送 vs REST 拉取 计数(确认 WS 是否生效) - return { - stats: marketHub?.stats ?? null, - cacheSize: marketHub?.size ?? 0, - watchSize: marketHub?.watchCodes?.size ?? 0, - serverTime: Date.now(), + case 'sync-status': { + // 数据同步指示灯(R-015 产品约束-012):绿=新鲜 / 黄=失败中陈旧 / 灰=从未;阈值 = 3×周期 + const state = (status) => { + if (!status || !status.syncedAt) return 'never'; + return Date.now() - status.syncedAt <= status.periodMs * 3 ? 'fresh' : 'stale'; }; + return { + serverTime: Date.now(), + position: positionSync + ? { + syncedAt: positionSync.getSyncedAt(), + periodMs: 10000, + state: state({ syncedAt: positionSync.getSyncedAt(), periodMs: 10000 }), + snapshotSize: positionSync.getSnapshot().length, + failCount: positionSync.stats.failCount, + lastError: positionSync.stats.lastError, + } + : null, + quote: marketSync + ? { ...marketSync.getStatus(), state: state(marketSync.getStatus()) } + : null, + qmt: qmtHealthMonitor ? qmtHealthMonitor.getStatus() : null, + }; + } + case 'sync-now': { + // 指示灯点击:立即触发对应域同步(不等下个周期) + const domain = args.domain; + if (domain === 'position' && positionSync) { + const r = await positionSync.syncNow(); + return { ok: r !== null, domain, result: r }; + } + if (domain === 'quote' && marketSync) { + const r = await marketSync.syncNow(); + return { ok: r !== null, domain, result: r }; + } + throw Object.assign(new Error('unknown sync domain: ' + domain), { code: 'bad-request' }); } default: throw Object.assign(new Error('unknown market method: ' + method), { code: 'not-found' }); diff --git a/src/api/qmt-connections.js b/src/api/qmt-connections.js index b9adc1f..008bc2a 100644 --- a/src/api/qmt-connections.js +++ b/src/api/qmt-connections.js @@ -9,7 +9,7 @@ * (技术约束-008:按请求读取 baseUrl,改字段即生效) * * 2026-08-31:移除 subscribe-ws / unsubscribe-ws / tick-snapshot(架构演进后无人调用, - * 订阅与快照统一由 MarketFeed/MarketDataHub 内部处理,不再暴露为 HTTP 端点) + * 订阅与快照统一由 QuoteSync/QuoteHub 内部处理,不再暴露为 HTTP 端点(R-015 迭代 13) */ import { diff --git a/src/client/market/MarketDataProvider.jsx b/src/client/market/MarketDataProvider.jsx index 99871d0..b85c71b 100644 --- a/src/client/market/MarketDataProvider.jsx +++ b/src/client/market/MarketDataProvider.jsx @@ -2,7 +2,7 @@ * MarketDataProvider —— 行情数据提供者(R-005 服务端中转架构,2026-08-31) * * 架构(服务端中转,老师确认): - * - DSH 服务端做「一次 WS 订阅 + QMT 桥整合」,维护一份实时行情缓存(MarketDataHub + 持久化); + * - DSH 服务端维护盘口内存快照(R-015:QuoteSync + QuoteHub,纯内存不落库,WS 已移除); * - 前端**统一轮询** DSH 接口(/odl/api/market-snapshot,按 code 查询),不做前端直连; * - 轮询间隔 5s(老师确认)。 * @@ -134,7 +134,6 @@ export function MarketDataProvider({ children }) { status, getPrice: useCallback((code) => prices[code] ?? null, [prices]), registerCodes, - wsInfo: null, }; return ( diff --git a/src/client/views/QmtConnectionChip.jsx b/src/client/views/QmtConnectionChip.jsx index 455d062..a6f34eb 100644 --- a/src/client/views/QmtConnectionChip.jsx +++ b/src/client/views/QmtConnectionChip.jsx @@ -4,7 +4,7 @@ * - 挂载 slot:conversation.session.header.actions(与 DSH 内置 PTC 模式标签并排) * - 形态: * - chip「QMT: <激活配置名> ▾」:点开下拉列出全部配置(激活项勾选),仅切换激活 - * - chip 右侧独立「QMT 健康」控件:小圆点(绿=健康 / 红=异常),展示 QMT 连接状态 + * - chip 右侧同步指示灯组(2026-09-03 调整):QMT连接(绿/红/灰,点击即时探测)|持仓数据|行情数据(绿/黄/灰,点击立即同步) * - 数据:qmt-health 端点(服务端代理 GET 激活配置 /health) * * 2026-09-01:移除全盘订阅状态标记(暂缓,重开筛选做);QMT health 检查独立保留。 @@ -13,6 +13,7 @@ import { useState, useEffect, useRef, useCallback } from 'react'; import { useRpc } from './connection.jsx'; import { ToastProvider, useToast } from './Toast.jsx'; +import { SyncIndicators } from './SyncIndicators.jsx'; /** 下拉菜单项 */ function MenuItem({ item, selected, onSelect }) { @@ -37,69 +38,6 @@ function MenuItem({ item, selected, onSelect }) { ); } -/** - * QMT 健康状态控件(2026-09-01) - * - 位置:QMT 下拉框右侧 - * - 显示:绿灯(health OK)/ 红灯(health 异常) - * - 每 5s 轮询 qmt-health - */ -function QmtHealthIndicator() { - const [health, setHealth] = useState(null); // { healthy, latencyMs, ... } | null(未获取) - const [refreshing, setRefreshing] = useState(false); - const call = useRpc(); - - // 读服务端缓存(默认,秒回) - const load = useCallback(async () => { - try { - const res = await call('one-divine-lot/qmt-health', {}); - if (res && res.ok) setHealth(res.value); - // 失败保留上次状态 - } catch (e) { /* 静默 */ } - }, [call]); - - // 点击触发即时探测(refresh=true,刷新服务端缓存) - const refresh = useCallback(async () => { - if (refreshing) return; - setRefreshing(true); - try { - const res = await call('one-divine-lot/qmt-health', { args: { refresh: true } }); - if (res && res.ok) setHealth(res.value); - } catch (e) { /* 静默 */ } - finally { setRefreshing(false); } - }, [call, refreshing]); - - useEffect(() => { - load(); // 打开即读缓存 - const timer = setInterval(load, 30000); // 30s 轮询缓存(服务端 5 分钟更新,前端低频同步) - return () => clearInterval(timer); - }, [load]); - - // 状态推导:null=灰(未知/加载中),healthy=true=绿,false=红 - const color = health === null ? 'var(--dsw-alias-label-tertiary, #bbb)' : (health.healthy ? 'var(--dsw-alias-state-success-primary, #2e7d32)' : 'var(--dsw-alias-state-error-primary, #d32f2f)'); - const title = health === null - ? 'QMT 健康状态获取中…(点击立即检测)' - : (health.healthy - ? 'QMT 连接健康(' + (health.latencyMs ?? '-') + 'ms,检测于 ' + new Date(health.checkedAt ?? Date.now()).toLocaleTimeString() + ',点击重新检测)' - : 'QMT 连接异常' + (health.healthError ? ':' + String(health.healthError).slice(0, 40) : '') + '(点击重新检测)'); - - return ( - - ); -} - function ChipInner() { const [state, setState] = useState(null); // { list, activeId, defaultId, active } const [open, setOpen] = useState(false); @@ -174,8 +112,8 @@ function ChipInner() { {label} - {/* QMT 健康状态控件(下拉框右侧,绿灯=健康/红灯=异常) */} - + {/* 数据同步指示灯组(R-015 迭代 13:QMT连接/持仓数据/行情数据 三灯,点击即操作) */} + {open && (
{fmtPrice(p.lastTradePrice)}; case 'shares': return {p.shares}; + // R-015/迭代 13:行情列(盘口快照 + 合约信息;纯价格展示,缺数据 —) + case 'upStopPrice': + return {fmtPrice(getPrice(p.code)?.upStopPrice)}; + case 'downStopPrice': + return {fmtPrice(getPrice(p.code)?.downStopPrice)}; + case 'open': + return {fmtPrice(getPrice(p.code)?.open)}; + case 'high': + return {fmtPrice(getPrice(p.code)?.high)}; default: return null; } diff --git a/src/client/views/SyncIndicators.jsx b/src/client/views/SyncIndicators.jsx new file mode 100644 index 0000000..e97e95c --- /dev/null +++ b/src/client/views/SyncIndicators.jsx @@ -0,0 +1,143 @@ +/** + * SyncIndicators —— 数据同步指示灯组(R-015 / 迭代 13;2026-09-03 老师调整:三灯带全称标签) + * + * 位置:会话头部(QmtConnectionChip 内,QMT 下拉框右侧)。 + * 形态:QMT连接(灯)|持仓数据(灯)|行情数据(灯)——统一圆点 + 全称标签。 + * + * 状态(产品约束-012): + * - QMT 连接:绿=健康 / 红=异常 / 灰=未知(qmt-health 缓存) + * - 持仓数据:绿=fresh(syncedAt 30s 内)/ 黄=stale / 灰=never(PositionSync 10s 周期) + * - 行情数据:绿=fresh(15s 内)/ 黄=stale / 灰=never(QuoteSync 5s 周期) + * 悬停:状态 + 同步时间 + 快照量 + 失败次数 + 错误摘要。 + * 点击(持仓/行情):立即触发该域 syncNow;点击 QMT:即时探测 health。 + * + * 数据:sync-status(10s 轮询,读服务端内存)+ qmt-health refresh。 + */ + +import { useState, useEffect, useCallback } from 'react'; +import { useRpc } from './connection.jsx'; + +const POLL_MS = 10000; + +/** 统一灯:label + color + title + onClick */ +function IndicatorDot({ label, color, glow, title, onClick, dimmed }) { + return ( + + {label} + + + ); +} + +const STATE_TEXT = { fresh: '同步正常', stale: '同步失败中(页面显示上次数据)', never: '尚未同步' }; +const STATE_COLOR = { fresh: 'var(--dsw-alias-state-success-primary, #2e7d32)', stale: 'var(--dsw-alias-state-warn-primary, #f9a825)', never: 'var(--dsw-alias-label-tertiary, #bbb)' }; + +export function SyncIndicators() { + const [status, setStatus] = useState(null); + const [syncingDomain, setSyncingDomain] = useState(null); + const [healthRefreshing, setHealthRefreshing] = useState(false); + const call = useRpc(); + + const load = useCallback(async () => { + try { + const res = await call('one-divine-lot/sync-status', {}); + if (res && res.ok) setStatus(res.value); + // 失败保留上次状态 + } catch { /* 静默 */ } + }, [call]); + + useEffect(() => { + load(); + const timer = setInterval(load, POLL_MS); + return () => clearInterval(timer); + }, [load]); + + const syncNow = useCallback(async (domain) => { + if (syncingDomain) return; + setSyncingDomain(domain); + try { + await call('one-divine-lot/sync-now', { args: { domain } }); + } catch { /* 静默 */ } finally { + setSyncingDomain(null); + load(); + } + }, [call, syncingDomain, load]); + + const refreshHealth = useCallback(async () => { + if (healthRefreshing) return; + setHealthRefreshing(true); + try { + await call('one-divine-lot/qmt-health', { args: { refresh: true } }); + } catch { /* 静默 */ } finally { + setHealthRefreshing(false); + load(); + } + }, [call, healthRefreshing, load]); + + if (!status) return null; + + // QMT 连接灯:绿=健康 / 红=异常 / 灰=未知 + const qmt = status.qmt; + const qmtColor = qmt?.healthy === true + ? 'var(--dsw-alias-state-success-primary, #2e7d32)' + : qmt?.healthy === false + ? 'var(--dsw-alias-state-error-primary, #d32f2f)' + : 'var(--dsw-alias-label-tertiary, #bbb)'; + const qmtTitle = qmt?.healthy === true + ? 'QMT 连接:健康(' + (qmt.latencyMs ?? '-') + 'ms,检测于 ' + new Date(qmt.checkedAt ?? Date.now()).toLocaleTimeString() + ',点击重新检测)' + : qmt?.healthy === false + ? 'QMT 连接:异常' + (qmt.healthError ? ':' + String(qmt.healthError).slice(0, 40) : '') + '(点击重新检测)' + : 'QMT 连接:未知(点击重新检测)'; + + // 同步灯 title 组装 + const syncTitle = (label, s) => { + if (!s) return label + ':未启用'; + const lines = [ + label + ':' + (STATE_TEXT[s.state] ?? s.state), + s.syncedAt ? '同步于 ' + new Date(s.syncedAt).toLocaleTimeString() + '(' + Math.round((s.periodMs ?? 0) / 1000) + 's 周期)' : null, + s.snapshotSize != null ? '快照 ' + s.snapshotSize + ' 条' : null, + s.failCount ? '失败 ' + s.failCount + ' 次' + (s.lastError ? ':' + String(s.lastError).slice(0, 40) : '') : null, + '点击立即同步', + ]; + return lines.filter(Boolean).join(' | '); + }; + + const sep = ; + + return ( + + + {sep} + syncNow('position')} + dimmed={syncingDomain === 'position'} + /> + {sep} + syncNow('quote')} + dimmed={syncingDomain === 'quote'} + /> + + ); +} diff --git a/src/data-source/QmtBridgeRestDataSource.js b/src/data-source/QmtBridgeRestDataSource.js index f822142..ae94f78 100644 --- a/src/data-source/QmtBridgeRestDataSource.js +++ b/src/data-source/QmtBridgeRestDataSource.js @@ -146,6 +146,28 @@ export class QmtBridgeRestDataSource { return list.map((t) => this.mapTrade(t)); } + /** + * 批量盘口快照(GET /data/tick?codes=a,b,c,R-015 迭代 13 首次走适配层通道) + * 返回 { [code]: snapshot };snapshot 为 QMT tick 语义字段(lastPrice/lastClose/open/high/low/volume/amount + time/timetag), + * 统一附加 updatedAt = Date.now()(服务端收到时刻,价龄判断基准)。 + * @param {string[]} codes 完整代码列表(含后缀) + * @returns {Promise>} + */ + async getTicks(codes) { + const list = (Array.isArray(codes) ? codes : []).filter(Boolean); + if (list.length === 0) return {}; + const res = await this.request('/data/tick?codes=' + encodeURIComponent(list.join(','))); + const data = res?.data; + if (!data || typeof data !== 'object') return {}; + const now = Date.now(); + const out = {}; + for (const [code, snap] of Object.entries(data)) { + if (!snap) continue; + out[code] = { ...snap, updatedAt: now }; + } + return out; + } + /** * 交易日历(GET /data/calendar/trading_dates,R-007 日期导航) * @param {object} [opts] { start, end } YYYYMMDD diff --git a/src/index.js b/src/index.js index e6df9fb..2034ce8 100644 --- a/src/index.js +++ b/src/index.js @@ -17,9 +17,10 @@ import { QmtBridgeRestDataSource } from './data-source/QmtBridgeRestDataSource.j import { registerSettings, resolveStartupConnection, getStrategies } from './settings.js'; import { DataStore } from './storage/DataStore.js'; import { PositionManager } from './position/PositionManager.js'; +import { PositionSync } from './position/PositionSync.js'; import { registerApi } from './api/index.js'; -import { MarketDataHub } from './market/MarketDataHub.js'; -import { MarketFeed } from './market/MarketFeed.js'; +import { QuoteHub } from './market/QuoteHub.js'; +import { QuoteSync } from './market/QuoteSync.js'; import { QmtHealthMonitor } from './data-source/QmtHealthMonitor.js'; import { TradeSync } from './trades/TradeSync.js'; @@ -58,22 +59,27 @@ async function apply(ctx, config) { dataDir: config?.dataDir, }); - // R-005:行情数据(MarketDataHub 缓存 + MarketFeed 获取,页面轮询读取) - const marketHub = new MarketDataHub({ logger, dataStore: storage }); - // 启动时从磁盘加载行情缓存(首屏快速展示,不依赖实盘订阅) - marketHub.warmup().catch((e) => logger.warn('[one-divine-lot] 行情预热失败: ' + e.message)); - const marketFeed = new MarketFeed({ hub: marketHub, runtime: { settings, dataSource }, logger }); + // R-015/迭代 13:盘口内存快照(QuoteSync 取数 + QuoteHub 存查;纯内存不落库,价格单一入口) + const marketHub = new QuoteHub({ logger }); + const marketFeed = new QuoteSync({ hub: marketHub, runtime: { settings, dataSource }, logger }); marketFeed.start(startup.baseUrl ?? config?.qmtBaseUrl); // QMT 连接健康检查(服务端定时探测 + 缓存,2026-09-01) const qmtHealthMonitor = new QmtHealthMonitor({ runtime: { settings, dataSource }, logger }); qmtHealthMonitor.start(); + // 持仓内存快照同步(2026-09-02 优化:10s 定时全量同步,以快照为准;失败保留上次快照) + // 幽灵持仓自动清仓:QMT 连续 3 轮(约 30s)消失的 code,本地当前持仓自动转历史(带账户身份守卫防误清) + const positionSync = new PositionSync({ runtime: { dataSource, storage }, logger }); + positionSync.start(); + // S5: 分仓逻辑(份额分配/查询) // R-013:注入策略自定义字段定义读取回调(updateHoldingValues 校验用;settings.getStrategies 已归一化 configSchema) + // 2026-09-02:注入 positionSync —— getAllPositions 改读内存快照(读穿透兜底),不再每次请求穿透 QMT const manager = new PositionManager({ dataSource, storage, + positionSync, getStrategySchema: (strategyId) => getStrategies(settings).find((s) => s.id === strategyId)?.configSchema ?? [], }); @@ -82,7 +88,7 @@ async function apply(ctx, config) { tradeSync.start(); // S6: 服务端 HTTP API(R-004:注入 dataSource 以编排激活热切换;R-009:注入 storage 供 trades/history) - registerApi(ctx, { manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor, storage, tradeSync }); + registerApi(ctx, { manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor, storage, tradeSync, positionSync }); // 启动时可用性检查(日志,不阻塞) dataSource.isAvailable().then((ok) => { @@ -94,8 +100,9 @@ async function apply(ctx, config) { ctx.effect(() => { logger.info('[one-divine-lot] 插件已加载(S1-S6 就绪)'); return () => { - marketFeed.stop(); + marketFeed.stop(); // 停止盘口定时同步(迭代 13:QuoteSync;内存快照保留供收尾读取) qmtHealthMonitor.stop(); + positionSync.stop(); // 停止持仓快照定时同步(2026-09-02;内存快照保留供收尾读取) tradeSync.stop(); // 停止交易记录定时同步(迭代 07) storage.close(); // 关闭 SQLite 连接(迭代 06) logger.info('[one-divine-lot] 插件已释放'); diff --git a/src/market/MarketDataHub.js b/src/market/MarketDataHub.js deleted file mode 100644 index c8ca409..0000000 --- a/src/market/MarketDataHub.js +++ /dev/null @@ -1,165 +0,0 @@ -/** - * MarketDataHub —— 行情缓存服务(R-005,2026-08-31 从 market-cache.js 拆分;R-006 加持久化) - * - * 职责(单一):行情缓存的存储与查询。 - * - 持有行情缓存 Map(code → snapshot); - * - getByCodes(codes):按 code 查询;miss 的 code 现场用 REST /data/tick 拉取补入(查询保证有值,调用方无感); - * - 写入入口:ingest(code→snapshot 批量写入,由 MarketFeed 调用)。 - * - * 2026-08-31 持久化(DataStore 为数据快速展现存在,老师确认): - * - 接入 DataStore 的 store.market.json:启动时读盘填充内存缓存(首屏即有上次价格); - * - 实盘增量更新(ingest)后写回磁盘(防抖),重启不丢价格; - * - 查询命中顺序:内存 → 磁盘 → REST 补拉。 - * - * 2026-09-01 修复:watchCodes 过滤(缓存膨胀 bug)—— - * - WS 推送是全市场(whole),若不过滤会把 5 万+ 只股票写进缓存(store.market.json 膨胀到 38MB); - * - 只保留「关注集合」(前端查询过的 code + 持仓 code)的行情; - * - 关注集合来源:getByCodes 查询的 code 自动加入;warmup/启动时加入持仓 code。 - */ - -import { resolveActiveBaseUrl } from '../api/common.js'; - -const REQUEST_TIMEOUT_MS = 10000; -const PERSIST_DEBOUNCE_MS = 3000; // 写盘防抖 - -export class MarketDataHub { - /** - * @param {object} opts - * @param {object} [opts.logger] - * @param {import('../storage/DataStore.js').DataStore} [opts.dataStore] 持久化存储(可选;提供则启用行情落盘) - */ - constructor({ logger, dataStore } = {}) { - this.logger = logger; - this.dataStore = dataStore; - /** 行情缓存:code → snapshot(内存,启动时从磁盘加载) */ - this.cache = new Map(); - /** 关注集合:只缓存/持久化这些 code 的行情(防止全市场推送膨胀) */ - this.watchCodes = new Set(); - this.stats = { totalTicks: 0, wsPushCount: 0, restFallbackCount: 0, diskHitCount: 0, lastTickAt: 0 }; - this._persistTimer = null; - } - - /** 加入关注集合(查询的 code / 持仓 code) */ - watch(codes) { - const list = Array.isArray(codes) ? codes : []; - for (const c of list) { - if (c) this.watchCodes.add(c); - } - } - - /** 启动时从磁盘加载行情缓存 + 关注集合 */ - async warmup() { - if (!this.dataStore) return; - try { - const all = await this.dataStore.loadMarket(); - let n = 0; - for (const [code, snap] of Object.entries(all)) { - // 只加载数据到内存缓存,不加入 watchCodes—— - // watchCodes 应由「持仓 + 前端查询」驱动,不能从磁盘历史认可(否则全市场又进缓存) - if (snap) { this.cache.set(code, snap); n++; } - } - if (n > 0) this.logger?.info?.('[one-divine-lot] MarketDataHub 从磁盘加载行情: ' + n + ' 条'); - } catch (e) { - this.logger?.debug?.('[one-divine-lot] MarketDataHub 行情预热失败: ' + (e?.message ?? e)); - } - } - - /** - * 按 code 列表查询缓存行情;miss 依次查磁盘、REST 补拉(保证返回尽可能全) - */ - async getByCodes(codes, { settings, dataSource } = {}) { - const list = Array.isArray(codes) ? codes : []; - // 查询的 code 加入关注集合 - this.watch(list); - const out = {}; - const miss = []; - for (const c of list) { - const v = this.cache.get(c); - if (v) { out[c] = v; continue; } - // 磁盘兜底(持久化缓存) - if (this.dataStore) { - const dv = await this.dataStore.getMarketQuote(c); - if (dv) { - this.cache.set(c, dv); - out[c] = dv; - this.stats.diskHitCount++; - continue; - } - } - miss.push(c); - } - // miss 补拉(REST /data/tick) - if (miss.length > 0) { - await this._fetchAndIngest(miss, { settings, dataSource }); - for (const c of miss) { - const v = this.cache.get(c); - if (v) out[c] = v; - } - } - return out; - } - - /** 批量写入缓存(MarketFeed 调用:WS 推送 / REST 定时刷新)——只保留关注集合内 */ - ingest(data) { - if (!data) return; - let dirty = false; - for (const [code, snap] of Object.entries(data)) { - // 只缓存关注集合内的 code(防止全市场推送膨胀) - if (!this.watchCodes.has(code)) continue; - if (snap) { - this.cache.set(code, { ...this.cache.get(code), ...snap }); - dirty = true; - } - } - if (dirty) this._schedulePersist(); - } - - /** 批量写入并统计(WS 推送专用) */ - ingestWs(data) { - this.ingest(data); - this.stats.totalTicks++; - this.stats.wsPushCount++; - this.stats.lastTickAt = Date.now(); - } - - /** 缓存大小(诊断用) */ - get size() { - return this.cache.size; - } - - /** 定期写盘(防抖) */ - _schedulePersist() { - if (!this.dataStore) return; - if (this._persistTimer) clearTimeout(this._persistTimer); - this._persistTimer = setTimeout(() => { - this._persistTimer = null; - const quotes = {}; - // 只写关注集合内的 code(防止 store.market.json 膨胀) - for (const code of this.watchCodes) { - const snap = this.cache.get(code); - if (snap) quotes[code] = snap; - } - this.dataStore.setMarketQuotes(quotes).catch((e) => { - this.logger?.debug?.('[one-divine-lot] 行情写盘失败: ' + (e?.message ?? e)); - }); - }, PERSIST_DEBOUNCE_MS); - } - - /** REST 拉取并写入缓存(miss 补拉 + 定时刷新共用) */ - async _fetchAndIngest(codes, { settings, dataSource } = {}) { - if (!Array.isArray(codes) || codes.length === 0) return; - try { - const baseUrl = resolveActiveBaseUrl(settings, dataSource); - const r = await fetch(baseUrl + '/data/tick?codes=' + encodeURIComponent(codes.join(',')), { - signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), - }); - const j = await r.json().catch(() => ({})); - if (j.ok && j.data) { - this.ingest(j.data); - this.stats.restFallbackCount++; - } - } catch (e) { - this.logger?.debug?.('[one-divine-lot] MarketDataHub REST 拉取失败: ' + (e?.message ?? e)); - } - } -} \ No newline at end of file diff --git a/src/market/MarketFeed.js b/src/market/MarketFeed.js deleted file mode 100644 index 83b6036..0000000 --- a/src/market/MarketFeed.js +++ /dev/null @@ -1,197 +0,0 @@ -/** - * MarketFeed —— 行情数据获取(R-005,2026-08-31 从 market-cache.js 拆分) - * - * 职责(单一):从 QMT Bridge 获取行情数据,写入 MarketDataHub。 - * - WS 订阅连接:ws://<激活host>:8610/ws,收全市场推送 → hub.ingestWs(); - * - REST 定时刷新:每 5s 对 hub 中已有 code 拉 /data/tick 刷新(WS 不可靠时的兜底); - * - 启动主动拉取(prime):服务启动即拉一次持仓股票的盘口,保证缓存有最新价 - * (2026-09-01 老师指出:服务端启动就该自动拿最新盘口,前端任何时候打开都有价,不依赖前端轮询触发); - * - 断线指数退避重连;setBaseUrl 热切换(激活配置变化时重建连接)。 - * - * 与 MarketDataHub 的分工: - * - MarketFeed 只管「取数据」,不直接对外提供查询; - * - 数据写入走 hub.ingest() / hub.ingestWs(),查询走 hub.getByCodes()。 - */ - -import { resolveActiveBaseUrl } from '../api/common.js'; - -const RECONNECT_BASE_MS = 1000; -const RECONNECT_MAX_MS = 15000; -const SNAPSHOT_REFRESH_MS = 5000; // REST 定时刷新间隔 -const REQUEST_TIMEOUT_MS = 10000; - -/** 从 baseUrl 解析 ws 地址 */ -function buildWsUrl(baseUrl) { - if (!baseUrl) return ''; - try { - const u = new URL(String(baseUrl).replace(/\/+$/, '')); - return 'ws://' + u.host + '/ws'; - } catch { - return ''; - } -} - -export class MarketFeed { - /** - * @param {object} opts - * @param {object} opts.hub MarketDataHub 实例(缓存服务,写入目标) - * @param {object} opts.runtime { settings, dataSource } —— 解析激活 baseUrl、取持仓 code 用 - * @param {object} [opts.logger] - */ - constructor({ hub, runtime, logger } = {}) { - this.hub = hub; - this.runtime = runtime; - this.logger = logger; - this.ws = null; - this.reconnectTimer = null; - this.reconnectAttempt = 0; - this.refreshTimer = null; - this.baseUrl = ''; - this.mounted = true; - } - - /** 启动:主动拉持仓盘口 + 连接 WS + 启动 REST 定时刷新 */ - start(baseUrl) { - this.baseUrl = String(baseUrl ?? '').replace(/\/+$/, ''); - if (!this.baseUrl) { - this.logger?.warn?.('[one-divine-lot] MarketFeed 无 baseUrl,不启动'); - return; - } - this.mounted = true; - // 启动主动拉一次持仓盘口(服务端预热,前端打开即有价) - this._primePositions(); - this._connect(); - this._startRefresh(); - } - - /** 热切换:激活配置变化时调用 */ - setBaseUrl(baseUrl) { - const next = String(baseUrl ?? '').replace(/\/+$/, ''); - if (next === this.baseUrl) return; - this.baseUrl = next; - if (!this.mounted) return; - this._teardownWs(); - this._connect(); - // 切换后也主动拉一次(新地址的盘口) - this._primePositions(); - } - - /** 停止(插件释放时) */ - stop() { - this.mounted = false; - this._teardownWs(); - if (this.refreshTimer) { - clearInterval(this.refreshTimer); - this.refreshTimer = null; - } - } - - /** - * 启动/切换时主动拉持仓股票的盘口(REST /data/tick)写入缓存。 - * 保证服务端启动即有最新价,前端任何时候打开页面都能从缓存拿到价格(不依赖前端轮询触发)。 - */ - async _primePositions() { - if (!this.mounted || !this.baseUrl) return; - try { - // 从数据源拿当前持仓 code(QMT 真实持仓) - const positions = await this.runtime.dataSource.getPositions(); - const codes = (positions || []).map((p) => p.code).filter(Boolean); - if (codes.length === 0) { - this.logger?.debug?.('[one-divine-lot] MarketFeed 无持仓可预热'); - return; - } - // 持仓 code 加入关注集合(ingest 只保留这些) - this.hub.watch(codes); - const baseUrl = resolveActiveBaseUrl(this.runtime.settings, this.runtime.dataSource); - for (let i = 0; i < codes.length; i += 50) { - const batch = codes.slice(i, i + 50); - const res = await fetch(baseUrl + '/data/tick?codes=' + encodeURIComponent(batch.join(',')), { - signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), - }); - const j = await res.json().catch(() => ({})); - if (j.ok && j.data) { - this.hub.ingest(j.data); - this.logger?.info?.('[one-divine-lot] MarketFeed 启动预热盘口: ' + Object.keys(j.data).length + ' 条'); - } - } - } catch (e) { - this.logger?.debug?.('[one-divine-lot] MarketFeed 启动预热失败: ' + (e?.message ?? e)); - } - } - - /** REST 定时刷新:对 hub 已有 code 拉 /data/tick 补数据 */ - _startRefresh() { - if (this.refreshTimer) clearInterval(this.refreshTimer); - const run = async () => { - if (!this.mounted || !this.baseUrl) return; - const knownCodes = [...this.hub.cache.keys()]; - if (knownCodes.length === 0) return; - try { - const baseUrl = resolveActiveBaseUrl(this.runtime.settings, this.runtime.dataSource); - for (let i = 0; i < knownCodes.length; i += 50) { - const batch = knownCodes.slice(i, i + 50); - const res = await fetch(baseUrl + '/data/tick?codes=' + encodeURIComponent(batch.join(',')), { - signal: AbortSignal.timeout(REQUEST_TIMEOUT_MS), - }); - const j = await res.json().catch(() => ({})); - if (j.ok && j.data) { - this.hub.ingest(j.data); - if (this.hub.stats) this.hub.stats.restFallbackCount++; // 轮询也计入统计 - } - } - } catch (e) { - this.logger?.debug?.('[one-divine-lot] MarketFeed REST 定时刷新失败: ' + (e?.message ?? e)); - } - }; - run(); - this.refreshTimer = setInterval(run, SNAPSHOT_REFRESH_MS); - } - - _teardownWs() { - if (this.reconnectTimer) { - clearTimeout(this.reconnectTimer); - this.reconnectTimer = null; - } - if (this.ws) { - try { this.ws.onclose = null; this.ws.close(); } catch { /* ignore */ } - this.ws = null; - } - } - - _connect() { - if (!this.mounted || !this.baseUrl) return; - const wsUrl = buildWsUrl(this.baseUrl); - if (!wsUrl) { - this.logger?.warn?.('[one-divine-lot] MarketFeed ws 地址无效: ' + this.baseUrl); - return; - } - - this.logger?.info?.('[one-divine-lot] MarketFeed 连接: ' + wsUrl); - const ws = new WebSocket(wsUrl); - this.ws = ws; - - ws.onopen = () => { - this.reconnectAttempt = 0; - this.logger?.info?.('[one-divine-lot] MarketFeed WS 已连接: ' + wsUrl); - }; - - ws.onmessage = (ev) => { - let msg = null; - try { msg = typeof ev.data === 'string' ? JSON.parse(ev.data) : null; } catch { msg = null; } - if (!msg || msg.type !== 'whole' || !msg.data) return; - this.hub.ingestWs(msg.data); - }; - - ws.onclose = () => { - if (!this.mounted) return; - this.logger?.warn?.('[one-divine-lot] MarketFeed WS 断开,重连中...'); - const attempt = this.reconnectAttempt++; - const delay = Math.min(RECONNECT_BASE_MS * Math.pow(2, attempt), RECONNECT_MAX_MS); - this.reconnectTimer = setTimeout(() => { - if (this.mounted && this.baseUrl) this._connect(); - }, delay); - }; - - ws.onerror = () => { /* onclose 触发重连 */ }; - } -} diff --git a/src/market/QuoteHub.js b/src/market/QuoteHub.js new file mode 100644 index 0000000..2eafcdf --- /dev/null +++ b/src/market/QuoteHub.js @@ -0,0 +1,165 @@ +/** + * QuoteHub —— 盘口内存快照存储与查询(R-015 / 迭代 13,2026-09-02) + * + * 职责(单一):盘口数据的「存 + 查」。取数由 QuoteSync 负责;本类只消费 QuoteSync.ingest / + * dataSource 读穿透,不再自己 fetch REST(修正旧 MarketDataHub 绕过适配层的层级破洞)。 + * + * 2026-09-02 R-015(老师拍板): + * - 纯内存,不落库:market_quotes_cache 表退役(DROP),「重启首屏有价」由 QuoteSync 启动 prime + + * 读穿透兜底保证(秒级有价);替换 MarketDataHub(三级命中改两级:内存 → 读穿透); + * - 价格单一入口:getQuote / getQuotes;涨停/跌停(/data/instrument 按交易日缓存)经合约信息缓存提供; + * - watchCodes 关注集合维持现状:只进不出、无上限(WS 移除后膨胀仅轻微浪费;切回旧 tab 价格秒显); + * - ingest 入口来源无关(source 标签:rest;WS 将来回归时加 source 即可,存储查询不动); + * + * 修复继承:旧 MarketDataHub.warmup 依赖 DataStore.loadMarket()(return this 残迹)导致磁盘预热从未生效—— + * 本类无磁盘层,bug 不复存在。 + */ + +export class QuoteHub { + /** + * @param {object} opts + * @param {object} [opts.logger] + */ + constructor({ logger } = {}) { + this.logger = logger; + /** 盘口内存快照:code → snapshot(全量 tick 字段 + updatedAt;只读约定) */ + this.quotes = new Map(); + /** 关注集合:定时刷新/prime 的 code 范围(只进不出,老师拍板维持现状) */ + this.watchCodes = new Set(); + /** 合约静态信息缓存:code → { upStopPrice, downStopPrice, preClose, tradingDay }(交易日内有效) */ + this.instruments = new Map(); + this.stats = { syncCount: 0, failCount: 0, lastError: '', lastSyncedAt: 0, restCount: 0, watchSize: 0 }; + } + + /** 加入关注集合(QuoteSync prime / api 查询共用) */ + watch(codes) { + const list = Array.isArray(codes) ? codes : []; + for (const c of list) { + if (c) this.watchCodes.add(c); + } + this.stats.watchSize = this.watchCodes.size; + } + + /** + * 统一写入入口(来源无关;QuoteSync 调用)。 + * @param {Object} data code → snapshot + * @param {object} [opts] { source: 'rest' }(来源标签;WS 回归时新增来源即可) + */ + ingest(data, { source = 'rest' } = {}) { + void source; + if (!data || typeof data !== 'object') return; + for (const [code, snap] of Object.entries(data)) { + // watch 过滤保留(防止未来全市场来源膨胀);读穿透路径先 watch 再 ingest,不受影响 + if (!this.watchCodes.has(code)) continue; + if (snap) this.quotes.set(code, { ...this.quotes.get(code), ...snap }); + } + } + + /** 同步成功标记(QuoteSync 每轮成功后调用,供指示灯) */ + markSynced() { + this.stats.syncCount++; + this.stats.lastSyncedAt = Date.now(); + } + + /** 同步失败标记 */ + markFailed(err) { + this.stats.failCount++; + this.stats.lastError = err?.message ?? String(err ?? ''); + } + + /** + * 批量查询:内存命中;miss 读穿透(dataSource.getTicks 补拉并回填)。 + * 与旧三级命中(内存→磁盘→REST)对齐为两级(磁盘层退役)。 + * @param {string[]} codes + * @param {object} [opts] { dataSource }(读穿透用) + * @returns {Promise>} + */ + async getQuotes(codes, { dataSource } = {}) { + const list = Array.isArray(codes) ? codes : []; + this.watch(list); // 查询的 code 加入关注集合(旧行为保留) + const out = {}; + const miss = []; + for (const c of list) { + const v = this.quotes.get(c); + if (v) { out[c] = v; continue; } + miss.push(c); + } + if (miss.length > 0 && dataSource?.getTicks) { + try { + const backfill = await dataSource.getTicks(miss); + this.ingest(backfill, { source: 'rest' }); + this.stats.restCount++; + } catch (e) { + this.logger?.debug?.('[one-divine-lot] QuoteHub 读穿透失败: ' + (e?.message ?? e)); + } + for (const c of miss) { + const v = this.quotes.get(c); + if (v) out[c] = v; + } + } + return out; + } + + /** 单码便捷查询(前端 getPrice 消费形态) */ + async getQuote(code, opts) { + if (!code) return null; + const out = await this.getQuotes([code], opts); + return out[code] ?? null; + } + + // ===== 合约静态信息(涨停/跌停;交易日缓存)===== + + /** 合约缓存是否可用(存在且交易日一致) */ + _instrumentFresh(code, tradingDay) { + const hit = this.instruments.get(code); + return hit && hit.tradingDay === tradingDay; + } + + /** + * 取合约静态信息(涨停/跌停/昨收);缓存失效(无缓存/换交易日)时经 fetcher 拉取。 + * @param {string} code + * @param {string|number} tradingDay YYYYMMDD(日切失效键) + * @param {(code: string) => Promise} fetcher dataSource.getInstrument 包装 + * @returns {Promise<{upStopPrice:number, downStopPrice:number, preClose:number}|null>} + */ + async getInstrumentCached(code, tradingDay, fetcher) { + if (!code || !fetcher) return null; + if (this._instrumentFresh(code, tradingDay)) { + const { upStopPrice, downStopPrice, preClose } = this.instruments.get(code); + return { upStopPrice, downStopPrice, preClose }; + } + try { + const inst = await fetcher(code); + if (!inst) return null; + this.instruments.set(code, { + upStopPrice: inst.upStopPrice ?? 0, + downStopPrice: inst.downStopPrice ?? 0, + preClose: inst.preClose ?? 0, + tradingDay: inst.tradingDay || tradingDay, + }); + const { upStopPrice, downStopPrice, preClose } = this.instruments.get(code); + return { upStopPrice, downStopPrice, preClose }; + } catch (e) { + this.logger?.debug?.('[one-divine-lot] QuoteHub 合约信息拉取失败: ' + code + ' ' + (e?.message ?? e)); + return null; + } + } + + /** 热切换/日切时清空合约缓存 */ + clearInstruments() { + this.instruments.clear(); + } + + /** 快照规模(诊断) */ + get size() { + return this.quotes.size; + } + + /** + * 兼容旧 MarketDataHub.getByCodes 签名(api/market.js 消费)。 + * @deprecated 用 getQuotes + */ + async getByCodes(codes, opts = {}) { + return this.getQuotes(codes, opts); + } +} diff --git a/src/market/QuoteSync.js b/src/market/QuoteSync.js new file mode 100644 index 0000000..842b07f --- /dev/null +++ b/src/market/QuoteSync.js @@ -0,0 +1,163 @@ +/** + * QuoteSync —— 盘口数据定时同步(R-015 / 迭代 13,2026-09-02) + * + * 职责(单一):把 QMT 盘口数据同步进 QuoteHub 内存快照。与 PositionSync(10s/持仓)、 + * TradeSync(60s/交易)三域同构:启动预热 + 定时 + 失败容忍(保留上次价)+ 手动 syncNow。 + * + * 老师拍板(R-015): + * - WS 数据通路移除(从无生效结论、全市场推送被 watch 过滤成本高、REST 5s 已覆盖); + * ingest 带 source 标签,将来 WS 回归 = 加一个来源,Hub/查询零改动; + * - 同步失败保留上次价(内存不动),恢复后 5s 内自动追上; + * - 涨停/跌停走 /data/instrument,按交易日缓存(当天不变,不跟 5s 周期;日切/热切换失效重拉)。 + */ + +import { resolveActiveBaseUrl } from '../api/common.js'; + +const SYNC_INTERVAL_MS = 5 * 1000; // 同步间隔(与旧行情轮询一致) +const TICK_BATCH = 50; // tick 批量(沿用旧 MarketFeed 分批) + +export class QuoteSync { + /** + * @param {object} opts + * @param {object} opts.hub QuoteHub 实例(写入目标) + * @param {object} opts.runtime { settings, dataSource } + * @param {object} [opts.logger] + * @param {number} [opts.intervalMs] 同步间隔(测试可调小) + */ + constructor({ hub, runtime, logger, intervalMs = SYNC_INTERVAL_MS } = {}) { + this.hub = hub; + this.runtime = runtime; + this.logger = logger; + this.intervalMs = intervalMs; + this.timer = null; + this.mounted = false; + this.syncing = false; + this._tradingDay = ""; + this.stats = { primeCount: 0, syncCount: 0, failCount: 0, lastError: "", lastSyncedAt: 0 }; + } + + /** 启动:主动拉持仓盘口(prime)+ 定时刷新(REST;WS 已移除) */ + start() { + if (this.mounted) return; + this.mounted = true; + this._primePositions(); + this.timer = setInterval(() => { this.syncNow().catch(() => {}); }, this.intervalMs); + this.logger?.info?.("[one-divine-lot] QuoteSync 启动(5s REST 同步;WS 通路已移除 R-015)"); + } + + /** 停止(插件释放时);内存快照保留 */ + stop() { + this.mounted = false; + if (this.timer) { + clearInterval(this.timer); + this.timer = null; + } + } + + /** 热切换:激活配置变化时调用(api/qmt-connections 编排)——清合约缓存 + 立即 prime */ + setBaseUrl() { + if (!this.mounted) return; + this.hub.clearInstruments(); // 新环境合约信息全部失效 + this._primePositions(); + } + + /** + * 启动/热切换时主动拉持仓股票的盘口,保证服务端启动即有最新价,前端打开页面即有价。 + */ + async _primePositions() { + if (!this.mounted) return; + try { + const positions = await this.runtime.dataSource.getPositions(); + const codes = (positions || []).map((p) => p.code).filter(Boolean); + if (codes.length === 0) { + this.logger?.debug?.("[one-divine-lot] QuoteSync 无持仓可预热"); + return; + } + this.hub.watch(codes); + const ticks = await this.runtime.dataSource.getTicks(codes); + this.hub.ingest(ticks, { source: "rest" }); + this.hub.markSynced(); + this.stats.primeCount++; + this.stats.lastSyncedAt = Date.now(); + this.logger?.info?.("[one-divine-lot] QuoteSync 启动预热盘口: " + Object.keys(ticks).length + " 条"); + } catch (e) { + this.hub.markFailed(e); + this.stats.failCount++; + this.stats.lastError = e?.message ?? String(e); + this.logger?.debug?.("[one-divine-lot] QuoteSync 启动预热失败: " + (e?.message ?? e)); + } + } + + /** + * 同步一轮:watch 集合分批拉 tick → hub.ingest;watch 中新 code 懒拉合约信息(涨停/跌停)。 + * 任何失败只记统计,不动内存(读方继续消费上次价)。 + */ + async syncNow() { + // mounted 不拦手动同步(测试/指示灯点击在 stop 后仍可用);syncing 只防重入(迭代 12 同款教训) + if (this.syncing) return null; + const { dataSource } = this.runtime; + if (!dataSource) return null; + this.syncing = true; + try { + const codes = [...this.hub.watchCodes]; + if (codes.length > 0) { + let ok = 0; + for (let i = 0; i < codes.length; i += TICK_BATCH) { + const batch = codes.slice(i, i + TICK_BATCH); + const ticks = await dataSource.getTicks(batch); + this.hub.ingest(ticks, { source: "rest" }); + ok += Object.keys(ticks).length; + } + this.hub.markSynced(); + this.stats.syncCount++; + this.stats.lastSyncedAt = Date.now(); + await this._syncInstruments(codes); + return { quotes: ok, watch: codes.length }; + } + // watch 空:不空转(watch 由 prime 登记,正常场景恒非空) + return { quotes: 0, watch: 0 }; + } catch (e) { + this.hub.markFailed(e); + this.stats.failCount++; + this.stats.lastError = e?.message ?? String(e); + this.logger?.debug?.("[one-divine-lot] QuoteSync 同步失败(保留旧价): " + this.stats.lastError); + return null; + } finally { + this.syncing = false; + } + } + + /** + * 合约信息懒拉:watch 集合中无缓存(或交易日变化)的 code 拉 /data/instrument。 + * 交易日取 asset.tradingDate(拿不到则沿用上次值兜底)。 + */ + async _syncInstruments(codes) { + const { dataSource } = this.runtime; + try { + const asset = await dataSource.getAsset(); + if (asset?.tradingDate) { + if (this._tradingDay && asset.tradingDate !== this._tradingDay) { + this.hub.clearInstruments(); // 日切:全部失效重拉 + } + this._tradingDay = asset.tradingDate; + } + } catch { /* getAsset 失败:沿用上次交易日兜底 */ } + const fetcher = (code) => dataSource.getInstrument(code); + for (const code of codes) { + if (this.hub._instrumentFresh(code, this._tradingDay)) continue; + await this.hub.getInstrumentCached(code, this._tradingDay, fetcher); // 失败静默,下轮重试 + } + } + + /** 指示灯状态读点 */ + getStatus() { + return { + syncedAt: this.stats.lastSyncedAt, + periodMs: this.intervalMs, + failCount: this.stats.failCount, + lastError: this.stats.lastError, + snapshotSize: this.hub.size, + watchSize: this.hub.watchCodes.size, + }; + } +} diff --git a/src/settings.js b/src/settings.js index 8a0df9a..4db1116 100644 --- a/src/settings.js +++ b/src/settings.js @@ -49,6 +49,11 @@ export const COLUMN_META = [ { key: 'avgPrice', label: '成本价', fixed: false }, { key: 'lastTradePrice', label: '最后一笔成交价', fixed: false }, { key: 'shares', label: '策略份额', fixed: false }, + // R-015/迭代 13:行情列(盘口快照 + 合约信息;defaultVisible=false 默认隐藏,列设置勾选开启) + { key: 'upStopPrice', label: '涨停价', fixed: false, defaultVisible: false }, + { key: 'downStopPrice', label: '跌停价', fixed: false, defaultVisible: false }, + { key: 'open', label: '今开', fixed: false, defaultVisible: false }, + { key: 'high', label: '最高', fixed: false, defaultVisible: false }, ]; /** 策略设置 schema(schemastery 定义,settings.register 要求 schema 对象) */ @@ -438,8 +443,8 @@ export function normalizeStrategyColumns(scope, strategyId) { const configSchema = Array.isArray(strategy?.configSchema) ? strategy.configSchema : []; const raw = (value.strategyColumns && value.strategyColumns[strategyId]) || []; - // 1. 默认全列(基础列全显默认序 + 字段列全显 configSchema 序) - const baseCols = COLUMN_META.map((c) => ({ key: c.key, label: c.label, fixed: false, kind: 'base' })); + // 1. 默认全列(基础列默认序 + 字段列 configSchema 序;显隐缺省跟随 defaultVisible,R-015 行情列默认隐藏) + const baseCols = COLUMN_META.map((c) => ({ key: c.key, label: c.label, fixed: false, kind: 'base', defaultVisible: c.defaultVisible !== false })); const fieldCols = configSchema.map((f) => ({ key: f.key, label: f.label || f.key, fixed: false, kind: 'field' })); const defaults = [...baseCols, ...fieldCols]; @@ -461,7 +466,8 @@ export function normalizeStrategyColumns(scope, strategyId) { } for (const d of defaults) { if (!placed.has(d.key)) { - ordered.push({ ...d, visible: visibleMap.has(d.key) ? visibleMap.get(d.key) : true }); + const dv = visibleMap.has(d.key) ? visibleMap.get(d.key) : d.defaultVisible !== false; + ordered.push({ ...d, visible: dv }); } } return ordered; diff --git a/src/storage/DataStore.js b/src/storage/DataStore.js index aff31d2..8f11838 100644 --- a/src/storage/DataStore.js +++ b/src/storage/DataStore.js @@ -2,8 +2,8 @@ * DataStore —— 数据存储门面(R-006 JSON → R-008 SQLite 迁移过渡) * * 迭代 06:存储引擎从 JSON data store 升级为 SQLite(node:sqlite)。 - * - 数据实际存于 SqliteStore(strategy_holdings / market_quotes_cache 表); - * - 本模块保留对外兼容 API(getDataset/getAllDatasets/getMarketQuotes 等),上层(PositionManager / MarketDataHub)调用点不变; + * - 数据实际存于 SqliteStore(strategy_holdings / trade_orders / trade_fills 表;market_quotes_cache 已退役 DROP,R-015); + * - 本模块保留对外兼容 API(getDataset/getAllDatasets 等),上层(PositionManager / QuoteSync)调用点不变; * - 启动时检测旧 JSON → 自动迁移(幂等),JSON 迁移后废弃(备份 .bak)。 * * 注:PositionManager 适配单票生命周期后,旧整策略读写 setDataset/removeDataset 已删除 @@ -19,7 +19,6 @@ export class DataStore { constructor({ dataDir } = {}) { this.sqlite = new SqliteStore({ dataDir }); this.loaded = false; - this.marketLoaded = false; } /** 确保初始化 + 自动迁移(幂等) */ @@ -42,13 +41,6 @@ export class DataStore { return { version: 1, strategies: this.getAllDatasetsSync() }; } - /** 加载行情(兼容旧调用) */ - async loadMarket() { - await this._ensure(); - this.marketLoaded = true; - return this; - } - // ===== 持仓生命周期(新 API,PositionManager 适配后使用)===== /** 建仓 */ @@ -122,26 +114,6 @@ export class DataStore { return [...map.entries()].map(([strategyId, dataset]) => ({ strategyId, dataset })); } - // ===== 行情(委托 SqliteStore)===== - - /** 单码行情 */ - async getMarketQuote(code) { - await this._ensure(); - return this.sqlite.getMarketQuote(code); - } - - /** 批量行情 */ - async getMarketQuotes(codes) { - await this._ensure(); - return this.sqlite.getMarketQuotes(codes); - } - - /** 写行情(整体替换防膨胀) */ - async setMarketQuotes(quotes) { - await this._ensure(); - return this.sqlite.setMarketQuotes(quotes); - } - // ===== 交易记录(R-009 迭代 07,委托 SqliteStore)===== /** UPSERT 委托(幂等) */ diff --git a/src/storage/SqliteStore.js b/src/storage/SqliteStore.js index d627353..1f9d1e6 100644 --- a/src/storage/SqliteStore.js +++ b/src/storage/SqliteStore.js @@ -6,14 +6,13 @@ * * 表结构(技术实现方案): * - strategy_holdings 策略持仓生命周期表(一笔 = 一次「建仓→清仓」,历史保留) - * - market_quotes_cache 行情快照缓存(code 维度一行一码,只存 UI 消费的 lastPrice/lastClose) * - trade_orders 交易委托(R-009:委托主行,order_id 主键,UPSERT 幂等,零冗余) * - trade_fills 交易成交(R-009:成交明细,trade_id 主键,order_id 关联委托) + * - (market_quotes_cache 已退役 DROP,R-015/迭代 13:行情改内存快照 QuoteHub,价格单一入口) * * 方法集(替代 DataStore 原 getDataset/setDataset/removeDataset/getAllDatasets 整体读写): * 持仓生命周期:openHolding / addShares / reduceShares / closeHolding * 持仓查询 :getCurrentHoldings / getHoldingHistory - * 行情读写 :getMarketQuote / getMarketQuotes / setMarketQuotes * 生命周期 :init / migrateJson / close */ @@ -33,6 +32,7 @@ const MARKET_FILE = 'store.market.json'; const LEGACY_FILE = 'allocations.json'; const SCHEMA_SQL = ` +DROP TABLE IF EXISTS market_quotes_cache; -- R-015:存量库清理(幂等;行情改内存快照) CREATE TABLE IF NOT EXISTS strategy_holdings ( holding_id INTEGER PRIMARY KEY AUTOINCREMENT, strategy_id TEXT NOT NULL, @@ -43,12 +43,7 @@ CREATE TABLE IF NOT EXISTS strategy_holdings ( ); CREATE UNIQUE INDEX IF NOT EXISTS idx_active_holding ON strategy_holdings (strategy_id, code) WHERE closed_at IS NULL; -CREATE TABLE IF NOT EXISTS market_quotes_cache ( - code TEXT PRIMARY KEY, - last_price REAL NOT NULL, - last_close REAL NOT NULL, - updated_at INTEGER NOT NULL DEFAULT 0 -); +-- market_quotes_cache 已退役(R-015/迭代 13:行情改内存快照 QuoteHub,价格单一入口不再落库) -- 交易委托(R-009 迭代 07:委托主行;order_id 唯一,UPSERT 幂等;零冗余 strategy_id/holding_id) CREATE TABLE IF NOT EXISTS trade_orders ( order_id TEXT PRIMARY KEY, @@ -223,28 +218,26 @@ export class SqliteStore { } } - /** 判断库是否为空(两表均无数据) */ + /** 判断库是否为空(R-015:行情表退役,只看 strategy_holdings) */ isEmpty() { this.init(); const h = this.db.prepare('SELECT COUNT(*) AS n FROM strategy_holdings').get(); - const m = this.db.prepare('SELECT COUNT(*) AS n FROM market_quotes_cache').get(); - return (h?.n ?? 0) === 0 && (m?.n ?? 0) === 0; + return (h?.n ?? 0) === 0; } // ===== 迁移(D6:一次性脚本 + 启动自动迁移,幂等)===== /** * 旧 JSON → SQLite 迁移(迁移前备份 JSON 为 *.bak;写库失败不删原文件,幂等) - * @returns {Promise<{holdings: number, quotes: number}>} 迁移的持仓行数 / 行情行数 + * @returns {Promise<{holdings: number}>} 迁移的持仓行数 */ async migrateJson() { this.init(); if (!this.isEmpty()) { - return { holdings: 0, quotes: 0, skipped: true }; // 已有数据,跳过 + return { holdings: 0, skipped: true }; // 已有数据,跳过 } let holdingsCount = 0; - let quotesCount = 0; const now = Date.now(); // 1. store.json → strategy_holdings(策略为中心平铺) @@ -274,30 +267,6 @@ export class SqliteStore { // ENOENT = 无 store.json,正常 } - // 2. store.market.json → market_quotes_cache(抽取 lastPrice/lastClose) - try { - const raw = await fsp.readFile(this.marketPath, 'utf8'); - const parsed = JSON.parse(raw); - const quotes = parsed?.quotes ?? {}; - const upsert = this.db.prepare( - 'INSERT INTO market_quotes_cache (code, last_price, last_close, updated_at) VALUES (?,?,?,?) ON CONFLICT(code) DO UPDATE SET last_price=excluded.last_price, last_close=excluded.last_close, updated_at=excluded.updated_at' - ); - this.db.exec('BEGIN'); - for (const [code, snap] of Object.entries(quotes)) { - if (!snap) continue; - const lastPrice = Number(snap.lastPrice ?? snap.last_price ?? 0); - const lastClose = Number(snap.lastClose ?? snap.last_close ?? 0); - upsert.run(code, lastPrice, lastClose, Number(snap.time ?? snap.updated_at ?? now)); - quotesCount++; - } - this.db.exec('COMMIT'); - await this._backupJson(this.marketPath); - } catch (err) { - if (err.code !== 'ENOENT') { - console.warn('[one-divine-lot] store.market.json 迁移失败: ' + err.message); - } - } - // 3. 旧 allocations.json(若仍存在,R-006 遗留迁移源,code 为中心)→ strategy_holdings try { const raw = await fsp.readFile(this.legacyPath, 'utf8'); @@ -329,7 +298,7 @@ export class SqliteStore { } } - return { holdings: holdingsCount, quotes: quotesCount }; + return { holdings: holdingsCount }; } /** 备份 JSON 文件为 .bak(D7:迁移前自动备份) */ @@ -433,60 +402,6 @@ export class SqliteStore { - // ===== 行情读写(market_quotes_cache)===== - - /** 单码行情 */ - getMarketQuote(code) { - this.init(); - const row = this.db.prepare( - 'SELECT code, last_price, last_close, updated_at FROM market_quotes_cache WHERE code=?' - ).get(code); - return row ? this._mapQuote(row) : undefined; - } - - /** 批量行情 */ - getMarketQuotes(codes) { - this.init(); - if (!codes || codes.length === 0) return {}; - const placeholders = codes.map(() => '?').join(','); - const rows = this.db.prepare( - 'SELECT code, last_price, last_close, updated_at FROM market_quotes_cache WHERE code IN (' + placeholders + ')' - ).all(...codes); - const out = {}; - for (const row of rows) out[row.code] = this._mapQuote(row); - return out; - } - - /** 写入行情(UPSERT;整体集合替换 —— 清除不在集合中的旧键,防膨胀) */ - setMarketQuotes(quotes) { - this.init(); - const upsert = this.db.prepare( - 'INSERT INTO market_quotes_cache (code, last_price, last_close, updated_at) VALUES (?,?,?,?) ON CONFLICT(code) DO UPDATE SET last_price=excluded.last_price, last_close=excluded.last_close, updated_at=excluded.updated_at' - ); - const now = Date.now(); - this.db.exec('BEGIN'); - try { - for (const [code, snap] of Object.entries(quotes || {})) { - if (!snap) continue; - upsert.run(code, Number(snap.lastPrice ?? snap.last_price ?? 0), Number(snap.lastClose ?? snap.last_close ?? 0), Number(snap.time ?? now)); - } - this.db.exec('COMMIT'); - } catch (err) { - this.db.exec('ROLLBACK'); - throw err; - } - return this.getMarketQuotes(Object.keys(quotes || {})); - } - - _mapQuote(row) { - return { - code: row.code, - lastPrice: row.last_price, - lastClose: row.last_close, - updatedAt: row.updated_at, - }; - } - // ===== 交易记录(trade_orders / trade_fills,R-009 迭代 07)===== /**