diff --git a/solstone/think/providers/openhands.py b/solstone/think/providers/openhands.py index d4a67bf1a..fb9596c12 100644 --- a/solstone/think/providers/openhands.py +++ b/solstone/think/providers/openhands.py @@ -75,6 +75,12 @@ from solstone.think.providers.shared import ( safe_raw, validate_generate_result_strict, ) +from solstone.think.responsiveness import ( + NON_RESPONSIVE_OUTPUT_FRAGMENT, + NON_RESPONSIVE_RAW_OUTPUT_CAP_CHARS, + NON_RESPONSIVE_REASON_CODE, + classify_output_responsiveness, +) from solstone.think.utils import get_journal, get_project_root, now_ms LOG = logging.getLogger("solstone.think.providers.openhands") @@ -1390,6 +1396,18 @@ def _raw_event(event: Any) -> list[dict[str, Any]]: return safe_raw([{"type": event.__class__.__name__, "repr": repr(event)}]) +def _non_responsive_raw_payload( + output: str, matched_signal: str | None +) -> list[dict[str, Any]]: + payload: dict[str, Any] = { + "reason_code": NON_RESPONSIVE_REASON_CODE, + "non_responsive_output": output[:NON_RESPONSIVE_RAW_OUTPUT_CAP_CHARS], + } + if matched_signal is not None: + payload["non_responsive_matched_signal"] = matched_signal + return safe_raw([payload]) + + def _tool_arguments(event: Any) -> dict[str, Any]: tool_call = getattr(event, "tool_call", None) raw_arguments = getattr(tool_call, "arguments", None) @@ -1833,31 +1851,49 @@ async def run_cogitate( raise terminal_error result = translator.result() + non_responsive = False + non_responsive_raw: list[dict[str, Any]] | None = None + if isinstance(result, str) and result.strip(): + responsiveness = classify_output_responsiveness(result) + if responsiveness.non_responsive: + non_responsive = True + non_responsive_raw = _non_responsive_raw_payload( + result[:NON_RESPONSIVE_RAW_OUTPUT_CAP_CHARS], + responsiveness.matched_signal, + ) + result = None usage = _usage_delta(usage_start, llm) if wall_clock_exceeded: has_partial = bool(result and result.strip()) - error_text = ( - "wall_clock_exceeded: cogitate run exceeded its wall-clock " - "deadline and was force-finished with a partial result preserved" - if has_partial - else "wall_clock_exceeded: cogitate run exceeded its wall-clock " - "deadline and was force-finished before emitting a final result" - ) + if non_responsive: + error_text = ( + "wall_clock_exceeded: cogitate run exceeded its wall-clock " + f"deadline after producing {NON_RESPONSIVE_OUTPUT_FRAGMENT}" + ) + else: + error_text = ( + "wall_clock_exceeded: cogitate run exceeded its wall-clock " + "deadline and was force-finished with a partial result preserved" + if has_partial + else "wall_clock_exceeded: cogitate run exceeded its wall-clock " + "deadline and was force-finished before emitting a final result" + ) conversation.close() - callback.emit( - { - "event": "error", - "error": error_text, - "reason_code": "wall_clock_exceeded", - "provider": provider, - "result": result, - "usage": usage, - "terminal": True, - "cli_session_id": str(conversation_id), - "ts": now_ms(), - } - ) - return result + error_event: dict[str, Any] = { + "event": "error", + "error": error_text, + "reason_code": "wall_clock_exceeded", + "provider": provider, + "result": result, + "usage": usage, + "terminal": True, + "cli_session_id": str(conversation_id), + "ts": now_ms(), + } + if non_responsive_raw is not None: + error_event["raw"] = non_responsive_raw + callback.emit(error_event) + return None if non_responsive else result if translator._cost_force_stopped or translator.max_turns_exhausted: reason_code = ( "token_budget_exceeded" @@ -1865,7 +1901,18 @@ async def run_cogitate( else "max_turns_exhausted" ) has_partial = bool(result and result.strip()) - if reason_code == "token_budget_exceeded": + if non_responsive: + if reason_code == "token_budget_exceeded": + error_text = ( + "token_budget_exceeded: cogitate run reached its per-run " + f"resource budget after producing {NON_RESPONSIVE_OUTPUT_FRAGMENT}" + ) + else: + error_text = ( + "max_turns_exhausted: cogitate run reached its turn budget " + f"after producing {NON_RESPONSIVE_OUTPUT_FRAGMENT}" + ) + elif reason_code == "token_budget_exceeded": error_text = ( "token_budget_exceeded: cogitate run reached its per-run " "resource budget and was force-finished with a partial result " @@ -1884,45 +1931,71 @@ async def run_cogitate( "and was force-finished before emitting a final result" ) conversation.close() - callback.emit( - { - "event": "error", - "error": error_text, - "reason_code": reason_code, - "provider": provider, - "result": result, - "usage": usage, - "terminal": True, - "cli_session_id": str(conversation_id), - "ts": now_ms(), - } - ) - return result + error_event = { + "event": "error", + "error": error_text, + "reason_code": reason_code, + "provider": provider, + "result": result, + "usage": usage, + "terminal": True, + "cli_session_id": str(conversation_id), + "ts": now_ms(), + } + if non_responsive_raw is not None: + error_event["raw"] = non_responsive_raw + callback.emit(error_event) + return None if non_responsive else result execution_status = _conversation_execution_status(conversation) if execution_status in {"stuck", "paused"}: has_partial = bool(result and result.strip()) - error_text = ( - "agent_stuck: cogitate run was interrupted/stuck with a partial " - "result preserved" - if has_partial - else "agent_stuck: cogitate run was interrupted/stuck before " - "emitting a final result" - ) + if non_responsive: + error_text = ( + "agent_stuck: cogitate run was interrupted/stuck after producing " + f"{NON_RESPONSIVE_OUTPUT_FRAGMENT}" + ) + else: + error_text = ( + "agent_stuck: cogitate run was interrupted/stuck with a partial " + "result preserved" + if has_partial + else "agent_stuck: cogitate run was interrupted/stuck before " + "emitting a final result" + ) conversation.close() + error_event = { + "event": "error", + "error": error_text, + "reason_code": "agent_stuck", + "provider": provider, + "result": result, + "usage": usage, + "terminal": True, + "cli_session_id": str(conversation_id), + "ts": now_ms(), + } + if non_responsive_raw is not None: + error_event["raw"] = non_responsive_raw + callback.emit(error_event) + return None if non_responsive else result + if non_responsive: callback.emit( { "event": "error", - "error": error_text, - "reason_code": "agent_stuck", + "error": ( + f"{NON_RESPONSIVE_REASON_CODE}: cogitate run produced " + f"{NON_RESPONSIVE_OUTPUT_FRAGMENT}" + ), + "reason_code": NON_RESPONSIVE_REASON_CODE, "provider": provider, - "result": result, "usage": usage, "terminal": True, "cli_session_id": str(conversation_id), "ts": now_ms(), + "raw": non_responsive_raw, } ) - return result + return None if wants_emit_final and not (result and result.strip()): callback.emit( { diff --git a/tests/test_local.py b/tests/test_local.py index 1b2f46e62..d9d3e1683 100644 --- a/tests/test_local.py +++ b/tests/test_local.py @@ -23,6 +23,7 @@ from solstone.think.models import ( get_model_provider, ) from solstone.think.providers.artifact_proof import ReadinessOutcome +from solstone.think.responsiveness import NON_RESPONSIVE_REASON_CODE from solstone.think.schema_prep import SCHEMA_TRUNCATE_KEY from solstone.think.talents import TalentHookError @@ -3573,6 +3574,46 @@ def test_run_cogitate_byo_acquires_permit_and_records_no_telemetry(monkeypatch): assert permit.slot_index == 0 +def test_run_cogitate_local_delegated_non_responsive_single_event( + monkeypatch, +): + provider = _provider() + monkeypatch.setattr( + provider, + "resolve_local_endpoint", + lambda: _byo_endpoint(parallel_slots=1), + ) + terminal_event = { + "event": "error", + "error": "non-responsive output", + "reason_code": NON_RESPONSIVE_REASON_CODE, + "provider": "local", + "terminal": True, + "raw": [{"reason_code": NON_RESPONSIVE_REASON_CODE}], + } + + async def fake_cogitate(*_args, on_event=None, slot_lease=None, **_kwargs): + assert slot_lease is not None + on_event(terminal_event) + return None + + monkeypatch.setattr( + "solstone.think.providers.openhands.run_cogitate", + fake_cogitate, + ) + events: list[dict] = [] + + result = asyncio.run( + provider.run_cogitate( + {"model": LOCAL_MODEL, "timeout_seconds": 1}, + on_event=events.append, + ) + ) + + assert result is None + assert events == [terminal_event] + + def test_run_cogitate_byo_keeps_permit_for_non_sol_work(monkeypatch): provider = _provider() monkeypatch.setattr( diff --git a/tests/test_openhands_provider.py b/tests/test_openhands_provider.py index 6c6babcd8..8b3b13130 100644 --- a/tests/test_openhands_provider.py +++ b/tests/test_openhands_provider.py @@ -20,6 +20,10 @@ from solstone.think.cogitate_policy import ( from solstone.think.providers import openhands from solstone.think.providers.local_admission import LocalAdmissionTimeout from solstone.think.providers.shared import USAGE_KEYS, JSONEventCallback +from solstone.think.responsiveness import ( + NON_RESPONSIVE_OUTPUT_FRAGMENT, + NON_RESPONSIVE_REASON_CODE, +) from solstone.think.talent import get_talent, get_talent_configs from tests.openhands_fakes import _REGISTERED_TOOLS, install_fake_openhands @@ -689,6 +693,56 @@ def test_run_cogitate_emits_finish_when_emit_final_has_content( ] == ["finish"] +def test_run_cogitate_non_responsive_finish_emits_terminal_error( + fake_openhands, + monkeypatch, + tmp_path, +): + refusal = "I cannot describe this screen." + + async def emit_final_with_usage(conversation): + _seed_usage(conversation) + for callback in conversation.callbacks: + callback(_emit_final_action(fake_openhands, refusal)) + + fake_openhands.Conversation.arun_impl = emit_final_with_usage + config = _run_config(monkeypatch, tmp_path, output_path=str(tmp_path / "out.md")) + events: list[dict] = [] + + result = asyncio.run(openhands.run_cogitate(config, events.append)) + + assert result is None + error_events = [event for event in events if event["event"] == "error"] + assert len(error_events) == 1 + assert error_events[0]["reason_code"] == NON_RESPONSIVE_REASON_CODE + assert error_events[0]["terminal"] is True + assert error_events[0]["usage"]["total_tokens"] > 0 + assert NON_RESPONSIVE_OUTPUT_FRAGMENT in error_events[0]["error"] + assert error_events[0]["raw"][0]["reason_code"] == NON_RESPONSIVE_REASON_CODE + assert error_events[0]["raw"][0]["non_responsive_output"] == refusal + assert [event for event in events if event["event"] == "finish"] == [] + + +def test_run_cogitate_non_responsive_does_not_write_output_path( + fake_openhands, + monkeypatch, + tmp_path, +): + _install_emit_final_arun(fake_openhands, "I cannot describe this screen.") + output_path = tmp_path / "out.md" + config = _run_config(monkeypatch, tmp_path, output_path=str(output_path)) + events: list[dict] = [] + + result = asyncio.run(openhands.run_cogitate(config, events.append)) + + assert result is None + assert not output_path.exists() + assert [event for event in events if event["event"] == "finish"] == [] + assert [event for event in events if event["event"] == "error"][0][ + "reason_code" + ] == NON_RESPONSIVE_REASON_CODE + + def test_run_cogitate_daily_no_output_finishes_when_emit_final_has_content( fake_openhands, monkeypatch, @@ -865,6 +919,42 @@ def test_run_cogitate_force_stop_emits_token_budget_exceeded( assert [event for event in events if event["event"] == "finish"] == [] +def test_run_cogitate_non_responsive_budget_exit_keeps_budget_reason( + fake_openhands, + fixed_time, + monkeypatch, + tmp_path, +): + refusal = "I cannot describe this screen." + + async def hit_cost_cap_with_refusal(conversation): + _seed_usage(conversation) + conversation.agent.llm.metrics.accumulated_cost = DEFAULT_RUN_COST_CAP_USD + for callback in conversation.callbacks: + callback(_agent_message(fake_openhands, refusal)) + for callback in conversation.callbacks: + callback(_sol_action(fake_openhands, "c1")) + for callback in conversation.callbacks: + callback(_sol_action(fake_openhands, "c2")) + + fake_openhands.Conversation.arun_impl = hit_cost_cap_with_refusal + config = _run_config(monkeypatch, tmp_path) + events: list[dict] = [] + + result = asyncio.run(openhands.run_cogitate(config, events.append)) + + assert result is None + error_events = [event for event in events if event["event"] == "error"] + assert len(error_events) == 1 + assert error_events[0]["reason_code"] == "token_budget_exceeded" + assert error_events[0]["terminal"] is True + assert error_events[0]["result"] is None + assert error_events[0]["raw"][0]["reason_code"] == NON_RESPONSIVE_REASON_CODE + assert error_events[0]["raw"][0]["non_responsive_output"] == refusal + assert NON_RESPONSIVE_OUTPUT_FRAGMENT in error_events[0]["error"] + assert [event for event in events if event["event"] == "finish"] == [] + + def test_run_cogitate_cost_force_stop_with_partial_logs_once( fake_openhands, fixed_time,