From 2abc062280c899e36460c5eed98effbabea501aa Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 19:26:41 +0000 Subject: [PATCH 1/7] fix(email): serialize urgency checkpoints --- src/builtin_actions.py | 216 +++++++++++++++---------- tests/test_email_urgency_checkpoint.py | 56 +++++++ 2 files changed, 188 insertions(+), 84 deletions(-) create mode 100644 tests/test_email_urgency_checkpoint.py diff --git a/src/builtin_actions.py b/src/builtin_actions.py index 68817467f..b560385db 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -20,6 +20,48 @@ from src.interactive_gate import wait_for_interactive_quiet logger = logging.getLogger(__name__) +def _run_email_urgency_state_transaction(state_path, lock_db_path, operation): + """Serialize one owner urgency checkpoint across worker processes. + + ``operation`` receives the latest JSON state and returns + ``(result, next_state)``. The scheduled-email database is used only as a + durable cross-process write lock; the UI-facing state remains the existing + JSON file. Running the operation in a worker thread keeps the async event + loop free while reminder delivery is in flight. + """ + import sqlite3 + import uuid + from pathlib import Path + + state_path = Path(state_path) + state_path.parent.mkdir(parents=True, exist_ok=True) + conn = sqlite3.connect(str(lock_db_path), timeout=120) + temp_path = state_path.with_name( + f".{state_path.name}.{uuid.uuid4().hex}.tmp" + ) + try: + conn.execute("BEGIN IMMEDIATE") + try: + prior = json.loads(state_path.read_text(encoding="utf-8")) \ + if state_path.exists() else {} + if not isinstance(prior, dict): + prior = {} + except Exception: + prior = {} + + result, next_state = operation(prior) + temp_path.write_text(json.dumps(next_state), encoding="utf-8") + temp_path.replace(state_path) + conn.commit() + return result + except Exception: + conn.rollback() + raise + finally: + temp_path.unlink(missing_ok=True) + conn.close() + + class TaskNoop(BaseException): """Raised by an action when it determined there's nothing to do. @@ -1878,6 +1920,7 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # filename for single-user installs (matches prior behaviour). _owner_slug = "".join(c if (c.isalnum() or c in "-_.@") else "_" for c in (owner or "default")) STATE_PATH = _P(DATA_DIR) / f"email_urgency_state_{_owner_slug}.json" + STATE_LOCK_DB = STATE_PATH.with_suffix(".lock.sqlite3") CACHE_DIR = _P(EMAIL_URGENCY_CACHE_DIR) CACHE_DIR.mkdir(parents=True, exist_ok=True) STATE_PATH.parent.mkdir(parents=True, exist_ok=True) @@ -2375,97 +2418,102 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: max_score = max((v.get("score", 0) for v in per_uid_scores.values()), default=0) total_urgent = len(urgent_keys) - # Load prior state to know which urgent UIDs we've already notified. - try: - prior = _json.loads(STATE_PATH.read_text(encoding="utf-8")) if STATE_PATH.exists() else {} - except Exception: - prior = {} - notified_uids = set(prior.get("notified_uids", [])) - - # ── 5. Fire reminder ONLY when a previously-unnotified UID scores urgent. - new_urgent = [k for k in urgent_keys if k not in notified_uids] + # ── 5. Fire a reminder only when a previously-unnotified UID scores + # urgent. The read, decision, delivery, and checkpoint are serialized + # below so two scheduler workers cannot both act on the same stale + # state or overwrite each other's successful checkpoint. newly_notified = set() notify_failed = set() - if new_urgent: - title = "Urgent email" if total_urgent == 1 else f"{total_urgent} urgent emails" - # Build a real listing — subject · sender · reason for each urgent - # one — so the reminder email tells you which messages to act on, - # not just "4 needing reply". Optional deep-link when the user has - # `app_public_url` configured in Settings (so the email row links - # straight into the Odysseus Email tab). - # Sort: highest-scored UIDs first; cap at 10 to keep the email tidy. - sorted_urgent = sorted( - ((k, per_uid_scores[k]) for k in urgent_keys), - key=lambda kv: kv[1].get("score", 0), reverse=True, - )[:10] - _pub = (settings.get("app_public_url") or "").strip().rstrip("/") - from urllib.parse import quote as _quote - lines = [f"{total_urgent} email" + ("" if total_urgent == 1 else "s") + " need an urgent reply:", ""] - for i, (k, v) in enumerate(sorted_urgent, 1): - subj = (v.get("subject") or "(no subject)")[:160] - frm = v.get("from") or "" - why = v.get("reason") or "" - uid_for_link = str(k).split(":", 1)[-1] - hash_link = f"#email={_quote('INBOX', safe='')}:{uid_for_link}" - open_link = f"{_pub}/{hash_link}" if _pub else hash_link - line = f"{i}. {subj}" - if frm: - line += f" — {frm}" - if why: - line += f" · {why}" - lines.append(line) - lines.append(f" Open email: {open_link}") - if total_urgent > len(sorted_urgent): - lines.append("") - lines.append(f"…and {total_urgent - len(sorted_urgent)} more.") - body = "\n".join(lines) - try: - # Call dispatch_reminder DIRECTLY (no HTTP/auth roundtrip — the - # endpoint version 401's the background scheduler because it - # has no session cookie). - from routes.note_routes import dispatch_reminder - dispatch_result = await dispatch_reminder( - title=title, note_body=body, note_id="urgent-email", - owner=owner or "", - ) - channel = (settings.get("reminder_channel") or "browser").strip().lower() - delivered = bool(dispatch_result.get("browser_sent")) - if channel == "email": - delivered = bool(dispatch_result.get("email_sent")) - elif channel == "ntfy": - delivered = bool(dispatch_result.get("ntfy_sent")) - elif channel == "webhook": - delivered = bool(dispatch_result.get("webhook_sent")) - if delivered: - newly_notified.update(new_urgent) - else: + title = "Urgent email" if total_urgent == 1 else f"{total_urgent} urgent emails" + # Build a real listing — subject · sender · reason for each urgent one + # — so the reminder tells the user which messages to act on. + sorted_urgent = sorted( + ((k, per_uid_scores[k]) for k in urgent_keys), + key=lambda kv: kv[1].get("score", 0), reverse=True, + )[:10] + _pub = (settings.get("app_public_url") or "").strip().rstrip("/") + from urllib.parse import quote as _quote + lines = [f"{total_urgent} email" + ("" if total_urgent == 1 else "s") + " need an urgent reply:", ""] + for i, (k, v) in enumerate(sorted_urgent, 1): + subj = (v.get("subject") or "(no subject)")[:160] + frm = v.get("from") or "" + why = v.get("reason") or "" + uid_for_link = str(k).split(":", 1)[-1] + hash_link = f"#email={_quote('INBOX', safe='')}:{uid_for_link}" + open_link = f"{_pub}/{hash_link}" if _pub else hash_link + line = f"{i}. {subj}" + if frm: + line += f" — {frm}" + if why: + line += f" · {why}" + lines.append(line) + lines.append(f" Open email: {open_link}") + if total_urgent > len(sorted_urgent): + lines.append("") + lines.append(f"…and {total_urgent - len(sorted_urgent)} more.") + body = "\n".join(lines) + + async def _dispatch_urgency_reminder(): + # Call dispatch_reminder directly: a scheduler has no browser + # session cookie with which to call the HTTP endpoint. + from routes.note_routes import dispatch_reminder + return await dispatch_reminder( + title=title, + note_body=body, + note_id="urgent-email", + owner=owner or "", + ) + + def _dispatch_and_checkpoint(prior): + notified_uids = set(prior.get("notified_uids", [])) + new_urgent = [k for k in urgent_keys if k not in notified_uids] + if new_urgent: + try: + dispatch_result = _aio.run(_dispatch_urgency_reminder()) + channel = (settings.get("reminder_channel") or "browser").strip().lower() + delivered = bool(dispatch_result.get("browser_sent")) + if channel == "email": + delivered = bool(dispatch_result.get("email_sent")) + elif channel == "ntfy": + delivered = bool(dispatch_result.get("ntfy_sent")) + elif channel == "webhook": + delivered = bool(dispatch_result.get("webhook_sent")) + if delivered: + newly_notified.update(new_urgent) + notified_uids.update(new_urgent) + else: + notify_failed.update(new_urgent) + logger.warning( + "urgency: reminder dispatch returned no successful " + f"delivery path: {dispatch_result}" + ) + except Exception as e: + logger.warning(f"urgency: reminder dispatch failed: {e}") notify_failed.update(new_urgent) - logger.warning(f"urgency: reminder dispatch returned no successful delivery path: {dispatch_result}") - except Exception as e: - logger.warning(f"urgency: reminder dispatch failed: {e}") - notify_failed.update(new_urgent) - # Mark only successfully delivered UIDs as notified so a transient - # SMTP/ntfy/browser failure retries instead of lying forever. - notified_uids.update(newly_notified) - # Prune notified_uids that aren't unread anymore (so a future re-urgent - # message with the same UID — rare but possible after archive→unarchive - # — can re-notify). Keep only UIDs still in `all_unread_keys`. - notified_uids = {u for u in notified_uids if u in all_unread_keys} + # Prune UIDs no longer unread, while retaining every successful + # checkpoint observed after this worker acquired the durable lock. + notified_uids.intersection_update(all_unread_keys) + next_state = { + "ts": _time.time(), + "owner": owner or "", + "total_unread": len(all_unread_keys), + "total_urgent": total_urgent, + "max_score": max_score, + "per_uid": per_uid_scores, + "notified_uids": sorted(notified_uids), + } + return notified_uids, next_state - state = { - "ts": _time.time(), - "owner": owner or "", - "total_unread": len(all_unread_keys), - "total_urgent": total_urgent, - "max_score": max_score, - "per_uid": per_uid_scores, - "notified_uids": sorted(notified_uids), - } try: - STATE_PATH.write_text(_json.dumps(state), encoding="utf-8") + await _aio.to_thread( + _run_email_urgency_state_transaction, + STATE_PATH, + STATE_LOCK_DB, + _dispatch_and_checkpoint, + ) except Exception as e: - logger.warning(f"urgency: state write failed: {e}") + logger.warning(f"urgency: state transaction failed: {e}") # ── 6. Activity-log summary — counts line on top, then per-tier # bulleted breakdown so the user can see WHICH emails ranked where diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py new file mode 100644 index 000000000..1790cc29d --- /dev/null +++ b/tests/test_email_urgency_checkpoint.py @@ -0,0 +1,56 @@ +import json +import threading +from concurrent.futures import ThreadPoolExecutor + + +def test_urgency_state_transaction_serializes_decision_and_checkpoint(tmp_path): + """A later worker must observe the first worker's delivered UID.""" + from src.builtin_actions import _run_email_urgency_state_transaction + + state_path = tmp_path / "email_urgency_state_alice.json" + lock_db = tmp_path / "scheduled-emails.db" + first_entered = threading.Event() + allow_first_to_finish = threading.Event() + second_entered = threading.Event() + deliveries = [] + + def operation(name): + def update(prior): + if name == "first": + first_entered.set() + assert allow_first_to_finish.wait(timeout=2) + else: + second_entered.set() + + notified = set(prior.get("notified_uids", [])) + delivered = "acct:42" not in notified + if delivered: + deliveries.append(name) + notified.add("acct:42") + return delivered, { + "owner": "alice", + "notified_uids": sorted(notified), + } + + return _run_email_urgency_state_transaction( + state_path, + lock_db, + update, + ) + + with ThreadPoolExecutor(max_workers=2) as executor: + first = executor.submit(operation, "first") + assert first_entered.wait(timeout=2) + second = executor.submit(operation, "second") + assert not second_entered.wait(timeout=0.2) + allow_first_to_finish.set() + + assert first.result(timeout=2) is True + assert second.result(timeout=2) is False + + assert deliveries == ["first"] + assert json.loads(state_path.read_text(encoding="utf-8")) == { + "owner": "alice", + "notified_uids": ["acct:42"], + } + From 49daa73106579f119b1bbc6dfef5618caf6eee79 Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 20:05:40 +0000 Subject: [PATCH 2/7] fix(email): preserve urgency transaction lifecycle --- src/builtin_actions.py | 316 +++++++++++++--- tests/test_email_urgency_checkpoint.py | 475 ++++++++++++++++++++++++- 2 files changed, 720 insertions(+), 71 deletions(-) diff --git a/src/builtin_actions.py b/src/builtin_actions.py index b560385db..936c6db60 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -20,41 +20,87 @@ from src.interactive_gate import wait_for_interactive_quiet logger = logging.getLogger(__name__) -def _run_email_urgency_state_transaction(state_path, lock_db_path, operation): - """Serialize one owner urgency checkpoint across worker processes. - - ``operation`` receives the latest JSON state and returns - ``(result, next_state)``. The scheduled-email database is used only as a - durable cross-process write lock; the UI-facing state remains the existing - JSON file. Running the operation in a worker thread keeps the async event - loop free while reminder delivery is in flight. - """ +def _acquire_email_urgency_state_lock( + state_path, + lock_db_path, + cancel_event, + timeout_seconds=120, +): + """Acquire the cross-process urgency lock without blocking the app loop.""" import sqlite3 - import uuid + import time from pathlib import Path state_path = Path(state_path) state_path.parent.mkdir(parents=True, exist_ok=True) - conn = sqlite3.connect(str(lock_db_path), timeout=120) - temp_path = state_path.with_name( - f".{state_path.name}.{uuid.uuid4().hex}.tmp" - ) - try: - conn.execute("BEGIN IMMEDIATE") + deadline = time.monotonic() + timeout_seconds + + while not cancel_event.is_set(): + remaining = deadline - time.monotonic() + if remaining <= 0: + raise sqlite3.OperationalError("timed out waiting for urgency state lock") + conn = sqlite3.connect( + str(lock_db_path), + timeout=min(0.25, max(0.01, remaining)), + check_same_thread=False, + ) try: - prior = json.loads(state_path.read_text(encoding="utf-8")) \ - if state_path.exists() else {} + conn.execute("BEGIN IMMEDIATE") + except sqlite3.OperationalError as exc: + conn.close() + if "locked" not in str(exc).lower(): + raise + cancel_event.wait(min(0.05, max(0.0, remaining))) + continue + except BaseException: + conn.close() + raise + + if cancel_event.is_set(): + conn.rollback() + conn.close() + return None, None + try: + prior = ( + json.loads(state_path.read_text(encoding="utf-8")) + if state_path.exists() + else {} + ) if not isinstance(prior, dict): prior = {} except Exception: prior = {} + return conn, prior - result, next_state = operation(prior) + return None, None + + +def _close_email_urgency_state_lock(conn): + if conn is None: + return + try: + try: + conn.rollback() + except Exception: + pass + finally: + conn.close() + + +def _commit_email_urgency_state(conn, state_path, next_state): + """Atomically publish JSON before releasing the SQLite write lock.""" + import uuid + from pathlib import Path + + state_path = Path(state_path) + temp_path = state_path.with_name( + f".{state_path.name}.{uuid.uuid4().hex}.tmp" + ) + try: temp_path.write_text(json.dumps(next_state), encoding="utf-8") temp_path.replace(state_path) conn.commit() - return result - except Exception: + except BaseException: conn.rollback() raise finally: @@ -62,6 +108,130 @@ def _run_email_urgency_state_transaction(state_path, lock_db_path, operation): conn.close() +async def _run_email_urgency_state_transaction( + state_path, + lock_db_path, + operation, +): + """Serialize one urgency decision while keeping async work on this loop. + + Only lock acquisition waits in a worker thread. ``operation`` is awaited + on the caller's long-lived event loop, where shared async clients, locks, + and the browser-notification queue belong. Cancellation rolls back the + SQLite transaction and never publishes a checkpoint. + """ + import asyncio + import threading + + loop = asyncio.get_running_loop() + cancel_event = threading.Event() + acquire_future = loop.run_in_executor( + None, + _acquire_email_urgency_state_lock, + state_path, + lock_db_path, + cancel_event, + ) + try: + conn, prior = await asyncio.shield(acquire_future) + except asyncio.CancelledError as cancelled: + cancel_event.set() + # The acquisition worker owns any connection until it returns. Wait + # for its short busy-poll to observe cancellation, then close a lock it + # may have won concurrently with the cancellation request. + while True: + try: + conn, _prior = await asyncio.shield(acquire_future) + break + except asyncio.CancelledError: + continue + except Exception: + conn = None + break + _close_email_urgency_state_lock(conn) + raise cancelled + + if conn is None: + raise asyncio.CancelledError + + try: + result, next_state = await operation(prior) + # Keep this small atomic publish synchronous. There is no await between + # the successful operation and commit, so cancellation cannot be + # observed and then followed by a checkpoint. + try: + _commit_email_urgency_state(conn, state_path, next_state) + finally: + conn = None + return result + except BaseException: + _close_email_urgency_state_lock(conn) + raise + + +def _email_urgency_account_key(message_key): + return str(message_key).split(":", 1)[0] + + +def _merge_email_urgency_state( + prior, + *, + owner, + per_uid_scores, + notified_uids, + all_unread_keys, + fully_scanned_account_ids, + timestamp, +): + """Merge current scan knowledge without deleting unscanned accounts.""" + prior_per_uid = prior.get("per_uid", {}) + if not isinstance(prior_per_uid, dict): + prior_per_uid = {} + complete = {str(account_id) for account_id in fully_scanned_account_ids} + + merged_per_uid = dict(prior_per_uid) + for key in list(merged_per_uid): + if _email_urgency_account_key(key) in complete: + merged_per_uid.pop(key, None) + # Partial scans may add or refresh facts, but absence from a partial scan + # is not evidence that another checkpoint or UI row is stale. + merged_per_uid.update(per_uid_scores) + + merged_notified = set(notified_uids) + for key in list(merged_notified): + if ( + _email_urgency_account_key(key) in complete + and key not in all_unread_keys + ): + merged_notified.discard(key) + + total_unread = 0 + total_urgent = 0 + max_score = 0 + for value in merged_per_uid.values(): + if not isinstance(value, dict): + continue + try: + score = max(0, min(3, int(value.get("score", 0)))) + except (TypeError, ValueError): + score = 0 + max_score = max(max_score, score) + if value.get("unread"): + total_unread += 1 + if score >= 2: + total_urgent += 1 + + return { + "ts": timestamp, + "owner": owner or "", + "total_unread": total_unread, + "total_urgent": total_urgent, + "max_score": max_score, + "per_uid": merged_per_uid, + "notified_uids": sorted(merged_notified), + } + + class TaskNoop(BaseException): """Raised by an action when it determined there's nothing to do. @@ -1972,6 +2142,7 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: failed_classifications = [] tag_write_details = [] scanned = 0 + fully_scanned_account_ids = set() def _heuristic_email_verdict(item: dict) -> dict: blob = ( @@ -2067,16 +2238,27 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: def _scan_one(account=acc, cache_uids=cache.get("uids", {})): """Sync IMAP work runs in a thread.""" results = [] + scan_complete = True conn = _imap_connect(account.id) try: - conn.select("INBOX", readonly=True) + select_status, _select_data = conn.select("INBOX", readonly=True) + if select_status != "OK": + return results, False # Tag recent inbox mail, not only unread mail. Urgency # reminders below still only notify for unread messages. since_str = AGE_CUTOFF.strftime("%d-%b-%Y") status, data = conn.uid("SEARCH", None, f'(SINCE {since_str})') - if status != "OK" or not data or not data[0]: - return results - uids = data[0].split()[-30:] + if status != "OK": + return results, False + if not data or not data[0]: + return results, True + matching_uids = data[0].split() + if len(matching_uids) > 30: + # The scale guard deliberately processes only the most + # recent 30. That is a partial account snapshot, so it + # cannot justify pruning older checkpoint facts. + scan_complete = False + uids = matching_uids[-30:] for uid_b in uids: uid = uid_b.decode() if isinstance(uid_b, bytes) else str(uid_b) key = f"{account.id}:{uid}" @@ -2084,12 +2266,41 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: cached_ok = isinstance(cached, dict) and cached.get("triage_version") == TRIAGE_VERSION results.append({"key": key, "uid": uid, "cached": cached if cached_ok else None}) if cached_ok: - # Already classified — skip the fetch. + # Cached verdicts still need a lightweight FLAGS + # refresh. Without it a cached unread message looks + # read and its successful notification checkpoint + # is pruned on the next pass. + try: + st, flag_data = conn.uid("FETCH", uid_b, "(UID FLAGS)") + if st != "OK" or not flag_data: + scan_complete = False + results.pop() + continue + flag_parts = [] + for part in flag_data: + if isinstance(part, (bytes, bytearray)): + flag_parts.append(bytes(part)) + elif ( + isinstance(part, tuple) + and part + and isinstance(part[0], (bytes, bytearray)) + ): + flag_parts.append(bytes(part[0])) + flags_blob = b" ".join(flag_parts) + results[-1]["unread"] = b"\\Seen" not in flags_blob + except Exception as _fe: + scan_complete = False + results.pop() + logger.debug( + f"urgency: flag fetch for uid {uid} failed: {_fe}" + ) continue # Pull headers + first ~800 chars of plaintext body. try: st, msg_data = conn.uid("FETCH", uid_b, "(UID FLAGS RFC822.HEADER BODY.PEEK[TEXT]<0.800>)") if st != "OK" or not msg_data: + scan_complete = False + results.pop() continue flags_blob = b" ".join( part[0] for part in msg_data @@ -2103,6 +2314,8 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: if isinstance(part, tuple) and part[1]: raw += part[1] + b"\n\n" if not raw: + scan_complete = False + results.pop() continue msg = _email_mod.message_from_bytes(raw) # Skip Odysseus-generated reminders so the scanner @@ -2158,17 +2371,21 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: "unread": is_unread, }) except Exception as _fe: + scan_complete = False + results.pop() logger.debug(f"urgency: header fetch for uid {uid} failed: {_fe}") finally: try: conn.logout() except Exception: pass - return results + return results, scan_complete try: - items = await _aio.to_thread(_scan_one) + items, scan_complete = await _aio.to_thread(_scan_one) except Exception as e: logger.warning(f"urgency: IMAP scan failed for account {acc.id}: {e}") continue + if scan_complete: + fully_scanned_account_ids.add(str(acc.id)) for item in items: scanned += 1 @@ -2305,13 +2522,13 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: logger.debug(f"urgency: LLM classify failed for {key}: {e}") continue - # ── Prune cache entries for UIDs that are no longer in the recent - # scan window. Read messages remain cached because tags are useful - # on read mail too; unread state is refreshed per scan above. - seen_uids = {it["uid"] for it in items} - cache_uids = cache.get("uids", {}) - for stale in [u for u in cache_uids if u not in seen_uids]: - cache_uids.pop(stale, None) + if scan_complete: + # Only a complete account scan proves a cached UID left the + # recent window. Partial/failing scans preserve prior facts. + seen_uids = {it["uid"] for it in items} + cache_uids = cache.get("uids", {}) + for stale in [u for u in cache_uids if u not in seen_uids]: + cache_uids.pop(stale, None) try: cache_file.write_text(_json.dumps(cache), encoding="utf-8") @@ -2415,7 +2632,6 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # ── 4. Aggregate state. urgent = score ≥ 2. urgent_keys = [k for k, v in per_uid_scores.items() if v.get("score", 0) >= 2 and v.get("unread")] - max_score = max((v.get("score", 0) for v in per_uid_scores.values()), default=0) total_urgent = len(urgent_keys) # ── 5. Fire a reminder only when a previously-unnotified UID scores @@ -2464,12 +2680,12 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: owner=owner or "", ) - def _dispatch_and_checkpoint(prior): + async def _dispatch_and_checkpoint(prior): notified_uids = set(prior.get("notified_uids", [])) new_urgent = [k for k in urgent_keys if k not in notified_uids] if new_urgent: try: - dispatch_result = _aio.run(_dispatch_urgency_reminder()) + dispatch_result = await _dispatch_urgency_reminder() channel = (settings.get("reminder_channel") or "browser").strip().lower() delivered = bool(dispatch_result.get("browser_sent")) if channel == "email": @@ -2491,23 +2707,19 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: logger.warning(f"urgency: reminder dispatch failed: {e}") notify_failed.update(new_urgent) - # Prune UIDs no longer unread, while retaining every successful - # checkpoint observed after this worker acquired the durable lock. - notified_uids.intersection_update(all_unread_keys) - next_state = { - "ts": _time.time(), - "owner": owner or "", - "total_unread": len(all_unread_keys), - "total_urgent": total_urgent, - "max_score": max_score, - "per_uid": per_uid_scores, - "notified_uids": sorted(notified_uids), - } + next_state = _merge_email_urgency_state( + prior, + owner=owner, + per_uid_scores=per_uid_scores, + notified_uids=notified_uids, + all_unread_keys=all_unread_keys, + fully_scanned_account_ids=fully_scanned_account_ids, + timestamp=_time.time(), + ) return notified_uids, next_state try: - await _aio.to_thread( - _run_email_urgency_state_transaction, + await _run_email_urgency_state_transaction( STATE_PATH, STATE_LOCK_DB, _dispatch_and_checkpoint, diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py index 1790cc29d..ae7ef1dc3 100644 --- a/tests/test_email_urgency_checkpoint.py +++ b/tests/test_email_urgency_checkpoint.py @@ -1,24 +1,164 @@ +import asyncio import json import threading -from concurrent.futures import ThreadPoolExecutor +from types import SimpleNamespace + +import pytest -def test_urgency_state_transaction_serializes_decision_and_checkpoint(tmp_path): +class _Column: + def __eq__(self, _other): + return True + + def __ne__(self, _other): + return True + + +class _Query: + def __init__(self, rows): + self._rows = rows + + def filter(self, *_args, **_kwargs): + return self + + def all(self): + return list(self._rows) + + +class _Db: + def __init__(self, rows): + self._rows = rows + + def query(self, _model): + return _Query(self._rows()) + + def close(self): + return None + + +class _EmailAccount: + enabled = _Column() + owner = _Column() + imap_user = _Column() + from_address = _Column() + id = _Column() + + +class _FakeImap: + def __init__(self, account_id, failures, seen_accounts): + self.account_id = account_id + self.failures = failures + self.seen_accounts = seen_accounts + + def select(self, *_args, **_kwargs): + if self.account_id in self.failures: + raise RuntimeError(f"{self.account_id} unavailable") + return "OK", [] + + def uid(self, command, *_args): + if self.account_id in self.failures: + raise RuntimeError(f"{self.account_id} unavailable") + if command == "SEARCH": + return "OK", [b"1"] + query = str(_args[-1]) if _args else "" + seen = "\\Seen" if self.account_id in self.seen_accounts else "" + flags = f"1 (UID 1 FLAGS ({seen}))".encode() + if query == "(UID FLAGS)": + return "OK", [flags] + raw = ( + f"From: Sender {self.account_id} \r\n" + f"Subject: Urgent request for {self.account_id}\r\n" + f"Message-ID: <{self.account_id}-1@example.com>\r\n" + "\r\n" + "Please reply immediately." + ).encode() + return "OK", [(flags, raw)] + + def logout(self): + return None + + +def _account(account_id): + return SimpleNamespace( + id=account_id, + enabled=True, + owner="alice", + imap_user="alice", + from_address="alice", + ) + + +def _configure_action(monkeypatch, tmp_path, account_ids): + from core import database + from routes import email_helpers + from src import builtin_actions, llm_core, settings, task_endpoint + + runtime = { + "accounts": list(account_ids), + "failures": set(), + "seen_accounts": set(), + "settings": { + "reminder_channel": "browser", + "reminder_llm_synthesis": False, + "app_public_url": "", + }, + } + + monkeypatch.setattr(builtin_actions, "DATA_DIR", str(tmp_path)) + monkeypatch.setattr( + builtin_actions, + "EMAIL_URGENCY_CACHE_DIR", + str(tmp_path / "urgency-cache"), + ) + monkeypatch.setattr(database, "EmailAccount", _EmailAccount) + monkeypatch.setattr( + database, + "SessionLocal", + lambda: _Db(lambda: [_account(value) for value in runtime["accounts"]]), + ) + monkeypatch.setattr( + task_endpoint, + "resolve_task_candidates", + lambda *args, **kwargs: [("http://llm", "model", {})], + ) + + async def fake_fallback(*_args, **_kwargs): + return '{"score": 3, "reason": "urgent"}' + + monkeypatch.setattr(llm_core, "llm_call_async_with_fallback", fake_fallback) + monkeypatch.setattr(settings, "load_settings", lambda: dict(runtime["settings"])) + monkeypatch.setattr( + email_helpers, + "SCHEDULED_DB", + tmp_path / "scheduled-emails.db", + ) + monkeypatch.setattr( + email_helpers, + "_imap_connect", + lambda account_id=None, **_kwargs: _FakeImap( + str(account_id), runtime["failures"], runtime["seen_accounts"] + ), + ) + return builtin_actions, runtime + + +@pytest.mark.asyncio +async def test_urgency_state_transaction_serializes_decision_and_checkpoint(tmp_path): """A later worker must observe the first worker's delivered UID.""" from src.builtin_actions import _run_email_urgency_state_transaction state_path = tmp_path / "email_urgency_state_alice.json" - lock_db = tmp_path / "scheduled-emails.db" - first_entered = threading.Event() - allow_first_to_finish = threading.Event() - second_entered = threading.Event() + lock_db = tmp_path / "urgency.lock.sqlite3" + first_entered = asyncio.Event() + allow_first_to_finish = asyncio.Event() + second_entered = asyncio.Event() deliveries = [] - def operation(name): - def update(prior): + async def operation(name): + async def update(prior): if name == "first": first_entered.set() - assert allow_first_to_finish.wait(timeout=2) + await allow_first_to_finish.wait() else: second_entered.set() @@ -32,25 +172,322 @@ def test_urgency_state_transaction_serializes_decision_and_checkpoint(tmp_path): "notified_uids": sorted(notified), } - return _run_email_urgency_state_transaction( + return await _run_email_urgency_state_transaction( state_path, lock_db, update, ) - with ThreadPoolExecutor(max_workers=2) as executor: - first = executor.submit(operation, "first") - assert first_entered.wait(timeout=2) - second = executor.submit(operation, "second") - assert not second_entered.wait(timeout=0.2) - allow_first_to_finish.set() - - assert first.result(timeout=2) is True - assert second.result(timeout=2) is False + first = asyncio.create_task(operation("first")) + await asyncio.wait_for(first_entered.wait(), timeout=2) + second = asyncio.create_task(operation("second")) + await asyncio.sleep(0.1) + assert not second_entered.is_set() + allow_first_to_finish.set() + assert await asyncio.wait_for(first, timeout=2) is True + assert await asyncio.wait_for(second, timeout=2) is False assert deliveries == ["first"] assert json.loads(state_path.read_text(encoding="utf-8")) == { "owner": "alice", "notified_uids": ["acct:42"], } + +@pytest.mark.asyncio +async def test_waiting_transaction_cancellation_does_not_leak_lock(tmp_path): + from src.builtin_actions import _run_email_urgency_state_transaction + + state_path = tmp_path / "email_urgency_state_alice.json" + lock_db = tmp_path / "urgency.lock.sqlite3" + first_entered = asyncio.Event() + release_first = asyncio.Event() + + async def first_operation(_prior): + first_entered.set() + await release_first.wait() + return None, {"notified_uids": ["acct:1"]} + + async def later_operation(prior): + return None, prior + + first = asyncio.create_task( + _run_email_urgency_state_transaction( + state_path, lock_db, first_operation + ) + ) + await asyncio.wait_for(first_entered.wait(), timeout=2) + waiting = asyncio.create_task( + _run_email_urgency_state_transaction( + state_path, lock_db, later_operation + ) + ) + await asyncio.sleep(0.05) + waiting.cancel() + with pytest.raises(asyncio.CancelledError): + await asyncio.wait_for(waiting, timeout=2) + + release_first.set() + await asyncio.wait_for(first, timeout=2) + await asyncio.wait_for( + _run_email_urgency_state_transaction( + state_path, lock_db, later_operation + ), + timeout=2, + ) + + +@pytest.mark.asyncio +async def test_action_dispatch_stays_on_app_loop_and_queues_browser_notification( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + from src import endpoint_resolver, llm_core + from src.task_scheduler import TaskScheduler + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + runtime["settings"]["reminder_llm_synthesis"] = True + monkeypatch.setattr(note_routes, "DATA_DIR", str(tmp_path)) + monkeypatch.setattr( + endpoint_resolver, + "resolve_endpoint", + lambda *_args, **_kwargs: ( + "https://api.openai.com/v1", + "utility-model", + {}, + ), + ) + + expected_loop = asyncio.get_running_loop() + expected_thread = threading.get_ident() + synthesis_loops = [] + shared_client = SimpleNamespace(is_closed=False) + monkeypatch.setattr(llm_core, "_http_client", shared_client) + monkeypatch.setattr(llm_core, "_response_cache", {}) + + class _Response: + is_success = True + status_code = 200 + text = "ok" + + @staticmethod + def json(): + return { + "choices": [ + {"message": {"content": "Synthesized urgency reminder."}} + ] + } + + async def fake_http_post(client, *_args, **_kwargs): + synthesis_loops.append(asyncio.get_running_loop()) + assert client is shared_client + return _Response() + + monkeypatch.setattr( + llm_core, + "httpx_post_kimi_aware_async", + fake_http_post, + ) + + scheduler = TaskScheduler(None) + notification_threads = [] + original_add = scheduler.add_notification + + def checked_add(*args, **kwargs): + notification_threads.append(threading.get_ident()) + return original_add(*args, **kwargs) + + monkeypatch.setattr(scheduler, "add_notification", checked_add) + monkeypatch.setattr(note_routes, "_scheduler_ref", scheduler) + + message, ok = await builtin_actions.action_check_email_urgency("alice") + + assert ok is True + assert "notified 1" in message + assert synthesis_loops == [expected_loop] + assert notification_threads == [expected_thread] + notifications = scheduler.pop_notifications(owner="alice") + assert len(notifications) == 1 + assert notifications[0]["body"] == "Synthesized urgency reminder." + + +@pytest.mark.asyncio +async def test_action_cancellation_rolls_back_without_checkpoint( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, _runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + state_path = tmp_path / "email_urgency_state_alice.json" + state_path.write_text( + json.dumps({"owner": "alice", "per_uid": {}, "notified_uids": []}), + encoding="utf-8", + ) + entered = threading.Event() + release = threading.Event() + dispatch_cancelled = threading.Event() + + async def blocked_dispatch(**_kwargs): + entered.set() + try: + while not release.is_set(): + await asyncio.sleep(0.01) + except asyncio.CancelledError: + dispatch_cancelled.set() + raise + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", blocked_dispatch) + task = asyncio.create_task( + builtin_actions.action_check_email_urgency("alice") + ) + for _ in range(200): + if entered.is_set(): + break + await asyncio.sleep(0.01) + assert entered.is_set() + + task.cancel() + with pytest.raises(asyncio.CancelledError): + await task + release.set() + await asyncio.sleep(0.1) + + assert dispatch_cancelled.is_set() + assert json.loads(state_path.read_text(encoding="utf-8")) == { + "owner": "alice", + "per_uid": {}, + "notified_uids": [], + } + + +@pytest.mark.asyncio +async def test_account_scoped_actions_merge_disjoint_checkpoints( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + + await builtin_actions.action_check_email_urgency( + "alice", prompt='{"account_id":"acct-a"}' + ) + runtime["accounts"] = ["acct-b"] + await builtin_actions.action_check_email_urgency( + "alice", prompt='{"account_id":"acct-b"}' + ) + runtime["accounts"] = ["acct-a"] + await builtin_actions.action_check_email_urgency( + "alice", prompt='{"account_id":"acct-a"}' + ) + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert len(deliveries) == 2 + assert "Urgent request for acct-a" in deliveries[0] + assert "Urgent request for acct-b" in deliveries[1] + assert state["notified_uids"] == ["acct-a:1", "acct-b:1"] + assert set(state["per_uid"]) == {"acct-a:1", "acct-b:1"} + assert state["total_unread"] == 2 + assert state["total_urgent"] == 2 + + +@pytest.mark.asyncio +async def test_failed_account_scan_preserves_checkpoint_until_recovery( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a", "acct-b"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + + await builtin_actions.action_check_email_urgency("alice") + runtime["failures"] = {"acct-b"} + await builtin_actions.action_check_email_urgency("alice") + + failed_state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert failed_state["notified_uids"] == ["acct-a:1", "acct-b:1"] + assert set(failed_state["per_uid"]) == {"acct-a:1", "acct-b:1"} + + runtime["failures"] = set() + runtime["accounts"] = ["acct-b"] + await builtin_actions.action_check_email_urgency( + "alice", prompt='{"account_id":"acct-b"}' + ) + + assert len(deliveries) == 1 + + +@pytest.mark.asyncio +async def test_cached_flags_refresh_prunes_checkpoint_after_message_is_read( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + + async def delivered(**_kwargs): + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + await builtin_actions.action_check_email_urgency("alice") + runtime["seen_accounts"] = {"acct-a"} + await builtin_actions.action_check_email_urgency("alice") + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert state["notified_uids"] == [] + assert state["per_uid"]["acct-a:1"]["unread"] is False + assert state["total_unread"] == 0 From ffe94a7edfe93011e169c6f751a014aaaffe0e6a Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 21:00:58 +0000 Subject: [PATCH 3/7] fix(email): fence stale urgency scans --- src/builtin_actions.py | 134 ++++++++++++++++--- tests/test_email_urgency_checkpoint.py | 177 ++++++++++++++++++++++++- 2 files changed, 288 insertions(+), 23 deletions(-) diff --git a/src/builtin_actions.py b/src/builtin_actions.py index 936c6db60..b62524727 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -20,6 +20,64 @@ from src.interactive_gate import wait_for_interactive_quiet logger = logging.getLogger(__name__) +def _read_email_urgency_state(state_path): + """Read one atomic urgency checkpoint, tolerating the legacy shape.""" + from pathlib import Path + + state_path = Path(state_path) + try: + state = ( + json.loads(state_path.read_text(encoding="utf-8")) + if state_path.exists() + else {} + ) + except Exception: + return {} + return state if isinstance(state, dict) else {} + + +def _email_urgency_account_generations(state): + """Return normalized per-account checkpoint/complete generations. + + Checkpoint generations fence every accepted state mutation. Complete + generations advance only for a non-stale complete scan. Missing metadata + is the legacy generation zero. + """ + raw = state.get("account_generations", {}) if isinstance(state, dict) else {} + if not isinstance(raw, dict): + return {} + + generations = {} + for account_id, value in raw.items(): + if isinstance(value, dict): + checkpoint = value.get("checkpoint", 0) + complete = value.get("complete", 0) + else: + # Tolerate an intermediate scalar representation as one completed + # checkpoint generation instead of discarding its fence. + checkpoint = value + complete = value + try: + checkpoint = max(0, int(checkpoint)) + except (TypeError, ValueError): + checkpoint = 0 + try: + complete = max(0, int(complete)) + except (TypeError, ValueError): + complete = 0 + generations[str(account_id)] = { + "checkpoint": checkpoint, + "complete": complete, + } + return generations + + +def _email_urgency_string_set(value): + if not isinstance(value, (list, tuple, set, frozenset)): + return set() + return {str(item) for item in value if isinstance(item, (str, int))} + + def _acquire_email_urgency_state_lock( state_path, lock_db_path, @@ -60,17 +118,7 @@ def _acquire_email_urgency_state_lock( conn.rollback() conn.close() return None, None - try: - prior = ( - json.loads(state_path.read_text(encoding="utf-8")) - if state_path.exists() - else {} - ) - if not isinstance(prior, dict): - prior = {} - except Exception: - prior = {} - return conn, prior + return conn, _read_email_urgency_state(state_path) return None, None @@ -181,29 +229,72 @@ def _merge_email_urgency_state( notified_uids, all_unread_keys, fully_scanned_account_ids, + base_account_generations, timestamp, ): - """Merge current scan knowledge without deleting unscanned accounts.""" + """Merge a scan without letting an older snapshot erase newer facts.""" prior_per_uid = prior.get("per_uid", {}) if not isinstance(prior_per_uid, dict): prior_per_uid = {} complete = {str(account_id) for account_id in fully_scanned_account_ids} + prior_generations = _email_urgency_account_generations(prior) + base_generations = _email_urgency_account_generations( + {"account_generations": base_account_generations} + ) + observed_accounts = { + _email_urgency_account_key(key) for key in per_uid_scores + } | complete + stale_accounts = { + account_id + for account_id in observed_accounts + if prior_generations.get(account_id, {}).get("checkpoint", 0) + != base_generations.get(account_id, {}).get("checkpoint", 0) + } + fresh_complete = complete - stale_accounts + changed_accounts = set(fresh_complete) merged_per_uid = dict(prior_per_uid) for key in list(merged_per_uid): - if _email_urgency_account_key(key) in complete: + account_id = _email_urgency_account_key(key) + if account_id in fresh_complete: merged_per_uid.pop(key, None) + changed_accounts.add(account_id) # Partial scans may add or refresh facts, but absence from a partial scan - # is not evidence that another checkpoint or UI row is stale. - merged_per_uid.update(per_uid_scores) + # is not evidence that another checkpoint or UI row is stale. When another + # worker committed after this scan captured its base generation, its value + # also wins collisions; the stale scan may still add facts it alone saw. + for key, value in per_uid_scores.items(): + account_id = _email_urgency_account_key(key) + if account_id in stale_accounts and key in merged_per_uid: + continue + if merged_per_uid.get(key) != value: + changed_accounts.add(account_id) + merged_per_uid[key] = value - merged_notified = set(notified_uids) + prior_notified = _email_urgency_string_set(prior.get("notified_uids", [])) + merged_notified = prior_notified | _email_urgency_string_set(notified_uids) + for key in merged_notified - prior_notified: + changed_accounts.add(_email_urgency_account_key(key)) for key in list(merged_notified): if ( - _email_urgency_account_key(key) in complete + _email_urgency_account_key(key) in fresh_complete and key not in all_unread_keys ): merged_notified.discard(key) + changed_accounts.add(_email_urgency_account_key(key)) + + next_generations = { + account_id: dict(value) + for account_id, value in prior_generations.items() + } + for account_id in changed_accounts: + generation = next_generations.setdefault( + account_id, + {"checkpoint": 0, "complete": 0}, + ) + generation["checkpoint"] += 1 + if account_id in fresh_complete: + generation["complete"] += 1 total_unread = 0 total_urgent = 0 @@ -229,6 +320,7 @@ def _merge_email_urgency_state( "max_score": max_score, "per_uid": merged_per_uid, "notified_uids": sorted(merged_notified), + "account_generations": next_generations, } @@ -2134,6 +2226,13 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: if not accounts: raise TaskNoop("no email accounts configured") + # Capture each account's checkpoint generation before any IMAP work. + # The locked merge compares this basis with the latest committed state, + # so an older scan cannot destructively replace a newer worker's facts. + base_account_generations = _email_urgency_account_generations( + _read_email_urgency_state(STATE_PATH) + ) + urgency_prompt = settings.get("urgent_email_prompt", "") per_uid_scores = {} # key = ":" → {"score": 0-3, "reason": "..."} all_unread_keys = set() @@ -2714,6 +2813,7 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: notified_uids=notified_uids, all_unread_keys=all_unread_keys, fully_scanned_account_ids=fully_scanned_account_ids, + base_account_generations=base_account_generations, timestamp=_time.time(), ) return notified_uids, next_state diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py index ae7ef1dc3..b9ed330a2 100644 --- a/tests/test_email_urgency_checkpoint.py +++ b/tests/test_email_urgency_checkpoint.py @@ -45,10 +45,11 @@ class _EmailAccount: class _FakeImap: - def __init__(self, account_id, failures, seen_accounts): + def __init__(self, account_id, failures, seen_accounts, search_uids): self.account_id = account_id self.failures = failures self.seen_accounts = seen_accounts + self.search_uids = tuple(search_uids.get(account_id, ("1",))) def select(self, *_args, **_kwargs): if self.account_id in self.failures: @@ -59,16 +60,18 @@ class _FakeImap: if self.account_id in self.failures: raise RuntimeError(f"{self.account_id} unavailable") if command == "SEARCH": - return "OK", [b"1"] + return "OK", [" ".join(self.search_uids).encode()] + uid = _args[0] + uid = uid.decode() if isinstance(uid, bytes) else str(uid) query = str(_args[-1]) if _args else "" seen = "\\Seen" if self.account_id in self.seen_accounts else "" - flags = f"1 (UID 1 FLAGS ({seen}))".encode() + flags = f"{uid} (UID {uid} FLAGS ({seen}))".encode() if query == "(UID FLAGS)": return "OK", [flags] raw = ( f"From: Sender {self.account_id} \r\n" - f"Subject: Urgent request for {self.account_id}\r\n" - f"Message-ID: <{self.account_id}-1@example.com>\r\n" + f"Subject: Urgent request for {self.account_id} uid {uid}\r\n" + f"Message-ID: <{self.account_id}-{uid}@example.com>\r\n" "\r\n" "Please reply immediately." ).encode() @@ -97,6 +100,7 @@ def _configure_action(monkeypatch, tmp_path, account_ids): "accounts": list(account_ids), "failures": set(), "seen_accounts": set(), + "search_uids": {}, "settings": { "reminder_channel": "browser", "reminder_llm_synthesis": False, @@ -136,7 +140,10 @@ def _configure_action(monkeypatch, tmp_path, account_ids): email_helpers, "_imap_connect", lambda account_id=None, **_kwargs: _FakeImap( - str(account_id), runtime["failures"], runtime["seen_accounts"] + str(account_id), + runtime["failures"], + runtime["seen_accounts"], + runtime["search_uids"], ), ) return builtin_actions, runtime @@ -194,6 +201,94 @@ async def test_urgency_state_transaction_serializes_decision_and_checkpoint(tmp_ } +def test_stale_complete_scan_preserves_newer_account_facts_and_checkpoints(): + from src.builtin_actions import _merge_email_urgency_state + + prior = { + "owner": "alice", + "per_uid": { + "acct:1": {"score": 3, "unread": True, "reason": "newer"}, + "acct:2": {"score": 2, "unread": True, "reason": "new UID"}, + }, + "notified_uids": ["acct:1", "acct:2"], + "account_generations": { + "acct": {"checkpoint": 1, "complete": 1}, + }, + } + + merged = _merge_email_urgency_state( + prior, + owner="alice", + per_uid_scores={ + "acct:1": {"score": 0, "unread": False, "reason": "older"}, + "acct:3": {"score": 3, "unread": True, "reason": "also observed"}, + }, + notified_uids={"acct:1", "acct:2", "acct:3"}, + all_unread_keys={"acct:3"}, + fully_scanned_account_ids={"acct"}, + base_account_generations={ + "acct": {"checkpoint": 0, "complete": 0}, + }, + timestamp=300.0, + ) + + assert merged["per_uid"]["acct:1"]["reason"] == "newer" + assert set(merged["per_uid"]) == {"acct:1", "acct:2", "acct:3"} + assert merged["notified_uids"] == ["acct:1", "acct:2", "acct:3"] + assert merged["account_generations"]["acct"] == { + "checkpoint": 2, + "complete": 1, + } + + +def test_newer_partial_checkpoint_fences_older_complete_snapshot(): + from src.builtin_actions import _merge_email_urgency_state + + legacy = { + "owner": "alice", + "per_uid": { + "acct:1": {"score": 3, "unread": True}, + }, + "notified_uids": ["acct:1"], + } + partial = _merge_email_urgency_state( + legacy, + owner="alice", + per_uid_scores={ + "acct:2": {"score": 3, "unread": True}, + }, + notified_uids={"acct:1", "acct:2"}, + all_unread_keys={"acct:2"}, + fully_scanned_account_ids=set(), + base_account_generations={}, + timestamp=200.0, + ) + assert partial["account_generations"]["acct"] == { + "checkpoint": 1, + "complete": 0, + } + + merged = _merge_email_urgency_state( + partial, + owner="alice", + per_uid_scores={ + "acct:1": {"score": 3, "unread": True}, + }, + notified_uids={"acct:1", "acct:2"}, + all_unread_keys={"acct:1"}, + fully_scanned_account_ids={"acct"}, + base_account_generations={}, + timestamp=300.0, + ) + + assert set(merged["per_uid"]) == {"acct:1", "acct:2"} + assert merged["notified_uids"] == ["acct:1", "acct:2"] + assert merged["account_generations"]["acct"] == { + "checkpoint": 1, + "complete": 0, + } + + @pytest.mark.asyncio async def test_waiting_transaction_cancellation_does_not_leak_lock(tmp_path): from src.builtin_actions import _run_email_urgency_state_transaction @@ -372,6 +467,76 @@ async def test_action_cancellation_rolls_back_without_checkpoint( } +@pytest.mark.asyncio +async def test_older_same_account_scan_cannot_erase_newer_committed_uid( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + runtime["search_uids"]["acct-a"] = ["1"] + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + original_transaction = builtin_actions._run_email_urgency_state_transaction + stale_scan_ready = asyncio.Event() + release_stale_scan = asyncio.Event() + transaction_count = 0 + + async def order_transactions(state_path, lock_db_path, operation): + nonlocal transaction_count + transaction_count += 1 + if transaction_count == 1: + stale_scan_ready.set() + await release_stale_scan.wait() + return await original_transaction(state_path, lock_db_path, operation) + + monkeypatch.setattr( + builtin_actions, + "_run_email_urgency_state_transaction", + order_transactions, + ) + + stale = asyncio.create_task( + builtin_actions.action_check_email_urgency("alice") + ) + await asyncio.wait_for(stale_scan_ready.wait(), timeout=2) + + runtime["search_uids"]["acct-a"] = ["1", "2"] + newer_result = await asyncio.wait_for( + builtin_actions.action_check_email_urgency("alice"), + timeout=2, + ) + release_stale_scan.set() + stale_result = await asyncio.wait_for(stale, timeout=2) + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert newer_result[1] is True + assert stale_result[1] is True + assert len(deliveries) == 1 + assert "uid 2" in deliveries[0] + assert set(state["per_uid"]) == {"acct-a:1", "acct-a:2"} + assert state["notified_uids"] == ["acct-a:1", "acct-a:2"] + assert state["account_generations"]["acct-a"] == { + "checkpoint": 1, + "complete": 1, + } + + @pytest.mark.asyncio async def test_account_scoped_actions_merge_disjoint_checkpoints( monkeypatch, From 60aabbd9e24022e1377ac0fb7c2342bc1beb7393 Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 21:12:37 +0000 Subject: [PATCH 4/7] fix(email): fence stale urgency delivery --- src/builtin_actions.py | 145 +++++++++++++++++-------- tests/test_email_urgency_checkpoint.py | 104 ++++++++++++++++-- 2 files changed, 189 insertions(+), 60 deletions(-) diff --git a/src/builtin_actions.py b/src/builtin_actions.py index b62524727..d4c1d0c59 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -221,6 +221,23 @@ def _email_urgency_account_key(message_key): return str(message_key).split(":", 1)[0] +def _email_urgency_stale_accounts( + prior, + base_account_generations, + account_ids, +): + prior_generations = _email_urgency_account_generations(prior) + base_generations = _email_urgency_account_generations( + {"account_generations": base_account_generations} + ) + return { + str(account_id) + for account_id in account_ids + if prior_generations.get(str(account_id), {}).get("checkpoint", 0) + != base_generations.get(str(account_id), {}).get("checkpoint", 0) + } + + def _merge_email_urgency_state( prior, *, @@ -238,18 +255,14 @@ def _merge_email_urgency_state( prior_per_uid = {} complete = {str(account_id) for account_id in fully_scanned_account_ids} prior_generations = _email_urgency_account_generations(prior) - base_generations = _email_urgency_account_generations( - {"account_generations": base_account_generations} - ) observed_accounts = { _email_urgency_account_key(key) for key in per_uid_scores } | complete - stale_accounts = { - account_id - for account_id in observed_accounts - if prior_generations.get(account_id, {}).get("checkpoint", 0) - != base_generations.get(account_id, {}).get("checkpoint", 0) - } + stale_accounts = _email_urgency_stale_accounts( + prior, + base_account_generations, + observed_accounts, + ) fresh_complete = complete - stale_accounts changed_accounts = set(fresh_complete) @@ -261,20 +274,26 @@ def _merge_email_urgency_state( changed_accounts.add(account_id) # Partial scans may add or refresh facts, but absence from a partial scan # is not evidence that another checkpoint or UI row is stale. When another - # worker committed after this scan captured its base generation, its value - # also wins collisions; the stale scan may still add facts it alone saw. + # worker committed after this scan captured its base generation, discard + # this account's whole stale snapshot. A key absent from the newer state + # may have been removed/read, so even a stale-only key is not safely + # additive without another fresh scan. for key, value in per_uid_scores.items(): account_id = _email_urgency_account_key(key) - if account_id in stale_accounts and key in merged_per_uid: + if account_id in stale_accounts: continue if merged_per_uid.get(key) != value: changed_accounts.add(account_id) merged_per_uid[key] = value prior_notified = _email_urgency_string_set(prior.get("notified_uids", [])) - merged_notified = prior_notified | _email_urgency_string_set(notified_uids) - for key in merged_notified - prior_notified: - changed_accounts.add(_email_urgency_account_key(key)) + merged_notified = set(prior_notified) + for key in _email_urgency_string_set(notified_uids) - prior_notified: + account_id = _email_urgency_account_key(key) + if account_id in stale_accounts: + continue + merged_notified.add(key) + changed_accounts.add(account_id) for key in list(merged_notified): if ( _email_urgency_account_key(key) in fresh_complete @@ -2731,7 +2750,6 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # ── 4. Aggregate state. urgent = score ≥ 2. urgent_keys = [k for k, v in per_uid_scores.items() if v.get("score", 0) >= 2 and v.get("unread")] - total_urgent = len(urgent_keys) # ── 5. Fire a reminder only when a previously-unnotified UID scores # urgent. The read, decision, delivery, and checkpoint are serialized @@ -2739,39 +2757,46 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # state or overwrite each other's successful checkpoint. newly_notified = set() notify_failed = set() - title = "Urgent email" if total_urgent == 1 else f"{total_urgent} urgent emails" - # Build a real listing — subject · sender · reason for each urgent one - # — so the reminder tells the user which messages to act on. - sorted_urgent = sorted( - ((k, per_uid_scores[k]) for k in urgent_keys), - key=lambda kv: kv[1].get("score", 0), reverse=True, - )[:10] - _pub = (settings.get("app_public_url") or "").strip().rstrip("/") - from urllib.parse import quote as _quote - lines = [f"{total_urgent} email" + ("" if total_urgent == 1 else "s") + " need an urgent reply:", ""] - for i, (k, v) in enumerate(sorted_urgent, 1): - subj = (v.get("subject") or "(no subject)")[:160] - frm = v.get("from") or "" - why = v.get("reason") or "" - uid_for_link = str(k).split(":", 1)[-1] - hash_link = f"#email={_quote('INBOX', safe='')}:{uid_for_link}" - open_link = f"{_pub}/{hash_link}" if _pub else hash_link - line = f"{i}. {subj}" - if frm: - line += f" — {frm}" - if why: - line += f" · {why}" - lines.append(line) - lines.append(f" Open email: {open_link}") - if total_urgent > len(sorted_urgent): - lines.append("") - lines.append(f"…and {total_urgent - len(sorted_urgent)} more.") - body = "\n".join(lines) - async def _dispatch_urgency_reminder(): + def _urgency_reminder_payload(reminder_keys): + total = len(reminder_keys) + title = "Urgent email" if total == 1 else f"{total} urgent emails" + sorted_urgent = sorted( + ((key, per_uid_scores[key]) for key in reminder_keys), + key=lambda item: item[1].get("score", 0), + reverse=True, + )[:10] + _pub = (settings.get("app_public_url") or "").strip().rstrip("/") + from urllib.parse import quote as _quote + lines = [ + f"{total} email" + ("" if total == 1 else "s") + + " need an urgent reply:", + "", + ] + for i, (key, value) in enumerate(sorted_urgent, 1): + subj = (value.get("subject") or "(no subject)")[:160] + frm = value.get("from") or "" + why = value.get("reason") or "" + uid_for_link = str(key).split(":", 1)[-1] + hash_link = f"#email={_quote('INBOX', safe='')}:{uid_for_link}" + open_link = f"{_pub}/{hash_link}" if _pub else hash_link + line = f"{i}. {subj}" + if frm: + line += f" — {frm}" + if why: + line += f" · {why}" + lines.append(line) + lines.append(f" Open email: {open_link}") + if total > len(sorted_urgent): + lines.append("") + lines.append(f"…and {total - len(sorted_urgent)} more.") + return title, "\n".join(lines) + + async def _dispatch_urgency_reminder(reminder_keys): # Call dispatch_reminder directly: a scheduler has no browser # session cookie with which to call the HTTP endpoint. from routes.note_routes import dispatch_reminder + title, body = _urgency_reminder_payload(reminder_keys) return await dispatch_reminder( title=title, note_body=body, @@ -2780,11 +2805,35 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: ) async def _dispatch_and_checkpoint(prior): - notified_uids = set(prior.get("notified_uids", [])) - new_urgent = [k for k in urgent_keys if k not in notified_uids] + notified_uids = _email_urgency_string_set( + prior.get("notified_uids", []) + ) + observed_accounts = { + _email_urgency_account_key(key) for key in per_uid_scores + } | fully_scanned_account_ids + stale_accounts = _email_urgency_stale_accounts( + prior, + base_account_generations, + observed_accounts, + ) + # Generation fencing must happen before delivery, not only during + # merge. A stale-only unread UID may have been removed, read, or + # downgraded by the newer completed scan. + deliverable_urgent = [ + key + for key in urgent_keys + if _email_urgency_account_key(key) not in stale_accounts + ] + new_urgent = [ + key + for key in deliverable_urgent + if key not in notified_uids + ] if new_urgent: try: - dispatch_result = await _dispatch_urgency_reminder() + dispatch_result = await _dispatch_urgency_reminder( + deliverable_urgent + ) channel = (settings.get("reminder_channel") or "browser").strip().lower() delivered = bool(dispatch_result.get("browser_sent")) if channel == "email": diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py index b9ed330a2..5392776ae 100644 --- a/tests/test_email_urgency_checkpoint.py +++ b/tests/test_email_urgency_checkpoint.py @@ -201,16 +201,16 @@ async def test_urgency_state_transaction_serializes_decision_and_checkpoint(tmp_ } -def test_stale_complete_scan_preserves_newer_account_facts_and_checkpoints(): +def test_stale_complete_scan_preserves_newer_state_and_discards_stale_only_facts(): from src.builtin_actions import _merge_email_urgency_state prior = { "owner": "alice", "per_uid": { - "acct:1": {"score": 3, "unread": True, "reason": "newer"}, + "acct:1": {"score": 0, "unread": True, "reason": "newer"}, "acct:2": {"score": 2, "unread": True, "reason": "new UID"}, }, - "notified_uids": ["acct:1", "acct:2"], + "notified_uids": ["acct:2"], "account_generations": { "acct": {"checkpoint": 1, "complete": 1}, }, @@ -233,10 +233,10 @@ def test_stale_complete_scan_preserves_newer_account_facts_and_checkpoints(): ) assert merged["per_uid"]["acct:1"]["reason"] == "newer" - assert set(merged["per_uid"]) == {"acct:1", "acct:2", "acct:3"} - assert merged["notified_uids"] == ["acct:1", "acct:2", "acct:3"] + assert set(merged["per_uid"]) == {"acct:1", "acct:2"} + assert merged["notified_uids"] == ["acct:2"] assert merged["account_generations"]["acct"] == { - "checkpoint": 2, + "checkpoint": 1, "complete": 1, } @@ -467,17 +467,27 @@ async def test_action_cancellation_rolls_back_without_checkpoint( } +@pytest.mark.parametrize( + ("stale_uids", "newer_uids"), + [ + (["1"], ["1", "2"]), + (["1", "2"], ["1"]), + ], + ids=["preserve-newer-addition", "reject-stale-only-delivery"], +) @pytest.mark.asyncio -async def test_older_same_account_scan_cannot_erase_newer_committed_uid( +async def test_stale_same_account_scan_cannot_override_newer_commit( monkeypatch, tmp_path, + stale_uids, + newer_uids, ): import routes.note_routes as note_routes builtin_actions, runtime = _configure_action( monkeypatch, tmp_path, ["acct-a"] ) - runtime["search_uids"]["acct-a"] = ["1"] + runtime["search_uids"]["acct-a"] = stale_uids deliveries = [] async def delivered(**kwargs): @@ -514,7 +524,7 @@ async def test_older_same_account_scan_cannot_erase_newer_committed_uid( ) await asyncio.wait_for(stale_scan_ready.wait(), timeout=2) - runtime["search_uids"]["acct-a"] = ["1", "2"] + runtime["search_uids"]["acct-a"] = newer_uids newer_result = await asyncio.wait_for( builtin_actions.action_check_email_urgency("alice"), timeout=2, @@ -528,15 +538,85 @@ async def test_older_same_account_scan_cannot_erase_newer_committed_uid( assert newer_result[1] is True assert stale_result[1] is True assert len(deliveries) == 1 - assert "uid 2" in deliveries[0] - assert set(state["per_uid"]) == {"acct-a:1", "acct-a:2"} - assert state["notified_uids"] == ["acct-a:1", "acct-a:2"] + for uid in newer_uids: + assert f"uid {uid}" in deliveries[0] + for uid in set(stale_uids) - set(newer_uids): + assert f"uid {uid}" not in deliveries[0] + expected_keys = {f"acct-a:{uid}" for uid in newer_uids} + assert set(state["per_uid"]) == expected_keys + assert state["notified_uids"] == sorted(expected_keys) assert state["account_generations"]["acct-a"] == { "checkpoint": 1, "complete": 1, } +@pytest.mark.asyncio +async def test_mixed_stale_and_fresh_accounts_exclude_stale_rows_from_delivery( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a", "acct-b"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + original_transaction = builtin_actions._run_email_urgency_state_transaction + stale_scan_ready = asyncio.Event() + release_stale_scan = asyncio.Event() + transaction_count = 0 + + async def order_transactions(state_path, lock_db_path, operation): + nonlocal transaction_count + transaction_count += 1 + if transaction_count == 1: + stale_scan_ready.set() + await release_stale_scan.wait() + return await original_transaction(state_path, lock_db_path, operation) + + monkeypatch.setattr( + builtin_actions, + "_run_email_urgency_state_transaction", + order_transactions, + ) + + stale = asyncio.create_task( + builtin_actions.action_check_email_urgency("alice") + ) + await asyncio.wait_for(stale_scan_ready.wait(), timeout=2) + + runtime["accounts"] = ["acct-a"] + await asyncio.wait_for( + builtin_actions.action_check_email_urgency("alice"), + timeout=2, + ) + release_stale_scan.set() + await asyncio.wait_for(stale, timeout=2) + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert len(deliveries) == 2 + assert "acct-a" in deliveries[0] + assert "acct-b" not in deliveries[0] + assert "acct-b" in deliveries[1] + assert "acct-a" not in deliveries[1] + assert set(state["per_uid"]) == {"acct-a:1", "acct-b:1"} + assert state["notified_uids"] == ["acct-a:1", "acct-b:1"] + + @pytest.mark.asyncio async def test_account_scoped_actions_merge_disjoint_checkpoints( monkeypatch, From acc87e8f56a4e181cf277c67922dbf9657035e10 Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 21:35:06 +0000 Subject: [PATCH 5/7] fix(email): retire stale urgency accounts --- src/builtin_actions.py | 143 +++++++++++++-- tests/test_email_urgency_checkpoint.py | 244 +++++++++++++++++++++++++ 2 files changed, 369 insertions(+), 18 deletions(-) diff --git a/src/builtin_actions.py b/src/builtin_actions.py index d4c1d0c59..5418e31ae 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -221,6 +221,21 @@ def _email_urgency_account_key(message_key): return str(message_key).split(":", 1)[0] +def _email_urgency_payload_account_ids(state): + """Return account IDs that still own user-visible urgency payload.""" + if not isinstance(state, dict): + return set() + + per_uid = state.get("per_uid", {}) + per_uid_keys = per_uid if isinstance(per_uid, dict) else {} + return { + _email_urgency_account_key(key) for key in per_uid_keys + } | { + _email_urgency_account_key(key) + for key in _email_urgency_string_set(state.get("notified_uids", [])) + } + + def _email_urgency_stale_accounts( prior, base_account_generations, @@ -248,6 +263,8 @@ def _merge_email_urgency_state( fully_scanned_account_ids, base_account_generations, timestamp, + retired_account_ids=(), + base_payload_account_ids=(), ): """Merge a scan without letting an older snapshot erase newer facts.""" prior_per_uid = prior.get("per_uid", {}) @@ -255,18 +272,39 @@ def _merge_email_urgency_state( prior_per_uid = {} complete = {str(account_id) for account_id in fully_scanned_account_ids} prior_generations = _email_urgency_account_generations(prior) + retire_requested = {str(account_id) for account_id in retired_account_ids} observed_accounts = { _email_urgency_account_key(key) for key in per_uid_scores - } | complete + } | complete | retire_requested stale_accounts = _email_urgency_stale_accounts( prior, base_account_generations, observed_accounts, ) - fresh_complete = complete - stale_accounts + prior_payload_accounts = _email_urgency_payload_account_ids(prior) + base_payload_accounts = { + str(account_id) for account_id in base_payload_account_ids + } + # A selected account can be absent from the base snapshot. If another + # worker creates its first payload before this transaction wins the lock, + # membership itself is a fence even when both snapshots normalize to the + # legacy generation zero. + retired_accounts = { + account_id + for account_id in retire_requested - stale_accounts + if not ( + account_id in prior_payload_accounts + and account_id not in base_payload_accounts + ) + } + fresh_complete = complete - stale_accounts - retired_accounts changed_accounts = set(fresh_complete) - merged_per_uid = dict(prior_per_uid) + merged_per_uid = { + key: value + for key, value in prior_per_uid.items() + if _email_urgency_account_key(key) not in retired_accounts + } for key in list(merged_per_uid): account_id = _email_urgency_account_key(key) if account_id in fresh_complete: @@ -280,17 +318,21 @@ def _merge_email_urgency_state( # additive without another fresh scan. for key, value in per_uid_scores.items(): account_id = _email_urgency_account_key(key) - if account_id in stale_accounts: + if account_id in stale_accounts or account_id in retired_accounts: continue if merged_per_uid.get(key) != value: changed_accounts.add(account_id) merged_per_uid[key] = value prior_notified = _email_urgency_string_set(prior.get("notified_uids", [])) - merged_notified = set(prior_notified) + merged_notified = { + key + for key in prior_notified + if _email_urgency_account_key(key) not in retired_accounts + } for key in _email_urgency_string_set(notified_uids) - prior_notified: account_id = _email_urgency_account_key(key) - if account_id in stale_accounts: + if account_id in stale_accounts or account_id in retired_accounts: continue merged_notified.add(key) changed_accounts.add(account_id) @@ -314,6 +356,19 @@ def _merge_email_urgency_state( generation["checkpoint"] += 1 if account_id in fresh_complete: generation["complete"] += 1 + for account_id in retired_accounts: + # Keep a generation-only tombstone. It prevents a scan that started + # before deletion/disable from resurrecting the retired account, while + # a later re-enabled account can capture this checkpoint and replace + # the tombstone normally. Do not advance an already-empty tombstone on + # every cleanup pass. + had_generation = account_id in prior_generations + if account_id in prior_payload_accounts or not had_generation: + generation = next_generations.setdefault( + account_id, + {"checkpoint": 0, "complete": 0}, + ) + generation["checkpoint"] += 1 total_unread = 0 total_urgent = 0 @@ -2216,16 +2271,14 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: "shopping", "social", "work", "personal", "legal", "support", "promo", } - # ── 1. Resolve LLM candidates (utility primary + utility fallbacks; fall - # through to default chat as a last resort). + # Resolve with the task owner as before, but defer the availability + # gate until after authoritative account cleanup. State retirement must + # still run when no model is configured. from src.task_endpoint import resolve_task_candidates candidates = resolve_task_candidates(owner=owner) - if not candidates: - return "No LLM endpoint available", False - target_account_id = _email_task_account_id(kwargs) - # ── 2. Enumerate enabled accounts. Match this task's owner AND fall + # ── 1. Enumerate enabled accounts. Match this task's owner AND fall # back to the legacy "unowned account whose imap_user / from_address # == this owner" pattern — same rule `_get_email_config` uses, so a # pre-multi-user account row still gets picked up for the seeded task. @@ -2242,15 +2295,69 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: accounts = q.all() finally: db.close() + + # Capture the checkpoint basis before any cleanup or IMAP work. A full + # owner-wide enumeration authoritatively retires payload belonging to + # accounts that are no longer enabled/visible. A scoped task may retire + # only its selected missing/disabled account. Existing accounts remain + # present here even if their later network scan fails, so transient + # IMAP failure never erases their last known state. + base_state = _read_email_urgency_state(STATE_PATH) + base_account_generations = _email_urgency_account_generations( + base_state + ) + base_payload_account_ids = _email_urgency_payload_account_ids(base_state) + enabled_account_ids = {str(account.id) for account in accounts} + if target_account_id: + retired_account_ids = ( + {str(target_account_id)} + if str(target_account_id) not in enabled_account_ids + else set() + ) + else: + retired_account_ids = ( + base_payload_account_ids - enabled_account_ids + ) + + if retired_account_ids: + async def _retire_accounts(prior): + next_state = _merge_email_urgency_state( + prior, + owner=owner, + per_uid_scores={}, + notified_uids=prior.get("notified_uids", []), + all_unread_keys=set(), + fully_scanned_account_ids=set(), + base_account_generations=base_account_generations, + timestamp=_time.time(), + retired_account_ids=retired_account_ids, + base_payload_account_ids=base_payload_account_ids, + ) + return None, next_state + + await _run_email_urgency_state_transaction( + STATE_PATH, + STATE_LOCK_DB, + _retire_accounts, + ) + # Cleanup may have advanced tombstone generations. Capture the + # exact committed basis that the subsequent scan must compare. + base_state = _read_email_urgency_state(STATE_PATH) + base_account_generations = _email_urgency_account_generations( + base_state + ) + base_payload_account_ids = _email_urgency_payload_account_ids( + base_state + ) + if not accounts: raise TaskNoop("no email accounts configured") - # Capture each account's checkpoint generation before any IMAP work. - # The locked merge compares this basis with the latest committed state, - # so an older scan cannot destructively replace a newer worker's facts. - base_account_generations = _email_urgency_account_generations( - _read_email_urgency_state(STATE_PATH) - ) + # ── 2. Account retirement above is state maintenance and does not + # depend on model availability. Scanning still requires the utility + # primary/fallback candidates resolved for this task owner. + if not candidates: + return "No LLM endpoint available", False urgency_prompt = settings.get("urgent_email_prompt", "") per_uid_scores = {} # key = ":" → {"score": 0-3, "reason": "..."} diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py index 5392776ae..2053a58f7 100644 --- a/tests/test_email_urgency_checkpoint.py +++ b/tests/test_email_urgency_checkpoint.py @@ -289,6 +289,104 @@ def test_newer_partial_checkpoint_fences_older_complete_snapshot(): } +def test_authoritative_retirement_prunes_payload_and_recomputes_api_totals(): + from src.builtin_actions import _merge_email_urgency_state + + prior = { + "owner": "alice", + "per_uid": { + "acct-a:1": {"score": 1, "unread": True}, + "acct-b:1": {"score": 3, "unread": True}, + }, + "notified_uids": ["acct-b:1"], + "account_generations": { + "acct-a": {"checkpoint": 1, "complete": 1}, + "acct-b": {"checkpoint": 1, "complete": 1}, + }, + } + + merged = _merge_email_urgency_state( + prior, + owner="alice", + per_uid_scores={}, + notified_uids=prior["notified_uids"], + all_unread_keys=set(), + fully_scanned_account_ids=set(), + base_account_generations=prior["account_generations"], + timestamp=300.0, + retired_account_ids={"acct-b"}, + base_payload_account_ids={"acct-a", "acct-b"}, + ) + + assert merged["per_uid"] == { + "acct-a:1": {"score": 1, "unread": True}, + } + assert merged["notified_uids"] == [] + assert merged["total_unread"] == 1 + assert merged["total_urgent"] == 0 + assert merged["max_score"] == 1 + assert merged["account_generations"]["acct-b"] == { + "checkpoint": 2, + "complete": 1, + } + + +def test_authoritative_retirement_respects_generation_and_membership_fences(): + from src.builtin_actions import _merge_email_urgency_state + + newer = { + "owner": "alice", + "per_uid": { + "acct-b:2": {"score": 3, "unread": True, "reason": "newer"}, + }, + "notified_uids": ["acct-b:2"], + "account_generations": { + "acct-b": {"checkpoint": 2, "complete": 2}, + }, + } + generation_fenced = _merge_email_urgency_state( + newer, + owner="alice", + per_uid_scores={}, + notified_uids=newer["notified_uids"], + all_unread_keys=set(), + fully_scanned_account_ids=set(), + base_account_generations={ + "acct-b": {"checkpoint": 1, "complete": 1}, + }, + timestamp=300.0, + retired_account_ids={"acct-b"}, + base_payload_account_ids={"acct-b"}, + ) + assert generation_fenced["per_uid"] == newer["per_uid"] + assert generation_fenced["notified_uids"] == ["acct-b:2"] + assert generation_fenced["account_generations"]["acct-b"] == { + "checkpoint": 2, + "complete": 2, + } + + legacy_first_write = { + "owner": "alice", + "per_uid": { + "acct-b:3": {"score": 2, "unread": True, "reason": "concurrent"}, + }, + "notified_uids": [], + } + membership_fenced = _merge_email_urgency_state( + legacy_first_write, + owner="alice", + per_uid_scores={}, + notified_uids=[], + all_unread_keys=set(), + fully_scanned_account_ids=set(), + base_account_generations={}, + timestamp=300.0, + retired_account_ids={"acct-b"}, + base_payload_account_ids=set(), + ) + assert membership_fenced["per_uid"] == legacy_first_write["per_uid"] + + @pytest.mark.asyncio async def test_waiting_transaction_cancellation_does_not_leak_lock(tmp_path): from src.builtin_actions import _run_email_urgency_state_transaction @@ -664,6 +762,146 @@ async def test_account_scoped_actions_merge_disjoint_checkpoints( assert state["total_urgent"] == 2 +@pytest.mark.asyncio +async def test_full_enumeration_retires_deleted_or_disabled_account( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a", "acct-b"] + ) + + async def delivered(**_kwargs): + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + await builtin_actions.action_check_email_urgency("alice") + + # The production query returns only enabled, owner-visible accounts. A + # deleted row and a disabled row are therefore the same authoritative + # absence at this boundary. + runtime["accounts"] = ["acct-a"] + await builtin_actions.action_check_email_urgency("alice") + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert set(state["per_uid"]) == {"acct-a:1"} + assert state["notified_uids"] == ["acct-a:1"] + assert state["total_unread"] == 1 + assert state["total_urgent"] == 1 + assert state["max_score"] == 3 + assert state["account_generations"]["acct-b"] == { + "checkpoint": 2, + "complete": 1, + } + + +@pytest.mark.asyncio +async def test_zero_enabled_accounts_cleanup_precedes_model_resolution( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + from src import task_endpoint + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + + async def delivered(**_kwargs): + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + await builtin_actions.action_check_email_urgency("alice") + runtime["accounts"] = [] + + model_resolution_owners = [] + + def no_model_available(*_args, **kwargs): + model_resolution_owners.append(kwargs.get("owner")) + return [] + + monkeypatch.setattr( + task_endpoint, + "resolve_task_candidates", + no_model_available, + ) + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency("alice") + assert model_resolution_owners == ["alice"] + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert state["per_uid"] == {} + assert state["notified_uids"] == [] + assert state["total_unread"] == 0 + assert state["total_urgent"] == 0 + assert state["max_score"] == 0 + assert state["account_generations"]["acct-a"] == { + "checkpoint": 2, + "complete": 1, + } + + +@pytest.mark.asyncio +async def test_scoped_missing_account_retires_only_selected_payload( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a", "acct-b"] + ) + + async def delivered(**_kwargs): + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + await builtin_actions.action_check_email_urgency("alice") + + # A production account-id filter returns no row when the selected account + # was deleted or disabled. Other accounts are outside this scoped query and + # must remain untouched. + runtime["accounts"] = [] + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency( + "alice", + prompt='{"account_id":"acct-b"}', + ) + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert set(state["per_uid"]) == {"acct-a:1"} + assert state["notified_uids"] == ["acct-a:1"] + assert state["total_unread"] == 1 + assert state["total_urgent"] == 1 + assert state["account_generations"]["acct-b"] == { + "checkpoint": 2, + "complete": 1, + } + + @pytest.mark.asyncio async def test_failed_account_scan_preserves_checkpoint_until_recovery( monkeypatch, @@ -696,6 +934,12 @@ async def test_failed_account_scan_preserves_checkpoint_until_recovery( ) assert failed_state["notified_uids"] == ["acct-a:1", "acct-b:1"] assert set(failed_state["per_uid"]) == {"acct-a:1", "acct-b:1"} + assert failed_state["total_unread"] == 2 + assert failed_state["total_urgent"] == 2 + assert failed_state["account_generations"]["acct-b"] == { + "checkpoint": 1, + "complete": 1, + } runtime["failures"] = set() runtime["accounts"] = ["acct-b"] From a7ca86495e1b6929188972601e2aa53eb918999e Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 21:52:26 +0000 Subject: [PATCH 6/7] fix(email): fence urgency account retirement --- src/builtin_actions.py | 130 ++++++++---- tests/test_email_urgency_checkpoint.py | 275 +++++++++++++++++++++++-- 2 files changed, 357 insertions(+), 48 deletions(-) diff --git a/src/builtin_actions.py b/src/builtin_actions.py index 5418e31ae..bdba9286b 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -236,6 +236,13 @@ def _email_urgency_payload_account_ids(state): } +def _email_urgency_known_account_ids(state): + """Return payload owners plus generation-only active/retired markers.""" + return _email_urgency_payload_account_ids(state) | set( + _email_urgency_account_generations(state) + ) + + def _email_urgency_stale_accounts( prior, base_account_generations, @@ -265,6 +272,7 @@ def _merge_email_urgency_state( timestamp, retired_account_ids=(), base_payload_account_ids=(), + known_account_ids=(), ): """Merge a scan without letting an older snapshot erase newer facts.""" prior_per_uid = prior.get("per_uid", {}) @@ -356,19 +364,22 @@ def _merge_email_urgency_state( generation["checkpoint"] += 1 if account_id in fresh_complete: generation["complete"] += 1 + for account_id in {str(value) for value in known_account_ids}: + next_generations.setdefault( + account_id, + {"checkpoint": 0, "complete": 0}, + ) for account_id in retired_accounts: - # Keep a generation-only tombstone. It prevents a scan that started - # before deletion/disable from resurrecting the retired account, while - # a later re-enabled account can capture this checkpoint and replace - # the tombstone normally. Do not advance an already-empty tombstone on - # every cleanup pass. - had_generation = account_id in prior_generations - if account_id in prior_payload_accounts or not had_generation: - generation = next_generations.setdefault( - account_id, - {"checkpoint": 0, "complete": 0}, - ) - generation["checkpoint"] += 1 + # Every authoritative absence advances its generation, even when the + # prior state is already a payload-empty tombstone. A re-enabled scan + # may have captured that previous tombstone immediately before the + # account was disabled/deleted again; monotonic advancement is what + # makes that in-flight scan stale. + generation = next_generations.setdefault( + account_id, + {"checkpoint": 0, "complete": 0}, + ) + generation["checkpoint"] += 1 total_unread = 0 total_urgent = 0 @@ -2282,32 +2293,84 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # back to the legacy "unowned account whose imap_user / from_address # == this owner" pattern — same rule `_get_email_config` uses, so a # pre-multi-user account row still gets picked up for the seeded task. - db = _SL() - try: - from sqlalchemy import and_ as _and, or_ as _or - q = db.query(_EA).filter(_EA.enabled == True) # noqa: E712 - if owner: - unowned = _or(_EA.owner == None, _EA.owner == "") # noqa: E711 - same_mailbox = _or(_EA.imap_user == owner, _EA.from_address == owner) - q = q.filter(_or(_EA.owner == owner, _and(unowned, same_mailbox))) - if target_account_id: - q = q.filter(_EA.id == target_account_id) - accounts = q.all() - finally: - db.close() + def _enumerate_enabled_accounts(): + db = _SL() + try: + from sqlalchemy import and_ as _and, or_ as _or + q = db.query(_EA).filter(_EA.enabled == True) # noqa: E712 + if owner: + unowned = _or(_EA.owner == None, _EA.owner == "") # noqa: E711 + same_mailbox = _or( + _EA.imap_user == owner, + _EA.from_address == owner, + ) + q = q.filter( + _or(_EA.owner == owner, _and(unowned, same_mailbox)) + ) + if target_account_id: + q = q.filter(_EA.id == target_account_id) + return q.all() + finally: + db.close() - # Capture the checkpoint basis before any cleanup or IMAP work. A full - # owner-wide enumeration authoritatively retires payload belonging to - # accounts that are no longer enabled/visible. A scoped task may retire + initial_accounts = _enumerate_enabled_accounts() + initial_account_ids = { + str(account.id) for account in initial_accounts + } + + # Register every account before IMAP work, including its first-ever + # scan. A concurrent zero-account cleanup can then advance this marker + # and fence delivery even before the scan has produced payload. + if initial_account_ids: + async def _register_accounts(prior): + next_state = _merge_email_urgency_state( + prior, + owner=owner, + per_uid_scores={}, + notified_uids=prior.get("notified_uids", []), + all_unread_keys=set(), + fully_scanned_account_ids=set(), + base_account_generations=( + _email_urgency_account_generations(prior) + ), + timestamp=_time.time(), + known_account_ids=initial_account_ids, + ) + return None, next_state + + await _run_email_urgency_state_transaction( + STATE_PATH, + STATE_LOCK_DB, + _register_accounts, + ) + + # Revalidate after registration. If deletion/disable and its cleanup + # completed before the marker was published, this second enumeration + # observes the absence and this action retires its own marker instead + # of starting IMAP. Accounts newly appearing between the two reads are + # left for the next pass rather than scanned without prior registration. + verified_accounts = _enumerate_enabled_accounts() + enabled_account_ids = { + str(account.id) for account in verified_accounts + } + accounts = [ + account + for account in verified_accounts + if str(account.id) in initial_account_ids + ] + + # Capture the checkpoint basis before cleanup or IMAP. A full + # owner-wide enumeration authoritatively retires all known state IDs + # absent from the current enabled/visible set. A scoped task may retire # only its selected missing/disabled account. Existing accounts remain - # present here even if their later network scan fails, so transient - # IMAP failure never erases their last known state. + # present even if their later network scan fails, so transient IMAP + # failure never erases their last known state. base_state = _read_email_urgency_state(STATE_PATH) base_account_generations = _email_urgency_account_generations( base_state ) base_payload_account_ids = _email_urgency_payload_account_ids(base_state) - enabled_account_ids = {str(account.id) for account in accounts} + known_state_account_ids = _email_urgency_known_account_ids(base_state) if target_account_id: retired_account_ids = ( {str(target_account_id)} @@ -2316,7 +2379,7 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: ) else: retired_account_ids = ( - base_payload_account_ids - enabled_account_ids + known_state_account_ids - enabled_account_ids ) if retired_account_ids: @@ -2346,9 +2409,6 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: base_account_generations = _email_urgency_account_generations( base_state ) - base_payload_account_ids = _email_urgency_payload_account_ids( - base_state - ) if not accounts: raise TaskNoop("no email accounts configured") diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py index 2053a58f7..80503c3ec 100644 --- a/tests/test_email_urgency_checkpoint.py +++ b/tests/test_email_urgency_checkpoint.py @@ -558,10 +558,254 @@ async def test_action_cancellation_rolls_back_without_checkpoint( await asyncio.sleep(0.1) assert dispatch_cancelled.is_set() - assert json.loads(state_path.read_text(encoding="utf-8")) == { - "owner": "alice", - "per_uid": {}, - "notified_uids": [], + state = json.loads(state_path.read_text(encoding="utf-8")) + assert state["owner"] == "alice" + assert state["per_uid"] == {} + assert state["notified_uids"] == [] + # The pre-scan active marker is not a delivered/checkpointed UID. It must + # survive cancellation so a concurrent deletion cleanup can fence this + # first-ever account scan. + assert state["account_generations"] == { + "acct-a": {"checkpoint": 0, "complete": 0}, + } + + +@pytest.mark.asyncio +async def test_first_scan_revalidates_account_deleted_before_registration( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + original_transaction = builtin_actions._run_email_urgency_state_transaction + + async def delete_before_registration(state_path, lock_db_path, operation): + if operation.__name__ == "_register_accounts": + # The initial DB read saw the account, but deletion commits before + # its first active marker. The post-registration enumeration must + # observe that absence and retire the marker without touching IMAP. + runtime["accounts"] = [] + return await original_transaction(state_path, lock_db_path, operation) + + monkeypatch.setattr( + builtin_actions, + "_run_email_urgency_state_transaction", + delete_before_registration, + ) + + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency("alice") + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert deliveries == [] + assert state["per_uid"] == {} + assert state["notified_uids"] == [] + assert state["account_generations"]["acct-a"] == { + "checkpoint": 1, + "complete": 0, + } + + +@pytest.mark.asyncio +async def test_first_scan_deleted_after_revalidation_is_fenced_by_cleanup_pass( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + original_transaction = builtin_actions._run_email_urgency_state_transaction + scanned = asyncio.Event() + release_scan = asyncio.Event() + paused = False + + async def pause_first_delivery(state_path, lock_db_path, operation): + nonlocal paused + if operation.__name__ == "_dispatch_and_checkpoint" and not paused: + paused = True + scanned.set() + await release_scan.wait() + return await original_transaction(state_path, lock_db_path, operation) + + monkeypatch.setattr( + builtin_actions, + "_run_email_urgency_state_transaction", + pause_first_delivery, + ) + + stale = asyncio.create_task( + builtin_actions.action_check_email_urgency("alice") + ) + await asyncio.wait_for(scanned.wait(), timeout=2) + registered = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert registered["account_generations"]["acct-a"] == { + "checkpoint": 0, + "complete": 0, + } + + # Public-action boundary: deletion itself does not mutate urgency state. + # The authoritative zero-account pass is what advances the active marker + # before the paused scan reaches any reminder channel. + runtime["accounts"] = [] + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency("alice") + release_scan.set() + await asyncio.wait_for(stale, timeout=2) + + state = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert deliveries == [] + assert state["per_uid"] == {} + assert state["notified_uids"] == [] + assert state["total_unread"] == 0 + assert state["total_urgent"] == 0 + assert state["account_generations"]["acct-a"] == { + "checkpoint": 1, + "complete": 0, + } + + +@pytest.mark.asyncio +async def test_payload_empty_tombstone_fences_reenable_redelete_and_can_recover( + monkeypatch, + tmp_path, +): + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + await builtin_actions.action_check_email_urgency("alice") + assert len(deliveries) == 1 + + runtime["accounts"] = [] + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency("alice") + first_tombstone = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert first_tombstone["per_uid"] == {} + assert first_tombstone["account_generations"]["acct-a"] == { + "checkpoint": 2, + "complete": 1, + } + + original_transaction = builtin_actions._run_email_urgency_state_transaction + scanned = asyncio.Event() + release_scan = asyncio.Event() + pause_reenabled = True + + async def pause_reenabled_delivery(state_path, lock_db_path, operation): + nonlocal pause_reenabled + if operation.__name__ == "_dispatch_and_checkpoint" and pause_reenabled: + pause_reenabled = False + scanned.set() + await release_scan.wait() + return await original_transaction(state_path, lock_db_path, operation) + + monkeypatch.setattr( + builtin_actions, + "_run_email_urgency_state_transaction", + pause_reenabled_delivery, + ) + deliveries.clear() + runtime["accounts"] = ["acct-a"] + stale_reenabled = asyncio.create_task( + builtin_actions.action_check_email_urgency("alice") + ) + await asyncio.wait_for(scanned.wait(), timeout=2) + + # Delete/disable again while the re-enabled scan is based on checkpoint 2. + # Discovery must include the payload-empty generation tombstone and advance + # it, otherwise the paused scan would deliver and resurrect acct-a:1. + runtime["accounts"] = [] + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency("alice") + release_scan.set() + await asyncio.wait_for(stale_reenabled, timeout=2) + + fenced = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert deliveries == [] + assert fenced["per_uid"] == {} + assert fenced["notified_uids"] == [] + assert fenced["account_generations"]["acct-a"] == { + "checkpoint": 3, + "complete": 1, + } + + # A later authoritative pass advances the empty tombstone again, fencing + # any scan that captured checkpoint 3 before this absence was confirmed. + with pytest.raises(builtin_actions.TaskNoop): + await builtin_actions.action_check_email_urgency("alice") + repeated = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert repeated["account_generations"]["acct-a"] == { + "checkpoint": 4, + "complete": 1, + } + + # Re-enabling after the latest tombstone captures checkpoint 4 and can + # publish fresh payload normally. + runtime["accounts"] = ["acct-a"] + await builtin_actions.action_check_email_urgency("alice") + recovered = json.loads( + (tmp_path / "email_urgency_state_alice.json").read_text(encoding="utf-8") + ) + assert len(deliveries) == 1 + assert set(recovered["per_uid"]) == {"acct-a:1"} + assert recovered["notified_uids"] == ["acct-a:1"] + assert recovered["account_generations"]["acct-a"] == { + "checkpoint": 5, + "complete": 2, } @@ -605,10 +849,11 @@ async def test_stale_same_account_scan_cannot_override_newer_commit( async def order_transactions(state_path, lock_db_path, operation): nonlocal transaction_count - transaction_count += 1 - if transaction_count == 1: - stale_scan_ready.set() - await release_stale_scan.wait() + if operation.__name__ == "_dispatch_and_checkpoint": + transaction_count += 1 + if transaction_count == 1: + stale_scan_ready.set() + await release_stale_scan.wait() return await original_transaction(state_path, lock_db_path, operation) monkeypatch.setattr( @@ -678,10 +923,11 @@ async def test_mixed_stale_and_fresh_accounts_exclude_stale_rows_from_delivery( async def order_transactions(state_path, lock_db_path, operation): nonlocal transaction_count - transaction_count += 1 - if transaction_count == 1: - stale_scan_ready.set() - await release_stale_scan.wait() + if operation.__name__ == "_dispatch_and_checkpoint": + transaction_count += 1 + if transaction_count == 1: + stale_scan_ready.set() + await release_stale_scan.wait() return await original_transaction(state_path, lock_db_path, operation) monkeypatch.setattr( @@ -697,7 +943,10 @@ async def test_mixed_stale_and_fresh_accounts_exclude_stale_rows_from_delivery( runtime["accounts"] = ["acct-a"] await asyncio.wait_for( - builtin_actions.action_check_email_urgency("alice"), + builtin_actions.action_check_email_urgency( + "alice", + prompt='{"account_id":"acct-a"}', + ), timeout=2, ) release_stale_scan.set() From 2ac68f683dfa3a0930f90b063072be62c5bbee86 Mon Sep 17 00:00:00 2001 From: RaresKeY <158580472+RaresKeY@users.noreply.github.com> Date: Mon, 20 Jul 2026 22:03:52 +0000 Subject: [PATCH 7/7] fix(email): retain urgency registration generation --- src/builtin_actions.py | 21 +++--- tests/test_email_urgency_checkpoint.py | 97 ++++++++++++++++++++++++++ 2 files changed, 108 insertions(+), 10 deletions(-) diff --git a/src/builtin_actions.py b/src/builtin_actions.py index bdba9286b..216bb0360 100644 --- a/src/builtin_actions.py +++ b/src/builtin_actions.py @@ -2321,6 +2321,7 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # Register every account before IMAP work, including its first-ever # scan. A concurrent zero-account cleanup can then advance this marker # and fence delivery even before the scan has produced payload. + registered_state = None if initial_account_ids: async def _register_accounts(prior): next_state = _merge_email_urgency_state( @@ -2336,9 +2337,12 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: timestamp=_time.time(), known_account_ids=initial_account_ids, ) - return None, next_state + # Return the exact state committed by registration. This is + # the scan's generation token: adopting a later checkpoint + # after account cleanup would let the stale scan appear fresh. + return next_state, next_state - await _run_email_urgency_state_transaction( + registered_state = await _run_email_urgency_state_transaction( STATE_PATH, STATE_LOCK_DB, _register_accounts, @@ -2365,7 +2369,11 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: # only its selected missing/disabled account. Existing accounts remain # present even if their later network scan fails, so transient IMAP # failure never erases their last known state. - base_state = _read_email_urgency_state(STATE_PATH) + base_state = ( + registered_state + if registered_state is not None + else _read_email_urgency_state(STATE_PATH) + ) base_account_generations = _email_urgency_account_generations( base_state ) @@ -2403,13 +2411,6 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]: STATE_LOCK_DB, _retire_accounts, ) - # Cleanup may have advanced tombstone generations. Capture the - # exact committed basis that the subsequent scan must compare. - base_state = _read_email_urgency_state(STATE_PATH) - base_account_generations = _email_urgency_account_generations( - base_state - ) - if not accounts: raise TaskNoop("no email accounts configured") diff --git a/tests/test_email_urgency_checkpoint.py b/tests/test_email_urgency_checkpoint.py index 80503c3ec..7fcec08a7 100644 --- a/tests/test_email_urgency_checkpoint.py +++ b/tests/test_email_urgency_checkpoint.py @@ -699,6 +699,103 @@ async def test_first_scan_deleted_after_revalidation_is_fenced_by_cleanup_pass( } +def test_scan_keeps_registration_generation_when_cleanup_precedes_basis( + monkeypatch, + tmp_path, +): + """Cleanup after verification cannot become the stale scan's baseline.""" + from core import database + import routes.note_routes as note_routes + + builtin_actions, runtime = _configure_action( + monkeypatch, tmp_path, ["acct-a"] + ) + deliveries = [] + + async def delivered(**kwargs): + deliveries.append(kwargs["note_body"]) + return { + "browser_sent": True, + "email_sent": False, + "ntfy_sent": False, + "webhook_sent": False, + } + + monkeypatch.setattr(note_routes, "dispatch_reminder", delivered) + verified_account_selected = threading.Event() + release_verified_account = threading.Event() + stale_query_count = 0 + query_count_lock = threading.Lock() + + class _PausingQuery(_Query): + def all(self): + nonlocal stale_query_count + rows = list(self._rows) + if threading.current_thread().name == "stale-urgency-scan": + with query_count_lock: + stale_query_count += 1 + should_pause = stale_query_count == 2 + if should_pause: + # The second enumeration has selected the enabled row, but + # the action has not yet retained/used its checkpoint basis. + verified_account_selected.set() + assert release_verified_account.wait(5) + return rows + + class _PausingDb(_Db): + def query(self, _model): + return _PausingQuery(self._rows()) + + monkeypatch.setattr( + database, + "SessionLocal", + lambda: _PausingDb( + lambda: [_account(value) for value in runtime["accounts"]] + ), + ) + + stale_result = {} + + def run_stale_scan(): + try: + stale_result["value"] = asyncio.run( + builtin_actions.action_check_email_urgency("alice") + ) + except BaseException as exc: + stale_result["error"] = exc + + worker = threading.Thread( + target=run_stale_scan, + name="stale-urgency-scan", + ) + worker.start() + assert verified_account_selected.wait(5) + + runtime["accounts"] = [] + with pytest.raises(builtin_actions.TaskNoop): + asyncio.run(builtin_actions.action_check_email_urgency("alice")) + state_path = tmp_path / "email_urgency_state_alice.json" + retired = json.loads(state_path.read_text(encoding="utf-8")) + assert retired["account_generations"]["acct-a"] == { + "checkpoint": 1, + "complete": 0, + } + + release_verified_account.set() + worker.join(timeout=5) + assert not worker.is_alive() + assert "error" not in stale_result + + state = json.loads(state_path.read_text(encoding="utf-8")) + assert deliveries == [] + assert state["per_uid"] == {} + assert state["notified_uids"] == [] + assert state["account_generations"]["acct-a"] == { + "checkpoint": 1, + "complete": 0, + } + + @pytest.mark.asyncio async def test_payload_empty_tombstone_fences_reenable_redelete_and_can_recover( monkeypatch,