Реализация защиты от потери связи (Dead Man’s Switch) для торгового бота на Python

Реализация защиты от потери связи (Dead Man's Switch) для торгового бота на Python ИНВЕСТИЦИИ, ВКЛАДЫ и СБЕРЕЖЕНИЯ
Пошаговое руководство по реализации отказоустойчивого механизма Dead Man's Switch (DMS) для торгового бота на Python с использованием асинхронного UDP-watchdog.
Суть: Реализация отказоустойчивого механизма Dead Man’s Switch (DMS) на Python с использованием асинхронного UDP-watchdog для автоматической отмены активных ордеров и экстренного закрытия позиций при потере связи торгового бота с биржей или сервером.

В высокочастотном и алгоритмическом трейдинге сетевой сбой, зависание процесса или отключение электроэнергии на стороне сервера могут привести к катастрофическим потерям. Если бот успел выставить лимитные ордера, а затем потерял связь, эти ордера остаются на бирже без контроля. Для решения этой проблемы применяется паттерн Dead Man’s Switch (DMS) — предохранитель, который автоматически ликвидирует риски, если бот перестает подавать признаки жизни.

В то время как институциональные игроки используют встроенные механизмы протоколов, о которых мы писали в статье про интеграцию Python с FIX API (QuickFIX) для HFT торговли, для розничных и полупрофессиональных систем на базе REST/WebSocket оптимальным решением является создание независимого локального или удаленного процесса-наблюдателя (Watchdog).

Исходный код

Ниже представлена промышленная реализация архитектуры DMS на Python с использованием библиотеки asyncio. Система состоит из двух компонентов: торгового бота, отправляющего легковесные UDP-хартбиты, и независимого демона-наблюдателя (Watchdog), который мониторит активность и выполняет экстренную отмену ордеров в случае зависания основного процесса.

import asyncio
import socket
import time
import logging

logging.basicConfig(
    level=logging.INFO,
    format='%(asctime)s [%(levelname)s] %(message)s'
)

class DeadMansSwitchWatchdog:
    def __init__(self, host: str = '127.0.0.1', port: int = 9999, timeout: float = 3.0):
        self.host = host
        self.port = port
        self.timeout = timeout
        self.last_heartbeat = time.time()
        self.is_active = True
        self.transport = None

    async def start(self):
        loop = asyncio.get_running_loop()
        self.transport, protocol = await loop.create_datagram_endpoint(
            lambda: WatchdogProtocol(self),
            local_addr=(self.host, self.port)
        )
        logging.info(f'Watchdog запущен на {self.host}:{self.port}')
        
        try:
            await self.monitor_loop()
        finally:
            self.stop()

    async def monitor_loop(self):
        while self.is_active:
            await asyncio.sleep(0.5)
            elapsed = time.time() - self.last_heartbeat
            if elapsed > self.timeout:
                logging.error(f'ВНИМАНИЕ: Хартбит отсутствует {elapsed:.2f} сек! Запуск экстренной отмены...')
                await self.trigger_emergency_action()
                self.is_active = False

    def feed(self):
        self.last_heartbeat = time.time()

    async def trigger_emergency_action(self):
        # Здесь вызывается API биржи для отмены всех ордеров
        # Например, отправка ордеров на Bybit V5
        logging.warning('ЭКСТРЕННЫЙ ВЫЗОВ: Все активные ордера аннулированы, позиции переведены в безопасный режим.')

    def stop(self):
        self.is_active = False
        if self.transport:
            self.transport.close()
            logging.info('Watchdog остановлен.')

class WatchdogProtocol(asyncio.DatagramProtocol):
    def __init__(self, watchdog: DeadMansSwitchWatchdog):
        self.watchdog = watchdog

    def datagram_received(self, data: bytes, addr):
        if data == b'PING':
            self.watchdog.feed()

class TradingBot:
    def __init__(self, watchdog_host: str = '127.0.0.1', watchdog_port: int = 9999, heartbeat_interval: float = 1.0):
        self.watchdog_host = watchdog_host
        self.watchdog_port = watchdog_port
        self.heartbeat_interval = heartbeat_interval
        self.is_running = True

    async def run(self):
        # Используем UDP-сокет для отправки неблокирующих хартбитов
        sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
        logging.info('Торговый бот запущен. Начало отправки хартбитов...')
        
        try:
            while self.is_running:
                sock.sendto(b'PING', (self.watchdog_host, self.watchdog_port))
                logging.info('Хартбит отправлен на Watchdog')
                
                # Имитация работы торговой логики
                await asyncio.sleep(self.heartbeat_interval)
        except asyncio.CancelledError:
            logging.info('Торговый бот остановлен пользователем.')
        finally:
            sock.close()

async def main():
    # Запуск Watchdog и Бота в едином event loop для демонстрации.
    # В продакшене это должны быть два абсолютно независимых процесса!
    watchdog = DeadMansSwitchWatchdog(timeout=3.0)
    bot = TradingBot(heartbeat_interval=1.0)

    # Запускаем задачи асинхронно
    watchdog_task = asyncio.create_task(watchdog.start())
    bot_task = asyncio.create_task(bot.run())

    # Симулируем 5 секунд нормальной работы
    await asyncio.sleep(5.0)
    
    # Симулируем зависание/падение бота
    logging.warning('Симуляция падения торгового бота...')
    bot_task.cancel()

    # Ожидаем срабатывания триггера Watchdog
    await asyncio.sleep(4.0)
    watchdog_task.cancel()

if __name__ == '__main__':
    try:
        asyncio.run(main())
    except KeyboardInterrupt:
        pass

Разбор параметров

  • host и port: Сетевой адрес и порт для UDP-сокета. Использование протокола UDP критически важно, так как он не требует установки соединения (handshake) и не блокирует основной поток выполнения торгового бота при сетевых задержках.
  • timeout: Максимально допустимое время отсутствия сигнала (хартбита) от бота в секундах. Для HFT-систем этот параметр может составлять 0.51.0 секунды, для среднесрочных стратегий — 3.05.0 секунд.
  • heartbeat_interval: Частота отправки сигналов ботом. Она должна быть как минимум в 2-3 раза меньше, чем timeout, чтобы избежать ложных срабатываний из-за кратковременных сетевых джиттеров.
  • trigger_emergency_action: Метод, выполняющий экстренное закрытие рисков. В реальных условиях здесь реализуется отправка запросов на биржу. Например, для минимизации проскальзываний и защиты лимитных заявок можно использовать логику, описанную в статье про создание Post-Only и Reduce-Only ордеров на Bybit V5 в Python.

Как запустить

Для тестирования механизма выполните следующие шаги:

1. Скопируйте представленный код в файл dms_protection.py.

2. Запустите скрипт с помощью интерпретатора Python:

python dms_protection.py

3. В консоли вы увидите логи нормального обмена сообщениями, после чего произойдет симуляция падения бота. Ровно через 3 секунды после последнего хартбита Watchdog зафиксирует аварию и вызовет метод экстренной очистки ордеров.

В промышленной эксплуатации рекомендуется запускать Watchdog на отдельном сервере (co-location), географически близком к бирже. Это гарантирует, что даже при полном отключении дата-центра, где запущен ваш бот, или при падении глобального интернета, внешний наблюдатель сможет оперативно очистить стакан от ваших заявок. Для анализа рыночной конъюнктуры перед перезапуском бота вы также можете использовать внешние метрики, например, настроив интеграцию Python с Glassnode API для получения SOPR и NVT.

Оцените статью
FinFluct