From 25fc02d07aa7a8d7ad7156606bf6bb2fde729b35 Mon Sep 17 00:00:00 2001 From: AIOSAI Date: Wed, 15 Jul 2026 11:14:05 -0700 Subject: [PATCH] =?UTF-8?q?feat:=20TG=20streaming=20v2=20polish=20(DPLAN-0?= =?UTF-8?q?229)=20=E2=80=94=20logs=5Fwere=5Factive=20+=20multi-chunk=20(>4?= =?UTF-8?q?096)=20finalize=20now=20honor=20the=20streaming=20flag:=20logs-?= =?UTF-8?q?active=20reconcile-edits=20the=20streamed=20message=20instead?= =?UTF-8?q?=20of=20Done.+fresh,=20multi-chunk=20edits=20chunk=201=20in=20p?= =?UTF-8?q?lace=20+=20sends=20[2/N]=20continuations,=20edit-fail=20falls?= =?UTF-8?q?=20back=20safely.=20Batch=20path=20verified=20zero-change=20via?= =?UTF-8?q?=20regression=20tests.=206=20new=20tests,=201077=20hooks=20gree?= =?UTF-8?q?n=20+=20seedgo=20100%=20devpulse-verified.=20Live=20streamed-tu?= =?UTF-8?q?rn=20proof=20pending=20Patrick's=20next=20streaming=20session.?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- CHANGELOG.md | 13 ++ .../notification/telegram_response.py | 16 ++- .../hooks/tests/test_telegram_response.py | 111 ++++++++++++++++++ 3 files changed, 137 insertions(+), 3 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 54f58a6b..c14e073a 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -66,6 +66,19 @@ PyPI version — not the changelog header. evidence now spans all three storm shapes: continuous fast (332/min), continuous moderate (257/min), bursty (191/min). +### Added + +- **TG streaming v2 polish (DPLAN-0229): the last two finalize paths now + honor the streaming flag.** v1 shipped with a deliberate gap — when logs + were active mid-turn or the final response exceeded 4096 chars, the Stop + hook fell back to "Done." + a fresh message, orphaning the streamed bubble. + @hooks threaded `streaming` through `_deliver_chunks`: logs-active now + reconcile-edits the streamed message with the final formatted response, and + multi-chunk edits chunk 1 in place then sends [2/N]+ as continuations. + Batch mode is verified zero-change (regression tests for both paths), plus + an edit-fail fallback. 6 new tests, 1077 green. Live streamed-turn proof + pending Patrick's next streaming session — honestly flagged, not faked. + ## [2026-07-14] ### Fixed diff --git a/src/aipass/hooks/apps/handlers/notification/telegram_response.py b/src/aipass/hooks/apps/handlers/notification/telegram_response.py index 10649022..c52b40d6 100644 --- a/src/aipass/hooks/apps/handlers/notification/telegram_response.py +++ b/src/aipass/hooks/apps/handlers/notification/telegram_response.py @@ -564,11 +564,19 @@ def handle(hook_data: dict) -> dict: response_text = _prepend_branch_prefix(response_text) logs_were_active = _check_log_streamer_active() + streaming = bool(pending_data.get("streaming", False)) chunks = chunk_text(response_text) - logger.info("[HOOKS] telegram: sending %d chunk(s) (logs_active=%s)", len(chunks), logs_were_active) + logger.info( + "[HOOKS] telegram: sending %d chunk(s) (logs_active=%s, streaming=%s)", + len(chunks), + logs_were_active, + streaming, + ) - all_sent, chunk_results = _deliver_chunks(chunks, bot_token, chat_id, processing_message_id, logs_were_active) + all_sent, chunk_results = _deliver_chunks( + chunks, bot_token, chat_id, processing_message_id, logs_were_active, streaming + ) _write_delivery_log(response_text, chunks, chunk_results, session_id) @@ -658,6 +666,7 @@ def _deliver_chunks( chat_id: int, processing_message_id: int | None, logs_were_active: bool, + streaming: bool = False, ) -> tuple[bool, list[dict]]: """Send all response chunks to Telegram. Returns (all_sent, per-chunk results).""" chunk_results: list[dict] = [] @@ -666,7 +675,8 @@ def _deliver_chunks( for i, chunk in enumerate(chunks): if i == 0 and processing_message_id: - if single and not logs_were_active: + should_reconcile = streaming or (single and not logs_were_active) + if should_reconcile: result = edit_telegram_message(bot_token, chat_id, processing_message_id, chunk) if result["ok"]: chunk_results.append({"idx": i, "method": "edit", **result}) diff --git a/src/aipass/hooks/tests/test_telegram_response.py b/src/aipass/hooks/tests/test_telegram_response.py index a9992c6b..e67e3ac0 100644 --- a/src/aipass/hooks/tests/test_telegram_response.py +++ b/src/aipass/hooks/tests/test_telegram_response.py @@ -1308,6 +1308,117 @@ class TestDeliverChunks: assert chunk_results[1]["message_id"] == 202 +# =========================================================================== +# _deliver_chunks streaming mode +# =========================================================================== + + +class TestDeliverChunksStreaming: + """Streaming flag changes reconcile behavior — batch unchanged.""" + + def test_streaming_logs_active_reconciles_instead_of_done(self): + """Streaming + logs_active: edit processing msg with response, not 'Done.'.""" + from aipass.hooks.apps.handlers.notification.telegram_response import _deliver_chunks + + with ( + patch(LOGGER_PATCH), + patch(f"{MOD}.edit_telegram_message", return_value=_ok_result()) as mock_edit, + ): + all_sent, chunk_results = _deliver_chunks(["Hello"], "tok", 123, 789, True, streaming=True) + + assert all_sent is True + assert chunk_results[0]["method"] == "edit" + mock_edit.assert_called_once_with("tok", 123, 789, "Hello") + + def test_streaming_multi_chunk_edits_first_sends_rest(self): + """Streaming + multi-chunk: edit chunk 1 into processing msg, send rest.""" + from aipass.hooks.apps.handlers.notification.telegram_response import _deliver_chunks + + sent_texts = [] + + def capture_send(bot_token, chat_id, text): + sent_texts.append(text) + return _ok_result() + + with ( + patch(LOGGER_PATCH), + patch(f"{MOD}.edit_telegram_message", return_value=_ok_result()) as mock_edit, + patch(f"{MOD}._send_with_retry", side_effect=capture_send), + ): + all_sent, chunk_results = _deliver_chunks( + ["Part A", "Part B", "Part C"], "tok", 123, 789, False, streaming=True + ) + + assert all_sent is True + mock_edit.assert_called_once_with("tok", 123, 789, "Part A") + assert chunk_results[0]["method"] == "edit" + assert len(sent_texts) == 2 + assert "[2/3]" in sent_texts[0] + assert "[3/3]" in sent_texts[1] + + def test_streaming_multi_chunk_edit_fails_sends_all(self): + """Streaming + multi-chunk: if edit fails, fall back to send for chunk 1.""" + from aipass.hooks.apps.handlers.notification.telegram_response import _deliver_chunks + + sent_texts = [] + + def capture_send(bot_token, chat_id, text): + sent_texts.append(text) + return _ok_result() + + with ( + patch(LOGGER_PATCH), + patch(f"{MOD}.edit_telegram_message", return_value=_fail_result()), + patch(f"{MOD}._send_with_retry", side_effect=capture_send), + ): + all_sent, chunk_results = _deliver_chunks(["Part A", "Part B"], "tok", 123, 789, False, streaming=True) + + assert all_sent is True + assert chunk_results[0]["method"] == "send" + assert len(sent_texts) == 2 + + def test_batch_logs_active_still_sends_done(self): + """Batch mode (no streaming): logs_active still sends 'Done.' — zero regression.""" + from aipass.hooks.apps.handlers.notification.telegram_response import _deliver_chunks + + with ( + patch(LOGGER_PATCH), + patch(f"{MOD}.edit_telegram_message", return_value=_ok_result()) as mock_edit, + patch(f"{MOD}._send_with_retry", return_value=_ok_result()), + ): + all_sent, chunk_results = _deliver_chunks(["Hello"], "tok", 123, 789, True, streaming=False) + + assert all_sent is True + mock_edit.assert_called_once_with("tok", 123, 789, "Done.") + + def test_batch_multi_chunk_still_sends_done(self): + """Batch mode (no streaming): multi-chunk still sends 'Done.' — zero regression.""" + from aipass.hooks.apps.handlers.notification.telegram_response import _deliver_chunks + + with ( + patch(LOGGER_PATCH), + patch(f"{MOD}.edit_telegram_message", return_value=_ok_result()) as mock_edit, + patch(f"{MOD}._send_with_retry", return_value=_ok_result()), + ): + _deliver_chunks(["Part A", "Part B"], "tok", 123, 789, False, streaming=False) + + mock_edit.assert_called_once_with("tok", 123, 789, "Done.") + + def test_streaming_single_no_logs_still_edits(self): + """Streaming + single chunk + no logs: same as batch — reconcile-edit.""" + from aipass.hooks.apps.handlers.notification.telegram_response import _deliver_chunks + + with ( + patch(LOGGER_PATCH), + patch(f"{MOD}.edit_telegram_message", return_value=_ok_result()) as mock_edit, + ): + all_sent, chunk_results = _deliver_chunks(["Hello"], "tok", 123, 789, False, streaming=True) + + assert all_sent is True + mock_edit.assert_called_once_with("tok", 123, 789, "Hello") + assert chunk_results[0]["method"] == "edit" + + # =========================================================================== # _advance_pending # ===========================================================================