From d90fca729374db24842ccdb159da55d0769be8f6 Mon Sep 17 00:00:00 2001 From: oliver Date: Thu, 2 Jul 2026 17:43:57 +0800 Subject: [PATCH] fix(whatsapp): add outbound retry backoff for failed sends Track send attempts and schedule retry windows for pending WhatsApp outbound messages, then mark as failed after max retries to avoid silent drops without infinite resend loops. Co-authored-by: Cursor --- svc/persistence/sqlite_store.py | 71 +++++++++++++++++++++++----- tests/test_channel_outbound_retry.py | 56 ++++++++++++++++++++++ 2 files changed, 116 insertions(+), 11 deletions(-) create mode 100644 tests/test_channel_outbound_retry.py diff --git a/svc/persistence/sqlite_store.py b/svc/persistence/sqlite_store.py index 560ab308..96e8dadd 100644 --- a/svc/persistence/sqlite_store.py +++ b/svc/persistence/sqlite_store.py @@ -741,6 +741,7 @@ class SqliteStore(ScheduledJobStoreMixin): ); """ ) + self._ensure_channel_outbound_retry_columns(conn) conn.execute( "CREATE INDEX IF NOT EXISTS idx_channel_outbound_pending ON channel_outbound_message(channel, account_id, status, created_at)" ) @@ -1397,6 +1398,7 @@ class SqliteStore(ScheduledJobStoreMixin): ) """ ) + self._ensure_channel_outbound_retry_columns(conn) conn.execute( "CREATE INDEX IF NOT EXISTS idx_channel_outbound_pending ON channel_outbound_message(channel, account_id, status, created_at)" ) @@ -5968,6 +5970,28 @@ class SqliteStore(ScheduledJobStoreMixin): if "phone" not in cols: conn.execute("ALTER TABLE whatsapp_access_pending ADD COLUMN phone TEXT NOT NULL DEFAULT ''") + def _ensure_channel_outbound_retry_columns(self, conn: Any) -> None: + if self._use_pg: + conn.execute( + "ALTER TABLE channel_outbound_message ADD COLUMN IF NOT EXISTS send_attempts INTEGER NOT NULL DEFAULT 0" + ) + conn.execute( + "ALTER TABLE channel_outbound_message ADD COLUMN IF NOT EXISTS next_attempt_at TEXT" + ) + conn.execute( + "ALTER TABLE channel_outbound_message ADD COLUMN IF NOT EXISTS last_attempt_at TEXT" + ) + return + cols = {row[1] for row in conn.execute("PRAGMA table_info(channel_outbound_message)").fetchall()} + if "send_attempts" not in cols: + conn.execute( + "ALTER TABLE channel_outbound_message ADD COLUMN send_attempts INTEGER NOT NULL DEFAULT 0" + ) + if "next_attempt_at" not in cols: + conn.execute("ALTER TABLE channel_outbound_message ADD COLUMN next_attempt_at TEXT") + if "last_attempt_at" not in cols: + conn.execute("ALTER TABLE channel_outbound_message ADD COLUMN last_attempt_at TEXT") + def create_whatsapp_access_pending( self, *, @@ -6369,10 +6393,11 @@ class SqliteStore(ScheduledJobStoreMixin): SELECT id, tenant_id, channel, account_id, chat_id, text, source, created_at FROM channel_outbound_message WHERE channel = ? AND account_id = ? AND status = 'pending' + AND (next_attempt_at IS NULL OR next_attempt_at = '' OR next_attempt_at <= ?) ORDER BY created_at ASC LIMIT ? """, - (str(channel), str(account_id), lim), + (str(channel), str(account_id), utc_now_iso(), lim), ).fetchall() return [ { @@ -6446,15 +6471,16 @@ class SqliteStore(ScheduledJobStoreMixin): error: str = "", stanza_id: str = "", ) -> bool: + max_attempts = max(1, int(os.getenv("AIA_OUTBOUND_MAX_RETRIES", "3"))) + backoff_schedule = [5, 15, 45, 120, 300] ts = utc_now_iso() - status = "sent" if ok else "failed" mid = str(message_id or "").strip() if not mid: return False with self._connect() as conn: row = conn.execute( """ - SELECT tenant_id, account_id, chat_id, text, source + SELECT tenant_id, account_id, chat_id, text, source, send_attempts FROM channel_outbound_message WHERE id = ? AND status = 'pending' """, @@ -6462,14 +6488,37 @@ class SqliteStore(ScheduledJobStoreMixin): ).fetchone() if not row: return False - cur = conn.execute( - """ - UPDATE channel_outbound_message - SET status = ?, sent_at = ?, error = ? - WHERE id = ? AND status = 'pending' - """, - (status, ts, str(error or ""), mid), - ) + attempt_no = int(row["send_attempts"] or 0) + 1 + if ok: + cur = conn.execute( + """ + UPDATE channel_outbound_message + SET status = 'sent', sent_at = ?, error = '', send_attempts = ?, last_attempt_at = ?, next_attempt_at = NULL + WHERE id = ? AND status = 'pending' + """, + (ts, attempt_no, ts, mid), + ) + else: + if attempt_no >= max_attempts: + cur = conn.execute( + """ + UPDATE channel_outbound_message + SET status = 'failed', sent_at = ?, error = ?, send_attempts = ?, last_attempt_at = ?, next_attempt_at = NULL + WHERE id = ? AND status = 'pending' + """, + (ts, str(error or ""), attempt_no, ts, mid), + ) + else: + delay_sec = backoff_schedule[min(attempt_no - 1, len(backoff_schedule) - 1)] + next_attempt_at = (datetime.now(timezone.utc) + timedelta(seconds=delay_sec)).isoformat() + cur = conn.execute( + """ + UPDATE channel_outbound_message + SET status = 'pending', error = ?, send_attempts = ?, last_attempt_at = ?, next_attempt_at = ? + WHERE id = ? AND status = 'pending' + """, + (str(error or ""), attempt_no, ts, next_attempt_at, mid), + ) changed = bool(cur.rowcount and cur.rowcount > 0) if changed and ok and str(stanza_id or "").strip(): source_raw = str(row["source"] or "").strip() diff --git a/tests/test_channel_outbound_retry.py b/tests/test_channel_outbound_retry.py new file mode 100644 index 00000000..4fd26e94 --- /dev/null +++ b/tests/test_channel_outbound_retry.py @@ -0,0 +1,56 @@ +from __future__ import annotations + +from svc.persistence.sqlite_store import SqliteStore + + +def _row(store: SqliteStore, msg_id: str) -> dict: + with store._connect() as conn: # noqa: SLF001 - test introspection + got = conn.execute( + """ + SELECT id, status, error, send_attempts, next_attempt_at, last_attempt_at + FROM channel_outbound_message + WHERE id = ? + """, + (msg_id,), + ).fetchone() + assert got is not None + return dict(got) + + +def test_channel_outbound_retry_transitions(tmp_path) -> None: + store = SqliteStore(str(tmp_path / "ops.sqlite")) + msg_id = store.enqueue_channel_outbound_message( + channel="whatsapp", + account_id="wa-default", + chat_id="8615601877957@s.whatsapp.net", + text="hello", + ) + + assert store.ack_channel_outbound_message(message_id=msg_id, ok=False, error="net down") is True + first = _row(store, msg_id) + assert first["status"] == "pending" + assert int(first["send_attempts"] or 0) == 1 + assert str(first["next_attempt_at"] or "").strip() + + assert store.ack_channel_outbound_message(message_id=msg_id, ok=False, error="still down") is True + assert store.ack_channel_outbound_message(message_id=msg_id, ok=False, error="final fail") is True + final = _row(store, msg_id) + assert final["status"] == "failed" + assert int(final["send_attempts"] or 0) == 3 + assert not str(final["next_attempt_at"] or "").strip() + + +def test_channel_outbound_retry_success_clears_schedule(tmp_path) -> None: + store = SqliteStore(str(tmp_path / "ops.sqlite")) + msg_id = store.enqueue_channel_outbound_message( + channel="whatsapp", + account_id="wa-default", + chat_id="8615601877957@s.whatsapp.net", + text="hello", + ) + assert store.ack_channel_outbound_message(message_id=msg_id, ok=False, error="temp") is True + assert store.ack_channel_outbound_message(message_id=msg_id, ok=True, stanza_id="stanza-1") is True + row = _row(store, msg_id) + assert row["status"] == "sent" + assert int(row["send_attempts"] or 0) == 2 + assert not str(row["next_attempt_at"] or "").strip()