diff --git a/application/rebalance_service.py b/application/rebalance_service.py index 96ac85c..e188f9c 100644 --- a/application/rebalance_service.py +++ b/application/rebalance_service.py @@ -737,6 +737,9 @@ def fetch_replanned_state(): execution_marker_key = _build_execution_marker_key(config=config, execution=execution) execution_state_store = getattr(config, "execution_state_store", None) + duplicate_marker_found = False + duplicate_report_found = False + duplicate_evidence_read_failed = False direct_live_routing_blocked = ( _direct_live_routing_requires_durable_command( execution=execution, @@ -761,8 +764,10 @@ def fetch_replanned_state(): ) elif not execution_already_recorded and execution_marker_key and execution_state_store: try: - execution_already_recorded = bool(execution_state_store.has_marker(execution_marker_key)) + duplicate_marker_found = bool(execution_state_store.has_marker(execution_marker_key)) + execution_already_recorded = duplicate_marker_found except Exception as exc: + duplicate_evidence_read_failed = True detail = ( "execution_marker_read_failed" if live_command_claimed @@ -774,9 +779,9 @@ def fetch_replanned_state(): ) if live_command_claimed: execution_already_recorded = True - if not execution_already_recorded and hasattr(execution_state_store, "has_prior_execution_report"): + if hasattr(execution_state_store, "has_prior_execution_report"): try: - execution_already_recorded = bool( + duplicate_report_found = bool( execution_state_store.has_prior_execution_report( platform="longbridge", strategy_profile=getattr(config, "strategy_profile", "") or "unknown", @@ -786,7 +791,9 @@ def fetch_replanned_state(): dry_run_only=bool(getattr(config, "dry_run_only", False)), ) ) + execution_already_recorded = execution_already_recorded or duplicate_report_found except Exception as exc: + duplicate_evidence_read_failed = True detail = ( "execution_report_dedup_read_failed" if live_command_claimed @@ -799,7 +806,23 @@ def fetch_replanned_state(): if live_command_claimed: execution_already_recorded = True + duplicate_execution_confirmed = bool( + live_command_claimed + and not duplicate_evidence_read_failed + and (duplicate_marker_found or duplicate_report_found) + ) + if execution_already_recorded: + if duplicate_execution_confirmed and not account_identity_blocked and not live_command_blocked: + try: + config.execution_command_store.append_event( + live_command, + next_state=ExecutionCommandState.CANCELLED, + details={"reason": "duplicate_execution_confirmed", "orders_count": 0}, + expected_previous_state=ExecutionCommandState.CLAIMED, + ) + except Exception: + pass if account_identity_blocked: message = _account_identity_blocked_message( findings=tuple(account_identity_decision.findings), diff --git a/tests/test_rebalance_service.py b/tests/test_rebalance_service.py index c4ba269..90cad0d 100644 --- a/tests/test_rebalance_service.py +++ b/tests/test_rebalance_service.py @@ -35,6 +35,7 @@ from application.runtime_dependencies import LongBridgeRebalanceConfig, LongBridgeRebalanceRuntime from quant_platform_kit.common.account_identity import BrokerAccountIdentity from quant_platform_kit.common.execution_commands import ExecutionCommandState, ExecutionCommandStore + from application.durable_execution_commands import build_live_execution_command from notifications.telegram import build_translator from decision_mapper import map_strategy_decision_to_plan from quant_platform_kit.common.strategy_contracts import PositionTarget, StrategyDecision @@ -3541,6 +3542,84 @@ def test_hybrid_heartbeat_hides_empty_semiconductor_fields_and_shows_benchmark_l self.assertNotIn("🏦 收入层锁定占比: ", sent_messages[0]) + def test_historical_sg_duplicate_closes_claim_only_with_confirmed_report(self): + for prior_report_outcome, expected_state in ((True, ExecutionCommandState.CANCELLED), (RuntimeError("unreadable"), ExecutionCommandState.CLAIMED)): + with self.subTest(prior_report_outcome=type(prior_report_outcome).__name__), TemporaryDirectory() as command_dir, TemporaryDirectory() as marker_dir: + old_plan = _build_plan( + strategy_symbols=("SOXL",), risk_symbols=("SOXL",), targets={"SOXL": 400.0}, + market_values={"SOXL": 0.0}, sellable_quantities={"SOXL": 0}, quantities={"SOXL": 0}, + current_min_trade=10.0, trade_threshold_value=10.0, investable_cash=500.0, + available_cash=500.0, total_strategy_equity=500.0, market_status="Risk on", + deploy_ratio_text="70.0%", income_ratio_text="0.0%", income_locked_ratio_text="0.0%", + signal_message="Historical SG duplicate", portfolio_rows=(("SOXL",),), + signal_date="2026-09-10", effective_date="2026-09-11", + ) + new_plan = {**old_plan, "execution": {**old_plan["execution"], "signal_date": "2026-09-11", "effective_date": "2026-09-12"}} + command_store = ExecutionCommandStore(local_dir=command_dir) + old_command = build_live_execution_command( + platform="longbridge", account_scope="SG", strategy_profile="soxl_soxx_trend_income", + physical_account_id="lb-sg-001", runtime_identity_digest="a" * 64, + execution=old_plan["execution"], allocation=old_plan["allocation"], + ) + self.assertTrue(command_store.enqueue(old_command)) + reads = [] + + class EvidenceStore: + def has_marker(self, key): + reads.append(("marker", key)) + return False + + def has_prior_execution_report(self, **kwargs): + reads.append(("report", kwargs["signal_date"], kwargs["effective_date"])) + if isinstance(prior_report_outcome, Exception): + raise prior_report_outcome + return prior_report_outcome + + snapshot = replace(_build_snapshot(old_plan), as_of="2026-09-11") + runtime = LongBridgeRebalanceRuntime( + bootstrap=lambda: ("quote", "trade", {}), + resolve_rebalance_plan=lambda **_kwargs: new_plan, + resolve_frozen_rebalance_plan=lambda *, allocation, execution, snapshot: old_plan, + market_data_port_factory=lambda _quote: CallableMarketDataPort( + quote_loader=lambda _symbol: (_ for _ in ()).throw(AssertionError("duplicate must not quote")) + ), + estimate_max_purchase_quantity=lambda *_args, **_kwargs: (_ for _ in ()).throw(AssertionError("duplicate must not estimate")), + notifications=CallableNotificationPort(lambda _message: None), + notify_issue=lambda _title, _detail: None, + portfolio_port_factory=lambda *_contexts: CallablePortfolioPort(lambda: snapshot), + execution_port_factory=lambda _context: CallableExecutionPort( + lambda _intent: (_ for _ in ()).throw(AssertionError("duplicate must not submit")) + ), + ) + config = LongBridgeRebalanceConfig( + limit_sell_discount=0.995, limit_buy_premium=1.005, separator="-", + translator=build_translator("en"), with_prefix=lambda message: message, + strategy_profile="soxl_soxx_trend_income", execution_state_account_scope="SG", + physical_account_id="lb-sg-001", dry_run_only=False, + execution_dedup_enabled=True, execution_state_store=EvidenceStore(), + durable_execution_command_live_enabled=True, execution_command_store=command_store, + durable_live_execution_session_authorized=True, + durable_execution_runtime_identity_digest="a" * 64, + notify_no_trade_cycles=False, + ) + + result = rebalance_service.run_strategy(runtime=runtime, config=config) + + self.assertFalse(result.action_done) + self.assertEqual(result.pending_orders, ()) + self.assertEqual(reads[0][0], "marker") + self.assertEqual(reads[1], ("report", "2026-09-10", "2026-09-11")) + self.assertEqual(command_store.current_state(old_command), expected_state) + self.assertEqual( + len([event for event in command_store.events(old_command) if event.state is ExecutionCommandState.CLAIMED]), + 1, + ) + self.assertEqual( + len([event for event in command_store.events(old_command) if event.state is ExecutionCommandState.CANCELLED]), + int(expected_state is ExecutionCommandState.CANCELLED), + ) + + class RequiredExecutionClaimTests(unittest.TestCase): def setUp(self): self.orders = []