Fix IB fill accounting by execution identity
Keep each commission callback's exact `Fill` paired with its report, so overlapping executions on one contract cannot borrow each other's fees. Cache positions by account and `conId`, queue commissioned fills by execution ID until a position arrives, and accept zero commissions. Reject mismatched execution/report IDs. Keep `Execution.execId` authoritative when flattening ledger entries; an empty commission ID must not erase a valid fill. Add regressions for empty commission IDs, reversed fee arrival, and positions arriving before or after overlapping executions. The human live-tested fills and reports the accounting fix works. Exact-boundary automated checks remain pending; the earlier nine passes covered the mixed working tree. Included Prompt-IO logs retain whole-session context, including separate roaming work. Prompt-IO: ai/prompt-io/opencode/20260916T194052Z_749ca1f0_prompt_io.md (this patch was generated in some part by `opencode` using `gpt-6-astra` (`openai`))wkt/ib_history_request_bounds
parent
749ca1f09c
commit
5eb4fc9d8a
|
|
@ -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.
|
||||||
|
|
@ -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.
|
||||||
|
|
@ -338,7 +338,8 @@ async def recv_trade_updates(
|
||||||
|
|
||||||
case 'commissionReportEvent':
|
case 'commissionReportEvent':
|
||||||
assert report
|
assert report
|
||||||
emit: CommissionReport = report
|
assert fill is not None
|
||||||
|
emit: tuple[Fill, CommissionReport] = (fill, report)
|
||||||
|
|
||||||
case 'execDetailsEvent':
|
case 'execDetailsEvent':
|
||||||
# execution details event
|
# execution details event
|
||||||
|
|
@ -1189,15 +1190,10 @@ async def deliver_trade_events(
|
||||||
Format and relay all trade events for a given client to emsd.
|
Format and relay all trade events for a given client to emsd.
|
||||||
|
|
||||||
'''
|
'''
|
||||||
# task local msg dialog tracking
|
# Positions belong to an account/contract; commissions belong to
|
||||||
clears: dict[
|
# executions. Never recover a Fill from just its Contract.
|
||||||
Contract,
|
positions: dict[tuple[str, int], IbPosition] = {}
|
||||||
list[
|
pending_fills: dict[tuple[str, int], dict[str, Fill]] = {}
|
||||||
IbPosition | None, # filled by positionEvent
|
|
||||||
Fill | None, # filled by order status and exec details
|
|
||||||
]
|
|
||||||
] = {}
|
|
||||||
execid2con: dict[str, Contract] = {}
|
|
||||||
|
|
||||||
# TODO: for some reason we can receive a ``None`` here when the
|
# TODO: for some reason we can receive a ``None`` here when the
|
||||||
# ib-gw goes down? Not sure exactly how that's happening looking
|
# 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
|
# translating them to ``Status`` msgs) if we can
|
||||||
# show the equivalent status events are no more latent.
|
# show the equivalent status events are no more latent.
|
||||||
case 'execDetailsEvent':
|
case 'execDetailsEvent':
|
||||||
# unpack attrs pep-0526 style.
|
# Accounting waits for commissionReportEvent, whose
|
||||||
trade: Trade
|
# callback supplies the exact Fill, including zero
|
||||||
fill: Fill
|
# commissions. Order-status handling relays fills
|
||||||
trade, fill = item
|
# independently above.
|
||||||
con: Contract = trade.contract
|
continue
|
||||||
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)
|
|
||||||
|
|
||||||
case 'commissionReportEvent':
|
case 'commissionReportEvent':
|
||||||
|
fill, cr = item
|
||||||
cr: CommissionReport = item
|
|
||||||
execid: str = cr.execId
|
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,
|
await emit_pp_update(
|
||||||
# - we haven't already
|
ems_stream,
|
||||||
# - the fill event has already arrived
|
accounts_def,
|
||||||
# but it didn't yet have a commision report
|
proxies,
|
||||||
# which we fill in now.
|
ledgers,
|
||||||
|
tables,
|
||||||
# placehold i guess until someone who know wtf
|
ibpos=pos,
|
||||||
# contract this is from can fill it in...
|
fill=fill,
|
||||||
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)
|
|
||||||
|
|
||||||
# always update with latest ib pos msg info since
|
# always update with latest ib pos msg info since
|
||||||
# we generally audit against it for sanity and
|
# 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)
|
bs_mktid, ppmsg = pack_position(pos, accounts_def)
|
||||||
log.info(f'New IB position msg: {ppmsg}')
|
log.info(f'New IB position msg: {ppmsg}')
|
||||||
|
|
||||||
_, fill = clears.setdefault(
|
key = (pos.account, con.conId)
|
||||||
con,
|
positions[key] = pos
|
||||||
[pos, None],
|
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
|
# only send a pos update once we've actually rxed
|
||||||
# the msg from IB since generally speaking we use
|
# the msg from IB since generally speaking we use
|
||||||
# their 'cumsize' as gospel.
|
# their 'cumsize' as gospel.
|
||||||
|
|
|
||||||
|
|
@ -429,9 +429,16 @@ def api_trades_to_ledger_entries(
|
||||||
for attr_name, val in fdict.items():
|
for attr_name, val in fdict.items():
|
||||||
match attr_name:
|
match attr_name:
|
||||||
# value is a `@dataclass` subtype
|
# value is a `@dataclass` subtype
|
||||||
case 'contract' | 'execution' | 'commissionReport':
|
case 'contract' | 'execution':
|
||||||
txn_dict.update(asdict(val))
|
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':
|
case 'time':
|
||||||
# ib has wack ns timestamps, or is that us?
|
# ib has wack ns timestamps, or is that us?
|
||||||
continue
|
continue
|
||||||
|
|
|
||||||
|
|
@ -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'),
|
||||||
|
]
|
||||||
Loading…
Reference in New Issue