mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-09 08:22:19 +02:00
fix(reminders): retry failed note delivery through the configured channel (#6566)
This commit is contained in:
+32
-9
@@ -2297,10 +2297,10 @@ async def action_audit_skills(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||
|
||||
async def action_ping_notes(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||
"""Background note-due scanner. Fires a reminder for any note whose
|
||||
`due_date` falls in the current ±5-minute window and hasn't been pinged
|
||||
within the last 25 minutes. Mirrors `action_ping_events` for calendar.
|
||||
`due_date` falls in the current ±90-second window and hasn't been delivered
|
||||
through its configured channel in the last 25 minutes.
|
||||
|
||||
State (`data/note_pings.json`): {note_id: iso_ts_of_last_ping}. Pruned
|
||||
Per-owner state: {note_id: {at, channel}}, with legacy timestamp support. Pruned
|
||||
on each run by dropping entries for notes that are gone/archived/replied.
|
||||
"""
|
||||
try:
|
||||
@@ -2309,6 +2309,11 @@ async def action_ping_notes(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||
from datetime import datetime as _dt, timezone as _tz, timedelta as _td
|
||||
from pathlib import Path as _P
|
||||
from core.database import SessionLocal as _SL, Note as _N
|
||||
from src.settings import load_settings
|
||||
|
||||
channel = load_settings().get("reminder_channel", "browser")
|
||||
external_channel = channel in ("email", "ntfy", "webhook")
|
||||
delivery_key = f"{channel}_sent" if external_channel else "browser_sent"
|
||||
|
||||
# Per-owner state file so cache-pruning doesn't cross-delete other
|
||||
# users' entries (review C4). Legacy path kept as fallback so a
|
||||
@@ -2373,20 +2378,27 @@ async def action_ping_notes(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||
due = _parse_due(n.due_date)
|
||||
if not due:
|
||||
continue
|
||||
# Inside the ±5min window?
|
||||
# Inside the due window?
|
||||
if abs((due - now).total_seconds()) > window.total_seconds():
|
||||
continue
|
||||
# Recently pinged? Skip.
|
||||
last = cache.get(n.id)
|
||||
browser_recent = False
|
||||
if last:
|
||||
try:
|
||||
last_channel = None
|
||||
if isinstance(last, dict):
|
||||
last_channel = last.get("channel")
|
||||
last = last.get("at")
|
||||
last_dt = _dt.fromisoformat(str(last))
|
||||
if last_dt.tzinfo is None:
|
||||
last_dt = last_dt.replace(tzinfo=_tz.utc)
|
||||
if last_dt >= reping_cutoff:
|
||||
continue
|
||||
if not external_channel or last_channel == channel:
|
||||
continue
|
||||
# Browser-only receipts do not prove external
|
||||
# delivery, but the fallback must not be queued twice.
|
||||
browser_recent = last_channel in (None, "browser")
|
||||
except Exception:
|
||||
pass
|
||||
# Compose + dispatch.
|
||||
@@ -2410,15 +2422,26 @@ async def action_ping_notes(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||
body = "\n\n".join(p for p in body_parts if p) or title
|
||||
try:
|
||||
from routes.note_routes import dispatch_reminder
|
||||
await dispatch_reminder(
|
||||
result = await dispatch_reminder(
|
||||
title=title, note_body=body, note_id=n.id,
|
||||
owner=n.owner or owner or "",
|
||||
queue_browser=not browser_recent,
|
||||
)
|
||||
cache[n.id] = now.isoformat()
|
||||
sent.append(title)
|
||||
if result.get("skipped"):
|
||||
continue
|
||||
if result.get(delivery_key):
|
||||
sent.append(title)
|
||||
else:
|
||||
logger.warning("ping_notes: %s delivery failed for %s", channel, n.id)
|
||||
except Exception as e:
|
||||
logger.warning(f"ping_notes: dispatch failed for {n.id}: {e}")
|
||||
|
||||
# Dispatch owns delivery receipts. Reload before pruning so this
|
||||
# scanner cannot replace a fresh channel receipt with a timestamp.
|
||||
try:
|
||||
cache = _json.loads(STATE.read_text(encoding="utf-8")) if STATE.exists() else {}
|
||||
except Exception:
|
||||
pass
|
||||
# Prune cache entries for notes that no longer exist.
|
||||
for stale in [k for k in cache if k not in seen_ids]:
|
||||
cache.pop(stale, None)
|
||||
@@ -2429,7 +2452,7 @@ async def action_ping_notes(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||
logger.warning(f"ping_notes: cache write failed: {e}")
|
||||
|
||||
if not sent:
|
||||
raise TaskNoop(f"scanned {len(notes)} note(s), none due in ±{WINDOW_SEC}s")
|
||||
raise TaskNoop(f"scanned {len(notes)} note(s), no reminders delivered")
|
||||
preview = "; ".join(sent[:3])
|
||||
extra = f" (+{len(sent) - 3} more)" if len(sent) > 3 else ""
|
||||
return f"Pinged {len(sent)} note(s): {preview}{extra}", True
|
||||
|
||||
@@ -0,0 +1,125 @@
|
||||
"""Background note reminders retry failed primary delivery without losing receipts."""
|
||||
import datetime
|
||||
import json
|
||||
from types import SimpleNamespace
|
||||
|
||||
import httpx
|
||||
import pytest
|
||||
|
||||
from tests.helpers.database import disposable_database
|
||||
|
||||
|
||||
@pytest.fixture
|
||||
def reminder_scan(tmp_path, monkeypatch):
|
||||
from core import database
|
||||
from routes import note_routes
|
||||
from src import builtin_actions, integrations, settings
|
||||
|
||||
instant = datetime.datetime(2026, 10, 7, 12, tzinfo=datetime.timezone.utc)
|
||||
|
||||
class Clock(datetime.datetime):
|
||||
@classmethod
|
||||
def now(cls, tz=None):
|
||||
return instant.astimezone(tz) if tz else instant.replace(tzinfo=None)
|
||||
|
||||
calls, notices = [], []
|
||||
statuses = [503, 200]
|
||||
channel = {"reminder_channel": "webhook", "reminder_webhook_integration_id": "audit", "reminder_webhook_payload_template": '{"message":"{{message}}"}'}
|
||||
|
||||
def deliver(request):
|
||||
calls.append(request)
|
||||
return httpx.Response(statuses.pop(0))
|
||||
|
||||
real_client = httpx.AsyncClient
|
||||
monkeypatch.setattr(httpx, "AsyncClient", lambda **kw: real_client(transport=httpx.MockTransport(deliver), **kw))
|
||||
monkeypatch.setattr(settings, "load_settings", lambda: channel)
|
||||
monkeypatch.setattr(integrations, "load_integrations", lambda: [{"id": "audit", "preset": "ntfy", "base_url": "https://reminder.example.test", "enabled": True}])
|
||||
monkeypatch.setattr(note_routes, "_scheduler_ref", SimpleNamespace(add_notification=lambda **kw: notices.append(kw)))
|
||||
monkeypatch.setattr(note_routes, "DATA_DIR", str(tmp_path))
|
||||
monkeypatch.setattr(builtin_actions, "DATA_DIR", str(tmp_path))
|
||||
monkeypatch.setattr(datetime, "datetime", Clock)
|
||||
monkeypatch.setattr("src.url_safety.check_outbound_url", lambda *a, **kw: (True, ""))
|
||||
with disposable_database(tmp_path) as factory:
|
||||
monkeypatch.setattr(database, "SessionLocal", factory)
|
||||
with factory() as db:
|
||||
db.add(database.Note(id="note-1", owner="alice", title="Due note", due_date=instant.isoformat()))
|
||||
db.commit()
|
||||
yield SimpleNamespace(
|
||||
scan=lambda: builtin_actions.action_ping_notes("alice"),
|
||||
path=tmp_path / "note_pings_alice.json", calls=calls, notices=notices,
|
||||
statuses=statuses, channel=channel, instant=instant, factory=factory,
|
||||
)
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
@pytest.mark.parametrize("channel", ["webhook", "ntfy"])
|
||||
async def test_failed_primary_is_retried_and_receipt_keeps_channel(reminder_scan, channel):
|
||||
from src.builtin_actions import TaskNoop
|
||||
scan = reminder_scan
|
||||
scan.channel["reminder_channel"] = channel
|
||||
with pytest.raises(TaskNoop):
|
||||
await scan.scan()
|
||||
assert len(scan.calls) == 1
|
||||
assert json.loads(scan.path.read_text())["note-1"]["channel"] == "browser"
|
||||
message, ok = await scan.scan()
|
||||
assert ok is True and "Pinged 1" in message
|
||||
assert len(scan.calls) == 2
|
||||
assert len(scan.notices) == 1
|
||||
assert json.loads(scan.path.read_text())["note-1"]["channel"] == channel
|
||||
with pytest.raises(TaskNoop):
|
||||
await scan.scan()
|
||||
assert len(scan.calls) == 2
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_legacy_browser_receipt_allows_external_retry(reminder_scan):
|
||||
scan = reminder_scan
|
||||
scan.path.write_text(json.dumps({"note-1": scan.instant.isoformat()}))
|
||||
scan.statuses[:] = [200]
|
||||
message, ok = await scan.scan()
|
||||
assert ok is True and "Pinged 1" in message
|
||||
assert len(scan.calls) == 1 and scan.notices == []
|
||||
assert json.loads(scan.path.read_text())["note-1"]["channel"] == "webhook"
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_browser_delivery_is_deduped_and_typed(reminder_scan):
|
||||
from src.builtin_actions import TaskNoop
|
||||
scan = reminder_scan
|
||||
scan.channel["reminder_channel"] = "browser"
|
||||
_, ok = await scan.scan()
|
||||
assert ok is True
|
||||
assert json.loads(scan.path.read_text())["note-1"]["channel"] == "browser"
|
||||
with pytest.raises(TaskNoop):
|
||||
await scan.scan()
|
||||
assert len(scan.notices) == 1 and scan.calls == []
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_total_delivery_failure_does_not_checkpoint(reminder_scan, monkeypatch):
|
||||
from routes import note_routes
|
||||
from src.builtin_actions import TaskNoop
|
||||
scan = reminder_scan
|
||||
monkeypatch.setattr(note_routes, "_scheduler_ref", None)
|
||||
with pytest.raises(TaskNoop):
|
||||
await scan.scan()
|
||||
assert "note-1" not in json.loads(scan.path.read_text())
|
||||
_, ok = await scan.scan()
|
||||
assert ok is True and len(scan.calls) == 2
|
||||
|
||||
|
||||
@pytest.mark.asyncio
|
||||
async def test_pruning_retains_dispatcher_receipt_and_other_due_note(reminder_scan):
|
||||
from core import database
|
||||
scan = reminder_scan
|
||||
with scan.factory() as db:
|
||||
db.add(database.Note(id="later", owner="alice", title="Later", due_date=(scan.instant + datetime.timedelta(days=1)).isoformat()))
|
||||
db.commit()
|
||||
scan.path.write_text(json.dumps({"gone": "old", "later": {"at": scan.instant.isoformat(), "channel": "email"}}))
|
||||
scan.statuses[:] = [200]
|
||||
_, ok = await scan.scan()
|
||||
assert ok is True
|
||||
cache = json.loads(scan.path.read_text())
|
||||
assert set(cache) == {"note-1", "later"}
|
||||
assert cache["note-1"]["channel"] == "webhook"
|
||||
assert cache["later"]["channel"] == "email"
|
||||
Reference in New Issue
Block a user