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 # ===========================================================================