feat: TG streaming v2 polish (DPLAN-0229) — logs_were_active + multi-chunk (>4096) finalize now honor the streaming flag: logs-active reconcile-edits the streamed message instead of Done.+fresh, multi-chunk edits chunk 1 in place + sends [2/N] continuations, edit-fail falls back safely. Batch path verified zero-change via regression tests. 6 new tests, 1077 hooks green + seedgo 100% devpulse-verified. Live streamed-turn proof pending Patrick's next streaming session.

This commit is contained in:
AIOSAI
2026-07-15 11:14:05 -07:00
parent 9dc2ecd604
commit 25fc02d07a
3 changed files with 137 additions and 3 deletions
+13
View File
@@ -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
@@ -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})
@@ -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
# ===========================================================================