""" WebSocket 事件监听器 监听 VeighNa 事件并推送到 WebSocket 客户端 """ import asyncio import logging from datetime import datetime from typing import Any from vnpy.trader.event import ( EVENT_TICK, EVENT_TRADE, EVENT_ORDER, EVENT_POSITION, EVENT_ACCOUNT, EVENT_LOG, EVENT_CONTRACT ) from vnpy.trader.object import ( TickData, TradeData, OrderData, PositionData, AccountData, LogData, ContractData ) from .manager import manager logger = logging.getLogger(__name__) def serialize_datetime(dt: datetime) -> str: """序列化 datetime 对象为 ISO 格式字符串""" if dt is None: return None return dt.isoformat() def serialize_tick_data(tick: TickData) -> dict: """序列化 TickData 为字典""" return { "vt_symbol": tick.vt_symbol, "symbol": tick.symbol, "exchange": tick.exchange.value if tick.exchange else None, "name": tick.name, "datetime": serialize_datetime(tick.datetime), "localtime": serialize_datetime(tick.localtime), "volume": tick.volume, "turnover": tick.turnover, "open_interest": tick.open_interest, "last_price": tick.last_price, "last_volume": tick.last_volume, "limit_up": tick.limit_up, "limit_down": tick.limit_down, "open_price": tick.open_price, "high_price": tick.high_price, "low_price": tick.low_price, "pre_close": tick.pre_close, "bid_price_1": tick.bid_price_1, "bid_price_2": tick.bid_price_2, "bid_price_3": tick.bid_price_3, "bid_price_4": tick.bid_price_4, "bid_price_5": tick.bid_price_5, "ask_price_1": tick.ask_price_1, "ask_price_2": tick.ask_price_2, "ask_price_3": tick.ask_price_3, "ask_price_4": tick.ask_price_4, "ask_price_5": tick.ask_price_5, "bid_volume_1": tick.bid_volume_1, "bid_volume_2": tick.bid_volume_2, "bid_volume_3": tick.bid_volume_3, "bid_volume_4": tick.bid_volume_4, "bid_volume_5": tick.bid_volume_5, "ask_volume_1": tick.ask_volume_1, "ask_volume_2": tick.ask_volume_2, "ask_volume_3": tick.ask_volume_3, "ask_volume_4": tick.ask_volume_4, "ask_volume_5": tick.ask_volume_5, "gateway_name": tick.gateway_name } def serialize_order_data(order: OrderData) -> dict: """序列化 OrderData 为字典""" return { "vt_orderid": order.vt_orderid, "vt_symbol": order.vt_symbol, "symbol": order.symbol, "exchange": order.exchange.value if order.exchange else None, "orderid": order.orderid, "type": order.type.value if order.type else None, "direction": order.direction.value if order.direction else None, "offset": order.offset.value if order.offset else None, "price": order.price, "volume": order.volume, "traded": order.traded, "status": order.status.value if order.status else None, "datetime": serialize_datetime(order.datetime), "reference": order.reference, "gateway_name": order.gateway_name } def serialize_trade_data(trade: TradeData) -> dict: """序列化 TradeData 为字典""" return { "vt_tradeid": trade.vt_tradeid, "vt_orderid": trade.vt_orderid, "vt_symbol": trade.vt_symbol, "symbol": trade.symbol, "exchange": trade.exchange.value if trade.exchange else None, "orderid": trade.orderid, "tradeid": trade.tradeid, "direction": trade.direction.value if trade.direction else None, "offset": trade.offset.value if trade.offset else None, "price": trade.price, "volume": trade.volume, "datetime": serialize_datetime(trade.datetime), "gateway_name": trade.gateway_name } def serialize_position_data(position: PositionData) -> dict: """序列化 PositionData 为字典""" return { "vt_positionid": position.vt_positionid, "vt_symbol": position.vt_symbol, "symbol": position.symbol, "exchange": position.exchange.value if position.exchange else None, "direction": position.direction.value if position.direction else None, "volume": position.volume, "frozen": position.frozen, "price": position.price, "pnl": position.pnl, "yd_volume": position.yd_volume, "gateway_name": position.gateway_name } def serialize_account_data(account: AccountData) -> dict: """序列化 AccountData 为字典""" return { "vt_accountid": account.vt_accountid, "accountid": account.accountid, "balance": account.balance, "frozen": account.frozen, "available": account.available, "gateway_name": account.gateway_name } def serialize_log_data(log: LogData) -> dict: """序列化 LogData 为字典""" return { "msg": log.msg, "level": log.level, "time": serialize_datetime(log.time), "gateway_name": log.gateway_name } def serialize_contract_data(contract: ContractData) -> dict: """序列化 ContractData 为字典""" return { "vt_symbol": contract.vt_symbol, "symbol": contract.symbol, "exchange": contract.exchange.value if contract.exchange else None, "name": contract.name, "product": contract.product.value if contract.product else None, "size": contract.size, "pricetick": contract.pricetick, "min_volume": contract.min_volume, "max_volume": contract.max_volume, "stop_supported": contract.stop_supported, "net_position": contract.net_position, "history_data": contract.history_data, "option_strike": contract.option_strike, "option_underlying": contract.option_underlying, "option_type": contract.option_type.value if contract.option_type else None, "option_listed": serialize_datetime(contract.option_listed), "option_expiry": serialize_datetime(contract.option_expiry), "option_portfolio": contract.option_portfolio, "option_index": contract.option_index, "gateway_name": contract.gateway_name } class TickEventMonitor: """行情事件监听器""" def __init__(self, event_engine): """ 初始化行情事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_TICK, self.on_tick) logger.info("TickEventMonitor registered") def on_tick(self, event): """处理行情事件""" tick: TickData = event.data try: tick_data = serialize_tick_data(tick) # 在事件循环中异步推送 loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_tick(tick_data)) else: logger.warning("Event loop not running, tick broadcast skipped") except Exception as e: logger.error(f"Error handling tick event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_TICK, self.on_tick) logger.info("TickEventMonitor stopped") class OrderEventMonitor: """订单事件监听器""" def __init__(self, event_engine): """ 初始化订单事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_ORDER, self.on_order) logger.info("OrderEventMonitor registered") def on_order(self, event): """处理订单事件""" order: OrderData = event.data try: order_data = serialize_order_data(order) loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_order(order_data)) else: logger.warning("Event loop not running, order broadcast skipped") except Exception as e: logger.error(f"Error handling order event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_ORDER, self.on_order) logger.info("OrderEventMonitor stopped") class TradeEventMonitor: """成交事件监听器""" def __init__(self, event_engine): """ 初始化成交事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_TRADE, self.on_trade) logger.info("TradeEventMonitor registered") def on_trade(self, event): """处理成交事件""" trade: TradeData = event.data try: trade_data = serialize_trade_data(trade) loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_trade(trade_data)) else: logger.warning("Event loop not running, trade broadcast skipped") except Exception as e: logger.error(f"Error handling trade event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_TRADE, self.on_trade) logger.info("TradeEventMonitor stopped") class PositionEventMonitor: """持仓事件监听器""" def __init__(self, event_engine): """ 初始化持仓事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_POSITION, self.on_position) logger.info("PositionEventMonitor registered") def on_position(self, event): """处理持仓事件""" position: PositionData = event.data try: position_data = serialize_position_data(position) loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_position(position_data)) else: logger.warning("Event loop not running, position broadcast skipped") except Exception as e: logger.error(f"Error handling position event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_POSITION, self.on_position) logger.info("PositionEventMonitor stopped") class AccountEventMonitor: """账户事件监听器""" def __init__(self, event_engine): """ 初始化账户事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_ACCOUNT, self.on_account) logger.info("AccountEventMonitor registered") def on_account(self, event): """处理账户事件""" account: AccountData = event.data try: account_data = serialize_account_data(account) loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_account(account_data)) else: logger.warning("Event loop not running, account broadcast skipped") except Exception as e: logger.error(f"Error handling account event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_ACCOUNT, self.on_account) logger.info("AccountEventMonitor stopped") class LogEventMonitor: """日志事件监听器""" def __init__(self, event_engine): """ 初始化日志事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_LOG, self.on_log) logger.info("LogEventMonitor registered") def on_log(self, event): """处理日志事件""" log: LogData = event.data try: log_data = serialize_log_data(log) loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_log(log_data)) else: logger.warning("Event loop not running, log broadcast skipped") except Exception as e: logger.error(f"Error handling log event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_LOG, self.on_log) logger.info("LogEventMonitor stopped") class ContractEventMonitor: """合约事件监听器""" def __init__(self, event_engine): """ 初始化合约事件监听器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.event_engine.register(EVENT_CONTRACT, self.on_contract) logger.info("ContractEventMonitor registered") def on_contract(self, event): """处理合约事件""" contract: ContractData = event.data try: contract_data = serialize_contract_data(contract) loop = asyncio.get_event_loop() if loop.is_running(): asyncio.create_task(manager.broadcast_contract(contract_data)) else: logger.warning("Event loop not running, contract broadcast skipped") except Exception as e: logger.error(f"Error handling contract event: {e}") def stop(self): """停止监听""" self.event_engine.unregister(EVENT_CONTRACT, self.on_contract) logger.info("ContractEventMonitor stopped") class EventMonitorManager: """事件监听器管理器""" def __init__(self, event_engine): """ 初始化事件监听器管理器 - **event_engine**: VeighNa 事件引擎 """ self.event_engine = event_engine self.monitors = [] def start_all(self): """启动所有事件监听器""" self.monitors = [ TickEventMonitor(self.event_engine), OrderEventMonitor(self.event_engine), TradeEventMonitor(self.event_engine), PositionEventMonitor(self.event_engine), AccountEventMonitor(self.event_engine), LogEventMonitor(self.event_engine), ContractEventMonitor(self.event_engine) ] logger.info("All event monitors started") def stop_all(self): """停止所有事件监听器""" for monitor in self.monitors: monitor.stop() self.monitors = [] logger.info("All event monitors stopped") __all__ = [ "EventMonitorManager", "TickEventMonitor", "OrderEventMonitor", "TradeEventMonitor", "PositionEventMonitor", "AccountEventMonitor", "LogEventMonitor", "ContractEventMonitor" ]