mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-08-05 02:45:28 +00:00
fix(email): retain urgency registration generation
This commit is contained in:
parent
a7ca86495e
commit
2ac68f683d
2 changed files with 108 additions and 10 deletions
|
|
@ -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
|
# Register every account before IMAP work, including its first-ever
|
||||||
# scan. A concurrent zero-account cleanup can then advance this marker
|
# scan. A concurrent zero-account cleanup can then advance this marker
|
||||||
# and fence delivery even before the scan has produced payload.
|
# and fence delivery even before the scan has produced payload.
|
||||||
|
registered_state = None
|
||||||
if initial_account_ids:
|
if initial_account_ids:
|
||||||
async def _register_accounts(prior):
|
async def _register_accounts(prior):
|
||||||
next_state = _merge_email_urgency_state(
|
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(),
|
timestamp=_time.time(),
|
||||||
known_account_ids=initial_account_ids,
|
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_PATH,
|
||||||
STATE_LOCK_DB,
|
STATE_LOCK_DB,
|
||||||
_register_accounts,
|
_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
|
# only its selected missing/disabled account. Existing accounts remain
|
||||||
# present even if their later network scan fails, so transient IMAP
|
# present even if their later network scan fails, so transient IMAP
|
||||||
# failure never erases their last known state.
|
# 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_account_generations = _email_urgency_account_generations(
|
||||||
base_state
|
base_state
|
||||||
)
|
)
|
||||||
|
|
@ -2403,13 +2411,6 @@ async def action_check_email_urgency(owner: str, **kwargs) -> Tuple[str, bool]:
|
||||||
STATE_LOCK_DB,
|
STATE_LOCK_DB,
|
||||||
_retire_accounts,
|
_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:
|
if not accounts:
|
||||||
raise TaskNoop("no email accounts configured")
|
raise TaskNoop("no email accounts configured")
|
||||||
|
|
||||||
|
|
|
||||||
|
|
@ -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
|
@pytest.mark.asyncio
|
||||||
async def test_payload_empty_tombstone_fences_reenable_redelete_and_can_recover(
|
async def test_payload_empty_tombstone_fences_reenable_redelete_and_can_recover(
|
||||||
monkeypatch,
|
monkeypatch,
|
||||||
|
|
|
||||||
Loading…
Add table
Reference in a new issue