import asyncio from pathlib import Path from frontend.web_runtime import WebRuntime from llm_panda.execution.cancellation import ( CancellationToken, ExecutionInterrupted, ) ROOT = Path(__file__).resolve().parents[1] def test_reload_snapshot_preserves_runtime_state_and_event_history(): runtime = WebRuntime(max_events=10) async def exercise(): await runtime.publish({"type": "log", "message": "planning"}) await runtime.publish( {"type": "plan", "plan": [{"step": 1}], "generation": 0} ) await runtime.publish( { "type": "state", "joints": [0.1, 0.2], "objects": {"banana": {"pos": [0.4, 0.1, 0.3]}}, }, journal=False, ) asyncio.run(exercise()) snapshot = runtime.snapshot() assert snapshot["state"]["joints"] == [0.1, 0.2] assert snapshot["plan"]["plan"] == [{"step": 1}] assert [event["message"] for event in snapshot["events"] if event["type"] == "log"] == [ "planning" ] def test_new_command_cancels_active_and_latest_pending_wins(): runtime = WebRuntime() first = runtime.submit("first") active = runtime.take_pending() second = runtime.submit("second") third = runtime.submit("third") assert active["task_id"] == first["task_id"] assert active["token"].requested pending = runtime.take_pending() assert pending["task_id"] == third["task_id"] assert pending["task_id"] != second["task_id"] def test_stale_result_is_journaled_but_does_not_replace_latest_result(): runtime = WebRuntime() first = runtime.submit("first") runtime.take_pending() second = runtime.submit("second") async def exercise(): await runtime.publish( { "type": "result", "generation": first["generation"], "success": False, } ) await runtime.publish( { "type": "result", "generation": second["generation"], "success": True, } ) asyncio.run(exercise()) assert runtime.latest_result["generation"] == second["generation"] assert runtime.latest_result["success"] def test_event_journal_is_bounded_and_replayable(): runtime = WebRuntime(max_events=3) async def exercise(): for index in range(5): await runtime.publish({"type": "log", "message": str(index)}) asyncio.run(exercise()) assert [event["message"] for event in runtime.events_after(0)] == [ "2", "3", "4", ] assert [event["message"] for event in runtime.events_after(4)] == ["4"] def test_cancellation_token_raises_only_after_request(): token = CancellationToken() token.checkpoint() token.request("new command") try: token.checkpoint() except ExecutionInterrupted as exc: assert "new command" in str(exc) else: raise AssertionError("checkpoint did not interrupt") def test_runtime_conversation_survives_client_disconnect(): runtime = WebRuntime() runtime.conversation_manager.observe( "move banana", [ { "function": "clarify", "args": { "message": "Where?", "missing_fields": ["destination"], "partial_intent": { "action": "pick_and_place", "object_id": "banana", }, }, } ], ) class Socket: def __init__(self): self.events = [] async def send_json(self, event): self.events.append(event) socket = Socket() async def connect_disconnect(): await runtime.subscribe(socket) await runtime.publish({"type": "log", "message": "waiting"}) await runtime.unsubscribe(socket) asyncio.run(connect_disconnect()) assert runtime.snapshot()["conversation_pending"] assert socket.events[0]["message"] == "waiting" def test_explicit_reset_clears_ui_history_but_keeps_monotonic_sequence(): runtime = WebRuntime() async def exercise(): await runtime.publish({"type": "log", "message": "old"}) asyncio.run(exercise()) old_seq = runtime.snapshot()["last_event_seq"] runtime.clear_task_history() asyncio.run(runtime.publish({"type": "log", "message": "new"})) snapshot = runtime.snapshot() assert [event["message"] for event in snapshot["events"]] == ["new"] assert snapshot["last_event_seq"] > old_seq def test_cancel_all_interrupts_active_and_discards_pending_command(): runtime = WebRuntime() runtime.submit("active") active = runtime.take_pending() runtime.submit("pending") assert runtime.cancel_all("reset") assert active["token"].requested assert not runtime.has_pending def test_repeated_standalone_cancel_is_idempotent(): runtime = WebRuntime() runtime.submit("active") active = runtime.take_pending() assert runtime.cancel_all("user_cancelled") assert runtime.cancel_all("user_cancelled") assert active["token"].requested assert runtime.active_generation == active["generation"] assert not runtime.has_pending def test_interruption_is_recorded_without_completing_the_goal(): runtime = WebRuntime() runtime.conversation_manager.record_interruption("move banana") last_turn = runtime.conversation_manager.recent_turns[-1] assert last_turn.kind == "interrupted" assert "move banana" in last_turn.content def test_web_pipeline_keeps_resolved_goal_across_replans_and_goal_check(): source = (ROOT / "frontend" / "interactive_server.py").read_text( encoding="utf-8" ) assert "resolved_goal = None" in source assert "resolved_goal.canonical_instruction" in source assert 'plan_kwargs["resolve_goal"] = False' in source assert "resolved_goal=resolved_goal" in source def test_web_replan_limit_uses_environment_configuration(): source = (ROOT / "frontend" / "interactive_server.py").read_text( encoding="utf-8" ) assert 'os.environ.get("MAX_REPLANS", "3")' in source assert "MAX_REPLANS = 999999" not in source def test_scene_reset_clears_cli_conversation_memory(): source = (ROOT / "main.py").read_text(encoding="utf-8") reset_body = source.split("def reset_objects", 1)[1].split( "def change_environment", 1, )[0] assert '"conversation_manager"' in reset_body assert "conversation_manager.clear()" in reset_body