Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
25 changes: 25 additions & 0 deletions application/cycle_result.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,25 @@
"""Structured strategy cycle results for InteractiveBrokersPlatform."""

from __future__ import annotations

from dataclasses import dataclass, field
from typing import Any


@dataclass(frozen=True)
class StrategyCycleResult:
"""Structured output from one strategy cycle."""

result: str
signal_metadata: dict[str, Any] = field(default_factory=dict)
target_weights: dict[str, float] | None = None
execution_summary: dict[str, Any] = field(default_factory=dict)
reconciliation_record: dict[str, Any] = field(default_factory=dict)
reconciliation_record_path: str | None = None


def coerce_strategy_cycle_result(value: StrategyCycleResult | str) -> StrategyCycleResult:
"""Keep request-handling tolerant of older string-only test doubles."""
if isinstance(value, StrategyCycleResult):
return value
return StrategyCycleResult(result=str(value))
222 changes: 121 additions & 101 deletions application/rebalance_service.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,14 +3,20 @@
from __future__ import annotations

from collections.abc import Mapping
from datetime import datetime, timezone
import json
import re

from application.cycle_result import StrategyCycleResult
from application.runtime_dependencies import IBKRRebalanceConfig, IBKRRebalanceRuntime
from application.reconciliation_service import (
build_reconciliation_record,
write_reconciliation_record,
)
from notifications.events import NotificationPublisher, RenderedNotification
from notifications.events import NotificationPublisher
from notifications import renderers as notification_renderers
from quant_platform_kit.common.models import PortfolioSnapshot, Position
from quant_platform_kit.common.port_adapters import CallableNotificationPort, CallablePortfolioPort


_ZH_REASON_REPLACEMENTS = (
Expand Down Expand Up @@ -443,29 +449,86 @@ def _build_compact_message(
return "\n".join(lines)


def _legacy_portfolio_snapshot(ib, *, get_current_portfolio) -> PortfolioSnapshot:
positions, account_values = get_current_portfolio(ib)
snapshot_positions = tuple(
Position(
symbol=str(symbol).strip().upper(),
quantity=float(details.get("quantity") or 0),
market_value=float(details.get("quantity") or 0) * float(details.get("avg_cost") or 0.0),
average_cost=float(details.get("avg_cost") or 0.0),
)
for symbol, details in dict(positions or {}).items()
)
return PortfolioSnapshot(
as_of=datetime.now(timezone.utc),
total_equity=float(account_values.get("equity") or 0.0),
buying_power=float(account_values.get("buying_power") or 0.0),
positions=snapshot_positions,
)


def _snapshot_to_portfolio_view(snapshot) -> tuple[dict[str, dict[str, float | int]], dict[str, float]]:
positions = {}
for position in getattr(snapshot, "positions", ()) or ():
positions[str(position.symbol).strip().upper()] = {
"quantity": int(position.quantity),
"avg_cost": float(position.average_cost or 0.0),
}
account_values = {
"equity": float(getattr(snapshot, "total_equity", 0.0) or 0.0),
"buying_power": float(getattr(snapshot, "buying_power", 0.0) or 0.0),
}
return positions, account_values


def run_strategy_core(
*,
connect_ib,
get_current_portfolio,
compute_signals,
execute_rebalance,
send_tg_message,
translator,
separator,
runtime: IBKRRebalanceRuntime | None = None,
config: IBKRRebalanceConfig | None = None,
connect_ib=None,
get_current_portfolio=None,
compute_signals=None,
execute_rebalance=None,
send_tg_message=None,
translator=None,
separator=None,
strategy_display_name=None,
reconciliation_output_path=None,
result_hook=None,
):
if runtime is None:
if not all((connect_ib, get_current_portfolio, compute_signals, execute_rebalance, send_tg_message)):
raise ValueError("Legacy IBKR rebalance call requires connect_ib/get_current_portfolio/compute_signals/execute_rebalance/send_tg_message")
runtime = IBKRRebalanceRuntime(
connect_ib=connect_ib,
portfolio_port_factory=lambda ib: CallablePortfolioPort(
lambda: _legacy_portfolio_snapshot(ib, get_current_portfolio=get_current_portfolio)
),
compute_signals=compute_signals,
execute_rebalance=execute_rebalance,
notifications=CallableNotificationPort(send_tg_message),
)
if config is None:
if translator is None or separator is None:
raise ValueError("IBKR rebalance config requires translator and separator")
config = IBKRRebalanceConfig(
translator=translator,
separator=separator,
strategy_display_name=strategy_display_name,
reconciliation_output_path=reconciliation_output_path,
)

notification_publisher = NotificationPublisher(
log_message=lambda message: print(message, flush=True),
send_message=send_tg_message,
send_message=runtime.notifications.send_text,
)
ib = None
try:
ib = connect_ib()
positions, account_values = get_current_portfolio(ib)
ib = runtime.connect_ib()
snapshot = runtime.portfolio_port_factory(ib).get_portfolio_snapshot()
positions, account_values = _snapshot_to_portfolio_view(snapshot)
current_holdings = set(positions.keys())
signal_result = compute_signals(ib, current_holdings)
signal_result = runtime.compute_signals(ib, current_holdings)
if len(signal_result) == 5:
target_weights, signal_desc, _is_emergency, status_desc, signal_metadata = signal_result
else:
Expand All @@ -474,17 +537,17 @@ def run_strategy_core(
allocation = _resolve_weight_allocation(signal_metadata, required=target_weights is not None)
resolved_target_weights = dict(allocation.get("targets") or {}) if target_weights is not None else None

dashboard = build_dashboard(
dashboard = notification_renderers.build_dashboard(
positions,
account_values,
signal_desc,
status_desc,
strategy_profile=signal_metadata.get("strategy_profile"),
strategy_display_name=strategy_display_name,
strategy_display_name=config.strategy_display_name,
target_weights=resolved_target_weights,
signal_metadata=signal_metadata,
translator=translator,
separator=separator,
translator=config.translator,
separator=config.separator,
status_icon=signal_metadata.get("status_icon", "🐤"),
)
strategy_dashboard = _strategy_dashboard_text(signal_metadata)
Expand All @@ -493,13 +556,13 @@ def run_strategy_core(
decision = signal_metadata.get("snapshot_guard_decision")
no_op_reason = signal_metadata.get("no_op_reason")
fail_reason = signal_metadata.get("fail_reason")
no_op_text = translator("no_trades")
no_op_text = config.translator("no_trades")
if decision:
no_op_text = f"{no_op_text} | {_localize_notification_text(f'decision={decision}', translator=translator)}"
no_op_text = f"{no_op_text} | {_localize_notification_text(f'decision={decision}', translator=config.translator)}"
if no_op_reason:
no_op_text = f"{no_op_text} | {_localize_notification_text(f'reason={no_op_reason}', translator=translator)}"
no_op_text = f"{no_op_text} | {_localize_notification_text(f'reason={no_op_reason}', translator=config.translator)}"
if fail_reason:
no_op_text = f"{no_op_text} | {_localize_notification_text(f'fail_reason={fail_reason}', translator=translator)}"
no_op_text = f"{no_op_text} | {_localize_notification_text(f'fail_reason={fail_reason}', translator=config.translator)}"
no_op_text = "\n".join(_split_labeled_text(no_op_text))
record = build_reconciliation_record(
strategy_profile=signal_metadata.get("strategy_profile"),
Expand All @@ -511,44 +574,35 @@ def run_strategy_core(
execution_summary=None,
no_op_reason=no_op_reason or fail_reason or decision,
)
record_path = write_reconciliation_record(record, output_path=reconciliation_output_path)
record_path = write_reconciliation_record(record, output_path=config.reconciliation_output_path)
print(
"reconciliation_record "
+ json.dumps({"path": str(record_path), "status": record.get("execution_status"), "no_op_reason": record.get("no_op_reason")}, ensure_ascii=False),
flush=True,
)
detailed_message = f"{translator('heartbeat_title')}\n{dashboard}\n{separator}\n{no_op_text}"
compact_message = _build_compact_message(
title=translator("heartbeat_title"),
strategy_display_name=strategy_display_name,
signal_desc=signal_desc,
status_desc=status_desc,
status_icon=signal_metadata.get("status_icon", "🐤"),
translator=translator,
separator=separator,
body_lines=[no_op_text],
dashboard_text=strategy_dashboard,
)
notification_publisher.publish(
RenderedNotification(
detailed_text=detailed_message,
compact_text=compact_message,
notification_renderers.render_heartbeat_notification(
dashboard=dashboard,
strategy_dashboard=strategy_dashboard,
no_op_text=no_op_text,
signal_desc=signal_desc,
status_desc=status_desc,
status_icon=signal_metadata.get("status_icon", "🐤"),
translator=config.translator,
separator=config.separator,
strategy_display_name=config.strategy_display_name,
)
)
if callable(result_hook):
result_hook(
{
"result": "OK - heartbeat",
"signal_metadata": dict(signal_metadata or {}),
"target_weights": None,
"execution_summary": None,
"reconciliation_record": dict(record),
"reconciliation_record_path": str(record_path),
}
)
return "OK - heartbeat"
return StrategyCycleResult(
result="OK - heartbeat",
signal_metadata=dict(signal_metadata or {}),
target_weights=None,
execution_summary={},
reconciliation_record=dict(record),
reconciliation_record_path=str(record_path),
)

execution_result = execute_rebalance(
execution_result = runtime.execute_rebalance(
ib,
resolved_target_weights,
positions,
Expand All @@ -570,7 +624,7 @@ def run_strategy_core(
target_weights=resolved_target_weights,
execution_summary=execution_summary,
)
record_path = write_reconciliation_record(record, output_path=reconciliation_output_path)
record_path = write_reconciliation_record(record, output_path=config.reconciliation_output_path)
print(
"reconciliation_record "
+ json.dumps(
Expand All @@ -585,62 +639,28 @@ def run_strategy_core(
),
flush=True,
)
if trade_logs:
notification_trade_lines = _build_notification_trade_lines(
trade_logs,
notification_publisher.publish(
notification_renderers.render_trade_notification(
dashboard=dashboard,
strategy_dashboard=strategy_dashboard,
trade_logs=trade_logs,
execution_summary=execution_summary,
translator=translator,
)
trade_lines = "\n".join(notification_trade_lines)
detailed_message = (
f"{translator('rebalance_title')}\n"
f"{dashboard}\n"
f"{separator}\n"
f"{trade_lines}"
)
compact_message = _build_compact_message(
title=translator("rebalance_title"),
strategy_display_name=strategy_display_name,
signal_desc=signal_desc,
status_desc=status_desc,
status_icon=signal_metadata.get("status_icon", "🐤"),
translator=translator,
separator=separator,
body_lines=notification_trade_lines,
dashboard_text=strategy_dashboard,
)
else:
detailed_message = f"{translator('heartbeat_title')}\n{dashboard}\n{separator}\n{translator('no_trades')}"
compact_message = _build_compact_message(
title=translator("heartbeat_title"),
strategy_display_name=strategy_display_name,
signal_desc=signal_desc,
status_desc=status_desc,
status_icon=signal_metadata.get("status_icon", "🐤"),
translator=translator,
separator=separator,
body_lines=[translator("no_trades")],
dashboard_text=strategy_dashboard,
)

notification_publisher.publish(
RenderedNotification(
detailed_text=detailed_message,
compact_text=compact_message,
translator=config.translator,
separator=config.separator,
strategy_display_name=config.strategy_display_name,
)
)
if callable(result_hook):
result_hook(
{
"result": "OK - executed",
"signal_metadata": dict(signal_metadata or {}),
"target_weights": dict(resolved_target_weights or {}),
"execution_summary": dict(execution_summary or {}),
"reconciliation_record": dict(record),
"reconciliation_record_path": str(record_path),
}
)
return "OK - executed"
return StrategyCycleResult(
result="OK - executed",
signal_metadata=dict(signal_metadata or {}),
target_weights=dict(resolved_target_weights or {}),
execution_summary=dict(execution_summary or {}),
reconciliation_record=dict(record),
reconciliation_record_path=str(record_path),
)
finally:
if ib is not None and ib.isConnected():
ib.disconnect()
Loading