From efa72980f8e54bd233b1632480d921f236451d8b Mon Sep 17 00:00:00 2001 From: goodboy Date: Thu, 17 Sep 2026 18:24:17 -0400 Subject: [PATCH] Raise on IB socket disconnects in proxy and trades Propagate the authoritative `disconnectedEvent` as a `ConnectionError` so proxy callers and trade consumers fail loudly instead of waiting on a dead socket. Preserve the last `apiError` text as the reason without treating numeric farm/order errors as socket failures. Deats, - Disconnect the one-shot callback before sending the exception. - Raise when `recv_trade_updates()` loses its socket, and relay disconnect exceptions through both Trio-side consumers. - Exercise real `eventkit` error/disconnect callbacks and preserve the supplied commission fill while testing bridge termination. - Record the aborted integration stage as historical provenance; its test passes do not validate this post-abort boundary. Prompt-IO: ai/prompt-io/opencode/20260916_ib_connectivity_integration.md (this patch was generated in some part by `opencode` using `gpt-6-astra` (`openai`)) --- .../20260916_ib_connectivity_integration.md | 61 +++++++++++++++++++ ...0260916_ib_connectivity_integration.raw.md | 25 ++++++++ piker/brokers/ib/api.py | 25 ++++++++ piker/brokers/ib/broker.py | 4 ++ tests/test_ib_method_proxy.py | 28 +++++++++ tests/test_ib_trade_events.py | 45 ++++++++++++++ 6 files changed, 188 insertions(+) create mode 100644 ai/prompt-io/opencode/20260916_ib_connectivity_integration.md create mode 100644 ai/prompt-io/opencode/20260916_ib_connectivity_integration.raw.md diff --git a/ai/prompt-io/opencode/20260916_ib_connectivity_integration.md b/ai/prompt-io/opencode/20260916_ib_connectivity_integration.md new file mode 100644 index 00000000..567e61a6 --- /dev/null +++ b/ai/prompt-io/opencode/20260916_ib_connectivity_integration.md @@ -0,0 +1,61 @@ +--- +model: openai/gpt-6-astra +service: opencode +scope: code +substantive: true +git_ref: 5eb4fc9d +raw_file: 20260916_ib_connectivity_integration.raw.md +--- + +## Prompt + +Investigate IB Gateway loss, silent chart survival, infected-actor +shutdown errors and leaked shared memory. Integrate pending tagged +address and broadcast-consumer compatibility work. The human required +isolated worktrees and preservation of the live setup, and selected +tractor piker_pin rather than main as the runtime fix base. + +## Response summary + +Created isolated piker and tractor worktrees and initially integrated +c07e3cbd and 4af885d9 without committing. That pending merge was later +aborted. The current piker scope is the original IB fail-loud patch +and its regression tests on 5eb4fc9d; neither upstream history is +merged. The live interpreter installation was not changed. + +Integration-era logging and Deribit compatibility fixes were archived +outside this patch and are absent from the current worktree. The raw +companion preserves that earlier response as historical output. + +## Generated files + +- piker/brokers/ib/api.py: propagate actual socket disconnection with + the last API error through the method-proxy relay. +- piker/brokers/ib/broker.py: raise on trade-event socket loss and + propagate the disconnect event to the trade consumer. +- tests/test_ib_method_proxy.py: exercise API error plus disconnect. +- tests/test_ib_trade_events.py: preserve callback fill identity and + raise when the callback bridge disconnects. + +Existing IB fail-loud edits and their tests were copied from the +earlier session, whose provenance remains in the fill commit. +The commit plan includes these four files and both provenance logs. + +## Verification + +Historically, twenty selected piker regression cases passed on the +integrated tree against isolated, patched tractor sources. Those +results do not validate the current post-abort tree or installed +pinned tractor. The new exact-tree commit-plan checks remain pending; +no project tests were run during this planning pass. No live +service/container tests or installed dependency changes were performed. +Original live-worktree diff matched SHA256 +cbdfa439e656ef3a5effaa0f0e1a7706f1070de863fc0e85222d36294eed8524. + +## Human contributions + +The human identified the dependency branches, corrected the proposed +tractor main base to piker_pin, and required a new worktree after the +agent initially started an uncommitted merge in the live checkout. +That merge was aborted and preservation was verified before isolated +work continued. The human supplied all live incident evidence. diff --git a/ai/prompt-io/opencode/20260916_ib_connectivity_integration.raw.md b/ai/prompt-io/opencode/20260916_ib_connectivity_integration.raw.md new file mode 100644 index 00000000..34b5a9a4 --- /dev/null +++ b/ai/prompt-io/opencode/20260916_ib_connectivity_integration.raw.md @@ -0,0 +1,25 @@ +--- +model: openai/gpt-6-astra +service: opencode +git_ref: 5eb4fc9d +diff_cmd: git diff HEAD +--- + +> git diff HEAD -- piker/data/feed.py piker/brokers/deribit/api.py piker/brokers/deribit/feed.py + +Use the standard logging predicate for INFO and repair import-time +Deribit message struct/annotation compatibility encountered while +validating the merged broadcast-consumer baseline. + +The isolated checks pass, including real shm cleanup: + +- All 20 selected piker integration cases pass after the narrow + compatibility fixes. +- Tractor's EOF tests pass for TCP and UDS, including an idle + owner and an unaffected peer. +- All four late-cancellation variants unlink their test-owned + data/first/last shm segments before parent teardown. + +The original live-worktree diff still matches its saved hash. +We're ready for a Gateway-container test of the full actor tree; +automatic reconnect is still a separate, unfinished step. diff --git a/piker/brokers/ib/api.py b/piker/brokers/ib/api.py index 6e0bd6ed..f9e736cd 100644 --- a/piker/brokers/ib/api.py +++ b/piker/brokers/ib/api.py @@ -1276,8 +1276,11 @@ class Client: self.ib.errorEvent.connect(push_err) api_err = self.ib.client.apiError + last_api_error: str = '' def report_api_err(msg: str) -> None: + nonlocal last_api_error + last_api_error = msg with remove_handler_on_err( api_err, report_api_err, @@ -1290,6 +1293,25 @@ class Client: api_err.connect(report_api_err) + disconnected = self.ib.disconnectedEvent + + def report_disconnect() -> None: + # IB.disconnectedEvent is authoritative for socket loss; + # numeric IB farm/order errors are a different namespace + # from OS errno values reported through Client.apiError. + disconnected.disconnect(report_disconnect) + with remove_handler_on_err( + disconnected, + report_disconnect, + ): + reason = last_api_error or 'IB API socket disconnected' + to_trio.send_nowait(( + 'disconnected', + ConnectionError(reason), + )) + + disconnected.connect(report_disconnect) + def positions( self, account: str = '', @@ -1851,6 +1873,9 @@ async def relay_client_proxy_messages( # proxy.reset() match msg: + case ('disconnected', ConnectionError() as err): + raise err + # Correlated response from a method call. case {'mid': _}: proxy._deliver_method_response(msg) diff --git a/piker/brokers/ib/broker.py b/piker/brokers/ib/broker.py index cb0b7a29..54a754f1 100644 --- a/piker/brokers/ib/broker.py +++ b/piker/brokers/ib/broker.py @@ -391,6 +391,7 @@ async def recv_trade_updates( # let the engine run and stream await client.ib.disconnectedEvent + raise ConnectionError('IB trade-event socket disconnected') async def update_and_audit_pos_msg( @@ -1361,6 +1362,9 @@ async def deliver_trade_events( ibpos=pos, ) + case 'disconnected': + raise item + case 'error': # NOTE: see impl deats in # `Client.inline_errors()::push_err()` diff --git a/tests/test_ib_method_proxy.py b/tests/test_ib_method_proxy.py index 051a7e59..928d22a0 100644 --- a/tests/test_ib_method_proxy.py +++ b/tests/test_ib_method_proxy.py @@ -5,6 +5,7 @@ IB asyncio method-proxy regressions. from types import SimpleNamespace from typing import Any +from eventkit import Event import pytest import trio @@ -63,6 +64,33 @@ class FakeChannel: await self._tx.aclose() +def test_socket_disconnect_raises_from_proxy() -> None: + ''' + Roaming reset the IB socket, but apiError was merely logged and + proxy callers could wait indefinitely. Emit the actual eventkit + apiError/disconnected sequence through Client.inline_errors(). + The single-reader proxy must raise ConnectionError carrying the + reset reason, while the preceding error alone remains nonfatal. + + ''' + async def main() -> None: + chan = FakeChannel() + ib = SimpleNamespace( + errorEvent=Event('errorEvent'), + disconnectedEvent=Event('disconnectedEvent'), + client=SimpleNamespace(apiError=Event('apiError')), + ) + client = ib_api.Client(ib=ib, config={}) + client.inline_errors(chan._tx) + proxy = MethodProxy(chan, {}, asyncio_ns=client) + ib.client.apiError.emit('[Errno 104] Connection reset by peer') + ib.disconnectedEvent.emit() + with pytest.raises(ConnectionError, match='Errno 104'): + await relay_client_proxy_messages(chan, proxy) + + trio.run(main) + + def test_method_proxy_routes_after_idle_status_events() -> None: ''' Idle IB status traffic must not break first-symbol qualification. diff --git a/tests/test_ib_trade_events.py b/tests/test_ib_trade_events.py index cdd2c5ee..963a7fd1 100644 --- a/tests/test_ib_trade_events.py +++ b/tests/test_ib_trade_events.py @@ -2,9 +2,11 @@ Execution identity across IB fill and commission callbacks. ''' +import asyncio from types import SimpleNamespace from bidict import bidict +from eventkit import Event from ib_async import ( CommissionReport, Contract, @@ -49,6 +51,49 @@ def test_execution_id_survives_empty_commission() -> None: ) +def test_callback_preserves_fill_and_disconnect() -> None: + ''' + The asyncio callback dropped commissionReportEvent's exact Fill, + forcing the Trio side to guess by contract. Exercise the actual + eventkit relay and assert it preserves the supplied object. After + that callback, emit disconnectedEvent: the bridge must raise + rather than return normally and leave its Trio owner asleep. + One asyncio checkpoint lets recv_trade_updates install handlers + before the test producer emits either event. + + ''' + async def main(): + ib = SimpleNamespace(**{ + name: Event(name) + for name in ( + 'orderStatusEvent', 'commissionReportEvent', + 'execDetailsEvent', 'positionEvent', + 'disconnectedEvent', + ) + }) + emitted = [] + chan = SimpleNamespace( + started_nowait=lambda _: None, + send_nowait=emitted.append, + ) + client = SimpleNamespace(ib=ib, inline_errors=lambda _: None) + fill = make_fill('A') + report = CommissionReport(execId='A', commission=.61) + + async def produce(): + await asyncio.sleep(0) + ib.commissionReportEvent.emit(None, fill, report) + ib.disconnectedEvent.emit() + + producer = asyncio.create_task(produce()) + with pytest.raises(ConnectionError, match='socket'): + await broker.recv_trade_updates(chan, client) + await producer + assert emitted == [('commissionReportEvent', (fill, report))] + + asyncio.run(main()) + + @pytest.mark.parametrize('position_first', [True, False]) def test_overlapping_execution_commissions( monkeypatch: pytest.MonkeyPatch,