Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
336 changes: 322 additions & 14 deletions api/streaming.py

Large diffs are not rendered by default.

9 changes: 6 additions & 3 deletions static/messages.js
Original file line number Diff line number Diff line change
Expand Up @@ -1228,11 +1228,14 @@ function attachLiveStream(activeSid, streamId, uploaded=[], options={}){
const isQuotaExhausted=d.type==='quota_exhausted';
const isAuthMismatch=d.type==='auth_mismatch';
const isModelNotFound=d.type==='model_not_found';
const isCancelled=d.type==='cancelled';
const isInterrupted=d.type==='interrupted';
const isNoResponse=d.type==='no_response'||d.type==='silent_failure';
const label=isQuotaExhausted?'Out of credits':isRateLimit?'Rate limit reached':isAuthMismatch?(typeof t==='function'?t('provider_mismatch_label'):'Provider mismatch'):isModelNotFound?(typeof t==='function'?t('model_not_found_label'):'Model not found'):isNoResponse?'No response received':'Error';
const label=isCancelled?'Task cancelled':isInterrupted?'Response interrupted':isQuotaExhausted?'Out of credits':isRateLimit?'Rate limit reached':isAuthMismatch?(typeof t==='function'?t('provider_mismatch_label'):'Provider mismatch'):isModelNotFound?(typeof t==='function'?t('model_not_found_label'):'Model not found'):isNoResponse?'No response from provider':'Error';
const hint=d.hint?`\n\n*${d.hint}*`:'';
const details=d.details?String(d.details).replace(/```/g,'`\u200b``'):'';
S.messages.push({role:'assistant',content:`**${label}:** ${d.message}${hint}`,provider_details:details});
const detailsLabel=isCancelled?'Cancellation details':isInterrupted?'Interruption details':undefined;
S.messages.push({role:'assistant',content:`**${label}:** ${d.message}${hint}`,provider_details:details,provider_details_label:detailsLabel});
}catch(_){
S.messages.push({role:'assistant',content:'**Error:** An error occurred. Check server logs.'});
}
Expand Down Expand Up @@ -1323,7 +1326,7 @@ function attachLiveStream(activeSid, streamId, uploaded=[], options={}){
// Fallback to local cancel message if API fails
if(S.session&&S.session.session_id===activeSid){
clearLiveToolCards();if(!assistantText)removeThinking();
S.messages.push({role:'assistant',content:'*Task cancelled.*'});renderMessages({preserveScroll:true});
S.messages.push({role:'assistant',content:'**Task cancelled:** Task cancelled.\n\n*The run was cancelled by the user before Skyly finished. No provider failure occurred.*',provider_details:'Task cancelled.',provider_details_label:'Cancellation details',_error:true});renderMessages({preserveScroll:true});
_markSessionViewed(activeSid, S.messages.length);
}
}
Expand Down
3 changes: 2 additions & 1 deletion static/ui.js
Original file line number Diff line number Diff line change
Expand Up @@ -5013,7 +5013,8 @@ function renderMessages(options){
}
let bodyHtml = isUser ? _renderUserFencedBlocks(displayContent) : renderMd(_stripXmlToolCallsDisplay(String(displayContent)));
if(!isUser&&m.provider_details){
bodyHtml += `<details class="provider-error-details"><summary>Provider details</summary><pre><code>${esc(String(m.provider_details))}</code></pre></details>`;
const summary=m.provider_details_label||'Provider details';
bodyHtml += `<details class="provider-error-details"><summary>${esc(String(summary))}</summary><pre><code>${esc(String(m.provider_details))}</code></pre></details>`;
}
const statusHtml = (!isUser&&m._statusCard) ? _statusCardHtml(m._statusCard) : '';
const isEditableUser=isUser&&rawIdx===lastUserRawIdx;
Expand Down
185 changes: 185 additions & 0 deletions tests/test_cancelled_turn_status.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,185 @@
"""Regression tests for accurate cancelled/interrupted turn status.

A user pressing Stop/Cancel must not be shown provider-empty guidance like
"No response from provider". Provider-empty remains valid only when there was
no explicit cancel/interruption signal.
"""
from __future__ import annotations

import pathlib

from api.streaming import (
_CANCEL_MARKER_PATTERNS,
_cancelled_turn_content,
_classify_provider_error,
_finalize_cancelled_turn,
)

REPO_ROOT = pathlib.Path(__file__).parent.parent.resolve()


def _read(rel_path: str) -> str:
return (REPO_ROOT / rel_path).read_text(encoding="utf-8")


class _DummySession:
def __init__(self, path: str = ''):
self.path = path
self.messages = []
self.active_stream_id = 'stream-1'
self.pending_user_message = 'hello'
self.pending_attachments = ['a.txt']
self.pending_started_at = 123
self.saved = 0

def save(self, *args, **kwargs):
self.saved += 1


class TestCancelledTurnClassification:
def test_user_cancelled_error_is_not_provider_no_response(self):
result = _classify_provider_error("Cancelled by user", Exception("Cancelled by user"))

assert result["type"] == "cancelled"
assert result["label"] == "Task cancelled"
assert "provider returned no content" not in result.get("hint", "").lower()
assert "rate limit" not in result.get("hint", "").lower()
assert "no provider failure" in result.get("hint", "").lower()

def test_string_only_cancelled_error_repr_is_cancelled(self):
result = _classify_provider_error("<CancelledError>", None, silent_failure=True)

assert result["type"] == "cancelled"
assert result["label"] == "Task cancelled"
assert "provider returned no content" not in result.get("hint", "").lower()

def test_interrupted_or_aborted_error_is_not_provider_no_response(self):
for text in (
"Interrupted by user",
"Operation aborted before provider response completed",
"AbortError: request was aborted",
):
result = _classify_provider_error(text, RuntimeError(text))
assert result["type"] == "interrupted", text
assert result["label"] == "Response interrupted", text
assert "provider returned no content" not in result.get("hint", "").lower()

def test_provider_empty_response_still_uses_no_response(self):
result = _classify_provider_error("", None, silent_failure=True)

assert result["type"] == "no_response"
assert result["label"] == "No response from provider"
assert "provider returned no content" in result.get("hint", "").lower()


class TestCancelledTurnFinalizer:
def test_persistent_cancel_finalizer_clears_pending_and_saves_cancel_marker(self):
session = _DummySession()

_finalize_cancelled_turn(session, ephemeral=False)

assert session.active_stream_id is None
assert session.pending_user_message is None
assert session.pending_attachments == []
assert session.pending_started_at is None
assert session.saved == 1
assert session.messages[-1]['content'] == _cancelled_turn_content('Task cancelled.')
assert '**Task cancelled:** Task cancelled.' in session.messages[-1]['content']
assert 'No provider failure occurred' in session.messages[-1]['content']
assert session.messages[-1]['provider_details'] == 'Task cancelled.'
assert session.messages[-1]['provider_details_label'] == 'Cancellation details'
assert session.messages[-1]['_error'] is True

def test_ephemeral_cancel_finalizer_unlinks_temp_session_without_saving_error_marker(self, tmp_path):
temp_session = tmp_path / 'btw-session.json'
temp_session.write_text('{}', encoding='utf-8')
session = _DummySession(str(temp_session))

_finalize_cancelled_turn(session, ephemeral=True)

assert session.active_stream_id is None
assert session.pending_user_message is None
assert session.pending_attachments == []
assert session.pending_started_at is None
assert session.saved == 0
assert session.messages == []
assert not temp_session.exists()


def test_message_renderer_allows_non_provider_details_label(self):
src = _read("static/ui.js")
assert "provider_details_label||'Provider details'" in src
assert "provider-error-details" in src


class TestCancelledTurnPersistenceGuards:
def test_cancel_marker_patterns_are_centralized_for_dedupe(self):
assert _CANCEL_MARKER_PATTERNS == ('task cancelled', 'task canceled', 'response interrupted')
src = _read("api/streaming.py")
assert "any(pattern in normalized for pattern in _CANCEL_MARKER_PATTERNS)" in src
assert "any(pattern in _content for pattern in _CANCEL_MARKER_PATTERNS)" in src

def test_silent_failure_path_checks_cancel_event_before_persisting_provider_error(self):
src = _read("api/streaming.py")
silent_idx = src.find("# ── Detect silent agent failure")
assert silent_idx != -1, "silent-failure block not found"
apperror_idx = src.find("put('apperror', _error_payload)", silent_idx)
assert apperror_idx != -1, "silent-failure apperror emission not found"
block = src[silent_idx:apperror_idx]

assert "cancel_event.is_set()" in block, (
"When a user cancels and the interrupted agent returns no assistant text, "
"the silent-failure path must not persist a provider no_response error."
)
assert "cancelled" in block.lower(), (
"The cancellation guard should persist/report a cancelled turn, not silently drop state."
)

def test_exception_path_classifies_after_cancel_event_before_generic_error(self):
src = _read("api/streaming.py")
except_idx = src.find("print('[webui] stream error:")
assert except_idx != -1, "stream exception handler not found"
classify_idx = src.find("_classify_provider_error", except_idx)
generic_idx = src.find("_exc_label, _exc_type, _exc_hint = 'Error', 'error', ''", except_idx)
assert classify_idx != -1 and generic_idx != -1
block = src[except_idx:generic_idx]

assert "cancel_event.is_set()" in block, (
"Exception handling must distinguish user-cancelled/aborted runs before generic errors."
)
assert "cancelled" in block.lower() or "interrupted" in block.lower()
assert "provider_details_label" in src
assert "Cancellation details" in src
assert "Interruption details" in src

def test_post_run_cancel_guard_runs_before_normal_success_merge(self):
src = _read("api/streaming.py")
run_idx = src.find("result = agent.run_conversation(")
merge_idx = src.find("_result_messages = result.get", run_idx)
assert run_idx != -1 and merge_idx != -1, "run/merge path not found"
block = src[run_idx:merge_idx]

assert "cancel_event.is_set()" in block, (
"If cancellation arrives after tokens streamed but before run_conversation returns, "
"the worker must emit/persist cancel before normal merge/save/completed handling."
)
assert "put('cancel'" in block
assert "_cleanup_ephemeral_cancelled_turn" in block or "_finalize_cancelled_turn" in block, (
"Ephemeral cancels must clean up their temporary session before returning."
)
assert "return" in block

def test_frontend_has_cancelled_and_interrupted_labels_for_apperror_fallbacks(self):
src = _read("static/messages.js")
start = src.find("source.addEventListener('apperror'")
end = src.find("source.addEventListener('warning'", start)
assert start != -1 and end != -1, "apperror handler not found"
block = src[start:end]

assert "d.type==='cancelled'" in block or 'd.type==="cancelled"' in block
assert "d.type==='interrupted'" in block or 'd.type==="interrupted"' in block
assert "Task cancelled" in block
assert "Response interrupted" in block
assert "No response from provider" in block
assert "Cancellation details" in block
assert "Interruption details" in block
72 changes: 72 additions & 0 deletions tests/test_issue1361_cancel_data_loss.py
Original file line number Diff line number Diff line change
Expand Up @@ -441,3 +441,75 @@ def test_materialize_helper_called_immediately_before_error_path_clears():
f"found {sites_with_helper}. PR #1760 / #1361 regression — re-wire the "
f"helper at the error-branch clear sites in api/streaming.py."
)



class TestCancelStreamIdempotentWithWorkerFinalizer:
"""The worker and explicit cancel endpoint can both finalize the same turn."""

def test_cancel_stream_does_not_duplicate_existing_worker_cancel_marker(self):
sid = "test_1361_idempotent"
stream_id = "stream_idempotent"
_make_session(
session_id=sid,
messages=[
{'role': 'user', 'content': 'Help me debug this', 'timestamp': 100},
{'role': 'assistant', 'content': '**Task cancelled:** Task cancelled.\n\n*The run was cancelled by the user before Skyly finished. No provider failure occurred.*', '_error': True, 'timestamp': 101},
],
)
_setup_cancel_state(sid, stream_id)
config.STREAM_PARTIAL_TEXT[stream_id] = "partial text before cancel"

cancel_stream(stream_id)

msgs = models.SESSIONS[sid].messages
cancel_markers = [
m for m in msgs
if isinstance(m, dict)
and m.get('role') == 'assistant'
and 'task cancelled' in str(m.get('content') or '').lower()
]
partial_idx = next(
i for i, m in enumerate(msgs)
if isinstance(m, dict) and m.get('_partial') and m.get('content') == 'partial text before cancel'
)
marker_idx = next(i for i, m in enumerate(msgs) if m in cancel_markers)

assert len(cancel_markers) == 1
assert partial_idx < marker_idx

def test_late_cancel_after_worker_finalized_does_not_add_cancel_marker(self):
sid = "test_1361_late_done"
stream_id = "stream_late_done"
s = Session(
session_id=sid,
title="Done Session",
messages=[
{'role': 'user', 'content': 'finish normally', 'timestamp': 100},
{'role': 'assistant', 'content': 'done normally', 'timestamp': 101},
],
)
s.active_stream_id = None
s.pending_user_message = None
s.pending_attachments = []
s.pending_started_at = None
s.save()
models.SESSIONS[sid] = s

q = queue.Queue()
config.STREAMS[stream_id] = q
config.CANCEL_FLAGS[stream_id] = threading.Event()
mock_agent = Mock()
mock_agent.session_id = sid
mock_agent.interrupt = Mock()
config.AGENT_INSTANCES[stream_id] = mock_agent
config.STREAM_PARTIAL_TEXT[stream_id] = 'stale partial snapshot'

assert cancel_stream(stream_id) is True

msgs = models.SESSIONS[sid].messages
assert msgs == [
{'role': 'user', 'content': 'finish normally', 'timestamp': 100},
{'role': 'assistant', 'content': 'done normally', 'timestamp': 101},
]
assert q.empty(), "late cancel must not emit a terminal cancel event after done"
9 changes: 5 additions & 4 deletions tests/test_issue893_cancel_preserves_partial.py
Original file line number Diff line number Diff line change
Expand Up @@ -3,7 +3,7 @@
assistant content rather than discarding it.

Before this fix, clicking Stop Generation threw away all streamed text. The
session was saved with only '*Task cancelled.*' appended, so the user lost
session was saved with only a cancellation marker appended, so the user lost
whatever the agent had produced up to that point.

After this fix:
Expand Down Expand Up @@ -118,7 +118,7 @@ def interrupt(self, _): pass
assert any('Python is a high-level programming language' in c for c in msg_contents), (
f"Partial text not found in session messages: {msg_contents}"
)
assert any('*Task cancelled.*' in c for c in msg_contents), (
assert any('Task cancelled:' in c for c in msg_contents), (
"Cancel marker missing from session messages"
)
# Partial message should NOT have _error=True (it's real content)
Expand All @@ -127,8 +127,9 @@ def interrupt(self, _): pass
assert partial_msg.get('_partial') is True
assert not partial_msg.get('_error')
# Cancel marker should have _error=True
cancel_msg = next(m for m in saved.messages if '*Task cancelled.*' in m.get('content', ''))
cancel_msg = next(m for m in saved.messages if 'Task cancelled:' in m.get('content', ''))
assert cancel_msg.get('_error') is True
assert cancel_msg.get('provider_details_label') == 'Cancellation details'

def test_cancel_stream_with_no_partial_text_still_saves_cancel_marker(self, tmp_path, monkeypatch):
"""If no tokens were streamed before cancel, only the cancel marker is saved."""
Expand Down Expand Up @@ -168,7 +169,7 @@ def interrupt(self, _): pass

saved = Session.load('sess_nopartial')
msg_contents = [m.get('content', '') for m in saved.messages]
assert any('*Task cancelled.*' in c for c in msg_contents)
assert any('Task cancelled:' in c for c in msg_contents)
# No extra partial message when there was nothing streamed
assert not any(m.get('_partial') for m in saved.messages), (
"Should not add partial message when no tokens were streamed"
Expand Down
7 changes: 4 additions & 3 deletions tests/test_pr1341_context_window_persistence.py
Original file line number Diff line number Diff line change
Expand Up @@ -38,11 +38,12 @@ def test_streaming_persists_context_fields_on_session_before_save():
# Save call follows shortly after
save_call = src.find("\n s.save()", block_start)
assert save_call != -1, "s.save() not found after the post-merge marker"
# Limit bumped to 8200 by turn-journal lifecycle events: the block now also
# records `assistant_started` immediately before the durable final save.
# Limit bumped to 9000 by cancellation finalization guards: the block now also
# checks for a late user cancel immediately before the durable final save,
# preventing a race that would otherwise save/emit a completed turn after Stop.
# The context_length fallback is still a single focused resolver call with
# arg-prep scaffold and commentary explaining the failure mode it prevents.
assert save_call - block_start < 8200, (
assert save_call - block_start < 9000, (
"s.save() should be close to the post-merge marker — block expanded unexpectedly. "
"If you've added a new pre-save mutation block here, bump this limit."
)
Expand Down
8 changes: 3 additions & 5 deletions tests/test_sprint36.py
Original file line number Diff line number Diff line change
Expand Up @@ -212,16 +212,14 @@ def test_cancel_marker_flagged_as_error_to_skip_in_api_history():
_error: True so _sanitize_messages_for_api() strips it from the
conversation_history sent to the agent on the next user message.

Without this flag, the LLM sees "*Task cancelled.*" as a prior assistant
Without this flag, the LLM sees "Task cancelled" as a prior assistant
turn and may reference it in subsequent responses ("As I mentioned, I was
cancelled...") — a behavioral regression introduced when this PR started
persisting the marker to the session.
"""
src = read("api/streaming.py")
idx = src.find("'content': '*Task cancelled.*'")
if idx == -1:
idx = src.find('"content": "*Task cancelled.*"')
assert idx != -1, "cancel marker content string not found in cancel_stream()"
idx = src.find("'content': _cancelled_turn_content(message)")
assert idx != -1, "cancel marker content writer not found in cancel_stream()"

# Walk back to the start of the dict literal (opening brace)
brace_open = src.rfind("{", 0, idx)
Expand Down
Loading