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] 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,