Merge branch 'simp/r2-gw-rest-stream' into simp/integration2

This commit is contained in:
Teknium
2026-09-02 17:05:22 -07:00
9 changed files with 3235 additions and 3083 deletions
+58 -1027
View File
File diff suppressed because it is too large Load Diff
+35
View File
@@ -0,0 +1,35 @@
"""Plumbing shared by the kanban notifier and dispatcher loops."""
from __future__ import annotations
import asyncio
import logging
from contextvars import Context
from typing import Any, Callable
# Keep the logger name run.py used so extracted log records are unchanged.
logger = logging.getLogger("gateway.run")
def _run_in_fresh_context(func: Callable[..., Any], /, *args: Any) -> Any:
"""Run *func* in an empty ``Context`` so request-local ContextVars stay behind.
``asyncio.to_thread`` copies the caller's context; a lingering
``delegate_task`` child marker would make ``write_txn`` false-trip for
these process-owned writers. An empty Context keeps the DB guard intact
for real children without exempting dispatcher writes.
"""
return Context().run(func, *args)
async def _to_thread_process_service(func: Callable[..., Any], /, *args: Any) -> Any:
"""Offload blocking process-service work without inheriting request ContextVars."""
return await asyncio.to_thread(_run_in_fresh_context, func, *args)
def _list_boards(kb: Any) -> list:
"""Enumerate live boards; fall back to the default board when listing fails."""
try:
return kb.list_boards(include_archived=False)
except Exception:
return [kb.read_board_metadata(kb.DEFAULT_BOARD)]
+375
View File
@@ -0,0 +1,375 @@
"""Embedded kanban dispatcher: settings resolution and per-tick board work.
``GatewayKanbanWatchersMixin._kanban_dispatcher_watcher`` owns the loop,
the singleton lock and the health telemetry; everything that only needs the
``kanban_db`` module and the resolved settings lives here.
"""
from __future__ import annotations
import contextlib
import os
import sqlite3
import time
from dataclasses import dataclass
from typing import Any, Optional
from gateway.kanban_watchers_common import _list_boards, logger
_CORRUPT_DB_MARKERS = ("file is not a database", "database disk image is malformed")
@dataclass
class _DispatcherSettings:
"""``kanban.*`` dispatch settings, read once at boot (restart to apply)."""
interval: float
max_spawn: Any
max_in_progress: Optional[int]
failure_limit: int
stale_timeout_seconds: int
reconcile_orphans: bool
default_assignee: Optional[str]
max_in_progress_per_profile: Optional[int]
def _positive_int_setting(kanban_cfg: dict, key: str) -> Optional[int]:
"""Parse an optional ``kanban.<key>`` int cap; None when unset or invalid (< 1 is invalid)."""
raw = kanban_cfg.get(key)
if raw is None:
return None
try:
value = int(raw)
except (TypeError, ValueError):
logger.warning("kanban dispatcher: invalid kanban.%s=%r; ignoring", key, raw)
return None
if value < 1:
logger.warning("kanban dispatcher: kanban.%s=%r is below 1; ignoring", key, raw)
return None
logger.info("kanban dispatcher: %s=%d", key, value)
return value
def _resolve_dispatcher_settings(kanban_cfg: dict, kb: Any) -> _DispatcherSettings:
"""Parse and log the dispatcher settings in their established order."""
try:
interval = float(kanban_cfg.get("dispatch_interval_seconds", 60) or 60)
except (ValueError, TypeError):
logger.warning(
"kanban dispatcher: invalid dispatch_interval_seconds=%r, using default 60",
kanban_cfg.get("dispatch_interval_seconds"),
)
interval = 60.0
interval = max(interval, 1.0) # sanity floor — tighter than this is a footgun
max_spawn = kanban_cfg.get("max_spawn")
if max_spawn is not None:
logger.info("kanban dispatcher: max_spawn=%s", max_spawn)
# Cap simultaneously running tasks so slow workers don't pile up and
# time out. Explicit config wins; otherwise a memory-derived default
# (unbounded fan-out swap-thrashes small hosts), or None where total
# memory can't be read.
max_in_progress = _positive_int_setting(kanban_cfg, "max_in_progress")
effective_max_in_progress = kb.resolve_max_in_progress(max_in_progress)
if max_in_progress is None and effective_max_in_progress is not None:
logger.info(
"kanban dispatcher: kanban.max_in_progress unset; using "
"memory-derived default max_in_progress=%d "
"(set kanban.max_in_progress in config.yaml to override)",
effective_max_in_progress,
)
raw_failure_limit = kanban_cfg.get("failure_limit", kb.DEFAULT_FAILURE_LIMIT)
try:
failure_limit = int(raw_failure_limit)
except (TypeError, ValueError):
logger.warning(
"kanban dispatcher: invalid kanban.failure_limit=%r; using default %d",
raw_failure_limit,
kb.DEFAULT_FAILURE_LIMIT,
)
failure_limit = kb.DEFAULT_FAILURE_LIMIT
if failure_limit < 1:
logger.warning(
"kanban dispatcher: kanban.failure_limit=%r is below 1; using default %d",
raw_failure_limit,
kb.DEFAULT_FAILURE_LIMIT,
)
failure_limit = kb.DEFAULT_FAILURE_LIMIT
# 0 disables stale detection.
raw_stale = kanban_cfg.get("dispatch_stale_timeout_seconds", 0)
try:
stale_timeout_seconds = int(raw_stale or 0)
except (TypeError, ValueError):
logger.warning(
"kanban dispatcher: invalid kanban.dispatch_stale_timeout_seconds=%r; "
"disabling stale detection",
raw_stale,
)
stale_timeout_seconds = 0
# Fallback profile for tasks created without an assignee (e.g. via the
# dashboard). Empty (the schema default) keeps skipping them.
default_assignee = (kanban_cfg.get("default_assignee") or "").strip() or None
if default_assignee:
logger.info(
"kanban dispatcher: default_assignee=%r (unassigned ready tasks "
"will route to this profile)",
default_assignee,
)
return _DispatcherSettings(
interval=interval,
max_spawn=max_spawn,
max_in_progress=effective_max_in_progress,
failure_limit=failure_limit,
stale_timeout_seconds=stale_timeout_seconds,
# Requeue 'running' cards with broken claim bookkeeping (zombie-card
# reconciliation); false keeps orphans frozen for manual forensics.
reconcile_orphans=bool(kanban_cfg.get("reconcile_orphans", True)),
default_assignee=default_assignee,
# Per-profile concurrency cap: no single profile's local model / API
# quota / browser pool gets overwhelmed by a fan-out.
max_in_progress_per_profile=_positive_int_setting(kanban_cfg, "max_in_progress_per_profile"),
)
class _KanbanDispatcher:
"""Per-tick board work for the embedded dispatcher (runs in worker threads).
Boards are enumerated every tick so a board created mid-run is picked up
without a restart. Corrupt-looking board DBs are quarantined per
fingerprint and retried after ``CORRUPT_BOARD_RETRY_AFTER_SECONDS``:
transient WAL/open races can look like "malformed" for one tick.
"""
CORRUPT_BOARD_RETRY_AFTER_SECONDS = 300
def __init__(self, kb: Any, settings: _DispatcherSettings) -> None:
self.kb = kb
self.settings = settings
self.disabled_corrupt_boards: dict[
str, tuple[tuple[str, int | None, int | None], float]
] = {}
def _board_slugs(self) -> list:
return [b.get("slug") or self.kb.DEFAULT_BOARD for b in _list_boards(self.kb)]
def board_db_fingerprint(self, slug: str) -> tuple[str, int | None, int | None]:
path = self.kb.kanban_db_path(slug)
try:
resolved = str(path.expanduser().resolve())
except Exception:
resolved = str(path)
try:
stat = path.stat()
except OSError:
return (resolved, None, None)
return (resolved, stat.st_mtime_ns, stat.st_size)
def is_corrupt_board_db_error(self, exc: Exception) -> bool:
corrupt_guard_error = getattr(self.kb, "KanbanDbCorruptError", None)
if corrupt_guard_error is not None and isinstance(exc, corrupt_guard_error):
return True
if not isinstance(exc, sqlite3.DatabaseError):
return False
msg = str(exc).lower()
return any(marker in msg for marker in _CORRUPT_DB_MARKERS)
def _quarantine_lifted(self, slug: str, fingerprint: tuple) -> bool:
"""Return False while *slug* stays quarantined; lift (and log) otherwise."""
disabled_entry = self.disabled_corrupt_boards.get(slug)
if disabled_entry is None:
return True
disabled_fingerprint, disabled_at = disabled_entry
age = time.monotonic() - disabled_at
if disabled_fingerprint == fingerprint and age < self.CORRUPT_BOARD_RETRY_AFTER_SECONDS:
return False
if disabled_fingerprint == fingerprint:
logger.info(
"kanban dispatcher: board %s database fingerprint unchanged "
"after %.0fs quarantine; retrying dispatch",
slug,
age,
)
else:
logger.info(
"kanban dispatcher: board %s database changed; retrying dispatch",
slug,
)
self.disabled_corrupt_boards.pop(slug, None)
return True
def tick_once_for_board(self, slug: str) -> Optional[object]:
"""Run one dispatch_once for a specific board.
The per-board DB is opened explicitly so boards never share a
connection or claim across each other.
"""
conn = None
fingerprint = self.board_db_fingerprint(slug)
if not self._quarantine_lifted(slug, fingerprint):
return None
s = self.settings
try:
# No explicit init_db(): connect() runs the migration once per
# process (see the matching note in the notifier collector).
conn = self.kb.connect(board=slug)
return self.kb.dispatch_once(
conn,
board=slug,
max_spawn=s.max_spawn,
max_in_progress=s.max_in_progress,
failure_limit=s.failure_limit,
stale_timeout_seconds=s.stale_timeout_seconds,
default_assignee=s.default_assignee,
max_in_progress_per_profile=s.max_in_progress_per_profile,
reconcile_orphans=s.reconcile_orphans,
)
except Exception as exc:
if self.is_corrupt_board_db_error(exc):
self.disabled_corrupt_boards[slug] = (fingerprint, time.monotonic())
logger.error(
"kanban dispatcher: board %s database %s is not a valid "
"SQLite database; pausing dispatch for this board until "
"the file changes, the gateway restarts, or the "
"quarantine timer expires. Move or restore the file, "
"then run `hermes kanban init` if you need a fresh board.",
slug,
fingerprint[0],
)
return None
logger.exception("kanban dispatcher: tick failed on board %s", slug)
return None
finally:
if conn is not None:
with contextlib.suppress(Exception):
conn.close()
def tick_once(self) -> list[tuple[str, Optional[object]]]:
"""Run one dispatch_once per board. Returns (slug, result) pairs."""
return [(slug, self.tick_once_for_board(slug)) for slug in self._board_slugs()]
def ready_nonempty(self) -> bool:
"""Is there a ready+assigned+unclaimed task on ANY board that the
dispatcher would actually spawn for?
Control-plane lanes (e.g. ``orion-cc``) are pulled by terminals
via ``claim_task`` and never spawnable — a queue full of those is
"correctly idle", not "stuck". The review column is probed only
when review dispatch is on (same gate as the dispatcher): a task
waiting for a human reviewer is idle, not stuck.
"""
kb = self.kb
_review_probe = kb.review_dispatch_enabled()
for slug in self._board_slugs():
conn = None
try:
conn = kb.connect(board=slug)
if kb.has_spawnable_ready(conn):
return True
if _review_probe and kb.has_spawnable_review(conn):
return True
except Exception:
continue
finally:
if conn is not None:
with contextlib.suppress(Exception):
conn.close()
return False
def auto_decompose_tick(self, auto_decompose_per_tick: int) -> int:
"""Auto-decompose up to N triage tasks across all boards into
ready workgraphs before dispatch fans out; the per-tick cap keeps
a bulk triage load from burst-spending the aux LLM. Returns the
number decomposed/specified.
"""
try:
from hermes_cli import kanban_decompose as _decomp
except Exception as exc: # pragma: no cover
logger.warning(
"kanban auto-decompose: import failed (%s); skipping", exc,
)
return 0
attempted = 0
successes = 0
for slug in self._board_slugs():
if attempted >= auto_decompose_per_tick:
break
# Pin the board via env for the call: the decomposer connects
# with no board kwarg (same pattern as the dashboard specify endpoint).
prev_env = os.environ.get("HERMES_KANBAN_BOARD")
try:
os.environ["HERMES_KANBAN_BOARD"] = slug
try:
triage_ids = _decomp.list_triage_ids()
except Exception as exc:
logger.debug(
"kanban auto-decompose: list_triage_ids failed on board %s (%s)",
slug, exc,
)
triage_ids = []
for tid in triage_ids:
if attempted >= auto_decompose_per_tick:
break
attempted += 1
successes += self._decompose_one(_decomp, slug, tid)
finally:
if prev_env is None:
os.environ.pop("HERMES_KANBAN_BOARD", None)
else:
os.environ["HERMES_KANBAN_BOARD"] = prev_env
return successes
@staticmethod
def _decompose_one(_decomp: Any, slug: str, tid: str) -> int:
"""Decompose one triage task; returns 1 on success, 0 otherwise."""
try:
outcome = _decomp.decompose_task(tid, author="auto-decomposer")
except Exception:
logger.exception(
"kanban auto-decompose: decompose_task crashed on %s",
tid,
)
return 0
if not outcome.ok:
# Common no-op reasons (no aux client) must not spam logs every tick.
logger.debug(
"kanban auto-decompose [%s]: %s skipped: %s",
slug, tid, outcome.reason,
)
return 0
if outcome.fanout and outcome.child_ids:
logger.info(
"kanban auto-decompose [%s]: %s → %d children",
slug, tid, len(outcome.child_ids),
)
else:
logger.info(
"kanban auto-decompose [%s]: %s → single task (no fanout)",
slug, tid,
)
return 1
def _log_spawn_results(results: Optional[list]) -> bool:
"""Log per-board spawn summaries; returns whether any board spawned."""
any_spawned = False
for slug, res in (results or []):
if res is not None and getattr(res, "spawned", None):
any_spawned = True
# Quiet by default: an idle gateway stays silent.
logger.info(
"kanban dispatcher [%s]: spawned=%d reclaimed=%d "
"crashed=%d timed_out=%d promoted=%d auto_blocked=%d",
slug,
len(res.spawned),
res.reclaimed,
len(res.crashed) if hasattr(res.crashed, "__len__") else 0,
len(res.timed_out) if hasattr(res.timed_out, "__len__") else 0,
res.promoted,
len(res.auto_blocked) if hasattr(res.auto_blocked, "__len__") else 0,
)
return any_spawned
+722
View File
@@ -0,0 +1,722 @@
"""Kanban notifier: claim terminal task events per subscription and deliver them.
``GatewayKanbanWatchersMixin._kanban_notifier_watcher`` owns the loop and
the GC cadence; the per-tick claim (``_notifier_collect``) and the
per-subscription delivery (``_KanbanNotification``) live here.
"""
from __future__ import annotations
import re
from pathlib import Path
from typing import Any, Callable, Optional
from agent.i18n import t
from gateway.kanban_watchers_common import _list_boards, _to_thread_process_service, logger
# "status" covers dashboard drag-drop and `_set_status_direct()`.
# ``review_requested`` wakes the origin like a block but is not one;
# the task is not archived so later review cycles keep notifying.
TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked", "block_loop_detected", "review_requested", "changes_requested")
# Kinds that hand a decision back to the origin, which must take a turn.
# status/archived/unblocked are bookkeeping.
_WAKE_KINDS = (
"completed", "gave_up", "crashed", "timed_out",
"blocked", "review_requested", "changes_requested",
"block_loop_detected",
)
# Consecutive send failures (adapter raised OR reported
# SendResult(success=False)) before a sub is dropped as a dead chat.
# 12 ≈ 60s at the 5s cadence: a transient API outage must not
# permanently unsubscribe a live review-gate channel.
MAX_SEND_FAILURES = 12
_LOCAL_PATH_RE = re.compile(
r"(?<![\w:/])(?:/(?:Users|home|private|tmp|var|etc|workspace)/[^\s,;]+|"
r"[A-Za-z]:\\[^\s,;]+)"
)
def _safe_review_reason(value: Any, limit: int = 160) -> str:
"""Return a mobile-friendly review reason safe for external delivery."""
from agent.redact import redact_sensitive_text
reason = redact_sensitive_text(
"" if value is None else str(value),
force=True,
redact_url_credentials=True,
)
reason = _LOCAL_PATH_RE.sub("[local path]", reason)
reason = " ".join(reason.split())
if len(reason) > limit:
reason = reason[: limit - 1].rstrip() + "…"
return reason
def _wake_scope_id(adapter: Any, sub: dict) -> Optional[str]:
"""Return the tenant scope (Slack workspace) a subscription's wake keys to.
``build_session_key()`` includes ``scope_id`` on multi-tenant platforms,
so the wake must carry the same scope as inbound messages. Persisted
``delivery_metadata`` wins (it records the creating scope); the adapter's
live chat → scope map only covers rows without metadata. ``None`` means
unscoped, matching an unscoped platform's key.
"""
delivery_meta = sub.get("delivery_metadata")
if isinstance(delivery_meta, dict):
for key in ("scope_id", "slack_team_id", "team_id"):
value = delivery_meta.get(key)
if value:
return str(value)
resolver = getattr(adapter, "scope_id_for_chat", None)
if callable(resolver):
try:
resolved = resolver(str(sub.get("chat_id") or ""))
except Exception as exc:
# An adapter-side lookup failure yields no scope, never an error.
logger.debug(
"kanban notifier: scope lookup failed for chat %s: %s",
sub.get("chat_id"),
exc,
exc_info=True,
)
return None
if resolved:
return str(resolved)
return None
def _platform_names(mapping: Any) -> set[str]:
"""Lower-cased platform names of an adapters mapping (Platform enums or strings)."""
return {getattr(platform, "value", str(platform)).lower() for platform in mapping}
# ---------------------------------------------------------------------------
# Collection (runs in a worker thread)
# ---------------------------------------------------------------------------
def _notifier_collect(
runner: Any,
kb: Any,
*,
notifier_profile: Optional[str],
gc_due: bool,
gc_retention_days: int,
) -> list[dict]:
"""Claim unseen terminal events for every owned subscription on every board.
Each gateway polls only subscriptions owned by profiles whose adapters it
hosts; legacy rows without a profile stamp are visible only to the process
holding the singleton dispatcher lock.
"""
deliveries: list[dict] = []
include_unowned = runner._owns_kanban_dispatcher_lock()
profile_adapters = getattr(runner, "_profile_adapters", {})
notifier_profiles = {notifier_profile}
notifier_profiles.update(
str(profile).strip() for profile in profile_adapters if str(profile).strip()
)
active_platforms = _platform_names(runner.adapters)
# Include every platform any secondary profile has live. This is only a
# coarse pre-filter; the precise per-profile check (_authorization_adapter,
# no default fallback) runs at delivery and rewinds the claim if it
# resolves to None. An unclaimed event never retries, so dropping a
# secondary-profile sub here would lose it.
for _profile_adapter_map in profile_adapters.values():
active_platforms.update(_platform_names(_profile_adapter_map))
if not active_platforms:
logger.debug("kanban notifier: no connected adapters; skipping tick")
return deliveries
# Poll each resolved DB path once: several slugs can map to one DB when
# HERMES_KANBAN_DB pins the board path.
seen_db_paths: set[str] = set()
for board_meta in _list_boards(kb):
slug = board_meta.get("slug") or kb.DEFAULT_BOARD
db_path = board_meta.get("db_path")
try:
resolved_db_path = str(Path(db_path).expanduser().resolve()) if db_path else str(kb.kanban_db_path(slug).resolve())
except Exception:
resolved_db_path = f"slug:{slug}"
if resolved_db_path in seen_db_paths:
logger.debug(
"kanban notifier: skipping duplicate board slug %s for DB %s",
slug, resolved_db_path,
)
continue
seen_db_paths.add(resolved_db_path)
_notifier_collect_board(
kb, slug, deliveries,
notifier_profile=notifier_profile,
notifier_profiles=notifier_profiles,
include_unowned=include_unowned,
profile_adapters=profile_adapters,
active_platforms=active_platforms,
gc_due=gc_due,
gc_retention_days=gc_retention_days,
)
return deliveries
def _notifier_collect_board(
kb: Any,
slug: str,
deliveries: list[dict],
*,
notifier_profile: Optional[str],
notifier_profiles: set,
include_unowned: bool,
profile_adapters: dict,
active_platforms: set[str],
gc_due: bool,
gc_retention_days: int,
) -> None:
"""Claim events on one board, appending delivery dicts to *deliveries*."""
# Cheap read-only probe before the writable connect() (schema init, WAL
# sidecars, checkpoints) — a board with no subscriptions has nothing to notify.
try:
if kb.count_notify_subs(
board=slug,
notifier_profiles=notifier_profiles,
include_unowned=include_unowned,
) == 0:
logger.debug(
"kanban notifier: board %s has no subscriptions owned by %s; skipping open",
slug, sorted(notifier_profiles),
)
return
except Exception as exc:
logger.debug(
"kanban notifier: read-only subscription probe failed "
"for board %s (%s); falling back to writable open",
slug, exc,
)
try:
conn = kb.connect(board=slug)
except Exception as exc:
logger.debug("kanban notifier: cannot open board %s: %s", slug, exc)
return
try:
if gc_due:
# Best-effort: a failed sweep never blocks delivery; the next
# hourly gate retries.
try:
_purged = kb.purge_stale_done_notify_subs(conn, max_age_days=gc_retention_days)
if _purged:
logger.info(
"kanban notifier: purged %d stale done/blocked-task subscription(s) on board %s (retention %dd)",
_purged, slug, gc_retention_days,
)
except Exception as _gc_exc:
logger.debug(
"kanban notifier: stale-sub GC failed for board %s: %s",
slug, _gc_exc,
)
# No explicit init_db(): connect() already runs the migration once per
# process, and init_db() would re-run it on a second connection racing
# the first.
subs = kb.list_notify_subs(
conn,
notifier_profiles=notifier_profiles,
include_unowned=include_unowned,
)
if not subs:
logger.debug("kanban notifier: board %s has no subscriptions", slug)
for sub in subs:
try:
owner_profile = sub.get("notifier_profile") or None
if owner_profile and owner_profile != notifier_profile and not profile_adapters.get(owner_profile):
logger.debug(
"kanban notifier: subscription for %s owned by profile %s; current profile %s has no adapter for it, skipping",
sub.get("task_id"), owner_profile, notifier_profile,
)
continue
platform = (sub.get("platform") or "").lower()
if platform not in active_platforms:
logger.debug(
"kanban notifier: subscription for %s on %s skipped; adapter not connected",
sub.get("task_id"), platform or "<missing>",
)
continue
old_cursor, cursor, events = kb.claim_unseen_events_for_sub(
conn,
task_id=sub["task_id"],
platform=sub["platform"],
chat_id=sub["chat_id"],
thread_id=sub.get("thread_id") or "",
kinds=TERMINAL_KINDS,
)
if not events:
continue
task = kb.get_task(conn, sub["task_id"])
logger.debug(
"kanban notifier: claimed %d event(s) for %s on board %s cursor %s→%s",
len(events), sub["task_id"], slug, old_cursor, cursor,
)
deliveries.append({
"sub": sub,
"old_cursor": old_cursor,
"cursor": cursor,
"events": events,
"task": task,
"board": slug,
})
except Exception as sub_exc:
# One bad subscription must not block the rest of the tick.
logger.warning(
"kanban notifier: subscription for %s on board %s failed: %s",
sub.get("task_id"), slug, sub_exc,
)
finally:
conn.close()
# ---------------------------------------------------------------------------
# Per-event message formatting: kind -> (msg, wake_handoff, wake_review_detail)
# ---------------------------------------------------------------------------
# ``None`` for handoff / review_detail leaves the accumulated wake value
# untouched. ``_payload(ev, key)`` is the shared "payload present and truthy" read.
def _payload(ev: Any, key: str) -> Any:
return ev.payload.get(key) if ev.payload and ev.payload.get(key) else None
def _first_line(text: str, limit: int) -> str:
lines = text.strip().splitlines()
return lines[0][:limit] if lines else text[:limit]
def _fmt_completed(ev, n) -> tuple:
# Prefer the run summary from the event payload; fall back to
# task.result for legacy rows.
handoff = ""
wake_handoff = None
payload_summary = _payload(ev, "summary")
if payload_summary:
wake_handoff = _first_line(str(payload_summary), 200)
handoff = f"\n{wake_handoff}"
elif n.task and n.task.result:
wake_handoff = _first_line(n.task.result, 160)
handoff = f"\n{wake_handoff}"
msg = (
f"✔ {n.board_tag}{n.tag}Kanban {n.task_id} done"
f" — {n.title}{handoff}"
)
return msg, wake_handoff, None
def _fmt_blocked(ev, n) -> tuple:
reason = _payload(ev, "reason")
reason = f": {str(reason)[:160]}" if reason else ""
return f"⏸ {n.board_tag}{n.tag}Kanban {n.task_id} blocked{reason}", None, None
def _fmt_gave_up(ev, n) -> tuple:
err = _payload(ev, "error")
err = f"\n{str(err)[:200]}" if err else ""
msg = (
f"✖ {n.board_tag}{n.tag}Kanban {n.task_id} gave up "
f"after repeated spawn failures{err}"
)
return msg, None, None
def _fmt_crashed(ev, n) -> tuple:
msg = (
f"✖ {n.board_tag}{n.tag}Kanban {n.task_id} worker crashed "
f"(pid gone); dispatcher will retry"
)
return msg, None, None
def _fmt_timed_out(ev, n) -> tuple:
limit = _payload(ev, "limit_seconds")
limit = int(limit) if limit else 0
msg = (
f"⏱ {n.board_tag}{n.tag}Kanban {n.task_id} timed out "
f"(max_runtime={limit}s); will retry"
)
return msg, None, None
def _fmt_status(ev, n) -> tuple:
new_status = _payload(ev, "status")
new_status = str(new_status) if new_status else ""
return f"🔄 {n.board_tag}{n.tag}Kanban {n.task_id} → {new_status}", None, None
def _fmt_review_requested(ev, n) -> tuple:
# Implementation done; task moved to the review lane. Carry the handoff
# into the wake turn like ``completed`` so the reviewer needn't re-read the board.
handoff = ""
wake_handoff = None
summary = _payload(ev, "summary")
if summary:
summary = str(summary)
handoff = f"\n{summary[:200]}"
wake_handoff = _first_line(summary, 200)
msg = (
f"👀 {n.board_tag}{n.tag}Kanban {n.task_id} ready for review"
f" — {n.title}{handoff}"
)
return msg, wake_handoff, None
def _fmt_changes_requested(ev, n) -> tuple:
payload = ev.payload or {}
reason = _safe_review_reason(payload.get("reason"))
reviewer = _safe_review_reason(payload.get("reviewer"), 48)
implementer = _safe_review_reason(payload.get("implementer"), 48)
reason_text = reason or "reviewer feedback requires changes"
provenance = ""
if reviewer:
provenance += f" — reviewer @{reviewer}"
if implementer:
provenance += f" → implementer @{implementer}"
msg = (
f"🛑 {n.board_tag}Kanban {n.task_id} review requested "
f"changes/BLOCK: {reason_text}{provenance}"
)
return msg, None, reason_text
def _fmt_block_loop_detected(ev, n) -> tuple:
# Re-blocked for the same cause past the limit and routed to `triage`
# for a human. It emits no blocked/status event, so ping loudly here.
reason = _payload(ev, "reason")
reason = f": {str(reason)[:160]}" if reason else ""
recurrences = ev.payload.get("recurrences") if ev.payload else None
rc = f" (blocked {recurrences}x for the same cause)" if recurrences else ""
msg = (
f"🛑 {n.board_tag}{n.tag}Kanban {n.task_id} routed to TRIAGE"
f" — needs a human decision{rc}{reason}"
)
return msg, None, None
# archived / unblocked are claimed (so the cursor advances past them) but
# intentionally silent (no formatter), and excluded from _WAKE_KINDS so they
# never wake the creator.
_EVENT_FORMATTERS: dict[str, Callable[[Any, "_KanbanNotification"], tuple]] = {
"completed": _fmt_completed,
"blocked": _fmt_blocked,
"gave_up": _fmt_gave_up,
"crashed": _fmt_crashed,
"timed_out": _fmt_timed_out,
"status": _fmt_status,
"review_requested": _fmt_review_requested,
"changes_requested": _fmt_changes_requested,
"block_loop_detected": _fmt_block_loop_detected,
}
# ---------------------------------------------------------------------------
# Delivery of one claimed batch (one subscription, N events)
# ---------------------------------------------------------------------------
class _KanbanNotification:
"""Deliver one subscription's claimed events, then settle the cursor.
Cursor advance ordering by adapter class:
* push + notify: the text send WAS the delivery → advance now; wake
injection stays best-effort.
* non-push or wake-only: the wake IS the delivery → it runs FIRST and the
cursor advances only after it succeeds; failure rewinds like a failed
send(). An unknown platform advances the cursor so it can't replay forever.
"""
def __init__(self, runner: Any, d: dict, *, platform_cls: Any, sub_fail_counts: dict) -> None:
self.runner = runner
self.d = d
self.platform_cls = platform_cls
self.sub_fail_counts = sub_fail_counts
self.sub = sub = d["sub"]
self.task = task = d["task"]
self.board_slug = d.get("board")
self.platform_str = (sub["platform"] or "").lower()
self.task_id = sub["task_id"]
self.sub_profile = sub.get("notifier_profile") or ""
self.title = (task.title if task else sub["task_id"])[:120]
self.board_tag = f"[{self.board_slug}] " if self.board_slug else ""
# Attribute the ping to the worker that did the work.
who = task.assignee if task and task.assignee else None
self.tag = f"@{who} " if who else ""
# The wake self-post path needs the key even when every event was skipped.
self.sub_key = (sub["task_id"], sub["platform"], sub["chat_id"], sub.get("thread_id") or "")
mode = sub.get("delivery_mode") or "notify"
self.wake_agent = mode in ("notify+wake", "wake")
self.send_passive = mode != "wake"
# Worker handoff carried into the synthetic wake turn so the woken
# creator doesn't re-decompose work already on the board.
self.wake_handoff = ""
self.wake_review_detail = ""
self.plat: Any = None
self.adapter: Any = None
self.is_push_adapter = True
self.wake_kinds: set = set()
self.session_key = ""
self.synth = ""
# -- cursor / subscription ops (blocking, run in a fresh-context thread) --
async def rewind(self) -> None:
await _to_thread_process_service(
self.runner._kanban_rewind, self.sub, self.d["cursor"], self.d.get("old_cursor", 0), self.board_slug,
)
async def advance(self) -> None:
await _to_thread_process_service(self.runner._kanban_advance, self.sub, self.d["cursor"], self.board_slug)
async def unsub(self) -> None:
await _to_thread_process_service(self.runner._kanban_unsub, self.sub, self.board_slug)
def clear_failures(self) -> None:
self.sub_fail_counts.pop(self.sub_key, None)
async def delivery_failed(self, fmt: str, prefix: tuple, drop_fmt: str, exc: Exception, exc_info: bool) -> None:
"""Bump the failure counter; drop the sub past the limit, else rewind the claim so the next tick retries."""
fails = self.sub_fail_counts.get(self.sub_key, 0) + 1
self.sub_fail_counts[self.sub_key] = fails
logger.warning(fmt, *prefix, fails, MAX_SEND_FAILURES, exc, exc_info=exc_info)
if fails >= MAX_SEND_FAILURES:
logger.warning(drop_fmt, self.task_id, self.platform_str, fails)
await self.unsub()
self.clear_failures()
else:
await self.rewind()
async def _wake_failed(self, fmt: str, exc: Exception) -> None:
await self.delivery_failed(
fmt, (self.task_id,),
"kanban notifier: dropping subscription %s on %s after %d consecutive wake failures",
exc, True,
)
# -- formatting --
def format_event(self, ev: Any) -> Optional[str]:
"""Render one event; accumulates wake handoff/review detail. None → silent kind."""
formatter = _EVENT_FORMATTERS.get(ev.kind)
if formatter is None:
return None
msg, handoff, review_detail = formatter(ev, self)
if handoff is not None:
self.wake_handoff = handoff
if review_detail is not None:
self.wake_review_detail = review_detail
return msg
def build_wake_text(self) -> None:
"""Set ``wake_kinds`` / ``session_key`` / ``synth`` for the wake paths."""
task, sub = self.task, self.sub
self.wake_kinds = (
{ev.kind for ev in self.d["events"] if ev.kind in _WAKE_KINDS}
if self.wake_agent
else set()
)
if not self.wake_kinds:
return
if self.is_push_adapter:
self.session_key = getattr(task, "session_id", None) or ""
else:
# Non-push wakes target sub["chat_id"] (the raw session id the
# subscriber registered). task.session_id may be a WORKER session
# for child tasks; use it only for legacy rows.
self.session_key = sub["chat_id"] or getattr(task, "session_id", None) or ""
# i18n keys: gateway.kanban.wake.<kind> for each _WAKE_KINDS entry.
_parts = [t(f"gateway.kanban.wake.{k}") for k in _WAKE_KINDS if k in self.wake_kinds]
_status = t("gateway.kanban.wake.status_joiner").join(_parts) or t("gateway.kanban.wake.status_default")
synth = t(
"gateway.kanban.wake.message",
task_id=sub["task_id"],
status=_status,
title=self.title,
assignee=task.assignee if task else "",
board=self.board_slug,
)
# Label as an automatic notification and carry the handoff so the
# creator inspects the board instead of re-decomposing.
if self.wake_handoff:
synth += "\n" + t("gateway.kanban.wake.handoff", summary=self.wake_handoff)
if self.wake_review_detail:
synth += "\n" + t("gateway.kanban.wake.review_detail", reason=self.wake_review_detail)
self.synth = synth + "\n\n" + t("gateway.kanban.wake.guidance")
def _log_woke(self) -> None:
logger.info(
"kanban notifier: woke agent for %s on %s/%s profile=%s events=%s",
self.task_id, self.platform_str, self.sub["chat_id"], self.sub_profile or "default", self.wake_kinds,
)
async def push_wake(self) -> None:
"""Wake the creator session behind a push adapter; raises on failure."""
from gateway.session import SessionSource
from gateway.wake import deliver_wake
sub = self.sub
# Rebuild the creator's real session scope from the persisted
# chat_type: build_session_key() keys DMs differently from
# group/thread, so a hardcoded "group" mis-routed DM/thread creators
# into a fresh session. Legacy rows may carry chat_type in
# delivery_metadata; last resort is "group". A mismatch only degrades
# to a fresh session.
_chat_type = str(sub.get("chat_type") or "").strip()
if not _chat_type:
_delivery_meta = sub.get("delivery_metadata")
if isinstance(_delivery_meta, dict):
_chat_type = str(_delivery_meta.get("chat_type") or "").strip()
_source = SessionSource(
platform=self.plat,
chat_id=sub["chat_id"],
chat_type=_chat_type or "group",
thread_id=sub.get("thread_id") or None,
user_id=sub.get("user_id"),
user_id_alt=sub.get("user_id_alt"),
profile=self.sub_profile or None,
scope_id=_wake_scope_id(self.adapter, sub),
)
await deliver_wake(self.adapter, text=self.synth, session_id=self.session_key, source=_source)
self._log_woke()
async def _send_event(self, ev: Any, msg: str) -> None:
"""Send one text ping; raises on adapter exception or SendResult(success=False)."""
sub, adapter = self.sub, self.adapter
delivery_metadata = sub.get("delivery_metadata")
metadata: dict[str, Any] = dict(delivery_metadata) if isinstance(delivery_metadata, dict) else {}
if sub.get("thread_id") and not metadata.get("thread_id"):
metadata["thread_id"] = sub["thread_id"]
_send_res = await adapter.send(sub["chat_id"], msg, metadata=metadata)
# SendResult(success=False) without an exception is a FAILED delivery
# (else the event is lost); None / non-SendResult keeps the
# "no exception == delivered" contract.
if getattr(_send_res, "success", True) is False:
raise RuntimeError(
"adapter send() reported failure: "
f"{getattr(_send_res, 'error', None) or 'unknown error'}"
)
logger.debug(
"kanban notifier: delivered %s event for %s to %s/%s on board %s",
ev.kind, self.task_id, self.platform_str, sub["chat_id"], self.board_slug,
)
# Upload artifact paths from the completion payload / legacy result as
# native files. Only on ``completed`` so retries never spam attachments.
if ev.kind == "completed":
try:
await self.runner._deliver_kanban_artifacts(
adapter=adapter,
chat_id=sub["chat_id"],
metadata=metadata,
event_payload=getattr(ev, "payload", None),
task=self.task,
)
except Exception as art_exc:
logger.debug(
"kanban notifier: artifact delivery for %s failed: %s",
self.task_id, art_exc,
)
async def _send_pings(self) -> bool:
"""Send every text ping; False when a send failed (claim already rewound/dropped)."""
for ev in self.d["events"]:
msg = self.format_event(ev)
if msg is None:
continue
# Non-push adapters (api_server) always report SendResult(success=False)
# from send(); treating that as failure would drop the sub forever and
# make the wake path unreachable. Skip the doomed send; the self-post
# IS the delivery and resolves the failure counter.
if not self.is_push_adapter and self.wake_agent:
logger.debug(
"kanban notifier: adapter %s has no push "
"channel; skipping text ping for %s, relying "
"on wake self-post instead",
self.platform_str, self.task_id,
)
continue
if not self.send_passive:
# Wake-only: the wake path is the sole delivery and resolves the counter.
continue
try:
await self._send_event(ev, msg)
self.clear_failures()
except Exception as exc:
await self.delivery_failed(
"kanban notifier: send failed for %s on %s (attempt %d/%d): %s",
(self.task_id, self.platform_str),
"kanban notifier: dropping subscription %s on %s after %d consecutive send failures",
exc, False,
)
return False
return True
async def deliver(self) -> None:
try:
self.plat = self.platform_cls(self.platform_str)
except ValueError:
await self.advance()
return
# Same chokepoint as authorization: a stamped profile is served by ITS
# same-platform adapter and never falls back to the default profile's
# bot (cross-profile mis-delivery). None only when the profile (or
# default) has no adapter.
adapter = self.runner._authorization_adapter(self.plat, self.sub_profile or None)
if adapter is None:
logger.debug(
"kanban notifier: adapter %s disconnected before delivery for %s; rewinding claim",
self.platform_str, self.task_id,
)
await self.rewind()
return
self.adapter = adapter
from gateway.wake import adapter_supports_push
self.is_push_adapter = adapter_supports_push(adapter)
if not await self._send_pings():
return
# All text pings delivered (or skipped for non-push / wake-only).
task_terminal = self.task and self.task.status == "archived"
self.build_wake_text()
wake_kinds, is_push = self.wake_kinds, self.is_push_adapter
if not is_push and wake_kinds and self.session_key:
# Self-post IS the delivery: must succeed BEFORE the cursor advances.
from gateway.wake import deliver_wake
try:
await deliver_wake(adapter, text=self.synth, session_id=self.session_key)
self._log_woke()
self.clear_failures()
except Exception as _wk_err:
await self._wake_failed("kanban notifier: wake self-post failed for %s (attempt %d/%d): %s", _wk_err)
return
if is_push and not self.send_passive and wake_kinds:
# Wake-only push sub: the wake is the sole delivery and must
# succeed BEFORE the cursor advances.
try:
await self.push_wake()
self.clear_failures()
except Exception as _wk_err:
await self._wake_failed("kanban notifier: wake-only delivery failed for %s (attempt %d/%d): %s", _wk_err)
return
# Delivery complete: advance the cursor (the dedup mechanism).
await self.advance()
if not is_push:
self.clear_failures()
if is_push and self.send_passive and wake_kinds:
# notify+wake: text ping was the delivery and the cursor has
# advanced; the wake stays best-effort, but log at WARNING so a
# persistently failing wake is visible.
try:
await self.push_wake()
except Exception as _wk_err:
logger.warning(
"kanban notifier: wakeup injection failed for %s: %s",
self.task_id, _wk_err, exc_info=True,
)
# Unsubscribe only on archive; ``done`` is reversible.
if task_terminal:
await self.unsub()
+689 -2056
View File
File diff suppressed because it is too large Load Diff
+489
View File
@@ -0,0 +1,489 @@
"""Fallback delivery for GatewayStreamConsumer: continuation sends after edits
stop working, chunking, cursor cleanup, commentary and silence retraction."""
from __future__ import annotations
import asyncio
import logging
from typing import Any, Callable, Optional
from gateway.platforms.base import BasePlatformAdapter as _BasePlatformAdapter
from gateway.stream_consumer_fences import ensure_closed_code_fences
logger = logging.getLogger("gateway.stream_consumer")
class StreamFallbackMixin:
"""Non-streaming delivery paths used once progressive edits fail or the turn ends oddly."""
async def _send_new_chunk(
self,
text: str,
reply_to_id: Optional[str],
*,
final: bool = False,
) -> Optional[str]:
"""Send a new chunk threaded to ``reply_to_id``; returns the new message_id."""
text = self._clean_for_display(text)
if not text.strip():
return reply_to_id
try:
result = await self.adapter.send(
chat_id=self.chat_id,
content=text,
reply_to=reply_to_id,
metadata=self._metadata_for_send(
final=final,
expect_edits=not final,
),
)
if result.success and result.message_id:
self._message_id = str(result.message_id)
self._track_preview_ids_from_result(result)
self._already_sent = True
self._last_sent_text = text
self._notify_new_message()
return str(result.message_id)
else:
self._edit_supported = False
return reply_to_id
except Exception as e:
logger.error("Stream send chunk error: %s", e)
return reply_to_id
def _visible_prefix(self) -> str:
"""Return the visible text already shown in the streamed message."""
prefix = self._last_sent_text or ""
if self.cfg.cursor and prefix.endswith(self.cfg.cursor):
prefix = prefix[:-len(self.cfg.cursor)]
return self._clean_for_display(prefix)
def _continuation_text(self, final_text: str) -> str:
"""Return only the part of final_text the user has not already seen."""
prefix = self._fallback_prefix or self._visible_prefix()
if prefix and final_text.startswith(prefix):
return final_text[len(prefix):].lstrip()
return final_text
@staticmethod
def _split_text_chunks(
text: str,
limit: int,
len_fn: "Callable[[str], int]" = len,
) -> list[str]:
"""Split text for fallback sends: newline-preferred, fence-balanced across chunks."""
from gateway.platforms.helpers import split_text_fence_aware
return split_text_fence_aware(
text,
limit,
len_fn,
prefer_paragraphs=False,
balance_fences=True,
)
def _truncate_for_stream(
self,
text: str,
limit: int,
len_fn: "Callable[[str], int]",
) -> list[str]:
"""Split via the adapter's canonical truncate_message (platform-specific rules).
Non-base test doubles / legacy adapters keep the two-argument call shape.
"""
truncate = getattr(self.adapter, "truncate_message", None)
if not callable(truncate):
return self._split_text_chunks(text, limit, len_fn)
if isinstance(self.adapter, _BasePlatformAdapter):
chunks = truncate(text, limit, len_fn=len_fn)
else:
chunks = truncate(text, limit)
if not isinstance(chunks, (list, tuple)) or not all(
isinstance(chunk, str) for chunk in chunks
):
return self._split_text_chunks(text, limit, len_fn)
return list(chunks)
async def _send_fallback_final(self, text: str) -> None:
"""Send the final continuation after streaming edits stop working.
Retries each chunk once on flood-control failures with a short delay.
"""
final_text = self._clean_for_display(text)
# Balance fences BEFORE computing the continuation so the closing
# fence reaches the user even when only the tail is delivered.
final_text = ensure_closed_code_fences(final_text)
continuation = self._continuation_text(final_text)
self._fallback_final_send = False
if not continuation.strip():
continuation = await self._fallback_when_nothing_unseen(final_text)
if continuation is None:
return
_len_fn, raw_limit = self._fallback_len_budget()
safe_limit = max(500, raw_limit - 100)
chunks = self._split_text_chunks(continuation, safe_limit, len_fn=_len_fn)
stale_message_id = self._message_id # partial message to clean up
last_message_id: Optional[str] = None
last_successful_chunk = ""
sent_any_chunk = False
for chunk in chunks:
result = await self._send_with_flood_retry(
content=chunk, retry_log="Flood control on fallback send, retrying in %.1fs",
)
if not result or not result.success:
if sent_any_chunk:
# Partial continuation landed: do NOT set _final_response_sent
# (gateway must still deliver the full answer); _already_sent
# only prevents a duplicate of the partial.
self._already_sent = True
self._message_id = last_message_id
self._last_sent_text = last_successful_chunk
self._fallback_prefix = ""
return
# Nothing landed — let the gateway final send try once more.
self._already_sent = False
self._message_id = None
self._last_sent_text = ""
self._fallback_prefix = ""
return
sent_any_chunk = True
last_successful_chunk = chunk
last_message_id = result.message_id or last_message_id
self._notify_new_message()
# Best-effort delete of the frozen partial — ONLY when the FULL final
# was re-sent. If only the missing tail went out, the partial IS the
# head of the answer ("sent only the second half" symptom).
if (
stale_message_id
and stale_message_id != last_message_id
and not self._fallback_preserve_partial_messages
and continuation == final_text
):
await self._delete_previews(
[stale_message_id], label="Fallback partial", skip_sentinel=False,
)
self._message_id = last_message_id
self._already_sent = True
# Recorder substitutes the unsplit ledger on a split turn.
self._mark_final_delivered(record=final_text)
self._last_sent_text = chunks[-1]
self._fallback_prefix = ""
self._fallback_preserve_partial_messages = False
async def _fallback_when_nothing_unseen(self, final_text: str) -> Optional[str]:
"""Fallback entered but the visible prefix already covers ``final_text``.
Returns the continuation to send (the whole final when the prefix is
from a *previous* segment), or None when the turn is settled here.
"""
# Telegram clients can lose (part of) a streamed preview after a
# failed final edit, so opt-in adapters commit a fresh final send.
if (
final_text.strip()
and final_text == self._visible_prefix()
and getattr(self.adapter, "RESEND_FINAL_ON_EMPTY_STREAM_FALLBACK", False) is True
):
delivery = await self._send_empty_fallback_final(final_text)
if delivery == "delivered":
return None
self._already_sent = True
self._fallback_prefix = ""
self._fallback_preserve_partial_messages = False
if delivery in {"ambiguous", "preview"}:
# Timeout: Telegram may have accepted the send. Flood
# rejection: the complete ACKed preview is authoritative.
# Keep duplicate suppression in both cases.
self._final_content_delivered = True
if delivery == "preview":
# Preview already shows the full final (checked above);
# record it so the gateway doesn't re-send next to it.
self._record_turn_final_payload(final_text)
else:
self._delivery_ambiguous = True
else:
# Confirmed failure: gateway performs its normal final send.
self._final_response_sent = False
self._final_content_delivered = False
return None
# The prefix may be from a *previous* segment (before a tool
# boundary), wrongly reading as "already shown" — send final_text as-is.
if final_text.strip() and final_text != self._visible_prefix():
return final_text
# Best-effort strip of a cursor left stuck by the edit failure that
# entered fallback mode.
if (
self._message_id
and self._last_sent_text
and self.cfg.cursor
and self._last_sent_text.endswith(self.cfg.cursor)
):
clean_text = self._last_sent_text[:-len(self.cfg.cursor)]
try:
result = await self._edit_message(message_id=self._message_id, content=clean_text)
if result.success:
self._last_sent_text = clean_text
except Exception:
pass
self._already_sent = True
# Recorder substitutes the full ledger on a split turn.
self._mark_final_delivered(record=final_text)
return None
def _fallback_len_budget(self) -> "tuple[Callable[[str], int], int]":
"""(len_fn, raw_limit) for fallback chunking — per-chat cap/unit on base adapters."""
raw_limit = getattr(self.adapter, "MAX_MESSAGE_LENGTH", 4096)
_len_fn: "Callable[[str], int]" = len
if isinstance(self.adapter, _BasePlatformAdapter):
_len_fn = self.adapter.message_len_fn
# Per-chat cap/unit (relay adapter fronting N platforms).
try:
raw_limit = self.adapter.max_message_length_for_chat(self.chat_id)
_len_fn = self.adapter.message_len_fn_for_chat(self.chat_id)
except Exception as e:
logger.debug("per-chat limit resolution failed: %s", e)
return _len_fn, raw_limit
async def _send_with_flood_retry(self, *, content: str, retry_log: str, reply_to=None):
"""adapter.send(final metadata) with ONE bounded retry on flood control.
Exceptions propagate (callers decide whether a raise means "ambiguous").
Returns the last SendResult (success or not).
"""
result = None
for attempt in range(2):
kwargs = dict(
chat_id=self.chat_id,
content=content,
metadata=self._metadata_for_send(final=True),
)
if reply_to is not None:
kwargs["reply_to"] = reply_to
result = await self.adapter.send(**kwargs)
if getattr(result, "success", False):
break
retry_delay = self._fallback_flood_retry_delay(result)
if attempt == 0 and retry_delay is not None:
logger.debug(retry_log, retry_delay)
await asyncio.sleep(retry_delay)
else:
break # non-flood error, long flood wait, or second failure
return result
async def _send_empty_fallback_final(self, final_text: str) -> str:
"""Commit a completed answer after Telegram finalization fails.
Returns "delivered", "failed" (gateway may retry), "ambiguous" (a
timeout may have reached the platform), or "preview" (flood control
leaves the complete streamed preview authoritative).
"""
# Segment-scoped only: never delete an earlier finalized preamble.
stale_ids = self._stale_preview_ids(segment_only=True)
try:
result = await self._send_with_flood_retry(
content=final_text,
reply_to=self._initial_reply_to_id,
retry_log="Flood control on empty fallback final send; retrying in %.1fs",
)
except Exception as exc:
logger.debug("Empty fallback final send failed: %s", exc)
return "ambiguous" if self._send_failure_may_have_delivered(exc) else "failed"
if not getattr(result, "success", False):
if self._is_flood_error(result):
return "preview"
return "ambiguous" if self._send_failure_may_have_delivered(result) else "failed"
new_message_id = getattr(result, "message_id", None)
# Telegram reports delete failure by returning False; the flood window
# that broke the finalize can reject this too — one bounded retry.
await self._delete_previews(
stale_ids, skip=new_message_id, label="Empty fallback", retry_on_false=True,
)
self._segment_preview_message_ids = set()
self._message_id = new_message_id or "__no_edit__"
self._already_sent = True
self._mark_final_delivered()
# Record VERBATIM, not via _record_turn_final_payload: the sealed
# previews were just deleted, so the ledger (which still holds sealed
# heads) would claim delivery for text this path removed.
self._delivered_final_text = ensure_closed_code_fences(
self._clean_for_display(final_text or "")
).strip()
self._last_sent_text = final_text
self._fallback_prefix = ""
self._fallback_preserve_partial_messages = False
self._notify_new_message()
return "delivered"
@staticmethod
def _send_failure_may_have_delivered(result_or_exc: Any) -> bool:
"""Return True for timeout failures where retrying may duplicate."""
if getattr(result_or_exc, "retryable", None) is True:
return False
error = str(getattr(result_or_exc, "error", None) or result_or_exc).lower()
name = result_or_exc.__class__.__name__.lower()
return "timeout" in error or "timed out" in error or "timeout" in name
def _fallback_flood_retry_delay(self, result: Any) -> float | None:
"""Return a bounded retry delay for a fallback send, if safe to retry."""
if not self._is_flood_error(result):
return None
try:
delay = float(getattr(result, "retry_after", None) or 3.0)
except (TypeError, ValueError):
delay = 3.0
if delay > self._max_fallback_flood_retry_seconds:
logger.debug(
"Flood control requests %.1fs; leaving final delivery to the gateway",
delay,
)
return None
return max(0.0, delay)
def _is_flood_error(self, result) -> bool:
"""Check if a SendResult failure is due to flood control / rate limiting."""
err = getattr(result, "error", "") or ""
err_lower = err.lower()
return "flood" in err_lower or "retry after" in err_lower or "rate" in err_lower
async def _flush_segment_tail_on_edit_failure(self) -> None:
"""Send the unseen tail after the delivered prefix as a new message before a segment reset.
Also best-effort strips the stuck cursor from the partial message.
"""
if not self._fallback_final_send:
await self._try_strip_cursor()
visible = self._fallback_prefix or self._visible_prefix()
tail = self._accumulated
if visible and tail.startswith(visible):
tail = tail[len(visible):].lstrip()
tail = self._clean_for_display(tail)
if not tail.strip():
return
try:
# Interim: must never seal a native stream (see _send_commentary).
_md = dict(self.metadata) if self.metadata else {}
_md["_interim_send"] = True
result = await self.adapter.send(
chat_id=self.chat_id,
content=tail,
metadata=_md,
)
if result.success:
self._already_sent = True
except Exception as e:
logger.error("Segment-break tail flush error: %s", e)
async def _try_strip_cursor(self) -> None:
"""Best-effort edit removing a stuck cursor when entering fallback mode."""
if not self._message_id or self._message_id == "__no_edit__":
return
prefix = self._visible_prefix()
if not prefix or not prefix.strip():
return
try:
result = await self._edit_message(
message_id=self._message_id,
content=prefix,
)
if getattr(result, "success", False):
self._last_sent_text = prefix
except Exception:
pass # best-effort — don't let this block the fallback path
async def _send_commentary(self, text: str) -> bool:
"""Send a completed interim assistant commentary message."""
text = self._clean_for_display(text)
if not text.strip():
return False
try:
# Interim: a stream-is-the-message adapter's seal-interception must
# not turn this into draft(final=true), which would seal the live
# stream with interim text and orphan the true final.
_md = self._metadata_for_send(final=False) or {}
_md["_interim_send"] = True
# reply_to only for reply-anchored threading; Discord/Telegram use
# thread_id metadata and reply_to on every commentary is spam.
_plat = getattr(getattr(self.adapter, "platform", None), "value", None)
_platform_name = str(_plat or getattr(self.adapter, "name", "")).lower()
_needs_reply_anchor = _platform_name in ("buzz", "slack", "mattermost", "feishu")
result = await self.adapter.send(
chat_id=self.chat_id,
content=text,
reply_to=self._initial_reply_to_id if _needs_reply_anchor else None,
metadata=_md,
)
# Do NOT set _already_sent: commentary is interim, and the flag
# would suppress the real final after multiple tool calls.
if result.success:
self._notify_new_message()
# Lets run.py confirm whether an interim send carried the final.
self._delivered_commentary_texts.append(text)
return result.success
except Exception as e:
logger.error("Commentary send error: %s", e)
return False
def _raw_message_limit(self) -> int:
"""Per-chat length budget (adapter ``message_len_fn`` units) before overflow splits.
Rich-capable adapters may raise it via ``streaming_overflow_limit`` so a
reply that fits one rich message isn't fragmented at the edit limit.
"""
base = getattr(self.adapter, "MAX_MESSAGE_LENGTH", 4096)
# isinstance gate keeps MagicMock adapters (mock attrs, not ints) on base.
if isinstance(self.adapter, _BasePlatformAdapter):
try:
base = self.adapter.max_message_length_for_chat(self.chat_id)
except Exception as e:
logger.debug("max_message_length_for_chat failed: %s", e)
try:
cap = self.adapter.streaming_overflow_limit()
except Exception as e:
logger.debug("streaming_overflow_limit check failed: %s", e)
cap = None
if isinstance(cap, int) and cap > base:
return cap
return base
async def _suppress_silence_marker(self) -> None:
"""Retract any streamed preview when the final reply is a bare silence marker.
Delivery flags and ``_already_sent`` are left False: nothing was
delivered, and the gateway's whole-response filter turns the marker
into "" so no fallback send happens either.
"""
# A native-stream bubble isn't a deletable message — close an open one
# (e.g. from an eager re-seed) with an empty finalize so it doesn't hang.
if self._native_stream_opened:
try:
await self._send_frame("", finalize=True)
except Exception as e:
logger.debug(
"Silence-marker native stream close failed: %s", e,
)
self._close_native_state()
self._reopen_seeded_eagerly = False
stale_ids = self._stale_preview_ids()
await self._delete_previews(stale_ids, label="Silence-marker")
self._preview_message_ids = set()
self._message_id = None
self._accumulated = ""
self._stream_ledger = ""
self._last_sent_text = ""
self._already_sent = False
self._final_response_sent = False
self._final_content_delivered = False
self._delivered_final_text = None
self._delivery_ambiguous = False
self._turn_split_delivery = False
logger.info(
"Suppressed streamed intentional-silence marker (chat=%s)",
self.chat_id,
)
+42
View File
@@ -0,0 +1,42 @@
"""Code-fence helpers shared by the stream consumer and the gateway final send."""
from __future__ import annotations
import re
def escape_code_fences_for_display(text: str) -> str:
"""Replace each ``` with \\`\\`\\` so text can be wrapped in an outer ``` block.
Reasoning content that quotes code would otherwise break the outer fence.
"""
if not isinstance(text, str) or "```" not in text:
return text
return text.replace("```", "\\`\\`\\`")
def ensure_closed_code_fences(text: str) -> str:
"""Append a closing ``` and/or ` if the text has orphaned code markers.
Output truncated mid-code-block (token limit, finish_reason="length") would
otherwise render everything after the orphan as one code block / inline
span. Trade-off: a spurious close creates a brief empty span at the end,
far less harmful than the alternative. Odd ``` count → append a fence on
its own line; then, with complete ```…``` regions stripped, odd ` count →
append a backtick.
"""
if not isinstance(text, str) or not text:
return text
if text.count("```") % 2 == 1:
text = text.rstrip("\n") + "\n```"
# Strip complete fenced regions (and any trailing unclosed ``` that leaks
# through) so their internal backticks don't pollute the standalone count.
without_fences = re.sub(r"```.*?```", "", text, flags=re.DOTALL)
without_fences = re.sub(r"```[^`]*$", "", without_fences)
if without_fences.count("`") % 2 == 1:
text = text + "`"
return text
+149
View File
@@ -0,0 +1,149 @@
"""Think-block filtering for GatewayStreamConsumer.
Some models emit inline <think>...</think> blocks in content. The agent strips
them from the final response, but intermediate edits go out before that, so this
mirrors the CLI's _stream_delta state machine."""
from __future__ import annotations
import logging
logger = logging.getLogger("gateway.stream_consumer")
class StreamThinkFilterMixin:
"""Progressive <think>-tag suppression over streamed deltas."""
# Must stay in sync with cli.py _OPEN_TAGS/_CLOSE_TAGS and
# run_agent.py _strip_think_blocks() tag variants.
_OPEN_THINK_TAGS = (
"<REASONING_SCRATCHPAD>", "<think>", "<reasoning>",
"<THINKING>", "<thinking>", "<thought>",
)
_CLOSE_THINK_TAGS = (
"</REASONING_SCRATCHPAD>", "</think>", "</reasoning>",
"</THINKING>", "</thinking>", "</thought>",
)
def _filter_and_accumulate(self, text: str) -> None:
"""Append a delta to the buffer, discarding think blocks.
Partial tags at buffer boundaries are held in ``_think_buffer`` until
enough characters arrive to decide.
"""
buf = self._think_buffer + text
self._think_buffer = ""
while buf:
# Case-insensitive: models emit <Think>, <THINKING>, …
lower_buf = buf.lower()
if self._in_think_block:
best_idx = -1
best_len = 0
for tag in self._CLOSE_THINK_TAGS:
idx = lower_buf.find(tag.lower())
if idx != -1 and (best_idx == -1 or idx < best_idx):
best_idx = idx
best_len = len(tag)
if best_len:
self._in_think_block = False
buf = buf[best_idx + best_len:]
else:
# Hold a tail that could be a partial close tag; discard the rest.
max_tag = max(len(t) for t in self._CLOSE_THINK_TAGS)
self._think_buffer = buf[-max_tag:] if len(buf) > max_tag else buf
return
else:
# Earliest opening tag at a block boundary (start of text, or
# newline + optional whitespace) — prose that merely *mentions*
# a tag must not trigger.
best_idx = -1
best_len = 0
for tag in self._OPEN_THINK_TAGS:
tag_lower = tag.lower()
search_start = 0
while True:
idx = lower_buf.find(tag_lower, search_start)
if idx == -1:
break
# Block-boundary check (mirrors cli.py logic)
if idx == 0:
is_boundary = (
not self._accumulated
or self._accumulated.endswith("\n")
)
else:
preceding = buf[:idx]
last_nl = preceding.rfind("\n")
if last_nl == -1:
is_boundary = (
(not self._accumulated
or self._accumulated.endswith("\n"))
and preceding.strip() == ""
)
else:
is_boundary = preceding[last_nl + 1:].strip() == ""
if is_boundary and (best_idx == -1 or idx < best_idx):
best_idx = idx
best_len = len(tag)
break # first boundary hit for this tag is enough
search_start = idx + 1
if best_len:
self._append_accumulated(buf[:best_idx])
self._in_think_block = True
buf = buf[best_idx + best_len:]
else:
# Hold back a partial open tag at the tail.
held_back = 0
for tag in self._OPEN_THINK_TAGS:
tag_lower = tag.lower()
for i in range(1, len(tag)):
if lower_buf.endswith(tag_lower[:i]) and i > held_back:
held_back = i
if held_back:
self._append_accumulated(buf[:-held_back])
self._think_buffer = buf[-held_back:]
else:
# An orphan </think> (thinking-mode toggle dropped the
# open, or incomplete upstream stripping) is noise.
self._append_accumulated(self._strip_orphan_close_tags(buf))
return
@classmethod
def _strip_orphan_close_tags(cls, text: str) -> str:
"""Remove close tags (plus trailing whitespace) that have no matching open.
Mirrors ``agent/think_scrubber.py::StreamingThinkScrubber`` so the
progressive display matches the post-stream scrubber.
"""
if "</" not in text:
return text
text_lower = text.lower()
out: list[str] = []
i = 0
while i < len(text):
matched = False
if text_lower[i:i + 2] == "</":
for tag in cls._CLOSE_THINK_TAGS:
tag_lower = tag.lower()
tag_len = len(tag_lower)
if text_lower[i:i + tag_len] == tag_lower:
j = i + tag_len
while j < len(text) and text[j] in " \t\n\r":
j += 1
i = j
matched = True
break
if not matched:
out.append(text[i])
i += 1
return "".join(out)
def _flush_think_buffer(self) -> None:
"""On stream end, flush text held back waiting for a possible open tag."""
if self._think_buffer and not self._in_think_block:
self._append_accumulated(self._strip_orphan_close_tags(self._think_buffer))
self._think_buffer = ""
+676
View File
@@ -0,0 +1,676 @@
"""Transport layer of GatewayStreamConsumer: native frames, drafts, edit/send.
Mixin methods use only ``self`` state; see gateway/stream_consumer.py for the
state model and the drain loop that calls into these."""
from __future__ import annotations
import asyncio
import inspect
import logging
import time
from typing import Any, Optional
from gateway.platforms.base import BasePlatformAdapter as _BasePlatformAdapter
from gateway.stream_consumer_fences import ensure_closed_code_fences
logger = logging.getLogger("gateway.stream_consumer")
class StreamTransportMixin:
"""Send/edit/frame primitives and the transport-ordered ``_send_or_edit``."""
_MIN_NEW_MSG_CHARS = 4
async def _edit_message(
self,
*,
message_id: str,
content: str,
finalize: bool = False,
):
"""Edit via the adapter, passing routing metadata when supported."""
# Contract: adapters must accept finalize= even when False (test-guarded).
kwargs = {
"chat_id": self.chat_id,
"message_id": message_id,
"content": content,
"finalize": finalize,
}
if self.metadata:
try:
params = inspect.signature(self.adapter.edit_message).parameters
if "metadata" in params or any(
param.kind is inspect.Parameter.VAR_KEYWORD
for param in params.values()
):
kwargs["metadata"] = self.metadata
except (TypeError, ValueError):
pass
return await self.adapter.edit_message(**kwargs)
async def _send_seed_frame(self):
"""Open a native stream with an empty seed frame (typing indicator before any token)."""
return await self.adapter.send_stream_frame(
"",
chat_id=self.chat_id,
reply_to=self._initial_reply_to_id,
turn_id=self._turn_id,
)
async def _send_frame(self, text: str, *, finalize: bool):
"""One native-stream frame; every frame carries the same chat/reply/turn routing."""
return await self.adapter.send_stream_frame(
text,
finalize=finalize,
chat_id=self.chat_id,
reply_to=self._initial_reply_to_id,
turn_id=self._turn_id,
)
def _close_native_state(self) -> None:
"""Mark the native stream closed (next content re-seeds or falls back)."""
self._native_stream_opened = False
self._native_last_pushed_len = 0
def _degrade_native_to_buffered_send(self) -> None:
"""Leave native mode; post-boundary output goes out as ONE send() at got_done.
buffer_only avoids mid-stream flushes that create multiple messages on
non-editable platforms.
"""
self._use_native_streaming = False
self._close_native_state()
self.cfg.buffer_only = True
def _draft_metadata(self) -> dict | None:
"""Draft-frame metadata.
Every frame must carry the same reply_to_message_id the final send gets
from _metadata_for_send: the relay adapter keys draft/seal state on it,
else the final can't find the open stream (flat DMs have no thread
metadata and would key on the bare chat).
"""
md = dict(self.metadata) if self.metadata else {}
if self._initial_reply_to_id:
md.setdefault("reply_to_message_id", self._initial_reply_to_id)
return md or None
def _stale_preview_ids(self, *, segment_only: bool = False) -> set:
"""Preview message ids a fresh final replaces.
``segment_only``: never delete an earlier finalized preamble.
"""
if segment_only:
stale_ids = set(self._segment_preview_message_ids)
if self._message_id and self._message_id != "__no_edit__":
stale_ids.add(str(self._message_id))
return stale_ids
stale_ids = set(self._preview_message_ids)
if self._message_id and self._message_id != "__no_edit__":
stale_ids.add(self._message_id)
return stale_ids
async def _delete_previews(
self, stale_ids, *, skip=None, label: str, retry_on_false: bool = False,
skip_sentinel: bool = True,
) -> None:
"""Best-effort delete of stale previews; never the message just sent (``skip``)."""
delete_fn = getattr(self.adapter, "delete_message", None)
if delete_fn is None:
return
for stale_id in stale_ids:
if not stale_id or stale_id == skip or (skip_sentinel and stale_id == "__no_edit__"):
continue
try:
deleted = await delete_fn(self.chat_id, stale_id)
if retry_on_false and deleted is False:
await asyncio.sleep(1.0)
await delete_fn(self.chat_id, stale_id)
except Exception as e:
logger.debug("%s preview cleanup failed (%s): %s", label, stale_id, e)
def _resolve_draft_streaming(self) -> bool:
"""Whether this run should use draft streaming per ``cfg.transport``.
"edit"/"off" → False. "draft"/"auto" → the adapter's
supports_draft_streaming probe (chat type, platform-version gates);
"draft" logs the downgrade when unsupported.
"""
transport = (self.cfg.transport or "edit").lower()
if transport in ("edit", "off"):
return False
# MagicMock test adapters default to edit.
if not isinstance(self.adapter, _BasePlatformAdapter):
return False
try:
try:
# Per-chat probe (relay adapters resolve through the CHAT's
# descriptor); older adapters without the kwarg keep the legacy probe.
supported = self.adapter.supports_draft_streaming(
chat_type=self.cfg.chat_type or None,
metadata=self.metadata,
chat_id=self.chat_id,
)
except TypeError:
supported = self.adapter.supports_draft_streaming(
chat_type=self.cfg.chat_type or None,
metadata=self.metadata,
)
except Exception:
logger.debug("supports_draft_streaming probe raised", exc_info=True)
supported = False
if not supported:
if transport == "draft":
logger.debug(
"Draft streaming requested but unsupported (chat=%s, type=%r) — "
"falling back to edit",
self.chat_id, self.cfg.chat_type,
)
return False
return True
def _resolve_native_streaming(self) -> bool:
"""Whether to use native streaming (adapter.send_stream_frame for ALL frames).
Requires a BasePlatformAdapter subclass with class-level
SUPPORTS_NATIVE_STREAMING and a truthy supports_native_streaming probe.
"""
if not isinstance(self.adapter, _BasePlatformAdapter):
return False
if not getattr(type(self.adapter), "SUPPORTS_NATIVE_STREAMING", False):
return False
probe = getattr(self.adapter, "supports_native_streaming", None)
if probe is None:
return False
try:
supported = probe(
chat_type=self.cfg.chat_type or None,
metadata=self.metadata,
)
except Exception:
logger.debug(
"supports_native_streaming probe raised", exc_info=True,
)
return False
return bool(supported)
async def _send_draft_frame(self, text: str) -> bool:
"""Emit one draft frame; any failure permanently disables drafts for this run.
Drafts have no message_id and clear on the client when the final
sendMessage lands.
"""
if self._draft_id is None:
# Should never happen (set in tandem with _use_draft_streaming in run()).
self._use_draft_streaming = False
return False
# Every frame must carry the same reply_to_message_id the final send
# gets from _metadata_for_send: the relay adapter keys draft/seal state
# on it, else the final can't find the open stream (flat DMs have no
# thread metadata and would key on the bare chat).
try:
result = await self.adapter.send_draft(
chat_id=self.chat_id,
draft_id=self._draft_id,
content=text,
metadata=self._draft_metadata(),
)
except Exception as e:
logger.debug(
"send_draft raised, disabling draft transport for this run: %s", e,
)
self._draft_failures += 1
self._use_draft_streaming = False
return False
if not getattr(result, "success", False):
logger.debug(
"send_draft returned success=False, disabling draft transport: %s",
getattr(result, "error", "unknown"),
)
self._draft_failures += 1
self._use_draft_streaming = False
return False
self._last_sent_text = text # parity with the edit-based no-op skip
return True
async def _abandon_native_stream(self) -> None:
"""Seal an orphaned draft stream in place on turn death (stale exit / cancel).
Otherwise the message keeps its live indicator forever and the
adapter's armed interception state leaks into the next turn. Never
sets delivery flags — the gateway's normal paths own what happens next.
"""
if not self._use_draft_streaming:
return
abandon = getattr(type(self.adapter), "abandon_open_draft", None)
if abandon is None:
return
try:
await self.adapter.abandon_open_draft(
self.chat_id,
self._last_sent_text or self._clean_for_display(self._accumulated),
metadata=self._draft_metadata(),
)
except Exception as e:
logger.debug("abandon_open_draft failed (best-effort): %s", e)
def _should_send_fresh_final(self) -> bool:
"""True when fresh-final is enabled and a real preview has been visible ≥ threshold."""
threshold = getattr(self.cfg, "fresh_final_after_seconds", 0.0) or 0.0
if threshold <= 0:
return False
if not self._message_id or self._message_id == "__no_edit__":
return False
if self._message_created_ts is None:
return False
age = time.monotonic() - self._message_created_ts
return age >= threshold
def _track_preview_id(self, message_id: Optional[str]) -> None:
"""Record a real preview message id for finalization cleanup."""
if message_id and message_id != "__no_edit__":
message_id = str(message_id)
self._preview_message_ids.add(message_id)
self._segment_preview_message_ids.add(message_id)
def _track_preview_ids_from_result(self, result: Any) -> None:
"""Record the primary id plus any continuation ids from an oversized split."""
self._track_preview_id(getattr(result, "message_id", None))
for mid in (getattr(result, "continuation_message_ids", None) or ()):
self._track_preview_id(mid)
raw = getattr(result, "raw_response", None) or {}
if isinstance(raw, dict):
for mid in (raw.get("message_ids") or ()):
self._track_preview_id(mid)
def _adapter_prefers_fresh_final(self, text: str) -> bool:
"""Adapter's prefers_fresh_final_streaming hook (e.g. Telegram's richer send path).
False when there's no real preview, no hook, or on any error.
"""
if not self._message_id or self._message_id == "__no_edit__":
return False
fn = getattr(self.adapter, "prefers_fresh_final_streaming", None)
if fn is None:
return False
try:
try:
# chat_id lets relay adapters decide via THIS chat's platform;
# otherwise a Slack-primary relay misroutes fronted chats
# through the fresh-send lane (duplicates: no delete op).
result = fn(text, metadata=self.metadata, chat_id=self.chat_id)
except TypeError:
try:
result = fn(text, metadata=self.metadata) # single-platform signature
except TypeError:
result = fn(text) # test doubles without the metadata kwarg
except Exception as e:
logger.debug("prefers_fresh_final_streaming check failed: %s", e)
return False
# ``is True`` keeps MagicMock auto-children from enabling fresh-final.
return result is True
async def _try_fresh_final(self, text: str, *, is_turn_final: bool = True) -> bool:
"""Send ``text`` as a fresh message and best-effort delete the preview(s).
Returns False on any failure so the caller falls back to edit.
``is_turn_final=False`` (interim segment at a tool boundary) leaves the
final-delivery flag unset so the gateway still delivers the real answer.
"""
# Replacing every tracked preview is only sound while ``text`` holds
# the whole answer; after a split, deleting sealed heads would erase
# delivered text — take the edit path instead.
if self._turn_split_delivery:
return False
stale_ids = self._stale_preview_ids()
try:
result = await self.adapter.send(
chat_id=self.chat_id,
content=text,
metadata=self._metadata_for_send(final=True),
)
except Exception as e:
logger.debug("Fresh-final send failed, falling back to edit: %s", e)
return False
if not getattr(result, "success", False):
return False
new_message_id = getattr(result, "message_id", None)
# Best-effort preview cleanup; never delete the message just sent.
await self._delete_previews(stale_ids, skip=new_message_id, label="Fresh-final")
self._preview_message_ids = set()
if new_message_id:
self._message_id = new_message_id
self._message_created_ts = time.monotonic()
else:
# No id returned: sentinel so we never try to edit it.
self._message_id = "__no_edit__"
self._message_created_ts = None
self._already_sent = True
self._last_sent_text = text
if is_turn_final:
self._final_response_sent = True
self._record_turn_final_payload(text)
return True
async def _send_or_edit(
self, text: str, *, finalize: bool = False, is_turn_final: bool = True,
) -> bool:
"""Send or edit the streaming message; True if delivered.
``finalize`` marks the last edit of a streaming sequence. Callers such
as the overflow split loop use the result to decide whether to advance.
Transport order: native frame → draft frame → edit existing → first
send; a transport returns None to fall through to the next.
"""
text = self._clean_for_display(text)
# Stream-is-the-message draft frames must stay prefix-stable: a closing
# ``` appended to a mid-code-block frame makes frame N not a prefix of
# N+1 and the connector re-appends the whole snapshot. The final
# message is still fence-closed below.
pre_fence_text = text
text = ensure_closed_code_fences(text)
# A bare cursor renders as a stray tofu box on some clients.
visible_stripped = text
if self.cfg.cursor:
visible_stripped = visible_stripped.replace(self.cfg.cursor, "")
visible_stripped = visible_stripped.strip()
if not visible_stripped:
# Native streams MUST still get a finalize frame (placeholder) to
# close the thinking bubble, e.g. for a MEDIA-only response.
if finalize and self._use_native_streaming and self._native_stream_opened:
try:
if await self._send_frame("✅", finalize=True):
self._mark_final_delivered()
except Exception as e:
logger.debug("Finalize empty stream failed: %s", e)
return True # cursor-only / whitespace-only update
if not text.strip():
return True # nothing to send is "success"
# Don't open a new message for 1-2 tokens + cursor (rapid tool-calling):
# if the cursor-strip edit is then rate-limited, "X ▉" stays forever.
# Only first sends are gated.
if (
self._message_id is None
and self.cfg.cursor
and self.cfg.cursor in text
and len(visible_stripped) < self._MIN_NEW_MSG_CHARS
):
return True # too short for a standalone message — accumulate more
if self._use_native_streaming:
ok = await self._native_push(text, finalize=finalize, is_turn_final=is_turn_final)
if ok is not None:
return ok
# Fall through so accumulated text still reaches the user via edit/send.
if self._use_draft_streaming and self._message_id is None:
ok = await self._draft_push(
text, pre_fence_text, finalize=finalize, is_turn_final=is_turn_final,
)
if ok is not None:
return ok
# Failure disabled drafts; fall through to edit/send.
self._last_edit_overflowed = False
try:
if self._message_id is None:
return await self._first_send(text, finalize=finalize)
if not self._edit_supported:
return False # edits unsupported; fallback path sends the final
return await self._edit_existing(text, finalize=finalize, is_turn_final=is_turn_final)
except Exception as e:
logger.error("Stream send/edit error: %s", e)
return False
async def _native_push(self, text: str, *, finalize: bool, is_turn_final: bool) -> Optional[bool]:
"""Native streaming: every frame goes through send_stream_frame().
The adapter's send/edit paths are not touched in this mode. Lazy
re-seed here after a boundary closed the stream. Returns None when
native was disabled (seed/frame failure) so the caller falls through.
"""
if not self._native_stream_opened and text:
try:
if await self._send_seed_frame():
self._native_stream_opened = True
self._awaiting_reopen_after_boundary = False
# Paired with the boundary-finalize INFO: typing-reappear latency.
logger.info(
"[latency] Re-opened native stream after boundary "
"(turn=%s, waited for first delta)",
self._turn_id,
)
else:
self._use_native_streaming = False
except Exception as e:
logger.debug("Re-seed failed, disabling native streaming: %s", e)
self._use_native_streaming = False
if not self._use_native_streaming:
return None
# WeCom renders each finalize as a separate bubble: only the
# turn-final and boundaries close the stream, not segment breaks.
if finalize and not is_turn_final:
finalize = False
if not finalize and text == self._last_sent_text:
return True # unchanged — skip
# Mark a finalize frame delivered OPTIMISTICALLY, before the ack wait:
# the bytes hit the wire (and WeCom renders them) before the ack, so a
# gateway join-cancel during the ack wait must not strand
# final_content_delivered=False and cause a duplicate normal send
# (docs/rca-wecom-stream-final-ack-timeout-duplicate.md). A definitive
# dispatch failure rolls the mark back below. Residual window (cancel
# between mark and wire write, sub-ms) is accepted.
if finalize:
# Recorded so a stale/partial frame can't suppress the corrective send.
self._mark_final_delivered(record=text)
try:
ok = await self._send_frame(text, finalize=finalize)
except Exception as e:
logger.debug("send_stream_frame raised, disabling native streaming: %s", e)
ok = False
if ok:
self._already_sent = True
self._last_sent_text = text
self._native_last_pushed_len = len(text)
if finalize:
self._mark_final_delivered()
return True
# Definitive failure: roll back the optimistic mark so the edit/send
# fallback delivers exactly once.
if finalize:
self._final_response_sent = False
self._final_content_delivered = False
self._delivered_final_text = None
# Subsequent frames take the edit/send fallback; the adapter marks the
# chat expired so it doesn't retry the dead stream.
self._use_native_streaming = False
# Best-effort close of an opened bubble (seed frame has zero length but
# still opens it — hence _native_stream_opened, not pushed_len). DO NOT
# mark delivered: the frame closes the bubble but WeCom may not render
# the content (errcode 6000 race).
if self._native_stream_opened:
try:
await self._send_frame(text, finalize=True)
logger.debug("Native fallback: finalized stream (best-effort close)")
except Exception as e:
logger.debug("Native fallback: failed to finalize stream: %s", e)
return None
async def _draft_push(
self, text: str, pre_fence_text: str, *, finalize: bool, is_turn_final: bool,
) -> Optional[bool]:
"""Draft frame while no message_id exists; None = not applicable / drafts just failed.
Drafts have no message_id: the final answer goes through the regular
send (which clears the draft client-side), so drafts are skipped when
finalizing. Exception: stream-is-the-message adapters keep ONE stream
per turn, so a segment-break finalize must NOT become a real send
(seal interception would seal the stream at every tool boundary);
only got_done seals.
"""
stream_is_msg = self._stream_is_message()
if finalize and not (stream_is_msg and not is_turn_final):
return None
frame_text = pre_fence_text if stream_is_msg else text
# Strip the cursor: native streams render their own indicator, and
# "...text▉" is never a prefix of "...text more▉", which forces the
# connector's whole-text re-append on EVERY tick (stacked copies).
if self.cfg.cursor and frame_text.endswith(self.cfg.cursor):
frame_text = frame_text[: -len(self.cfg.cursor)]
if frame_text == self._last_sent_text:
return True
if await self._send_draft_frame(frame_text):
# Deliberately NOT _already_sent: the gateway's fallback final send
# must still fire so the user gets a real message.
return True
return None
async def _first_send(self, text: str, *, finalize: bool) -> bool:
"""First send, threaded to the user's message (correct topic/thread)."""
result = await self.adapter.send(
chat_id=self.chat_id,
content=text,
reply_to=self._initial_reply_to_id,
metadata=self._metadata_for_send(final=finalize, expect_edits=not finalize),
)
if not result.success:
self._edit_supported = False
return False
if result.message_id:
self._message_id = result.message_id
self._message_created_ts = time.monotonic()
self._track_preview_ids_from_result(result)
else:
self._edit_supported = False
self._already_sent = True
self._last_sent_text = text
if not result.message_id:
self._fallback_prefix = self._visible_prefix()
self._fallback_final_send = True
# Sentinel: no editable id, don't re-enter first-send on every
# delta/tool boundary.
self._message_id = "__no_edit__"
self._notify_new_message()
return True
async def _edit_existing(self, text: str, *, finalize: bool, is_turn_final: bool) -> bool:
"""Edit the live preview (or replace it via fresh-final when finalizing)."""
# REQUIRES_EDIT_FINALIZE adapters need the finalize=True edit even
# when unchanged; everyone else short-circuits.
if text == self._last_sent_text and not (finalize and self._adapter_requires_finalize):
return True
# Fresh-final: replace a long-lived preview with a fresh message
# (timestamp reflects completion), or whenever the adapter prefers it
# (Telegram's send path renders richer markdown than its edit path).
# An explicit hook returning False must NOT be overridden by the time
# threshold — on Telegram both messages would stay on screen since the
# delete is best-effort. Check the CLASS (MagicMock auto-creates
# attrs) plus instance __dict__ (test doubles assign the hook explicitly).
has_prefers_hook = (
hasattr(type(self.adapter), "prefers_fresh_final_streaming")
or "prefers_fresh_final_streaming" in getattr(self.adapter, "__dict__", {})
)
prefers_fresh = self._adapter_prefers_fresh_final(text)
if (
finalize
and (prefers_fresh or (not has_prefers_hook and self._should_send_fresh_final()))
and await self._try_fresh_final(text, is_turn_final=is_turn_final)
):
return True
result = await self._edit_message(
message_id=self._message_id, content=text, finalize=finalize,
)
if not result.success:
return await self._on_edit_failure(result, text, finalize=finalize, is_turn_final=is_turn_final)
self._already_sent = True
self._track_preview_ids_from_result(result)
# Oversized edit split across continuations: message_id is now the
# LAST continuation, which holds only the final chunk — retarget edits
# and reset skip-if-same. getattr keeps SimpleNamespace test mocks working.
continuation_ids = getattr(result, "continuation_message_ids", ()) or ()
if continuation_ids and result.message_id and result.message_id != self._message_id:
self._last_edit_overflowed = True
self._turn_split_delivery = True
self._message_id = str(result.message_id)
self._message_created_ts = time.monotonic()
self._last_sent_text = ""
self._notify_new_message()
else:
self._last_sent_text = text
self._flood_strikes = 0
return True
async def _on_edit_failure(self, result, text: str, *, finalize: bool, is_turn_final: bool) -> bool:
"""Classify a failed edit: partial overflow, flood backoff, or fallback mode.
Returns the _send_or_edit result (always False here; the caller's
finalize path may still deliver the tail via _send_fallback_final).
"""
immediate_final_fallback = False
if (
finalize
and is_turn_final
and self.cfg.cursor
and self._last_sent_text.endswith(self.cfg.cursor)
and self._visible_prefix() == text
):
# Cosmetic final edit was rate-limited but the full answer is
# already on screen (cursor stuck): mark delivered so the gateway
# doesn't send it twice, and record the on-screen payload.
self._final_content_delivered = True
self._record_turn_final_payload(text)
raw_response = getattr(result, "raw_response", None)
if isinstance(raw_response, dict) and raw_response.get("partial_overflow"):
# Some overflow chunks landed but not the whole response: preserve
# the visible prefix so got_done sends the missing tail.
self._message_id = str(
raw_response.get("last_message_id") or result.message_id or self._message_id
)
delivered_prefix = raw_response.get("delivered_prefix")
if isinstance(delivered_prefix, str) and delivered_prefix:
self._last_sent_text = delivered_prefix
self._fallback_prefix = delivered_prefix
self._fallback_preserve_partial_messages = text.startswith(delivered_prefix)
else:
self._fallback_prefix = self._visible_prefix()
self._fallback_preserve_partial_messages = False
self._fallback_final_send = True
self._edit_supported = False
self._already_sent = True
if getattr(result, "continuation_message_ids", ()):
self._notify_new_message()
return False
# Flood control: adaptive backoff (double the interval); disable edits
# only after _MAX_FLOOD_STRIKES in a row.
if self._is_flood_error(result):
self._flood_strikes += 1
self._current_edit_interval = min(self._current_edit_interval * 2, 10.0)
logger.debug(
"Flood control on edit (strike %d/%d), backoff interval → %.1fs",
self._flood_strikes, self._MAX_FLOOD_STRIKES, self._current_edit_interval,
)
immediate_final_fallback = (
finalize
and is_turn_final
and getattr(self.adapter, "FALLBACK_ON_FINAL_EDIT_FLOOD", False) is True
)
if self._flood_strikes < self._MAX_FLOOD_STRIKES and not immediate_final_fallback:
self._last_edit_time = time.monotonic() # honor the new interval
return False
if immediate_final_fallback:
logger.debug("Turn-final edit hit flood control; entering fallback immediately")
# Fallback mode: send only the missing tail at got_done.
logger.debug("Edit failed (strikes=%d), entering fallback mode", self._flood_strikes)
self._fallback_prefix = self._visible_prefix()
self._fallback_final_send = True
self._edit_supported = False
self._already_sent = True
# A turn-final flood skips the cosmetic cursor strip: it would burn the
# same flood budget and delay the answer.
if not immediate_final_fallback:
await self._try_strip_cursor()
return False