見出し画像

超高頻度(HFT)botの一例

目的


概要

実際に月間30万件程度約定させてたやつ
トレイル機能付きのこれの上位版もある(サルベージ必須)
コードを完璧に読める人にしかこれの価値はわからないので説明はしない
構成を正しく理解した上で当然気づくはずの不足部品を、作者の意図を汲んで寸分の狂いもなく再構築できなければ実行自体不可能なのだが、マジで絶対になにがあろうともこれは実行すべきではない超絶危険なものなので要注意!!!

画像
約2週間で352966回の約定履歴
急騰急落は入出金

今ははるかに効率化したものを作るし実際に作ってあるのだが、これくらい雑なものでも数十万件程度の約定は実現できる、という好例
ほぼ執行オンリー、気合いで殴るタイプ
これを晒したとしてもエッジ的なものは一切消失することはないため、誰にも迷惑をかけることなく私の半年以上前の能力は少しくらいは証明できるはずのpythonでつくったbot本体全文コード


ちなみにガチの履歴なのでOTC業者は余裕で本人特定できるでしょうし、悪いやつら対策&駆逐要員としてどうですか?(フルリモート)
私は叩ける側なので、それなりに役に立つ可能性はありますよね


全文コード中、最も好きな一文はこれ

nextupdate_main = (int(time.time() // (60*int(interval)) + 1) * (60*int(interval)))+1

2021年くらいのpythonはじめましての頃に何気なくtwitterで時間処理についてつぶやいたら、フォロワーさんが教えてくれたもの
あれから5年くらい作り続けているほぼ全てのbotに入っていますよ
その節はありがとうございました

コード

import asyncio
import platform
# Check if the operating system is Windows
if platform.system() == 'Windows':
    asyncio.set_event_loop_policy(asyncio.WindowsSelectorEventLoopPolicy())
import zmq.asyncio
from collections import deque
from typing import Optional, Dict, Any
import sys,os
from decimal import Decimal
import pandas as pd
import time
import datetime
import logging
logger = logging.getLogger()
logger.setLevel(20)
sh = logging.StreamHandler()
logger.addHandler(sh)
fh = logging.FileHandler('log_.log')
logger.addHandler(fh)
formatter = logging.Formatter('%(asctime)s:%(lineno)d:%(levelname)s:%(message)s')
fh.setFormatter(formatter)
sh.setFormatter(formatter)

from dataclasses import dataclass
from typing import Dict, Optional

import json

# UDP socket settings port
import socket
udp_port=6113


HOST = "127.0.0.1"
T_PORT = 5580
OHLCV_PORT = 5581
OHLCV_REQUEST_PORT = 5582

POSITION_PORT = 5583
ACCOUNT_PORT = 5584
ORDER_PORT_FAST = 5585      # Script経由(遅延 ~10-20ms)【MT4のみ】
RESULT_PORT_FAST = 5586     # 【MT4のみ】


###################################################
static_qty=Decimal('0.01')#重要
###################################################

c_n_of=Decimal('0')
#--------------------------------
MORE_T_MODE=False
more_c_n_of=Decimal('0')
more_set_odr_qty=Decimal('0.1')
#--------------------------------
JPY_MODE=False

p_list=[
        {'acc':1,'p_num':'<sub_sim_name>','sym':'<sym_name>','CLS_ONLY':False},
        ]

max_lev=20# 重要:リスク設定
##########################################################
USE_PV=True
DATAONLY=False
PRINT_MODE=False
MAX_MIN=True#

# bal設定
mt_base_bal=Decimal('100000')
limit_bal_per=Decimal('0.40')# 証拠金の何パーセントまで許容できるか

from zoneinfo import ZoneInfo
async def main():

    # 曜日取得(tyベース)
    ima = datetime.datetime.now(tz=ZoneInfo('Asia/Tokyo'))
    logger.info(f'ima:{ima}')
    youbi = ima.weekday()
    now_hour=datetime.datetime.fromtimestamp(time.time(),tz=ZoneInfo('Asia/Tokyo')).hour
    if (youbi==5 and now_hour>=4) or (youbi==6) or (youbi==0 and now_hour<8):
        if youbi==0:
            next_monday=ima
        else:
            days_until_monday =(0 - youbi + 7) % 7
            next_monday = ima + datetime.timedelta(days=days_until_monday)
        
        logger.info(f'next_monday:{next_monday}')
        
        getuyoubi = datetime.datetime(next_monday.year,next_monday.month,next_monday.day,8,30,0,tzinfo=ZoneInfo('Asia/Tokyo'))
        logger.info(f'getuyoubi:{getuyoubi}')
        ts_getuyou= getuyoubi.timestamp()
        taiki_jikan=ts_getuyou-time.time()
        logger.info(f'取引不可、再開まで待機:{taiki_jikan}秒')
        await asyncio.sleep(taiki_jikan)

    sym=p_list[0]['sym']
    p_num=p_list[0]['p_num']
    CLS_ONLY=p_list[0]['CLS_ONLY']

    ##############################################################################
    # mt4 settings
    mtstore = MetaTraderZMQ(use_fast_orders=True)  # MT4の場合 True, MT5の場合 False
    
    await mtstore.connect()

    # バックグラウンドタスク開始
    tasks = [
        asyncio.create_task(mtstore._receive_t()),
        asyncio.create_task(mtstore._receive_position(sym)),
        asyncio.create_task(mtstore._receive_account()),
    ]
    
    if mtstore.use_fast_orders:
        tasks.append(asyncio.create_task(mtstore._receive_order_results_fast()))
    
    # 初期データ待機
    await mtstore.wait_for_data()

    ##############################################################################

    #---
    localstore=LocalDataList(sym=sym,p_num=p_num)

    # jiku用 data
    asyncio.create_task(localstore.udp_client_prc_handler())

    while len(localstore.ask_list)<5 or len(localstore.bid_list)<5:
        await asyncio.sleep(0.01)

    await localstore.set_nortional_data(JPY_MODE)

    asyncio.create_task(main_loop(localstore,sym,CLS_ONLY,mtstore))

    await asyncio.Event().wait()


async def main_loop(localstore,sym,CLS_ONLY,mtstore,interval=1):
    logger.info('main_loop開始')

    while mtstore.bal==None or mtstore.mid==None:
        await asyncio.sleep(0.1)

    ct_permit_time=time.time()+3

    nextupdate_main = (int(time.time() // (60*int(interval)) + 1) * (60*int(interval)))+1
    while True:
        await asyncio.sleep(0)

        if JPY_MODE:
            localstore.mt_max_pos=round((mtstore.bal*max_lev)/mtstore.mid,-3)-(round((mtstore.bal*max_lev)/mtstore.mid,-3)%1000)
            localstore.max_pos=(localstore.mt_max_pos/100000)-((localstore.mt_max_pos/100000)%localstore.min_step)
        else:
            localstore.mt_max_pos=round(mtstore.bal/Decimal('160')*max_lev/mtstore.mid,-3)-(round(mtstore.bal/Decimal('160')*max_lev/mtstore.mid,-3)%1000)
            localstore.max_pos=(localstore.mt_max_pos/100000)-((localstore.mt_max_pos/100000)%localstore.min_step)

        if MAX_MIN:
            localstore.max_pos=static_qty

        # 数量は手動設定、固定である必要あり
        set_odr_qty=static_qty
        cls_set_odr_qty=static_qty

        s_ct_cond=localstore.b_cond_num>c_n_of and localstore.a_cond_num<0
        l_ct_cond=localstore.a_cond_num>c_n_of and localstore.b_cond_num<0

        s_cls_cond=mtstore.pos_size>0
        l_cls_cond=mtstore.pos_size<0

        if MORE_T_MODE:
            more_s_ct_cond=localstore.b_cond_num>more_c_n_of and localstore.a_cond_num<0
            more_l_ct_cond=localstore.a_cond_num>more_c_n_of and localstore.b_cond_num<0
        else:
            more_s_ct_cond=False
            more_l_ct_cond=False

        ############################################################################
        use_bal_ratio=mtstore.bal_ratio

        if USE_PV:
            ALL_STOP=False
            # 最大発注制限
            TRADE_MODE=True
            if abs(mtstore.pos_size)>=localstore.max_pos:
                TRADE_MODE=False
        else:
            ALL_STOP=True
            TRADE_MODE=False

        # spread limit
        use_spread_limit=True
        sp_limit_cnt=0
        if use_spread_limit:
            if mtstore.sp>localstore.min_t*sp_limit_cnt:
                TRADE_MODE=False

        ############################################################################
        # 生存確認用ログ
        if time.time()>nextupdate_main:
            logger.info(f'[{sym}]定期生存確認:[prc]l:({localstore.udp_ask} {localstore.udp_bid}) b:({mtstore.ask} {mtstore.bid})')
            logger.info(f'[{sym}]s_ct_cond:{s_ct_cond} l_ct_cond:{l_ct_cond} s_cls_cond:{s_cls_cond} l_cls_cond:{l_cls_cond}')
            
            if mtstore.pl!=None and localstore.lt_pl!=None:
                logger.info(f'!!!!! - [{sym}]bal:[b:l]({mtstore.bal}) avail:[b:l]({mtstore.avail_bal}) pos:[b:l]({mtstore.pos_size}) pos_prc:({mtstore.pos_prc}) pl:[b:l]({mtstore.pl})::: set_odr_qty -> {set_odr_qty} cls_set_odr_qty -> {cls_set_odr_qty} max_pos -> {localstore.max_pos} ::: bal_ratio [{mtstore.bal_ratio}]')
                logger.info(f'[{sym}] ||| MODE Cond ||| ALL_STOP:{ALL_STOP} TRADE_MODE:{TRADE_MODE} CT_PERMIT:{localstore.CT_PERMIT}||| ')           
            else:
                logger.info(f'!!!!! - [{sym}]bal:[b:l]({mtstore.bal}) avail:[b:l]({mtstore.avail_bal}) pos:[b:l]({mtstore.pos_size}) pos_prc:({mtstore.pos_prc}) pl:[b:l]({mtstore.pl})::: set_odr_qty -> {set_odr_qty} cls_set_odr_qty -> {cls_set_odr_qty} max_pos -> {localstore.max_pos} ::: bal_ratio [{mtstore.bal_ratio}]')
                logger.info(f'[{sym}] ||| MODE Cond ||| ALL_STOP:{ALL_STOP} TRADE_MODE:{TRADE_MODE} CT_PERMIT:{localstore.CT_PERMIT}||| ')           
            nextupdate_main = (int(time.time() // (60*int(interval)) + 1) * (60*int(interval)))+1

        ############################################################################

        if not ALL_STOP and not DATAONLY:

            ############################################################################
            # 発注処理 Create & Close
            ############################################################################
            cls_pos_cond=abs(mtstore.pos_size)>=Decimal(f"0.01")

            if time.time()>ct_permit_time:
                localstore.CT_PERMIT=True
            ##############################################################################
            if mtstore.pos_size>0 and s_cls_cond and cls_set_odr_qty!=0 and cls_pos_cond:
                logger.info(f'[{sym}]Upper LONG Close処理開始:T -> {localstore.b_cond_num} | now prc -> {mtstore.bid}')

                main_side='SELL'

                ticket=mtstore.ticket_list[0]['ticket']
                odr_size=mtstore.ticket_list[0]['size']
                mtodr_res=await mtstore.place_order(cmd="CLOSE",symbol=sym,volume=float(odr_size),use_fast=True,ticket=ticket)
                logger.info(f'order sended:{mtodr_res}')
                await asyncio.sleep(0.1)
                localstore.CT_PERMIT=True

            elif localstore.CT_PERMIT and not localstore.STOP_ALL_TRADE and not CLS_ONLY and mtstore.pos_size>=0 and TRADE_MODE and s_ct_cond and use_bal_ratio>limit_bal_per and set_odr_qty!=0:
                logger.info(f'[{sym}]Upper LONG Create処理開始:T -> {localstore.b_cond_num} | now sp -> {mtstore.ask}')

                main_side='BUY'

                if more_s_ct_cond:
                    mtodr_res=await mtstore.place_order(cmd=main_side,symbol=sym,volume=float(more_set_odr_qty),use_fast=True)
                else:
                    mtodr_res=await mtstore.place_order(cmd=main_side,symbol=sym,volume=float(set_odr_qty),use_fast=True)
                logger.info(f'order sended:{mtodr_res}')
                await asyncio.sleep(0.1)
                ct_permit_time=time.time()+3
                localstore.CT_PERMIT=False

            if mtstore.pos_size<0 and l_cls_cond and cls_set_odr_qty!=0 and cls_pos_cond:
                logger.info(f'[{sym}]Lower SHORT Close処理開始:T -> {localstore.a_cond_num} | now sp -> {mtstore.ask}')

                main_side='BUY'

                ticket=mtstore.ticket_list[0]['ticket']
                odr_size=mtstore.ticket_list[0]['size']

                mtodr_res=await mtstore.place_order(cmd="CLOSE",symbol=sym,volume=float(odr_size),use_fast=True,ticket=ticket)
                logger.info(f'order sended:{mtodr_res}')
                await asyncio.sleep(0.1)
                localstore.CT_PERMIT=True

            elif localstore.CT_PERMIT and not localstore.STOP_ALL_TRADE and not CLS_ONLY and mtstore.pos_size<=0 and TRADE_MODE and l_ct_cond and use_bal_ratio>limit_bal_per and set_odr_qty!=0:
                logger.info(f'[{sym}]LOWER SHORT Create処理開始:T -> {localstore.a_cond_num} | now sp -> {mtstore.bid}')

                main_side='SELL'
                if more_l_ct_cond:
                    mtodr_res=await mtstore.place_order(cmd=main_side,symbol=sym,volume=float(more_set_odr_qty),use_fast=True)
                else:
                    mtodr_res=await mtstore.place_order(cmd=main_side,symbol=sym,volume=float(set_odr_qty),use_fast=True)
                logger.info(f'order sended:{mtodr_res}')
                await asyncio.sleep(0.1)
                ct_permit_time=time.time()+3
                localstore.CT_PERMIT=False

            ##############################################################################
            
            
            ##############################################################################
                

        ############################################################################
        
        if PRINT_MODE:
           sys.stdout.flush()
           print("\n"
               + "\n"
               + "\n"
               + str(mtstore.ask)+" "+str(mtstore.bid)+" "+str(mtstore.sp)+" "
               + "\n"
               + "\n"
               + "----------------------------------------------------------------"
               + "\n"
               + "\n"
               + "\n"
               + "\n"
               + "\n"
               + "----------------------------------------------------------------"
               + "\n"
               + "\n"
               + "\n"
               + str("s_ct_cond")+" "+str(s_ct_cond)+" "+str("l_ct_cond")+" "+str(l_ct_cond)+" "
               + "\n"
               + str('pos_size')+" "+str(mtstore.pos_size)+" "+ str('pos_prc')+" "+str(mtstore.pos_prc)+"                             "
               + "\n"
               + "\n"
               + "\n"
               + "\n"
               + str('TRADE_MODE')+":"+str(TRADE_MODE)+"                             "
               + "\033[18A",end="")

        
        await asyncio.sleep(0.01)

@dataclass
class LocalDataList:
    def __init__(self,sym=None,p_num=None):

        self.ask_list,self.bid_list=[],[]
        self.sym=sym
        self.p_num=p_num
        self.udp_ask,self.udp_bid=None,None

        # nortional data
        self.min_nortional,self.min_step,self.min_t=None,None,None

        # etc hensuu
        self.max_pos=0
        self.mt_max_pos
        self.set_odr_qty=0
        self.CT_PERMIT=True

        # jikoku seigen
        self.STOP_ALL_TRADE=False
        #asyncio.create_task(self.jikoku_seigen_task())

    async def jikoku_seigen_task(self):
        from zoneinfo import ZoneInfo
        no_trade_hour_st=3
        no_trade_hour_ed=7

        while True:
            now_hour=datetime.datetime.fromtimestamp(time.time(),tz=ZoneInfo('Asia/Tokyo')).hour
            if no_trade_hour_st<=now_hour<=no_trade_hour_ed:
                self.STOP_ALL_TRADE=True
                logger.info('取引不可能時刻のためすべての処理を停止')
            else:
                self.STOP_ALL_TRADE=False
                logger.info('時刻に問題なし')

            nextupdate = (int(time.time() // (60*int(60)) + 1) * (60*int(60)))
            sleep_time=nextupdate-time.time()
            logger.info(f'今{now_hour}時、次回時刻確認まで sleep:{sleep_time}秒')
            await asyncio.sleep(sleep_time)

    async def set_nortional_data(self,JPY_MODE):
        logger.info('set_nortional start')

        #- mt 一時手動設定(EURUSD)
        if JPY_MODE:
            self.min_nortional=Decimal("0.01")
            self.min_step=Decimal("0.01")
            self.min_t=Decimal("0.001")
        else:
            self.min_nortional=Decimal("0.01")
            self.min_step=Decimal("0.01")
            self.min_t=Decimal("0.00001")

    async def udp_client_prc_handler(self):

        udp = AsyncUDPReceiver(udp_port)
        await udp.start()
        
        while True:
            await asyncio.sleep(0.01)

            try:
                data = udp.c_data

                self.udp_ask=Decimal(f"{data['data1']}")
                self.udp_bid=Decimal(f"{data['data2']}")
                self.udp_sp=self.udp_ask-self.udp_bid
                self.udp_mid=(self.udp_ask+self.udp_bid)/2

                if 00000 and 00000:
                    self.a_cond_num=00000
                    self.b_cond_num=00000
                else:
                    self.a_cond_num=0
                    self.b_cond_num=0

            except:
                pass

class AsyncUDPReceiver:
    def __init__(self, port: int = 6111):
        self.port = port
        self.sock = None
        
        # 最新の受信データ
        self.c_data = None
        self.sensor_data = None
        # 必要に応じて追加
        
    async def start(self):
        """受信開始"""
        self.sock = socket.socket(socket.AF_INET, socket.SOCK_DGRAM)
        self.sock.bind(("0.0.0.0", self.port))
        self.sock.setblocking(False)
        
        # 受信タスク開始
        asyncio.create_task(self._receive())
        
    async def _receive(self):
        """受信処理(データが来た時だけ動く)"""
        loop = asyncio.get_event_loop()
        
        while True:
            raw_data, _ = await loop.sock_recvfrom(self.sock, 65535)
            
            # パース
            message = raw_data.decode()
            topic, json_str = message.split(' ', 1)
            data = json.loads(json_str)
            
            # トピック別に格納
            if topic == "c_data":
                self.c_data = data
            elif topic == "sensor":
                self.sensor_data = data

########################################################################################################################
########################################################################################################################
########################################################################################################################

# データ格納
@dataclass
class MarketData:
    symbol: str
    bid: float
    ask: float
    volume: int
    time: int
    timestamp: float = None
    
    def __post_init__(self):
        if self.timestamp is None:
            self.timestamp = asyncio.get_event_loop().time()

@dataclass
class Position:
    ticket: int
    symbol: str
    type: int
    volume: float
    entry_price: float
    current_price: float
    profit: float
    profit_pips: float
    sl: float
    tp: float
    open_time: int

@dataclass
class AccountInfo:
    balance: float
    equity: float
    margin: float
    margin_free: float
    margin_level: float
    profit: float
    timestamp: int

@dataclass
class OrderResult:
    ticket: int
    status: str
    message: str
    time: int
    received_at: float = None
    
    def __post_init__(self):
        if self.received_at is None:
            self.received_at = asyncio.get_event_loop().time()

class MetaTraderZMQ:
    """MetaTrader ZMQ 統合クラス"""
    
    def __init__(self, host: str = HOST, use_fast_orders: bool = True):
        self.host = host
        self.use_fast_orders = use_fast_orders
        
        # データ格納
        self.latest_t: Optional[MarketData] = None
        self.latest_positions: list[Position] = []
        self.latest_account: Optional[AccountInfo] = None
        
        # 発注追跡
        self.pending_orders: Dict[int, asyncio.Future] = {}
        self.order_timeout = 5.0  # 秒
        
        # パフォーマンス測定
        self.t_latencies = deque(maxlen=100)
        self.position_latencies = deque(maxlen=100)
        self.account_latencies = deque(maxlen=100)
        
        # ZMQ コンテキスト
        self.context: Optional[zmq.asyncio.Context] = None
        self.sockets: Dict[str, zmq.asyncio.Socket] = {}

        # mt基本データ
        self.vol_to_int=100000
        
        self.ask,self.bid,self.mid,self.sp=None,None,None,None
        self.pos_size,self.pos_prc=Decimal('0'),None
        self.real_pos_size=Decimal('0')
        self.ticket_list=[]
        self.bal,self.avail_bal,self.pl=None,None,None
        self.bal_ratio=None

        self.MTODR_PERMIT=True
        
        logger.info("[INFO] MetaTrader ZMQ initialized")
    
    async def request_ohlcv(self, symbol, period=1, count=5000):
        """Request OHLCV data"""
        socket = self.sockets['req_ohlcv']
        request = {
            "symbol": symbol,
            "period": period,
            "count": count
        }
        await socket.send_string(json.dumps(request))
        logger.info(f"[OHLCV REQUEST] {symbol} {period}min {count} bars")

        await self.receive_ohlcv()
        
        # レスポンス待機(3秒のタイムアウト)
        #import time
        #bk_time=time.time()+60
        #while not self.ohlcv_data:
        #    if time.time()>bk_time:
        #        break
        #    await asyncio.sleep(0.5)

        if self.ohlcv_data:
            res_df=pd.DataFrame(self.ohlcv_data.get('bars', []))
            res_df.columns=['0','1','2','3','4','5']
            res_df['0']=(res_df['0']-60*60*2)*1000
            res_df['0']=res_df['0'].astype('int64')
            return res_df
        else:
            logger.info('ohlcv取得失敗')
            sys.exit()
            #return None
    
    async def receive_ohlcv(self):
        """Receive OHLCV data"""
        socket = self.sockets['ohlcv']

        try:
            msg = await socket.recv_string()
            data = json.loads(msg)
            if data.get('t') == 'OHLCV_DATA':
                self.ohlcv_data = data
                logger.info(f"[OHLCV] Received {data.get('count')} bars for {data.get('symbol')}")
                return self.ohlcv_data
        except Exception as e:
            logger.info(f"[OHLCV ERROR] {e}")
            sys.exit()
    
    async def connect(self):
        """全ての接続を確立"""
        self.context = zmq.asyncio.Context()
        
        try:
            self.sockets['ohlcv'] = self.context.socket(zmq.SUB)
            self.sockets['ohlcv'].connect(f"tcp://{HOST}:{OHLCV_PORT}")
            self.sockets['ohlcv'].setsockopt(zmq.SUBSCRIBE, b"")

            self.sockets['req_ohlcv'] = self.context.socket(zmq.PUSH)
            self.sockets['req_ohlcv'].connect(f"tcp://{HOST}:{OHLCV_REQUEST_PORT}")

            # TICK データ受信
            self.sockets['tick'] = self.context.socket(zmq.SUB)
            self.sockets['tick'].connect(f"tcp://{self.host}:{T_PORT}")
            self.sockets['tick'].setsockopt(zmq.SUBSCRIBE, b"")
            logger.info(f"[TICK] Connected to tcp://{self.host}:{T_PORT}")
            
            # POSITION データ受信
            self.sockets['position'] = self.context.socket(zmq.SUB)
            self.sockets['position'].connect(f"tcp://{self.host}:{POSITION_PORT}")
            self.sockets['position'].setsockopt(zmq.SUBSCRIBE, b"")
            logger.info(f"[POSITION] Connected to tcp://{self.host}:{POSITION_PORT}")
            
            # ACCOUNT データ受信
            self.sockets['account'] = self.context.socket(zmq.SUB)
            self.sockets['account'].connect(f"tcp://{self.host}:{ACCOUNT_PORT}")
            self.sockets['account'].setsockopt(zmq.SUBSCRIBE, b"")
            logger.info(f"[ACCOUNT] Connected to tcp://{self.host}:{ACCOUNT_PORT}")
            
            # 発注送信(通常)
            ###self.sockets['order_slow'] = self.context.socket(zmq.PUSH)
            ###self.sockets['order_slow'].connect(f"tcp://{self.host}:{ORDER_PORT_SLOW}")
            ###logger.info(f"[ORDER_SLOW] Connected to tcp://{self.host}:{ORDER_PORT_SLOW}")
            ###
            #### 発注結果受信(通常)
            ###self.sockets['result_slow'] = self.context.socket(zmq.SUB)
            ###self.sockets['result_slow'].connect(f"tcp://{self.host}:{RESULT_PORT_SLOW}")
            ###self.sockets['result_slow'].setsockopt(zmq.SUBSCRIBE, b"")
            ###logger.info(f"[RESULT_SLOW] Connected to tcp://{self.host}:{RESULT_PORT_SLOW}")
            
            # 発注送信(高速)【MT4のみ】
            if self.use_fast_orders:
                self.sockets['order_fast'] = self.context.socket(zmq.PUSH)
                self.sockets['order_fast'].connect(f"tcp://{self.host}:{ORDER_PORT_FAST}")
                logger.info(f"[ORDER_FAST] Connected to tcp://{self.host}:{ORDER_PORT_FAST}")
                
                # 発注結果受信(高速)
                self.sockets['result_fast'] = self.context.socket(zmq.SUB)
                self.sockets['result_fast'].connect(f"tcp://{self.host}:{RESULT_PORT_FAST}")
                self.sockets['result_fast'].setsockopt(zmq.SUBSCRIBE, b"")
                logger.info(f"[RESULT_FAST] Connected to tcp://{self.host}:{RESULT_PORT_FAST}")
        
        except Exception as e:
            logger.info(f"[ERROR] Connection failed: {e}")
            raise
    
    async def disconnect(self):
        """全ての接続を閉じる"""
        for name, socket in self.sockets.items():
            try:
                socket.close()
            except:
                pass
        
        if self.context:
            self.context.term()
        
        logger.info("[INFO] All connections closed")
    
    async def _receive_t(self):
        """TICK データを継続受信"""
        socket = self.sockets['tick']
        
        while True:
            try:
                msg = await socket.recv_string()
                data = json.loads(msg)
                
                if data.get('t') == 'TICK':
                    self.ask=Decimal(f"{data.get('a')}")
                    self.bid=Decimal(f"{data.get('b')}")
                    self.mid=(self.ask+self.bid)/2
                    self.sp =self.ask-self.bid

                    self.latest_t = MarketData(
                        symbol=data.get('s'),
                        bid=float(data.get('b')),
                        ask=float(data.get('a')),
                        volume=int(data.get('v')),
                        time=int(data.get('time'))
                    )
                    
                    # レイテンシ計測
                    #latency = (asyncio.get_event_loop().time() - self.latest_t.timestamp) * 1000
                    #self.t_latencies.append(latency)
            
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.info(f"[TICK ERROR] {e}")
                await asyncio.sleep(0.1)
    
    async def _receive_position(self,sym):
        """POSITION データを継続受信"""
        socket = self.sockets['position']
        
        while True:
            try:
                msg = await socket.recv_string()
                data = json.loads(msg)
                
                if data.get('t') == 'POSITIONS_SNAPSHOT':
                    positions_data = data.get('positions', [])
                    self.latest_positions = [
                        Position(
                            ticket=int(p.get('ticket')),
                            symbol=p.get('symbol'),
                            type=int(p.get('type')),
                            volume=float(p.get('volume')),
                            entry_price=float(p.get('entry_price')),
                            current_price=float(p.get('current_price')),
                            profit=float(p.get('profit')),
                            profit_pips=float(p.get('profit_pips')),
                            sl=float(p.get('sl')),
                            tp=float(p.get('tp')),
                            open_time=int(p.get('open_time'))
                        )
                        for p in positions_data
                    ]

                    # pos
                    if len([x for x in self.latest_positions if x.symbol==sym]):
                        self.pos_size=Decimal(f"{sum([Decimal(f'{x.volume}') if x.type==0 else Decimal(f'{-x.volume}') for x in self.latest_positions if x.symbol==sym])}")
                        #self.pos_size=Decimal(f"{sum([x.volume*self.vol_to_int if x.type==0 else -x.volume*self.vol_to_int for x in self.latest_positions if x.symbol==sym])}")
                        self.pos_prc =sum([Decimal(f"{x.entry_price}")*Decimal(f"{x.volume}") for x in self.latest_positions if x.symbol==sym])/sum([Decimal(f"{x.volume}") for x in self.latest_positions if x.symbol==sym])
                        self.ticket_list=[{'ticket':x.ticket,'side':x.type,'size':x.volume} for x in self.latest_positions if x.symbol==sym]
                        self.real_pos_size=Decimal(f'{sum([Decimal(f"{x.volume}") for x in self.latest_positions if x.symbol==sym])}')
                    else:
                        self.pos_size=Decimal('0')
                        self.pos_prc =None
                        self.ticket_list=[]
                        self.real_pos_size=Decimal('0')

                    # レイテンシ計測
                    #latency = (asyncio.get_event_loop().time() - data.get('timestamp', 0)) * 1000
                    #self.position_latencies.append(latency)
            
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.info(f"[POSITION ERROR] {e}")
                await asyncio.sleep(0.1)
    
    async def _receive_account(self):
        """ACCOUNT データを継続受信"""
        socket = self.sockets['account']
        
        while True:
            try:
                msg = await socket.recv_string()
                data = json.loads(msg)
                
                if data.get('t') == 'ACCOUNT_INFO':

                    self.bal=Decimal(f"{data.get('balance')}")
                    self.avail_bal=Decimal(f"{data.get('margin_free')}")

                    #self.pl=self.bal-mt_base_bal

                    if self.bal!=0:
                        self.bal_ratio=round(self.avail_bal/self.bal,3)
                    else:
                        self.bal_ratio=0

                    #self.latest_account = AccountInfo(
                    #    balance=float(data.get('balance')),
                    #    equity=float(data.get('equity')),
                    #    margin=float(data.get('margin')),
                    #    margin_free=float(data.get('margin_free')),
                    #    margin_level=float(data.get('margin_level')),
                    #    profit=float(data.get('profit')),
                    #    timestamp=int(data.get('timestamp'))
                    #)
                    
                    # レイテンシ計測
                    #latency = (asyncio.get_event_loop().time() - self.latest_account.timestamp) * 1000
                    #self.account_latencies.append(latency)
            
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.info(f"[ACCOUNT ERROR] {e}")
                await asyncio.sleep(0.1)
    
    async def _receive_order_results_slow(self):
        """発注結果を継続受信(通常)"""
        socket = self.sockets['result_slow']
        
        while True:
            try:
                msg = await socket.recv_string()
                data = json.loads(msg)
                
                if data.get('t') == 'ORDER_RESULT':
                    ticket = int(data.get('ticket'))
                    result = OrderResult(
                        ticket=ticket,
                        status=data.get('status'),
                        message=data.get('message'),
                        time=int(data.get('time'))
                    )
                    
                    # ペンディング発注の完了を通知
                    if ticket in self.pending_orders:
                        future = self.pending_orders.pop(ticket)
                        if not future.done():
                            future.set_result(result)
            
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.info(f"[RESULT_SLOW ERROR] {e}")
                await asyncio.sleep(0.1)
    
    async def _receive_order_results_fast(self):
        """発注結果を継続受信(高速)【MT4のみ】"""
        if not self.use_fast_orders:
            return
        
        socket = self.sockets['result_fast']
        
        while True:
            try:
                msg = await socket.recv_string()
                data = json.loads(msg)
                
                if data.get('t') == 'ORDER_RESULT':
                    
                    self.MTODR_PERMIT=True

                    ticket = int(data.get('ticket'))
                    result = OrderResult(
                        ticket=ticket,
                        status=data.get('status'),
                        message=data.get('message'),
                        time=int(data.get('time'))
                    )
                    
                    # ペンディング発注の完了を通知
                    if ticket in self.pending_orders:
                        future = self.pending_orders.pop(ticket)
                        if not future.done():
                            future.set_result(result)
            
            except asyncio.CancelledError:
                break
            except Exception as e:
                logger.info(f"[RESULT_FAST ERROR] {e}")
                await asyncio.sleep(0.1)
    
    async def place_order(self, cmd: str, symbol: str, volume: float, 
                         sl: float = 0.0, tp: float = 0.0, 
                         ticket: int = 0, use_fast: bool = None) -> OrderResult:
        """
        発注を送信して結果を待つ

        # CT(BUY or SELL)
        await mtstore.place_order(cmd="BUY",symbol=sym,volume=0.01,use_fast=True)
        # CLS(CLOSE)
        await mtstore.place_order(cmd="CLOSE",symbol=sym,volume=0.01,use_fast=True,ticket=ticket)
        
        Args:
            cmd: "BUY", "SELL", "CLOSE"
            symbol: シンボル(例: "EURUSD-cd")
            volume: ロット数
            sl: ストップロス
            tp: テイクプロフィット
            ticket: CLOSE時のチケット番号
            use_fast: MT4高速発注を使用(Noneの場合は設定値に従う)
        
        Returns:
            OrderResult: 発注結果
        """
        # use_fast を決定
        if use_fast is None:
            use_fast = self.use_fast_orders
        
        order_data = {
            "cmd": cmd,
            "symbol": symbol,
            "volume": volume,
            "sl": sl,
            "tp": tp,
        }
        
        if cmd == "CLOSE" and ticket > 0:
            order_data["ticket"] = ticket
        
        # ペンディング発注を追跡するための Future を作成
        # チケット番号は未知なので、タイムアウト機能を使う
        #result_future = asyncio.Future()
        
        try:
            # 発注を送信
            if use_fast and 'order_fast' in self.sockets:
                await self.sockets['order_fast'].send_string(json.dumps(order_data))
                logger.info(f"[ORDER_FAST] {cmd} {symbol} {volume}x (SL:{sl}, TP:{tp})")
                self.MTODR_PERMIT=False
            else:
                await self.sockets['order_slow'].send_string(json.dumps(order_data))
                logger.info(f"[ORDER_SLOW] {cmd} {symbol} {volume}x (SL:{sl}, TP:{tp})")
                self.MTODR_PERMIT=False
            
            # 結果を待つ(タイムアウト付き)
            #try:
            #    # すべてのペンディング発注を確認(チケット番号で追跡)
            #    # 最新の結果を返す
            #    await asyncio.wait_for(asyncio.sleep(self.order_timeout), timeout=self.order_timeout)
            #except asyncio.TimeoutError:
            #    return OrderResult(
            #        ticket=0,
            #        status="TIMEOUT",
            #        message="Order execution timeout",
            #        time=int(asyncio.get_event_loop().time())
            #    )
        
        except Exception as e:
            logger.info(f"[ORDER ERROR] {e}")
            return OrderResult(
                ticket=0,
                status="FAILED",
                message=str(e),
                time=int(asyncio.get_event_loop().time())
            )
        
        # 最後に受け取った結果を返す
        return OrderResult(
            ticket=0,
            status="OK",
            message="Order sent",
            time=int(asyncio.get_event_loop().time())
        )
    
    async def wait_for_data(self, timeout: float = 5.0):
        """初期データ到着を待つ"""
        start_time = asyncio.get_event_loop().time()
        
        while asyncio.get_event_loop().time() - start_time < timeout:
            if self.latest_t and self.latest_positions and self.latest_account:
                logger.info("[INFO] All data received")
                return True
            await asyncio.sleep(0.1)
        
        logger.info("[WARNING] Timeout waiting for initial data")
        return False
    
    def get_stats(self) -> Dict[str, Any]:
        """統計情報を取得"""
        return {
            't_latency_avg': sum(self.t_latencies) / len(self.t_latencies) if self.t_latencies else 0,
            't_latency_max': max(self.t_latencies) if self.t_latencies else 0,
            'position_latency_avg': sum(self.position_latencies) / len(self.position_latencies) if self.position_latencies else 0,
            'account_latency_avg': sum(self.account_latencies) / len(self.account_latencies) if self.account_latencies else 0,
        }

########################################################################################################################
########################################################################################################################
########################################################################################################################

if __name__ == '__main__':

    logger.info('!!![START]!!!')
    # hennsu

    pyname=os.path.basename(__file__)
    logger.info(f'{pyname}:START')
    try:
        loop = asyncio.get_event_loop()
        loop.run_until_complete(main())
    except KeyboardInterrupt:
        sys.exit()
    except Exception as e:
        logger.info(e)
        sys.exit()


画像
うちの犬 having fun on The Mactan Newtown Beach

いいなと思ったら応援しよう!

コメント

コメントするには、 ログイン または 会員登録 をお願いします。
超高頻度(HFT)botの一例|bufujini
word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word word

mmMwWLliI0fiflO&1
mmMwWLliI0fiflO&1
mmMwWLliI0fiflO&1
mmMwWLliI0fiflO&1
mmMwWLliI0fiflO&1
mmMwWLliI0fiflO&1
mmMwWLliI0fiflO&1