Stream research answer tokens while preserving final draft reconciliation

This commit is contained in:
pewdiepie-archdaemon
2026-09-17 22:10:59 +00:00
parent 4de1a4b9bb
commit f576406515
4 changed files with 56 additions and 5 deletions
@@ -81,6 +81,8 @@ async function send(page, prompt) {
const metrics = events.find(event => event.type === 'metrics')?.data;
const observation = {
prompt, seconds: (performance.now() - started) / 1000,
streamed_text_chunks: events.filter(event => typeof event.delta === 'string' && event.delta.length).length,
final_replacement_count: events.filter(event => event.type === 'final_response').length,
rounds: metrics?.agent_rounds ?? null,
tool_execution_timings: metrics?.tool_execution_timings || [],
runtime_seconds: metrics?.response_time ?? null,
+5 -4
View File
@@ -4031,7 +4031,9 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
official_source_retry_attempted = False
note_search_recovery_attempted = False
replace_streamed_draft_on_finish = False
buffer_completion_drafts = broad_current_web_request(direct_user_text) or requested_web_source_links(direct_user_text)
# Stream model text immediately. A canonical final event reconciles any
# draft that completion/research checks subsequently replace.
finalize_search_answer = broad_current_web_request(direct_user_text) or requested_web_source_links(direct_user_text)
usage_in = usage_out = 0
has_real_usage = False
first_request_tokens = last_request_tokens = 0
@@ -4253,7 +4255,6 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
if (
not prior_summary_answer
and not progressive_thinking
and not buffer_completion_drafts
):
yield event({'delta': text})
for fragment in delta.get('tool_calls') or []:
@@ -4264,7 +4265,7 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
call['function'][key] += (fragment.get('function') or {}).get(key) or ''
if progressive_thinking:
content = visible_content_after_qwen_thinking(content)
if content and not prior_summary_answer and not buffer_completion_drafts:
if content and not prior_summary_answer:
yield event({'delta': content})
proposed = [pending[i] for i in sorted(pending)]
proposed = serialize_required_email_attachment_chain(
@@ -4550,7 +4551,7 @@ async def stream_preview(*, endpoint_url, model, messages, headers, turn_contrac
yield event({'delta': suffix})
if not content:
yield event({'delta': 'The test model returned no answer. No substitute answer was generated.'})
elif replace_streamed_draft_on_finish or buffer_completion_drafts:
elif replace_streamed_draft_on_finish or finalize_search_answer:
yield event({'type': 'final_response', 'content': content})
break
# Treat a model-proposed call batch atomically for preview
+5 -1
View File
@@ -1313,10 +1313,14 @@ async def test_stream_retries_an_obviously_truncated_broad_web_answer(monkeypatc
and 'fuller evidence-based briefing' in event.get('content', '')
for event in events
)
assert not any(
assert any(
event.get('delta') == 'Current AI news includes reports about U.'
for event in events
)
# Live drafts may be visible, but the final canonical answer replaces the
# incomplete draft rather than persisting both as one answer.
final = [event['content'] for event in events if event.get('type') == 'final_response'][-1]
assert 'Current AI news includes reports about U.' not in final
def test_task_renderer_honors_few_and_filters_confirmed_morning_schedule():
+44
View File
@@ -177,6 +177,50 @@ def test_short_search_results_are_unchanged():
assert preview_tool_result_text({'output': text}, 'web_search', {}) == text
@pytest.mark.asyncio
@pytest.mark.parametrize('prompt', ['latest AI news', 'Explain these findings with sources'])
async def test_search_answer_streams_before_upstream_completion(monkeypatch, prompt):
import src.clean_agent_preview as runtime
from src.tool_policy import ToolPolicy
from src.turn_contract import resolve_full_inventory_contract
consumed = []
class Response:
async def __aenter__(self): return self
async def __aexit__(self, *args): pass
def raise_for_status(self): pass
async def aiter_lines(self):
for part in ['Supported finding. ', 'Source: https://example.org/report']:
consumed.append(part)
yield 'data: ' + json.dumps({'choices': [{'delta': {'content': part}}]})
consumed.append('DONE')
yield 'data: [DONE]'
class Client:
def __init__(self, **kwargs): pass
async def __aenter__(self): return self
async def __aexit__(self, *args): pass
def stream(self, *args, **kwargs): return Response()
monkeypatch.setattr(runtime.httpx, 'AsyncClient', Client)
contract = resolve_full_inventory_contract(schemas=[], policy=ToolPolicy())
events = []
async for chunk in runtime.stream_preview(
endpoint_url='http://test', model='test', headers={},
messages=[{'role': 'user', 'content': prompt}], turn_contract=contract,
session_id='test', owner='test', disabled_tools=set(), tool_policy=ToolPolicy(), max_rounds=1,
):
if '[DONE]' in chunk:
continue
event = json.loads(chunk[6:])
events.append(event)
if event.get('delta') == 'Supported finding. ':
assert consumed == ['Supported finding. '], 'First chunk was buffered until model completion'
assert [e['delta'] for e in events if e.get('delta')] == [
'Supported finding. ', 'Source: https://example.org/report',
]
assert [e['content'] for e in events if e.get('type') == 'final_response'] == [
'Supported finding. Source: https://example.org/report',
]
def test_failure_status_is_not_lost_to_search_compaction():
result = {'output': 'Partial evidence. ' * 1000, 'error': 'fetch failed', 'exit_code': 1}
output = preview_tool_result_text(result, 'web_search', {})