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`))wkt/ib_history_request_bounds
parent
5eb4fc9d8a
commit
efa72980f8
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -1276,8 +1276,11 @@ class Client:
|
||||||
self.ib.errorEvent.connect(push_err)
|
self.ib.errorEvent.connect(push_err)
|
||||||
|
|
||||||
api_err = self.ib.client.apiError
|
api_err = self.ib.client.apiError
|
||||||
|
last_api_error: str = ''
|
||||||
|
|
||||||
def report_api_err(msg: str) -> None:
|
def report_api_err(msg: str) -> None:
|
||||||
|
nonlocal last_api_error
|
||||||
|
last_api_error = msg
|
||||||
with remove_handler_on_err(
|
with remove_handler_on_err(
|
||||||
api_err,
|
api_err,
|
||||||
report_api_err,
|
report_api_err,
|
||||||
|
|
@ -1290,6 +1293,25 @@ class Client:
|
||||||
|
|
||||||
api_err.connect(report_api_err)
|
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(
|
def positions(
|
||||||
self,
|
self,
|
||||||
account: str = '',
|
account: str = '',
|
||||||
|
|
@ -1851,6 +1873,9 @@ async def relay_client_proxy_messages(
|
||||||
# proxy.reset()
|
# proxy.reset()
|
||||||
|
|
||||||
match msg:
|
match msg:
|
||||||
|
case ('disconnected', ConnectionError() as err):
|
||||||
|
raise err
|
||||||
|
|
||||||
# Correlated response from a method call.
|
# Correlated response from a method call.
|
||||||
case {'mid': _}:
|
case {'mid': _}:
|
||||||
proxy._deliver_method_response(msg)
|
proxy._deliver_method_response(msg)
|
||||||
|
|
|
||||||
|
|
@ -391,6 +391,7 @@ async def recv_trade_updates(
|
||||||
|
|
||||||
# let the engine run and stream
|
# let the engine run and stream
|
||||||
await client.ib.disconnectedEvent
|
await client.ib.disconnectedEvent
|
||||||
|
raise ConnectionError('IB trade-event socket disconnected')
|
||||||
|
|
||||||
|
|
||||||
async def update_and_audit_pos_msg(
|
async def update_and_audit_pos_msg(
|
||||||
|
|
@ -1361,6 +1362,9 @@ async def deliver_trade_events(
|
||||||
ibpos=pos,
|
ibpos=pos,
|
||||||
)
|
)
|
||||||
|
|
||||||
|
case 'disconnected':
|
||||||
|
raise item
|
||||||
|
|
||||||
case 'error':
|
case 'error':
|
||||||
# NOTE: see impl deats in
|
# NOTE: see impl deats in
|
||||||
# `Client.inline_errors()::push_err()`
|
# `Client.inline_errors()::push_err()`
|
||||||
|
|
|
||||||
|
|
@ -5,6 +5,7 @@ IB asyncio method-proxy regressions.
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
from typing import Any
|
from typing import Any
|
||||||
|
|
||||||
|
from eventkit import Event
|
||||||
import pytest
|
import pytest
|
||||||
import trio
|
import trio
|
||||||
|
|
||||||
|
|
@ -63,6 +64,33 @@ class FakeChannel:
|
||||||
await self._tx.aclose()
|
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:
|
def test_method_proxy_routes_after_idle_status_events() -> None:
|
||||||
'''
|
'''
|
||||||
Idle IB status traffic must not break first-symbol qualification.
|
Idle IB status traffic must not break first-symbol qualification.
|
||||||
|
|
|
||||||
|
|
@ -2,9 +2,11 @@
|
||||||
Execution identity across IB fill and commission callbacks.
|
Execution identity across IB fill and commission callbacks.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
|
import asyncio
|
||||||
from types import SimpleNamespace
|
from types import SimpleNamespace
|
||||||
|
|
||||||
from bidict import bidict
|
from bidict import bidict
|
||||||
|
from eventkit import Event
|
||||||
from ib_async import (
|
from ib_async import (
|
||||||
CommissionReport,
|
CommissionReport,
|
||||||
Contract,
|
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])
|
@pytest.mark.parametrize('position_first', [True, False])
|
||||||
def test_overlapping_execution_commissions(
|
def test_overlapping_execution_commissions(
|
||||||
monkeypatch: pytest.MonkeyPatch,
|
monkeypatch: pytest.MonkeyPatch,
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue