11"""IBKR strategy runner for shared us_equity strategy profiles."""
22import os
3+ import threading
34import time
45import traceback
56from datetime import datetime , timezone
5657app = Flask (__name__ )
5758ensure_event_loop = ibkr_ensure_event_loop
5859NEW_YORK_TZ = ZoneInfo ("America/New_York" )
60+ STRATEGY_RUN_LOCK = threading .Lock ()
5961
6062
6163# ---------------------------------------------------------------------------
@@ -160,6 +162,32 @@ def get_ib_connect_timeout_seconds():
160162 return timeout_seconds
161163
162164
165+ def get_positive_int_env (name , default ):
166+ raw_value = os .getenv (name , str (default ))
167+ try :
168+ parsed = int (raw_value )
169+ except (TypeError , ValueError ):
170+ print (f"Invalid { name } ={ raw_value !r} ; using { default } " , flush = True )
171+ return default
172+ if parsed <= 0 :
173+ print (f"Invalid { name } ={ raw_value !r} ; using { default } " , flush = True )
174+ return default
175+ return parsed
176+
177+
178+ def get_non_negative_float_env (name , default ):
179+ raw_value = os .getenv (name , str (default ))
180+ try :
181+ parsed = float (raw_value )
182+ except (TypeError , ValueError ):
183+ print (f"Invalid { name } ={ raw_value !r} ; using { default } " , flush = True )
184+ return default
185+ if parsed < 0 :
186+ print (f"Invalid { name } ={ raw_value !r} ; using { default } " , flush = True )
187+ return default
188+ return parsed
189+
190+
163191def _env_flag (name : str ) -> bool :
164192 return str (os .getenv (name ) or "" ).strip ().lower () in {"1" , "true" , "yes" , "on" }
165193
@@ -172,6 +200,9 @@ def _env_flag(name: str) -> bool:
172200IB_PORT = get_ib_port ()
173201IB_CLIENT_ID = RUNTIME_SETTINGS .ib_client_id
174202IB_CONNECT_TIMEOUT_SECONDS = get_ib_connect_timeout_seconds ()
203+ IB_CONNECT_ATTEMPTS = get_positive_int_env ("IBKR_CONNECT_ATTEMPTS" , 3 )
204+ IB_CONNECT_RETRY_DELAY_SECONDS = get_non_negative_float_env ("IBKR_CONNECT_RETRY_DELAY_SECONDS" , 5.0 )
205+ IB_CLIENT_ID_RETRY_OFFSET = get_positive_int_env ("IBKR_CLIENT_ID_RETRY_OFFSET" , 100 )
175206STRATEGY_PROFILE = RUNTIME_SETTINGS .strategy_profile
176207STRATEGY_DISPLAY_NAME = RUNTIME_SETTINGS .strategy_display_name
177208ACCOUNT_GROUP = RUNTIME_SETTINGS .account_group
@@ -255,6 +286,8 @@ def t(key, **kwargs):
255286 "strategy_artifact_dir" : RUNTIME_SETTINGS .strategy_artifact_dir ,
256287 "strategy_display_name" : STRATEGY_DISPLAY_NAME ,
257288 "strategy_display_name_localized" : strategy_display_name ,
289+ "ib_connect_attempts" : IB_CONNECT_ATTEMPTS ,
290+ "ib_client_id_retry_offset" : IB_CLIENT_ID_RETRY_OFFSET ,
258291 },
259292)
260293LAST_CYCLE_DETAILS : dict [str , object ] = {}
@@ -271,20 +304,39 @@ def send_tg_message(message):
271304
272305def connect_ib ():
273306 host = get_ib_host ()
274- print (
275- "Connecting to IB gateway "
276- f"{ host } :{ IB_PORT } "
277- f"(mode={ RUNTIME_SETTINGS .ib_gateway_mode } , "
278- f"client_id={ IB_CLIENT_ID } , "
279- f"timeout={ IB_CONNECT_TIMEOUT_SECONDS } s)" ,
280- flush = True ,
281- )
282- return ibkr_connect_ib (
283- host ,
284- IB_PORT ,
285- IB_CLIENT_ID ,
286- timeout = IB_CONNECT_TIMEOUT_SECONDS ,
287- )
307+ last_error = None
308+ for attempt in range (1 , IB_CONNECT_ATTEMPTS + 1 ):
309+ client_id = IB_CLIENT_ID + ((attempt - 1 ) * IB_CLIENT_ID_RETRY_OFFSET )
310+ print (
311+ "Connecting to IB gateway "
312+ f"{ host } :{ IB_PORT } "
313+ f"(mode={ RUNTIME_SETTINGS .ib_gateway_mode } , "
314+ f"client_id={ client_id } , "
315+ f"attempt={ attempt } /{ IB_CONNECT_ATTEMPTS } , "
316+ f"timeout={ IB_CONNECT_TIMEOUT_SECONDS } s)" ,
317+ flush = True ,
318+ )
319+ try :
320+ return ibkr_connect_ib (
321+ host ,
322+ IB_PORT ,
323+ client_id ,
324+ timeout = IB_CONNECT_TIMEOUT_SECONDS ,
325+ )
326+ except (ConnectionError , TimeoutError , OSError ) as exc :
327+ last_error = exc
328+ print (
329+ "IB gateway connection attempt failed "
330+ f"(attempt={ attempt } /{ IB_CONNECT_ATTEMPTS } , "
331+ f"client_id={ client_id } , "
332+ f"error_type={ type (exc ).__name__ } , "
333+ f"error={ exc } )" ,
334+ flush = True ,
335+ )
336+ if attempt < IB_CONNECT_ATTEMPTS and IB_CONNECT_RETRY_DELAY_SECONDS > 0 :
337+ time .sleep (IB_CONNECT_RETRY_DELAY_SECONDS )
338+
339+ raise last_error
288340
289341
290342def log_runtime_event (log_context , event , ** fields ):
@@ -568,13 +620,27 @@ def handle_request():
568620 LAST_CYCLE_DETAILS = {}
569621 log_context = build_request_log_context ()
570622 report = build_execution_report (log_context )
623+ lock_acquired = STRATEGY_RUN_LOCK .acquire (blocking = False )
571624 try :
572625 log_runtime_event (
573626 log_context ,
574627 "strategy_cycle_received" ,
575628 message = "Received strategy execution request" ,
576629 http_method = request .method ,
577630 )
631+ if not lock_acquired :
632+ log_runtime_event (
633+ log_context ,
634+ "strategy_cycle_already_running" ,
635+ message = "Another strategy execution is already running; skip overlapping request" ,
636+ severity = "WARNING" ,
637+ )
638+ finalize_runtime_report (
639+ report ,
640+ status = "skipped" ,
641+ diagnostics = {"skip_reason" : "already_running" },
642+ )
643+ return "Already Running" , 200
578644 if not is_market_open_today ():
579645 log_runtime_event (
580646 log_context ,
@@ -654,6 +720,8 @@ def handle_request():
654720 print (error_msg , flush = True )
655721 return "Error" , 500
656722 finally :
723+ if lock_acquired :
724+ STRATEGY_RUN_LOCK .release ()
657725 try :
658726 report_path = persist_execution_report (report )
659727 print (f"execution_report { report_path } " , flush = True )
0 commit comments