Skip to content

Commit f9252c1

Browse files
committed
decouple ibkr runtime wiring
1 parent 7ca8f74 commit f9252c1

13 files changed

Lines changed: 1810 additions & 413 deletions

application/cycle_result.py

Lines changed: 25 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,25 @@
1+
"""Structured strategy cycle results for InteractiveBrokersPlatform."""
2+
3+
from __future__ import annotations
4+
5+
from dataclasses import dataclass, field
6+
from typing import Any
7+
8+
9+
@dataclass(frozen=True)
10+
class StrategyCycleResult:
11+
"""Structured output from one strategy cycle."""
12+
13+
result: str
14+
signal_metadata: dict[str, Any] = field(default_factory=dict)
15+
target_weights: dict[str, float] | None = None
16+
execution_summary: dict[str, Any] = field(default_factory=dict)
17+
reconciliation_record: dict[str, Any] = field(default_factory=dict)
18+
reconciliation_record_path: str | None = None
19+
20+
21+
def coerce_strategy_cycle_result(value: StrategyCycleResult | str) -> StrategyCycleResult:
22+
"""Keep request-handling tolerant of older string-only test doubles."""
23+
if isinstance(value, StrategyCycleResult):
24+
return value
25+
return StrategyCycleResult(result=str(value))

application/rebalance_service.py

Lines changed: 121 additions & 101 deletions
Original file line numberDiff line numberDiff line change
@@ -3,14 +3,20 @@
33
from __future__ import annotations
44

55
from collections.abc import Mapping
6+
from datetime import datetime, timezone
67
import json
78
import re
89

10+
from application.cycle_result import StrategyCycleResult
11+
from application.runtime_dependencies import IBKRRebalanceConfig, IBKRRebalanceRuntime
912
from application.reconciliation_service import (
1013
build_reconciliation_record,
1114
write_reconciliation_record,
1215
)
13-
from notifications.events import NotificationPublisher, RenderedNotification
16+
from notifications.events import NotificationPublisher
17+
from notifications import renderers as notification_renderers
18+
from quant_platform_kit.common.models import PortfolioSnapshot, Position
19+
from quant_platform_kit.common.port_adapters import CallableNotificationPort, CallablePortfolioPort
1420

1521

1622
_ZH_REASON_REPLACEMENTS = (
@@ -443,29 +449,86 @@ def _build_compact_message(
443449
return "\n".join(lines)
444450

445451

452+
def _legacy_portfolio_snapshot(ib, *, get_current_portfolio) -> PortfolioSnapshot:
453+
positions, account_values = get_current_portfolio(ib)
454+
snapshot_positions = tuple(
455+
Position(
456+
symbol=str(symbol).strip().upper(),
457+
quantity=float(details.get("quantity") or 0),
458+
market_value=float(details.get("quantity") or 0) * float(details.get("avg_cost") or 0.0),
459+
average_cost=float(details.get("avg_cost") or 0.0),
460+
)
461+
for symbol, details in dict(positions or {}).items()
462+
)
463+
return PortfolioSnapshot(
464+
as_of=datetime.now(timezone.utc),
465+
total_equity=float(account_values.get("equity") or 0.0),
466+
buying_power=float(account_values.get("buying_power") or 0.0),
467+
positions=snapshot_positions,
468+
)
469+
470+
471+
def _snapshot_to_portfolio_view(snapshot) -> tuple[dict[str, dict[str, float | int]], dict[str, float]]:
472+
positions = {}
473+
for position in getattr(snapshot, "positions", ()) or ():
474+
positions[str(position.symbol).strip().upper()] = {
475+
"quantity": int(position.quantity),
476+
"avg_cost": float(position.average_cost or 0.0),
477+
}
478+
account_values = {
479+
"equity": float(getattr(snapshot, "total_equity", 0.0) or 0.0),
480+
"buying_power": float(getattr(snapshot, "buying_power", 0.0) or 0.0),
481+
}
482+
return positions, account_values
483+
484+
446485
def run_strategy_core(
447486
*,
448-
connect_ib,
449-
get_current_portfolio,
450-
compute_signals,
451-
execute_rebalance,
452-
send_tg_message,
453-
translator,
454-
separator,
487+
runtime: IBKRRebalanceRuntime | None = None,
488+
config: IBKRRebalanceConfig | None = None,
489+
connect_ib=None,
490+
get_current_portfolio=None,
491+
compute_signals=None,
492+
execute_rebalance=None,
493+
send_tg_message=None,
494+
translator=None,
495+
separator=None,
455496
strategy_display_name=None,
456497
reconciliation_output_path=None,
457-
result_hook=None,
458498
):
499+
if runtime is None:
500+
if not all((connect_ib, get_current_portfolio, compute_signals, execute_rebalance, send_tg_message)):
501+
raise ValueError("Legacy IBKR rebalance call requires connect_ib/get_current_portfolio/compute_signals/execute_rebalance/send_tg_message")
502+
runtime = IBKRRebalanceRuntime(
503+
connect_ib=connect_ib,
504+
portfolio_port_factory=lambda ib: CallablePortfolioPort(
505+
lambda: _legacy_portfolio_snapshot(ib, get_current_portfolio=get_current_portfolio)
506+
),
507+
compute_signals=compute_signals,
508+
execute_rebalance=execute_rebalance,
509+
notifications=CallableNotificationPort(send_tg_message),
510+
)
511+
if config is None:
512+
if translator is None or separator is None:
513+
raise ValueError("IBKR rebalance config requires translator and separator")
514+
config = IBKRRebalanceConfig(
515+
translator=translator,
516+
separator=separator,
517+
strategy_display_name=strategy_display_name,
518+
reconciliation_output_path=reconciliation_output_path,
519+
)
520+
459521
notification_publisher = NotificationPublisher(
460522
log_message=lambda message: print(message, flush=True),
461-
send_message=send_tg_message,
523+
send_message=runtime.notifications.send_text,
462524
)
463525
ib = None
464526
try:
465-
ib = connect_ib()
466-
positions, account_values = get_current_portfolio(ib)
527+
ib = runtime.connect_ib()
528+
snapshot = runtime.portfolio_port_factory(ib).get_portfolio_snapshot()
529+
positions, account_values = _snapshot_to_portfolio_view(snapshot)
467530
current_holdings = set(positions.keys())
468-
signal_result = compute_signals(ib, current_holdings)
531+
signal_result = runtime.compute_signals(ib, current_holdings)
469532
if len(signal_result) == 5:
470533
target_weights, signal_desc, _is_emergency, status_desc, signal_metadata = signal_result
471534
else:
@@ -474,17 +537,17 @@ def run_strategy_core(
474537
allocation = _resolve_weight_allocation(signal_metadata, required=target_weights is not None)
475538
resolved_target_weights = dict(allocation.get("targets") or {}) if target_weights is not None else None
476539

477-
dashboard = build_dashboard(
540+
dashboard = notification_renderers.build_dashboard(
478541
positions,
479542
account_values,
480543
signal_desc,
481544
status_desc,
482545
strategy_profile=signal_metadata.get("strategy_profile"),
483-
strategy_display_name=strategy_display_name,
546+
strategy_display_name=config.strategy_display_name,
484547
target_weights=resolved_target_weights,
485548
signal_metadata=signal_metadata,
486-
translator=translator,
487-
separator=separator,
549+
translator=config.translator,
550+
separator=config.separator,
488551
status_icon=signal_metadata.get("status_icon", "🐤"),
489552
)
490553
strategy_dashboard = _strategy_dashboard_text(signal_metadata)
@@ -493,13 +556,13 @@ def run_strategy_core(
493556
decision = signal_metadata.get("snapshot_guard_decision")
494557
no_op_reason = signal_metadata.get("no_op_reason")
495558
fail_reason = signal_metadata.get("fail_reason")
496-
no_op_text = translator("no_trades")
559+
no_op_text = config.translator("no_trades")
497560
if decision:
498-
no_op_text = f"{no_op_text} | {_localize_notification_text(f'decision={decision}', translator=translator)}"
561+
no_op_text = f"{no_op_text} | {_localize_notification_text(f'decision={decision}', translator=config.translator)}"
499562
if no_op_reason:
500-
no_op_text = f"{no_op_text} | {_localize_notification_text(f'reason={no_op_reason}', translator=translator)}"
563+
no_op_text = f"{no_op_text} | {_localize_notification_text(f'reason={no_op_reason}', translator=config.translator)}"
501564
if fail_reason:
502-
no_op_text = f"{no_op_text} | {_localize_notification_text(f'fail_reason={fail_reason}', translator=translator)}"
565+
no_op_text = f"{no_op_text} | {_localize_notification_text(f'fail_reason={fail_reason}', translator=config.translator)}"
503566
no_op_text = "\n".join(_split_labeled_text(no_op_text))
504567
record = build_reconciliation_record(
505568
strategy_profile=signal_metadata.get("strategy_profile"),
@@ -511,44 +574,35 @@ def run_strategy_core(
511574
execution_summary=None,
512575
no_op_reason=no_op_reason or fail_reason or decision,
513576
)
514-
record_path = write_reconciliation_record(record, output_path=reconciliation_output_path)
577+
record_path = write_reconciliation_record(record, output_path=config.reconciliation_output_path)
515578
print(
516579
"reconciliation_record "
517580
+ json.dumps({"path": str(record_path), "status": record.get("execution_status"), "no_op_reason": record.get("no_op_reason")}, ensure_ascii=False),
518581
flush=True,
519582
)
520-
detailed_message = f"{translator('heartbeat_title')}\n{dashboard}\n{separator}\n{no_op_text}"
521-
compact_message = _build_compact_message(
522-
title=translator("heartbeat_title"),
523-
strategy_display_name=strategy_display_name,
524-
signal_desc=signal_desc,
525-
status_desc=status_desc,
526-
status_icon=signal_metadata.get("status_icon", "🐤"),
527-
translator=translator,
528-
separator=separator,
529-
body_lines=[no_op_text],
530-
dashboard_text=strategy_dashboard,
531-
)
532583
notification_publisher.publish(
533-
RenderedNotification(
534-
detailed_text=detailed_message,
535-
compact_text=compact_message,
584+
notification_renderers.render_heartbeat_notification(
585+
dashboard=dashboard,
586+
strategy_dashboard=strategy_dashboard,
587+
no_op_text=no_op_text,
588+
signal_desc=signal_desc,
589+
status_desc=status_desc,
590+
status_icon=signal_metadata.get("status_icon", "🐤"),
591+
translator=config.translator,
592+
separator=config.separator,
593+
strategy_display_name=config.strategy_display_name,
536594
)
537595
)
538-
if callable(result_hook):
539-
result_hook(
540-
{
541-
"result": "OK - heartbeat",
542-
"signal_metadata": dict(signal_metadata or {}),
543-
"target_weights": None,
544-
"execution_summary": None,
545-
"reconciliation_record": dict(record),
546-
"reconciliation_record_path": str(record_path),
547-
}
548-
)
549-
return "OK - heartbeat"
596+
return StrategyCycleResult(
597+
result="OK - heartbeat",
598+
signal_metadata=dict(signal_metadata or {}),
599+
target_weights=None,
600+
execution_summary={},
601+
reconciliation_record=dict(record),
602+
reconciliation_record_path=str(record_path),
603+
)
550604

551-
execution_result = execute_rebalance(
605+
execution_result = runtime.execute_rebalance(
552606
ib,
553607
resolved_target_weights,
554608
positions,
@@ -570,7 +624,7 @@ def run_strategy_core(
570624
target_weights=resolved_target_weights,
571625
execution_summary=execution_summary,
572626
)
573-
record_path = write_reconciliation_record(record, output_path=reconciliation_output_path)
627+
record_path = write_reconciliation_record(record, output_path=config.reconciliation_output_path)
574628
print(
575629
"reconciliation_record "
576630
+ json.dumps(
@@ -585,62 +639,28 @@ def run_strategy_core(
585639
),
586640
flush=True,
587641
)
588-
if trade_logs:
589-
notification_trade_lines = _build_notification_trade_lines(
590-
trade_logs,
642+
notification_publisher.publish(
643+
notification_renderers.render_trade_notification(
644+
dashboard=dashboard,
645+
strategy_dashboard=strategy_dashboard,
646+
trade_logs=trade_logs,
591647
execution_summary=execution_summary,
592-
translator=translator,
593-
)
594-
trade_lines = "\n".join(notification_trade_lines)
595-
detailed_message = (
596-
f"{translator('rebalance_title')}\n"
597-
f"{dashboard}\n"
598-
f"{separator}\n"
599-
f"{trade_lines}"
600-
)
601-
compact_message = _build_compact_message(
602-
title=translator("rebalance_title"),
603-
strategy_display_name=strategy_display_name,
604-
signal_desc=signal_desc,
605-
status_desc=status_desc,
606-
status_icon=signal_metadata.get("status_icon", "🐤"),
607-
translator=translator,
608-
separator=separator,
609-
body_lines=notification_trade_lines,
610-
dashboard_text=strategy_dashboard,
611-
)
612-
else:
613-
detailed_message = f"{translator('heartbeat_title')}\n{dashboard}\n{separator}\n{translator('no_trades')}"
614-
compact_message = _build_compact_message(
615-
title=translator("heartbeat_title"),
616-
strategy_display_name=strategy_display_name,
617648
signal_desc=signal_desc,
618649
status_desc=status_desc,
619650
status_icon=signal_metadata.get("status_icon", "🐤"),
620-
translator=translator,
621-
separator=separator,
622-
body_lines=[translator("no_trades")],
623-
dashboard_text=strategy_dashboard,
624-
)
625-
626-
notification_publisher.publish(
627-
RenderedNotification(
628-
detailed_text=detailed_message,
629-
compact_text=compact_message,
651+
translator=config.translator,
652+
separator=config.separator,
653+
strategy_display_name=config.strategy_display_name,
630654
)
631655
)
632-
if callable(result_hook):
633-
result_hook(
634-
{
635-
"result": "OK - executed",
636-
"signal_metadata": dict(signal_metadata or {}),
637-
"target_weights": dict(resolved_target_weights or {}),
638-
"execution_summary": dict(execution_summary or {}),
639-
"reconciliation_record": dict(record),
640-
"reconciliation_record_path": str(record_path),
641-
}
642-
)
643-
return "OK - executed"
656+
return StrategyCycleResult(
657+
result="OK - executed",
658+
signal_metadata=dict(signal_metadata or {}),
659+
target_weights=dict(resolved_target_weights or {}),
660+
execution_summary=dict(execution_summary or {}),
661+
reconciliation_record=dict(record),
662+
reconciliation_record_path=str(record_path),
663+
)
644664
finally:
645665
if ib is not None and ib.isConnected():
646666
ib.disconnect()

0 commit comments

Comments
 (0)