diff --git a/ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.md b/ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.md new file mode 100644 index 00000000..f39858eb --- /dev/null +++ b/ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.md @@ -0,0 +1,37 @@ +--- +model: openai/gpt-6-astra +service: opencode +session: unavailable +timestamp: 2026-09-16T19:40:52Z +git_ref: 749ca1f0 +scope: code +substantive: true +raw_file: 20260916T194052Z_749ca1f0_prompt_io.raw.md +--- + +## Prompt + +Fix wrong IB fill/commission pairing and empty ledger normalization; +retain and investigate silent socket loss during network roaming. + +## Response summary + +Preserve execution identity from callbacks into accounting, prevent +empty commission IDs from erasing ledger identity, and propagate actual +socket disconnection as a transport exception. The existing data-farm +reset mechanism is not a socket-session reconnect implementation. + +## Files changed + +- `piker/brokers/ib/broker.py`: fill correlation and socket termination. +- `piker/brokers/ib/ledger.py`: authoritative execution ID. +- `piker/brokers/ib/api.py`: disconnect event propagation. +- `tests/test_ib_trade_events.py`: accounting event regressions. +- `tests/test_ib_method_proxy.py`: socket-disconnect regression. + +## Human edits + +The human supplied live traceback and debugger values proving the +execution mismatch, requested a regression rather than a KeyError +guard, and required continued attention to roaming recovery. No human +source edits to this patch have been observed; acceptance is pending. diff --git a/ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.raw.md b/ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.raw.md new file mode 100644 index 00000000..4751137f --- /dev/null +++ b/ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.raw.md @@ -0,0 +1,45 @@ +--- +model: openai/gpt-6-astra +service: opencode +timestamp: 2026-09-16T19:40:52Z +git_ref: 749ca1f0 +diff_cmd: git diff HEAD -- piker/brokers/ib tests/test_ib_trade_events.py tests/test_ib_method_proxy.py +--- + +## Human instructions and evidence + +Fix silent IB socket loss after network roaming, reported as +`[Errno 104] Connection reset by peer`. Also investigate a live +MNQ accounting crash in `emit_pp_update()` with empty normalized +ledger entries. The human supplied debugger evidence that a commission +for execution 4261 was paired with execution 4262 on the same contract +and account, and explicitly requested implementation and regression +tests. They asked to retain the roaming task and continue working +through follow-up messages. + +## Generated changes + +> `git diff HEAD -- piker/brokers/ib/broker.py` + +Preserve the callback-supplied fill with its commission. Track positions +by account and contract ID and pending commissioned fills by execution. +Reject mismatched commission/execution IDs. Propagate socket closure +instead of returning normally from the trade-event bridge. + +> `git diff HEAD -- piker/brokers/ib/ledger.py` + +Keep Execution.execId authoritative when flattening commission fields. + +> `git diff HEAD -- piker/brokers/ib/api.py` + +Relay the authoritative disconnected event as ConnectionError, retaining +the most recent API error text, and raise it from the proxy reader. +Automatic session reconnect/resynchronization remains unimplemented. + +> `git diff HEAD -- tests/test_ib_trade_events.py tests/test_ib_method_proxy.py` + +Cover overlapping executions, zero commissions, position-before/after +commission ordering, empty commission IDs, and eventkit socket loss. +Nine focused tests passed, including the actual eventkit callback +bridge and its disconnect exception. Focused Ruff checks and +`git diff --check` also passed. diff --git a/piker/brokers/ib/broker.py b/piker/brokers/ib/broker.py index 1fb377f3..cb0b7a29 100644 --- a/piker/brokers/ib/broker.py +++ b/piker/brokers/ib/broker.py @@ -338,7 +338,8 @@ async def recv_trade_updates( case 'commissionReportEvent': assert report - emit: CommissionReport = report + assert fill is not None + emit: tuple[Fill, CommissionReport] = (fill, report) case 'execDetailsEvent': # execution details event @@ -1189,15 +1190,10 @@ async def deliver_trade_events( Format and relay all trade events for a given client to emsd. ''' - # task local msg dialog tracking - clears: dict[ - Contract, - list[ - IbPosition | None, # filled by positionEvent - Fill | None, # filled by order status and exec details - ] - ] = {} - execid2con: dict[str, Contract] = {} + # Positions belong to an account/contract; commissions belong to + # executions. Never recover a Fill from just its Contract. + positions: dict[tuple[str, int], IbPosition] = {} + pending_fills: dict[tuple[str, int], dict[str, Fill]] = {} # TODO: for some reason we can receive a ``None`` here when the # ib-gw goes down? Not sure exactly how that's happening looking @@ -1296,105 +1292,38 @@ async def deliver_trade_events( # translating them to ``Status`` msgs) if we can # show the equivalent status events are no more latent. case 'execDetailsEvent': - # unpack attrs pep-0526 style. - trade: Trade - fill: Fill - trade, fill = item - con: Contract = trade.contract - execu: Execution = fill.execution - execid: str = execu.execId - report: CommissionReport = fill.commissionReport - - # always fill in id to con map so when commissions - # arrive we can maybe fire the pos update.. - execid2con[execid] = con - - # TODO: - # - normalize out commissions details? - # - this is the same as the unpacking loop above in - # ``trades_to_ledger_entries()`` no? - - # 2 cases: - # - fill comes first or - # - commission report comes first - clear: tuple = clears.setdefault( - con, - [None, fill], - ) - pos, _fill = clear - - # NOTE: we have to handle the case where a pos msg - # has already been set (bc we already relayed rxed - # one before both the exec-deats AND the - # comms-report?) but the comms-report hasn't yet - # arrived, so we fill in the fill (XD) and wait for - # the cost to show up before relaying the pos msg - # to the EMS.. - if _fill is None: - clear[1] = fill - - cost: float = report.commission - if ( - pos - and fill - and cost - ): - await emit_pp_update( - ems_stream, - accounts_def, - proxies, - ledgers, - tables, - - ibpos=pos, - fill=fill, - ) - clears.pop(con) + # Accounting waits for commissionReportEvent, whose + # callback supplies the exact Fill, including zero + # commissions. Order-status handling relays fills + # independently above. + continue case 'commissionReportEvent': - - cr: CommissionReport = item + fill, cr = item execid: str = cr.execId + if execid != fill.execution.execId: + raise ValueError( + f'IB commission/fill execution mismatch:\n' + f'{cr!r}\n{fill!r}' + ) + key = ( + fill.execution.acctNumber, + fill.contract.conId, + ) + fill = fill._replace(commissionReport=cr) + if (pos := positions.get(key)) is None: + pending_fills.setdefault(key, {})[execid] = fill + continue - # only fire a pp msg update if, - # - we haven't already - # - the fill event has already arrived - # but it didn't yet have a commision report - # which we fill in now. - - # placehold i guess until someone who know wtf - # contract this is from can fill it in... - con: Contract | None = execid2con.setdefault(execid, None) - if ( - con - and (clear := clears.get(con)) - ): - pos, fill = clear - if ( - pos - and fill - ): - now_cr: CommissionReport = fill.commissionReport - if (now_cr != cr): - log.warning( - 'UhhHh ib updated the commission report mid-fill..?\n' - f'was: {pformat(cr)}\n' - f'now: {pformat(now_cr)}\n' - ) - - await emit_pp_update( - ems_stream, - accounts_def, - proxies, - ledgers, - tables, - - ibpos=pos, - fill=fill, - ) - clears.pop(con) - # TODO: should we clean this? - # execid2con.pop(execid) + await emit_pp_update( + ems_stream, + accounts_def, + proxies, + ledgers, + tables, + ibpos=pos, + fill=fill, + ) # always update with latest ib pos msg info since # we generally audit against it for sanity and @@ -1407,10 +1336,18 @@ async def deliver_trade_events( bs_mktid, ppmsg = pack_position(pos, accounts_def) log.info(f'New IB position msg: {ppmsg}') - _, fill = clears.setdefault( - con, - [pos, None], - ) + key = (pos.account, con.conId) + positions[key] = pos + for fill in pending_fills.pop(key, {}).values(): + await emit_pp_update( + ems_stream, + accounts_def, + proxies, + ledgers, + tables, + ibpos=pos, + fill=fill, + ) # only send a pos update once we've actually rxed # the msg from IB since generally speaking we use # their 'cumsize' as gospel. diff --git a/piker/brokers/ib/ledger.py b/piker/brokers/ib/ledger.py index 31dc4ab9..e3476592 100644 --- a/piker/brokers/ib/ledger.py +++ b/piker/brokers/ib/ledger.py @@ -429,9 +429,16 @@ def api_trades_to_ledger_entries( for attr_name, val in fdict.items(): match attr_name: # value is a `@dataclass` subtype - case 'contract' | 'execution' | 'commissionReport': + case 'contract' | 'execution': txn_dict.update(asdict(val)) + case 'commissionReport': + report_fields = asdict(val) + # Execution.execId owns ledger identity. An empty + # CommissionReport.execId must not erase it. + report_fields.pop('execId') + txn_dict.update(report_fields) + case 'time': # ib has wack ns timestamps, or is that us? continue diff --git a/tests/test_ib_trade_events.py b/tests/test_ib_trade_events.py new file mode 100644 index 00000000..cdd2c5ee --- /dev/null +++ b/tests/test_ib_trade_events.py @@ -0,0 +1,115 @@ +''' +Execution identity across IB fill and commission callbacks. + +''' +from types import SimpleNamespace + +from bidict import bidict +from ib_async import ( + CommissionReport, + Contract, + Execution, + Fill, + Position, +) +import pytest +import trio + +from piker.brokers.ib import broker +from piker.brokers.ib.ledger import api_trades_to_ledger_entries + + +def make_fill(tid: str) -> Fill: + return Fill( + Contract(conId=815824267, secType='FUT', symbol='MNQ'), + Execution( + execId=tid, + acctNumber='DU_TEST', + time=1789586554.0, + ), + CommissionReport(), + 1789586554.0, + ) + + +def test_execution_id_survives_empty_commission() -> None: + ''' + Flattening a default CommissionReport erased Execution.execId, + causing a real fill to be skipped as an ID-less adjustment. Pass + a fill with a valid execution and empty commission through ledger + conversion and assert that its account and execution survive. + + ''' + entries = api_trades_to_ledger_entries( + bidict({'DU_TEST': 'ib.algopaper'}), + [make_fill('execution-1')], + ) + assert entries['ib.algopaper']['execution-1']['execId'] == ( + 'execution-1' + ) + + +@pytest.mark.parametrize('position_first', [True, False]) +def test_overlapping_execution_commissions( + monkeypatch: pytest.MonkeyPatch, + position_first: bool, +) -> None: + ''' + Two executions on one contract shared a single pending Fill. + A commission for execution A could therefore account execution B. + Queue both executions before either commission, delivering fees + in reverse order, and exercise positions before and after fees. + The deterministic channel order verifies every execution reaches + accounting once with its own fee, including a zero commission. + + ''' + fills = [make_fill('A'), make_fill('B')] + pos = Position('DU_TEST', fills[0].contract, 2, 58000) + recorded = [] + + async def capture(*args, ibpos, fill=None): + if fill is not None: + recorded.append(( + fill.execution.execId, + fill.commissionReport.execId, + fill.commissionReport.commission, + ibpos.account, + )) + + monkeypatch.setattr(broker, 'emit_pp_update', capture) + monkeypatch.setattr( + broker, 'pack_position', lambda *args: ('815824267', {}) + ) + + async def main(): + send, recv = trio.open_memory_channel(10) + events = [ + ('execDetailsEvent', ( + SimpleNamespace(contract=fill.contract), fill, + )) + for fill in fills + ] + events.extend([ + ('commissionReportEvent', ( + fills[1], CommissionReport(execId='B', commission=0), + )), + ('commissionReportEvent', ( + fills[0], CommissionReport(execId='A', commission=.61), + )), + ]) + events.insert(0 if position_first else len(events), ( + 'positionEvent', pos, + )) + async with send, recv: + for event in events: + send.send_nowait(event) + await send.aclose() + await broker.deliver_trade_events( + recv, None, {}, {}, {}, {}, None, + ) + + trio.run(main) + assert recorded == [ + ('B', 'B', 0, 'DU_TEST'), + ('A', 'A', .61, 'DU_TEST'), + ]