From e3965b23637a2accfa4af6c7dc261c4dfc8b5129 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:28:07 -0300 Subject: [PATCH 01/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#79:=20Preserve=20part?= =?UTF-8?q?ial=20stream=20text=20on=20mid-response=20timeout?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/providers/generic.py | 164 +++++++++++++++++--------------- 1 file changed, 86 insertions(+), 78 deletions(-) diff --git a/pycodeloop/providers/generic.py b/pycodeloop/providers/generic.py index d2b3531..3b09f0f 100644 --- a/pycodeloop/providers/generic.py +++ b/pycodeloop/providers/generic.py @@ -137,7 +137,7 @@ class GenericProvider(Provider): "model": "my-model", "api_key_env": "MY_API_KEY", "headers": {"X-Custom": "value"}, - "timeout": 60, + "timeout": 180, "context_window": 4096, "response_paths": { "text": "choices.0.message.content", @@ -191,7 +191,7 @@ def __init__( auth_prefix: str = "Bearer ", request_builder: RequestBuilder | None = None, response_parser: ResponseParser | None = None, - timeout: float = 60.0, + timeout: float = 180.0, repetition_min_period: int = _REPETITION_MIN_PERIOD, repetition_max_period: int = _REPETITION_MAX_PERIOD, repetition_repeats: int = _REPETITION_REPEATS, @@ -258,7 +258,7 @@ def _build_from_json(cls, path: str | Path) -> GenericProvider: auth_prefix=data.get("auth_prefix", "Bearer "), request_builder=request_builder, response_parser=response_parser, - timeout=data.get("timeout", 60.0), + timeout=data.get("timeout", 180.0), context_window=data.get("context_window"), supports_openai_sse=response_shape != "anthropic", include_usage_in_stream=data.get("include_usage_in_stream", True), @@ -406,82 +406,90 @@ def _stream( saw_terminal_marker = False usage = Usage() - with self._open(body, config) as response: - for raw_line in response: - if cancel_event is not None and cancel_event.is_set(): - stop_reason = "cancelled" - saw_terminal_marker = True - break - line = raw_line.decode().strip() - if not line or not line.startswith("data: "): - continue - payload = line[len("data: ") :] - if payload == "[DONE]": - saw_terminal_marker = True - break - try: - chunk = json.loads(payload) - except json.JSONDecodeError: - if stop_reason is None: - stop_reason = "malformed_stream" - break - - if chunk.get("usage"): - usage = Usage( - input_tokens=chunk["usage"].get("prompt_tokens", 0), - output_tokens=chunk["usage"].get( - "completion_tokens", 0 - ), - ) - - choices = chunk.get("choices") or [] - if not choices: - continue - choice = choices[0] - delta = choice.get("delta") or {} - - if delta.get("content"): - candidate = text + delta["content"] - if _is_repeating( - candidate, - self.repetition_min_period, - self.repetition_max_period, - self.repetition_repeats, - ): - stop_reason = "repetition" + try: + with self._open(body, config) as response: + for raw_line in response: + if cancel_event is not None and cancel_event.is_set(): + stop_reason = "cancelled" + saw_terminal_marker = True + break + line = raw_line.decode().strip() + if not line or not line.startswith("data: "): + continue + payload = line[len("data: ") :] + if payload == "[DONE]": + saw_terminal_marker = True break - text = candidate - on_delta(delta["content"]) - - for tc in delta.get("tool_calls") or []: - index = tc.get("index", 0) - acc = pending.setdefault( - index, - { - "id": None, - "name": None, - "arguments": "", - "extra": {}, - }, - ) - if tc.get("id"): - acc["id"] = tc["id"] - function = tc.get("function") or {} - if function.get("name"): - acc["name"] = function["name"] - if function.get("arguments"): - acc["arguments"] += function["arguments"] - acc["extra"].update( - { - k: v - for k, v in tc.items() - if k not in ("index", "id", "type", "function") - } - ) - - if choice.get("finish_reason"): - stop_reason = choice["finish_reason"] - saw_terminal_marker = True + try: + chunk = json.loads(payload) + except json.JSONDecodeError: + if stop_reason is None: + stop_reason = "malformed_stream" + break + + if chunk.get("usage"): + usage = Usage( + input_tokens=chunk["usage"].get( + "prompt_tokens", 0 + ), + output_tokens=chunk["usage"].get( + "completion_tokens", 0 + ), + ) + + choices = chunk.get("choices") or [] + if not choices: + continue + choice = choices[0] + delta = choice.get("delta") or {} + + if delta.get("content"): + candidate = text + delta["content"] + if _is_repeating( + candidate, + self.repetition_min_period, + self.repetition_max_period, + self.repetition_repeats, + ): + stop_reason = "repetition" + break + text = candidate + on_delta(delta["content"]) + + for tc in delta.get("tool_calls") or []: + index = tc.get("index", 0) + acc = pending.setdefault( + index, + { + "id": None, + "name": None, + "arguments": "", + "extra": {}, + }, + ) + if tc.get("id"): + acc["id"] = tc["id"] + function = tc.get("function") or {} + if function.get("name"): + acc["name"] = function["name"] + if function.get("arguments"): + acc["arguments"] += function["arguments"] + acc["extra"].update( + { + k: v + for k, v in tc.items() + if k not in ("index", "id", "type", "function") + } + ) + + if choice.get("finish_reason"): + stop_reason = choice["finish_reason"] + saw_terminal_marker = True + except (TimeoutError, ConnectionError, urllib.error.URLError): + if not text and not pending: + raise + stop_reason = "connection_lost" + saw_terminal_marker = False if stop_reason is None: stop_reason = "stop" if saw_terminal_marker else "connection_lost" From 9c680cf93c4411c87bd4e3b7afa7ebacb73ad28d Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:28:07 -0300 Subject: [PATCH 02/10] =?UTF-8?q?=E2=9D=A4=EF=B8=8F=20TEST-#79:=20Cover=20?= =?UTF-8?q?mid-stream=20timeout=20preserving/reraising=20behavior?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/providers/test_generic.py | 74 +++++++++++++++++++++++++++++++++ 1 file changed, 74 insertions(+) diff --git a/tests/providers/test_generic.py b/tests/providers/test_generic.py index 4fbd9dc..877d7f9 100644 --- a/tests/providers/test_generic.py +++ b/tests/providers/test_generic.py @@ -21,6 +21,37 @@ def __exit__(self, *exc): return False +class _TimeoutAfterLinesResponse(io.BytesIO): + """Like `_FakeResponse`, but raises `TimeoutError` once iteration + passes `raise_after` lines — simulates the model going quiet + mid-stream (long reasoning) and the socket's read timeout firing + before any terminal marker or `finish_reason` arrives.""" + + def __init__(self, data: bytes, raise_after: int): + super().__init__(data) + self._raise_after = raise_after + self._yielded = 0 + + def __enter__(self): + return self + + def __exit__(self, *exc): + self.close() + return False + + def __iter__(self): + return self + + def __next__(self): + if self._yielded >= self._raise_after: + raise TimeoutError("timed out") + line = super().readline() + if not line: + raise StopIteration + self._yielded += 1 + return line + + class GenericProviderTestCase(unittest.TestCase): def setUp(self): tmpdir = tempfile.TemporaryDirectory() @@ -356,6 +387,49 @@ def test_streaming_flags_a_connection_dropped_mid_response(self): self.assertEqual(result.text, "cut off mid") self.assertEqual(result.stop_reason, "connection_lost") + def test_streaming_returns_partial_text_on_mid_stream_timeout(self): + """Regression: a `TimeoutError` raised mid-read (model silent for + longer than the socket timeout while "thinking") used to + propagate out of `_stream()` uncaught, discarding whatever text + had already streamed in and losing the assistant's turn + entirely instead of returning it as a partial response.""" + path = self._write_config( + {"url": "http://fake/v1/chat/completions", "model": "my-model"} + ) + provider = GenericProvider.from_json(path) + + chunks = [ + {"choices": [{"delta": {"content": "partial "}}]}, + {"choices": [{"delta": {"content": "answer"}}]}, + ] + sse_body = "".join(f"data: {json.dumps(c)}\n" for c in chunks).encode() + + with mock.patch( + "pycodeloop.providers.generic.urllib.request.urlopen", + return_value=_TimeoutAfterLinesResponse(sse_body, raise_after=1), + ): + result = provider.complete("sys", [], [], on_delta=lambda _: None) + + self.assertEqual(result.text, "partial ") + self.assertEqual(result.stop_reason, "connection_lost") + + def test_streaming_reraises_timeout_when_nothing_was_streamed_yet(self): + """A timeout before any content or tool-call delta arrived means + nothing was generated to preserve — the exception should still + propagate so `Agent._complete()`'s existing retry logic kicks + in, instead of being swallowed into an empty response.""" + path = self._write_config( + {"url": "http://fake/v1/chat/completions", "model": "my-model"} + ) + provider = GenericProvider.from_json(path) + + with mock.patch( + "pycodeloop.providers.generic.urllib.request.urlopen", + return_value=_TimeoutAfterLinesResponse(b"", raise_after=0), + ): + with self.assertRaises(TimeoutError): + provider.complete("sys", [], [], on_delta=lambda _: None) + def test_streaming_stops_promptly_when_cancel_event_is_set(self): """Regression: cancel_event was accepted nowhere in the streaming read loop, so pressing Esc/Cancel mid-response did nothing until From 2390ae9ddb18d4201a6d9c10c2c229e2bdfdfd82 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:28:17 -0300 Subject: [PATCH 03/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#79:=20Bump=20example?= =?UTF-8?q?=20provider=20timeout=20to=20match=20new=20180s=20default?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- docs/examples/provider.example.json | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/docs/examples/provider.example.json b/docs/examples/provider.example.json index bf326ee..64464d7 100644 --- a/docs/examples/provider.example.json +++ b/docs/examples/provider.example.json @@ -5,7 +5,7 @@ "headers": { "X-Org": "acme" }, - "timeout": 60, + "timeout": 180, "_comment": "response_paths is optional. Omit it entirely if the API already returns the OpenAI chat-completions shape (choices[0].message.content, usage.prompt_tokens, ...). Keep it only to remap a different shape, like the example below.", "response_paths": { From 01a3c6258ae4fe84ccd50bc00c58ffffa9383de0 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:41:51 -0300 Subject: [PATCH 04/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#80:=20Widen=20stream?= =?UTF-8?q?=20except=20clause=20to=20preserve=20partial=20output=20on=20an?= =?UTF-8?q?y=20mid-read=20failure?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/providers/generic.py | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/pycodeloop/providers/generic.py b/pycodeloop/providers/generic.py index 3b09f0f..77466f0 100644 --- a/pycodeloop/providers/generic.py +++ b/pycodeloop/providers/generic.py @@ -485,7 +485,7 @@ def _stream( if choice.get("finish_reason"): stop_reason = choice["finish_reason"] saw_terminal_marker = True - except (TimeoutError, ConnectionError, urllib.error.URLError): + except Exception: if not text and not pending: raise stop_reason = "connection_lost" From d25165404319e16ea9d6d980720dae666b45ff5a Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:41:51 -0300 Subject: [PATCH 05/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#80:=20Flush=20streame?= =?UTF-8?q?d=20text=20buffer=20before=20showing=20a=20turn=20error?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/cli/chat.py | 14 ++++++++++++++ 1 file changed, 14 insertions(+) diff --git a/pycodeloop/cli/chat.py b/pycodeloop/cli/chat.py index da323d4..d5c313b 100644 --- a/pycodeloop/cli/chat.py +++ b/pycodeloop/cli/chat.py @@ -431,6 +431,20 @@ async def _run_turn( ) except Exception as exc: self.call_from_thread(self._stop_thinking) + if self._text_buffer.strip(): + self.call_from_thread( + self._log, + Panel( + Markdown(self._text_buffer), + border_style="grey50", + title="[dim]interrupted[/dim]", + subtitle=( + "[bold white on grey30] Agent [/bold white on grey30]" + ), + subtitle_align="right", + ), + ) + self._text_buffer = "" self.call_from_thread( self._log, self._styled("[bold white]✗ Error:[/bold white] ", str(exc)), From 16a338656a7fd4e3b4b04072a7a3e8332e8e5848 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:41:51 -0300 Subject: [PATCH 06/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#80:=20Harden=20serve?= =?UTF-8?q?=20mode=20against=20broken=20pipes,=20silent=20drops=20and=20st?= =?UTF-8?q?uck=20confirms?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/cli/serve.py | 53 +++++++++++++++++++++++++++++++++++++---- 1 file changed, 48 insertions(+), 5 deletions(-) diff --git a/pycodeloop/cli/serve.py b/pycodeloop/cli/serve.py index 82efbd4..d19bf2f 100644 --- a/pycodeloop/cli/serve.py +++ b/pycodeloop/cli/serve.py @@ -9,6 +9,7 @@ import queue import sys import threading +import time import uuid import typer @@ -31,6 +32,9 @@ response, ) +_HEARTBEAT_INTERVAL = 15.0 +_CONFIRM_TIMEOUT = 120.0 + class RpcServer: """Wires `Agent` callbacks to JSON-RPC notifications instead of the @@ -52,13 +56,28 @@ def __init__( self._confirm_waiters: dict[str, queue.Queue] = {} self._cancel_event: threading.Event | None = None self._chat_thread: threading.Thread | None = None + self._disconnected = False self._wire_callbacks() def _send(self, message: dict) -> None: + """Best-effort write of one NDJSON line to stdout. The client + (editor extension) can disconnect mid-turn — closing its end of + the pipe — at any point, including while a background thread is + still streaming deltas for an in-flight turn. Once that happens + every further write raises the same broken-pipe error, so this + marks the server disconnected and gives up quietly instead of + raising out of a callback (which would otherwise abort whatever + turn/tool loop is in progress) or crashing a second time from + inside an error handler that itself calls `_send`.""" + if self._disconnected: + return line = json.dumps(message) - with self._out_lock: - sys.stdout.write(line + "\n") - sys.stdout.flush() + try: + with self._out_lock: + sys.stdout.write(line + "\n") + sys.stdout.flush() + except OSError: + self._disconnected = True def _notify(self, method: str, params: dict) -> None: self._send(notification(method, params)) @@ -137,7 +156,10 @@ def confirm(name: str, preview: str) -> bool | str: {"id": request_id, "name": name, "preview": preview}, ) try: - return answer_queue.get() + return answer_queue.get(timeout=_CONFIRM_TIMEOUT) + except queue.Empty: + self._notify("chat/confirmTimeout", {"id": request_id}) + return False finally: self._confirm_waiters.pop(request_id, None) @@ -153,8 +175,23 @@ def confirm(name: str, preview: str) -> bool | str: agent.on_compact_end = on_compact_end agent.confirm = confirm + def _run_heartbeat(self, stop: threading.Event) -> None: + """Emits `chat/heartbeat` every `_HEARTBEAT_INTERVAL` seconds + while a turn is in flight. Long reasoning or a long-running + tool can otherwise leave the client with no message at all for + minutes; a client with its own read timeout may then conclude + the process died and drop the connection, losing the turn even + though the server was still working on it.""" + while not stop.wait(_HEARTBEAT_INTERVAL): + self._notify("chat/heartbeat", {}) + def _run_chat(self, request_id, params: dict) -> None: self._cancel_event = threading.Event() + heartbeat_stop = threading.Event() + heartbeat_thread = threading.Thread( + target=self._run_heartbeat, args=(heartbeat_stop,), daemon=True + ) + heartbeat_thread.start() try: result = self.flow.run( params.get("prompt", ""), @@ -165,6 +202,8 @@ def _run_chat(self, request_id, params: dict) -> None: self._respond(request_id, {"text": result}) except Exception as exc: self._respond_error(request_id, SERVER_ERROR, str(exc)) + finally: + heartbeat_stop.set() def _run_ask(self, request_id, params: dict) -> None: try: @@ -259,7 +298,11 @@ def serve_forever(self) -> None: continue try: request = json.loads(line) - except json.JSONDecodeError: + except json.JSONDecodeError as exc: + console.print( + f"[dim]⚠ dropped malformed request line ({exc}): " + f"{line[:200]!r}[/dim]" + ) continue self.handle(request) From cf125ec9c90cdafc19d757917c267acdf14261d6 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:41:51 -0300 Subject: [PATCH 07/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#80:=20Isolate=20Agent?= =?UTF-8?q?=20callback=20failures=20from=20the=20turn/tool=20loop?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/core/agent.py | 111 +++++++++++++++++++++++---------------- 1 file changed, 67 insertions(+), 44 deletions(-) diff --git a/pycodeloop/core/agent.py b/pycodeloop/core/agent.py index 546eb7f..cba5939 100644 --- a/pycodeloop/core/agent.py +++ b/pycodeloop/core/agent.py @@ -116,7 +116,29 @@ def __init__( def _trace(self, event_type: str, **fields) -> None: if self.on_trace_event: - self.on_trace_event({"type": event_type, **fields}) + try: + self.on_trace_event({"type": event_type, **fields}) + except Exception: + pass + + def _safe_call(self, callback: Callable | None, *args) -> None: + """Invokes a consumer-supplied `on_*` callback (UI rendering, + storage persistence, etc.) without letting a bug on that side + abort the turn/tool loop still in progress. Before this, an + exception from e.g. `on_message` (a storage write failing) or + `on_tool_result` (a rendering bug) propagated straight out of + `Agent.run()`/`_run_tool_calls()`, killing the rest of the turn + over what should have been a self-contained side effect.""" + if callback is None: + return + try: + callback(*args) + except Exception as exc: + self._trace( + "callback_error", + callback=getattr(callback, "__name__", repr(callback)), + error=str(exc), + ) def _complete(self, **kwargs) -> ProviderResponse: """`provider.complete()` with retry + exponential backoff on @@ -155,8 +177,7 @@ def _complete(self, **kwargs) -> ProviderResponse: delay=delay, error=str(exc), ) - if self.on_retry: - self.on_retry(attempt + 1, delay, exc) + self._safe_call(self.on_retry, attempt + 1, delay, exc) time.sleep(delay) delay *= 2 @@ -172,8 +193,7 @@ def _notify_message(self) -> None: caller persist incrementally instead of only after the whole (possibly long, multi-tool-call) turn finishes, so a crash mid-turn doesn't lose everything already done in it.""" - if self.on_message: - self.on_message() + self._safe_call(self.on_message) def _tool_schemas(self) -> list[dict]: return [tool.schema() for tool in self.tools.values()] @@ -252,8 +272,7 @@ def _run_tool_calls( ) -> None: for call in calls: self._trace("tool_call", name=call.name, arguments=call.arguments) - if self.on_tool_call: - self.on_tool_call(call.name, call.arguments) + self._safe_call(self.on_tool_call, call.name, call.arguments) if cancel_event and cancel_event.is_set(): results = {call.id: ("Cancelled by user.", True) for call in calls} @@ -302,8 +321,9 @@ def _run_tool_calls( is_error=is_error, result_len=len(result_text), ) - if self.on_tool_result: - self.on_tool_result(call.name, result_text, is_error) + self._safe_call( + self.on_tool_result, call.name, result_text, is_error + ) session.add_tool_result(call.id, result_text) self._notify_message() @@ -332,8 +352,7 @@ def _compact(self, session: Session) -> None: if len(turn_starts) <= _COMPACT_KEEP_RECENT_TURNS: return - if self.on_compact_start: - self.on_compact_start() + self._safe_call(self.on_compact_start) before_count = len(history) cutoff = turn_starts[-_COMPACT_KEEP_RECENT_TURNS] @@ -365,8 +384,9 @@ def _compact(self, session: Session) -> None: self._trace( "compact", before=before_count, after=len(session.messages) ) - if self.on_compact_end: - self.on_compact_end(before_count, len(session.messages)) + self._safe_call( + self.on_compact_end, before_count, len(session.messages) + ) def run( self, @@ -403,8 +423,9 @@ def run( self._compact(session) tools = self._tool_schemas() - if self.on_request: - self.on_request(len(session.history()), len(tools)) + self._safe_call( + self.on_request, len(session.history()), len(tools) + ) started_at = time.perf_counter() response = self._complete( @@ -418,13 +439,13 @@ def run( if response.stop_reason == "cancelled": self.usage = self.usage + response.usage - if self.on_usage: - self.on_usage(response.usage, self.usage, elapsed) + self._safe_call( + self.on_usage, response.usage, self.usage, elapsed + ) if response.text.strip(): session.add_assistant(response.text) self._notify_message() - if self.on_turn_end: - self.on_turn_end() + self._safe_call(self.on_turn_end) self._trace("run_end", reason="cancelled") return "Cancelled by user." @@ -435,8 +456,9 @@ def run( and empty_retries < _MAX_EMPTY_RESPONSE_RETRIES ): self.usage = self.usage + response.usage - if self.on_usage: - self.on_usage(response.usage, self.usage, elapsed) + self._safe_call( + self.on_usage, response.usage, self.usage, elapsed + ) empty_retries += 1 self._trace( @@ -444,15 +466,15 @@ def run( model=self.provider.model, attempt=empty_retries, ) - if self.on_retry: - self.on_retry( - empty_retries, - 0.0, - RuntimeError( - f"{self.provider.model} returned an empty " - "response with no tool calls" - ), - ) + self._safe_call( + self.on_retry, + empty_retries, + 0.0, + RuntimeError( + f"{self.provider.model} returned an empty " + "response with no tool calls" + ), + ) started_at = time.perf_counter() response = self._complete( system_prompt=self.system_prompt, @@ -465,20 +487,21 @@ def run( if response.stop_reason == "cancelled": self.usage = self.usage + response.usage - if self.on_usage: - self.on_usage(response.usage, self.usage, elapsed) + self._safe_call( + self.on_usage, response.usage, self.usage, elapsed + ) if response.text.strip(): session.add_assistant(response.text) self._notify_message() - if self.on_turn_end: - self.on_turn_end() + self._safe_call(self.on_turn_end) self._trace("run_end", reason="cancelled") return "Cancelled by user." if not response.text.strip() and not response.tool_calls: self.usage = self.usage + response.usage - if self.on_usage: - self.on_usage(response.usage, self.usage, elapsed) + self._safe_call( + self.on_usage, response.usage, self.usage, elapsed + ) error_text = ( f"{self.provider.model} returned an empty response " @@ -489,14 +512,14 @@ def run( ) session.add_assistant(error_text) self._notify_message() - if self.on_turn_end: - self.on_turn_end() + self._safe_call(self.on_turn_end) self._trace("run_end", reason="empty_response") return error_text self.usage = self.usage + response.usage - if self.on_usage: - self.on_usage(response.usage, self.usage, elapsed) + self._safe_call( + self.on_usage, response.usage, self.usage, elapsed + ) self._trace( "turn", @@ -509,8 +532,9 @@ def run( ) session.update_last_context_tokens(response.usage.input_tokens) - if self.on_context: - self.on_context(response.usage.input_tokens, context_window) + self._safe_call( + self.on_context, response.usage.input_tokens, context_window + ) tool_calls = [ { @@ -523,8 +547,7 @@ def run( ] session.add_assistant(response.text, tool_calls=tool_calls or None) self._notify_message() - if self.on_turn_end: - self.on_turn_end() + self._safe_call(self.on_turn_end) if not response.tool_calls: self._trace("run_end", reason="done") From efbd3d92137036135b3b138d9950cf9f46f5d1b5 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:41:51 -0300 Subject: [PATCH 08/10] =?UTF-8?q?=E2=9D=A4=EF=B8=8F=20TEST-#80:=20Cover=20?= =?UTF-8?q?message-loss=20hardening=20in=20provider,=20chat,=20serve=20and?= =?UTF-8?q?=20agent?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/cli/test_chat.py | 52 +++++++++++++++++++++ tests/cli/test_serve.py | 99 ++++++++++++++++++++++++++++++++++++++++ tests/core/test_agent.py | 63 +++++++++++++++++++++++++ 3 files changed, 214 insertions(+) create mode 100644 tests/cli/test_serve.py diff --git a/tests/cli/test_chat.py b/tests/cli/test_chat.py index 35aeb5b..0e6004c 100644 --- a/tests/cli/test_chat.py +++ b/tests/cli/test_chat.py @@ -6,6 +6,9 @@ from types import SimpleNamespace from unittest import mock +from rich.markdown import Markdown +from rich.panel import Panel + from pycodeloop.cli.chat import CodeLoopApp @@ -112,5 +115,54 @@ def test_expired_stale_flag_is_not_drained_forever(self): self.assertFalse(app._stale_confirm_answer) +class TestRunTurnPreservesStreamedText(unittest.TestCase): + """`_text_buffer` accumulates streamed text as it arrives; if + `flow.run()` then raises (a failure not covered by the provider's + own partial-response handling), that already-streamed text used to + vanish — only the generic error was shown. It must now be flushed + to the log first.""" + + def test_partial_text_is_logged_before_the_error(self): + app = _fake_app() + + def failing_run(*args, **kwargs): + app._text_buffer = "here is what I had so far" + raise RuntimeError("connection died") + + app.flow.run = failing_run + + import asyncio + + asyncio.run(app._run_turn("do something")) + + logged = [call.args[0] for call in app._log.call_args_list] + markdown_bodies = [ + entry.renderable.markup + for entry in logged + if isinstance(entry, Panel) and isinstance(entry.renderable, Markdown) + ] + self.assertTrue( + any( + "here is what I had so far" in body + for body in markdown_bodies + ) + ) + self.assertEqual(app._text_buffer, "") + + def test_no_buffer_only_logs_the_error(self): + app = _fake_app() + app.flow.run = mock.Mock(side_effect=RuntimeError("boom")) + + import asyncio + + asyncio.run(app._run_turn("do something")) + + logged = [call.args[0] for call in app._log.call_args_list] + self.assertTrue(any("boom" in str(entry) for entry in logged)) + self.assertFalse( + any("interrupted" in str(entry).lower() for entry in logged) + ) + + if __name__ == "__main__": unittest.main() diff --git a/tests/cli/test_serve.py b/tests/cli/test_serve.py new file mode 100644 index 0000000..33a5374 --- /dev/null +++ b/tests/cli/test_serve.py @@ -0,0 +1,99 @@ +"""Unit tests for RpcServer's resilience against a disconnected client +and an unresponsive one — no subprocess, no real stdin/stdout, so these +stay fast and don't share the flakiness of the end-to-end serve tests.""" + +import io +import json +import queue +import time +import unittest +from types import SimpleNamespace +from unittest import mock + +from pycodeloop.cli.serve import RpcServer + + +def _fake_server(): + agent = SimpleNamespace( + on_request=None, + on_text_delta=None, + on_turn_end=None, + on_tool_call=None, + on_tool_result=None, + on_usage=None, + on_context=None, + on_retry=None, + on_compact_start=None, + on_compact_end=None, + confirm=None, + ) + flow = SimpleNamespace(agent=agent, config=SimpleNamespace(storage=None)) + return RpcServer(flow, "generic", "my-model") + + +class TestSendBrokenPipe(unittest.TestCase): + def test_broken_pipe_marks_disconnected_instead_of_raising(self): + server = _fake_server() + + with mock.patch( + "sys.stdout", new=mock.Mock(write=mock.Mock(side_effect=BrokenPipeError)) + ): + server._send({"jsonrpc": "2.0", "method": "chat/heartbeat", "params": {}}) + + self.assertTrue(server._disconnected) + + def test_further_sends_are_skipped_once_disconnected(self): + server = _fake_server() + server._disconnected = True + stdout = mock.Mock() + + with mock.patch("sys.stdout", new=stdout): + server._send({"jsonrpc": "2.0", "method": "chat/heartbeat", "params": {}}) + + stdout.write.assert_not_called() + + +class TestConfirmTimeout(unittest.TestCase): + def test_confirm_times_out_and_declines(self): + server = _fake_server() + server._send = mock.Mock() + + import pycodeloop.cli.serve as serve_module + + with mock.patch.object(serve_module, "_CONFIRM_TIMEOUT", 0.05): + server._wire_callbacks() + answer = server.flow.agent.confirm("bash", "$ echo hi") + + self.assertFalse(answer) + methods = [call.args[0]["method"] for call in server._send.call_args_list] + self.assertIn("chat/confirmTimeout", methods) + + def test_confirm_returns_answer_when_it_arrives_in_time(self): + server = _fake_server() + server._send = mock.Mock( + side_effect=lambda msg: ( + server._confirm_waiters[msg["params"]["id"]].put(True) + if msg["method"] == "chat/confirmRequest" + else None + ) + ) + server._wire_callbacks() + + answer = server.flow.agent.confirm("bash", "$ echo hi") + + self.assertTrue(answer) + + +class TestMalformedInputLine(unittest.TestCase): + def test_malformed_line_is_logged_and_skipped_not_silently_dropped(self): + server = _fake_server() + server._send = mock.Mock() + server.handle = mock.Mock() + + with mock.patch( + "sys.stdin", new=io.StringIO("not json at all\n") + ), mock.patch("pycodeloop.cli.serve.console.print") as mock_print: + server.serve_forever() + + mock_print.assert_called_once() + server.handle.assert_not_called() diff --git a/tests/core/test_agent.py b/tests/core/test_agent.py index fa30120..3ba4f75 100644 --- a/tests/core/test_agent.py +++ b/tests/core/test_agent.py @@ -1159,5 +1159,68 @@ def run(self) -> ToolResult: self.assertEqual(len(tool_message.content), 10_000) +class TestSafeCallback(unittest.TestCase): + """A buggy consumer callback (UI rendering, storage persistence, + etc.) must not abort the turn/tool loop still in progress — only + that one side effect should fail.""" + + def test_on_message_raising_does_not_abort_the_turn(self): + provider = FakeProvider([ProviderResponse(text="done")]) + agent = Agent(provider=provider, tools=[]) + agent.on_message = mock.Mock(side_effect=RuntimeError("storage down")) + session = Session(system_prompt="sys") + + result = agent.run("go", session=session) + + self.assertEqual(result, "done") + self.assertTrue(agent.on_message.called) + + def test_on_tool_result_raising_still_records_the_tool_result(self): + provider = FakeProvider( + [ + ProviderResponse( + text="", + tool_calls=[ + ToolCall(id="1", name="echo", arguments={"text": "hi"}) + ], + ), + ProviderResponse(text="done"), + ] + ) + agent = Agent(provider=provider, tools=[EchoTool()]) + agent.on_tool_result = mock.Mock(side_effect=RuntimeError("render bug")) + session = Session(system_prompt="sys") + + result = agent.run("go", session=session) + + self.assertEqual(result, "done") + tool_message = next(m for m in session.messages if m.role == "tool") + self.assertEqual(tool_message.content, "echo: hi") + + def test_failing_callback_is_traced_without_raising(self): + provider = FakeProvider([ProviderResponse(text="done")]) + agent = Agent(provider=provider, tools=[]) + agent.on_usage = mock.Mock(side_effect=ValueError("boom")) + events = [] + agent.on_trace_event = events.append + session = Session(system_prompt="sys") + + agent.run("go", session=session) + + callback_errors = [e for e in events if e["type"] == "callback_error"] + self.assertEqual(len(callback_errors), 1) + self.assertEqual(callback_errors[0]["error"], "boom") + + def test_trace_event_callback_raising_does_not_propagate(self): + provider = FakeProvider([ProviderResponse(text="done")]) + agent = Agent(provider=provider, tools=[]) + agent.on_trace_event = mock.Mock(side_effect=RuntimeError("bad sink")) + session = Session(system_prompt="sys") + + result = agent.run("go", session=session) + + self.assertEqual(result, "done") + + if __name__ == "__main__": unittest.main() From 827ab6b3e35b6d6f670d919ebf971d7fbe7790cc Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:47:30 -0300 Subject: [PATCH 09/10] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#80:=20Bump=20shipped?= =?UTF-8?q?=20provider=20templates=20timeout=20to=20match=20new=20180s=20d?= =?UTF-8?q?efault?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/providers/templates/anthropic.json | 2 +- pycodeloop/providers/templates/aws.json | 2 +- pycodeloop/providers/templates/deepseek.json | 2 +- pycodeloop/providers/templates/gemini.json | 2 +- pycodeloop/providers/templates/grok.json | 2 +- pycodeloop/providers/templates/groq.json | 2 +- pycodeloop/providers/templates/kimi.json | 2 +- pycodeloop/providers/templates/llama.json | 2 +- pycodeloop/providers/templates/nvidia.json | 2 +- pycodeloop/providers/templates/openai.json | 2 +- pycodeloop/providers/templates/qwen.json | 2 +- pycodeloop/providers/templates/reference.json | 2 +- templates/anthropic.json | 2 +- templates/aws.json | 2 +- templates/deepseek.json | 2 +- templates/gemini.json | 2 +- templates/grok.json | 2 +- templates/groq.json | 2 +- templates/kimi.json | 2 +- templates/llama.json | 2 +- templates/nvidia.json | 2 +- templates/openai.json | 2 +- templates/qwen.json | 2 +- templates/reference.json | 2 +- 24 files changed, 24 insertions(+), 24 deletions(-) diff --git a/pycodeloop/providers/templates/anthropic.json b/pycodeloop/providers/templates/anthropic.json index 3c9abc4..b135c86 100644 --- a/pycodeloop/providers/templates/anthropic.json +++ b/pycodeloop/providers/templates/anthropic.json @@ -7,7 +7,7 @@ "headers": { "anthropic-version": "2023-06-01" }, - "timeout": 60, + "timeout": 180, "request": { "body_paths": { "system": "system" diff --git a/pycodeloop/providers/templates/aws.json b/pycodeloop/providers/templates/aws.json index a3f4ee8..8b5ba23 100644 --- a/pycodeloop/providers/templates/aws.json +++ b/pycodeloop/providers/templates/aws.json @@ -2,5 +2,5 @@ "url": "https://bedrock-mantle.us-east-1.api.aws/v1/chat/completions", "model": "openai.gpt-oss-120b", "api_key_env": "AWS_BEARER_TOKEN_BEDROCK", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/deepseek.json b/pycodeloop/providers/templates/deepseek.json index 1acd68d..2142b55 100644 --- a/pycodeloop/providers/templates/deepseek.json +++ b/pycodeloop/providers/templates/deepseek.json @@ -2,5 +2,5 @@ "url": "https://api.deepseek.com/chat/completions", "model": "deepseek-v4-pro", "api_key_env": "DEEPSEEK_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/gemini.json b/pycodeloop/providers/templates/gemini.json index 9a2648f..52b7800 100644 --- a/pycodeloop/providers/templates/gemini.json +++ b/pycodeloop/providers/templates/gemini.json @@ -2,5 +2,5 @@ "url": "https://generativelanguage.googleapis.com/v1beta/openai/chat/completions", "model": "gemini-3.6-flash", "api_key_env": "GEMINI_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/grok.json b/pycodeloop/providers/templates/grok.json index f3f53fd..577966b 100644 --- a/pycodeloop/providers/templates/grok.json +++ b/pycodeloop/providers/templates/grok.json @@ -2,5 +2,5 @@ "url": "https://api.x.ai/v1/chat/completions", "model": "grok-4.5", "api_key_env": "XAI_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/groq.json b/pycodeloop/providers/templates/groq.json index d7cbb14..4739f2d 100644 --- a/pycodeloop/providers/templates/groq.json +++ b/pycodeloop/providers/templates/groq.json @@ -2,5 +2,5 @@ "url": "https://api.groq.com/openai/v1/chat/completions", "model": "openai/gpt-oss-120b", "api_key_env": "GROQ_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/kimi.json b/pycodeloop/providers/templates/kimi.json index 8f28ffd..5c9230a 100644 --- a/pycodeloop/providers/templates/kimi.json +++ b/pycodeloop/providers/templates/kimi.json @@ -2,5 +2,5 @@ "url": "https://api.moonshot.ai/v1/chat/completions", "model": "kimi-k3", "api_key_env": "MOONSHOT_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/llama.json b/pycodeloop/providers/templates/llama.json index efc8a9c..1b40edd 100644 --- a/pycodeloop/providers/templates/llama.json +++ b/pycodeloop/providers/templates/llama.json @@ -2,5 +2,5 @@ "url": "https://api.together.ai/v1/chat/completions", "model": "meta-llama/Llama-3.3-70B-Instruct-Turbo", "api_key_env": "TOGETHER_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/nvidia.json b/pycodeloop/providers/templates/nvidia.json index c43fa42..3ee90c3 100644 --- a/pycodeloop/providers/templates/nvidia.json +++ b/pycodeloop/providers/templates/nvidia.json @@ -2,5 +2,5 @@ "url": "https://integrate.api.nvidia.com/v1/chat/completions", "model": "meta/llama-3.3-70b-instruct", "api_key_env": "NVIDIA_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/openai.json b/pycodeloop/providers/templates/openai.json index 020b66a..d9c143a 100644 --- a/pycodeloop/providers/templates/openai.json +++ b/pycodeloop/providers/templates/openai.json @@ -2,5 +2,5 @@ "url": "https://api.openai.com/v1/chat/completions", "model": "gpt-5.6", "api_key_env": "OPENAI_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/qwen.json b/pycodeloop/providers/templates/qwen.json index 9470f38..77fcabc 100644 --- a/pycodeloop/providers/templates/qwen.json +++ b/pycodeloop/providers/templates/qwen.json @@ -2,5 +2,5 @@ "url": "https://dashscope-us.aliyuncs.com/compatible-mode/v1/chat/completions", "model": "qwen-max", "api_key_env": "DASHSCOPE_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/pycodeloop/providers/templates/reference.json b/pycodeloop/providers/templates/reference.json index 3db96ce..f795eda 100644 --- a/pycodeloop/providers/templates/reference.json +++ b/pycodeloop/providers/templates/reference.json @@ -8,7 +8,7 @@ "headers": { "X-Org": "acme" }, - "timeout": 60, + "timeout": 180, "request": { "message_shape": "openai", "tool_schema": "openai", diff --git a/templates/anthropic.json b/templates/anthropic.json index 3c9abc4..b135c86 100644 --- a/templates/anthropic.json +++ b/templates/anthropic.json @@ -7,7 +7,7 @@ "headers": { "anthropic-version": "2023-06-01" }, - "timeout": 60, + "timeout": 180, "request": { "body_paths": { "system": "system" diff --git a/templates/aws.json b/templates/aws.json index a3f4ee8..8b5ba23 100644 --- a/templates/aws.json +++ b/templates/aws.json @@ -2,5 +2,5 @@ "url": "https://bedrock-mantle.us-east-1.api.aws/v1/chat/completions", "model": "openai.gpt-oss-120b", "api_key_env": "AWS_BEARER_TOKEN_BEDROCK", - "timeout": 60 + "timeout": 180 } diff --git a/templates/deepseek.json b/templates/deepseek.json index 1acd68d..2142b55 100644 --- a/templates/deepseek.json +++ b/templates/deepseek.json @@ -2,5 +2,5 @@ "url": "https://api.deepseek.com/chat/completions", "model": "deepseek-v4-pro", "api_key_env": "DEEPSEEK_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/gemini.json b/templates/gemini.json index 9a2648f..52b7800 100644 --- a/templates/gemini.json +++ b/templates/gemini.json @@ -2,5 +2,5 @@ "url": "https://generativelanguage.googleapis.com/v1beta/openai/chat/completions", "model": "gemini-3.6-flash", "api_key_env": "GEMINI_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/grok.json b/templates/grok.json index f3f53fd..577966b 100644 --- a/templates/grok.json +++ b/templates/grok.json @@ -2,5 +2,5 @@ "url": "https://api.x.ai/v1/chat/completions", "model": "grok-4.5", "api_key_env": "XAI_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/groq.json b/templates/groq.json index d7cbb14..4739f2d 100644 --- a/templates/groq.json +++ b/templates/groq.json @@ -2,5 +2,5 @@ "url": "https://api.groq.com/openai/v1/chat/completions", "model": "openai/gpt-oss-120b", "api_key_env": "GROQ_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/kimi.json b/templates/kimi.json index 8f28ffd..5c9230a 100644 --- a/templates/kimi.json +++ b/templates/kimi.json @@ -2,5 +2,5 @@ "url": "https://api.moonshot.ai/v1/chat/completions", "model": "kimi-k3", "api_key_env": "MOONSHOT_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/llama.json b/templates/llama.json index efc8a9c..1b40edd 100644 --- a/templates/llama.json +++ b/templates/llama.json @@ -2,5 +2,5 @@ "url": "https://api.together.ai/v1/chat/completions", "model": "meta-llama/Llama-3.3-70B-Instruct-Turbo", "api_key_env": "TOGETHER_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/nvidia.json b/templates/nvidia.json index c43fa42..3ee90c3 100644 --- a/templates/nvidia.json +++ b/templates/nvidia.json @@ -2,5 +2,5 @@ "url": "https://integrate.api.nvidia.com/v1/chat/completions", "model": "meta/llama-3.3-70b-instruct", "api_key_env": "NVIDIA_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/openai.json b/templates/openai.json index 020b66a..d9c143a 100644 --- a/templates/openai.json +++ b/templates/openai.json @@ -2,5 +2,5 @@ "url": "https://api.openai.com/v1/chat/completions", "model": "gpt-5.6", "api_key_env": "OPENAI_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/qwen.json b/templates/qwen.json index 9470f38..77fcabc 100644 --- a/templates/qwen.json +++ b/templates/qwen.json @@ -2,5 +2,5 @@ "url": "https://dashscope-us.aliyuncs.com/compatible-mode/v1/chat/completions", "model": "qwen-max", "api_key_env": "DASHSCOPE_API_KEY", - "timeout": 60 + "timeout": 180 } diff --git a/templates/reference.json b/templates/reference.json index 3db96ce..f795eda 100644 --- a/templates/reference.json +++ b/templates/reference.json @@ -8,7 +8,7 @@ "headers": { "X-Org": "acme" }, - "timeout": 60, + "timeout": 180, "request": { "message_shape": "openai", "tool_schema": "openai", From fe5068f7c4f76f11ce2709c3ac20847b4ec4d72c Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Mon, 24 Aug 2026 18:47:31 -0300 Subject: [PATCH 10/10] =?UTF-8?q?=F0=9F=93=9D=20PEP8-#80:=20Fix=20ruff=20l?= =?UTF-8?q?int/format=20violations=20flagged=20by=20CI?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/cli/serve.py | 1 - pycodeloop/core/agent.py | 9 +++------ tests/cli/test_chat.py | 6 +++--- tests/cli/test_serve.py | 25 +++++++++++++++---------- tests/core/test_agent.py | 4 +++- tests/providers/test_generic.py | 12 +++++++----- 6 files changed, 31 insertions(+), 26 deletions(-) diff --git a/pycodeloop/cli/serve.py b/pycodeloop/cli/serve.py index d19bf2f..5b5f30b 100644 --- a/pycodeloop/cli/serve.py +++ b/pycodeloop/cli/serve.py @@ -9,7 +9,6 @@ import queue import sys import threading -import time import uuid import typer diff --git a/pycodeloop/core/agent.py b/pycodeloop/core/agent.py index cba5939..02040cc 100644 --- a/pycodeloop/core/agent.py +++ b/pycodeloop/core/agent.py @@ -2,6 +2,7 @@ from __future__ import annotations +import contextlib import threading import time from collections.abc import Callable @@ -116,10 +117,8 @@ def __init__( def _trace(self, event_type: str, **fields) -> None: if self.on_trace_event: - try: + with contextlib.suppress(Exception): self.on_trace_event({"type": event_type, **fields}) - except Exception: - pass def _safe_call(self, callback: Callable | None, *args) -> None: """Invokes a consumer-supplied `on_*` callback (UI rendering, @@ -517,9 +516,7 @@ def run( return error_text self.usage = self.usage + response.usage - self._safe_call( - self.on_usage, response.usage, self.usage, elapsed - ) + self._safe_call(self.on_usage, response.usage, self.usage, elapsed) self._trace( "turn", diff --git a/tests/cli/test_chat.py b/tests/cli/test_chat.py index 0e6004c..470c24b 100644 --- a/tests/cli/test_chat.py +++ b/tests/cli/test_chat.py @@ -139,12 +139,12 @@ def failing_run(*args, **kwargs): markdown_bodies = [ entry.renderable.markup for entry in logged - if isinstance(entry, Panel) and isinstance(entry.renderable, Markdown) + if isinstance(entry, Panel) + and isinstance(entry.renderable, Markdown) ] self.assertTrue( any( - "here is what I had so far" in body - for body in markdown_bodies + "here is what I had so far" in body for body in markdown_bodies ) ) self.assertEqual(app._text_buffer, "") diff --git a/tests/cli/test_serve.py b/tests/cli/test_serve.py index 33a5374..37820a0 100644 --- a/tests/cli/test_serve.py +++ b/tests/cli/test_serve.py @@ -3,9 +3,6 @@ stay fast and don't share the flakiness of the end-to-end serve tests.""" import io -import json -import queue -import time import unittest from types import SimpleNamespace from unittest import mock @@ -36,9 +33,12 @@ def test_broken_pipe_marks_disconnected_instead_of_raising(self): server = _fake_server() with mock.patch( - "sys.stdout", new=mock.Mock(write=mock.Mock(side_effect=BrokenPipeError)) + "sys.stdout", + new=mock.Mock(write=mock.Mock(side_effect=BrokenPipeError)), ): - server._send({"jsonrpc": "2.0", "method": "chat/heartbeat", "params": {}}) + server._send( + {"jsonrpc": "2.0", "method": "chat/heartbeat", "params": {}} + ) self.assertTrue(server._disconnected) @@ -48,7 +48,9 @@ def test_further_sends_are_skipped_once_disconnected(self): stdout = mock.Mock() with mock.patch("sys.stdout", new=stdout): - server._send({"jsonrpc": "2.0", "method": "chat/heartbeat", "params": {}}) + server._send( + {"jsonrpc": "2.0", "method": "chat/heartbeat", "params": {}} + ) stdout.write.assert_not_called() @@ -65,7 +67,9 @@ def test_confirm_times_out_and_declines(self): answer = server.flow.agent.confirm("bash", "$ echo hi") self.assertFalse(answer) - methods = [call.args[0]["method"] for call in server._send.call_args_list] + methods = [ + call.args[0]["method"] for call in server._send.call_args_list + ] self.assertIn("chat/confirmTimeout", methods) def test_confirm_returns_answer_when_it_arrives_in_time(self): @@ -90,9 +94,10 @@ def test_malformed_line_is_logged_and_skipped_not_silently_dropped(self): server._send = mock.Mock() server.handle = mock.Mock() - with mock.patch( - "sys.stdin", new=io.StringIO("not json at all\n") - ), mock.patch("pycodeloop.cli.serve.console.print") as mock_print: + with ( + mock.patch("sys.stdin", new=io.StringIO("not json at all\n")), + mock.patch("pycodeloop.cli.serve.console.print") as mock_print, + ): server.serve_forever() mock_print.assert_called_once() diff --git a/tests/core/test_agent.py b/tests/core/test_agent.py index 3ba4f75..ce4d065 100644 --- a/tests/core/test_agent.py +++ b/tests/core/test_agent.py @@ -1188,7 +1188,9 @@ def test_on_tool_result_raising_still_records_the_tool_result(self): ] ) agent = Agent(provider=provider, tools=[EchoTool()]) - agent.on_tool_result = mock.Mock(side_effect=RuntimeError("render bug")) + agent.on_tool_result = mock.Mock( + side_effect=RuntimeError("render bug") + ) session = Session(system_prompt="sys") result = agent.run("go", session=session) diff --git a/tests/providers/test_generic.py b/tests/providers/test_generic.py index 877d7f9..96a1368 100644 --- a/tests/providers/test_generic.py +++ b/tests/providers/test_generic.py @@ -423,12 +423,14 @@ def test_streaming_reraises_timeout_when_nothing_was_streamed_yet(self): ) provider = GenericProvider.from_json(path) - with mock.patch( - "pycodeloop.providers.generic.urllib.request.urlopen", - return_value=_TimeoutAfterLinesResponse(b"", raise_after=0), + with ( + mock.patch( + "pycodeloop.providers.generic.urllib.request.urlopen", + return_value=_TimeoutAfterLinesResponse(b"", raise_after=0), + ), + self.assertRaises(TimeoutError), ): - with self.assertRaises(TimeoutError): - provider.complete("sys", [], [], on_delta=lambda _: None) + provider.complete("sys", [], [], on_delta=lambda _: None) def test_streaming_stops_promptly_when_cancel_event_is_set(self): """Regression: cancel_event was accepted nowhere in the streaming