Files
odysseus/tests/test_headless_completion_consumers.py
T
Alexandre Teixeira d49071bbec fix: close Wave 1.1 completion-gate audit findings
- Headless consumers (task scheduler, background follow-up) now treat a
  completion-gate final_response as the authoritative answer instead of
  collecting deltas only. A gated replacement no longer leaves scheduled
  output empty, which used to trigger an extra, ungated grace-summary
  model call.
- The scheduler closes the agent stream with contextlib.aclosing, so the
  approval-pause break unwinds the gate's journal and teacher-takeover
  context in its own task. Chained runs no longer inherit a stale
  parent_run_id, and later finalization no longer raises ContextVar
  reset errors.
- On provider error, the completion gate applies the live answer's
  statement filter to persisted round_texts. Diagnostics and the failure
  note survive; claims rejected by the gate cannot reappear on reload.
2026-10-01 14:49:32 +01:00

159 lines
6.0 KiB
Python

"""Headless consumers present the completion gate's answer and close its stream."""
import asyncio
import json
import sys
import types
from types import SimpleNamespace
import pytest
from src.agent_runtime.completion import with_completion_gate
from src.agent_runtime.journal import current_journal
from src.teacher_escalation import with_teacher_takeover
CLAIM = 'I created report.md and all tests passed.'
def _event(payload):
return 'data: ' + json.dumps(payload) + '\n\n'
def _task():
return SimpleNamespace(
crew_member_id=None, endpoint_url='http://ep/v1', model='m',
session_id='s', owner='admin', prompt='create report.md and run the tests',
name='job', max_steps=5, character_id=None,
)
def _gated_loop(released):
"""Real gate and takeover adapters around a loop that over-claims."""
@with_teacher_takeover
@with_completion_gate
async def stream_agent_loop(*args, messages=None, client_runtime_context=None, **kwargs):
yield _event({'delta': 'Inspected the layout. ' + CLAIM})
yield _event({'type': 'metrics', 'data': {}})
yield 'data: [DONE]\n\n'
async def recording(*args, **kwargs):
async for chunk in stream_agent_loop(*args, **kwargs):
if chunk.startswith('data: {') and '"final_response"' in chunk:
released.append(json.loads(chunk[6:])['content'])
yield chunk
return recording
async def test_scheduler_result_is_the_gated_replacement_without_grace_call(monkeypatch):
from src.task_scheduler import TaskScheduler
released = []
grace_calls = []
async def grace(*args, **kwargs):
grace_calls.append(kwargs)
return 'ungated summary: all tests passed'
monkeypatch.setattr('src.agent_loop.stream_agent_loop', _gated_loop(released))
monkeypatch.setattr('src.task_endpoint.resolve_task_candidates', lambda **kwargs: [])
monkeypatch.setattr('src.task_endpoint.task_llm_call_async', grace)
result = await TaskScheduler(session_manager=None)._run_agent_loop(
'http://ep/v1', 'model', _task(), 's')
assert len(released) == 1
assert result == released[0].strip()
assert 'Inspected the layout.' in result
assert 'tests passed' not in result
assert grace_calls == []
def test_background_followup_prose_is_the_gated_replacement(monkeypatch):
from src import bg_monitor
released = []
agent_loop = types.ModuleType('src.agent_loop')
agent_loop.stream_agent_loop = _gated_loop(released)
monkeypatch.setitem(sys.modules, 'src.agent_loop', agent_loop)
sess = SimpleNamespace(endpoint_url='http://example.test', model='model',
headers=None, context_length=0, id='s1', owner='owner')
full, _ = asyncio.run(bg_monitor._drain_agent(
sess, [{'role': 'user', 'content': 'create report.md and run the tests'}]))
assert len(released) == 1
assert full == released[0]
assert 'tests passed' not in full
@pytest.mark.parametrize('consumer', ['scheduler', 'background'])
def test_later_answer_supersedes_earlier_replacement(monkeypatch, consumer):
async def stream_agent_loop(*args, **kwargs):
yield _event({'type': 'final_response', 'content': 'Earlier summary.'})
yield _event({'delta': 'Final '})
yield _event({'delta': 'answer.'})
yield 'data: [DONE]\n\n'
if consumer == 'scheduler':
from src.task_scheduler import TaskScheduler
monkeypatch.setattr('src.agent_loop.stream_agent_loop', stream_agent_loop)
monkeypatch.setattr('src.task_endpoint.resolve_task_candidates', lambda **kwargs: [])
result = asyncio.run(TaskScheduler(session_manager=None)._run_agent_loop(
'http://ep/v1', 'model', _task(), 's'))
else:
from src import bg_monitor
agent_loop = types.ModuleType('src.agent_loop')
agent_loop.stream_agent_loop = stream_agent_loop
monkeypatch.setitem(sys.modules, 'src.agent_loop', agent_loop)
sess = SimpleNamespace(endpoint_url='http://example.test', model='model',
headers=None, context_length=0, id='s1')
result, _ = asyncio.run(bg_monitor._drain_agent(sess, []))
assert result == 'Final answer.'
async def test_scheduler_approval_pause_closes_gated_stream_in_its_own_context(monkeypatch):
from src.task_scheduler import TaskScheduler
closed = []
lineage = []
@with_teacher_takeover
@with_completion_gate
async def paused_loop(*args, messages=None, client_runtime_context=None, **kwargs):
try:
yield _event({'type': 'tool_output', 'tool': 'bash', 'output': 'Waiting for an exact user approval.',
'ask_user': {'kind': 'tool_approval', 'approval_id': 'missing'}})
yield _event({'delta': 'not reached'})
finally:
closed.append(current_journal() is not None)
@with_teacher_takeover
@with_completion_gate
async def later_loop(*args, messages=None, client_runtime_context=None, **kwargs):
yield _event({'delta': 'Later run.'})
yield _event({'type': 'metrics', 'data': {}})
yield 'data: [DONE]\n\n'
async def later_run():
async for chunk in later_loop(messages=[{'role': 'user', 'content': 'x'}]):
if chunk.startswith('data: {') and '"metrics"' in chunk:
lineage.append(json.loads(chunk[6:])['data']['parent_run_id'])
monkeypatch.setattr('src.agent_loop.stream_agent_loop', paused_loop)
monkeypatch.setattr('src.task_endpoint.resolve_task_candidates', lambda **kwargs: [])
result = await TaskScheduler(session_manager=None)._run_agent_loop(
'http://ep/v1', 'model', _task(), 's')
assert 'paused safely' in result
# Closed during the pause, while its own journal was still bound.
assert closed == [True]
assert current_journal() is None
# Neither a chained task (which copies this context) nor a later run in
# this task inherits the paused run's journal as its parent.
chained = asyncio.create_task(later_run())
await chained
await later_run()
assert lineage == [None, None]