feat: 迭代03-05 交易记录接入、行情实时、QMT连接配置、data store 标准化

- R-004 QMT连接配置:多配置管理 + 会话头部快捷切换(03-QMT连接配置)
- R-005 WS盘中价格实时更新:服务端中转+缓存+前端轮询(04-WS盘中价格实时更新)
- R-006 策略数据 JSON data store 标准化(store.json/market.json/schema.json)
- R-007 交易记录接入:当日订单+成交单表合并(委托主行+展开成交明细)、时间段查询(今日/本周/本月+手动起止)、费用计算(佣金费率万m_dCommission+最低5元、印花税卖出万5、过户费双向万0.1)
- 架构治理:api 按领域拆分、component 目录、QmtHealthMonitor
This commit is contained in:
2026-09-01 13:21:49 +08:00
parent 5a723d514e
commit bb3f8bd28d
57 changed files with 5244 additions and 397 deletions
+42
View File
@@ -0,0 +1,42 @@
/**
* 服务端 HTTP API 公共工具
*
* 为什么不用 RPCconnection.rpc.intercept('/api')):
* DSH 的 /api 通道只能有一个 interceptorapi-gateway 已占用),
* 再 intercept 会抛 "already has an interceptor" 导致插件 apply 失败。
* 因此改用 webServer 自开路由(与 dsh-better-sidebar 的 /sidebar/api 同模式)。
*/
export const API_PREFIX = '/odl/api/';
/** 响应 JSON */
export function writeJson(res, status, data) {
res.writeHead(status, { 'Content-Type': 'application/json' });
res.end(JSON.stringify(data));
}
/** 读取 JSON body */
export async function readJsonBody(req) {
const chunks = [];
for await (const chunk of req) chunks.push(chunk);
const raw = Buffer.concat(chunks).toString('utf8');
if (!raw) return {};
try {
return JSON.parse(raw);
} catch {
return {};
}
}
/** 从激活配置解析 baseUrl(连接列表为空时回退数据源当前 baseUrl) */
export function resolveActiveBaseUrl(settings, dataSource) {
const state = getQmtConnections(settings);
const active = getActiveQmtConnection(state);
return String(active?.baseUrl ?? dataSource?.baseUrl ?? '').replace(/\/+$/, '');
}
// 延迟引入 settings 函数避免循环依赖
import {
getQmtConnections,
getActiveQmtConnection,
} from '../settings.js';
+102
View File
@@ -0,0 +1,102 @@
/**
* 服务端 HTTP API 入口(2026-08-31 从 api.js 按领域拆分)
*
* 结构:
* src/api/
* ├── index.js # 路由入口:合并各领域 + 公共处理(405/404/500
* ├── common.js # 公共工具(writeJson/readJsonBody/resolveActiveBaseUrl
* ├── positions.js # 持仓域(positions
* ├── strategies.js # 策略/份额/tabs 域
* ├── qmt-connections.js # QMT 连接配置域(含热切换编排)
* └── market.js # 行情域(market-snapshot
*
* 端点(POST /odl/api/<method>body = { args: {...} }):
* positions → 全量持仓
* strategies → 策略列表(含 visible)
* strategies/update → 更新策略列表(整表)
* strategies/add → 新增策略 { name }
* strategies/remove → 删除策略 { strategyId },联动清份额
* strategies/move → 排序 { strategyId, dir }
* strategy-positions → 某策略持仓 { strategyId }
* unallocated → 未分配持仓
* summary → 分仓摘要
* add-shares → 添加份额 { code, strategyId, shares }
* move-all-shares → 全部移入 { code, strategyId }
* remove-shares → 移出份额 { code, strategyId, shares }
* remove-all-shares → 一键清零 { strategyId }
* tabs / tabs/update → 通用 tab 显隐配置
* qmt-connections → QMT 连接配置(R-004
* qmt-connections/* → 增删改/激活/默认/测试/active-host/订阅代理
* market-snapshot → 按 code 查询行情缓存(R-005
*/
import { API_PREFIX, writeJson, readJsonBody } from './common.js';
import { POSITION_METHODS, handlePosition } from './positions.js';
import { STRATEGY_METHODS, handleStrategy } from './strategies.js';
import { QMT_METHODS, handleQmt } from './qmt-connections.js';
import { MARKET_METHODS, handleMarket } from './market.js';
import { TRADE_METHODS, handleTrade } from './trades.js';
/** 全部端点方法表(合并各领域) */
const METHODS = new Set([
...POSITION_METHODS,
...STRATEGY_METHODS,
...QMT_METHODS,
...MARKET_METHODS,
...TRADE_METHODS,
]);
/** 领域分发:method → handler */
const HANDLERS = [
{ set: POSITION_METHODS, handle: handlePosition },
{ set: STRATEGY_METHODS, handle: handleStrategy },
{ set: QMT_METHODS, handle: handleQmt },
{ set: MARKET_METHODS, handle: handleMarket },
{ set: TRADE_METHODS, handle: handleTrade },
];
/**
* 注册 HTTP API
* @param {import('@deepseek-ai/cordis').Context} ctx
* @param {object} runtime { manager, settings, dataSource, marketCache }
*/
export function registerApi(ctx, { manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor }) {
const runtime = { ctx, manager, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor };
ctx.effect(() => ctx.webServer.register({
kind: 'prefix',
path: '/odl/api',
handler: async (req, res) => {
if (req.method !== 'POST') {
writeJson(res, 405, { ok: false, error: { code: 'method-error', message: 'method not allowed' } });
return;
}
const pathname = new URL(req.url ?? '/', 'http://dsh.internal').pathname;
const method = pathname.startsWith(API_PREFIX) ? pathname.slice(API_PREFIX.length) : void 0;
if (!method || !METHODS.has(method)) {
writeJson(res, 404, { ok: false, error: { code: 'not-found', message: 'unknown method: ' + method } });
return;
}
try {
const body = await readJsonBody(req);
const args = body?.args ?? {};
// 按领域分发
let value;
for (const { set, handle } of HANDLERS) {
if (set.has(method)) {
value = await handle(method, args, runtime);
break;
}
}
writeJson(res, 200, { ok: true, value });
} catch (error) {
writeJson(res, 500, {
ok: false,
error: { code: error?.code ?? 'internal', message: error?.message ?? String(error), details: {} },
});
}
},
}), 'one-divine-lot.api');
ctx.logger?.info('[one-divine-lot] HTTP API 已注册 (/odl/api/*)');
}
+58
View File
@@ -0,0 +1,58 @@
/**
* 服务端 API:行情域(2026-08-31 拆分;2026-09-01 精简)
*
* 端点:
* market-snapshot → 按 code 查询行情缓存(miss 自动补拉在 MarketDataHub 内部,此处只做编排)
* qmt-health → QMT 连接健康检查(读服务端缓存;refresh=true 触发即时探测)
*
* 2026-09-01QMT health 改为「服务端定时探测(5 分钟)+ 缓存」机制(QmtHealthMonitor),
* 前端打开读缓存秒回,点击/悬停传 refresh=true 触发即时探测。
*/
import { resolveActiveBaseUrl } from './common.js';
/** 行情端点方法表 */
export const MARKET_METHODS = new Set([
'market-snapshot',
'qmt-health',
'market-stats', // 临时诊断:WS/REST 通道统计
]);
/**
* 处理行情端点
* @param {string} method
* @param {object} args
* @param {object} runtime { settings, dataSource, marketHub, marketFeed, qmtHealthMonitor }
*/
export async function handleMarket(method, args, { ctx, settings, dataSource, marketHub, marketFeed, qmtHealthMonitor }) {
switch (method) {
case 'market-snapshot': {
const codes = Array.isArray(args.codes) ? args.codes : [];
if (!marketHub) return {};
return await marketHub.getByCodes(codes, { settings, dataSource });
}
case 'qmt-health': {
// QMT 连接健康检查(服务端缓存 + 即时探测,2026-09-01
if (!qmtHealthMonitor) {
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(),
};
}
default:
throw Object.assign(new Error('unknown market method: ' + method), { code: 'not-found' });
}
}
+26
View File
@@ -0,0 +1,26 @@
/**
* 服务端 API:持仓域(从 api.js 拆分,2026-08-31
*
* 端点:
* positions → 全量持仓
*/
/** 持仓端点方法表 */
export const POSITION_METHODS = new Set([
'positions',
]);
/**
* 处理持仓端点
* @param {string} method
* @param {object} args
* @param {object} runtime { manager }
*/
export async function handlePosition(method, args, { manager }) {
switch (method) {
case 'positions':
return await manager.getAllPositions();
default:
throw Object.assign(new Error('unknown position method: ' + method), { code: 'not-found' });
}
}
+138
View File
@@ -0,0 +1,138 @@
/**
* 服务端 APIQMT 连接配置域(从 api.js 拆分,2026-08-31
*
* 端点:
* qmt-connections / add / update / remove / activate / set-default / test
* active-host
*
* 热切换编排:激活/更新/删除连接导致生效地址变化时,同步 dataSource.setBaseUrl + marketFeed.setBaseUrl
* (技术约束-008:按请求读取 baseUrl,改字段即生效)
*
* 2026-08-31:移除 subscribe-ws / unsubscribe-ws / tick-snapshot(架构演进后无人调用,
* 订阅与快照统一由 MarketFeed/MarketDataHub 内部处理,不再暴露为 HTTP 端点)
*/
import {
getQmtConnections, getActiveQmtConnection,
addQmtConnection, updateQmtConnection, removeQmtConnection,
activateQmtConnection, setDefaultQmtConnection,
} from '../settings.js';
/** QMT 连接端点方法表 */
export const QMT_METHODS = new Set([
'qmt-connections',
'qmt-connections/add',
'qmt-connections/update',
'qmt-connections/remove',
'qmt-connections/activate',
'qmt-connections/set-default',
'qmt-connections/test',
'qmt-connections/active-host',
]);
/** 服务端代理测试连接:GET {baseUrl}/health(浏览器跨域规避,Q3 手动测试) */
async function testQmtConnection(ctx, args) {
const baseUrl = String(args?.baseUrl ?? '').trim().replace(/\/+$/, '');
if (!/^https?:\/\//i.test(baseUrl)) {
throw Object.assign(new Error('HTTP 地址格式无效(需以 http:// 或 https:// 开头)'), { code: 'qmt-invalid-url' });
}
const controller = new AbortController();
const timer = setTimeout(() => controller.abort(), 15000);
const startedAt = Date.now();
try {
const res = await fetch(baseUrl + '/health', { signal: controller.signal });
const latencyMs = Date.now() - startedAt;
let status = null;
try { status = (await res.json())?.status ?? null; } catch { /* 非 JSON 响应按 status=null */ }
return { ok: res.ok, latencyMs, httpStatus: res.status, status, reachable: res.ok };
} catch (err) {
return {
ok: false,
latencyMs: Date.now() - startedAt,
reachable: false,
error: err.name === 'AbortError' ? '请求超时(15s' : (err?.message ?? String(err)),
};
} finally {
clearTimeout(timer);
}
}
/**
* 处理 QMT 连接端点
* @param {string} method
* @param {object} args
* @param {object} runtime { manager, settings, dataSource, marketCache }
*/
export async function handleQmt(method, args, { ctx, settings, dataSource, marketFeed, qmtHealthMonitor }) {
switch (method) {
case 'qmt-connections': {
const state = getQmtConnections(settings);
return { ...state, active: getActiveQmtConnection(state) };
}
case 'qmt-connections/add':
return await addQmtConnection(settings, { name: args.name, baseUrl: args.baseUrl });
case 'qmt-connections/update': {
const before = getQmtConnections(settings);
const value = await updateQmtConnection(settings, args.id, { name: args.name, baseUrl: args.baseUrl });
// 编辑的是激活配置且地址变化 → 立即热切换(技术约束-008)
if (before.activeId === args.id && value.baseUrl !== dataSource.baseUrl) {
dataSource.setBaseUrl(value.baseUrl);
marketFeed?.setBaseUrl(value.baseUrl);
qmtHealthMonitor?.probe(); // 切换后立即补探健康
ctx?.logger?.info('[one-divine-lot] 激活配置地址已变更,数据源热切换 → ' + value.baseUrl);
}
return value;
}
case 'qmt-connections/remove': {
const result = await removeQmtConnection(settings, args.id);
// 生效连接因删除发生变化 → 立即热切换(产品约束-005 删除边界)
const afterState = getQmtConnections(settings);
const active = getActiveQmtConnection(afterState);
if (active && active.baseUrl !== dataSource.baseUrl) {
dataSource.setBaseUrl(active.baseUrl);
marketFeed?.setBaseUrl(active.baseUrl);
qmtHealthMonitor?.probe(); // 切换后立即补探健康
ctx?.logger?.info('[one-divine-lot] 激活配置已删除,自动切换 → ' + active.name + ' (' + active.baseUrl + ')');
}
return { ...afterState, active, removedWasActive: result.removedWasActive, removedWasDefault: result.removedWasDefault };
}
case 'qmt-connections/activate': {
const state = await activateQmtConnection(settings, args.id);
const active = getActiveQmtConnection(state);
if (active && active.baseUrl !== dataSource.baseUrl) {
dataSource.setBaseUrl(active.baseUrl);
marketFeed?.setBaseUrl(active.baseUrl);
qmtHealthMonitor?.probe(); // 切换后立即补探健康
ctx?.logger?.info('[one-divine-lot] QMT 连接激活: ' + active.name + ' (' + active.baseUrl + ')');
}
return { ...state, active };
}
case 'qmt-connections/set-default':
return await setDefaultQmtConnection(settings, args.id);
case 'qmt-connections/test': {
// 支持传未保存的 baseUrl(表单先测后存),或传 id 用已存配置地址
let url = args.baseUrl;
if (!url && args.id) {
const state = getQmtConnections(settings);
url = state.list.find((c) => c.id === args.id)?.baseUrl;
}
return await testQmtConnection(ctx, { baseUrl: url });
}
case 'qmt-connections/active-host': {
// R-005:返回激活配置的 host 与 ws 地址
const state = getQmtConnections(settings);
const active = getActiveQmtConnection(state) ?? { baseUrl: dataSource.baseUrl, name: '(fallback)' };
const baseUrl = String(active.baseUrl ?? '').replace(/\/+$/, '');
let host = '';
let wsUrl = '';
try {
const u = new URL(baseUrl);
host = u.host; // 含端口,如 100.110.38.78:8610
wsUrl = 'ws://' + u.host + '/ws';
} catch { /* baseUrl 非法时返回空 */ }
return { name: active.name, baseUrl, host, wsUrl, activeId: state.activeId ?? '' };
}
default:
throw Object.assign(new Error('unknown qmt method: ' + method), { code: 'not-found' });
}
}
+78
View File
@@ -0,0 +1,78 @@
/**
* 服务端 API:策略 / 份额 / tabs 域(从 api.js 拆分,2026-08-31
*
* 端点:
* strategies / strategies/update / strategies/add / strategies/remove / strategies/move
* strategy-positions / unallocated / summary
* add-shares / move-all-shares / remove-shares / remove-all-shares
* tabs / tabs/update
*/
import {
getStrategies, updateStrategies, addStrategy, removeStrategy, moveStrategy,
getTabConfig, updateTabConfig,
} from '../settings.js';
/** 策略/份额/tabs 端点方法表 */
export const STRATEGY_METHODS = new Set([
'strategies',
'strategies/update',
'strategies/add',
'strategies/remove',
'strategies/move',
'strategy-positions',
'unallocated',
'summary',
'add-shares',
'move-all-shares',
'remove-shares',
'remove-all-shares',
'tabs',
'tabs/update',
]);
/**
* 处理策略/份额/tabs 端点
* @param {string} method
* @param {object} args
* @param {object} runtime { manager, settings }
*/
export async function handleStrategy(method, args, { manager, settings }) {
switch (method) {
case 'strategies':
return getStrategies(settings);
case 'strategies/update':
await updateStrategies(settings, args.strategies);
return getStrategies(settings);
case 'strategies/add':
return await addStrategy(settings, args.name);
case 'strategies/remove':
// 删除策略 + 联动清空该策略下所有份额(份额回未分配)
await removeStrategy(settings, args.strategyId);
await manager.clearStrategyShares(args.strategyId);
return getStrategies(settings);
case 'strategies/move':
return await moveStrategy(settings, args.strategyId, args.dir);
case 'tabs':
return getTabConfig(settings);
case 'tabs/update':
await updateTabConfig(settings, args.tabs);
return getTabConfig(settings);
case 'strategy-positions':
return await manager.getStrategyPositions(args.strategyId);
case 'unallocated':
return await manager.getUnallocatedPositions();
case 'summary':
return await manager.getSummary();
case 'add-shares':
return await manager.addToStrategy(args.code, args.strategyId, args.shares);
case 'move-all-shares':
return await manager.moveAllUnallocatedToStrategy(args.code, args.strategyId);
case 'remove-shares':
return await manager.removeFromStrategy(args.code, args.strategyId, args.shares);
case 'remove-all-shares':
return await manager.clearStrategyShares(args.strategyId);
default:
throw Object.assign(new Error('unknown strategy method: ' + method), { code: 'not-found' });
}
}
+48
View File
@@ -0,0 +1,48 @@
/**
* 服务端 API:交易记录域(R-0072026-09-01
*
* 端点:
* orders → 当日委托(支持 code/status 过滤)
* trades → 当日成交
* trading-dates → 交易日历(日期导航)
*
* 数据流:前端轮询 → /odl/api/orders|trades → dataSource.getOrders/getTrades
* (服务端直连 QMT Bridge REST,技术约束-003
*/
/** 交易记录端点方法表 */
export const TRADE_METHODS = new Set([
'orders',
'trades',
'trading-dates',
]);
/**
* 处理交易记录端点
* @param {string} method
* @param {object} args
* @param {object} runtime { dataSource }
*/
export async function handleTrade(method, args, { dataSource }) {
switch (method) {
case 'orders':
return await dataSource.getOrders({
code: args.code,
status: args.status,
start: args.start,
end: args.end,
});
case 'trades':
return await dataSource.getTrades({
start: args.start,
end: args.end,
});
case 'trading-dates':
return await dataSource.getTradingDates({
start: args.start,
end: args.end,
});
default:
throw Object.assign(new Error('unknown trade method: ' + method), { code: 'not-found' });
}
}