分阶段开发计划 · 测试用例
阶段总览
核心原则:每个阶段结束时,系统必须能独立跑通一个真实场景,不依赖下一阶段的功能。可观测产出优先——
能发出 Discord 通知、能在测试网下单、能看到缓存命中日志,比"代码写完了"更有说服力。
| 阶段 | 核心交付 | 可观测产出 | 完成标准(代码) | 完成标准(功能) |
|---|---|---|---|---|
| P1 | 单系统离线流水线 M1→M2→M3→M4 短路流水线,本地 fixture 数据驱动 |
给定 OHLCV JSON → 控制台打印 TradePlan | pipeline 单元测试 + e2e fixture 测试全绿 | 5 组 fixture(2 有信号、2 短路、1 盈亏比不足),全部结果符合预期 |
| P2 | 实时行情 + 自动调度 CCXT 拉 K 线、K 线收盘触发、Discord 推送 |
无人值守 24h,BTC/USDT 有信号时 Discord 收到完整通知 | mock CCXT + 调度器测试全绿 | 接入真实行情,手动确认一条 Discord 通知内容格式正确 |
| P3 | 自动下单 + 持仓监控 M5A OrderManager + M5B PositionMonitor + M5C ReviewLogger |
测试网完成 3 笔完整交易(入场 + 止盈/止损出场 + 日志归档) | mock exchange 下单 + SL/TP 触发测试全绿 | 测试网跑通 3 笔,ReviewLogger 可查询胜率 |
| P4 | 多系统 + 共享数据层 MultiSystemRunner + MarketDataHub 去重 + TrendResultCache + KeyLevelCache + SymbolStrategyRegistry |
2 套系统并行,日志显示 BTC/4H K 线只拉一次 | 去重 + 缓存失效 + 路由测试全绿 | 2 系统跑 1h,CCXT 调用次数 = 预期次数(无重复拉取) |
| P5 | 风控层 RiskController 熔断 + RevengeController 跨系统配额 |
模拟连续亏损触发熔断,日志显示后续订单被拦截 | 熔断 + 复仇配额测试全绿 | 在测试网模拟 3 连亏,确认熔断后系统停止开仓 |
Phase 1
单系统 · 离线流水线
先让流水线在 fixture 数据上跑通,不涉及任何网络和真实资金
本阶段构建
SystemConfig数据类(基础字段)TrendModuleBase抽象类 + 首个实现(三步曲 EMA)KeyLevelModuleBase+ 首个实现(水平关键位)SignalModuleBase+ 首个实现(金K/针形)PlanModule(固定逻辑,不可插拔)pipeline.run(klines, config)短路编排
MVP 定义 & 完成标准
MVP:调用 pipeline.run(fixture_klines, config),能正确返回 TradePlan 或 None。
完成标准:下列 5 组 fixture 全部结果符合预期,pytest 全绿,无跳过用例。
- BTC 上涨趋势 + 关键位 + 信号 →
TradePlan - BTC 震荡趋势 →
None(M1 短路) - 无关键位 →
None(M2 短路) - 无合格信号 →
None(M3 短路) - 信号存在但盈亏比不足 →
None(M4 短路)
测试用例
tests/test_trend_module.py
import pytest
from alphaquant.modules.trend import SanbuquTrendModule, TrendDirection, TrendStage
from tests.fixtures import load_klines
def test_trend_bullish_returns_correct_fields():
klines = load_klines("btc_4h_bullish_20candles.json")
result = SanbuquTrendModule().analyze(klines, config=base_config)
assert result.direction == TrendDirection.BULLISH
assert result.stage in (TrendStage.IMPULSE, TrendStage.PULLBACK)
assert result.extreme_price > 0
assert 0 <= result.pullback_depth <= 1 # 回调深度是 0~1 的比例
def test_trend_bearish():
klines = load_klines("btc_4h_bearish_20candles.json")
result = SanbuquTrendModule().analyze(klines, config=base_config)
assert result.direction == TrendDirection.BEARISH
def test_trend_ranging_direction():
klines = load_klines("btc_4h_ranging_20candles.json")
result = SanbuquTrendModule().analyze(klines, config=base_config)
assert result.direction == TrendDirection.RANGING
def test_trend_module_is_pluggable():
# 验证任何实现了 TrendModuleBase 的类都能替换
from alphaquant.modules.trend import TrendModuleBase
assert issubclass(SanbuquTrendModule, TrendModuleBase)
tests/test_keylevel_module.py
from alphaquant.modules.keylevel import HorizontalKeyLevelModule
def test_keylevel_returns_list():
klines = load_klines("btc_4h_with_structure.json")
result = HorizontalKeyLevelModule().identify(klines, config=base_config)
assert isinstance(result.key_levels, list)
assert len(result.key_levels) > 0
# 每个关键位是一个区间,有低点和高点
for kl in result.key_levels:
assert kl.low < kl.high
assert kl.touch_count >= 2
def test_keylevel_empty_on_insufficient_data():
klines = load_klines("btc_4h_only_5candles.json") # 数据不足
result = HorizontalKeyLevelModule().identify(klines, config=base_config)
assert result.key_levels == []
tests/test_plan_module.py
from alphaquant.modules.plan import PlanModule
from alphaquant.types import SignalResult, KeyLevel, Direction
def test_plan_rejects_low_rr():
# entry=100, sl=99 (risk=1), tp=102 (reward=2) → RR=2.0,min_rr=2.5 → 拒绝
signal = SignalResult(entry=100.0, sl=99.0, direction=Direction.LONG)
kl = [KeyLevel(low=101.8, high=102.2)] # 最近关键位作为 TP 目标
config = base_config._replace(min_rr_ratio=2.5)
plan = PlanModule().generate(signal, kl, config)
assert plan is None
def test_plan_accepts_good_rr():
# entry=100, sl=99 (risk=1), tp=103.5 (reward=3.5) → RR=3.5 ≥ min_rr=2.0 → 通过
signal = SignalResult(entry=100.0, sl=99.0, direction=Direction.LONG)
kl = [KeyLevel(low=103.2, high=103.8)]
config = base_config._replace(min_rr_ratio=2.0)
plan = PlanModule().generate(signal, kl, config)
assert plan is not None
assert plan.rr_ratio >= 2.0
assert plan.entry_price == 100.0
assert plan.stop_loss == 99.0
def test_plan_position_size_respects_risk_pct():
# 账户 10000 USDT,风险 1%=100 USDT,止损距离 1 USDT → 仓位大小 = 100/1 = 100 合约
signal = SignalResult(entry=100.0, sl=99.0, direction=Direction.LONG)
kl = [KeyLevel(low=103.2, high=103.8)]
config = base_config._replace(account_balance=10000, risk_pct=0.01, min_rr_ratio=2.0)
plan = PlanModule().generate(signal, kl, config)
assert plan.position_size == pytest.approx(100.0, rel=0.01)
tests/test_pipeline.py
from alphaquant.pipeline import run_pipeline
def test_pipeline_returns_plan_when_signal_found():
klines = load_klines("btc_4h_signal_found.json") # 精心构造的有信号数据
plan = run_pipeline(klines, base_config)
assert plan is not None
assert plan.rr_ratio >= base_config.min_rr_ratio
def test_pipeline_short_circuits_on_ranging():
klines = load_klines("btc_4h_ranging.json")
plan = run_pipeline(klines, base_config)
assert plan is None # M1 判断震荡,流水线短路
def test_pipeline_short_circuits_on_no_keylevels():
klines = load_klines("btc_4h_no_structure.json")
plan = run_pipeline(klines, base_config)
assert plan is None # M2 返回空列表,流水线短路
def test_pipeline_short_circuits_on_low_rr():
# fixture 中价格进关键位、形态也对,但 TP 太近,盈亏比不足
klines = load_klines("btc_4h_signal_low_rr.json")
config = base_config._replace(min_rr_ratio=3.0)
plan = run_pipeline(klines, config)
assert plan is None
@pytest.mark.parametrize("fixture,expected_direction", [
("btc_4h_long_signal.json", Direction.LONG),
("btc_4h_short_signal.json", Direction.SHORT),
])
def test_pipeline_direction(fixture, expected_direction):
klines = load_klines(fixture)
plan = run_pipeline(klines, base_config)
assert plan is not None
assert plan.direction == expected_direction
Phase 2
单系统 · 实时行情 + 自动调度
接入 CCXT,K 线收盘自动触发流水线,Discord 推送完整通知
本阶段新增
MarketDataHub(CCXT 拉 K 线,写内存缓存)KLineCloseScheduler(轮询检测收盘时间)DiscordNotifier(发送 webhook 通知)- 入口脚本
main.py
MVP 定义 & 完成标准
MVP:系统无人值守运行,BTC/USDT 4H 每根 K 线收盘后自动分析,发现信号时 Discord 收到包含入场价/止损/止盈/盈亏比的通知。
完成标准:mock 测试全绿 + 手动跑 24h,至少收到 1 条格式正确的 Discord 通知。
测试用例
tests/test_market_data_hub.py
import asyncio
from unittest.mock import AsyncMock, patch
from alphaquant.data import MarketDataHub
async def test_fetch_returns_klines(mock_exchange):
mock_exchange.fetch_ohlcv = AsyncMock(return_value=sample_raw_ohlcv())
hub = MarketDataHub(exchange=mock_exchange)
klines = await hub.get_klines("BTC/USDT", "4h", limit=100)
assert len(klines) == 100
assert klines[0].symbol == "BTC/USDT"
mock_exchange.fetch_ohlcv.assert_called_once()
async def test_hub_caches_second_call(mock_exchange):
mock_exchange.fetch_ohlcv = AsyncMock(return_value=sample_raw_ohlcv())
hub = MarketDataHub(exchange=mock_exchange)
await hub.get_klines("BTC/USDT", "4h")
await hub.get_klines("BTC/USDT", "4h") # 第二次命中缓存
assert mock_exchange.fetch_ohlcv.call_count == 1
tests/test_scheduler.py
from alphaquant.scheduler import KLineCloseScheduler
async def test_scheduler_calls_handler_on_close():
handler_called = asyncio.Event()
scheduler = KLineCloseScheduler(timeframe="4h")
scheduler.on_close = lambda ts: handler_called.set()
# 模拟当前时间恰好是 4H 收盘时刻
with patch("alphaquant.scheduler.utcnow", return_value=candle_close_ts("4h")):
await scheduler._tick()
assert handler_called.is_set()
async def test_scheduler_does_not_double_fire():
call_count = [0]
scheduler = KLineCloseScheduler(timeframe="4h")
scheduler.on_close = lambda ts: call_count.__setitem__(0, call_count[0] + 1)
close_ts = candle_close_ts("4h")
with patch("alphaquant.scheduler.utcnow", return_value=close_ts):
await scheduler._tick()
await scheduler._tick() # 同一根 K 线收盘只触发一次
assert call_count[0] == 1
tests/test_discord_notifier.py
from alphaquant.notifier import DiscordNotifier
async def test_notification_contains_required_fields(httpx_mock):
httpx_mock.add_response(url=WEBHOOK_URL, method="POST", status_code=204)
notifier = DiscordNotifier(webhook_url=WEBHOOK_URL)
await notifier.send_plan(sample_trade_plan())
body = httpx_mock.get_request().content.decode()
import json
msg = json.loads(body)["content"]
assert "BTC" in msg
assert "入场" in msg or "Entry" in msg
assert "止损" in msg or "SL" in msg
assert "RR" in msg or "盈亏比" in msg
async def test_notifier_no_call_when_auto_trade_off(httpx_mock):
# auto_trade=False 时仍然推通知(通知与下单无关)
notifier = DiscordNotifier(webhook_url=WEBHOOK_URL)
await notifier.send_plan(sample_trade_plan(auto_trade=False))
assert len(httpx_mock.get_requests()) == 1 # 通知照样发
Phase 3
单系统 · 自动下单 + 持仓监控
接入真实交易所(测试网),M5A 下单 + M5B WebSocket 监控止盈止损 + M5C 归档
本阶段新增
OrderManager(CCXT 下单,幂等处理)PositionMonitor(WebSocket 价格推送,SL/TP 触发)ReviewLogger(平仓后写 SQLite / JSON)- WebSocket 重连机制
auto_trade开关(False→只通知,True→真实下单)
MVP 定义 & 完成标准
MVP:在测试网完成一笔完整生命周期:信号触发 → 下单成交 → 触达止损/止盈 → 自动平仓 → ReviewLogger 记录。
完成标准:mock 测试全绿 + 测试网至少 3 笔完整交易,ReviewLogger.get_summary() 可查胜率和平均盈亏比。
测试用例
tests/test_order_manager.py
from alphaquant.execution import OrderManager
async def test_order_placed_when_auto_trade_enabled(mock_exchange):
config = base_config._replace(auto_trade=True)
mgr = OrderManager(exchange=mock_exchange, config=config)
await mgr.execute(sample_plan())
mock_exchange.create_order.assert_called_once()
# 验证下单参数正确
call_kwargs = mock_exchange.create_order.call_args[1]
assert call_kwargs["symbol"] == "BTC/USDT"
assert call_kwargs["side"] == "buy"
assert call_kwargs["amount"] == sample_plan().position_size
async def test_no_order_when_auto_trade_disabled(mock_exchange):
config = base_config._replace(auto_trade=False)
mgr = OrderManager(exchange=mock_exchange, config=config)
await mgr.execute(sample_plan())
mock_exchange.create_order.assert_not_called()
async def test_order_manager_idempotent(mock_exchange):
# 同一 plan_id 重复调用,只下一次单
config = base_config._replace(auto_trade=True)
mgr = OrderManager(exchange=mock_exchange, config=config)
plan = sample_plan()
await mgr.execute(plan)
await mgr.execute(plan) # 重复
assert mock_exchange.create_order.call_count == 1
tests/test_position_monitor.py
from alphaquant.execution import PositionMonitor
from alphaquant.types import ActivePosition
async def test_stop_loss_triggers_close(mock_exchange):
monitor = PositionMonitor(exchange=mock_exchange)
pos = ActivePosition(symbol="BTC/USDT", entry=100.0,
sl=98.0, tp=106.0, size=0.1, direction=Direction.LONG)
monitor.track(pos)
await monitor._on_price_tick(97.5) # 低于 SL=98.0
mock_exchange.close_position.assert_called_once_with(pos)
async def test_take_profit_triggers_close(mock_exchange):
monitor = PositionMonitor(exchange=mock_exchange)
pos = ActivePosition(symbol="BTC/USDT", entry=100.0,
sl=98.0, tp=106.0, size=0.1, direction=Direction.LONG)
monitor.track(pos)
await monitor._on_price_tick(106.5) # 高于 TP=106.0
mock_exchange.close_position.assert_called_once_with(pos)
async def test_no_close_within_range(mock_exchange):
monitor = PositionMonitor(exchange=mock_exchange)
pos = ActivePosition(symbol="BTC/USDT", entry=100.0,
sl=98.0, tp=106.0, size=0.1, direction=Direction.LONG)
monitor.track(pos)
for price in [99.5, 100.5, 102.0, 103.8]:
await monitor._on_price_tick(price)
mock_exchange.close_position.assert_not_called()
async def test_websocket_reconnect_on_disconnect(mock_ws):
mock_ws.recv = AsyncMock(side_effect=["price:100", ConnectionError("disconnected"), "price:101"])
monitor = PositionMonitor(exchange=mock_exchange, reconnect_delay=0.01)
# 应在断连后自动重连,不抛异常
await asyncio.wait_for(monitor._ws_loop(), timeout=1.0)
tests/test_review_logger.py
from alphaquant.logger import ReviewLogger
from alphaquant.types import ClosedTrade
def test_record_and_retrieve(tmp_path):
logger = ReviewLogger(storage_path=tmp_path / "trades.db")
trade = ClosedTrade(plan=sample_plan(), actual_entry=100.0,
actual_exit=106.0, close_reason="TP")
logger.record(trade)
records = logger.get_all()
assert len(records) == 1
assert records[0].close_reason == "TP"
def test_summary_win_rate(tmp_path):
logger = ReviewLogger(storage_path=tmp_path / "trades.db")
logger.record(ClosedTrade(close_reason="TP", pnl=60.0, ...))
logger.record(ClosedTrade(close_reason="SL", pnl=-20.0, ...))
logger.record(ClosedTrade(close_reason="TP", pnl=45.0, ...))
summary = logger.get_summary()
assert summary.win_rate == pytest.approx(2/3, abs=0.01)
assert summary.total_pnl == pytest.approx(85.0)
Phase 4
多系统 · 共享数据层
MultiSystemRunner + 四层共享缓存,多套系统并行时每个交易对/方法只计算一次
本阶段新增
MultiSystemRunner(并发调度多套系统)MarketDataHub多系统去重(同 symbol+timeframe 合并请求)TrendResultCache(key = symbol+tf+method+params_hash)KeyLevelCache(key 含 source_close_time)SymbolStrategyRegistry(路由表)
MVP 定义 & 完成标准
MVP:2 套系统(SystemA: BTC 4H EMA,SystemB: BTC 4H 裸K)并行运行,BTC/4H K 线每根收盘只拉一次,EMA 趋势结果只算一次(两系统共享),裸K趋势各算各的。
完成标准:缓存测试全绿 + 日志验证 CCXT call_count 符合预期。
测试用例
tests/test_market_data_hub_dedup.py
async def test_concurrent_requests_deduplicated(mock_exchange):
mock_exchange.fetch_ohlcv = AsyncMock(return_value=sample_raw_ohlcv())
hub = MarketDataHub(exchange=mock_exchange)
# 两套系统同时请求同一数据(并发)
results = await asyncio.gather(
hub.get_klines("BTC/USDT", "4h"),
hub.get_klines("BTC/USDT", "4h"),
)
assert mock_exchange.fetch_ohlcv.call_count == 1 # 只拉一次
assert results[0] is results[1] # 同一对象
async def test_different_timeframes_fetch_separately(mock_exchange):
mock_exchange.fetch_ohlcv = AsyncMock(return_value=sample_raw_ohlcv())
hub = MarketDataHub(exchange=mock_exchange)
await hub.get_klines("BTC/USDT", "4h")
await hub.get_klines("BTC/USDT", "1h") # 不同周期,应各自拉
assert mock_exchange.fetch_ohlcv.call_count == 2
tests/test_trend_result_cache.py
from alphaquant.cache import TrendResultCache
def test_cache_hit_same_key():
cache = TrendResultCache()
key = ("BTC/USDT", "4h", "ema_cross", "abc123")
cache.set(key, trend_result_bullish(), candle_ts=1000)
assert cache.get(key, candle_ts=1000) is not None
def test_cache_miss_different_method():
cache = TrendResultCache()
cache.set(("BTC/USDT", "4h", "ema_cross", "abc123"), trend_result_bullish(), candle_ts=1000)
# 不同方法(裸K vs EMA),不应命中
result = cache.get(("BTC/USDT", "4h", "naked_candle", "abc123"), candle_ts=1000)
assert result is None
def test_cache_invalidated_on_new_candle():
cache = TrendResultCache()
key = ("BTC/USDT", "4h", "ema_cross", "abc123")
cache.set(key, trend_result_bullish(), candle_ts=1000)
# 新 K 线收盘(ts=2000),旧缓存应失效
result = cache.get(key, candle_ts=2000)
assert result is None
tests/test_symbol_registry.py
from alphaquant.registry import SymbolStrategyRegistry
def test_routing_returns_correct_systems():
registry = SymbolStrategyRegistry()
registry.register("BTC/USDT", [system_a, system_b])
registry.register("ETH/USDT", [system_c])
assert registry.get_systems("BTC/USDT") == [system_a, system_b]
assert registry.get_systems("ETH/USDT") == [system_c]
assert registry.get_systems("SOL/USDT") == []
def test_disable_system_for_symbol():
registry = SymbolStrategyRegistry()
registry.register("BTC/USDT", [system_a, system_b])
registry.disable("BTC/USDT", system_b)
assert registry.get_systems("BTC/USDT") == [system_a]
Phase 5
风控层
RiskController 多级熔断 + RevengeController 跨系统复仇配额
本阶段新增
RiskController:单交易对每日亏损上限RiskController:全账户每日亏损上限RevengeController:跨系统复仇配额- 浮亏不加仓硬规则(
PositionMonitor扩展) MultiSystemRunner集成风控拦截
MVP 定义 & 完成标准
MVP:模拟 3 笔连续亏损超出单交易对每日上限,后续下单请求被 RiskController.allow_new_order() 拒绝,日志可见 "CIRCUIT_BREAKER_TRIGGERED"。
完成标准:所有风控测试全绿 + 在测试网上触发一次真实熔断。
测试用例
tests/test_risk_controller.py
from alphaquant.risk import RiskController
def test_symbol_daily_limit_blocks_order():
rc = RiskController(max_daily_loss_per_symbol=100, max_daily_account_loss=500)
rc.record_loss("BTC/USDT", 120) # 超出 100 的上限
assert not rc.allow_new_order("BTC/USDT")
assert rc.allow_new_order("ETH/USDT") # 其他交易对不受影响
def test_account_daily_limit_blocks_all():
rc = RiskController(max_daily_loss_per_symbol=300, max_daily_account_loss=500)
rc.record_loss("BTC/USDT", 300)
rc.record_loss("ETH/USDT", 250) # 累计 550 > 500
assert not rc.allow_new_order("BTC/USDT")
assert not rc.allow_new_order("SOL/USDT") # 全账户熔断,所有交易对冻结
def test_loss_resets_at_day_boundary():
rc = RiskController(max_daily_loss_per_symbol=100, max_daily_account_loss=500)
rc.record_loss("BTC/USDT", 150)
assert not rc.allow_new_order("BTC/USDT")
rc.reset_daily() # 模拟 UTC 00:00 重置
assert rc.allow_new_order("BTC/USDT")
def test_floating_loss_blocks_add_position(mock_exchange):
# 当前持有 BTC 多单浮亏,禁止在同交易对加仓
rc = RiskController(max_daily_loss_per_symbol=100, max_daily_account_loss=500)
rc.register_open_position(symbol="BTC/USDT", unrealized_pnl=-50)
assert not rc.allow_new_order("BTC/USDT", check_floating=True)
tests/test_revenge_controller.py
from alphaquant.risk import RevengeController, RevengeEvent
from alphaquant.types import Direction
def test_first_system_gets_quota():
rc = RevengeController(quota_per_event=1)
event = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=50000, timestamp=1000)
assert rc.try_acquire(event, system="SystemA")
def test_second_system_blocked_same_event():
rc = RevengeController(quota_per_event=1)
base_event = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=50000, timestamp=1000)
similar_event = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=50180, # 偏差 < 0.4%,视为同一事件
timestamp=1000 + 2*3600) # 时间差 < 4h
rc.try_acquire(base_event, system="SystemA")
assert not rc.try_acquire(similar_event, system="SystemB")
def test_different_key_level_gets_own_quota():
rc = RevengeController(quota_per_event=1)
event_a = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=50000, timestamp=1000)
event_b = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=48000, # 偏差 > 0.4%,视为不同关键位
timestamp=1000)
rc.try_acquire(event_a, system="SystemA")
assert rc.try_acquire(event_b, system="SystemB") # 允许,不同关键位
def test_event_quota_expires_after_4h():
rc = RevengeController(quota_per_event=1)
event_a = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=50000, timestamp=1000)
event_b = RevengeEvent(symbol="BTC/USDT", direction=Direction.LONG,
key_level_mid=50050, timestamp=1000 + 5*3600) # 超过 4h
rc.try_acquire(event_a, system="SystemA")
# 超过 4h 时间窗,视为新事件,允许配额
assert rc.try_acquire(event_b, system="SystemB")