API Reference¶
Use this page for signatures and public entry points. Use the User Guide for workflows, rollout order, and broker/feed selection.
Recommended Imports¶
from ml4t.live import (
AlpacaBroker,
AlpacaDataFeed,
AsyncBrokerProtocol,
BarAggregator,
BarBuffer,
BrokerProtocol,
BrokerOrderContractError,
AcceptedOrderPersistenceError,
AuditJournalError,
CanonicalOrderRequest,
ConcurrentStateWriterError,
CorruptStateError,
DataFeedProtocol,
FeedContinuityError,
FeedContractError,
FeedOverflowError,
FeedQueueSnapshot,
IBBroker,
IBDataFeed,
LiveEngine,
LiveRiskConfig,
OKXFundingFeed,
OrderReplacementGapError,
OrderValidationError,
PersistenceSafetyError,
ReconciliationMismatchError,
RiskLimitError,
RiskState,
RuntimeCleanupError,
RuntimeErrorContext,
RuntimeFailureError,
RuntimeState,
RuntimeTransition,
SafeBroker,
ThreadSafeBrokerWrapper,
UnsafePersistencePathError,
VirtualPortfolio,
runtime_error_context,
)
Public Surface At A Glance¶
| Group | Primary symbols |
|---|---|
| Engine | LiveEngine, RuntimeState, RuntimeTransition, RuntimeErrorContext, runtime_error_context, RuntimeFailureError, RuntimeCleanupError |
| Brokers | IBBroker, AlpacaBroker |
| Stable-supported feeds | OKXFundingFeed |
| Experimental feeds | AlpacaDataFeed, IBDataFeed, DataBentoFeed, CryptoFeed, ExperimentalFeedError, ExperimentalFeedWarning |
| Feed helpers | BarAggregator, BarBuffer, FeedContractError, FeedContinuityError, FeedOverflowError, FeedQueueSnapshot |
| Orders | CanonicalOrderRequest, OrderValidationError, BrokerOrderContractError |
| Safety | LiveRiskConfig, SafeBroker, RiskState, RiskLimitError, PersistenceSafetyError, AuditJournalError, AcceptedOrderPersistenceError, VirtualPortfolio |
| Sync/async bridge | ThreadSafeBrokerWrapper |
| Protocols | BrokerProtocol, AsyncBrokerProtocol, DataFeedProtocol |
Engine¶
LiveEngine
¶
LiveEngine(
strategy,
broker,
feed,
*,
on_error=None,
feed_silence_seconds=None,
watchdog_poll_seconds=1.0,
halt_on_unhealthy=False,
auto_recover=False,
recovery_cooldown_seconds=5.0,
max_recovery_attempts=3,
max_event_age_seconds=None,
on_health_change=None,
strategy_callback_timeout_seconds=5.0,
lifecycle_version=V1,
execution_policy=None,
strategy_config=None,
)
Async live trading engine.
Bridges async infrastructure with sync Strategy.on_data().
Initialize LiveEngine.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
strategy
|
Strategy
|
Strategy instance to execute. |
required |
broker
|
AsyncBrokerProtocol
|
Async broker implementation. |
required |
feed
|
DataFeedProtocol
|
Data feed providing timestamp, data, context tuples. |
required |
on_error
|
Callable[[Exception, datetime, dict], None] | None
|
Custom error handler callback. |
None
|
feed_silence_seconds
|
float | None
|
Optional threshold for degraded feed reporting. |
None
|
watchdog_poll_seconds
|
float
|
Poll interval for runtime health monitoring. |
1.0
|
halt_on_unhealthy
|
bool
|
Stop the engine when watchdog detects a degraded state. |
False
|
auto_recover
|
bool
|
Attempt reconnect/restart when watchdog detects a recoverable state. |
False
|
recovery_cooldown_seconds
|
float
|
Delay between recovery attempts. |
5.0
|
max_recovery_attempts
|
int
|
Maximum recovery attempts before stopping. |
3
|
max_event_age_seconds
|
float | None
|
Maximum provider-event age before dispatch. When omitted, use the supported feed's declared limit if present. |
None
|
on_health_change
|
Callable[[str, dict[str, Any]], None] | None
|
Optional callback invoked when runtime health changes. |
None
|
strategy_callback_timeout_seconds
|
float
|
Maximum callback duration. A callback that exceeds this duration is allowed to become quiescent before a typed timeout aborts the run. |
5.0
|
lifecycle_version
|
LifecycleVersion | str
|
Portable strategy lifecycle version. |
V1
|
execution_policy
|
ExecutionPolicy | None
|
Explicit live execution capabilities and behavior. |
None
|
strategy_config
|
BacktestConfig | None
|
Backtest strategy configuration supplied to |
None
|
Source code in src/ml4t/live/engine.py
261 262 263 264 265 266 267 268 269 270 271 272 273 274 275 276 277 278 279 280 281 282 283 284 285 286 287 288 289 290 291 292 293 294 295 296 297 298 299 300 301 302 303 304 305 306 307 308 309 310 311 312 313 314 315 316 317 318 319 320 321 322 323 324 325 326 327 328 329 330 331 332 333 334 335 336 337 338 339 340 341 342 343 344 345 346 347 348 349 350 351 352 353 354 355 356 357 358 359 360 361 362 363 364 365 366 367 368 369 370 371 372 373 374 375 376 377 378 379 380 381 | |
runtime_transitions
property
¶
Return retained state transitions in occurrence order.
connect
async
¶
Acquire the broker and feed transactionally and become ready.
Source code in src/ml4t/live/engine.py
run
async
¶
Main async loop - receives bars and dispatches to strategy.
Source code in src/ml4t/live/engine.py
449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 520 521 522 523 524 525 526 527 528 529 530 531 532 533 534 535 536 537 538 539 540 541 542 543 544 545 546 547 548 549 550 551 552 553 554 555 556 557 558 559 560 561 562 563 564 565 566 567 568 569 570 571 572 573 574 575 576 577 578 579 580 581 582 583 584 585 586 587 588 589 590 591 592 593 594 595 596 597 598 599 600 601 602 603 604 605 606 607 608 609 610 611 612 613 614 615 616 617 618 619 620 621 622 623 624 625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 705 706 707 708 709 710 711 712 713 714 715 716 717 718 719 720 721 722 723 724 725 726 727 728 729 730 731 732 733 734 735 736 737 | |
stop
async
¶
Request shutdown and release resources exactly once.
Source code in src/ml4t/live/engine.py
runtime_status
¶
Return engine runtime health and session context.
Source code in src/ml4t/live/engine.py
1463 1464 1465 1466 1467 1468 1469 1470 1471 1472 1473 1474 1475 1476 1477 1478 1479 1480 1481 1482 1483 1484 1485 1486 1487 1488 1489 1490 1491 1492 1493 1494 1495 1496 1497 1498 1499 1500 1501 1502 1503 1504 1505 1506 1507 1508 1509 1510 1511 1512 1513 1514 1515 1516 1517 1518 1519 1520 1521 1522 1523 1524 1525 1526 1527 1528 1529 1530 1531 1532 1533 1534 1535 1536 1537 1538 1539 1540 1541 1542 1543 1544 1545 1546 1547 1548 | |
RuntimeState
¶
Bases: StrEnum
One explicit phase of engine resource and strategy ownership.
RuntimeTransition
dataclass
¶
Structured evidence for one runtime state change.
RuntimeErrorContext
dataclass
¶
Redacted operator context attached to a runtime exception.
to_dict
¶
Return machine-readable context without exception text.
Source code in src/ml4t/live/engine.py
runtime_error_context
¶
Return structured runtime context when the engine attached it.
Source code in src/ml4t/live/engine.py
RuntimeFailureError
¶
Bases: RuntimeError
Raised when an asynchronous runtime failure reaches a terminal state.
Source code in src/ml4t/live/engine.py
RuntimeCleanupError
¶
Bases: RuntimeError
Raised when runtime finalization cannot release every acquired resource.
Source code in src/ml4t/live/engine.py
Safety And Rollout¶
CanonicalOrderRequest
dataclass
¶
One unsigned venue request used unchanged for checks and submission.
validate_result
¶
Require the adapter result to describe this exact request.
Source code in src/ml4t/live/orders.py
LiveRiskConfig
dataclass
¶
LiveRiskConfig(
max_position_value=50000.0,
max_position_shares=1000.0,
max_total_exposure=200000.0,
max_positions=20,
max_order_value=10000.0,
max_order_shares=500.0,
max_orders_per_minute=10,
max_daily_loss=5000.0,
max_drawdown_pct=0.05,
max_price_deviation_pct=0.05,
max_data_staleness_seconds=60.0,
dedup_window_seconds=1.0,
allowed_assets=set(),
blocked_assets=set(),
shadow_mode=False,
execution_mode=None,
kill_switch_enabled=False,
allow_reducing_risk_when_killed=True,
halt_on_reducing_risk_failure=True,
fail_on_reconciliation_mismatch=False,
state_file=".ml4t_risk_state.json",
journal_file=None,
fail_on_journal_error=True,
)
Risk configuration for live trading.
Multiple layers of protection. Set a limit to None to disable that
specific check. NaN and infinity are always invalid.
Example
Conservative configuration¶
config = LiveRiskConfig( max_position_value=25_000.0, max_daily_loss=2_000.0, execution_mode="shadow", )
Disable a specific check explicitly¶
config = LiveRiskConfig( max_position_value=None, max_daily_loss=10_000.0, # Only daily loss limit )
Safety Recommendations
- Always start with execution_mode="shadow"
- Graduate to paper trading
- Use small positions when going live
- Set conservative risk limits
__post_init__
¶
Validate configuration parameters.
Source code in src/ml4t/live/safety.py
require_execution_mode
¶
Return the explicit mode or reject ambiguous external execution.
Source code in src/ml4t/live/safety.py
SafeBroker
¶
Risk-controlled wrapper with state persistence.
Safety Features: 1. Pre-trade validation against all risk limits 2. Order rate limiting 3. Drawdown monitoring with kill switch 4. Fat finger protection (price deviation check) 5. Stale data protection 6. Duplicate order filter 7. Shadow mode with VirtualPortfolio (realistic paper trading) 8. Owner-only, versioned state and chained audit persistence across restarts
Example
broker = IBBroker() await broker.connect()
safe = SafeBroker( broker=broker, config=LiveRiskConfig( max_position_value=25000, execution_mode="shadow", ) )
Use safe in strategy¶
engine = LiveEngine(strategy, safe, feed)
Initialize SafeBroker.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
broker
|
AsyncBrokerProtocol
|
Async broker implementation (IBBroker, AlpacaBroker, etc.) |
required |
config
|
LiveRiskConfig
|
Risk configuration |
required |
Source code in src/ml4t/live/safety.py
positions
property
¶
Get current positions.
In shadow mode, returns virtual positions. In live mode, returns broker positions.
reconciliation_report
property
¶
Return the latest startup reconciliation report.
persistence_status
property
¶
Return the current state and journal health without secret values.
execution_capabilities
property
¶
Return capabilities declared by the wrapped venue.
close_persistence
¶
assert_paper_trading
¶
Fail unless this wrapper and its provider are configured for paper execution.
Source code in src/ml4t/live/safety.py
assert_live_trading
¶
Fail unless this wrapper and its provider are configured for live execution.
Source code in src/ml4t/live/safety.py
load_portable_strategy_state
¶
save_portable_strategy_state
¶
Persist target and position-rule state with the safety state.
get_position
¶
Get position for specific asset.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
Returns:
| Type | Description |
|---|---|
Position | None
|
Position object or None |
Source code in src/ml4t/live/safety.py
get_positions
¶
get_account_value_async
async
¶
Get total account value (async).
Returns:
| Type | Description |
|---|---|
float
|
Total account value in base currency |
Source code in src/ml4t/live/safety.py
get_cash_async
async
¶
Get available cash (async).
Returns:
| Type | Description |
|---|---|
float
|
Available cash in base currency |
Source code in src/ml4t/live/safety.py
cancel_order_async
async
¶
Cancel pending order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
order_id
|
str
|
ID of order to cancel |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if cancel request submitted |
Source code in src/ml4t/live/safety.py
close_position_async
async
¶
Close entire position.
Close positions bypass normal limits (safety feature).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol to close |
required |
Returns:
| Type | Description |
|---|---|
Order | None
|
Order object if position exists |
Source code in src/ml4t/live/safety.py
reduce_position_async
async
¶
Submit an explicitly reducing order under the configured kill-switch policy.
Source code in src/ml4t/live/safety.py
replace_order_async
async
¶
Replace a pending order via cancel-and-resubmit.
Source code in src/ml4t/live/safety.py
submit_order_async
async
¶
submit_order_async(
asset,
quantity,
side=None,
order_type=MARKET,
limit_price=None,
stop_price=None,
**kwargs,
)
Submit order with full risk validation.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
quantity
|
float
|
Signed shares/contracts when side is omitted; positive unsigned shares/contracts when side is provided |
required |
side
|
OrderSide | None
|
Order side (BUY/SELL), inferred from signed quantity if omitted |
None
|
order_type
|
OrderType
|
Type of order |
MARKET
|
limit_price
|
float | None
|
Limit price for LIMIT/STOP_LIMIT orders |
None
|
stop_price
|
float | None
|
Stop price for STOP/STOP_LIMIT orders |
None
|
**kwargs
|
Any
|
Additional broker-specific parameters |
{}
|
Returns:
| Type | Description |
|---|---|
Order
|
Order object |
Raises:
| Type | Description |
|---|---|
RiskLimitError
|
If order violates any risk limit |
Source code in src/ml4t/live/safety.py
record_market_snapshot
¶
Cache a single price observation for the staleness guard.
The supported way to keep the cache fresh is the streaming path:
a Feed (e.g. IBDataFeed) emits ticks, LiveEngine shuttles
them into _record_market_data on every bar, and the cache stays
current automatically. Use that for any continuous-loop deployment.
This method is the non-streaming escape hatch for one-shot flows that legitimately have no tick stream in front of the broker:
- A CLI flatten tool that takes a position list and submits MOC/MARKET closeouts.
- A REST-only broker adapter that fetches quotes synchronously per request rather than via a streaming feed.
- A test harness setting up controlled state.
It is not the right tool inside a notebook that already runs a
live engine - there the streaming path covers staleness implicitly.
Reaching for record_market_snapshot from inside a tick-driven
flow is a code smell: it usually means the engine isn't actually
wired up, and the snapshot will go stale silently.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Symbol to record. |
required |
price
|
float
|
Reference price (e.g. last close, mid quote, snapshot top-of-book mid). Must be > 0. |
required |
timestamp
|
datetime | None
|
Bar/quote timestamp. Defaults to now (UTC). The
|
None
|
Source code in src/ml4t/live/safety.py
enable_kill_switch
¶
Manually enable kill switch.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
reason
|
str
|
Reason for activation (default: "Manual") |
'Manual'
|
disable_kill_switch
¶
Manually disable kill switch (use with caution!).
Source code in src/ml4t/live/safety.py
close_all_positions
async
¶
Emergency close all positions.
Returns:
| Type | Description |
|---|---|
list[Order]
|
List of close orders |
Source code in src/ml4t/live/safety.py
record_event
¶
Append a structured runtime event to the execution journal.
Source code in src/ml4t/live/safety.py
preview_reconciliation_async
async
¶
Build the current runtime reconciliation report without mutating state.
Source code in src/ml4t/live/safety.py
preflight_async
async
¶
Probe broker reachability and startup reconciliation without persisting state.
Source code in src/ml4t/live/safety.py
connect
async
¶
Connect to broker and reconcile persisted state.
Source code in src/ml4t/live/safety.py
disconnect
async
¶
Disconnect from broker and save state.
Source code in src/ml4t/live/safety.py
is_connected_async
async
¶
get_positions_async
async
¶
Get all positions (async).
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbol to Position |
Source code in src/ml4t/live/safety.py
get_pending_orders_async
async
¶
Get pending orders (async).
Returns:
| Type | Description |
|---|---|
list[Order]
|
List of pending orders |
Source code in src/ml4t/live/safety.py
get_position_async
async
¶
Get position (async).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
Returns:
| Type | Description |
|---|---|
Position | None
|
Position object or None |
Source code in src/ml4t/live/safety.py
ThreadSafeBrokerWrapper
¶
Wraps an async broker for use from sync strategy code.
This wrapper is passed to Strategy.on_data() instead of the raw broker. It bridges the sync/async boundary by scheduling coroutines on the main event loop and blocking the worker thread until they complete.
Thread Safety: - Every portable strategy callback runs on one dedicated worker thread - Broker methods run on main event loop - run_coroutine_threadsafe() handles the cross-thread communication
Timeouts (from design review): - Getters (get_cash, get_account_value): 5s - Order operations (submit, cancel, close): 30s
Example
LiveEngine creates this wrapper¶
loop = asyncio.get_running_loop() wrapped = ThreadSafeBrokerWrapper(ib_broker, loop)
Strategy uses it like a normal sync broker¶
order = wrapped.submit_order('AAPL', 100, OrderSide.BUY)
Note
This class implements BrokerProtocol but does not inherit from it. It provides a sync interface backed by async operations.
Initialize thread-safe wrapper.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
async_broker
|
AsyncBrokerProtocol
|
Async broker implementation (IBBroker, etc.) |
required |
loop
|
AbstractEventLoop
|
Main event loop (from asyncio.get_running_loop()) |
required |
Source code in src/ml4t/live/wrappers.py
positions
property
¶
Get current positions (thread-safe read).
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbol to Position |
pending_orders
property
¶
Get pending orders (thread-safe read).
Returns:
| Type | Description |
|---|---|
list[Order]
|
List of pending Order objects |
is_connected
property
¶
Check if broker is connected.
Returns:
| Type | Description |
|---|---|
bool
|
True if connected and ready to trade |
get_position
¶
Get position for specific asset.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol (e.g., "AAPL") |
required |
Returns:
| Type | Description |
|---|---|
Position | None
|
Position object if holding position, None otherwise |
Raises:
| Type | Description |
|---|---|
TimeoutError
|
If operation times out |
RuntimeError
|
If broker error occurs |
Source code in src/ml4t/live/wrappers.py
get_positions
¶
get_account_value
¶
Get total account value (cash + positions).
Returns:
| Type | Description |
|---|---|
float
|
Total account value in base currency |
Raises:
| Type | Description |
|---|---|
TimeoutError
|
If operation times out (5s) |
RuntimeError
|
If broker error occurs |
Source code in src/ml4t/live/wrappers.py
get_cash
¶
Get available cash balance.
Returns:
| Type | Description |
|---|---|
float
|
Available cash in base currency |
Raises:
| Type | Description |
|---|---|
TimeoutError
|
If operation times out (5s) |
RuntimeError
|
If broker error occurs |
Source code in src/ml4t/live/wrappers.py
get_pending_orders
¶
Get pending orders, optionally filtered by asset.
Source code in src/ml4t/live/wrappers.py
submit_order
¶
submit_order(
asset,
quantity,
side=None,
order_type=MARKET,
limit_price=None,
stop_price=None,
**kwargs,
)
Submit order for execution.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol (e.g., "AAPL") |
required |
quantity
|
float
|
Signed shares/contracts when side is omitted; positive unsigned shares/contracts when side is provided |
required |
side
|
OrderSide | None
|
Order side (BUY/SELL), inferred from signed quantity if omitted |
None
|
order_type
|
OrderType
|
Type of order (MARKET, LIMIT, STOP, etc.) |
MARKET
|
limit_price
|
float | None
|
Limit price for LIMIT/STOP_LIMIT orders |
None
|
stop_price
|
float | None
|
Stop price for STOP/STOP_LIMIT orders |
None
|
**kwargs
|
Any
|
Additional broker-specific parameters |
{}
|
Returns:
| Type | Description |
|---|---|
Order
|
Order object with order_id and initial status |
Raises:
| Type | Description |
|---|---|
TimeoutError
|
If operation times out (30s) |
ValueError
|
If order parameters are invalid |
RuntimeError
|
If broker is not connected or error occurs |
Source code in src/ml4t/live/wrappers.py
cancel_order
¶
Cancel pending order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
order_id
|
str
|
ID of order to cancel |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if cancel request submitted, False if order not found |
Raises:
| Type | Description |
|---|---|
TimeoutError
|
If operation times out (30s) |
RuntimeError
|
If broker error occurs |
Source code in src/ml4t/live/wrappers.py
replace_order
¶
Replace a pending order with updated parameters.
Source code in src/ml4t/live/wrappers.py
close_position
¶
Close entire position in asset.
Convenience method that submits a closing order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol to close |
required |
Returns:
| Type | Description |
|---|---|
Order | None
|
Order object if position exists, None if no position |
Raises:
| Type | Description |
|---|---|
TimeoutError
|
If operation times out (30s) |
RuntimeError
|
If broker error occurs |
Source code in src/ml4t/live/wrappers.py
register_target_intent
¶
Register a persistent target intent during a causal callback.
Source code in src/ml4t/live/wrappers.py
register_position_rule_policy
¶
Bind a portable position-rule policy to its client implementation.
get_target_intents
¶
get_child_order_intents
¶
get_intent_reconciliations
¶
export_target_intent_state
¶
set_position_rules
¶
Set client-evaluated position rules globally or for one asset.
clear_position_rules
¶
Clear client-evaluated position rules globally or for one asset.
update_position_context
¶
Merge portable context used by position-rule evaluation.
Brokers¶
IBBroker
¶
Interactive Brokers implementation.
Design: - All broker operations are async - Uses asyncio.Lock for thread safety - Event handlers use put_nowait() (non-blocking) - Reconnection handled externally
Connection Ports: - TWS Paper: 7497 - TWS Live: 7496 - Gateway Paper: 4002 - Gateway Live: 4001
Example
broker = IBBroker(port=7497) # Paper trading await broker.connect() positions = await broker.get_positions_async() await broker.disconnect()
Initialize IBBroker.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
host
|
str
|
IB Gateway/TWS host (default: '127.0.0.1') |
'127.0.0.1'
|
port
|
int
|
IB Gateway/TWS port (default: 7497 for paper) |
7497
|
client_id
|
int
|
Unique client ID (default: 1) |
1
|
account
|
str | None
|
IB account ID (default: use first account) |
None
|
market_data_type
|
int | None
|
IB market-data type to request after connect.
|
None
|
Source code in src/ml4t/live/brokers/ib.py
execution_capabilities
property
¶
Return order behaviors implemented by this adapter.
positions
property
¶
Return a thread-safe snapshot of positions.
Note: This is called from worker thread via ThreadSafeBrokerWrapper. The lock prevents RuntimeError during dict iteration if IB callback modifies positions concurrently.
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbols to Position objects |
pending_orders
property
¶
Get list of pending orders.
Returns:
| Type | Description |
|---|---|
list[Order]
|
List of pending Order objects |
connect
async
¶
Connect to IB Gateway/TWS.
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If connection fails |
TimeoutError
|
If connection times out |
Source code in src/ml4t/live/brokers/ib.py
166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 201 202 203 204 205 206 207 208 209 210 211 212 213 214 215 216 217 218 219 220 221 222 223 224 225 226 227 228 229 230 231 232 233 234 235 236 237 238 239 240 241 242 243 244 245 246 247 248 249 250 251 252 253 | |
disconnect
async
¶
Disconnect from IB.
Source code in src/ml4t/live/brokers/ib.py
is_connected_async
async
¶
assert_paper_trading
¶
Fail unless the connected endpoint and managed account identify IB paper trading.
Source code in src/ml4t/live/brokers/ib.py
assert_live_trading
¶
Fail unless the connected live endpoint matches an explicitly selected account.
Source code in src/ml4t/live/brokers/ib.py
get_position
¶
Thread-safe single position access.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
Returns:
| Type | Description |
|---|---|
Position | None
|
Position object if exists, None otherwise |
get_positions_async
async
¶
Async thread-safe position access with lock.
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbols to Position objects |
Source code in src/ml4t/live/brokers/ib.py
get_position_async
async
¶
Return one position from the synchronized adapter snapshot.
get_pending_orders_async
async
¶
Return pending orders, optionally filtered by asset.
Source code in src/ml4t/live/brokers/ib.py
get_account_value_async
async
¶
Get Net Liquidation Value.
Returns:
| Type | Description |
|---|---|
float
|
Account net liquidation value in USD |
Source code in src/ml4t/live/brokers/ib.py
get_cash_async
async
¶
Get available funds.
Returns:
| Type | Description |
|---|---|
float
|
Available funds in USD |
Source code in src/ml4t/live/brokers/ib.py
submit_order_async
async
¶
submit_order_async(
asset,
quantity,
side=None,
order_type=MARKET,
limit_price=None,
stop_price=None,
**kwargs,
)
Submit order to IB.
TASK-013: Full order submission implementation with IB order tracking.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
quantity
|
float
|
Signed shares/contracts when side is omitted; positive unsigned shares/contracts when side is provided |
required |
side
|
OrderSide | None
|
BUY or SELL, inferred from signed quantity if omitted |
None
|
order_type
|
OrderType
|
Market, limit, stop, or stop-limit |
MARKET
|
limit_price
|
float | None
|
Limit price for limit orders |
None
|
stop_price
|
float | None
|
Stop price for stop orders |
None
|
Returns:
| Type | Description |
|---|---|
Order
|
Order object |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If not connected |
ValueError
|
If order parameters are invalid |
Source code in src/ml4t/live/brokers/ib.py
403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 | |
cancel_order_async
async
¶
Cancel pending order.
TASK-016: Full order cancellation implementation.
This method finds the IB order ID from our tracking map and cancels the order via the IB API. Handles edge cases like order not found or order already filled.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
order_id
|
str
|
Order ID to cancel (e.g., 'ML4T-1') |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if cancellation request sent successfully, False otherwise |
Note
The actual cancellation is confirmed via _on_order_status callback when IB sends the 'Cancelled' status update.
Source code in src/ml4t/live/brokers/ib.py
replace_order_async
async
¶
Replace a pending order via cancel-and-resubmit.
Source code in src/ml4t/live/brokers/ib.py
close_position_async
async
¶
Close position in asset.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
Returns:
| Type | Description |
|---|---|
Order | None
|
Order object if position exists, None otherwise |
Raises:
| Type | Description |
|---|---|
NotImplementedError
|
Depends on TASK-013 |
Source code in src/ml4t/live/brokers/ib.py
AlpacaBroker
¶
Alpaca Markets broker implementation.
Design (matching IBBroker patterns): - All broker operations are async - Uses asyncio.Lock for thread safety - WebSocket stream for real-time order updates - REST API for account/position queries and order submission
Paper vs Live: - paper=True (default): Uses paper trading endpoint - paper=False: Uses live trading endpoint (USE WITH CAUTION)
Example
broker = AlpacaBroker( api_key='PKXXXXXXXX', secret_key='XXXXXXXXXX', paper=True, # Always start with paper trading! ) await broker.connect() positions = await broker.get_positions_async() await broker.disconnect()
Initialize AlpacaBroker.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
api_key
|
str
|
Alpaca API key (from https://app.alpaca.markets) |
required |
secret_key
|
str
|
Alpaca secret key |
required |
paper
|
bool
|
Use paper trading endpoint (default: True) |
True
|
Source code in src/ml4t/live/brokers/alpaca.py
execution_capabilities
property
¶
Return order behaviors implemented by this adapter.
positions
property
¶
Thread-safe position access (shallow copy).
Note: This is called from worker thread via ThreadSafeBrokerWrapper. The shallow copy prevents RuntimeError during dict iteration.
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbols to Position objects |
pending_orders
property
¶
Get list of pending orders.
Returns:
| Type | Description |
|---|---|
list[Order]
|
List of pending Order objects |
connect
async
¶
Connect to Alpaca and sync initial state.
Steps: 1. Create TradingClient (REST) 2. Create TradingStream (WebSocket) 3. Verify connection by fetching account 4. Register trade update callback 5. Sync positions and open orders 6. Start WebSocket stream for order updates
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If connection fails |
Source code in src/ml4t/live/brokers/alpaca.py
118 119 120 121 122 123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 200 | |
disconnect
async
¶
Disconnect from Alpaca.
Source code in src/ml4t/live/brokers/alpaca.py
is_connected_async
async
¶
assert_paper_trading
¶
Fail unless the connected client is authenticated through Alpaca's paper endpoint.
Source code in src/ml4t/live/brokers/alpaca.py
assert_live_trading
¶
Fail unless the connected client is authenticated through Alpaca's live endpoint.
Source code in src/ml4t/live/brokers/alpaca.py
get_position
¶
Thread-safe single position access.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol (e.g., 'AAPL' or 'BTC/USD') |
required |
Returns:
| Type | Description |
|---|---|
Position | None
|
Position object if exists, None otherwise |
Source code in src/ml4t/live/brokers/alpaca.py
get_positions_async
async
¶
Async thread-safe position access with lock.
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbols to Position objects |
Source code in src/ml4t/live/brokers/alpaca.py
get_position_async
async
¶
Return one position from the synchronized adapter snapshot.
Source code in src/ml4t/live/brokers/alpaca.py
get_pending_orders_async
async
¶
Return pending orders, optionally filtered by asset.
Source code in src/ml4t/live/brokers/alpaca.py
get_account_value_async
async
¶
Get portfolio value (equity).
Returns:
| Type | Description |
|---|---|
float
|
Total account equity in USD |
Source code in src/ml4t/live/brokers/alpaca.py
get_cash_async
async
¶
Get the signed cash balance.
Returns:
| Type | Description |
|---|---|
float
|
Cash balance in USD. Margin accounts may report a negative balance. |
Source code in src/ml4t/live/brokers/alpaca.py
submit_order_async
async
¶
submit_order_async(
asset,
quantity,
side=None,
order_type=MARKET,
limit_price=None,
stop_price=None,
**kwargs,
)
Submit order to Alpaca.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol (e.g., 'AAPL' or 'BTC/USD') |
required |
quantity
|
float
|
Signed shares/units when side is omitted; positive unsigned shares/units when side is provided |
required |
side
|
OrderSide | None
|
BUY or SELL, inferred from signed quantity if omitted |
None
|
order_type
|
OrderType
|
Market, limit, stop, or stop-limit |
MARKET
|
limit_price
|
float | None
|
Limit price for limit orders |
None
|
stop_price
|
float | None
|
Stop price for stop orders |
None
|
**kwargs
|
Any
|
Additional parameters (ignored) |
{}
|
Returns:
| Type | Description |
|---|---|
Order
|
Order object |
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If not connected |
ValueError
|
If order parameters are invalid |
Source code in src/ml4t/live/brokers/alpaca.py
381 382 383 384 385 386 387 388 389 390 391 392 393 394 395 396 397 398 399 400 401 402 403 404 405 406 407 408 409 410 411 412 413 414 415 416 417 418 419 420 421 422 423 424 425 426 427 428 429 430 431 432 433 434 435 436 437 438 439 440 441 442 443 444 445 446 447 448 449 450 451 452 453 454 455 456 457 458 459 460 461 462 463 464 465 466 467 468 469 470 471 472 473 474 475 476 477 478 479 480 481 482 483 484 485 486 487 488 489 490 491 492 493 494 495 496 497 498 499 500 501 502 503 504 505 506 507 508 509 510 511 512 513 514 515 516 517 518 519 | |
cancel_order_async
async
¶
Cancel pending order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
order_id
|
str
|
Order ID to cancel (e.g., 'ML4T-1') |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if cancellation request sent successfully, False otherwise |
Source code in src/ml4t/live/brokers/alpaca.py
replace_order_async
async
¶
Replace a pending order via cancel-and-resubmit.
Source code in src/ml4t/live/brokers/alpaca.py
close_position_async
async
¶
Close position in asset.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol |
required |
Returns:
| Type | Description |
|---|---|
Order | None
|
Order object if position exists, None otherwise |
Source code in src/ml4t/live/brokers/alpaca.py
Data Feeds And Aggregation¶
IBDataFeed
¶
IBDataFeed(
ib,
symbols,
*,
exchange="SMART",
currency="USD",
tick_throttle_ms=100,
queue_capacity=1024,
experimental=False,
)
Bases: DataFeedProtocol
Experimental real-time market data feed from Interactive Brokers.
Subscribes to tick-by-tick market data for specified symbols. Emits validated trade and quote events.
Data Format
MarketEvent trade and quote snapshots with UTC provider or receipt time.
Note
- IB must be connected before creating feed
- Requires market data subscription for symbols
- Throttles rapid ticks to avoid overwhelming strategy
Example
ib = IB() await ib.connectAsync('127.0.0.1', 7497, clientId=1)
feed = IBDataFeed(ib, symbols=['SPY', 'QQQ', 'IWM'], experimental=True) await feed.start()
Use directly or wrap with BarAggregator¶
aggregator = BarAggregator(feed, bar_size_minutes=1)
Initialize IB data feed.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
ib
|
IB
|
Connected IB instance |
required |
symbols
|
list[str]
|
List of symbols to subscribe to |
required |
exchange
|
str
|
IB exchange (default: SMART routing) |
'SMART'
|
currency
|
str
|
Currency (default: USD) |
'USD'
|
tick_throttle_ms
|
int
|
Minimum milliseconds between tick emissions (prevents overwhelming strategy with rapid ticks) |
100
|
queue_capacity
|
int
|
Maximum pending events before a fail-closed overflow. |
1024
|
experimental
|
bool
|
Must be true to acknowledge the unsupported feed contract. |
False
|
Source code in src/ml4t/live/feeds/ib_feed.py
stats
property
¶
Get feed statistics.
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
Dict with keys: |
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
start
async
¶
Subscribe to market data for all symbols.
Creates contracts and subscribes to real-time tick data.
Raises:
| Type | Description |
|---|---|
RuntimeError
|
If IB not connected |
Source code in src/ml4t/live/feeds/ib_feed.py
stop
¶
Unsubscribe from market data.
Cancels all market data subscriptions and stops feed.
Source code in src/ml4t/live/feeds/ib_feed.py
__aiter__
async
¶
Async iterator yielding market data.
Yields:
| Type | Description |
|---|---|
AsyncIterator[MarketEvent]
|
Validated trade and quote events. |
Stops when
- stop() is called (None sentinel)
- Feed is not running
Source code in src/ml4t/live/feeds/ib_feed.py
AlpacaDataFeed
¶
AlpacaDataFeed(
api_key,
secret_key,
symbols,
*,
data_type="bars",
feed="iex",
queue_capacity=1024,
experimental=False,
)
Bases: DataFeedProtocol
Experimental real-time market data feed from Alpaca Markets.
Subscribes to real-time data for specified symbols. Supports both stocks and crypto.
Data Types
bars: OHLCV minute bars (default, recommended for strategies) quotes: Bid/ask quotes (for spread-sensitive strategies) trades: Individual trades (highest frequency)
Data Feeds
iex: Free tier (limited data) sip: Premium (full market data, requires subscription)
Data Format
Validated MarketEvent bars, quotes, or trades with UTC event and receipt times.
Example
Stocks only¶
feed = AlpacaDataFeed( api_key='PKXXXXXXXX', secret_key='XXXXXXXXXX', symbols=['AAPL', 'MSFT'], experimental=True, )
Mixed stocks and crypto¶
feed = AlpacaDataFeed( api_key='PKXXXXXXXX', secret_key='XXXXXXXXXX', symbols=['AAPL', 'BTC/USD', 'ETH/USD'], experimental=True, )
await feed.start()
async for event in feed: consume(event)
Initialize Alpaca data feed.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
api_key
|
str
|
Alpaca API key |
required |
secret_key
|
str
|
Alpaca secret key |
required |
symbols
|
list[str]
|
List of symbols (e.g., ['AAPL', 'BTC/USD']) |
required |
data_type
|
str
|
Type of data - 'bars' (default), 'quotes', or 'trades' |
'bars'
|
feed
|
str
|
Data feed type - 'iex' (free) or 'sip' (premium) |
'iex'
|
queue_capacity
|
int
|
Maximum pending events before a fail-closed overflow. |
1024
|
experimental
|
bool
|
Must be true to acknowledge the unsupported feed contract. |
False
|
Source code in src/ml4t/live/feeds/alpaca_feed.py
stats
property
¶
Get feed statistics.
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
Dict with keys: |
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
dict[str, Any]
|
|
start
async
¶
Subscribe to market data for all symbols.
Creates streams and subscribes to real-time data.
Source code in src/ml4t/live/feeds/alpaca_feed.py
stop
¶
Stop data feed.
Closes all streams and signals consumer to exit.
Source code in src/ml4t/live/feeds/alpaca_feed.py
__aiter__
async
¶
Async iterator yielding market data.
Yields:
| Type | Description |
|---|---|
AsyncIterator[MarketEvent]
|
Validated bar, quote, or trade events. |
Stops when
- stop() is called (None sentinel)
- Feed is not running
Source code in src/ml4t/live/feeds/alpaca_feed.py
__anext__
async
¶
Get next data item.
Returns:
| Type | Description |
|---|---|
MarketEvent
|
A validated bar, quote, or trade event. |
Raises:
| Type | Description |
|---|---|
StopAsyncIteration
|
When feed is stopped |
Source code in src/ml4t/live/feeds/alpaca_feed.py
OKXFundingFeed
¶
Bases: DataFeedProtocol
OKX funding rate feed with OHLCV bars.
Combines price data with funding rate information for ML strategies that trade crypto perpetual futures based on funding rate signals.
Data Flow
- Poll /market/candles for latest OHLCV bar
- Poll /public/funding-rate for current funding rate
- Emit each causal record as its own event
Symbol Format
OKX perpetual swaps use format: BTC-USDT-SWAP, ETH-USDT-SWAP
Initialize OKX funding rate feed.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
symbols
|
list[str]
|
List of perpetual swap symbols (e.g., ['BTC-USDT-SWAP']) |
required |
timeframe
|
str
|
OHLCV bar timeframe ('1m', '1H', '4H', '1D') |
'1H'
|
poll_interval_seconds
|
float
|
How often to poll for new data |
60.0
|
queue_capacity
|
int
|
Maximum pending events before a fail-closed overflow. |
256
|
Source code in src/ml4t/live/feeds/okx_feed.py
start
async
¶
Start the OKX data feed.
Begins polling for OHLCV and funding rate data.
Source code in src/ml4t/live/feeds/okx_feed.py
stop
¶
Stop the data feed.
Source code in src/ml4t/live/feeds/okx_feed.py
close
async
¶
Close HTTP client.
Source code in src/ml4t/live/feeds/okx_feed.py
__aiter__
¶
__anext__
async
¶
Get next bar with funding data.
Returns:
| Type | Description |
|---|---|
MarketEvent
|
A validated bar or funding event. |
Raises:
| Type | Description |
|---|---|
StopAsyncIteration
|
When feed stops |
Source code in src/ml4t/live/feeds/okx_feed.py
BarAggregator
¶
BarAggregator(
source_feed,
bar_size_minutes=1,
assets=None,
flush_timeout_seconds=2.0,
queue_capacity=256,
)
Aggregates raw ticks or 5-second bars into minute bars.
Addresses aggregation and finalization requirements: 1. "If IBDataFeed pushes a tick to Strategy.on_data, the strategy might trigger 60x more often than intended." - Buffer incoming data. 2. "The 15:59 bar is never emitted because no 16:00 tick arrives." - Background flush checker emits bars on timeout.
The aggregator buffers incoming data and emits when: - A bar boundary is crossed (new tick arrives in next minute) - OR timeout expires (2s past bar end with no new data)
Example
raw_feed = IBTickFeed(ib, assets=['AAPL']) aggregated_feed = BarAggregator(raw_feed, bar_size_minutes=1)
async for event in aggregated_feed: consume(event)
Initialize BarAggregator.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
source_feed
|
DataFeedProtocol
|
Raw tick or sub-minute bar feed |
required |
bar_size_minutes
|
int
|
Output bar size in minutes (default: 1) |
1
|
assets
|
list[str] | None
|
List of assets to track (default: all from source) |
None
|
flush_timeout_seconds
|
float
|
Seconds after bar end before forcing emit (default: 2.0) |
2.0
|
queue_capacity
|
int
|
Maximum pending bars before a fail-closed overflow. |
256
|
Source code in src/ml4t/live/feeds/aggregator.py
start
async
¶
Start aggregation.
Source code in src/ml4t/live/feeds/aggregator.py
stop
¶
__aiter__
¶
__anext__
async
¶
Return the next validated bar or finish after the queue drains.
Source code in src/ml4t/live/feeds/aggregator.py
BarBuffer
dataclass
¶
Accumulates ticks into OHLCV bar.
Attributes:
| Name | Type | Description |
|---|---|---|
open |
float | None
|
Opening price (first tick) |
high |
float
|
Highest price seen |
low |
float
|
Lowest price seen |
close |
float
|
Most recent price |
volume |
float
|
Total volume accumulated |
update
¶
Add a tick to the bar.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
price
|
float
|
Trade price |
required |
size
|
float
|
Trade size (defaults to 0 for quote ticks) |
0
|
Source code in src/ml4t/live/feeds/aggregator.py
update_bar
¶
Merge one validated OHLCV payload without discarding its range.
Source code in src/ml4t/live/feeds/aggregator.py
to_dict
¶
Convert to OHLCV dict.
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
Dictionary with keys: open, high, low, close, volume |
dict[str, Any]
|
If no ticks received, uses close price as fallback for OHLC |
Source code in src/ml4t/live/feeds/aggregator.py
FeedContractError
¶
Bases: ValueError
Raised when provider data cannot cross the portable feed boundary.
FeedContinuityError
¶
Bases: FeedContractError
Raised when event history cannot support another causal decision.
Source code in src/ml4t/live/feeds/events.py
FeedOverflowError
¶
Bases: RuntimeError
Raised when continuing after a full queue would hide lost market data.
Source code in src/ml4t/live/feeds/queue.py
FeedQueueSnapshot
dataclass
¶
FeedQueueSnapshot(
capacity,
occupancy,
high_watermark,
overflow_count,
oldest_event_lag_seconds,
failed,
finished,
)
Observable state for one supported feed queue.
Experimental Feed Opt-In¶
These imports are available for deliberate evaluation. They are not part of the stable support
contract and require experimental=True.
from ml4t.live import (
AlpacaDataFeed,
CryptoFeed,
DataBentoFeed,
ExperimentalFeedError,
ExperimentalFeedWarning,
IBDataFeed,
)
ExperimentalFeedError
¶
Bases: RuntimeError
Raised when an experimental feed is constructed without deliberate opt-in.
ExperimentalFeedWarning
¶
Bases: UserWarning
Reports guarantees that do not apply to an experimental feed.
DataBentoFeed
¶
Bases: DataFeedProtocol
Experimental market data feed from DataBento.
Supports selected historical replay and live record shapes after explicit opt-in.
Historical Mode
- Reads from .dbn files (DataBento native format)
- Replays at historical speed or accelerated
- Intended for custom evaluation under the experimental limitations
Real-time Mode
- Streams live market data
- Supports multiple datasets (GLBX, XNAS, OPRA, etc.)
- No stable latency or throughput guarantee
Data Format
Experimental typed MarketEvent bars, trades, and quotes with UTC timestamps.
Example Historical
feed = DataBentoFeed.from_file( 'ES_202401.dbn', symbols=['ES.FUT'], replay_speed=10.0, # 10x speed experimental=True, )
Example Real-time
feed = DataBentoFeed.from_live( api_key=os.getenv('DATABENTO_API_KEY'), dataset='GLBX.MDP3', schema='ohlcv-1s', symbols=['ES.c.0', 'NQ.c.0'], experimental=True, )
Initialize DataBento feed.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
client
|
_HistoricalRecordSource | _LiveRecordSource
|
DataBento client (Historical or Live) |
required |
symbols
|
list[str]
|
List of symbols to subscribe to |
required |
mode
|
str
|
'historical' or 'live' |
'historical'
|
replay_speed
|
float
|
Playback speed multiplier (historical only) 1.0 = real-time, 10.0 = 10x speed |
1.0
|
experimental
|
bool
|
Must be true to acknowledge the unsupported feed contract. |
False
|
Source code in src/ml4t/live/feeds/databento_feed.py
from_file
classmethod
¶
Create feed from historical .dbn file.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
file_path
|
str | Path
|
Path to .dbn file |
required |
symbols
|
list[str]
|
Non-empty symbol selection |
required |
replay_speed
|
float
|
Playback speed (1.0 = real-time) |
1.0
|
experimental
|
bool
|
Must be true to acknowledge the unsupported feed contract. |
False
|
Returns:
| Type | Description |
|---|---|
DataBentoFeed
|
DataBentoFeed configured for historical replay |
Source code in src/ml4t/live/feeds/databento_feed.py
from_live
classmethod
¶
Create feed for real-time streaming.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
api_key
|
str
|
DataBento API key |
required |
dataset
|
str
|
Dataset code (e.g., 'GLBX.MDP3', 'XNAS.ITCH') |
required |
schema
|
str
|
Data schema (e.g., 'ohlcv-1s', 'mbp-10', 'trades') |
required |
symbols
|
list[str]
|
Symbols to subscribe to |
required |
experimental
|
bool
|
Must be true to acknowledge the unsupported feed contract. |
False
|
Returns:
| Type | Description |
|---|---|
DataBentoFeed
|
DataBentoFeed configured for live streaming |
Source code in src/ml4t/live/feeds/databento_feed.py
start
async
¶
Start data feed.
Historical mode: Begins replay task Live mode: Starts streaming subscription
Source code in src/ml4t/live/feeds/databento_feed.py
stop
¶
Stop data feed.
Source code in src/ml4t/live/feeds/databento_feed.py
__aiter__
async
¶
Async iterator yielding market data.
Yields:
| Type | Description |
|---|---|
AsyncIterator[MarketEvent]
|
Typed experimental market event. |
Source code in src/ml4t/live/feeds/databento_feed.py
CryptoFeed
¶
CryptoFeed(
exchange,
symbols,
*,
timeframe="1m",
stream_trades=False,
stream_ohlcv=True,
api_key=None,
api_secret=None,
api_passphrase=None,
experimental=False,
)
Bases: DataFeedProtocol
Experimental cryptocurrency market data feed via asynchronous CCXT.
Supports async REST polling and uses CCXT Pro websocket methods when available. No exchange, overload, reconnect, or performance guarantee is included in the stable support contract.
Data Format
Experimental typed MarketEvent bars and trades with UTC timestamps.
Exchange Symbols
- Binance: 'BTC/USDT', 'ETH/USDT'
- Coinbase: 'BTC-USD', 'ETH-USD'
- Kraken: 'BTC/USD', 'ETH/USD'
Timeframes
'1m', '5m', '15m', '1h', '4h', '1d'
Example WebSocket (Real-time): feed = CryptoFeed( exchange='binance', symbols=['BTC/USDT', 'ETH/USDT'], stream_trades=True, # Stream trades (fastest) experimental=True, )
Example Authenticated
feed = CryptoFeed( exchange='binance', symbols=['BTC/USDT'], api_key='your-key', api_secret='your-secret', experimental=True, )
Initialize crypto feed.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
exchange
|
str
|
Exchange ID (e.g., 'binance', 'coinbasepro', 'kraken') |
required |
symbols
|
list[str]
|
Trading pairs (e.g., ['BTC/USDT', 'ETH/USDT']) |
required |
timeframe
|
str
|
OHLCV timeframe ('1m', '5m', '1h', etc.) |
'1m'
|
stream_trades
|
bool
|
Stream trade ticks (faster updates) |
False
|
stream_ohlcv
|
bool
|
Stream OHLCV candles |
True
|
api_key
|
str | None
|
API key (for authenticated endpoints) |
None
|
api_secret
|
str | None
|
API secret |
None
|
api_passphrase
|
str | None
|
API passphrase (Coinbase only) |
None
|
experimental
|
bool
|
Must be true to acknowledge the unsupported feed contract. |
False
|
Source code in src/ml4t/live/feeds/crypto_feed.py
123 124 125 126 127 128 129 130 131 132 133 134 135 136 137 138 139 140 141 142 143 144 145 146 147 148 149 150 151 152 153 154 155 156 157 158 159 160 161 162 163 164 165 166 167 168 169 170 171 172 173 174 175 176 177 178 179 180 181 182 183 184 185 186 187 188 189 190 191 192 193 194 195 196 197 198 199 | |
start
async
¶
Start streaming market data.
Initiates WebSocket subscriptions for all symbols.
Source code in src/ml4t/live/feeds/crypto_feed.py
stop
¶
Stop streaming; use close to release the exchange connection.
Source code in src/ml4t/live/feeds/crypto_feed.py
__aiter__
async
¶
Async iterator yielding market data.
Yields:
| Type | Description |
|---|---|
AsyncIterator[MarketEvent]
|
Typed experimental market event. |
Source code in src/ml4t/live/feeds/crypto_feed.py
Protocols¶
BrokerProtocol
¶
Bases: Protocol
Synchronous broker protocol for Strategy.on_data().
This is the interface strategies interact with. It must be synchronous because Strategy.on_data() is synchronous (matches backtest behavior).
ThreadSafeBrokerWrapper implements this protocol by wrapping an AsyncBrokerProtocol and using run_coroutine_threadsafe().
Example
class MyStrategy(Strategy): def on_data(self, timestamp, data, context, broker: BrokerProtocol): # broker is sync - no async/await needed pos = broker.get_position("AAPL") if pos is None: broker.submit_order("AAPL", 100)
positions
property
¶
Get all current positions.
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbol to Position |
pending_orders
property
¶
Get all pending orders.
Returns:
| Type | Description |
|---|---|
list[Order]
|
List of pending Order objects |
is_connected
property
¶
Check if broker is connected.
Returns:
| Type | Description |
|---|---|
bool
|
True if connected and ready to trade |
get_position
¶
Get position for specific asset.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol (e.g., "AAPL") |
required |
Returns:
| Type | Description |
|---|---|
Position | None
|
Position object if holding position, None otherwise |
get_positions
¶
get_account_value
¶
Get total account value (cash + positions).
Returns:
| Type | Description |
|---|---|
float
|
Total account value in base currency |
get_cash
¶
get_pending_orders
¶
submit_order
¶
submit_order(
asset,
quantity,
side=None,
order_type=MARKET,
limit_price=None,
stop_price=None,
**kwargs,
)
Submit order for execution.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol (e.g., "AAPL") |
required |
quantity
|
float
|
Signed shares/contracts when side is omitted; positive unsigned shares/contracts when side is provided |
required |
side
|
OrderSide | None
|
Order side (BUY/SELL), inferred from signed quantity if omitted |
None
|
order_type
|
OrderType
|
Type of order (MARKET, LIMIT, STOP, etc.) |
MARKET
|
limit_price
|
float | None
|
Limit price for LIMIT/STOP_LIMIT orders |
None
|
stop_price
|
float | None
|
Stop price for STOP/STOP_LIMIT orders |
None
|
**kwargs
|
Any
|
Additional broker-specific parameters |
{}
|
Returns:
| Type | Description |
|---|---|
Order
|
Order object with order_id and initial status |
Raises:
| Type | Description |
|---|---|
ValueError
|
If order parameters are invalid |
RuntimeError
|
If broker is not connected |
Source code in src/ml4t/live/protocols.py
cancel_order
¶
Cancel pending order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
order_id
|
str
|
ID of order to cancel |
required |
Returns:
| Type | Description |
|---|---|
bool
|
True if cancel request submitted, False if order not found |
replace_order
¶
Replace a pending order with updated parameters.
close_position
¶
Close entire position in asset.
Convenience method that submits a closing order.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
asset
|
str
|
Asset symbol to close |
required |
Returns:
| Type | Description |
|---|---|
Order | None
|
Order object if position exists, None if no position |
Source code in src/ml4t/live/protocols.py
register_target_intent
¶
Register an idempotent target for a causal opening phase.
register_position_rule_policy
¶
get_target_intents
¶
get_child_order_intents
¶
get_intent_reconciliations
¶
export_target_intent_state
¶
set_position_rules
¶
clear_position_rules
¶
update_position_context
¶
AsyncBrokerProtocol
¶
Bases: Protocol
Asynchronous broker protocol for broker implementations.
All broker implementations (IBBroker, AlpacaBroker, etc.) must implement this protocol. The async design enables efficient I/O without blocking.
ThreadSafeBrokerWrapper wraps this protocol to provide BrokerProtocol for strategies.
Example
class IBBroker: async def connect(self): await self.ib.connectAsync(...)
async def submit_order_async(self, asset, quantity, ...):
# Async I/O to IB
trade = await self.ib.placeOrderAsync(...)
return order
execution_capabilities
property
¶
Return the order capabilities implemented by the adapter.
connect
async
¶
Connect to broker and sync initial state.
This should: 1. Establish connection to broker API 2. Sync current positions 3. Sync pending orders 4. Register event callbacks
disconnect
async
¶
is_connected_async
async
¶
assert_paper_trading
¶
assert_live_trading
¶
get_positions_async
async
¶
get_pending_orders_async
async
¶
get_position_async
async
¶
get_account_value_async
async
¶
get_cash_async
async
¶
submit_order_async
async
¶
submit_order_async(
asset,
quantity,
side=None,
order_type=MARKET,
limit_price=None,
stop_price=None,
**kwargs,
)
Submit an order using the same quantity contract as submit_order.
Source code in src/ml4t/live/protocols.py
cancel_order_async
async
¶
replace_order_async
async
¶
Replace a pending order with updated parameters (async version).
Source code in src/ml4t/live/protocols.py
DataFeedProtocol
¶
Bases: Protocol
Protocol for real-time data feeds.
Stable-supported feeds provide an async iterator of validated MarketEvent objects.
The legacy tuple member of FeedItem is temporary compatibility for experimental feeds.
The feed handles:
- Subscribing to market data
- Aggregating ticks to bars (if needed)
- Emitting data on schedule (e.g., every minute)
Example
class IBDataFeed: async def start(self): await self._subscribe_to_ticks()
async def __aiter__(self):
return self
async def __anext__(self):
return await self._event_queue.get()
def stop(self):
self._running = False
start
async
¶
Start the data feed.
This should: 1. Subscribe to market data 2. Start internal aggregation/buffering 3. Begin emitting data
stop
¶
Stop the data feed gracefully.
Should be non-blocking. The feed should stop after the current iteration completes.
__aiter__
¶
Return async iterator.
Yields:
| Type | Description |
|---|---|
AsyncIterator[FeedItem]
|
A validated market event, or a temporary legacy tuple from an experimental feed. |
__anext__
async
¶
Get next market event.
Returns:
| Type | Description |
|---|---|
FeedItem
|
A validated market event, or a temporary legacy tuple from an experimental feed. |
Raises:
| Type | Description |
|---|---|
StopAsyncIteration
|
When feed ends |
Safety Types¶
RiskLimitError
¶
Bases: Exception
Raised when an order violates risk limits.
ReconciliationMismatchError
¶
Bases: RuntimeError
Raised when startup reconciliation is configured to fail closed.
RiskState
dataclass
¶
RiskState(
date,
daily_loss=0.0,
orders_placed=0,
high_water_mark=0.0,
session_start_equity=None,
persisted_positions=dict(),
persisted_pending_orders=list(),
portable_strategy_state=dict(),
replacement_gaps=dict(),
shadow_portfolio=dict(),
execution_mode=None,
kill_switch_activated=False,
kill_switch_reason="",
)
Persisted risk state - survives restarts.
This state is validated and written in a versioned integrity envelope after every accepted order and on shutdown.
Example
state = RiskState(date="2023-10-15", daily_loss=1500.0) state.kill_switch_activated = True state.kill_switch_reason = "Max daily loss exceeded"
Note
The kill switch state persists across restarts and must be reset explicitly.
from_dict
classmethod
¶
Create RiskState from dictionary (for JSON loading).
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
data
|
dict
|
Dictionary with RiskState fields |
required |
Returns:
| Type | Description |
|---|---|
RiskState
|
RiskState instance |
Source code in src/ml4t/live/safety.py
validate
¶
Validate persisted field types and numeric boundaries.
Source code in src/ml4t/live/safety.py
to_dict
¶
Convert to dictionary (for JSON saving).
Returns:
| Type | Description |
|---|---|
dict[str, Any]
|
Dictionary with all fields |
Source code in src/ml4t/live/safety.py
save_atomic
staticmethod
¶
Save state with atomic write (write to .tmp then os.replace).
This prevents corruption if process dies mid-write.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
state
|
RiskState
|
RiskState to save |
required |
filepath
|
str
|
Path to save to |
required |
Raises:
| Type | Description |
|---|---|
OSError
|
If write fails |
Source code in src/ml4t/live/safety.py
load
staticmethod
¶
Load state from file.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
filepath
|
str
|
Path to load from |
required |
Returns:
| Type | Description |
|---|---|
RiskState | None
|
RiskState if the file exists, otherwise None |
Raises:
| Type | Description |
|---|---|
PersistenceSafetyError
|
If the file is unsafe, corrupt, or incompatible |
Source code in src/ml4t/live/safety.py
VirtualPortfolio
¶
Manages internal accounting for Shadow Mode (Paper Trading).
Tracks prior virtual fills so repeated shadow orders observe current positions
Problem: In shadow mode, returning fake Order objects without updating position state causes strategies to keep buying forever because get_position() always returns None.
Solution: Track shadow positions locally. When shadow_mode=True: - submit_order() updates this virtual portfolio - positions/get_position() return from this portfolio - Strategy sees realistic position state
Handles: - New positions - Position increases (weighted avg cost basis) - Position decreases (partial close) - Position close (quantity = 0) - Position flip (long -> short or vice versa)
Example
portfolio = VirtualPortfolio(initial_cash=100_000.0)
Simulate buy order fill¶
order = Order( asset="AAPL", side=OrderSide.BUY, quantity=100, filled_price=150.0, filled_quantity=100, ... ) portfolio.process_fill(order)
Check position¶
pos = portfolio.positions.get("AAPL") assert pos.quantity == 100 assert pos.entry_price == 150.0
Simulate sell order fill (close)¶
sell_order = Order( asset="AAPL", side=OrderSide.SELL, quantity=100, filled_price=155.0, filled_quantity=100, ... ) portfolio.process_fill(sell_order) assert "AAPL" not in portfolio.positions
Initialize virtual portfolio.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
initial_cash
|
float
|
Starting cash balance (default: 100,000) |
100000.0
|
Source code in src/ml4t/live/safety.py
positions
property
¶
Get current positions (returns copy for safety).
Returns:
| Type | Description |
|---|---|
dict[str, Position]
|
Dictionary mapping asset symbol to Position |
account_value
property
¶
Get total account value (cash + position market value).
Returns:
| Type | Description |
|---|---|
float
|
Total account value |
process_fill
¶
Update state based on filled shadow order.
Handles: - Weighted average cost basis for position increases - Position flipping (long -> short or vice versa) - Partial and full closes
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
order
|
Order
|
Filled Order object (must have filled_quantity and filled_price) |
required |
Source code in src/ml4t/live/safety.py
625 626 627 628 629 630 631 632 633 634 635 636 637 638 639 640 641 642 643 644 645 646 647 648 649 650 651 652 653 654 655 656 657 658 659 660 661 662 663 664 665 666 667 668 669 670 671 672 673 674 675 676 677 678 679 680 681 682 683 684 685 686 687 688 689 690 691 692 693 694 695 696 697 698 699 700 701 702 703 704 | |
update_prices
¶
Update current prices for accurate account value.
Parameters:
| Name | Type | Description | Default |
|---|---|---|---|
prices
|
dict[str, float]
|
Dictionary mapping asset symbol to current price |
required |
Source code in src/ml4t/live/safety.py
to_state
¶
Return restart-safe shadow cash and position state.
Source code in src/ml4t/live/safety.py
restore_state
¶
Restore shadow state before runtime reconciliation.