Skip to content
Open
8 changes: 7 additions & 1 deletion docs/market-data-architecture.md
Original file line number Diff line number Diff line change
Expand Up @@ -86,7 +86,13 @@ BarService federates:

- vendor K-lines from the embedded provider adapters;
- broker/exchange K-lines exposed through UTA;
- source metadata such as capability and freshness.
- source metadata such as capability and freshness;
- the trading SESSION a broker series covers.

Session is resolved by UTA per instrument (stocks/options default to regular
hours, FX/futures/crypto to the continuous tape), never by the caller, and the
effective value plus a `sessionForced` flag ride back on `BarMeta`. See the
IBKR adapter README for the default table and the broker capability fallback.

UTA source discovery and Broker Pack installation are independent. `asVendor`
controls whether a configured UTA joins default K-line/contract discovery;
Expand Down
10 changes: 6 additions & 4 deletions docs/uta-live-testing.md
Original file line number Diff line number Diff line change
Expand Up @@ -250,12 +250,14 @@ cancel. *Guards: editOrder venue quirks, id truncation.*
`placeOrderWithTpSl` override this must REFUSE loudly (never place a naked
entry). On a verified venue: after fill, confirm BOTH protective legs exist
on the exchange — including the trigger/algo namespace — before calling it
working. On a native-bracket venue (Alpaca): the push result must carry
working. On a native-bracket venue (Alpaca, IBKR): the push result must carry
`legs` ids, and after the entry fills `order list` must show BOTH legs as
tracked orders. The held SL leg never appears in the venue's open-orders
listing (Alpaca holds it while the TP works) — place-time is the ONLY
tracked orders. Alpaca's held SL leg never appears in the venue's open-orders
listing (it stays `held` while the TP works): place-time is the ONLY
moment Alice can learn it exists, so a venue listing diff can NOT recover
a missed leg. *Guards: the silent unprotected-position failure (okx,
a missed leg. IBKR children do appear (an untriggered stop is
`PreSubmitted`) and share an OCA group with `ocaType=1`; the parent is not
in that group. *Guards: the silent unprotected-position failure (okx,
ledger lied protected) and its mirror, the naked ledger (alpaca, ledger
blind to real protection) — both fatal to "trust the log".*

Expand Down
25 changes: 21 additions & 4 deletions packages/ibkr/src/client/base.ts
Original file line number Diff line number Diff line change
Expand Up @@ -88,6 +88,9 @@ export class EClient {
}

reset(): void {
// A batched drain may still hold frames from this connection, and this
// instance is reused on reconnect.
this.reader?.stop()
this.decoder = null
this.conn = null
this.host = null
Expand Down Expand Up @@ -307,10 +310,24 @@ export class EClient {
throw error
}

if (frame.kind === 'protobuf') {
this.decoder.processProtoBuf(frame.payload, frame.msgId)
} else {
this.decoder.interpret(frame.msgId, readFields(frame.payload))
try {
if (frame.kind === 'protobuf') {
this.decoder.processProtoBuf(frame.payload, frame.msgId)
} else {
this.decoder.interpret(frame.msgId, readFields(frame.payload))
}
} catch (error) {
// Lost framing alignment is unrecoverable, so rethrow and let the reader
// tear the connection down. Anything else is a consumer-side defect on one
// already-parsed message.
if (error instanceof BadMessage) throw error
this.wrapper.error(
NO_VALID_ID,
currentTimeMillis(),
errors.BAD_MESSAGE.code(),
`Handler failure for ${frame.kind} msgId=${frame.msgId}`,
'',
)
}
}

Expand Down
5 changes: 4 additions & 1 deletion packages/ibkr/src/client/orders.ts
Original file line number Diff line number Diff line change
Expand Up @@ -349,7 +349,10 @@ export function applyOrders(Client: typeof EClient): void {
flds.push(makeField(order.action))

if (this.serverVersion() >= SV.MIN_SERVER_VER_FRACTIONAL_POSITIONS) {
flds.push(makeField(order.totalQuantity))
// A monetary-value order leaves `totalQuantity` unset, and a raw
// sentinel goes out as a 39-digit share count that TWS rejects with
// `320 Error reading request. Unable to parse field`.
flds.push(makeFieldHandleEmpty(order.totalQuantity))
} else {
flds.push(makeField(order.totalQuantity.toNumber() | 0))
}
Expand Down
13 changes: 11 additions & 2 deletions packages/ibkr/src/comm.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,15 @@ import { UNSET_INTEGER, UNSET_DOUBLE, UNSET_DECIMAL, DOUBLE_INFINITY, INFINITY_S
import { ClientException, isAsciiPrintable } from './utils.js'
import { INVALID_SYMBOL } from './errors.js'

/**
* Realm-safe `Decimal` detection. A broker pack can load a second `decimal.js`
* copy, and `instanceof` would then miss the UNSET sentinel and transmit
* 2^127-1 as a real value.
*/
function isDecimal(val: unknown): val is Decimal {
return Decimal.isDecimal(val)
}

/**
* Wrap protobuf data with 4-byte big-endian length prefix and msgId.
* Wire format: [4-byte total length][4-byte msgId BE][protobuf bytes]
Expand Down Expand Up @@ -60,7 +69,7 @@ export function makeField(val: unknown): string {

// Decimal: use toFixed() to avoid scientific notation on small values
// (Decimal.toString() uses '1e-8' by default; TWS wire expects '0.00000001').
if (val instanceof Decimal) {
if (isDecimal(val)) {
return val.toFixed() + '\0'
}

Expand Down Expand Up @@ -89,7 +98,7 @@ export function makeFieldHandleEmpty(val: unknown): string {
throw new Error('Cannot send None to TWS')
}

if (val instanceof Decimal) {
if (isDecimal(val)) {
if (val.equals(UNSET_DECIMAL)) return makeField('')
return makeField(val)
}
Expand Down
63 changes: 63 additions & 0 deletions packages/ibkr/src/connection.spec.ts
Original file line number Diff line number Diff line change
@@ -0,0 +1,63 @@
import { describe, it, expect, vi, afterEach } from 'vitest'
import net from 'node:net'
import { EventEmitter } from 'node:events'
import { Connection } from './connection.js'

/** A net.Socket stand-in that never completes its TCP connect, as a recreated
* IB Gateway container looks while its port is not yet reachable. */
function stubPendingSocket(): EventEmitter & { destroy: () => void } {
const socket = new EventEmitter() as EventEmitter & {
connect: (port: number, host: string, cb: () => void) => void
destroy: () => void
write: () => boolean
}
socket.connect = () => { /* stays in SYN_SENT: no 'connect', no 'error' */ }
socket.destroy = () => { socket.emit('close') }
socket.write = () => true
// Callers use `new net.Socket()`, so the implementation must be
// constructible and an arrow function will not do.
vi.spyOn(net, 'Socket').mockImplementation(function () { return socket as unknown as net.Socket })
return socket
}

describe('Connection.connect — terminal settlement', () => {
afterEach(() => { vi.restoreAllMocks() })

// destroy() emits 'close' but not 'error', so a promise listening only for
// 'connect'/'error' never settles.
it('rejects a still-connecting attempt when disconnect() destroys the socket', async () => {
stubPendingSocket()
const conn = new Connection('127.0.0.1', 4002)
conn.wrapper = { error: () => {}, connectionClosed: () => {} }

const attempt = conn.connect()
const settled = vi.fn()
void attempt.then(settled, settled)
await Promise.resolve()
expect(settled).not.toHaveBeenCalled()

conn.disconnect()
await expect(attempt).rejects.toThrow()
})

it('rejects when the peer closes the socket before the connect callback', async () => {
const socket = stubPendingSocket()
const conn = new Connection('127.0.0.1', 4002)
conn.wrapper = { error: () => {}, connectionClosed: () => {} }

const attempt = conn.connect()
socket.emit('close')
await expect(attempt).rejects.toThrow()
})

it('still resolves normally once the socket connects', async () => {
const socket = stubPendingSocket() as EventEmitter & {
connect: (port: number, host: string, cb: () => void) => void
}
socket.connect = (_port, _host, cb) => { cb() }
const conn = new Connection('127.0.0.1', 4002)
conn.wrapper = { error: () => {}, connectionClosed: () => {} }

await expect(conn.connect()).resolves.toBeUndefined()
})
})
16 changes: 14 additions & 2 deletions packages/ibkr/src/connection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -32,6 +32,15 @@ export class Connection extends EventEmitter {

connect(): Promise<void> {
return new Promise<void>((resolve, reject) => {
// The attempt must settle on every terminal socket outcome: destroying a
// socket still in SYN_SENT emits 'close' without 'error', and a promise
// watching only 'connect'/'error' would hang every awaiting caller.
let settled = false
const settle = (fn: () => void): void => {
if (settled) return
settled = true
fn()
}
try {
this.socket = new net.Socket()
} catch {
Expand All @@ -47,6 +56,9 @@ export class Connection extends EventEmitter {
})

this.socket.on('close', () => {
// Reject before the cleanup guard below, which returns early when
// disconnect() already nulled the field.
settle(() => reject(new Error(CONNECT_FAIL.msg())))
// Guard: if socket is already null, disconnect() already handled cleanup.
// Without this check, connectionClosed() would be called twice when
// disconnect() is invoked (once by disconnect, once by the close event).
Expand All @@ -69,14 +81,14 @@ export class Connection extends EventEmitter {
})

this.socket.connect(this.port, this.host, () => {
resolve()
settle(resolve)
})

this.socket.once('error', (err: Error) => {
if (this.wrapper) {
this.wrapper.error(NO_VALID_ID, currentTimeMillis(), CONNECT_FAIL.code(), CONNECT_FAIL.msg())
}
reject(err)
settle(() => reject(err))
})
})
}
Expand Down
109 changes: 87 additions & 22 deletions packages/ibkr/src/reader.ts
Original file line number Diff line number Diff line change
Expand Up @@ -3,16 +3,27 @@
* Mirrors: ibapi/reader.py
*
* Node.js adaptation: Python uses a background thread + queue. Here we use
* socket 'data' events → buffer accumulation → message extraction → callback.
* No threads, no queue.
* socket 'data' events → buffer accumulation → frame extraction → bounded
* drain → callback. Dispatch is batched so a large inbound burst cannot
* monopolise the event loop and starve other layers' timers.
*/

import { readMsg } from './comm.js'
import type { Connection } from './connection.js'

/** Messages dispatched per macrotask before yielding back to the event loop. */
export const MAX_MESSAGES_PER_TURN = 200

export class EReader {
private conn: Connection
private buf: Buffer = Buffer.alloc(0)
private queue: Buffer[] = []
/** Read cursor into `queue`. Removing from the head per message would make a
* burst quadratic, so the dispatched prefix is dropped once per batch. */
private head = 0
private draining = false
private stopped = false
private dataListener: (() => void) | null = null
private onMessage: (msg: Buffer) => void
private onError?: (error: unknown) => void

Expand All @@ -30,41 +41,95 @@ export class EReader {
* Start listening for incoming data.
*/
start(): void {
this.conn.on('data', () => {
if (this.stopped || this.dataListener) return
this.dataListener = (): void => {
this.processData()
})
}
this.conn.on('data', this.dataListener)
}

/**
* Detaches from the socket and discards everything still queued. Idempotent
* and terminal: `start()` will not re-arm a stopped reader.
*/
stop(): void {
this.stopped = true
this.draining = false
this.queue.length = 0
this.head = 0
this.buf = Buffer.alloc(0)
if (this.dataListener) {
this.conn.off('data', this.dataListener)
this.dataListener = null
}
}

/**
* Process accumulated socket data, extracting complete messages.
*/
private processData(): void {
if (this.stopped) return

// Consume whatever has accumulated in the connection buffer
const incoming = this.conn.consumeBuffer()
if (incoming.length === 0) return

this.buf = Buffer.concat([this.buf, incoming])

// Extract as many complete messages as possible
// Framing must stay ordered and synchronous; only dispatch is deferred.
while (this.buf.length > 0) {
const [size, msg, rest] = readMsg(this.buf)
if (msg.length > 0) {
this.buf = rest
try {
this.onMessage(msg)
} catch (error) {
// A decoder failure means field alignment is no longer trustworthy.
// Drop every buffered successor and let the client replace this
// connection instead of continuing from an uncertain boundary.
this.buf = Buffer.alloc(0)
if (!this.onError) throw error
this.onError(error)
break
}
} else {
// Incomplete message — wait for more data
break
const [, msg, rest] = readMsg(this.buf)
// Incomplete message: wait for more data
if (msg.length === 0) break
this.buf = rest
this.queue.push(msg)
}

if (!this.draining) {
this.draining = true
this.drain()
}
}

/**
* Dispatch queued frames in order, at most MAX_MESSAGES_PER_TURN per
* macrotask. Yielding with setImmediate lets pending timers fire between
* batches instead of after the whole burst.
*/
private drain(): void {
let dispatched = 0

while (this.head < this.queue.length) {
if (dispatched >= MAX_MESSAGES_PER_TURN) {
this.queue = this.queue.slice(this.head)
this.head = 0
setImmediate(() => {
if (this.stopped) {
this.draining = false
return
}
this.drain()
})
return
}

const msg = this.queue[this.head]!
this.head++
dispatched++
try {
this.onMessage(msg)
} catch (error) {
// Field alignment is no longer trustworthy, so drop every buffered
// successor and let the client replace this connection.
this.stop()
if (!this.onError) throw error
this.onError(error)
return
}
}

this.queue.length = 0
this.head = 0
this.draining = false
}
}
Loading