From cfe00f67b54c41a400b6d8a49167d4502190e1f8 Mon Sep 17 00:00:00 2001 From: Joe Tretter Date: Wed, 1 Apr 2026 23:28:51 -0500 Subject: [PATCH] Automatic flex query syncing --- .gitignore | 9 +- backend/__pycache__/flex_lots.cpython-313.pyc | Bin 10167 -> 10190 bytes .../__pycache__/flex_query.cpython-313.pyc | Bin 7265 -> 8390 bytes backend/__pycache__/ib_client.cpython-313.pyc | Bin 34904 -> 34927 bytes backend/app.py | 119 ++++++++++++++++-- backend/flex_query.py | 33 ++++- 6 files changed, 145 insertions(+), 16 deletions(-) diff --git a/.gitignore b/.gitignore index f1ac045..4729068 100644 --- a/.gitignore +++ b/.gitignore @@ -1,2 +1,9 @@ .env -data/ \ No newline at end of file +data/ + +# Ignore Python bytecode files +*.pyc +__pycache__/ +*.pyo +*.pyd +*$py.class \ No newline at end of file diff --git a/backend/__pycache__/flex_lots.cpython-313.pyc b/backend/__pycache__/flex_lots.cpython-313.pyc index 7c5b30a5f23d1c29b80c06d2b6807307abffb9ce..a7d699f000cff1588e78f751e0c59a43efd4783a 100644 GIT binary patch delta 49 zcmdn)f6kxlGcPX}0}z~ezLD!9qlC7zRZK`~Zb3{%jHi>XOJZ?GQhs7lO3davjJH(* DpC=KI delta 26 gcmX@-zulkfGcPX}0}#x8w2|u~BcsUX&y2TK0eBG!4FCWD diff --git a/backend/__pycache__/flex_query.cpython-313.pyc b/backend/__pycache__/flex_query.cpython-313.pyc index 20b345cac95976683a4e34fc6c274661f25457bf..21fce279b5392a567b496dfa299c55e8b125a6be 100644 GIT binary patch delta 2890 zcmZWrYit|G5#Hn7@%<9t;zOhyJxo|n$dn~QY9uR>D;rTPo#?8e+zUmXBqn@Fd&fjJ z(8_7iz%5Y3>Y|7eq>X^0fQ|lW0ry9K1Vvq-E()Y61p<+9(2w*_QKKy&L#mOc1=`u8 zDJjhYJ2yM~&CcxZ?9Ban=;65i&}K6uXfuERV}8TCX%8^}+3fuJwE#MU1X~aZ_62W( zr2vGmPw)zkAbN$ve!(X=H7+3dg(i(_76O7xWi-KhgtAZ7yckL2?R|Q531Ts>t`yblgn-#a_*07u|-3!o8v>D;Xu96#<6xixowZ zMG<>p?hrasyC#wgbD+PlzA$yM`?-v~vRKaG+`?ie`zlyjSSm>CB6x#u^sL^%22jyR zGcfW9+2lGxHO{|xwyVmuZ1Xj)Yp=bd%0;&1Pq+>oCLeNBLL1eg2B+`w*hdF`N&=M7 zWpIcR3cz7XB9ufMglZr1K;A?uXRITlI zTN~`Dr8R#j(Qs!dLbKa;f1Xe4?A;)cX(+^EF&rTe3_d;qqb{;!mQMNw8o(8t$cJVSnD|FMNS3Y|Ipj3HND17y{K zA@Kg;7$84!w6JW7A&+@)h)*zpgEq?0P1Bas!Qx2JlO}znmrJpNAvupw#)429dcCE% zhMvTvZ6{LF>TG3<uijebZD&HOQ50R+44a zq}|&@IVX`{JL6uizpwASYMYgoBrKJ(5}o**JaKjp7*yT!bEmy0q)b7@&MF6B!( zMOs&IKV6b9lwX%H)~=X_P+U)Dnu0HxK~~0M78X)%b-oCOGs}`{#|<;uv`inFcIZz} z-7T(Wuq=Jy$>BVJ($MAKk-JT!U)%Eq-!5(zx3AQEeY?KFUtZkxj_!CzYu;yXoA;U{ z??>K=Y=5iPJh!WCoe)EcIz7S~Ne&}H4NV%reftmMO^sIOURAQj((xwJP;K;Pe_ng4d zHO!JT?b>An=nhLd;$FZCy0f}t3B!{#Psvta|g`#N~Xo$`og0Ihw|qhl#U-RfZ* z`zsmJ?ZCn81`dpp*EPXS_yAq5i zdpMS+L2UTJt!`lyvib;&DP3d96ffz>U%g|oD1z(a$N5Xl3vWS@)c7T435~K*w8X?6#O_Z79?|@z zG6Qi|WfYZJQ<>NBH$X&Y{L_%{w*BG36^h^&Z$N>I2e6ilwITxcL0)_{m&YoPGo@uo ztBl$k#@EPS{B7ZSzMx^Pa!S4^sk~f}q*bkqlBPflM_&mU3xsA=wjh;MJ!B>Nwyby@ zCa`ANQ^+(5HmyXWmSt4K)8ec$qRL9^S=C0fzEqS-inv}Z;4Jw;;M^c}v7U2IXCBXz zq-0l4YK5kSHvy|aht%`8Kw{%lZ}Y9`y->8u1@_&6t){KX-SEgxc;th%o$%Nv?(y4( z1D)BF+>Z_XYHHuvzUK?SUDzz#_5C7xFZzMJOHuy!Q62_yovMdlq6#kEY@)?~qFh?aFO!j= zdE?Nv7DrQO RM=b_s^q-zSMz8%M{SS^fnt}iT delta 1866 zcmZ8iTW=Fb6yEXf+Uwie@!cCIF$4lDNP-$76qA zs+3AqTSNg353~}3KJ=#cp;Gz_+N$b9FCZj}R`da>>N8U41AU=$#+Xz+lE0iebI$BH zb1wOD=Y!|fm2fy9FnRpzFU`9HSJVvp>&k|&UN+)JSQCy*BxNKFm1}7uX+*e|F;Yg9 zYgr?03~)_jele~MHbgzMyb)<0bdb)9uSyM7kI-*LB7Z3qbcJ@LgPYu>Y1&TJY1YiD zhhR;S4xvCJudGt)-o6FT+n3_heZM) z36KKF0KjP^3!nky`T)_k%-~(Ym6V=GZux?@;?dW`mdGHV&+{lZ+G7sW zvDqNwYmB)F0>fDHf}0dzVW+Jl~> z7eiYHftMR;V5j#1c#^J7W|(zy`CjNd&%`cpPzbg=ZOg`-G-KB1<~NTyPRC^Nq`(~a z(VrsM18^~R1GKV8+tJU$qk>QgG7}eiGn%H816n|yMm-CrH=-HjrCSEJy&?iXkU9%k z4KjSsoYDN~R?j*?R|>Az=#kAcf7fSsYc*`!^uw4Q2^aAqCfKUs*>)YDr9a2U*$4iK z4aRivm0Ru8nDA`q65SlvG?nxN*b~>WEit#=CFz^-Y~To6mf6pKkKT;$(C(|bYo!}U zKRWuw%$=d_zYopdnwh_&9=@d4GmYosvWI?}I<~d~1}>}5TM6}TWud6jlj$Ui(g%UKC>b)no=*E>hG2=6IDI2M zOb=)Z9m)8z=s)pL@LZn8$9H%P7CRwMdwOTpsv(BJLU+=lE4Bb-`ms;<)Zl*pT z`DEl9`?m7JiY$(MR|GNZyDcTp?R*#zl36-4bVufcke(gR&z=VgQoUJobT8S<7Mv*A z$F#F5m@DiyoAV5g67VG*2D;R17|_Xi`p0lO$Q$VC`pIx=I_LUKv)-#c z`N;=jFL{s=0?+fme+aJf;!?ZbZQ&VmiS5Ebv@g>2c_S;M%tL{}iYR!34?RqK6h+8? Z`3Sw9PsbFr_EE}*cKn?fMPB}B{{oiTp!EO% diff --git a/backend/__pycache__/ib_client.cpython-313.pyc b/backend/__pycache__/ib_client.cpython-313.pyc index 61b5bb1ca827b6ea3a4934af3dcbca75f99040ec..864a4bc310ccf084731ba4d30c88a8b6283f97b0 100644 GIT binary patch delta 59 zcmcaHf$9AOCa%xCyj%=GaN_w!u7B(b+Rj!nA*s0qF%>bMPP#6M#TiNYiA5kEbe diff --git a/backend/app.py b/backend/app.py index c29d6dc..624cc2f 100644 --- a/backend/app.py +++ b/backend/app.py @@ -16,6 +16,19 @@ from flex_query import FlexQueryError, fetch_flex_statement, get_flex_status from ib_client import IBClient import time +FLEX_AUTO_SYNC_MAX_AGE_SECONDS = 8 * 60 * 60 +_flex_sync_lock = threading.Lock() +_flex_startup_sync_started = False +_flex_sync_state = { + "last_attempt_at": None, + "last_success_at": None, + "last_result": None, + "last_error": None, + "last_saved_to": None, + "last_reference_code": None, + "last_trigger": None, +} + def load_dotenv_file() -> None: env_path = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", ".env")) @@ -41,18 +54,92 @@ load_dotenv_file() # Use an absolute path for the static folder so Flask can reliably serve files STATIC_FOLDER = os.path.abspath(os.path.join(os.path.dirname(__file__), "..", "frontend")) -app = Flask(__name__, static_folder=STATIC_FOLDER, static_url_path="/static") - -# instantiate the IB client with a process-unique client id to avoid -# collisions when Flask's reloader spawns multiple processes. -ib = IBClient(client_id=os.getpid()) +app = Flask(__name__, static_folder=STATIC_FOLDER, static_url_path="/static") + + +def should_auto_sync_flex(status: dict) -> bool: + if not status.get("configured"): + return False + latest_timestamp = status.get("latest_timestamp") + if latest_timestamp is None: + return True + return (time.time() - latest_timestamp) > FLEX_AUTO_SYNC_MAX_AGE_SECONDS + + +def get_flex_status_with_sync_state() -> dict: + status = get_flex_status() + status.update(_flex_sync_state) + return status + + +def run_flex_sync_if_needed(reason: str) -> bool: + status = get_flex_status() + if not should_auto_sync_flex(status): + return False + + if not _flex_sync_lock.acquire(blocking=False): + return False + + try: + _flex_sync_state["last_attempt_at"] = time.time() + _flex_sync_state["last_trigger"] = reason + _flex_sync_state["last_result"] = "running" + _flex_sync_state["last_error"] = None + + # Re-check after taking the lock so concurrent callers do not duplicate work. + status = get_flex_status() + if not should_auto_sync_flex(status): + _flex_sync_state["last_result"] = "skipped" + return False + result = fetch_flex_statement() + _flex_sync_state["last_success_at"] = time.time() + _flex_sync_state["last_result"] = "success" + _flex_sync_state["last_saved_to"] = result["saved_to"] + _flex_sync_state["last_reference_code"] = result["reference_code"] + print( + f"Flex sync completed ({reason}): saved_to={result['saved_to']} " + f"reference_code={result['reference_code']}" + ) + return True + except FlexQueryError as exc: + _flex_sync_state["last_result"] = "error" + _flex_sync_state["last_error"] = str(exc) + print(f"Flex sync skipped/failed ({reason}): {exc}") + return False + except Exception as exc: + _flex_sync_state["last_result"] = "error" + _flex_sync_state["last_error"] = str(exc) + print(f"Unexpected Flex sync failure ({reason}): {exc}") + return False + finally: + _flex_sync_lock.release() + + +def start_flex_startup_sync_if_needed() -> None: + global _flex_startup_sync_started + + if _flex_startup_sync_started: + return + _flex_startup_sync_started = True + + thread = threading.Thread( + target=run_flex_sync_if_needed, + args=("startup",), + daemon=True, + ) + thread.start() + +# instantiate the IB client with a process-unique client id to avoid +# collisions when Flask's reloader spawns multiple processes. +ib = IBClient(client_id=os.getpid()) # Start connection only in the main process (Werkzeug sets WERKZEUG_RUN_MAIN) if os.environ.get("WERKZEUG_RUN_MAIN") == "true" or not os.environ.get("WERKZEUG_RUN_MAIN"): # attempt to start the connection; failures are logged by the client - try: - ib.connect_and_start() - except Exception: - pass + try: + ib.connect_and_start() + except Exception: + pass + start_flex_startup_sync_if_needed() @app.route("/") @@ -216,7 +303,7 @@ def api_executions_debug(): @app.route("/api/flex/status") def api_flex_status(): - return jsonify(get_flex_status()) + return jsonify(get_flex_status_with_sync_state()) @app.route("/api/flex/lots-debug") @@ -236,11 +323,23 @@ def api_flex_lots_debug(): @app.route("/api/flex/sync", methods=["POST"]) def api_flex_sync(): try: + _flex_sync_state["last_attempt_at"] = time.time() + _flex_sync_state["last_trigger"] = "manual" + _flex_sync_state["last_result"] = "running" + _flex_sync_state["last_error"] = None result = fetch_flex_statement() + _flex_sync_state["last_success_at"] = time.time() + _flex_sync_state["last_result"] = "success" + _flex_sync_state["last_saved_to"] = result["saved_to"] + _flex_sync_state["last_reference_code"] = result["reference_code"] return jsonify({"ok": True, **result}) except FlexQueryError as exc: + _flex_sync_state["last_result"] = "error" + _flex_sync_state["last_error"] = str(exc) return jsonify({"ok": False, "error": str(exc)}), 400 except Exception as exc: + _flex_sync_state["last_result"] = "error" + _flex_sync_state["last_error"] = str(exc) return jsonify({"ok": False, "error": str(exc)}), 500 diff --git a/backend/flex_query.py b/backend/flex_query.py index 7dbafb1..6a40733 100644 --- a/backend/flex_query.py +++ b/backend/flex_query.py @@ -29,12 +29,18 @@ class FlexConfig: data_dir: Path poll_interval_seconds: float = 2.0 max_polls: int = 30 + send_request_retry_seconds: float = 10.0 + max_send_request_attempts: int = 6 class FlexQueryError(RuntimeError): pass +class FlexQueryRetryableError(FlexQueryError): + pass + + def load_flex_config() -> FlexConfig | None: token = os.environ.get("IB_FLEX_TOKEN", "").strip() query_id = os.environ.get("IB_FLEX_QUERY_ID", "").strip() @@ -86,7 +92,11 @@ def _parse_send_request(xml_bytes: bytes) -> tuple[str, str]: root = ET.fromstring(xml_bytes) status = (root.findtext("Status") or "").strip() if status.lower() != "success": - raise FlexQueryError(root.findtext("ErrorMessage") or "Flex send-request failed") + error_code = (root.findtext("ErrorCode") or "").strip() + error_message = (root.findtext("ErrorMessage") or "Flex send-request failed").strip() + if error_code == "1004": + raise FlexQueryRetryableError(error_message) + raise FlexQueryError(error_message) reference_code = (root.findtext("ReferenceCode") or "").strip() if not reference_code: raise FlexQueryError("Flex send-request returned no ReferenceCode") @@ -110,10 +120,23 @@ def fetch_flex_statement() -> dict: config.data_dir.mkdir(parents=True, exist_ok=True) - _, reference_code = _parse_send_request(_http_get( - FLEX_SEND_REQUEST_URL, - {"t": config.token, "q": config.query_id, "v": "3"}, - )) + reference_code = None + last_retryable_error = None + for attempt in range(1, config.max_send_request_attempts + 1): + try: + _, reference_code = _parse_send_request(_http_get( + FLEX_SEND_REQUEST_URL, + {"t": config.token, "q": config.query_id, "v": "3"}, + )) + break + except FlexQueryRetryableError as exc: + last_retryable_error = str(exc) + if attempt >= config.max_send_request_attempts: + raise FlexQueryError(last_retryable_error) from exc + time.sleep(config.send_request_retry_seconds) + + if not reference_code: + raise FlexQueryError(last_retryable_error or "Flex send-request returned no ReferenceCode") statement_xml = None for _ in range(config.max_polls):