From ba29e379e45296acfaf5197aa2ec892a21d0f8a0 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Sat, 15 Aug 2026 23:59:01 -0300 Subject: [PATCH 1/4] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#62:=20Persist=20already?= =?UTF-8?q?-shown=20text=20when=20an=20SSE=20chunk=20is=20malformed?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/providers/generic.py | 14 ++++++++++++-- 1 file changed, 12 insertions(+), 2 deletions(-) diff --git a/pycodeloop/providers/generic.py b/pycodeloop/providers/generic.py index f2854ec..54d2fbb 100644 --- a/pycodeloop/providers/generic.py +++ b/pycodeloop/providers/generic.py @@ -383,7 +383,8 @@ def _stream( body = {**body, "stream": True} text = "" pending: dict[int, dict] = {} - stop_reason = "stop" + stop_reason: str | None = None + saw_terminal_marker = False usage = Usage() with self._open(body, config) as response: @@ -393,8 +394,13 @@ def _stream( continue payload = line[len("data: ") :] if payload == "[DONE]": + saw_terminal_marker = True + break + try: + chunk = json.loads(payload) + except json.JSONDecodeError: + stop_reason = "malformed_stream" break - chunk = json.loads(payload) if chunk.get("usage"): usage = Usage( @@ -451,6 +457,10 @@ def _stream( if choice.get("finish_reason"): stop_reason = choice["finish_reason"] + saw_terminal_marker = True + + if stop_reason is None: + stop_reason = "stop" if saw_terminal_marker else "connection_lost" tool_calls = [ ToolCall( From 6eaa76345567f5c81b7d61f6f81022158ea79a38 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Sat, 15 Aug 2026 23:59:01 -0300 Subject: [PATCH 2/4] =?UTF-8?q?=E2=9D=A4=EF=B8=8F=20TEST-#62:=20Cover=20ma?= =?UTF-8?q?lformed=20chunk=20and=20dropped=20connection=20during=20streami?= =?UTF-8?q?ng?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/providers/test_generic.py | 42 +++++++++++++++++++++++++++++++++ 1 file changed, 42 insertions(+) diff --git a/tests/providers/test_generic.py b/tests/providers/test_generic.py index 2c71abc..b6d8dec 100644 --- a/tests/providers/test_generic.py +++ b/tests/providers/test_generic.py @@ -267,6 +267,48 @@ def test_streaming_response_also_captures_extra_tool_call_fields(self): {"extra_content": {"google": {"thought_signature": "xyz789"}}}, ) + def test_streaming_keeps_already_shown_text_on_malformed_chunk(self): + path = self._write_config( + {"url": "http://fake/v1/chat/completions", "model": "my-model"} + ) + provider = GenericProvider.from_json(path) + + chunks = [{"choices": [{"delta": {"content": "hello there"}}]}] + sse_body = ( + "".join(f"data: {json.dumps(c)}\n" for c in chunks) + + "data: {not valid json\n" + + "data: [DONE]\n" + ).encode() + + deltas = [] + with mock.patch( + "pycodeloop.providers.generic.urllib.request.urlopen", + return_value=_FakeResponse(sse_body), + ): + result = provider.complete("sys", [], [], on_delta=deltas.append) + + self.assertEqual(result.text, "hello there") + self.assertEqual(result.stop_reason, "malformed_stream") + self.assertEqual("".join(deltas), "hello there") + + def test_streaming_flags_a_connection_dropped_mid_response(self): + path = self._write_config( + {"url": "http://fake/v1/chat/completions", "model": "my-model"} + ) + provider = GenericProvider.from_json(path) + + chunks = [{"choices": [{"delta": {"content": "cut off mid"}}]}] + 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=_FakeResponse(sse_body), + ): + result = provider.complete("sys", [], [], on_delta=lambda _: None) + + self.assertEqual(result.text, "cut off mid") + self.assertEqual(result.stop_reason, "connection_lost") + def test_streaming_cuts_a_looping_response_short(self): path = self._write_config( {"url": "http://fake/v1/chat/completions", "model": "my-model"} From ede73035a1a074884d808153e3d101ad98639f15 Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Sun, 16 Aug 2026 00:25:11 -0300 Subject: [PATCH 3/4] =?UTF-8?q?=F0=9F=AA=B2=20BUG-#62:=20Don't=20let=20a?= =?UTF-8?q?=20trailing=20malformed=20chunk=20clobber=20a=20clean=20finish?= =?UTF-8?q?=5Freason?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- pycodeloop/providers/generic.py | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/pycodeloop/providers/generic.py b/pycodeloop/providers/generic.py index 54d2fbb..754e5cd 100644 --- a/pycodeloop/providers/generic.py +++ b/pycodeloop/providers/generic.py @@ -399,7 +399,8 @@ def _stream( try: chunk = json.loads(payload) except json.JSONDecodeError: - stop_reason = "malformed_stream" + if stop_reason is None: + stop_reason = "malformed_stream" break if chunk.get("usage"): From a18230edf3deb16f9ed86db8b4501e0e51432d9a Mon Sep 17 00:00:00 2001 From: Fernando Celmer Date: Sun, 16 Aug 2026 00:25:13 -0300 Subject: [PATCH 4/4] =?UTF-8?q?=E2=9D=A4=EF=B8=8F=20TEST-#62:=20Cover=20fi?= =?UTF-8?q?nish=5Freason=20precedence=20over=20trailing=20malformed=20chun?= =?UTF-8?q?k?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tests/providers/test_generic.py | 29 +++++++++++++++++++++++++++++ 1 file changed, 29 insertions(+) diff --git a/tests/providers/test_generic.py b/tests/providers/test_generic.py index b6d8dec..903846b 100644 --- a/tests/providers/test_generic.py +++ b/tests/providers/test_generic.py @@ -291,6 +291,35 @@ def test_streaming_keeps_already_shown_text_on_malformed_chunk(self): self.assertEqual(result.stop_reason, "malformed_stream") self.assertEqual("".join(deltas), "hello there") + def test_streaming_prefers_finish_reason_over_trailing_malformed_chunk( + self, + ): + path = self._write_config( + {"url": "http://fake/v1/chat/completions", "model": "my-model"} + ) + provider = GenericProvider.from_json(path) + + chunks = [ + { + "choices": [ + {"delta": {"content": "done"}, "finish_reason": "stop"} + ] + } + ] + sse_body = ( + "".join(f"data: {json.dumps(c)}\n" for c in chunks) + + "data: {junk after finish\n" + ).encode() + + with mock.patch( + "pycodeloop.providers.generic.urllib.request.urlopen", + return_value=_FakeResponse(sse_body), + ): + result = provider.complete("sys", [], [], on_delta=lambda _: None) + + self.assertEqual(result.text, "done") + self.assertEqual(result.stop_reason, "stop") + def test_streaming_flags_a_connection_dropped_mid_response(self): path = self._write_config( {"url": "http://fake/v1/chat/completions", "model": "my-model"}