mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-06 15:02:20 +02:00
fix(runtime): preserve provider error terminal ordering
This commit is contained in:
@@ -8,6 +8,7 @@ export PYTHON_DOTENV_DISABLED=1
|
||||
export ODYSSEUS_DATA_DIR="${ODYSSEUS_DATA_DIR:-/tmp/odysseus-runtime-decomposition-test-state}"
|
||||
exec "${ODYSSEUS_TEST_PYTHON:-python3}" -m pytest -q -p no:cacheprovider \
|
||||
tests/test_runtime_evidence_contract.py tests/test_agent_evidence.py \
|
||||
tests/test_completion_boundary.py \
|
||||
tests/test_agent_evidence_loop.py tests/test_agent_render_ownership.py \
|
||||
tests/test_agent_runs_terminal_order.py tests/test_agent_loop.py \
|
||||
tests/test_tool_task_cancelled_on_disconnect.py tests/test_turn_contract.py \
|
||||
|
||||
+3
-1
@@ -26258,6 +26258,9 @@ async def stream_agent_loop(
|
||||
# next model request. Do not expose a transient provider
|
||||
# error or terminate the turn before that retry.
|
||||
break
|
||||
# Let the completion gate retain the original failure even
|
||||
# when earlier tool evidence supplies useful fallback prose.
|
||||
yield chunk
|
||||
terminal_status = None
|
||||
try:
|
||||
error_line = next(
|
||||
@@ -26547,7 +26550,6 @@ async def stream_agent_loop(
|
||||
else "The model provider returned no usable output. No workspace change was made."
|
||||
)
|
||||
yield f'data: {json.dumps({"type": "final_response", "content": _failure_text})}\n\n'
|
||||
yield chunk
|
||||
# A terminal provider/request failure is not a completed Agent
|
||||
# round. Stop before empty-response synthesis, metrics,
|
||||
# teacher escalation, post-processing, or a success [DONE].
|
||||
|
||||
@@ -126,7 +126,7 @@ def with_completion_gate(func):
|
||||
done = False
|
||||
awaiting = False
|
||||
exhausted = False
|
||||
provider_error = False
|
||||
provider_error: str | None = None
|
||||
with bind_journal(journal):
|
||||
async with aclosing(func(*args, **kwargs)) as stream:
|
||||
async for chunk in stream:
|
||||
@@ -139,7 +139,11 @@ def with_completion_gate(func):
|
||||
data = None
|
||||
if not isinstance(data, dict):
|
||||
if chunk.startswith('event: error'):
|
||||
provider_error = True
|
||||
# The inner stream may still emit failed-terminal
|
||||
# diagnostics. Hold the original error until those
|
||||
# and the buffered answer have been released.
|
||||
provider_error = provider_error or chunk
|
||||
continue
|
||||
yield chunk
|
||||
continue
|
||||
kind = data.get('type')
|
||||
@@ -196,12 +200,16 @@ def with_completion_gate(func):
|
||||
continue
|
||||
yield chunk
|
||||
if provider_error and not answer_events and not metrics_events:
|
||||
yield provider_error
|
||||
return
|
||||
ledger = EvidenceLedger.from_tool_events(journal.evidence_events(), requirements)
|
||||
decision = ledger.evaluate(exhausted=exhausted, awaiting_user=awaiting)
|
||||
if provider_error:
|
||||
decision = replace(decision, status=CompletionStatus.FAILED,
|
||||
can_complete=False, reason='Model request failed')
|
||||
# Exhaustion limits execution; factual source synthesis can remain
|
||||
# useful and must not be replaced merely because the budget ended.
|
||||
presentation_decision = ledger.evaluate(awaiting_user=awaiting) if exhausted else decision
|
||||
presentation_decision = ledger.evaluate(awaiting_user=awaiting) if exhausted and not provider_error else decision
|
||||
safe_answer, reason = completion_answer(answer, ledger, presentation_decision)
|
||||
# Evaluate each earlier draft as well as the final replacement.
|
||||
# Never replay an unsupported intermediate success claim.
|
||||
@@ -216,7 +224,8 @@ def with_completion_gate(func):
|
||||
decision = CompletionDecision(CompletionStatus.UNVERIFIED, False, reason,
|
||||
decision.evidence_ids, decision.missing_artifacts)
|
||||
released_at = perf_counter()
|
||||
yield _event({'type': 'completion_decision', 'data': decision.to_dict()})
|
||||
if not provider_error:
|
||||
yield _event({'type': 'completion_decision', 'data': decision.to_dict()})
|
||||
replaced_answer = bool(reason or unsafe_draft or safe_answer != answer)
|
||||
if replaced_answer:
|
||||
reasoning = [event for event in answer_events if event.get('thinking') is True]
|
||||
@@ -230,6 +239,8 @@ def with_completion_gate(func):
|
||||
else:
|
||||
for event in answer_events:
|
||||
yield _event(event)
|
||||
if provider_error:
|
||||
yield _event({'type': 'completion_decision', 'data': decision.to_dict()})
|
||||
for event in metrics_events:
|
||||
metadata = event.setdefault('data', {})
|
||||
metadata.update(completion_decision=decision.to_dict(), evidence_events=ledger.to_list(),
|
||||
@@ -241,7 +252,8 @@ def with_completion_gate(func):
|
||||
'answer_replaced': replaced_answer,
|
||||
}
|
||||
if replaced_answer:
|
||||
metadata['round_texts'] = [safe_answer]
|
||||
if not provider_error:
|
||||
metadata['round_texts'] = [safe_answer]
|
||||
metadata['completion_gate_reason'] = reason or unsafe_draft or 'receipt_summary'
|
||||
if isinstance(metadata.get('thinking'), str):
|
||||
_, unsafe_thinking = completion_answer(metadata['thinking'], ledger,
|
||||
@@ -249,6 +261,9 @@ def with_completion_gate(func):
|
||||
if unsafe_thinking:
|
||||
metadata.pop('thinking')
|
||||
yield _event(event)
|
||||
if provider_error:
|
||||
yield provider_error
|
||||
return
|
||||
if done:
|
||||
yield 'data: [DONE]\n\n'
|
||||
|
||||
|
||||
@@ -89,7 +89,7 @@ def _patch_loop(monkeypatch, responses, captured_kwargs=None):
|
||||
return lambda: call_index
|
||||
|
||||
|
||||
def _run(instruction, *, max_rounds=4, relevant_tools=None, runtime_context=None):
|
||||
def _run_chunks(instruction, *, max_rounds=4, relevant_tools=None, runtime_context=None):
|
||||
async def collect():
|
||||
return [
|
||||
chunk
|
||||
@@ -104,7 +104,11 @@ def _run(instruction, *, max_rounds=4, relevant_tools=None, runtime_context=None
|
||||
)
|
||||
]
|
||||
|
||||
return _events(asyncio.run(collect()))
|
||||
return asyncio.run(collect())
|
||||
|
||||
|
||||
def _run(instruction, **kwargs):
|
||||
return _events(_run_chunks(instruction, **kwargs))
|
||||
|
||||
|
||||
def test_failed_workspace_mutation_attempts_are_not_hidden_by_successful_probe():
|
||||
@@ -497,7 +501,7 @@ def test_verified_artifact_survives_provider_error_during_finish_round(monkeypat
|
||||
|
||||
monkeypatch.setattr(agent_loop, "stream_llm_with_fallback", stream)
|
||||
|
||||
events = _run(
|
||||
chunks = _run_chunks(
|
||||
"Create /workspace/output.html",
|
||||
max_rounds=5,
|
||||
relevant_tools={"write_file", "private_browser"},
|
||||
@@ -510,12 +514,17 @@ def test_verified_artifact_survives_provider_error_during_finish_round(monkeypat
|
||||
},
|
||||
)
|
||||
|
||||
events = _events(chunks)
|
||||
assert calls == 2
|
||||
assert chunks[-1] == 'event: error\ndata: {"status": 504, "error": "stream timeout"}\n\n'
|
||||
assert not any(chunk.strip() == 'data: [DONE]' for chunk in chunks)
|
||||
decision = next(event['data'] for event in events if event.get('type') == 'completion_decision')
|
||||
assert decision['can_complete'] is False
|
||||
assert decision['status'] == 'failed'
|
||||
assert not any(event.get("type") == "agent_terminal" for event in events)
|
||||
final = next(event for event in events if event.get("type") == "final_response")
|
||||
assert "output.html" in final["content"]
|
||||
assert "Output available" in final["content"]
|
||||
assert "No passing executable test result" in final["content"]
|
||||
assert final['content'].startswith('The task is incomplete:')
|
||||
|
||||
|
||||
def test_uninspected_artifact_still_fails_on_provider_error(monkeypatch):
|
||||
@@ -540,7 +549,7 @@ def test_uninspected_artifact_still_fails_on_provider_error(monkeypatch):
|
||||
|
||||
monkeypatch.setattr(agent_loop, "stream_llm_with_fallback", stream)
|
||||
|
||||
events = _run(
|
||||
chunks = _run_chunks(
|
||||
"Create /workspace/answer.json",
|
||||
max_rounds=4,
|
||||
relevant_tools={"write_file"},
|
||||
@@ -553,7 +562,13 @@ def test_uninspected_artifact_still_fails_on_provider_error(monkeypatch):
|
||||
},
|
||||
)
|
||||
|
||||
events = _events(chunks)
|
||||
assert calls == 2
|
||||
assert chunks[-1] == 'event: error\ndata: {"status": 504, "error": "stream timeout"}\n\n'
|
||||
assert not any(chunk.strip() == 'data: [DONE]' for chunk in chunks)
|
||||
decision = next(event['data'] for event in events if event.get('type') == 'completion_decision')
|
||||
assert decision['can_complete'] is False
|
||||
assert decision['status'] == 'failed'
|
||||
terminal = next(
|
||||
(event for event in events if event.get("type") == "agent_terminal"),
|
||||
None,
|
||||
@@ -561,6 +576,7 @@ def test_uninspected_artifact_still_fails_on_provider_error(monkeypatch):
|
||||
assert terminal is not None, events
|
||||
assert terminal["data"]["failed"] is True
|
||||
assert terminal["data"]["failure"]["status"] == 504
|
||||
assert '[Agent stopped: Model request failed (HTTP 504)]' in terminal['data']['round_texts'][-1]
|
||||
|
||||
|
||||
def test_verified_artifact_gets_only_one_finish_nudge(monkeypatch):
|
||||
|
||||
@@ -0,0 +1,227 @@
|
||||
"""Provider failure is the final frame, after gated output and diagnostics."""
|
||||
import asyncio
|
||||
from inspect import signature
|
||||
import json
|
||||
|
||||
import pytest
|
||||
|
||||
from src.agent_runtime.completion import with_completion_gate
|
||||
from src.agent_runtime.journal import current_journal
|
||||
from src.tool_types import ToolBlock
|
||||
from tests.runtime_evidence_helpers import authoritative_executor
|
||||
|
||||
|
||||
ERROR = 'event: error\ndata: {"status": 504, "error": {"message": "stream timeout"}, "fallback_eligible": false}\n\n'
|
||||
DONE = 'data: [DONE]\n\n'
|
||||
|
||||
|
||||
def _event(payload):
|
||||
return 'data: ' + json.dumps(payload) + '\n\n'
|
||||
|
||||
|
||||
def _frames(chunks):
|
||||
"""Decode network chunks without losing named error frames or [DONE]."""
|
||||
pending = ''
|
||||
for chunk in chunks:
|
||||
pending += chunk
|
||||
while '\n\n' in pending:
|
||||
frame, pending = pending.split('\n\n', 1)
|
||||
lines = frame.splitlines()
|
||||
event = next((line[7:] for line in lines if line.startswith('event: ')), 'message')
|
||||
payload = '\n'.join(line[6:] for line in lines if line.startswith('data: '))
|
||||
yield event, payload if payload == '[DONE]' else json.loads(payload)
|
||||
assert not pending, 'incomplete SSE frame'
|
||||
|
||||
|
||||
def _labels(chunks):
|
||||
return [event if event != 'message' else (
|
||||
'done' if data == '[DONE]' else data.get('type', 'delta')
|
||||
) for event, data in _frames(chunks)]
|
||||
|
||||
|
||||
def _decision(chunks):
|
||||
return next(data['data'] for event, data in _frames(chunks)
|
||||
if event == 'message' and isinstance(data, dict)
|
||||
and data.get('type') == 'completion_decision')
|
||||
|
||||
|
||||
@authoritative_executor
|
||||
async def _successful_tool(block):
|
||||
return block.tool_type, {'exit_code': 0, 'output': 'OK'}
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_bare_error_preserves_original_frame_without_success_output():
|
||||
@with_completion_gate
|
||||
async def stream(messages):
|
||||
yield ERROR
|
||||
yield DONE
|
||||
|
||||
assert [chunk async for chunk in stream([])] == [ERROR]
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize('partial', ['', 'The parser checks the header first.'])
|
||||
async def test_provider_error_releases_partial_then_decision_terminal_and_original_error(partial):
|
||||
closed = []
|
||||
|
||||
@with_completion_gate
|
||||
async def stream(messages):
|
||||
try:
|
||||
yield _event({'type': 'tool_start', 'tool': 'read_file'})
|
||||
if partial:
|
||||
yield _event({'delta': partial})
|
||||
yield ERROR
|
||||
yield _event({'type': 'agent_terminal', 'data': {
|
||||
'failed': True, 'failure': {'status': 504},
|
||||
'round_texts': ['Earlier diagnostic', partial + '\n[Agent stopped]'],
|
||||
}})
|
||||
yield DONE
|
||||
finally:
|
||||
closed.append(current_journal() is not None)
|
||||
|
||||
chunks = [chunk async for chunk in stream([])]
|
||||
assert _labels(chunks) == [
|
||||
'tool_start', 'final_response', 'completion_decision', 'agent_terminal', 'error',
|
||||
], _labels(chunks)
|
||||
assert chunks[-1] == ERROR
|
||||
assert DONE not in chunks
|
||||
assert _decision(chunks)['can_complete'] is False
|
||||
assert _decision(chunks)['status'] == 'failed'
|
||||
final = next(data for event, data in _frames(chunks)
|
||||
if event == 'message' and data.get('type') == 'final_response')
|
||||
assert final['content'].startswith('The task is incomplete:')
|
||||
assert partial in final['content']
|
||||
assert closed == [True]
|
||||
assert current_journal() is None
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize('successful_tool', [False, True])
|
||||
@pytest.mark.parametrize('earlier_status', [None, 'awaiting_user', 'exhausted'])
|
||||
async def test_provider_failure_overrides_even_successful_execution(successful_tool, earlier_status):
|
||||
@with_completion_gate
|
||||
async def stream(messages):
|
||||
if successful_tool:
|
||||
await _successful_tool(ToolBlock('bash', 'python -m unittest'))
|
||||
if earlier_status:
|
||||
yield _event({'type': 'completion_decision', 'data': {'status': earlier_status}})
|
||||
yield _event({'delta': 'The response is partial.'})
|
||||
yield ERROR
|
||||
yield _event({'type': 'metrics', 'data': {}})
|
||||
|
||||
chunks = [chunk async for chunk in stream([])]
|
||||
decision = _decision(chunks)
|
||||
assert decision['can_complete'] is False, decision
|
||||
assert decision['status'] == 'failed'
|
||||
if successful_tool:
|
||||
metrics = next(data['data'] for event, data in _frames(chunks)
|
||||
if event == 'message' and data.get('type') == 'metrics')
|
||||
assert any(e['authoritative'] and e['success'] for e in metrics['evidence_events'])
|
||||
assert chunks[-1] == ERROR
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_error_after_final_response_does_not_add_calls_or_success_done():
|
||||
invocations = []
|
||||
|
||||
@with_completion_gate
|
||||
async def stream(messages, workspace=None, client_runtime_context=None):
|
||||
invocations.append(1)
|
||||
yield _event({'type': 'final_response', 'content': 'The header contains three fields.'})
|
||||
yield DONE
|
||||
yield ERROR
|
||||
|
||||
chunks = [chunk async for chunk in stream([])]
|
||||
assert _labels(chunks) == ['final_response', 'completion_decision', 'error']
|
||||
assert invocations == [1]
|
||||
assert str(signature(stream)) == '(messages, workspace=None, client_runtime_context=None)'
|
||||
assert DONE not in chunks
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize('terminal_kind', ['agent_terminal', 'metrics'])
|
||||
async def test_failed_terminal_diagnostics_survive_answer_replacement(terminal_kind):
|
||||
diagnostics = ['Earlier tool failure and retry', 'All tests passed.\n[Agent stopped: HTTP 504]']
|
||||
|
||||
@with_completion_gate
|
||||
async def stream(messages):
|
||||
yield _event({'delta': 'All tests passed.'})
|
||||
yield ERROR
|
||||
yield _event({'type': terminal_kind, 'data': {
|
||||
'failed': True, 'failure': {'status': 504, 'message': 'Model request failed'},
|
||||
'round_texts': diagnostics, 'round_models': ['first-model', 'failed-model'],
|
||||
}})
|
||||
|
||||
chunks = [chunk async for chunk in stream([])]
|
||||
terminal = next(data['data'] for event, data in _frames(chunks)
|
||||
if event == 'message' and data.get('type') == terminal_kind)
|
||||
assert terminal['round_texts'] == diagnostics
|
||||
assert terminal['round_models'] == ['first-model', 'failed-model']
|
||||
assert terminal['failure'] == {'status': 504, 'message': 'Model request failed'}
|
||||
assert terminal['failed'] is True
|
||||
assert terminal['completion_decision'] == _decision(chunks)
|
||||
assert terminal['completion_gate']['answer_replaced'] is True
|
||||
assert terminal['completion_gate']['additional_provider_calls'] == 0
|
||||
assert _labels(chunks).index(terminal_kind) < _labels(chunks).index('error')
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_error_boundary_is_independent_of_network_chunking():
|
||||
@with_completion_gate
|
||||
async def stream(messages):
|
||||
yield _event({'delta': 'Partial explanation.'})
|
||||
yield ERROR
|
||||
yield _event({'type': 'agent_terminal', 'data': {'failed': True}})
|
||||
|
||||
chunks = [chunk async for chunk in stream([])]
|
||||
wire = ''.join(chunks)
|
||||
expected = list(_frames(chunks))
|
||||
for delivered in [chunks, [wire], list(wire)]:
|
||||
# A client stops consuming on the first error, regardless of chunking.
|
||||
visible = []
|
||||
for frame in _frames(delivered):
|
||||
visible.append(frame)
|
||||
if frame[0] == 'error':
|
||||
break
|
||||
assert visible == expected
|
||||
assert visible[-2][1]['type'] == 'agent_terminal'
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize('after_error', [False, True])
|
||||
async def test_cancellation_closes_inner_stream_without_releasing_completion(after_error):
|
||||
progress_seen = asyncio.Event()
|
||||
closed = []
|
||||
chunks = []
|
||||
|
||||
@with_completion_gate
|
||||
async def stream(messages):
|
||||
try:
|
||||
yield _event({'delta': 'Tests passed.'})
|
||||
if after_error:
|
||||
yield ERROR
|
||||
yield _event({'type': 'tool_start', 'tool': 'bash'})
|
||||
await asyncio.Event().wait()
|
||||
finally:
|
||||
closed.append(current_journal() is not None)
|
||||
|
||||
async def collect():
|
||||
async for chunk in stream([]):
|
||||
chunks.append(chunk)
|
||||
if chunk == _event({'type': 'tool_start', 'tool': 'bash'}):
|
||||
progress_seen.set()
|
||||
|
||||
task = asyncio.create_task(collect())
|
||||
try:
|
||||
await asyncio.wait_for(progress_seen.wait(), timeout=5)
|
||||
task.cancel()
|
||||
with pytest.raises(asyncio.CancelledError):
|
||||
await task
|
||||
finally:
|
||||
if not task.done():
|
||||
task.cancel()
|
||||
await asyncio.gather(task, return_exceptions=True)
|
||||
assert _labels(chunks) == ['tool_start']
|
||||
assert closed == [True]
|
||||
assert current_journal() is None
|
||||
Reference in New Issue
Block a user