mirror of
https://github.com/hansjone/oclaw.git
synced 2026-10-09 03:13:19 +08:00
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 <cursoragent@cursor.com>
This commit is contained in:
parent
a8d3d2d5ec
commit
d90fca7293
2 changed files with 116 additions and 11 deletions
|
|
@ -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()
|
||||
|
|
|
|||
56
tests/test_channel_outbound_retry.py
Normal file
56
tests/test_channel_outbound_retry.py
Normal file
|
|
@ -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()
|
||||
Loading…
Add table
Add a link
Reference in a new issue