diff --git a/README.md b/README.md index f14159f..d6f58d1 100644 --- a/README.md +++ b/README.md @@ -36,6 +36,27 @@ Provides streaming orchestration that delivers LLM responses token-by-token for - Progressive rendering - Interruptible generation +### Goals changed during a turn + +Hosts may set `coordinator.session_state["goal"]` before execution or while a +conversation turn is running. The original goal fields (`condition`, +`turns_used`, `last_reason`, `cap`) remain supported; the loop initializes its +additional bookkeeping whenever it consumes a goal. A goal introduced during +a turn receives the same bounded tool/response evidence and completion +ordering as one set before execution. + +If the host clears, replaces, or revises a goal while its evaluator, stall +judge, or terminal summary is awaiting a response, the old result is discarded. +The current answer is finalized and returned; a successor goal is preserved +for the next explicit execution, without replaying the input or starting work +for that successor. Hosts that revise a goal in place must change its condition, +cap, or task revision. Cancellation stops automatic continuation and clears +only the goal being cancelled. + +`orchestrator:complete` is deferred until the goal decision is known. Continuing +progress hooks run before an intermediate completion so a host can pause there; +the final completion then remains marked `goal_final: true`. + ### Durable completed-tool checkpoints A host may register a zero-argument `session.durable_checkpoint` capability on diff --git a/amplifier_module_loop_streaming/__init__.py b/amplifier_module_loop_streaming/__init__.py index e66eaa8..b1f488f 100644 --- a/amplifier_module_loop_streaming/__init__.py +++ b/amplifier_module_loop_streaming/__init__.py @@ -1411,22 +1411,82 @@ async def execute( self._goal_model_basis = None self._retention_capability_warned = False - # Peek at goal state *before* the first turn. Goal state can only be - # set (by the app layer's /goal command) before execute() is called, - # and can only be cleared (never newly set) from within this method's - # own goal loop below -- so whether a goal is active for this whole - # execute() invocation is a stable fact determinable once, up front. - # This lets us tell _execute_one_turn whether ORCHESTRATOR_COMPLETE - # emission must be deferred (see its `goal_turn` param and - # _flush_pending_complete below) -- deferred because "was this turn - # the *final* one" isn't knowable until the evaluator judges its - # result, which happens only after the turn (and its completion - # event) would otherwise already have fired. + # Hosts may create, revise, or pause a goal during a turn. Normalize + # the goal whenever it is consumed, not just at execute() entry. initial_goal = coordinator.session_state.get("goal") if coordinator else None if initial_goal: self._ensure_goal_defaults(initial_goal) goal_turn = (initial_goal["turns_used"] + 1) if initial_goal else None + def goal_version(goal): + # Identity covers replacement; these fields also cover an in-place + # edit while an evaluator or judge is awaiting its provider. + return tuple( + goal.get(key) + for key in ("condition", "task_id", "task_revision", "cap") + ) + + def current(goal, version): + return ( + coordinator.session_state.get("goal") is goal + and goal_version(goal) == version + ) + + async def cancel_goal(goal, version): + cancelled_current = current(goal, version) + if cancelled_current: + coordinator.session_state["goal"] = None + # Best-effort diagnostics must not mask the original cancellation. + try: + await self._flush_pending_complete(goal_final=True) + except (Exception, asyncio.CancelledError): + logger.warning("Failed to emit cancelled goal completion") + if cancelled_current and not coordinator.session_state.get("goal"): + try: + await hooks.emit( + "orchestrator:goal_progress", + self._goal_progress_payload( + goal, state="cancelled", + reason=f"condition was: {goal['condition']}", + ), + ) + except (Exception, asyncio.CancelledError): + logger.warning("Failed to emit cancelled goal progress") + + async def stop_if_changed(goal, version): + if not current(goal, version): + await self._flush_pending_complete(goal_final=True) + return True + if coordinator.cancellation.is_cancelled: + await cancel_goal(goal, version) + return True + return False + + async def finish_goal(goal, version, state, reason, **details): + summary = None + if self._goal_run_needs_summary(state): + try: + summary = await self._summarize_goal_run( + goal, providers, hooks, coordinator, state, + **({"error_detail": reason} if state == "error" else {}), + ) + except asyncio.CancelledError: + await cancel_goal(goal, version) + raise + if await stop_if_changed(goal, version): + return + coordinator.session_state["goal"] = None + await self._flush_pending_complete(goal_final=True) + # A completion hook may install a successor goal. An old result + # must not mark that successor achieved or clear it. + if coordinator.session_state.get("goal"): + return + await hooks.emit( + "orchestrator:goal_progress", + self._goal_progress_payload( + goal, state=state, reason=reason, summary=summary, **details), + ) + async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: try: return await self._execute_one_turn( @@ -1434,7 +1494,7 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: goal_turn=goal_turn, ) except (Exception, asyncio.CancelledError) as error: - if goal_turn is not None and coordinator is not None: + if coordinator is not None: failed_goal = coordinator.session_state.get("goal") coordinator.session_state["goal"] = None # Each diagnostic is best effort. Neither may replace the @@ -1468,6 +1528,8 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: # part of that count, since it isn't a re-prompt. Set True right # after each continuation _execute_one_turn call below. is_continuation_turn = False + evaluated_goal = initial_goal + evaluated_version = None while True: goal = coordinator.session_state.get("goal") @@ -1479,19 +1541,13 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: await self._flush_pending_complete(goal_final=True) return full_response - if coordinator.cancellation.is_cancelled: - coordinator.session_state["goal"] = None - await self._flush_pending_complete(goal_final=True) - # No fast-model summary for a user-initiated cancellation -- - # nothing to explain. - await hooks.emit( - "orchestrator:goal_progress", - self._goal_progress_payload( - goal, - state="cancelled", - reason=f"condition was: {goal['condition']}", - ), - ) + self._ensure_goal_defaults(goal) + version = goal_version(goal) + if goal is not evaluated_goal or version != evaluated_version: + is_continuation_turn = False + evaluated_goal, evaluated_version = goal, version + + if await stop_if_changed(goal, version): return full_response goal["turns_used"] += 1 @@ -1521,34 +1577,18 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: satisfied, reason = await self._evaluate_goal( goal["condition"], context, providers, hooks, coordinator ) + except asyncio.CancelledError: + await cancel_goal(goal, version) + raise except Exception as e: - # FAIL LOUD: never silently keep going, never silently - # declare success -- regardless of whether this turn also - # hit the cap. (Previously, an eval failure exactly at the - # cap boundary was swallowed into a bare "cap_hit" with no - # reason via a separate "final evaluation at cap" call; - # now there's only one evaluation call per turn, so a - # failure here is always reported honestly as "error".) - coordinator.session_state["goal"] = None - await self._flush_pending_complete(goal_final=True) - summary = ( - await self._summarize_goal_run( - goal, - providers, - hooks, - coordinator, - "error", - error_detail=str(e), - ) - if self._goal_run_needs_summary("error") - else None - ) - await hooks.emit( - "orchestrator:goal_progress", - self._goal_progress_payload( - goal, state="error", reason=str(e), summary=summary - ), - ) + # A result (including an error) for an obsolete goal cannot + # clear, complete, or continue the user's newer goal. + if await stop_if_changed(goal, version): + return full_response + await finish_goal(goal, version, "error", str(e)) + return full_response + + if await stop_if_changed(goal, version): return full_response goal["last_reason"] = reason @@ -1556,24 +1596,7 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: self._record_goal_evidence(goal, reason) if satisfied: - # Achieved regardless of cap_hit -- the cap merely stops the - # loop from re-checking again, it never fails a goal that - # was, in fact, satisfied on its last permitted turn. - coordinator.session_state["goal"] = None - await self._flush_pending_complete(goal_final=True) - summary = ( - await self._summarize_goal_run( - goal, providers, hooks, coordinator, "achieved" - ) - if self._goal_run_needs_summary("achieved") - else None - ) - await hooks.emit( - "orchestrator:goal_progress", - self._goal_progress_payload( - goal, state="achieved", reason=reason, summary=summary - ), - ) + await finish_goal(goal, version, "achieved", reason) return full_response # Stall bookkeeping runs for every completed continuation turn @@ -1646,24 +1669,16 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: coordinator, trigger=stall_trigger, ) + except asyncio.CancelledError: + await cancel_goal(goal, version) + raise except Exception as e: + if await stop_if_changed(goal, version): + return full_response logger.exception("/goal: progress assessment failed") - coordinator.session_state["goal"] = None - await self._flush_pending_complete(goal_final=True) - summary = await self._summarize_goal_run( - goal, - providers, - hooks, - coordinator, - "error", - error_detail=str(e), - ) - await hooks.emit( - "orchestrator:goal_progress", - self._goal_progress_payload( - goal, state="error", reason=str(e), summary=summary - ), - ) + await finish_goal(goal, version, "error", str(e)) + return full_response + if await stop_if_changed(goal, version): return full_response if progress_verdict == "demonstrated": self._snapshot_goal_progress_anchors(goal) @@ -1675,34 +1690,9 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: if is_stalled and ( stall_trigger == "recovery" or goal["escalated"] or cap_hit ): - # Either this is the second trip (escalation already used - # and it stalled again), or it's the first trip but there's - # no cap budget left to offer the one-shot rescue turn. - # Either way: hard stop, reported as "stalled" rather than - # "cap_hit" -- we now know definitively the run is stuck, - # which is more informative than "ran out of turns", even - # when the cap also happened to run out on this same turn. - # This is a LOUD failure state, never mistakable for success. - coordinator.session_state["goal"] = None - await self._flush_pending_complete(goal_final=True) - summary = ( - await self._summarize_goal_run( - goal, providers, hooks, coordinator, "stalled" - ) - if self._goal_run_needs_summary("stalled") - else None - ) - await hooks.emit( - "orchestrator:goal_progress", - self._goal_progress_payload( - goal, - state="stalled", - reason=reason, - stall_detail=stall_detail, - stall_verdict=stall_verdict, - progress_verdict=progress_verdict, - summary=summary, - ), + await finish_goal( + goal, version, "stalled", reason, stall_detail=stall_detail, + stall_verdict=stall_verdict, progress_verdict=progress_verdict, ) return full_response @@ -1715,7 +1705,8 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: goal["recovery_pending"] = { "before": list(goal.get("progress_evidence", [])[-1:]) } - await self._flush_pending_complete(goal_final=False) + # Let host controls veto continuation before publishing an + # intermediate completion for the saved answer. await hooks.emit( "orchestrator:goal_progress", self._goal_progress_payload( @@ -1726,6 +1717,11 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: progress_verdict=progress_verdict, ), ) + if await stop_if_changed(goal, version): + return full_response + await self._flush_pending_complete(goal_final=False) + if await stop_if_changed(goal, version): + return full_response goal["continuations"] += 1 stall_prompt = self._goal_stall_escalation_prompt( goal, @@ -1741,29 +1737,11 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: continue if cap_hit: - # Not stalled (or the mechanical streak hadn't reached - # threshold this turn) -- the cap simply ran out. `reason` - # is already known from the evaluation above; no separate - # "final" evaluation call is needed since evaluation now - # always runs before the cap is checked. - coordinator.session_state["goal"] = None - await self._flush_pending_complete(goal_final=True) - summary = await self._summarize_goal_run( - goal, providers, hooks, coordinator, "cap_hit" - ) - await hooks.emit( - "orchestrator:goal_progress", - self._goal_progress_payload( - goal, - state="cap_hit", - reason=reason, - progress_verdict=progress_verdict, - summary=summary, - ), + await finish_goal( + goal, version, "cap_hit", reason, progress_verdict=progress_verdict, ) return full_response - await self._flush_pending_complete(goal_final=False) await hooks.emit( "orchestrator:goal_progress", self._goal_progress_payload( @@ -1774,6 +1752,11 @@ async def run_turn(turn_prompt: str, *, goal_turn: int | None) -> str: ), ) + if await stop_if_changed(goal, version): + return full_response + await self._flush_pending_complete(goal_final=False) + if await stop_if_changed(goal, version): + return full_response goal["continuations"] += 1 full_response = await run_turn( reason, @@ -1803,9 +1786,9 @@ async def _execute_one_turn( goal_turn: When this turn is part of an active /goal auto-continue pursuit (spike: docs/designs/goal-command.md), the 1-based goal-turn number; ``None`` when no goal is active. When - ``None`` (the default), ``ORCHESTRATOR_COMPLETE`` is emitted - immediately as before -- zero behavior change. When set, the - caller (``execute()``'s goal loop) doesn't yet know whether + ``None`` (the default), completion is immediate unless a + goal was introduced during this turn. When a goal is active, + the caller (``execute()``'s goal loop) doesn't yet know whether this is the *final* turn of the pursuit (that depends on the evaluator judging this turn's result, which hasn't happened yet), so emission is deferred via @@ -1824,7 +1807,7 @@ async def _execute_one_turn( # increment it) so execute()'s stall detection can read an accurate # "did this turn run any tools" count once this method returns. self._tool_calls_this_turn = 0 - self._goal_turn_evidence = {"tools": []} if goal_turn is not None else None + self._goal_turn_evidence = {"tools": []} if coordinator is not None else None self._llm_calls_this_turn = 0 # Layer 1 call-budget bookkeeping, reset per turn alongside # _tool_calls_this_turn above (spec: 298-replacement). @@ -1871,6 +1854,9 @@ async def _execute_one_turn( # payload below. None when no goal is active -- same discriminator # pattern as goal_turn/goal_final. goal_state = coordinator.session_state.get("goal") if coordinator else None + if goal_state: + self._ensure_goal_defaults(goal_state) + goal_turn = goal_state["turns_used"] + 1 payload = { "orchestrator": "loop-streaming", diff --git a/tests/test_live_goal_lifecycle.py b/tests/test_live_goal_lifecycle.py new file mode 100644 index 0000000..97b00e4 --- /dev/null +++ b/tests/test_live_goal_lifecycle.py @@ -0,0 +1,277 @@ +"""Live hosts may change a goal while an ordinary turn or utility call awaits. + +Real execute() flow, deterministic providers, and hooks; no network/model calls. +""" +import asyncio + +import pytest + +from .test_goal_loop import ( + FakeProvider, MockContext, MockCoordinator, MockHooks, MockTurnResponse, + MockTool, MockToolCall, _make_orchestrator, +) + + +def goal(condition="original", **extra): + return {"condition": condition, "turns_used": 0, "last_reason": None, + "cap": None, **extra} + + +class ChangingProvider(FakeProvider): + def __init__(self): + super().__init__() + self.callbacks = {} + self.calls = [] + + async def complete(self, request, **kwargs): + text = next((m.content for m in request.messages if m.role == "system"), "") + text = text if isinstance(text, str) else "" + kind = ("eval" if "tool-less evaluator" in text else + "judge" if "tool-less judge" in text else + "summary" if "single, short line for a developer" in text else "turn") + self.calls.append(kind) + # Yield at the same boundary at which real network calls permit host + # control changes, then apply one-shot deterministic host actions. + await asyncio.sleep(0) + action = self.callbacks.pop(kind, None) + if action: + action() + return await super().complete(request, **kwargs) + + +async def run(provider, coordinator, hooks=None, config=None, tools=None): + hooks = hooks or MockHooks() + engine = _make_orchestrator(config) + result = await engine.execute("work", MockContext(), {"main": provider}, + tools or {}, hooks, coordinator) + return engine, hooks, result + + +@pytest.mark.asyncio +@pytest.mark.parametrize("already_present", [False, True]) +async def test_partial_goal_at_entry_or_created_mid_turn_gets_full_defaults(already_present): + coordinator, provider, hooks = MockCoordinator(), ChangingProvider(), MockHooks() + active = goal(task_id="task", task_revision=1) + if already_present: + coordinator.session_state["goal"] = active + else: + provider.callbacks["turn"] = lambda: coordinator.session_state.update(goal=active) + # Verify the answer's complete event is deferred until its new goal has + # been evaluated, not prematurely published as the final completion. + provider.callbacks["eval"] = lambda: assert_no_completions(hooks) + provider.turn_queue.append(MockTurnResponse(text="answer")) + provider.eval_queue.append((True, "done")) + engine, hooks, result = await run(provider, coordinator, hooks) + assert result == "answer" + assert active["progress_evidence"] + assert active["turns_used"] == 1 + assert coordinator.session_state["goal"] is None + assert engine._pending_orchestrator_complete is None + assert [e["state"] for e in hooks.goal_progress_events()] == ["achieved"] + assert [(e["goal_turn"], e["goal_final"]) for e in hooks.orchestrator_complete_events()] == [(1, True)] + + +def assert_no_completions(hooks): + assert hooks.orchestrator_complete_events() == [] + + +@pytest.mark.asyncio +async def test_goal_added_during_tool_turn_keeps_evidence_and_continues(): + coordinator, provider = MockCoordinator(), ChangingProvider() + active = goal() + provider.callbacks["turn"] = lambda: coordinator.session_state.update(goal=active) + provider.turn_queue.extend([MockTurnResponse(tool_calls=[MockToolCall()]), + MockTurnResponse(text="first"), MockTurnResponse(text="second")]) + provider.eval_queue.extend([(False, "needs another step"), (True, "finished")]) + _, hooks, result = await run(provider, coordinator, tools={"mock_tool": MockTool()}) + assert result == "second" + assert active["continuations"] == 1 + assert active["turns_used"] == 2 + assert active["progress_evidence"][0]["tools"] + assert [e["goal_final"] for e in hooks.orchestrator_complete_events()] == [False, True] + + +@pytest.mark.asyncio +@pytest.mark.parametrize("outcome", [True, False, "error"]) +@pytest.mark.parametrize("change", ["replace", "revise", "clear", "pause"]) +async def test_obsolete_evaluation_cannot_complete_clear_or_continue_new_goal(outcome, change): + coordinator, provider = MockCoordinator(), ChangingProvider() + original = goal(task_id="task", task_revision=1) + replacement = goal("revised", task_id="task", task_revision=2) + coordinator.session_state["goal"] = original + provider.turn_queue.append(MockTurnResponse(text="already saved answer")) + provider.eval_queue.extend([(bool(outcome), "obsolete"), (True, "current result")]) + + def update(): + if change in {"clear", "pause"}: + coordinator.session_state["goal"] = None + elif change == "replace": + coordinator.session_state["goal"] = replacement + else: + original.update(condition="revised", task_revision=2) + if outcome == "error": + provider.eval_queue.pop(0) + raise RuntimeError("obsolete evaluator failed") + + provider.callbacks["eval"] = update + _, hooks, result = await run(provider, coordinator) + assert result == "already saved answer" + assert "obsolete" not in original["reasons"] + assert provider.calls.count("turn") == 1 + assert hooks.goal_progress_events() == [] + assert provider.calls.count("eval") == 1 + if change in {"clear", "pause"}: + assert coordinator.session_state["goal"] is None + else: + expected = replacement if change == "replace" else original + assert coordinator.session_state["goal"] is expected + assert expected.get("reasons", []) == [] + assert [e["goal_final"] for e in hooks.orchestrator_complete_events()] == [True] + + +@pytest.mark.asyncio +@pytest.mark.parametrize("change", ["replace", "clear"]) +@pytest.mark.parametrize("boundary", ["judge", "summary"]) +async def test_stale_judge_and_terminal_summary_do_not_end_successor_goal(change, boundary): + coordinator, provider = MockCoordinator(), ChangingProvider() + original = goal(cap=1 if boundary == "summary" else None) + successor = goal("successor") + coordinator.session_state["goal"] = original + provider.turn_queue.extend([MockTurnResponse(text="first"), MockTurnResponse(text="second")]) + provider.eval_queue.extend([(False, "blocked"), (False, "blocked"), (True, "successor done")] + if boundary == "judge" else [(False, "blocked"), (True, "successor done")]) + provider.judge_queue.append((True, "stalled")) + provider.callbacks[boundary] = lambda: coordinator.session_state.update( + goal=successor if change == "replace" else None) + _, hooks, _ = await run(provider, coordinator, config={"goal_stall_threshold": 1}) + events = hooks.goal_progress_events() + assert not any(e["state"] in {"stalled", "cap_hit", "error"} for e in events) + if change == "replace": + assert coordinator.session_state["goal"] is successor + assert successor.get("reasons", []) == [] + else: + assert coordinator.session_state["goal"] is None + assert hooks.orchestrator_complete_events()[-1]["goal_final"] is True + + +@pytest.mark.asyncio +async def test_pause_in_continuing_progress_hook_stops_before_next_turn(): + coordinator, provider = MockCoordinator(), ChangingProvider() + coordinator.session_state["goal"] = goal() + provider.turn_queue.append(MockTurnResponse(text="first")) + provider.eval_queue.append((False, "needs work")) + + class PausingHooks(MockHooks): + async def emit(self, event_name, payload=None): + result = await super().emit(event_name, payload) + if event_name == "orchestrator:goal_progress": + coordinator.session_state["goal"] = None + return result + + _, hooks, _ = await run(provider, coordinator, PausingHooks()) + assert provider.calls == ["turn", "eval"] + assert coordinator.session_state["goal"] is None + assert [e["goal_final"] for e in hooks.orchestrator_complete_events()] == [True] + + +@pytest.mark.asyncio +@pytest.mark.parametrize("boundary", ["eval", "judge", "summary"]) +@pytest.mark.parametrize("mode", ["flag", "raise"]) +async def test_cancellation_during_utility_call_stops_and_clears_current_goal(boundary, mode): + coordinator, provider, hooks = MockCoordinator(), ChangingProvider(), MockHooks() + active = goal(cap=1 if boundary == "summary" else None) + coordinator.session_state["goal"] = active + provider.turn_queue.extend([MockTurnResponse(text="first"), MockTurnResponse(text="second")]) + provider.eval_queue.extend([(False, "blocked"), (False, "blocked")]) + provider.judge_queue.append((True, "stalled")) + + def cancel(): + if mode == "raise": + raise asyncio.CancelledError() + coordinator.cancellation.is_cancelled = True + + provider.callbacks[boundary] = cancel + if mode == "raise": + with pytest.raises(asyncio.CancelledError): + await run(provider, coordinator, hooks, config={"goal_stall_threshold": 1}) + else: + await run(provider, coordinator, hooks, config={"goal_stall_threshold": 1}) + assert coordinator.session_state["goal"] is None + assert hooks.goal_progress_events()[-1]["state"] == "cancelled" + assert hooks.orchestrator_complete_events()[-1]["goal_final"] is True + assert provider.calls.count("turn") == (2 if boundary == "judge" else 1) + if boundary == "eval": + assert active["reasons"] == [] + + +@pytest.mark.asyncio +async def test_cancellation_of_obsolete_evaluator_does_not_clear_successor(): + coordinator, provider, hooks = MockCoordinator(), ChangingProvider(), MockHooks() + coordinator.session_state["goal"] = goal() + successor = goal("replacement") + provider.turn_queue.append(MockTurnResponse(text="saved answer")) + + def replace_and_cancel(): + coordinator.session_state["goal"] = successor + raise asyncio.CancelledError() + + provider.callbacks["eval"] = replace_and_cancel + with pytest.raises(asyncio.CancelledError): + await run(provider, coordinator, hooks) + assert coordinator.session_state["goal"] is successor + assert hooks.goal_progress_events() == [] + assert hooks.orchestrator_complete_events()[-1]["goal_final"] is True + + +@pytest.mark.asyncio +@pytest.mark.parametrize("change", ["replace", "revise"]) +async def test_goal_revised_during_conversation_is_evaluated_with_current_condition(change): + coordinator, provider = MockCoordinator(), ChangingProvider() + original = goal(task_revision=1) + replacement = goal("revised", task_revision=2) + coordinator.session_state["goal"] = original + + def revise(): + if change == "replace": + coordinator.session_state["goal"] = replacement + else: + original.update(condition="revised", task_revision=2) + + provider.callbacks["turn"] = revise + provider.turn_queue.append(MockTurnResponse(text="answer")) + provider.eval_queue.append((True, "revised condition met")) + _, hooks, result = await run(provider, coordinator) + assert result == "answer" + active = replacement if change == "replace" else original + assert active["progress_evidence"] + assert active["reasons"] == ["revised condition met"] + assert hooks.goal_progress_events()[-1]["condition"] == "revised" + assert len(hooks.orchestrator_complete_events()) == 1 + assert coordinator.session_state["goal"] is None + + +@pytest.mark.asyncio +@pytest.mark.parametrize("event", ["orchestrator:complete", "orchestrator:goal_progress"]) +async def test_utility_cancellation_preserves_cancellation_when_diagnostics_fail(event): + coordinator, provider = MockCoordinator(), ChangingProvider() + coordinator.session_state["goal"] = goal() + provider.turn_queue.append(MockTurnResponse(text="saved answer")) + cancellation = asyncio.CancelledError() + + def cancel(): + raise cancellation + + class FailingHooks(MockHooks): + async def emit(self, event_name, payload=None): + result = await super().emit(event_name, payload) + if event_name == event: + raise RuntimeError("diagnostic hook failed") + return result + + hooks = FailingHooks() + provider.callbacks["eval"] = cancel + with pytest.raises(asyncio.CancelledError) as raised: + await run(provider, coordinator, hooks) + assert raised.value is cancellation + assert coordinator.session_state["goal"] is None + assert hooks.goal_progress_events()[-1]["state"] == "cancelled"