超高頻度(HFT)botの一例
目的
概要
実際に月間30万件程度約定させてたやつ
トレイル機能付きのこれの上位版もある(サルベージ必須)
コードを完璧に読める人にしかこれの価値はわからないので説明はしない
構成を正しく理解した上で当然気づくはずの不足部品を、作者の意図を汲んで寸分の狂いもなく再構築できなければ実行自体不可能なのだが、マジで絶対になにがあろうともこれは実行すべきではない超絶危険なものなので要注意!!!
急騰急落は入出金
今ははるかに効率化したものを作るし実際に作ってあるのだが、これくらい雑なものでも数十万件程度の約定は実現できる、という好例
ほぼ執行オンリー、気合いで殴るタイプ
これを晒したとしてもエッジ的なものは一切消失することはないため、誰にも迷惑をかけることなく私の半年以上前の能力は少しくらいは証明できるはずの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()

コメント