WIP: perf(yahoo): respect Yahoo/AOL IMAP limits & avoid the 1-hour auth lockout #472

Draft
JMR-dev wants to merge 5 commits from feat-362-yahoo-imap-limits into main
35 changed files with 2510 additions and 219 deletions
Showing only changes of commit b1489e35da - Show all commits
+11 -3
View File
@@ -365,10 +365,18 @@ def main() -> int:
if not booted:
raise RuntimeError("API 37 preview emulator failed to boot after 2 attempts.")
# 4. Dismiss the keyguard, then run the instrumented/E2E suite against the booted emulator.
subprocess.run(cmd(adb, "shell", "input", "keyevent", "82"), check=False)
# 4. Force the emulator to grant the app window focus, then GATE on it (the SAME shared
# helper CI's e2e / e2e-preview jobs invoke, issue #468), before running the suite: wake
# the display, dismiss + disable the keyguard, keep the screen on, disable animations, and
# wait for a focused window. Replaces the lone `input keyevent 82`. Best-effort: fall back
# to that legacy nudge if the shared helper is somehow missing.
repo_root = Path(__file__).resolve().parents[3]
focus_gate = repo_root / ".github" / "scripts" / "emulator_focus_gate.py"
if focus_gate.is_file():
subprocess.run([sys.executable, str(focus_gate), "--adb", adb], check=False)
else:
subprocess.run(cmd(adb, "shell", "input", "keyevent", "82"), check=False)
gradlew = repo_root / ("gradlew.bat" if IS_WINDOWS else "gradlew")
print(f"Running :app:connectedDebugAndroidTest against {AVD_NAME}...")
test_exit = subprocess.run(
+12 -5
View File
@@ -324,10 +324,17 @@ def wait_for_boot(adb: str, proc: subprocess.Popen, timeout: int) -> bool:
return False
def dismiss_keyguard(adb: str) -> None:
# Dismiss the keyguard (mirrors CI + api37_e2e.py). Best-effort: a cold -wipe-data boot
# rarely needs it, and the input service can lose a race right after boot.
_run_quiet(cmd(adb, "-s", SERIAL, "shell", "input", "keyevent", "82"))
def prepare_focus(adb: str) -> None:
# Force the emulator to grant the app window focus BEFORE the suite, then gate on it -- the
# SAME shared helper CI's e2e / e2e-preview jobs invoke (issue #468), so local preflight
# exercises the identical fix. It wakes the display, dismisses + disables the keyguard, keeps
# the screen on, disables animations, and waits for a focused window. Best-effort: fall back
# to the legacy `input keyevent 82` nudge if the shared helper is somehow missing.
gate = REPO_ROOT / ".github" / "scripts" / "emulator_focus_gate.py"
if gate.is_file():
subprocess.run([sys.executable, str(gate), "--adb", adb, "--serial", SERIAL], check=False)
else:
_run_quiet(cmd(adb, "-s", SERIAL, "shell", "input", "keyevent", "82"))
def run_tests(test_classes: str) -> int:
@@ -420,7 +427,7 @@ def main() -> int:
proc = start_emulator(emulator)
if wait_for_boot(adb, proc, BOOT_TIMEOUT):
print("Emulator booted.")
dismiss_keyguard(adb)
prepare_focus(adb)
_STATE.test_exit = run_tests(args.test_classes)
else:
warn(f"Emulator did not reach sys.boot_completed within {BOOT_TIMEOUT}s.")
+277
View File
@@ -0,0 +1,277 @@
#!/usr/bin/env python3
# SPDX-License-Identifier: GPL-3.0-or-later
"""emulator_focus_gate.py -- make a booted emulator reliably grant the app window focus
BEFORE an instrumented UI suite runs, then GATE on that state (issue #468).
WHY THIS EXISTS (issue #468 -- the environmental app-window-focus flake)
------------------------------------------------------------------------
Intermittently, on the CI emulator the launched activity window has
``has-window-focus=false`` for the WHOLE instrumented run, so Espresso's ``RootViewPicker``
(used by ``onView(...).check()``, ``Intents.intended()``, ``Espresso.pressBack()`` and
focus-dependent clipboard reads) waits 10s for a focused root and times out --
``RootViewWithoutFocusException``. It fails EVERY window-focus-dependent test at once while
the ~280 pure-Compose semantics tests (which do not need window focus) pass. Root-cause
evidence from a failing ``E2E (35)`` leg (PR #470, run 28985259521): across the entire
captured logcat ``has-window-focus=true`` appears ZERO times and both the first attempt and
the once-retry fail identically -- i.e. the window NEVER gains focus for the session, a
persistent environmental state, not a per-test transient.
The prior mitigation was a single fire-and-forget ``adb shell input keyevent 82`` (MENU)
right after ``sys.boot_completed=1``. On modern Android (API 30+) MENU does NOT reliably
dismiss the keyguard, and when it is delivered before SystemUI/keyguard finishes coming up it
is simply dropped ("no focused window"). The insecure keyguard / non-interactive display then
persists and no app window ever takes focus -- hence the intermittent, whole-leg flake.
WHAT THIS DOES
--------------
A single shared mechanism invoked identically by every E2E job (the ``e2e`` API 29-36 matrix
AND the ``e2e-preview`` API 37 job in ``.github/workflows/ci.yml``) and by the local preflight
runners (``local_instrumented.py`` / ``api37_e2e.py``), so the fix cannot drift between them:
1. PREPARE the device so an app window CAN take focus, and keep it that way for the whole
run (all best-effort; a missing service right after boot must never abort the leg):
* ``input keyevent WAKEUP`` (224) -- force the display INTERACTIVE (never toggles it
off the way POWER would).
* ``wm dismiss-keyguard`` -- dismiss the (insecure) keyguard now.
* ``locksettings set-disabled true`` -- disable the lock screen for the session so it
cannot re-curtain the app window mid-run.
* ``svc power stayon true`` + a max ``screen_off_timeout`` -- never sleep during the run.
* ``input keyevent 82`` (MENU) -- legacy nudge, kept harmless for parity with #454.
* zero the three animation scales -- deterministic UI tests (this also gives the
``e2e-preview`` job the animation-disable the matrix already had -- uniformly).
2. GATE: poll ``dumpsys power`` + ``dumpsys window`` until the device is interactive
(``mWakefulness=Awake``) AND a real window holds input focus (``mCurrentFocus`` is a
``Window{...}``, not ``null``) -- i.e. the exact precondition ``RootViewPicker`` needs --
re-issuing the wake / dismiss-keyguard nudges each iteration so a lost race self-heals.
The gate is SOFT: it waits up to ``--timeout`` seconds and then proceeds regardless, printing a
GitHub ``::warning::`` annotation and the final device state if it never confirmed focus (the
determinism comes from the PREPARE actions + the wait; a parsing quirk on some API level must
not convert an otherwise-fine leg into a hard failure -- the real tests remain the arbiter).
It always prints the final ``mWakefulness`` / ``mCurrentFocus`` / keyguard state so a genuine
environmental failure is diagnosable from the step log without downloading artifacts.
Pure standard library, cross-platform (Windows / Linux / macOS): ``adb`` is invoked via
subprocess. The readiness parser (``evaluate_readiness``) is a pure function, unit-tested by
``test_emulator_focus_gate.py`` (run by the ``traffic-control-tests`` CI job).
"""
from __future__ import annotations
import argparse
import re
import shutil
import subprocess
import sys
import time
from typing import NamedTuple
# WAKEUP (not POWER): guarantees the display ends up INTERACTIVE. POWER (26) toggles, so it
# would turn an already-on display OFF. MENU (82) is kept only as a legacy parity nudge.
KEYCODE_WAKEUP = "224"
KEYCODE_MENU = "82"
# Max int -- effectively "never" auto-sleep the screen during the suite.
SCREEN_OFF_TIMEOUT_MS = "2147483647"
DEFAULT_TIMEOUT_S = 90
POLL_INTERVAL_S = 2
class Readiness(NamedTuple):
"""Outcome of parsing ``dumpsys power`` + ``dumpsys window`` for focus readiness."""
ready: bool
awake: bool
focus_state: str # 'focused' | 'unfocused' | 'unknown'
focus_value: str # the mCurrentFocus / mFocusedWindow token, or ''
keyguard_state: str # 'showing' | 'not_showing' | 'unknown'
@property
def summary(self) -> str:
return (
f"awake={self.awake} focus={self.focus_state}"
f"({self.focus_value or '-'}) keyguard={self.keyguard_state}"
)
def _is_awake(power_out: str) -> bool:
"""True if ``dumpsys power`` reports an INTERACTIVE display. ``mWakefulness=Awake`` is the
stable signal across API 29-37; ``Display Power: state=ON`` / ``mInteractive=true`` are
accepted as fallbacks for dump-format drift."""
return bool(
re.search(r"mWakefulness=Awake\b", power_out)
or re.search(r"Display Power:\s*state=ON\b", power_out)
or re.search(r"mInteractive=true\b", power_out)
)
def _focus(window_out: str) -> tuple[str, str]:
"""Classify the current input focus from ``dumpsys window``.
Returns ``(state, value)`` where state is 'focused' (a non-null ``Window{...}`` holds
focus -- what RootViewPicker needs), 'unfocused' (focus is explicitly ``null`` -- asleep /
keyguard-curtained / no focusable window), or 'unknown' (the field is absent on this dump
format). ``mCurrentFocus`` is preferred; ``mFocusedWindow`` is the fallback field name."""
tokens = re.findall(r"mCurrentFocus=(\S+)", window_out)
if not tokens:
tokens = re.findall(r"mFocusedWindow=(\S+)", window_out)
if not tokens:
return ("unknown", "")
non_null = [t for t in tokens if t != "null"]
if non_null:
return ("focused", non_null[0])
return ("unfocused", "null")
def _keyguard(window_out: str) -> str:
"""Best-effort keyguard state from ``dumpsys window``: 'showing' / 'not_showing' /
'unknown'. Informational for the summary, plus a fallback readiness signal when the focus
field is absent. Field names vary by API level, so several are accepted."""
match = re.search(
r"(?:mShowingLockscreen|mDreamingLockscreen|isKeyguardShowing|"
r"mKeyguardShowing|keyguardShowing|mKeyguardOccluded)=(true|false)",
window_out,
)
if not match:
return "unknown"
return "showing" if match.group(1) == "true" else "not_showing"
def evaluate_readiness(power_out: str, window_out: str) -> Readiness:
"""Pure decision core (unit-tested). The device is READY for a focus-dependent UI suite
when it is interactive AND a real window holds input focus. When the focus field is absent
on a given dump format, fall back to "interactive AND keyguard explicitly not showing" so a
format quirk cannot hang the gate forever."""
awake = _is_awake(power_out)
focus_state, focus_value = _focus(window_out)
keyguard_state = _keyguard(window_out)
ready = awake and (
focus_state == "focused"
or (focus_state == "unknown" and keyguard_state == "not_showing")
)
return Readiness(ready, awake, focus_state, focus_value, keyguard_state)
def _adb_base(adb: str, serial: str | None) -> list[str]:
return [adb, "-s", serial] if serial else [adb]
def _adb_quiet(adb: str, serial: str | None, *args: str) -> None:
"""Run an ``adb`` command, swallowing output and any error -- every prepare nudge is
best-effort (a service can lose a race right after boot; a missing tool must not abort)."""
try:
subprocess.run(
_adb_base(adb, serial) + list(args),
stdout=subprocess.DEVNULL,
stderr=subprocess.DEVNULL,
check=False,
timeout=30,
)
except (OSError, subprocess.SubprocessError):
pass
def _adb_capture(adb: str, serial: str | None, *args: str) -> str:
try:
return (
subprocess.run(
_adb_base(adb, serial) + list(args),
capture_output=True,
text=True,
check=False,
timeout=30,
).stdout
or ""
)
except (OSError, subprocess.SubprocessError):
return ""
def nudge_focus(adb: str, serial: str | None) -> None:
"""Wake the display + dismiss the keyguard. Cheap and idempotent, so it is re-issued every
poll iteration to self-heal a nudge that lost the post-boot race with SystemUI/keyguard."""
_adb_quiet(adb, serial, "shell", "input", "keyevent", KEYCODE_WAKEUP)
_adb_quiet(adb, serial, "shell", "wm", "dismiss-keyguard")
def prepare_device(adb: str, serial: str | None) -> None:
"""One-time device preparation: disable the lock screen for the session, keep the screen on
for the whole run, zero the animation scales for deterministic UI tests, and issue the first
wake / dismiss-keyguard nudge. All best-effort."""
print("focus-gate: preparing device (wake + dismiss-keyguard + stay-awake + no-animations)")
nudge_focus(adb, serial)
_adb_quiet(adb, serial, "shell", "input", "keyevent", KEYCODE_MENU) # legacy #454 parity
_adb_quiet(adb, serial, "shell", "locksettings", "set-disabled", "true")
_adb_quiet(adb, serial, "shell", "svc", "power", "stayon", "true")
_adb_quiet(adb, serial, "shell", "settings", "put", "system",
"screen_off_timeout", SCREEN_OFF_TIMEOUT_MS)
for scale in ("window_animation_scale", "transition_animation_scale",
"animator_duration_scale"):
_adb_quiet(adb, serial, "shell", "settings", "put", "global", scale, "0.0")
def probe(adb: str, serial: str | None) -> Readiness:
power_out = _adb_capture(adb, serial, "shell", "dumpsys", "power")
window_out = _adb_capture(adb, serial, "shell", "dumpsys", "window")
return evaluate_readiness(power_out, window_out)
def wait_for_focus(adb: str, serial: str | None, timeout: int, label: str) -> Readiness:
"""Prepare the device, then poll (re-nudging each iteration) until it is interactive with a
focused window, or ``timeout`` seconds elapse. Returns the final Readiness (SOFT gate: the
caller proceeds regardless -- see the module docstring)."""
tag = f" [{label}]" if label else ""
prepare_device(adb, serial)
deadline = time.monotonic() + timeout
last = probe(adb, serial)
attempt = 0
while True:
if last.ready:
elapsed = timeout - max(0, int(deadline - time.monotonic()))
print(f"focus-gate{tag}: READY after ~{elapsed}s -- {last.summary}")
return last
if time.monotonic() >= deadline:
print(f"::warning::focus-gate{tag}: window focus NOT confirmed within {timeout}s "
f"-- proceeding anyway -- {last.summary}")
return last
attempt += 1
if attempt % 5 == 0:
print(f"focus-gate{tag}: waiting for window focus -- {last.summary}")
nudge_focus(adb, serial)
time.sleep(POLL_INTERVAL_S)
last = probe(adb, serial)
def main(argv: list[str] | None = None) -> int:
parser = argparse.ArgumentParser(
prog="emulator_focus_gate.py",
description=(
"Force a booted emulator to grant the app window focus (wake + dismiss-keyguard + "
"stay-awake + no-animations) and gate on that state before an instrumented UI "
"suite runs. Shared by CI's e2e / e2e-preview jobs and the local preflight runners "
"(issue #468)."
),
)
parser.add_argument("--serial", default=None,
help="adb device serial (default: the single attached device).")
parser.add_argument("--adb", default=None,
help="Path to adb (default: resolve from PATH). For callers that resolve "
"adb from the SDK rather than PATH (e.g. api37_e2e.py).")
parser.add_argument("--timeout", type=int, default=DEFAULT_TIMEOUT_S,
help=f"Max seconds to wait for window focus (default {DEFAULT_TIMEOUT_S}).")
parser.add_argument("--label", default="",
help="Label for log lines (e.g. an API level), for multi-leg runs.")
args = parser.parse_args(argv)
adb = args.adb or shutil.which("adb")
if not adb:
# Non-fatal by contract: never turn a missing-tool hiccup into a red leg. The suite that
# follows will surface a genuinely broken device.
print("::warning::focus-gate: adb not on PATH -- skipping focus preparation/gate")
return 0
wait_for_focus(adb, args.serial, args.timeout, args.label)
return 0
if __name__ == "__main__":
sys.exit(main())
+111
View File
@@ -0,0 +1,111 @@
# SPDX-License-Identifier: GPL-3.0-or-later
"""Unit tests for the pure readiness parser of emulator_focus_gate.py (no adb, no emulator).
Covers the decision core that decides whether a booted emulator is ready for a focus-dependent
instrumented UI suite (issue #468): interactive (``mWakefulness=Awake``) AND a real window holds
input focus (``mCurrentFocus`` is a non-null ``Window{...}``). The window-focus flake this guards
against is exactly the "awake but mCurrentFocus=null" state, so that case must read NOT ready."""
from __future__ import annotations
import unittest
import emulator_focus_gate as gate
# A ``dumpsys power`` where the display is interactive vs. asleep.
POWER_AWAKE = "Power Manager State:\n mWakefulness=Awake\n mWakefulnessChanging=false\n"
POWER_ASLEEP = "Power Manager State:\n mWakefulness=Asleep\n mWakefulnessChanging=false\n"
# ``dumpsys window`` with a focused app window (the healthy state RootViewPicker needs)...
WINDOW_FOCUSED = (
" mCurrentFocus=Window{23e192a u0 org.libremail.app/org.libremail.MainActivity}\n"
" mFocusedApp=ActivityRecord{a1 u0 org.libremail.app/.MainActivity t9}\n"
" mDreamingLockscreen=false\n"
)
# ...and the flake state: interactive-parse aside, NO window holds focus.
WINDOW_NO_FOCUS = " mCurrentFocus=null\n mFocusedApp=null\n mDreamingLockscreen=true\n"
class AwakeParsingTests(unittest.TestCase):
def test_mwakefulness_awake(self) -> None:
self.assertTrue(gate._is_awake(POWER_AWAKE))
def test_mwakefulness_asleep(self) -> None:
self.assertFalse(gate._is_awake(POWER_ASLEEP))
def test_display_power_state_on_fallback(self) -> None:
self.assertTrue(gate._is_awake("Display Power: state=ON"))
def test_minteractive_fallback(self) -> None:
self.assertTrue(gate._is_awake("mInteractive=true"))
def test_empty_is_not_awake(self) -> None:
self.assertFalse(gate._is_awake(""))
class FocusParsingTests(unittest.TestCase):
def test_non_null_current_focus(self) -> None:
state, value = gate._focus(WINDOW_FOCUSED)
self.assertEqual(state, "focused")
self.assertTrue(value.startswith("Window{"))
def test_null_current_focus(self) -> None:
self.assertEqual(gate._focus(WINDOW_NO_FOCUS), ("unfocused", "null"))
def test_focused_window_fallback_field(self) -> None:
state, value = gate._focus("mFocusedWindow=Window{deadbeef u0 launcher}\n")
self.assertEqual(state, "focused")
self.assertEqual(value, "Window{deadbeef")
def test_absent_focus_field_is_unknown(self) -> None:
self.assertEqual(gate._focus("no focus fields here"), ("unknown", ""))
class KeyguardParsingTests(unittest.TestCase):
def test_showing(self) -> None:
self.assertEqual(gate._keyguard("mDreamingLockscreen=true"), "showing")
def test_not_showing(self) -> None:
self.assertEqual(gate._keyguard("isKeyguardShowing=false"), "not_showing")
def test_unknown(self) -> None:
self.assertEqual(gate._keyguard("nothing relevant"), "unknown")
class EvaluateReadinessTests(unittest.TestCase):
def test_awake_and_focused_is_ready(self) -> None:
result = gate.evaluate_readiness(POWER_AWAKE, WINDOW_FOCUSED)
self.assertTrue(result.ready)
self.assertTrue(result.awake)
self.assertEqual(result.focus_state, "focused")
def test_the_flake_awake_but_no_focus_is_not_ready(self) -> None:
# The exact issue #468 signature: display parses/awake but no window has focus.
result = gate.evaluate_readiness(POWER_AWAKE, WINDOW_NO_FOCUS)
self.assertFalse(result.ready)
def test_asleep_even_with_focus_is_not_ready(self) -> None:
result = gate.evaluate_readiness(POWER_ASLEEP, WINDOW_FOCUSED)
self.assertFalse(result.ready)
def test_unknown_focus_but_awake_and_keyguard_gone_is_ready(self) -> None:
# Fallback so a dump format without mCurrentFocus can't hang the gate forever.
result = gate.evaluate_readiness(POWER_AWAKE, "mDreamingLockscreen=false")
self.assertTrue(result.ready)
def test_unknown_focus_and_keyguard_showing_is_not_ready(self) -> None:
result = gate.evaluate_readiness(POWER_AWAKE, "mDreamingLockscreen=true")
self.assertFalse(result.ready)
def test_unknown_focus_and_keyguard_unknown_is_not_ready(self) -> None:
result = gate.evaluate_readiness(POWER_AWAKE, "")
self.assertFalse(result.ready)
def test_summary_is_human_readable(self) -> None:
summary = gate.evaluate_readiness(POWER_AWAKE, WINDOW_FOCUSED).summary
self.assertIn("awake=True", summary)
self.assertIn("focus=focused", summary)
if __name__ == "__main__":
unittest.main()
+21 -7
View File
@@ -632,12 +632,20 @@ jobs:
for attempt in 1 2; do boot_emulator "$attempt" && { booted=1; break; }; done
[ "$booted" = "1" ] || { echo "::error::API ${{ matrix.api-level }} emulator failed to boot after 2 attempts"; exit 1; }
# NON-FATAL unlock (the boot-race fix) + disable animations for deterministic UI tests
# (parity with the replaced android-emulator-runner `disable-animations: true`).
adb shell input keyevent 82 || true
adb shell settings put global window_animation_scale 0.0 || true
adb shell settings put global transition_animation_scale 0.0 || true
adb shell settings put global animator_duration_scale 0.0 || true
# Focus/readiness gate (issue #468): the durable fix for the intermittent app-window-
# focus flake. Occasionally the launched activity window has has-window-focus=false for
# the WHOLE run, so Espresso's RootViewPicker (onView().check(), Intents.intended(),
# pressBack(), focus-dependent clipboard) times out after 10s and fails EVERY
# focus-dependent test at once while the ~280 pure-Compose tests pass (evidence: a
# failing E2E leg's logcat had has-window-focus=true ZERO times, both attempt + retry).
# The single `input keyevent 82` here was too weak (MENU no longer dismisses the modern
# keyguard, and races SystemUI coming up). This shared helper WAKES the display, dismisses
# + disables the keyguard, keeps the screen on, disables animations (parity with the
# replaced `disable-animations: true`), then WAITS until a real window holds input focus
# before the suite runs. It is the IDENTICAL mechanism the e2e-preview job and the local
# preflight runners invoke, so the fix cannot drift between jobs. Non-fatal (`|| true`),
# preserving #454's guarantee that the unlock never aborts the boot.
python3 .github/scripts/emulator_focus_gate.py --label "api${{ matrix.api-level }}" || true
# WEDGE (hang) smoking-gun capture (#404, restored to the matrix by #421). On the wrapper
# `timeout` below (exit 124), grab the smoking gun WHILE this hand-provisioned emulator is
@@ -968,7 +976,13 @@ jobs:
for attempt in 1 2; do boot_emulator "$attempt" && { booted=1; break; }; done
[ "$booted" = "1" ] || { echo "::error::API 37 preview emulator failed to boot after 2 attempts"; exit 1; }
adb shell input keyevent 82 || true
# Focus/readiness gate (issue #468) -- the IDENTICAL shared mechanism the `e2e` matrix job
# (and the local preflight runners) invoke, so the fix can't drift: wake the display,
# dismiss + disable the keyguard, keep the screen on, disable animations, then WAIT until a
# real window holds input focus before the suite runs. This replaces the lone `input
# keyevent 82` and, applied UNIFORMLY, also gives this preview job the animation-disable
# the matrix already had. Non-fatal (`|| true`) -- the unlock must never abort the boot.
python3 .github/scripts/emulator_focus_gate.py --label "api37-shard${{ matrix.shard }}" || true
# WEDGE (hang) smoking-gun capture (#404). On the wrapper `timeout` below (exit 124), grab
# the smoking gun WHILE this hand-provisioned emulator is still alive (it stays up until the
@@ -0,0 +1,71 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import androidx.test.ext.junit.runners.AndroidJUnit4
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.async
import kotlinx.coroutines.awaitAll
import kotlinx.coroutines.runBlocking
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertTrue
import org.junit.Test
import org.junit.runner.RunWith
import org.libremail.domain.model.Account
import org.libremail.domain.model.MailProvider
/**
* On-device proof of issue #361's Gmail bandwidth pacing across the CI API matrix. Production feeds
* [GmailBandwidthTracker] from `MailRepositoryImpl.prefetchMessage`, which can race across
* concurrently-syncing accounts under the real dispatcher, so this proves the tracker's
* concurrent-map-backed accounting holds up under genuine concurrent updates on real threads rather
* than coroutines-test virtual time (the JVM [GmailBandwidthTrackerTest] covers the day-rollover and
* threshold-crossing logic in detail). Also proves [GmailSyncLimits.appliesTo] resolves the real
* [MailProvider] presets identically to production. Deliberately mock-free — no `mockk`, no framework
* `Context` — the tracker's only collaborator is the real wall clock (mirrors
* `BackfillPacerInstrumentedTest`'s mock-free idiom).
*/
@RunWith(AndroidJUnit4::class)
class GmailBandwidthTrackerInstrumentedTest {
@Test
fun concurrentDownloadsForOneAccountAllLandWithoutLosingAnUpdate() = runBlocking {
val tracker = GmailBandwidthTracker()
val perTask = 1_000L
val jobs = (1..CONCURRENT_TASKS).map {
async(Dispatchers.Default) { tracker.recordDownload("acct", perTask) }
}
jobs.awaitAll()
assertEquals(CONCURRENT_TASKS * perTask, tracker.bytesDownloadedToday("acct"))
}
@Test
fun accountsStayIsolatedUnderConcurrentRecording() = runBlocking {
val tracker = GmailBandwidthTracker()
val heavy = async(Dispatchers.Default) {
tracker.recordDownload("heavy", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
}
val light = async(Dispatchers.Default) { tracker.recordDownload("light", 1L) }
heavy.await()
light.await()
assertTrue(tracker.isOverDailyBudget("heavy"))
assertFalse(tracker.isOverDailyBudget("light"))
}
@Test
fun appliesToResolvesTheRealGmailPresetOnDevice() {
val gmail = MailProvider.GMAIL.createAccount("user@gmail.com")
val outlook = Account.outlook("user@outlook.com")
assertTrue(GmailSyncLimits.appliesTo(gmail))
assertFalse(GmailSyncLimits.appliesTo(outlook))
}
private companion object {
const val CONCURRENT_TASKS = 50
}
}
@@ -0,0 +1,95 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import androidx.test.ext.junit.runners.AndroidJUnit4
import kotlinx.coroutines.runBlocking
import org.json.JSONArray
import org.json.JSONObject
import org.junit.Assert.assertEquals
import org.junit.Assert.assertFalse
import org.junit.Assert.assertTrue
import org.junit.Test
import org.junit.runner.RunWith
import org.libremail.data.sync.AccountThrottleGate
/**
* Exercises the issue #364 Graph throttling toolkit on a real device/emulator (the CI E2E matrix is
* authoritative for this): the `$batch` call-volume reduction, the chunked upload session, and the 429
* composition with the shared #360 [AccountThrottleGate] all run under the Android runtime and its real
* `android.util.Log` (no JVM stub to mock). Deliberately avoids the retry-delay path so it needs no
* coroutines-test virtual clock (unavailable on the androidTest classpath) — the delay/backoff schedule
* is proven under virtual time in the JVM suite (GraphThrottleTest).
*/
@RunWith(AndroidJUnit4::class)
class GraphThrottleInstrumentedTest {
private val accountId = "outlook:user@example.org"
/** A no-network [GraphHttpClient] that records requests and returns scripted responses. */
private class FakeClient(private val responder: (GraphRequest, Int) -> GraphResponse) : GraphHttpClient() {
val requests = mutableListOf<GraphRequest>()
override suspend fun execute(request: GraphRequest): GraphResponse {
val index = requests.size
requests += request
return responder(request, index)
}
}
@Test
fun batch_collapses_many_operations_into_few_calls() = runBlocking {
val client = FakeClient { request, _ ->
val requested = JSONObject(String(request.body!!, Charsets.UTF_8)).getJSONArray("requests")
val responses = JSONArray()
for (i in 0 until requested.length()) {
responses.put(JSONObject().put("id", requested.getJSONObject(i).getString("id")).put("status", 200))
}
GraphResponse(status = 200, body = JSONObject().put("responses", responses).toString())
}
val requests = (1..25).map { GraphSubRequest(id = it.toString(), method = "GET", url = "/me/messages/$it") }
val responses = GraphBatch(GraphThrottle(AccountThrottleGate())).execute(accountId, client, requests)
assertEquals(2, client.requests.size)
assertEquals(25, responses.size)
}
@Test
fun upload_session_uploads_over_threshold_content_in_chunks() = runBlocking {
val chunk = 320 * 1024
val client = FakeClient { request, _ ->
if (request.method == "POST") {
GraphResponse(200, JSONObject().put("uploadUrl", "https://upload.example/1").toString())
} else {
GraphResponse(202, "")
}
}
val total = 4 * 1024 * 1024 + 500
val content = ByteArray(total) { (it % 251).toByte() }
val session = GraphUploadSession(GraphThrottle(AccountThrottleGate()))
assertTrue(session.requiresUploadSession(total.toLong()))
session.upload(accountId, client, "https://graph/createUploadSession", JSONObject(), content, chunk)
val puts = client.requests.drop(1)
assertEquals((total + chunk - 1) / chunk, puts.size)
assertTrue(puts.all { it.method == "PUT" })
}
@Test
fun a_429_records_and_a_success_clears_the_shared_gate() = runBlocking {
val gate = AccountThrottleGate()
val throttle = GraphThrottle(gate)
// maxRetries = 0 → no backoff delay; the 429 is recorded and returned.
val client429 = FakeClient { _, _ -> GraphResponse(429, "") }
val throttled = throttle.execute(accountId, client429, sendMail(), maxRetries = 0)
assertEquals(429, throttled.status)
assertTrue("an unrecovered 429 backs the account off", gate.isThrottled(accountId))
val ok = throttle.execute(accountId, FakeClient { _, _ -> GraphResponse(202, "") }, sendMail())
assertEquals(202, ok.status)
assertFalse("a later success clears the backoff", gate.isThrottled(accountId))
}
private fun sendMail() = GraphRequest("POST", "https://graph.microsoft.com/v1.0/me/sendMail")
}
@@ -1,10 +1,11 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.ui.accountsetup
import android.app.Activity
import android.app.Instrumentation
import android.content.Intent
import android.net.Uri
import androidx.activity.ComponentActivity
import androidx.compose.runtime.CompositionLocalProvider
import androidx.compose.ui.platform.LocalUriHandler
import androidx.compose.ui.platform.UriHandler
import androidx.compose.ui.test.assertIsDisplayed
import androidx.compose.ui.test.junit4.v2.createAndroidComposeRule
import androidx.compose.ui.test.onAllNodesWithText
@@ -13,14 +14,9 @@ import androidx.compose.ui.test.performClick
import androidx.compose.ui.test.performScrollTo
import androidx.compose.ui.test.performTextInput
import androidx.lifecycle.SavedStateHandle
import androidx.test.espresso.intent.Intents
import androidx.test.espresso.intent.matcher.IntentMatchers.hasAction
import androidx.test.espresso.intent.matcher.IntentMatchers.hasData
import androidx.test.espresso.intent.matcher.UriMatchers.hasHost
import androidx.test.ext.junit.runners.AndroidJUnit4
import jakarta.mail.AuthenticationFailedException
import org.hamcrest.CoreMatchers.allOf
import org.hamcrest.CoreMatchers.equalTo
import org.junit.Assert.assertEquals
import org.junit.Rule
import org.junit.Test
import org.junit.runner.RunWith
@@ -35,8 +31,17 @@ import org.libremail.ui.theme.LibreMailTheme
* [AppPasswordSetupScreen] + [AppPasswordViewModel] over a [FakeAccountRepository] for the preset
* Gmail vendor: the provider-specific chrome renders, entering an email + app password and tapping
* "Test & add" persists through the repository and reports the new account id, and tapping the
* "create an app password" help link fires the browser intent. That launch is asserted with
* Espresso-Intents (mirroring `AccountPickerScreenTest`'s Outlook test), so no real browser opens.
* "create an app password" help link opens the provider's help page.
*
* The outbound help links are verified by injecting a recording [UriHandler] for [LocalUriHandler]
* and asserting the URL the screen asked to open — deliberately NOT via Espresso-Intents. The two
* approaches verify the same behaviour, but `Intents.intended(...)` runs an `onView(isRoot())` view
* assertion whose `RootViewPicker` waits up to 10s for a window-focused root; on the CI emulator the
* activity window intermittently reports `has-window-focus=false`, so that assertion flakes with
* `RootViewWithoutFocusException` (an infra flake that fails every `intended()`-based E2E test on the
* affected leg and forces a costly 9-min retry). Driving the link through a fake [UriHandler] keeps
* the whole test on Compose interactions, which do not depend on window focus, so it is deterministic
* — while still asserting the exact provider page the tap opens.
*/
@RunWith(AndroidJUnit4::class)
class AppPasswordSetupScreenTest {
@@ -46,6 +51,9 @@ class AppPasswordSetupScreenTest {
private val provider = MailProvider.GMAIL
// Captures the URL the screen hands to LocalUriHandler instead of launching a real browser.
private val uriHandler = RecordingUriHandler()
private fun string(resId: Int, vararg args: Any) = composeTestRule.activity.getString(resId, *args)
private fun setContent(
@@ -57,8 +65,10 @@ class AppPasswordSetupScreenTest {
repository,
)
composeTestRule.setContent {
LibreMailTheme(darkTheme = false, dynamicColor = false) {
AppPasswordSetupScreen(onBack = {}, onAccountAdded = onAccountAdded, viewModel = viewModel)
CompositionLocalProvider(LocalUriHandler provides uriHandler) {
LibreMailTheme(darkTheme = false, dynamicColor = false) {
AppPasswordSetupScreen(onBack = {}, onAccountAdded = onAccountAdded, viewModel = viewModel)
}
}
}
}
@@ -90,34 +100,29 @@ class AppPasswordSetupScreenTest {
}
/**
* Tapping the "create an app password" link opens the provider's help page via
* [androidx.compose.ui.platform.UriHandler], which starts an `ACTION_VIEW` intent. Stubbing that
* intent both proves the tap launched it and stops a real browser from opening on the device.
* Tapping the "create an app password" link opens the provider's help page via [LocalUriHandler].
* Asserting the URL captured by [RecordingUriHandler] proves the tap requested the right page
* without launching a real browser (and without the window-focus-dependent Espresso-Intents
* assertion that flakes on CI — see the class comment).
*/
@Test
fun tappingCreateAppPasswordPage_launchesBrowserIntentToHelpUrl() {
setContent()
Intents.init()
try {
Intents.intending(hasAction(Intent.ACTION_VIEW))
.respondWith(Instrumentation.ActivityResult(Activity.RESULT_CANCELED, null))
composeTestRule.onNodeWithText(string(R.string.app_password_open_page, provider.displayName))
.performScrollTo()
.performClick()
composeTestRule.onNodeWithText(string(R.string.app_password_open_page, provider.displayName))
.performScrollTo()
.performClick()
Intents.intended(allOf(hasAction(Intent.ACTION_VIEW), hasData(provider.appPasswordHelpUrl)))
} finally {
Intents.release()
}
assertEquals(provider.appPasswordHelpUrl, uriHandler.lastUri)
}
/**
* When the connection test fails specifically because IMAP is disabled (Gmail's "not enabled for
* IMAP use"), the screen surfaces the actionable "turn on IMAP" dialog instead of a generic error,
* and its help link opens the provider's enable-IMAP page (#390). Driving the failure through a
* [FakeAccountRepository] exercises the real classification + dialog wiring end to end on device.
* [FakeAccountRepository] exercises the real classification + dialog wiring end to end on device;
* the help link's target is verified through the injected [RecordingUriHandler] (see the class
* comment for why not Espresso-Intents).
*/
@Test
fun imapDisabledFailure_showsThePrompt_andHelpLinkOpensTheProviderPage() {
@@ -138,16 +143,20 @@ class AppPasswordSetupScreenTest {
}
composeTestRule.onNodeWithText(string(R.string.imap_disabled_message, provider.displayName)).assertIsDisplayed()
Intents.init()
try {
Intents.intending(hasAction(Intent.ACTION_VIEW))
.respondWith(Instrumentation.ActivityResult(Activity.RESULT_CANCELED, null))
composeTestRule.onNodeWithText(string(R.string.imap_disabled_help)).performClick()
composeTestRule.onNodeWithText(string(R.string.imap_disabled_help)).performClick()
// Mirrors the previous Espresso hasHost(...) check: the Gmail enable-IMAP page is on Google's
// support host. Verifying the exact host keeps the assertion strength without any focus wait.
assertEquals("support.google.com", Uri.parse(uriHandler.lastUri).host)
}
Intents.intended(allOf(hasAction(Intent.ACTION_VIEW), hasData(hasHost(equalTo("support.google.com")))))
} finally {
Intents.release()
/** A [UriHandler] that records the last opened URL instead of starting a real `ACTION_VIEW` intent. */
private class RecordingUriHandler : UriHandler {
var lastUri: String? = null
private set
override fun openUri(uri: String) {
lastUri = uri
}
}
}
@@ -4,23 +4,22 @@ package org.libremail.ui.onboarding
import android.app.Activity
import android.app.Instrumentation
import android.content.Context
import android.content.Intent
import android.net.Uri
import androidx.activity.ComponentActivity
import androidx.compose.runtime.CompositionLocalProvider
import androidx.compose.ui.platform.LocalUriHandler
import androidx.compose.ui.platform.UriHandler
import androidx.compose.ui.test.assertIsDisplayed
import androidx.compose.ui.test.junit4.v2.createAndroidComposeRule
import androidx.compose.ui.test.onNodeWithText
import androidx.compose.ui.test.performClick
import androidx.compose.ui.test.performScrollTo
import androidx.test.espresso.intent.Intents
import androidx.test.espresso.intent.matcher.IntentMatchers.hasAction
import androidx.test.espresso.intent.matcher.IntentMatchers.hasComponent
import androidx.test.espresso.intent.matcher.IntentMatchers.hasData
import androidx.test.espresso.intent.matcher.UriMatchers.hasHost
import androidx.test.ext.junit.runners.AndroidJUnit4
import androidx.test.platform.app.InstrumentationRegistry
import net.openid.appauth.AuthorizationManagementActivity
import org.hamcrest.CoreMatchers.allOf
import org.hamcrest.CoreMatchers.equalTo
import org.junit.Assert.assertEquals
import org.junit.Rule
import org.junit.Test
import org.junit.runner.RunWith
@@ -33,11 +32,18 @@ import org.libremail.ui.theme.LibreMailTheme
/**
* End-to-end UI test for the pre-auth Outlook IMAP-enablement notice (#411). Drives the real
* [OutlookImapNoticeScreen] + [AccountSetupViewModel] over a [FakeAccountRepository]: the IMAP
* question and both outbound links render, tapping a help link fires the browser `ACTION_VIEW`
* intent, and tapping the bottom "Sign in" button starts the existing Microsoft OAuth (AppAuth)
* flow. Both launches are asserted with Espresso-Intents (mirroring `AccountPickerScreenTest`), so no
* real browser ever opens; the interstitial → OAuth navigation in the full onboarding graph is
* covered by `OnboardingFlowTest`.
* question and both outbound links render, tapping the "How to enable IMAP" help link opens
* Microsoft's help article, and tapping the bottom "Sign in" button starts the existing Microsoft
* OAuth (AppAuth) flow.
*
* The help link is verified by injecting a recording [UriHandler] for [LocalUriHandler] and asserting
* the opened URL — not via Espresso-Intents, whose `intended(...)` runs an `onView(isRoot())`
* assertion that waits for a window-focused root and flakes with `RootViewWithoutFocusException` on
* the CI emulator (see `AppPasswordSetupScreenTest` for the full write-up). The "Sign in" launch has
* no [UriHandler] seam — AppAuth calls `startActivity` directly — so it stays on Espresso-Intents,
* matched by AppAuth's [AuthorizationManagementActivity] component and stubbed so no real browser
* opens; the interstitial → OAuth navigation in the full onboarding graph is covered by
* `OnboardingFlowTest`.
*/
@RunWith(AndroidJUnit4::class)
class OutlookImapNoticeScreenTest {
@@ -48,13 +54,18 @@ class OutlookImapNoticeScreenTest {
private val context: Context =
InstrumentationRegistry.getInstrumentation().targetContext.applicationContext
// Captures the URL the screen hands to LocalUriHandler instead of launching a real browser.
private val uriHandler = RecordingUriHandler()
private fun string(resId: Int) = composeTestRule.activity.getString(resId)
private fun setContent(onAccountAdded: (String) -> Unit = {}) {
val viewModel = AccountSetupViewModel(OutlookAuthManager(context), FakeAccountRepository())
composeTestRule.setContent {
LibreMailTheme(darkTheme = false, dynamicColor = false) {
OutlookImapNoticeScreen(onBack = {}, onAccountAdded = onAccountAdded, viewModel = viewModel)
CompositionLocalProvider(LocalUriHandler provides uriHandler) {
LibreMailTheme(darkTheme = false, dynamicColor = false) {
OutlookImapNoticeScreen(onBack = {}, onAccountAdded = onAccountAdded, viewModel = viewModel)
}
}
}
}
@@ -70,27 +81,20 @@ class OutlookImapNoticeScreenTest {
}
/**
* Tapping the "How to enable IMAP" link opens Microsoft's help article via
* [androidx.compose.ui.platform.UriHandler], which starts an `ACTION_VIEW` intent. Stubbing that
* intent both proves the tap launched it and stops a real browser from opening on the device.
* Tapping the "How to enable IMAP" link opens Microsoft's help article via [LocalUriHandler].
* Asserting the URL captured by [RecordingUriHandler] proves the tap requested the right page
* without launching a real browser (and without the window-focus-dependent Espresso-Intents
* assertion that flakes on CI — see the class comment).
*/
@Test
fun tappingImapHelpLink_opensTheMicrosoftArticle() {
setContent()
Intents.init()
try {
Intents.intending(hasAction(Intent.ACTION_VIEW))
.respondWith(Instrumentation.ActivityResult(Activity.RESULT_CANCELED, null))
composeTestRule.onNodeWithText(string(R.string.outlook_imap_help))
.performScrollTo()
.performClick()
composeTestRule.onNodeWithText(string(R.string.outlook_imap_help))
.performScrollTo()
.performClick()
Intents.intended(allOf(hasAction(Intent.ACTION_VIEW), hasData(hasHost(equalTo("support.microsoft.com")))))
} finally {
Intents.release()
}
assertEquals("support.microsoft.com", Uri.parse(uriHandler.lastUri).host)
}
/**
@@ -118,4 +122,14 @@ class OutlookImapNoticeScreenTest {
Intents.release()
}
}
/** A [UriHandler] that records the last opened URL instead of starting a real `ACTION_VIEW` intent. */
private class RecordingUriHandler : UriHandler {
var lastUri: String? = null
private set
override fun openUri(uri: String) {
lastUri = uri
}
}
}
@@ -39,6 +39,8 @@ import org.libremail.data.local.toOutgoingAttachments
import org.libremail.data.local.toOutgoingAttachmentsJson
import org.libremail.data.settings.AccountSettingsRepository
import org.libremail.data.settings.SignatureRepository
import org.libremail.data.sync.GmailBandwidthTracker
import org.libremail.data.sync.GmailSyncLimits
import org.libremail.data.sync.InteractiveImapGate
import org.libremail.data.sync.MailConnectionFactory
import org.libremail.data.sync.SendScheduler
@@ -81,6 +83,7 @@ class MailRepositoryImpl @Inject constructor(
private val signatureRepository: SignatureRepository,
private val attachmentUriGrants: AttachmentUriGrants,
private val interactiveGate: InteractiveImapGate,
private val bandwidthTracker: GmailBandwidthTracker,
) : MailRepository {
// Application-lifetime scope for fire-and-forget server pushes that must outlive the caller — e.g.
@@ -253,7 +256,7 @@ class MailRepositoryImpl @Inject constructor(
// fetch just omits that image, leaving a broken <img> rather than failing the open.
val file = runCatching {
ensureAttachmentFile(messageId, routing.accountId, routing.folder, row.partIndex, row.filename)
}.getOrNull() ?: return@mapNotNull null
}.getOrNull()?.file ?: return@mapNotNull null
InlineImage(contentId = row.contentId!!, mimeType = row.mimeType, bytes = file.readBytes())
}
}
@@ -275,11 +278,18 @@ class MailRepositoryImpl @Inject constructor(
routing.folder,
partIndex,
meta?.filename ?: "attachment",
)
).file
}
}
}
/**
* [ensureAttachmentFile]'s outcome: the cached [file] plus the bytes actually pulled over the
* network THIS call — `0` on a cache hit. [downloadedBytes] feeds Gmail's bandwidth accounting
* (issue #361, see [prefetchMessage]); a cache hit costs nothing so it must not be double-counted.
*/
private class AttachmentFetch(val file: File, val downloadedBytes: Long)
/**
* Returns the on-disk file for one attachment part, downloading and caching it on first use so it
* then opens instantly and offline. Takes the message's already-resolved account/folder so a batch
@@ -292,16 +302,16 @@ class MailRepositoryImpl @Inject constructor(
folder: String,
partIndex: Int,
filename: String,
): File {
): AttachmentFetch {
val target = attachmentFile(messageId, partIndex, filename)
// Reuse a previously downloaded (or pre-fetched) file so it opens instantly and offline.
if (target.exists() && target.length() > 0L) return target
if (target.exists() && target.length() > 0L) return AttachmentFetch(target, downloadedBytes = 0L)
val account = accountDao.getById(accountId)?.toDomain() ?: error("Account not found")
val params = connectionFactory.imapParamsFor(account)
val downloaded = imapClient.fetchAttachment(params, folder, uidOf(messageId), partIndex)
target.parentFile?.mkdirs()
target.outputStream().use { it.write(downloaded.bytes) }
return target
return AttachmentFetch(target, downloadedBytes = downloaded.bytes.size.toLong())
}
override suspend fun downloadedAttachmentParts(messageId: String): Set<Int> = withContext(Dispatchers.IO) {
@@ -317,12 +327,17 @@ class MailRepositoryImpl @Inject constructor(
override suspend fun prefetchMessage(messageId: String): Result<Unit> = runCatching {
val routing = messageDao.getRouting(messageId) ?: return@runCatching
val account = accountDao.getById(routing.accountId)?.toDomain() ?: return@runCatching
// Bytes actually pulled over the network this call (0 on an all-cache-hit prefetch), fed to
// Gmail's daily download-budget tracker below (issue #361) — the proactive pacing that composes
// with #360's reactive AccountThrottleGate and #356's BackfillPacer without modifying either.
var downloadedBytes = 0L
// Cache the body (peek, so prefetching never marks the message read) and its attachment metadata.
if (!routing.bodyFetched) {
val params = connectionFactory.imapParamsFor(account)
val content = imapClient.fetchBodyPeek(params, routing.folder, uidOf(messageId))
messageDao.updateBody(messageId, content.body, content.isHtml, Snippet.of(content.body, content.isHtml))
attachmentDao.replaceForMessage(messageId, content.attachments.map { it.toEntity(messageId) })
downloadedBytes += content.body.toByteArray(Charsets.UTF_8).size.toLong()
}
// Auto-download every attachment's bytes into the persistent per-part cache (skips ones present).
// This is BACKGROUND work driven by the backfill, so it goes straight to ensureAttachmentFile and
@@ -338,7 +353,13 @@ class MailRepositoryImpl @Inject constructor(
attachment.partIndex,
attachment.filename,
)
}
}.getOrNull()?.let { downloadedBytes += it.downloadedBytes }
}
// Gmail-specific bandwidth accounting (issue #361): only tracked for Gmail, since that is the
// only provider whose daily download budget is enforced today (see GmailSyncLimits.appliesTo /
// MailBackfiller.prefetchIfEnabled / MailSyncer.prefetchIfEnabled for the deferral this feeds).
if (downloadedBytes > 0L && GmailSyncLimits.appliesTo(account)) {
bandwidthTracker.recordDownload(account.id, downloadedBytes)
}
}
@@ -0,0 +1,88 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import org.libremail.reporting.AppLog
import org.libremail.reporting.accountLogRef
import java.util.concurrent.ConcurrentHashMap
import javax.inject.Inject
import javax.inject.Singleton
/**
* Per-account, per-day running total of bytes downloaded by background prefetch, tracked against
* Gmail's documented [GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES] (issue #361). The stateful
* counterpart to the pure [GmailSyncLimits]:
* [org.libremail.data.repository.MailRepositoryImpl.prefetchMessage] feeds it bytes actually pulled
* over the network via [recordDownload], and [MailBackfiller] / [MailSyncer] consult
* [isOverDailyBudget] before starting a fresh prefetch batch for an account, deferring the rest of the
* day's prefetch once the budget is reached.
*
* **Proactive, not reactive — orthogonal to [AccountThrottleGate].** The #360 gate only fires once a
* provider actually rejects a request; this tracker heads that off by pacing our OWN traffic against a
* budget Gmail documents but does not necessarily announce hitting. It has the same relationship to
* [AccountThrottleGate] that [InteractiveImapGate] already documents having with it: a separate,
* composing mechanism, not a duplicate or a replacement.
*
* **Self-healing.** A day boundary (a wall-clock day number derived from [nowMillis]) resets an
* account's tracked total, so a deferred account automatically resumes full prefetch the next day with
* no explicit reset needed — mirrors how [AccountThrottleGate]'s backoff window elapses on its own.
*
* **Interactive traffic is deliberately NOT tracked here.** Only background prefetch (issue #361's
* "sync/backfill" scope) feeds this tracker — opening a message, downloading a tapped attachment, and
* loading inline images are never deferred by a budget (the same interactive-priority principle
* #355/#360 already apply), so counting their bytes here would only make the tracker's *deferral*
* decision — which exclusively affects background prefetch — less representative of what it can
* actually still influence.
*
* State lives only in-process (a `@Singleton`); a process restart clears it, which simply means a
* fresh process re-earns its budget for the (partial) remainder of the day — a conservative direction
* to fail in, same as [AccountThrottleGate]'s reset-on-restart. Every log line is PII-free:
* [accountLogRef] for the account, and byte counts only.
*/
@Singleton
class GmailBandwidthTracker internal constructor(private val nowMillis: () -> Long) {
/** Production wiring: the real wall clock. */
@Inject constructor() : this(nowMillis = System::currentTimeMillis)
/** One account's running download total for [dayEpoch] (whole days since the epoch). */
private data class Window(val dayEpoch: Long, val bytes: Long)
private val windows = ConcurrentHashMap<String, Window>()
/**
* Adds [bytes] to [accountId]'s running total for today, starting a fresh window if the day has
* rolled over since the last call (yesterday's total is simply discarded, not carried forward). A
* no-op for `bytes <= 0`. Logs once, PII-free, the moment this call carries the account from under
* budget to at-or-over it — not on every call, so an account that stays over budget for the rest of
* a sync/backfill pass doesn't spam the log.
*/
fun recordDownload(accountId: String, bytes: Long) {
if (bytes <= 0L) return
val day = currentDayEpoch()
val before = windows[accountId]?.takeIf { it.dayEpoch == day }?.bytes ?: 0L
val updated = windows.compute(accountId) { _, previous ->
val carried = if (previous != null && previous.dayEpoch == day) previous.bytes else 0L
Window(day, carried + bytes)
}!!
val budget = GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES
if (before < budget && updated.bytes >= budget) {
AppLog.w(TAG, "daily download budget reached ${accountLogRef(accountId)}")
}
}
/** [accountId]'s tracked download bytes so far today, or 0 when untracked or the day has rolled over. */
fun bytesDownloadedToday(accountId: String): Long {
val window = windows[accountId] ?: return 0L
return if (window.dayEpoch == currentDayEpoch()) window.bytes else 0L
}
/** True once [accountId]'s tracked downloads for today reach [GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES]. */
fun isOverDailyBudget(accountId: String): Boolean =
bytesDownloadedToday(accountId) >= GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES
private fun currentDayEpoch(): Long = nowMillis() / MILLIS_PER_DAY
private companion object {
const val TAG = "GmailBandwidth"
const val MILLIS_PER_DAY = 24 * 60 * 60 * 1000L
}
}
@@ -0,0 +1,76 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import org.libremail.domain.model.Account
import org.libremail.domain.model.MailProvider
/**
* Gmail's documented IMAP connection and bandwidth ceilings (issue #361) — the Gmail-specific config
* that feeds the shared, provider-agnostic pacing machinery ([AccountThrottleGate]'s reactive backoff,
* issue #360; [BackfillPacer]'s proactive inter-slice cooldown, issue #356) instead of reinventing
* either. Per-account, per-provider connection/bandwidth caps were deliberately deferred out of both
* (see [org.libremail.mail.ImapConnectionCache]'s "separate effort #356/#360-#364" note); this is
* Gmail's slice of that follow-up. Kept additive and provider-scoped — the sibling Yahoo (#362), iCloud
* (#363), and Outlook/Graph (#364) tickets land their own provider config independently.
*
* Values are Google's documented Gmail IMAP limits (as referenced by issue #361):
* - 15 max simultaneous IMAP connections per account.
* - 2,500 MB/day download, 500 MB/day upload.
* - 10,000 messages per label, 10,000 labels.
*
* Pure data + provider detection only — no state, no logging (mirrors how [ThrottleBackoff] and
* [ThrottleClassifier] stay pure while [AccountThrottleGate] carries the state and logging). The
* stateful counterpart that actually tracks bytes against [DAILY_DOWNLOAD_BUDGET_BYTES] is
* [GmailBandwidthTracker].
*/
object GmailSyncLimits {
/** Gmail's documented simultaneous-IMAP-connection ceiling, per account. */
const val MAX_IMAP_CONNECTIONS = 15
/**
* Connections proactively reserved for interactive use (never spent by background sync/backfill),
* mirroring the "interactive-request priority over backfill" lever from the #360 umbrella issue.
* LibreMail's IMAP layer already keeps at most ONE reused connection per account
* ([org.libremail.mail.ImapConnectionCache], issue #125/#357) plus one dedicated IMAP-IDLE
* connection ([org.libremail.mail.ImapClient.idle]) — 2 total, well inside the resulting headroom
* regardless of this value — so today this is documented config for any future pooling work rather
* than something that needs active enforcement (see `GmailSyncLimitsTest` for the invariant that
* ties the two together).
*/
const val INTERACTIVE_RESERVED_CONNECTIONS = 1
/** [MAX_IMAP_CONNECTIONS] minus [INTERACTIVE_RESERVED_CONNECTIONS] — background sync/backfill's budget. */
const val MAX_BACKGROUND_IMAP_CONNECTIONS = MAX_IMAP_CONNECTIONS - INTERACTIVE_RESERVED_CONNECTIONS
private const val BYTES_PER_MB = 1024L * 1024L
/**
* Gmail's documented daily download budget. [GmailBandwidthTracker] accumulates bytes actually
* pulled by background prefetch (see
* [org.libremail.data.repository.MailRepositoryImpl.prefetchMessage]) against this; [MailBackfiller]
* and [MailSyncer] defer further prefetch for an account once it is reached, resuming automatically
* the next day. Interactive fetches (opening a message, downloading a tapped attachment, loading
* inline images) are never gated by it — the same interactive-priority principle #355/#360 already
* apply elsewhere.
*/
const val DAILY_DOWNLOAD_BUDGET_BYTES = 2_500L * BYTES_PER_MB
/**
* Gmail's documented daily upload budget (SMTP send). Captured here for completeness against
* issue #361's documented limits; sending/composing is a separate path from this issue's
* sync/backfill scope, so it is not enforced by this change.
*/
const val DAILY_UPLOAD_BUDGET_BYTES = 500L * BYTES_PER_MB
/** Gmail's documented per-label message ceiling. */
const val MAX_MESSAGES_PER_LABEL = 10_000
/** Gmail's documented total-labels ceiling. */
const val MAX_LABELS = 10_000
/**
* True when [account]'s IMAP host resolves to [MailProvider.GMAIL] (including its legacy
* `imap.googlemail.com` alias) — the single place issue #361's caps decide "is this Gmail".
*/
fun appliesTo(account: Account): Boolean = MailProvider.forImapHost(account.imap.host) == MailProvider.GMAIL
}
@@ -60,6 +60,7 @@ class MailBackfiller @Inject constructor(
private val maintenanceGate: MailMaintenanceGate,
private val throttleGate: AccountThrottleGate,
private val interactiveGate: InteractiveImapGate,
private val bandwidthTracker: GmailBandwidthTracker,
private val authGate: AuthThrottleGate,
) {
/** One folder's slice outcome: pages fetched, and whether an immediate follow-up slice has work to do. */
@@ -233,7 +234,7 @@ class MailBackfiller @Inject constructor(
}
beforeUid = nextBeforeUid
backfillProgressDao.upsert(BackfillProgressEntity(account.id, folder, beforeUid, complete = false))
prefetchIfEnabled(entities.map { it.id })
prefetchIfEnabled(account, entities.map { it.id })
// Breathe between pages so a large mailbox doesn't hammer the server.
delay(BACKFILL_BATCH_DELAY_MS)
}
@@ -312,7 +313,7 @@ class MailBackfiller @Inject constructor(
* Best-effort and cancellable between messages so an interruption stops promptly; anything not
* fetched is filled in lazily when the message is opened.
*/
private suspend fun prefetchIfEnabled(ids: List<String>) {
private suspend fun prefetchIfEnabled(account: Account, ids: List<String>) {
// Debug-only fetch gate (issue #393): pause proactive body prefetch so a later open is a genuine
// uncached fetch. Header paging above is untouched (its own gate is the BackfillWorker entry), so
// history still lands; a skipped body is filled in lazily on open. Compiled out of release
@@ -327,6 +328,17 @@ class MailBackfiller @Inject constructor(
battery = batteryStatusProvider.current(),
)
if (!shouldPrefetch) return
// Gmail-specific proactive bandwidth pacing (#361): once this account's tracked downloads for
// today reach Gmail's documented daily budget, defer body/attachment prefetch for the rest of
// the day instead of continuing to spend it — header paging above is unaffected, and a fresh
// cycle resumes automatically once the day rolls over (GmailBandwidthTracker). Orthogonal to
// the #360 throttle skip above (which only fires once the provider actually rejects a request)
// and #356's BackfillPacer (which paces slice cadence, not bytes) — same graceful-degradation
// shape as both, composing rather than duplicating either.
if (GmailSyncLimits.appliesTo(account) && bandwidthTracker.isOverDailyBudget(account.id)) {
AppLog.i(TAG, "prefetch deferred ${accountLogRef(account.id)}: Gmail daily download budget reached")
return
}
for (id in ids) {
currentCoroutineContext().ensureActive()
mailRepository.prefetchMessage(id)
@@ -41,6 +41,7 @@ class MailSyncer @Inject constructor(
private val notifier: MailNotifier,
private val mailRepository: MailRepository,
private val throttleGate: AccountThrottleGate,
private val bandwidthTracker: GmailBandwidthTracker,
) : Syncer {
// Serializes all syncing: syncAll/syncAccount/syncFolder are invoked concurrently by the periodic
// worker, pull-to-refresh, one-shot syncs, folder opens, and one IDLE watcher per account. Without
@@ -187,6 +188,17 @@ class MailSyncer @Inject constructor(
battery = batteryStatusProvider.current(),
)
if (!shouldPrefetch) return
// Gmail-specific proactive bandwidth pacing (#361): once this account's tracked downloads for
// today reach Gmail's documented daily budget, defer body/attachment prefetch for the rest of
// the day instead of continuing to spend it — header sync above is unaffected, and a fresh
// cycle resumes automatically once the day rolls over (GmailBandwidthTracker). Orthogonal to
// #360's reactive AccountThrottleGate (only fires once the provider actually rejects a request)
// and #356's BackfillPacer (paces backfill slice cadence, not bytes) — same graceful-degradation
// shape as both, composing rather than duplicating either.
if (GmailSyncLimits.appliesTo(account) && bandwidthTracker.isOverDailyBudget(account.id)) {
AppLog.i(TAG, "prefetch deferred ${accountLogRef(account.id)}: Gmail daily download budget reached")
return
}
for (id in messageDao.getUnfetchedIds(account.id, folder)) {
currentCoroutineContext().ensureActive()
mailRepository.prefetchMessage(id) // best-effort; swallows its own per-message failures
@@ -7,10 +7,13 @@ import kotlinx.coroutines.withContext
import org.json.JSONArray
import org.json.JSONObject
import org.libremail.domain.model.OutgoingMessage
import org.libremail.mail.graph.GraphHttpClient
import org.libremail.mail.graph.GraphRequest
import org.libremail.mail.graph.GraphThrottle
import org.libremail.mail.graph.GraphTransportException
import org.libremail.mail.graph.isHttpSuccess
import org.libremail.reporting.AppLog
import java.io.IOException
import java.net.HttpURLConnection
import java.net.URL
import org.libremail.reporting.accountLogRef
import java.util.Base64
import javax.inject.Inject
import javax.inject.Singleton
@@ -26,9 +29,15 @@ class GraphSendException(message: String, val mayHaveSent: Boolean, cause: Throw
/**
* Sends mail via Microsoft Graph `me/sendMail` — Microsoft's preferred send path for Outlook /
* Microsoft 365, used in place of SMTP. Authenticated with a Graph access token (Bearer).
*
* The request goes through [GraphThrottle] (issue #364), so a Graph 429/503 is honored — its
* `Retry-After` is respected and the send retried after that wait — and recorded against the account's
* shared reactive backoff gate (#360), which also cools that account's IMAP background work down. Only
* after the honored retry is exhausted does a throttled or rejected response surface as a
* [GraphSendException] (`mayHaveSent = false`, safe for the outbox to fall back to SMTP).
*/
@Singleton
class GraphSender @Inject constructor() {
class GraphSender @Inject constructor(private val httpClient: GraphHttpClient, private val throttle: GraphThrottle) {
suspend fun send(
accessToken: String,
@@ -38,62 +47,60 @@ class GraphSender @Inject constructor() {
// Guard before any attachment is read into memory: Graph sendMail carries attachment bytes inline
// (base64) in a single ~4 MB request, so an oversized file would blow that request limit and risk
// an OOM from readBytes(). Fail with mayHaveSent=false so the outbox falls back to SMTP, which
// streams attachments and handles far larger files (#298).
// streams attachments and handles far larger files (#298). (An over-4 MB attachment sent *through*
// Graph would need a draft + createUploadSession chunked upload — see GraphUploadSession — which
// requires the Mail.ReadWrite scope this send-only token does not hold; SMTP fallback is simpler.)
attachments.firstOrNull { it.file.length() > MAX_ATTACHMENT_BYTES }?.let {
AppLog.w(TAG, "Attachment over Graph sendMail size limit; not sending via Graph")
throw GraphSendException("Attachment exceeds the Graph sendMail size limit", mayHaveSent = false)
}
val payload = buildSendMailPayload(message, attachments)
val connection = (URL(SEND_MAIL_URL).openConnection() as HttpURLConnection).apply {
requestMethod = "POST"
connectTimeout = TIMEOUT_MS
readTimeout = TIMEOUT_MS
doOutput = true
setRequestProperty("Authorization", "Bearer $accessToken")
setRequestProperty("Content-Type", "application/json; charset=utf-8")
val request = GraphRequest(
method = "POST",
url = SEND_MAIL_URL,
headers = mapOf(
"Authorization" to "Bearer $accessToken",
"Content-Type" to "application/json; charset=utf-8",
),
body = payload.toByteArray(Charsets.UTF_8),
)
val response = try {
throttle.execute(message.accountId, httpClient, request, maxRetries = SEND_MAX_RETRIES)
} catch (e: GraphTransportException) {
// No HTTP response at all. A lost response means Graph may already have accepted+sent the
// message, so callers must NOT fall back (it would duplicate); a transmit failure is safe.
throw GraphSendException(
if (e.mayHaveSent) {
"Graph sendMail sent but no response received"
} else {
"Graph sendMail could not be transmitted"
},
mayHaveSent = e.mayHaveSent,
cause = e,
)
}
try {
// Failure here means the request never reached Graph — safe to fall back/retry.
try {
connection.outputStream.use { it.write(payload.toByteArray(Charsets.UTF_8)) }
} catch (e: IOException) {
throw GraphSendException("Graph sendMail could not be transmitted", mayHaveSent = false, cause = e)
}
// The request was fully sent; if we can't read the response, Graph may already have
// accepted and sent it — do not fall back to SMTP or the message would be duplicated.
val code = try {
connection.responseCode
} catch (e: IOException) {
throw GraphSendException(
"Graph sendMail sent but no response received",
mayHaveSent = true,
cause = e,
)
}
if (code !in HTTP_OK_MIN..HTTP_OK_MAX) {
val body = (connection.errorStream ?: connection.inputStream)
?.bufferedReader()?.use { it.readText() }.orEmpty()
// An explicit non-2xx means Graph rejected (did not send) — safe to fall back.
throw GraphSendException(
"Graph sendMail failed (HTTP $code): ${body.take(ERROR_BODY_LIMIT)}",
mayHaveSent = false,
)
}
} finally {
connection.disconnect()
if (!isHttpSuccess(response.status)) {
// An explicit non-2xx (including a 429 the retry budget couldn't clear) means Graph did not
// send — safe for the outbox to fall back to SMTP.
throw GraphSendException(
"Graph sendMail failed (HTTP ${response.status}): ${response.body.take(ERROR_BODY_LIMIT)}",
mayHaveSent = false,
)
}
AppLog.i(TAG, "graph sendMail ok ${accountLogRef(message.accountId)}")
}
private companion object {
const val TAG = "GraphSender"
const val SEND_MAIL_URL = "https://graph.microsoft.com/v1.0/me/sendMail"
const val TIMEOUT_MS = 15_000
// Per-file ceiling kept below Graph sendMail's ~4 MB whole-request cap, so one attachment can never
// exceed the request limit or OOM when read into the base64 payload; larger files fall back to SMTP.
const val MAX_ATTACHMENT_BYTES = 3L * 1024 * 1024
const val HTTP_OK_MIN = 200
const val HTTP_OK_MAX = 299
// One honored Retry-After retry for a user-initiated send: respect the provider's explicit
// "slow down" once, then fall back to SMTP rather than block the outbox drain for long.
const val SEND_MAX_RETRIES = 1
const val ERROR_BODY_LIMIT = 500
}
}
@@ -0,0 +1,141 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import org.json.JSONArray
import org.json.JSONObject
import org.libremail.reporting.AppLog
import org.libremail.reporting.accountLogRef
import javax.inject.Inject
import javax.inject.Singleton
/**
* One request inside a Graph `$batch`. [url] is **relative** to the service root (e.g.
* `/me/messages/{id}`), as the batch envelope requires. [body] is the optional JSON payload for a
* write sub-request; [headers] carries any per-op header the sub-request needs.
*/
data class GraphSubRequest(
val id: String,
val method: String,
val url: String,
val headers: Map<String, String> = emptyMap(),
val body: JSONObject? = null,
)
/**
* The result of one `$batch` sub-request, correlated back to its [GraphSubRequest.id]. [status] is the
* sub-response's own HTTP status (a batch call can return 200 overall while an individual op is 429),
* [body] its JSON payload if any, and [retryAfterMillis] the parsed per-op `Retry-After` when throttled.
*/
data class GraphSubResponse(
val id: String,
val status: Int,
val body: JSONObject? = null,
val retryAfterMillis: Long? = null,
)
/**
* Multiplexes many Microsoft Graph reads/writes through the `$batch` endpoint (issue #364) so N calls
* collapse to `ceil(N / 20)` HTTP round-trips — the documented cap is 20 operations per batch. This is
* the lever for the 10,000-requests-per-10-minutes window: a backfill page that would otherwise fan out
* one request per message becomes a single batch, staying far under the budget and the 4-concurrent cap.
*
* Every batch call runs through [GraphThrottle], so the envelope shares the app-wide Graph throttle
* policy; a 429 surfaced **inside** the envelope (per-op) is fed back to the shared account backoff too.
* Logging is PII-free (account ref, counts, statuses only).
*/
@Singleton
class GraphBatch @Inject constructor(private val throttle: GraphThrottle) {
/**
* Executes [requests] for [accountId] via [client], chunked into batches of at most
* [MAX_BATCH_OPERATIONS]. Returns every sub-response, correlated by id and preserving request order.
* Empty in, empty out (no HTTP call).
*/
suspend fun execute(
accountId: String,
client: GraphHttpClient,
requests: List<GraphSubRequest>,
): List<GraphSubResponse> {
if (requests.isEmpty()) return emptyList()
val chunks = requests.chunked(MAX_BATCH_OPERATIONS)
val responses = ArrayList<GraphSubResponse>(requests.size)
for (chunk in chunks) {
val request = GraphRequest(
method = "POST",
url = BATCH_URL,
headers = mapOf(HEADER_CONTENT_TYPE to CONTENT_TYPE_JSON),
body = buildBatchPayload(chunk).toByteArray(Charsets.UTF_8),
)
val parsed = parseBatchResponses(throttle.execute(accountId, client, request).body)
parsed.forEach { sub ->
if (!isHttpSuccess(sub.status)) {
throttle.recordSubResponseThrottle(accountId, sub.status, sub.retryAfterMillis)
}
}
responses += parsed
}
AppLog.d(
TAG,
"graph batch ${accountLogRef(accountId)} ops=${requests.size} calls=${chunks.size}",
)
return responses
}
private companion object {
const val TAG = "GraphBatch"
/** Graph accepts at most 20 operations per `$batch` request. */
const val MAX_BATCH_OPERATIONS = 20
const val BATCH_URL = "https://graph.microsoft.com/v1.0/\$batch"
const val HEADER_CONTENT_TYPE = "Content-Type"
const val CONTENT_TYPE_JSON = "application/json"
}
}
/** Builds the `{"requests":[...]}` JSON body for a chunk of sub-requests (pure, so it is unit-testable). */
internal fun buildBatchPayload(requests: List<GraphSubRequest>): String {
val array = JSONArray()
requests.forEach { request ->
val obj = JSONObject()
.put("id", request.id)
.put("method", request.method)
.put("url", request.url)
if (request.headers.isNotEmpty()) {
val headers = JSONObject()
request.headers.forEach { (name, value) -> headers.put(name, value) }
obj.put("headers", headers)
}
request.body?.let { obj.put("body", it) }
array.put(obj)
}
return JSONObject().put("requests", array).toString()
}
/**
* Parses a `$batch` response body's `{"responses":[...]}` array into [GraphSubResponse]s (pure, so it is
* unit-testable). A per-op `Retry-After` header (case-insensitive) is read into
* [GraphSubResponse.retryAfterMillis]. A malformed/empty body yields an empty list rather than throwing.
*/
internal fun parseBatchResponses(body: String): List<GraphSubResponse> {
if (body.isBlank()) return emptyList()
val responses = runCatching { JSONObject(body).optJSONArray("responses") }.getOrNull() ?: return emptyList()
val out = ArrayList<GraphSubResponse>(responses.length())
for (index in 0 until responses.length()) {
val item = responses.optJSONObject(index) ?: continue
out += GraphSubResponse(
id = item.optString("id"),
status = item.optInt("status"),
body = item.optJSONObject("body"),
retryAfterMillis = subResponseRetryAfterMillis(item),
)
}
return out
}
/** Reads a case-insensitive `Retry-After` from a sub-response's `headers` object, parsed to millis. */
private fun subResponseRetryAfterMillis(item: JSONObject): Long? {
val headers = item.optJSONObject("headers") ?: return null
val key = headers.keys().asSequence().firstOrNull { it.equals("Retry-After", ignoreCase = true) } ?: return null
return parseRetryAfterMillis(headers.optString(key), System.currentTimeMillis())
}
@@ -0,0 +1,137 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.withContext
import java.io.IOException
import java.net.HttpURLConnection
import java.net.URL
import java.time.ZonedDateTime
import java.time.format.DateTimeFormatter
import javax.inject.Inject
import javax.inject.Singleton
/** Base of the HTTP success range (200), reused across the Graph layer. */
internal const val HTTP_SUCCESS_MIN = 200
/** Top of the HTTP success range (299), reused across the Graph layer. */
internal const val HTTP_SUCCESS_MAX = 299
/** True for any 2xx status — the shared "the request was accepted" predicate for the Graph layer. */
internal fun isHttpSuccess(status: Int): Boolean = status in HTTP_SUCCESS_MIN..HTTP_SUCCESS_MAX
/**
* A single Microsoft Graph HTTP request. [url] is absolute for a top-level call (`me/sendMail`, the
* `$batch` endpoint, an upload-session URL) and [body], when non-null, is the raw bytes to transmit —
* JSON for most calls, a binary slice for an upload-session chunk. [headers] carries the bearer
* `Authorization` plus any per-call header (`Content-Type`, `Content-Range`).
*
* Intentionally a plain class, not a `data class`: it carries a [ByteArray] (value-equality would be a
* footgun and detekt's `ArrayInDataClass` forbids it) and is never compared or destructured — callers
* only read its properties.
*/
class GraphRequest(
val method: String,
val url: String,
val headers: Map<String, String> = emptyMap(),
val body: ByteArray? = null,
)
/**
* A completed Graph HTTP response: the [status] code, the decoded [body] text, and — when the server
* sent a `Retry-After` header — the parsed minimum wait in [retryAfterMillis]. A response object means
* the server answered (any status, including 429/503); a request that never got an answer surfaces as
* a [GraphTransportException] instead so the caller can reason about whether it may already have sent.
*/
data class GraphResponse(val status: Int, val body: String, val retryAfterMillis: Long? = null)
/**
* Thrown when a Graph request produced **no HTTP response** at all. [mayHaveSent] is true only when the
* request was fully transmitted but the response could not be read — the server may already have acted
* on it, so a `sendMail` caller must NOT blindly retry or fall back (it could duplicate the message).
* A failure to even transmit the body sets it false (safe to retry / fall back).
*/
class GraphTransportException(message: String, val mayHaveSent: Boolean, cause: Throwable? = null) :
Exception(message, cause)
/**
* The single Microsoft Graph HTTP transport seam. Deliberately thin: it opens one [HttpURLConnection],
* writes the body, reads the status + body + `Retry-After`, and always disconnects — it does **no**
* retrying, backoff, or throttle bookkeeping (that is [GraphThrottle]'s job, layered on top so every
* Graph caller — send, `$batch`, chunked upload — shares one throttle policy).
*
* `open` with an injectable no-arg constructor so unit/instrumented tests substitute an in-memory fake
* (no network, no process-wide URL handler) while production talks to `graph.microsoft.com`.
*/
@Singleton
open class GraphHttpClient @Inject constructor() {
/**
* Executes [request], returning the server's [GraphResponse] for **any** status it answered with.
* Throws [GraphTransportException] when there was no response: [GraphTransportException.mayHaveSent]
* distinguishes a body that never left the device (false) from one fully sent but whose response was
* lost (true), preserving the send path's no-duplicate guarantee.
*/
open suspend fun execute(request: GraphRequest): GraphResponse = withContext(Dispatchers.IO) {
val connection = (URL(request.url).openConnection() as HttpURLConnection).apply {
requestMethod = request.method
connectTimeout = TIMEOUT_MS
readTimeout = TIMEOUT_MS
request.headers.forEach { (name, value) -> setRequestProperty(name, value) }
if (request.body != null) doOutput = true
}
try {
writeBody(connection, request.body)
val status = readStatus(connection)
val stream = if (isHttpSuccess(status)) connection.inputStream else connection.errorStream
val body = stream?.bufferedReader()?.use { it.readText() }.orEmpty()
val retryAfterHeader = connection.getHeaderField(HEADER_RETRY_AFTER)
val retryAfter = parseRetryAfterMillis(retryAfterHeader, System.currentTimeMillis())
GraphResponse(status = status, body = body, retryAfterMillis = retryAfter)
} finally {
connection.disconnect()
}
}
/** Writes the request body, mapping a transmit failure to a safe-to-retry [GraphTransportException]. */
private fun writeBody(connection: HttpURLConnection, body: ByteArray?) {
if (body == null) return
try {
connection.outputStream.use { it.write(body) }
} catch (e: IOException) {
throw GraphTransportException("Graph request could not be transmitted", mayHaveSent = false, cause = e)
}
}
/** Reads the status code, mapping a lost response to a maybe-sent [GraphTransportException]. */
private fun readStatus(connection: HttpURLConnection): Int = try {
connection.responseCode
} catch (e: IOException) {
throw GraphTransportException("Graph request sent but no response received", mayHaveSent = true, cause = e)
}
private companion object {
const val TIMEOUT_MS = 15_000
const val HEADER_RETRY_AFTER = "Retry-After"
}
}
/** Milliseconds in one second — the unit of a numeric `Retry-After` (delta-seconds) header. */
private const val MILLIS_PER_SECOND = 1000L
/**
* Parses an HTTP `Retry-After` header ([headerValue]) into a non-negative millisecond wait, or null
* when it is absent/unparseable. Graph normally sends the RFC 7231 delta-seconds form (an integer
* count of seconds); the HTTP-date form is also accepted and converted to a wait relative to
* [nowMillis]. A past date (or a negative/garbage value) clamps to a valid wait rather than going
* negative. Pure and clock-injected so it is deterministically unit-testable.
*/
internal fun parseRetryAfterMillis(headerValue: String?, nowMillis: Long): Long? {
val trimmed = headerValue?.trim().orEmpty()
if (trimmed.isEmpty()) return null
trimmed.toLongOrNull()?.let { seconds -> return (seconds * MILLIS_PER_SECOND).coerceAtLeast(0L) }
return runCatching {
val instant = ZonedDateTime.parse(trimmed, DateTimeFormatter.RFC_1123_DATE_TIME).toInstant()
(instant.toEpochMilli() - nowMillis).coerceAtLeast(0L)
}.getOrNull()
}
@@ -0,0 +1,101 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import kotlinx.coroutines.delay
import kotlinx.coroutines.sync.Semaphore
import kotlinx.coroutines.sync.withPermit
import org.libremail.data.sync.AccountThrottleGate
import org.libremail.data.sync.ThrottleClassifier
import org.libremail.reporting.AppLog
import org.libremail.reporting.accountLogRef
import javax.inject.Inject
import javax.inject.Singleton
/**
* The Microsoft Graph throttle policy layered over the raw [GraphHttpClient] transport — the Graph-side
* embodiment of issue #364, composed with the shared reactive backoff of issue #360
* ([AccountThrottleGate]). Every Graph call — `me/sendMail`, `$batch`, an upload-session chunk — runs
* through [execute], so all of them share one throttle policy:
*
* - **Concurrency cap.** Graph rejects a mailbox's 5th concurrent request with a 429, so at most
* [MAX_CONCURRENT_REQUESTS] Graph calls run at once (a [Semaphore]); the rest queue rather than
* provoke the limit. Held only around the network round-trip, released before any backoff sleep.
* - **Honor `Retry-After` on 429/503.** A throttled response is classified via
* [ThrottleClassifier.classifyHttpStatus] and recorded against the account in the shared gate, whose
* backoff honors the server's `Retry-After` as a floor. The call then waits that long and retries, up
* to [maxRetries] extra attempts, before surfacing the last throttled response to the caller.
* - **Cross-path cooperation.** Because the gate is keyed by account id and shared with the IMAP
* backfill/sync paths (#360), a Graph 429 also cools that account's background IMAP work down, and a
* clean Graph response clears any lingering backoff — the send and receive paths never fight the same
* provider limit from two directions.
*
* All logging is PII-free ([accountLogRef], statuses, and durations only).
*/
@Singleton
class GraphThrottle @Inject constructor(private val throttleGate: AccountThrottleGate) {
/** At most four Graph requests in flight at once — Graph 429s the 5th concurrent request per mailbox. */
private val concurrency = Semaphore(MAX_CONCURRENT_REQUESTS)
/**
* Runs [request] through [client] for [accountId], honoring Graph throttling. A 429/503 is recorded
* against the account's shared backoff gate and retried after the honored wait, up to [maxRetries]
* additional attempts; a 2xx clears any backoff. Returns the final [GraphResponse] — including a
* still-throttled one once retries are exhausted — so the caller maps it to its own result.
* [GraphTransportException] (no response at all) propagates unretried: the caller owns the
* may-have-sent decision.
*/
suspend fun execute(
accountId: String,
client: GraphHttpClient,
request: GraphRequest,
maxRetries: Int = DEFAULT_MAX_RETRIES,
): GraphResponse {
var retries = 0
while (true) {
val response = concurrency.withPermit { client.execute(request) }
val signal = ThrottleClassifier.classifyHttpStatus(response.status, response.retryAfterMillis)
if (signal == null) {
if (isHttpSuccess(response.status)) throttleGate.onSuccess(accountId)
return response
}
val backoff = throttleGate.onThrottle(accountId, signal)
if (retries >= maxRetries) {
AppLog.w(
TAG,
"graph throttled ${accountLogRef(accountId)} status=${response.status} " +
"retries exhausted ($retries/$maxRetries)",
)
return response
}
retries++
AppLog.w(
TAG,
"graph throttled ${accountLogRef(accountId)} status=${response.status} " +
"backoff=${backoff}ms retry=$retries/$maxRetries",
)
delay(backoff)
}
}
/**
* Records a throttle observed **inside** a `$batch` envelope — Graph answers the batch call 200 but
* can mark individual sub-responses 429/503 with their own `Retry-After`. Feeds the shared gate so a
* multiplexed read that hit the limit backs the account off exactly like a top-level 429 would.
* Returns the resulting backoff in ms, or null when [status] was not a throttle.
*/
fun recordSubResponseThrottle(accountId: String, status: Int, retryAfterMillis: Long?): Long? {
val signal = ThrottleClassifier.classifyHttpStatus(status, retryAfterMillis) ?: return null
return throttleGate.onThrottle(accountId, signal)
}
private companion object {
const val TAG = "GraphThrottle"
/** Graph caps a single mailbox at four concurrent requests before 429-ing the rest. */
const val MAX_CONCURRENT_REQUESTS = 4
/** Default extra attempts after the first, for background/multiplexed calls (send overrides lower). */
const val DEFAULT_MAX_RETRIES = 2
}
}
@@ -0,0 +1,109 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import org.json.JSONObject
import org.libremail.reporting.AppLog
import org.libremail.reporting.accountLogRef
import javax.inject.Inject
import javax.inject.Singleton
/**
* Uploads large content to Microsoft Graph in chunks via an upload session (issue #364). Graph caps a
* one-shot (inline / single-request) upload at ~4 MB; anything larger must go through
* `createUploadSession` and a series of ranged `PUT`s, which also keeps each request well inside the
* 150 MB-per-5-minutes bandwidth window instead of one giant spike.
*
* Each `PUT` (and the session-create `POST`) runs through [GraphThrottle], so chunked uploads obey the
* same concurrency cap and honor `Retry-After` between chunks — a throttle mid-upload pauses and
* resumes rather than failing the whole transfer. Chunk size is a multiple of [CHUNK_MULTIPLE_BYTES]
* (320 KiB), as Graph requires for every non-final chunk. Logging is PII-free.
*/
@Singleton
class GraphUploadSession @Inject constructor(private val throttle: GraphThrottle) {
/** True when [sizeBytes] exceeds Graph's one-shot ceiling and must use a chunked upload session. */
fun requiresUploadSession(sizeBytes: Long): Boolean = sizeBytes > LARGE_ATTACHMENT_THRESHOLD_BYTES
/**
* Creates an upload session at [createSessionUrl] (with [createSessionBody], e.g. the
* `{"AttachmentItem": …}` descriptor) and uploads [content] to the returned `uploadUrl` in
* [chunkSize]-byte ranges. Returns the final chunk's [GraphResponse] (Graph answers the last chunk
* 200/201 and intermediate chunks 202). A non-2xx session-create, or a non-2xx chunk, stops the
* upload and returns that response so the caller can react.
*
* Requires non-empty [content] and a [chunkSize] that is a positive multiple of
* [CHUNK_MULTIPLE_BYTES] (both enforced by `require`), and a session response carrying an `uploadUrl`.
*/
suspend fun upload(
accountId: String,
client: GraphHttpClient,
createSessionUrl: String,
createSessionBody: JSONObject,
content: ByteArray,
chunkSize: Int = DEFAULT_CHUNK_SIZE_BYTES,
): GraphResponse {
require(content.isNotEmpty()) { "cannot upload empty content" }
require(chunkSize > 0 && chunkSize % CHUNK_MULTIPLE_BYTES == 0) {
"chunk size must be a positive multiple of $CHUNK_MULTIPLE_BYTES bytes"
}
val create = throttle.execute(
accountId,
client,
GraphRequest(
method = "POST",
url = createSessionUrl,
headers = mapOf(HEADER_CONTENT_TYPE to CONTENT_TYPE_JSON),
body = createSessionBody.toString().toByteArray(Charsets.UTF_8),
),
)
if (!isHttpSuccess(create.status)) {
AppLog.w(TAG, "graph upload session create failed ${accountLogRef(accountId)} status=${create.status}")
return create
}
val uploadUrl = runCatching { JSONObject(create.body).optString("uploadUrl") }.getOrNull()
require(!uploadUrl.isNullOrBlank()) { "upload session response carried no uploadUrl" }
val total = content.size
var offset = 0
var chunks = 0
var last = create
while (offset < total) {
val end = minOf(offset + chunkSize, total)
last = throttle.execute(
accountId,
client,
GraphRequest(
method = "PUT",
url = uploadUrl,
headers = mapOf(HEADER_CONTENT_RANGE to "bytes $offset-${end - 1}/$total"),
body = content.copyOfRange(offset, end),
),
)
chunks++
if (!isHttpSuccess(last.status)) {
AppLog.w(TAG, "graph upload chunk failed ${accountLogRef(accountId)} status=${last.status} at=$chunks")
return last
}
offset = end
}
AppLog.i(TAG, "graph upload ok ${accountLogRef(accountId)} bytes=$total chunks=$chunks")
return last
}
private companion object {
const val TAG = "GraphUpload"
/** Graph's one-shot upload ceiling (~4 MB); larger content must use a chunked session. */
const val LARGE_ATTACHMENT_THRESHOLD_BYTES = 4L * 1024 * 1024
/** Every non-final upload chunk must be a multiple of 320 KiB, per Graph. */
const val CHUNK_MULTIPLE_BYTES = 320 * 1024
/** Default ~4.7 MB chunk (a 320 KiB multiple), inside Graph's recommended 5–10 MiB range. */
const val DEFAULT_CHUNK_SIZE_BYTES = CHUNK_MULTIPLE_BYTES * 15
const val HEADER_CONTENT_TYPE = "Content-Type"
const val HEADER_CONTENT_RANGE = "Content-Range"
const val CONTENT_TYPE_JSON = "application/json"
}
}
@@ -16,6 +16,7 @@ import org.libremail.data.local.dao.DraftDao
import org.libremail.data.local.dao.OutboxDao
import org.libremail.data.local.entity.DraftEntity
import org.libremail.data.local.entity.OutboxEntity
import org.libremail.data.sync.GmailBandwidthTracker
import org.libremail.data.sync.InteractiveImapGate
import java.nio.file.Files
@@ -47,6 +48,7 @@ class MailRepositoryGrantsTest {
signatureRepository = mockk(relaxed = true),
attachmentUriGrants = attachmentUriGrants,
interactiveGate = InteractiveImapGate(),
bandwidthTracker = GmailBandwidthTracker(),
)
@Test
@@ -40,6 +40,7 @@ import org.libremail.data.local.entity.OutboxEntity
import org.libremail.data.local.entity.ServerConfigEmbedded
import org.libremail.data.settings.AccountSettingsRepository
import org.libremail.data.settings.SignatureRepository
import org.libremail.data.sync.GmailBandwidthTracker
import org.libremail.data.sync.InteractiveImapGate
import org.libremail.data.sync.MailConnectionFactory
import org.libremail.data.sync.SendScheduler
@@ -99,6 +100,7 @@ class MailRepositoryImplCoverageTest {
signatureRepository = signatureRepository,
attachmentUriGrants = mockk<AttachmentUriGrants>(relaxed = true),
interactiveGate = InteractiveImapGate(),
bandwidthTracker = GmailBandwidthTracker(),
)
// openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain
@@ -40,6 +40,7 @@ import org.libremail.data.local.entity.MessageSummary
import org.libremail.data.local.entity.ServerConfigEmbedded
import org.libremail.data.settings.AccountSettingsRepository
import org.libremail.data.settings.SignatureRepository
import org.libremail.data.sync.GmailBandwidthTracker
import org.libremail.data.sync.InteractiveImapGate
import org.libremail.data.sync.MailConnectionFactory
import org.libremail.domain.model.AccountSettings
@@ -78,6 +79,9 @@ class MailRepositoryImplTest {
// A real gate (cheap, no deps) so tests can observe the interactive-fetch counter it raises (#355).
private val interactiveGate = InteractiveImapGate()
// A real tracker (cheap, no deps) so tests can observe Gmail's daily download-budget accounting (#361).
private val bandwidthTracker = GmailBandwidthTracker()
private val repository = MailRepositoryImpl(
context = context,
messageDao = messageDao,
@@ -94,6 +98,7 @@ class MailRepositoryImplTest {
// Grant-release wiring (deleteDraft / cancelOutboxMessage) is covered by MailRepositoryGrantsTest.
attachmentUriGrants = mockk(relaxed = true),
interactiveGate = interactiveGate,
bandwidthTracker = bandwidthTracker,
)
// openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain
@@ -411,6 +416,48 @@ class MailRepositoryImplTest {
assertEquals("Reply to <ada@example.org>: 3 < 5", snippet.captured)
}
// --- issue #361: Gmail bandwidth-aware prefetch pacing ----------------------------------------
@Test
fun `prefetchMessage records the fetched body and attachment bytes for a gmail account`() = runTest {
val cache = Files.createTempDirectory("attach").toFile()
every { context.cacheDir } returns cache
val id = "acct:INBOX:22"
val body = "Hello world"
coEvery { messageDao.getRouting(id) } returns messageRouting(id, "INBOX")
coEvery { accountDao.getById("acct") } returns
accountEntity().copy(imap = ServerConfigEmbedded("imap.gmail.com", 993, "SSL_TLS"))
coEvery { connectionFactory.imapParamsFor(any()) } returns imapParams()
coEvery { imapClient.fetchBodyPeek(any(), "INBOX", "22") } returns MessageContent(body, isHtml = false)
coEvery { messageDao.updateBody(id, any(), any(), any()) } just Runs
coEvery { attachmentDao.getForMessage(id) } returns listOf(attachmentEntity(id, 0, "photo.jpg"))
coEvery { imapClient.fetchAttachment(any(), "INBOX", "22", 0) } returns
DownloadedAttachment("photo.jpg", "image/jpeg", ByteArray(500))
repository.prefetchMessage(id)
val expectedBytes = body.toByteArray(Charsets.UTF_8).size + 500
assertEquals(expectedBytes.toLong(), bandwidthTracker.bytesDownloadedToday("acct"))
}
@Test
fun `prefetchMessage does not track bytes for a non-gmail account`() = runTest {
val cache = Files.createTempDirectory("attach").toFile()
every { context.cacheDir } returns cache
val id = "acct:INBOX:23"
coEvery { messageDao.getRouting(id) } returns messageRouting(id, "INBOX")
coEvery { accountDao.getById("acct") } returns accountEntity() // non-Gmail host (imap.example.org)
coEvery { connectionFactory.imapParamsFor(any()) } returns imapParams()
coEvery { imapClient.fetchBodyPeek(any(), "INBOX", "23") } returns
MessageContent("Hello world", isHtml = false)
coEvery { messageDao.updateBody(id, any(), any(), any()) } just Runs
coEvery { attachmentDao.getForMessage(id) } returns emptyList()
repository.prefetchMessage(id)
assertEquals(0L, bandwidthTracker.bytesDownloadedToday("acct"))
}
@Test
fun `archive moves messages to the account's archive folder and drops the local rows`() = runTest {
val id = "acct:INBOX:5"
@@ -0,0 +1,131 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import android.util.Log
import io.mockk.every
import io.mockk.mockkStatic
import io.mockk.unmockkAll
import org.junit.After
import org.junit.Before
import org.junit.Test
import org.libremail.reporting.AppLog
import org.libremail.reporting.RingLogBuffer
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
/**
* [GmailBandwidthTracker] (issue #361) must accumulate an account's daily download bytes, isolate
* accounts from one another, roll its window over at the day boundary, report the daily-budget
* threshold accurately, and log only a PII-free, once-per-crossing breadcrumb. Mirrors
* [AccountThrottleGateTest]'s virtual-clock idiom (a plain injected `nowMillis` instead of
* coroutines-test virtual time, since the tracker itself is not a suspend API).
*/
class GmailBandwidthTrackerTest {
private val logBuffer = RingLogBuffer()
/** A manual clock for the day-rollover tests; [tracker] reads it live, so tests advance it by hand. */
private var now = 0L
private fun tracker() = GmailBandwidthTracker(nowMillis = { now })
@Before
fun setUp() {
// AppLog forwards to android.util.Log, a throwing no-op stub under plain JVM unit tests.
mockkStatic(Log::class)
every { Log.w(any<String>(), any<String>()) } returns 0
AppLog.install(logBuffer)
}
@After
fun tearDown() = unmockkAll()
@Test
fun `bytes accumulate across calls for the same account and day`() {
val tracker = tracker()
tracker.recordDownload("acct", 100L)
tracker.recordDownload("acct", 250L)
assertEquals(350L, tracker.bytesDownloadedToday("acct"))
}
@Test
fun `an untracked account reports zero bytes and is not over budget`() {
val tracker = tracker()
assertEquals(0L, tracker.bytesDownloadedToday("acct"))
assertFalse(tracker.isOverDailyBudget("acct"))
}
@Test
fun `crossing the daily download budget marks the account over budget`() {
val tracker = tracker()
tracker.recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES - 1)
assertFalse(tracker.isOverDailyBudget("acct"), "one byte under budget must not trip it")
tracker.recordDownload("acct", 1L)
assertTrue(tracker.isOverDailyBudget("acct"), "reaching the budget exactly must trip it")
}
@Test
fun `a tracked account never affects another account's budget`() {
val tracker = tracker()
tracker.recordDownload("heavy", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
assertTrue(tracker.isOverDailyBudget("heavy"))
assertFalse(tracker.isOverDailyBudget("light"))
assertEquals(0L, tracker.bytesDownloadedToday("light"))
}
@Test
fun `a new day resets the tracked total instead of carrying it forward`() {
val tracker = tracker()
val oneDayMs = 24 * 60 * 60 * 1000L
tracker.recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
assertTrue(tracker.isOverDailyBudget("acct"))
now += oneDayMs
assertFalse(tracker.isOverDailyBudget("acct"), "a new day must clear yesterday's total")
assertEquals(0L, tracker.bytesDownloadedToday("acct"))
tracker.recordDownload("acct", 10L)
assertEquals(10L, tracker.bytesDownloadedToday("acct"), "today's total starts fresh, not carried over")
}
@Test
fun `recording zero or negative bytes is a no-op`() {
val tracker = tracker()
tracker.recordDownload("acct", 0L)
tracker.recordDownload("acct", -5L)
assertEquals(0L, tracker.bytesDownloadedToday("acct"))
}
@Test
fun `crossing the budget logs exactly once and stays PII-free`() {
val tracker = tracker()
val accountId = "imap:user@example.org"
tracker.recordDownload(accountId, GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES) // crosses
tracker.recordDownload(accountId, 10L) // still over budget; must not log again
val messages = logBuffer.snapshot().map { it.message }
assertEquals(1, messages.count { it.contains("daily download budget reached") })
messages.forEach { assertFalse(it.contains("user@example.org"), it) }
}
@Test
fun `staying under budget never logs`() {
val tracker = tracker()
tracker.recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES - 1)
assertTrue(logBuffer.snapshot().isEmpty())
}
}
@@ -0,0 +1,58 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import org.junit.Test
import org.libremail.domain.model.Account
import org.libremail.domain.model.MailProvider
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertTrue
/**
* Locks down Gmail's documented IMAP connection/bandwidth ceilings (issue #361) — easy to mistype and
* painful to debug on-device, so the exact figures are asserted here rather than trusted to a code
* review (mirrors [org.libremail.domain.model.MailProviderTest]'s rationale for the provider presets).
*/
class GmailSyncLimitsTest {
@Test
fun `documented connection and bandwidth ceilings match Gmail's published limits`() {
assertEquals(15, GmailSyncLimits.MAX_IMAP_CONNECTIONS)
assertEquals(1, GmailSyncLimits.INTERACTIVE_RESERVED_CONNECTIONS)
assertEquals(14, GmailSyncLimits.MAX_BACKGROUND_IMAP_CONNECTIONS)
assertEquals(2_500L * 1024 * 1024, GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
assertEquals(500L * 1024 * 1024, GmailSyncLimits.DAILY_UPLOAD_BUDGET_BYTES)
assertEquals(10_000, GmailSyncLimits.MAX_MESSAGES_PER_LABEL)
assertEquals(10_000, GmailSyncLimits.MAX_LABELS)
}
/**
* Ties Gmail's documented ceiling to today's actual architecture: [org.libremail.mail.ImapConnectionCache]
* (#125/#357) keeps at most ONE reused connection per account, plus `ImapClient.idle`'s own dedicated
* IDLE connection — 2 total, regardless of provider. Asserting that invariant against the real
* headroom-adjusted cap means a future change that grows per-account concurrency (e.g. a real
* connection pool) trips this test well before it could ever approach Gmail's actual ceiling.
*/
@Test
fun `today's architecture keeps concurrent connections per account well inside the background budget`() {
val knownConcurrentConnectionsPerAccount = 2 // one reused IMAP connection + one dedicated IDLE connection
assertTrue(knownConcurrentConnectionsPerAccount <= GmailSyncLimits.MAX_BACKGROUND_IMAP_CONNECTIONS)
}
@Test
fun `appliesTo is true for a gmail account, including the legacy googlemail host`() {
val gmail = MailProvider.GMAIL.createAccount("user@gmail.com")
assertTrue(GmailSyncLimits.appliesTo(gmail))
val legacyHost = gmail.copy(imap = gmail.imap.copy(host = "imap.googlemail.com"))
assertTrue(GmailSyncLimits.appliesTo(legacyHost))
}
@Test
fun `appliesTo is false for a non-gmail account`() {
assertFalse(GmailSyncLimits.appliesTo(MailProvider.YAHOO.createAccount("user@yahoo.com")))
assertFalse(GmailSyncLimits.appliesTo(MailProvider.ICLOUD.createAccount("user@icloud.com")))
assertFalse(GmailSyncLimits.appliesTo(MailProvider.AOL.createAccount("user@aol.com")))
assertFalse(GmailSyncLimits.appliesTo(Account.outlook("user@outlook.com")))
}
}
@@ -639,6 +639,57 @@ class MailBackfillerTest {
coVerify(atLeast = 1) { imapClient.fetchOlderThan(any(), any(), any(), any()) }
}
// --- issue #361: Gmail bandwidth-aware prefetch pacing ----------------------------------------
/**
* Once a Gmail account's tracked downloads for today reach the documented daily budget, backfill's
* body/attachment prefetch is deferred for the rest of the day — but header paging (the history
* itself) is untouched, the same "prefetch-only" gating shape as the low-battery (#89) and
* fetch-gate (#393) tests above.
*/
@Test
fun `gmail prefetch defers once the daily download budget is reached, but header paging is unaffected`() = runTest {
appendMessages(60)
seedForegroundWindow()
val gmailAccount = accountEntity.copy(imap = ServerConfigEmbedded("imap.gmail.com", 993, "SSL_TLS"))
val tracker = GmailBandwidthTracker().apply {
recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
}
backfiller(
AccountSettings("acct"),
fetchPolicy = FetchPolicy.ALWAYS,
bandwidthTracker = tracker,
account = gmailAccount,
).runBackfill()
assertEquals(60, distinctCachedUids().size, "header paging itself is not budget-gated")
coVerify(exactly = 0) { requireNotNull(lastMailRepository).prefetchMessage(any()) }
assertTrue(
logBuffer.snapshot().any { it.message.startsWith("prefetch deferred acct:") },
"a PII-free deferral breadcrumb is recorded",
)
}
/**
* The Gmail bandwidth budget is provider-scoped, not a blanket cap: an over-budget tracker entry
* for the same account id must not affect a non-Gmail account's prefetch.
*/
@Test
fun `a non-gmail account's prefetch is unaffected by an over-budget gmail bandwidth tracker`() = runTest {
appendMessages(60)
seedForegroundWindow()
val tracker = GmailBandwidthTracker().apply {
recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
}
// account defaults to the fixture's non-Gmail (127.0.0.1) host.
backfiller(AccountSettings("acct"), fetchPolicy = FetchPolicy.ALWAYS, bandwidthTracker = tracker)
.runBackfill()
coVerify(atLeast = 1) { requireNotNull(lastMailRepository).prefetchMessage(any()) }
}
// --- issue #329: AppLog breadcrumbs ---------------------------------------------------------
@Test
@@ -719,10 +770,12 @@ class MailBackfillerTest {
imapClient: ImapClient = client,
throttleGate: AccountThrottleGate = AccountThrottleGate(),
interactiveGate: InteractiveImapGate = InteractiveImapGate(),
bandwidthTracker: GmailBandwidthTracker = GmailBandwidthTracker(),
authGate: AuthThrottleGate = AuthThrottleGate(),
account: AccountEntity = accountEntity,
): MailBackfiller {
val accountDao = mockk<AccountDao>()
coEvery { accountDao.getAll() } returns listOf(accountEntity)
coEvery { accountDao.getAll() } returns listOf(account)
val messageDao = mockk<MessageDao>(relaxed = true)
coEvery { messageDao.insertNew(any()) } answers {
@@ -777,6 +830,7 @@ class MailBackfillerTest {
maintenanceGate = MailMaintenanceGate(),
throttleGate = throttleGate,
interactiveGate = interactiveGate,
bandwidthTracker = bandwidthTracker,
authGate = authGate,
).also {
lastMessageDao = messageDao
@@ -203,6 +203,7 @@ class MailMaintenanceGateTest {
maintenanceGate = gate,
throttleGate = AccountThrottleGate(),
interactiveGate = InteractiveImapGate(),
bandwidthTracker = GmailBandwidthTracker(),
authGate = AuthThrottleGate(),
)
}
@@ -332,6 +332,7 @@ class MailSyncConcurrencyTest {
notifier = mockk(relaxed = true),
mailRepository = mockk(relaxed = true),
throttleGate = AccountThrottleGate(),
bandwidthTracker = GmailBandwidthTracker(),
)
}
@@ -363,6 +364,7 @@ class MailSyncConcurrencyTest {
maintenanceGate = MailMaintenanceGate(),
throttleGate = AccountThrottleGate(),
interactiveGate = InteractiveImapGate(),
bandwidthTracker = GmailBandwidthTracker(),
authGate = AuthThrottleGate(),
)
}
@@ -95,9 +95,11 @@ class MailSyncerTest {
battery: BatteryStatus = BatteryStatus(percent = 100, isCharging = false),
fetched: List<FetchedMessage> = emptyList(),
throttleGate: AccountThrottleGate = AccountThrottleGate(),
bandwidthTracker: GmailBandwidthTracker = GmailBandwidthTracker(),
accountEntity: AccountEntity = account,
): MailSyncer {
val accountDao = mockk<AccountDao>()
coEvery { accountDao.getById("acct") } returns account
coEvery { accountDao.getById("acct") } returns accountEntity
val messageDao = mockk<MessageDao>(relaxed = true)
coEvery { messageDao.getSyncedIds(any(), any()) } returns emptyList()
coEvery { messageDao.getUnfetchedIds("acct", "INBOX") } returns listOf("acct:INBOX:1")
@@ -124,6 +126,7 @@ class MailSyncerTest {
notifier = mockk<MailNotifier>(relaxed = true),
mailRepository = mailRepository,
throttleGate = throttleGate,
bandwidthTracker = bandwidthTracker,
)
}
@@ -282,6 +285,43 @@ class MailSyncerTest {
coVerify { repo.prefetchMessage("acct:INBOX:1") }
}
// --- issue #361: Gmail bandwidth-aware prefetch pacing ----------------------------------------
@Test
fun `gmail prefetch defers once the daily download budget is reached, leaving the header sync untouched`() =
runTest {
val repo = mockk<MailRepository>(relaxed = true)
val gmailAccount = account.copy(imap = ServerConfigEmbedded("imap.gmail.com", 993, "SSL_TLS"))
val tracker = GmailBandwidthTracker().apply {
recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
}
val result = syncer(
FetchPolicy.ALWAYS,
repo,
accountEntity = gmailAccount,
bandwidthTracker = tracker,
).syncFolder("acct", "INBOX")
assertEquals(0, result.getOrNull()) // header sync still ran and succeeded
coVerify(exactly = 0) { repo.prefetchMessage(any()) }
assertTrue(logBuffer.snapshot().any { it.message.startsWith("prefetch deferred acct:") })
}
@Test
fun `a non-gmail account's prefetch is unaffected by an over-budget gmail bandwidth tracker`() = runTest {
val repo = mockk<MailRepository>()
coEvery { repo.prefetchMessage(any()) } returns Result.success(Unit)
val tracker = GmailBandwidthTracker().apply {
recordDownload("acct", GmailSyncLimits.DAILY_DOWNLOAD_BUDGET_BYTES)
}
// accountEntity defaults to the fixture's non-Gmail (imap.example.org) host.
syncer(FetchPolicy.ALWAYS, repo, bandwidthTracker = tracker).syncFolder("acct", "INBOX")
coVerify { repo.prefetchMessage("acct:INBOX:1") }
}
@Test
fun `notifies for new mail when both global and per-account notifications are enabled`() = runTest {
val notifier = mockk<MailNotifier>(relaxed = true)
@@ -340,6 +380,7 @@ class MailSyncerTest {
notifier = notifier,
mailRepository = mockk(relaxed = true),
throttleGate = AccountThrottleGate(),
bandwidthTracker = GmailBandwidthTracker(),
)
}
@@ -423,6 +464,7 @@ class MailSyncerTest {
notifier = mockk(relaxed = true),
mailRepository = mockk(relaxed = true),
throttleGate = AccountThrottleGate(),
bandwidthTracker = GmailBandwidthTracker(),
)
}
@@ -8,77 +8,65 @@ import kotlinx.coroutines.test.runTest
import org.junit.After
import org.junit.Before
import org.junit.Test
import org.libremail.data.sync.AccountThrottleGate
import org.libremail.domain.model.OutgoingMessage
import java.io.ByteArrayOutputStream
import org.libremail.mail.graph.FakeGraphHttpClient
import org.libremail.mail.graph.GraphResponse
import org.libremail.mail.graph.GraphThrottle
import org.libremail.mail.graph.GraphTransportException
import java.io.File
import java.io.IOException
import java.io.InputStream
import java.io.OutputStream
import java.io.RandomAccessFile
import java.net.HttpURLConnection
import java.net.URL
import java.net.URLConnection
import java.net.URLStreamHandler
import java.net.URLStreamHandlerFactory
import java.util.concurrent.atomic.AtomicReference
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertFalse
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
* Exercises the [GraphSender.send] transport path (its JSON payload builder is unit-tested separately
* in [GraphSenderTest]). The Graph endpoint is a fixed https URL the sender news up itself, so the
* test routes https through a process-wide [URLStreamHandlerFactory] to a per-test fake connection —
* no network, no production seam. The behaviour pinned: a 2xx succeeds; a non-2xx is a safe-to-retry
* rejection; a lost response is flagged [GraphSendException.mayHaveSent] (must NOT retry/fall back);
* a transmit failure is not; and the connection is always disconnected.
* Exercises [GraphSender.send] over a fake [org.libremail.mail.graph.GraphHttpClient] (its JSON payload
* builder is unit-tested separately in [GraphSenderTest], and the raw transport in
* [org.libremail.mail.graph.GraphHttpClientTest]). Pinned behaviour: a 2xx succeeds; a non-2xx is a
* safe-to-retry rejection; a lost response is flagged [GraphSendException.mayHaveSent] (must NOT
* retry/fall back); a transmit failure is not; an oversized attachment fails before any request; and a
* 429 is honored via [GraphThrottle] (retried after Retry-After) rather than immediately failing over.
*/
class GraphSenderSendTest {
@Before
fun setUp() {
// send() now breadcrumbs through AppLog on the oversized-attachment guard; android.util.Log is a
// no-op stub under plain JVM tests, so mock it (fully qualified, so this file never imports it).
// send() and the throttle/gate it composes with log via AppLog → android.util.Log, a throwing
// no-op stub under plain JVM tests; mock it fully-qualified so this file never imports it.
mockkStatic(android.util.Log::class)
every { android.util.Log.w(any<String>(), any<String>()) } returns 0
every { android.util.Log.i(any<String>(), any<String>()) } returns 0
every { android.util.Log.d(any<String>(), any<String>()) } returns 0
}
@After
fun tearDown() {
armed.set(null)
unmockkAll()
}
fun tearDown() = unmockkAll()
private val message =
OutgoingMessage(accountId = "outlook:me@x.com", to = "bob@example.org", subject = "Hi", body = "Body")
private fun arm(
status: Int = 202,
body: String = "",
failOutput: Boolean = false,
failResponse: Boolean = false,
): AtomicReference<FakeGraphConnection?> {
val last = AtomicReference<FakeGraphConnection?>(null)
armed.set { u -> FakeGraphConnection(u, status, body, failOutput, failResponse).also { last.set(it) } }
return last
private fun sender(client: FakeGraphHttpClient): Pair<GraphSender, AccountThrottleGate> {
val gate = AccountThrottleGate(nowMillis = { 0L }, random = { 0.0 })
return GraphSender(client, GraphThrottle(gate)) to gate
}
@Test
fun `a 2xx response completes the send`() = runTest {
val last = arm(status = 202)
val client = FakeGraphHttpClient.always(GraphResponse(status = 202, body = ""))
GraphSender().send("token", message)
sender(client).first.send("token", message)
assertTrue(last.get()!!.disconnected, "the connection must be disconnected when done")
assertEquals(1, client.callCount)
}
@Test
fun `a non-2xx response is a safe-to-retry rejection`() = runTest {
arm(status = 400, body = "{\"error\":\"bad request\"}")
val client = FakeGraphHttpClient.always(GraphResponse(status = 400, body = "{\"error\":\"bad request\"}"))
val ex = assertFailsWith<GraphSendException> { GraphSender().send("token", message) }
val ex = assertFailsWith<GraphSendException> { sender(client).first.send("token", message) }
assertFalse(ex.mayHaveSent, "an explicit rejection means Graph did not send")
assertTrue(ex.message!!.contains("HTTP 400"), ex.message!!)
@@ -86,25 +74,25 @@ class GraphSenderSendTest {
@Test
fun `a lost response is flagged as maybe-sent`() = runTest {
arm(failResponse = true)
val client = FakeGraphHttpClient { _, _ -> throw GraphTransportException("lost", mayHaveSent = true) }
val ex = assertFailsWith<GraphSendException> { GraphSender().send("token", message) }
val ex = assertFailsWith<GraphSendException> { sender(client).first.send("token", message) }
assertTrue(ex.mayHaveSent, "the request was fully sent, so it may already have delivered")
}
@Test
fun `a transmit failure is not maybe-sent`() = runTest {
arm(failOutput = true)
val client = FakeGraphHttpClient { _, _ -> throw GraphTransportException("no transmit", mayHaveSent = false) }
val ex = assertFailsWith<GraphSendException> { GraphSender().send("token", message) }
val ex = assertFailsWith<GraphSendException> { sender(client).first.send("token", message) }
assertFalse(ex.mayHaveSent, "the request never reached Graph, so a retry is safe")
}
@Test
fun `an oversized attachment fails safe-to-fall-back before opening a connection`() = runTest {
val last = arm() // armed, but the guard must trip before any connection is opened
val client = FakeGraphHttpClient.always(GraphResponse(status = 202, body = ""))
val big = File.createTempFile("graph-big", ".bin")
try {
// 4 MiB, over the 3 MiB per-file cap. setLength allocates the size without writing the bytes,
@@ -112,16 +100,43 @@ class GraphSenderSendTest {
RandomAccessFile(big, "rw").use { it.setLength(4L * 1024 * 1024) }
val ex = assertFailsWith<GraphSendException> {
GraphSender().send("token", message, listOf(SendableAttachment(big)))
sender(client).first.send("token", message, listOf(SendableAttachment(big)))
}
assertFalse(ex.mayHaveSent, "oversized never reached Graph, so SMTP fallback is safe")
assertNull(last.get(), "the guard must trip before any connection is opened")
assertEquals(0, client.callCount, "the guard must trip before any request is made")
} finally {
big.delete()
}
}
@Test
fun `a 429 is retried after Retry-After then succeeds`() = runTest {
val client = FakeGraphHttpClient.sequence(
GraphResponse(status = 429, body = "", retryAfterMillis = 1_000L),
GraphResponse(status = 202, body = ""),
)
sender(client).first.send("token", message)
assertEquals(2, client.callCount, "the 429 is honored once, then the retry sends")
}
@Test
fun `a persistent 429 surfaces as a safe-to-fall-back rejection`() = runTest {
val client = FakeGraphHttpClient.always(GraphResponse(status = 429, body = "quota", retryAfterMillis = 1_000L))
val (graphSender, gate) = sender(client)
val ex = assertFailsWith<GraphSendException> { graphSender.send("token", message) }
// One initial attempt + one honored retry (SEND_MAX_RETRIES = 1), then fall back.
assertEquals(2, client.callCount)
assertFalse(ex.mayHaveSent, "a 429 means Graph did not send, so SMTP fallback is safe")
assertTrue(ex.message!!.contains("HTTP 429"), ex.message!!)
// The throttle is recorded so the account's background IMAP work backs off too (#360 composition).
assertTrue(gate.isThrottled(message.accountId))
}
@Test
fun `GraphSendException carries its message, flag and cause`() {
val cause = IOException("boom")
@@ -131,46 +146,4 @@ class GraphSenderSendTest {
assertTrue(ex.mayHaveSent)
assertEquals(cause, ex.cause)
}
private class FakeGraphConnection(
url: URL,
private val status: Int,
private val body: String,
private val failOutput: Boolean,
private val failResponse: Boolean,
) : HttpURLConnection(url) {
var disconnected = false
override fun connect() = Unit
override fun disconnect() {
disconnected = true
}
override fun usingProxy() = false
override fun getOutputStream(): OutputStream =
if (failOutput) throw IOException("cannot transmit") else ByteArrayOutputStream()
override fun getResponseCode(): Int = if (failResponse) throw IOException("no response") else status
override fun getInputStream(): InputStream = body.byteInputStream()
override fun getErrorStream(): InputStream? = if (body.isEmpty()) null else body.byteInputStream()
}
companion object {
private val armed = AtomicReference<((URL) -> HttpURLConnection)?>(null)
// Set once per JVM: route https opens to whatever the running test armed. Only this test opens
// https in the unit-test JVM (ReportUploadWorker's endpoint is blank and never opened), so a
// permanent https handler is safe here.
init {
URL.setURLStreamHandlerFactory(
URLStreamHandlerFactory { protocol ->
if (protocol == "https") {
object : URLStreamHandler() {
override fun openConnection(u: URL): URLConnection =
armed.get()?.invoke(u) ?: throw IOException("no fake connection armed")
}
} else {
null
}
},
)
}
}
}
@@ -0,0 +1,32 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
/**
* A no-network [GraphHttpClient] for unit tests: it records every [GraphRequest] it is asked to execute
* and returns whatever the supplied [responder] produces for that call (the request plus its 0-based
* call index), so a test can script a 429-then-200 sequence, echo a `$batch` body, or throw a
* [GraphTransportException] — all without opening a socket.
*/
class FakeGraphHttpClient(private val responder: (request: GraphRequest, callIndex: Int) -> GraphResponse) :
GraphHttpClient() {
/** Every request passed to [execute], in call order — the assertion surface for call count / headers. */
val requests = mutableListOf<GraphRequest>()
val callCount: Int get() = requests.size
override suspend fun execute(request: GraphRequest): GraphResponse {
val index = requests.size
requests += request
return responder(request, index)
}
companion object {
/** Replays [responses] in order; the last one repeats if execute is called more times than supplied. */
fun sequence(vararg responses: GraphResponse): FakeGraphHttpClient =
FakeGraphHttpClient { _, index -> responses[minOf(index, responses.size - 1)] }
/** Always answers with the same [response]. */
fun always(response: GraphResponse): FakeGraphHttpClient = FakeGraphHttpClient { _, _ -> response }
}
}
@@ -0,0 +1,151 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import io.mockk.every
import io.mockk.mockkStatic
import io.mockk.unmockkAll
import kotlinx.coroutines.test.runTest
import org.json.JSONArray
import org.json.JSONObject
import org.junit.After
import org.junit.Before
import org.junit.Test
import org.libremail.data.sync.AccountThrottleGate
import kotlin.test.assertEquals
import kotlin.test.assertTrue
/**
* [GraphBatch] must collapse many operations into `ceil(N / 20)` `$batch` HTTP calls (the core
* call-volume lever of issue #364), correlate every sub-response back to its request id, make no call
* for an empty list, and feed a per-op 429 back into the shared #360 backoff gate. Plus pure coverage
* of the payload builder and response parser.
*/
class GraphBatchTest {
private val accountId = "outlook:user@example.org"
@Before
fun setUp() {
mockkStatic(android.util.Log::class)
every { android.util.Log.d(any<String>(), any<String>()) } returns 0
every { android.util.Log.w(any<String>(), any<String>()) } returns 0
}
@After
fun tearDown() = unmockkAll()
private fun gate() = AccountThrottleGate(nowMillis = { 0L }, random = { 0.0 })
/** A fake that echoes one 200 sub-response per sub-request in the batch payload it receives. */
private fun echoClient() = FakeGraphHttpClient { request, _ ->
val requested = JSONObject(String(request.body!!, Charsets.UTF_8)).getJSONArray("requests")
val responses = JSONArray()
for (index in 0 until requested.length()) {
responses.put(JSONObject().put("id", requested.getJSONObject(index).getString("id")).put("status", 200))
}
GraphResponse(status = 200, body = JSONObject().put("responses", responses).toString())
}
private fun subRequests(count: Int) =
(1..count).map { GraphSubRequest(id = it.toString(), method = "GET", url = "/me/messages/$it") }
@Test
fun `twenty-five operations collapse to two batch calls`() = runTest {
val client = echoClient()
val responses = GraphBatch(GraphThrottle(gate())).execute(accountId, client, subRequests(25))
assertEquals(2, client.callCount, "25 ops / 20-per-batch = 2 HTTP round-trips instead of 25")
assertEquals(25, responses.size)
assertEquals((1..25).map { it.toString() }, responses.map { it.id }, "order + id correlation preserved")
assertTrue(responses.all { it.status == 200 })
}
@Test
fun `exactly twenty operations are a single batch call`() = runTest {
val client = echoClient()
GraphBatch(GraphThrottle(gate())).execute(accountId, client, subRequests(20))
assertEquals(1, client.callCount)
}
@Test
fun `an empty operation list makes no HTTP call`() = runTest {
val client = echoClient()
val responses = GraphBatch(GraphThrottle(gate())).execute(accountId, client, emptyList())
assertTrue(responses.isEmpty())
assertEquals(0, client.callCount)
}
@Test
fun `a per-op 429 inside the envelope backs the account off`() = runTest {
val gate = gate()
val client = FakeGraphHttpClient.always(
GraphResponse(
status = 200,
body = JSONObject().put(
"responses",
JSONArray().put(
JSONObject().put("id", "1").put("status", 429)
.put("headers", JSONObject().put("Retry-After", "60")),
),
).toString(),
),
)
val responses = GraphBatch(GraphThrottle(gate)).execute(accountId, client, subRequests(1))
assertEquals(429, responses.single().status)
assertEquals(60_000L, responses.single().retryAfterMillis)
assertTrue(gate.isThrottled(accountId), "a throttled sub-response cools the account down like a top-level 429")
}
@Test
fun `buildBatchPayload carries id, method, url, headers and body`() {
val payload = JSONObject(
buildBatchPayload(
listOf(
GraphSubRequest(
id = "a",
method = "POST",
url = "/me/sendMail",
headers = mapOf("Content-Type" to "application/json"),
body = JSONObject().put("saveToSentItems", true),
),
),
),
)
val request = payload.getJSONArray("requests").getJSONObject(0)
assertEquals("a", request.getString("id"))
assertEquals("POST", request.getString("method"))
assertEquals("/me/sendMail", request.getString("url"))
assertEquals("application/json", request.getJSONObject("headers").getString("Content-Type"))
assertTrue(request.getJSONObject("body").getBoolean("saveToSentItems"))
}
@Test
fun `parseBatchResponses correlates id, status and body`() {
val parsed = parseBatchResponses(
JSONObject().put(
"responses",
JSONArray()
.put(JSONObject().put("id", "1").put("status", 200).put("body", JSONObject().put("ok", true)))
.put(JSONObject().put("id", "2").put("status", 404)),
).toString(),
)
assertEquals(listOf("1", "2"), parsed.map { it.id })
assertEquals(200, parsed[0].status)
assertTrue(parsed[0].body!!.getBoolean("ok"))
assertEquals(404, parsed[1].status)
}
@Test
fun `parseBatchResponses tolerates a blank or malformed body`() {
assertTrue(parseBatchResponses("").isEmpty())
assertTrue(parseBatchResponses("not json").isEmpty())
assertTrue(parseBatchResponses("{}").isEmpty())
}
}
@@ -0,0 +1,207 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import kotlinx.coroutines.test.runTest
import org.junit.After
import org.junit.Test
import java.io.ByteArrayOutputStream
import java.io.IOException
import java.io.InputStream
import java.io.OutputStream
import java.net.HttpURLConnection
import java.net.URL
import java.net.URLConnection
import java.net.URLStreamHandler
import java.net.URLStreamHandlerFactory
import java.util.concurrent.atomic.AtomicReference
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertFalse
import kotlin.test.assertNull
import kotlin.test.assertTrue
/**
* Exercises the real [GraphHttpClient] transport and the pure [parseRetryAfterMillis] helper. The Graph
* endpoints are fixed https URLs the client opens itself, so — like the old GraphSender transport test —
* https is routed through a process-wide [URLStreamHandlerFactory] to a per-test fake connection: no
* network, no production seam. Pinned behaviour: any answered status (2xx or not) returns a
* [GraphResponse] with its body and parsed `Retry-After`; a transmit failure and a lost response each
* throw [GraphTransportException] with the right [GraphTransportException.mayHaveSent] flag; and the
* connection is always disconnected.
*/
class GraphHttpClientTest {
@After
fun tearDown() {
armed.set(null)
}
private fun arm(
status: Int = 202,
body: String = "",
retryAfter: String? = null,
failOutput: Boolean = false,
failResponse: Boolean = false,
): AtomicReference<FakeConnection?> {
val last = AtomicReference<FakeConnection?>(null)
armed.set { u -> FakeConnection(u, status, body, retryAfter, failOutput, failResponse).also { last.set(it) } }
return last
}
private fun request(body: ByteArray? = ByteArray(0)) =
GraphRequest("POST", "https://graph.microsoft.com/v1.0/me/sendMail", mapOf("Authorization" to "Bearer t"), body)
@Test
fun `a 2xx response returns its status and body and disconnects`() = runTest {
val last = arm(status = 200, body = "{\"ok\":true}")
val response = GraphHttpClient().execute(request())
assertEquals(200, response.status)
assertEquals("{\"ok\":true}", response.body)
assertTrue(last.get()!!.disconnected)
}
@Test
fun `a non-2xx response returns its error body and disconnects`() = runTest {
val last = arm(status = 400, body = "{\"error\":\"bad\"}")
val response = GraphHttpClient().execute(request())
assertEquals(400, response.status)
assertEquals("{\"error\":\"bad\"}", response.body)
assertTrue(last.get()!!.disconnected)
}
@Test
fun `a Retry-After header is parsed into milliseconds`() = runTest {
arm(status = 429, body = "", retryAfter = "30")
val response = GraphHttpClient().execute(request())
assertEquals(429, response.status)
assertEquals(30_000L, response.retryAfterMillis)
}
@Test
fun `no Retry-After header yields a null wait`() = runTest {
arm(status = 429, body = "")
assertNull(GraphHttpClient().execute(request()).retryAfterMillis)
}
@Test
fun `a transmit failure is thrown as not-maybe-sent`() = runTest {
val last = arm(failOutput = true)
val ex = assertFailsWith<GraphTransportException> { GraphHttpClient().execute(request()) }
assertFalse(ex.mayHaveSent, "the body never left the device, so a retry is safe")
assertTrue(last.get()!!.disconnected)
}
@Test
fun `a lost response is thrown as maybe-sent`() = runTest {
arm(failResponse = true)
val ex = assertFailsWith<GraphTransportException> { GraphHttpClient().execute(request()) }
assertTrue(ex.mayHaveSent, "the request was fully sent, so it may already have been acted on")
}
@Test
fun `a request without a body is not transmitted but still reads the response`() = runTest {
val last = arm(status = 200, body = "ok")
val response = GraphHttpClient().execute(request(body = null))
assertEquals(200, response.status)
assertFalse(last.get()!!.wroteBody, "a null body must not open the output stream")
}
// --- parseRetryAfterMillis (pure) -----------------------------------------------------------
@Test
fun `parseRetryAfterMillis reads delta-seconds`() {
assertEquals(120_000L, parseRetryAfterMillis("120", nowMillis = 0L))
}
@Test
fun `parseRetryAfterMillis reads an HTTP-date relative to now`() {
// 60 seconds after the epoch, evaluated as of the epoch → 60_000 ms.
assertEquals(60_000L, parseRetryAfterMillis("Thu, 01 Jan 1970 00:01:00 GMT", nowMillis = 0L))
}
@Test
fun `parseRetryAfterMillis clamps a past date to zero`() {
assertEquals(0L, parseRetryAfterMillis("Thu, 01 Jan 1970 00:00:00 GMT", nowMillis = 120_000L))
}
@Test
fun `parseRetryAfterMillis returns null for blank or garbage`() {
assertNull(parseRetryAfterMillis(null, 0L))
assertNull(parseRetryAfterMillis(" ", 0L))
assertNull(parseRetryAfterMillis("soon", 0L))
}
@Test
fun `parseRetryAfterMillis clamps a negative delta to zero`() {
assertEquals(0L, parseRetryAfterMillis("-5", 0L))
}
@Test
fun `isHttpSuccess spans the 2xx range only`() {
assertTrue(isHttpSuccess(200))
assertTrue(isHttpSuccess(202))
assertTrue(isHttpSuccess(299))
assertFalse(isHttpSuccess(199))
assertFalse(isHttpSuccess(300))
assertFalse(isHttpSuccess(429))
}
private class FakeConnection(
url: URL,
private val status: Int,
private val body: String,
private val retryAfter: String?,
private val failOutput: Boolean,
private val failResponse: Boolean,
) : HttpURLConnection(url) {
var disconnected = false
var wroteBody = false
override fun connect() = Unit
override fun disconnect() {
disconnected = true
}
override fun usingProxy() = false
override fun getOutputStream(): OutputStream =
if (failOutput) throw IOException("cannot transmit") else ByteArrayOutputStream().also { wroteBody = true }
override fun getResponseCode(): Int = if (failResponse) throw IOException("no response") else status
override fun getInputStream(): InputStream = body.byteInputStream()
override fun getErrorStream(): InputStream? = if (body.isEmpty()) null else body.byteInputStream()
override fun getHeaderField(name: String): String? =
if (name.equals("Retry-After", ignoreCase = true)) retryAfter else super.getHeaderField(name)
}
companion object {
private val armed = AtomicReference<((URL) -> HttpURLConnection)?>(null)
// Set once per JVM: route https opens to whatever the running test armed. Only this test opens
// real https in the unit-test JVM (every other Graph test uses FakeGraphHttpClient), so a
// permanent https handler is safe here.
init {
URL.setURLStreamHandlerFactory(
URLStreamHandlerFactory { protocol ->
if (protocol == "https") {
object : URLStreamHandler() {
override fun openConnection(u: URL): URLConnection =
armed.get()?.invoke(u) ?: throw IOException("no fake connection armed")
}
} else {
null
}
},
)
}
}
}
@@ -0,0 +1,136 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import io.mockk.every
import io.mockk.mockkStatic
import io.mockk.unmockkAll
import kotlinx.coroutines.ExperimentalCoroutinesApi
import kotlinx.coroutines.test.runTest
import org.junit.After
import org.junit.Before
import org.junit.Test
import org.libremail.data.sync.AccountThrottleGate
import org.libremail.data.sync.ThrottleKind
import org.libremail.data.sync.ThrottleSignal
import org.libremail.reporting.AppLog
import org.libremail.reporting.RingLogBuffer
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertFalse
import kotlin.test.assertTrue
/**
* [GraphThrottle] must honor a Graph 429/503 (retry after the server's `Retry-After`, composed with the
* shared #360 backoff gate), bound its retries, clear the gate on success, leave a non-throttle error
* untouched, not retry a transport failure, and log only PII-free breadcrumbs — all under coroutines-test
* virtual time, with no real sleeps.
*/
@OptIn(ExperimentalCoroutinesApi::class)
class GraphThrottleTest {
private val logBuffer = RingLogBuffer()
private val accountId = "outlook:user@example.org"
@Before
fun setUp() {
// GraphThrottle (and the gate it drives) log via AppLog → android.util.Log, a throwing no-op stub
// under plain JVM tests; mock it fully-qualified so this file never imports android.util.Log.
mockkStatic(android.util.Log::class)
every { android.util.Log.w(any<String>(), any<String>()) } returns 0
every { android.util.Log.i(any<String>(), any<String>()) } returns 0
every { android.util.Log.d(any<String>(), any<String>()) } returns 0
AppLog.install(logBuffer)
}
@After
fun tearDown() = unmockkAll()
private fun ok(body: String = "") = GraphResponse(status = 202, body = body)
private fun throttled(retryAfterMillis: Long?) =
GraphResponse(status = 429, body = "", retryAfterMillis = retryAfterMillis)
private fun request() = GraphRequest("POST", "https://graph.microsoft.com/v1.0/me/sendMail")
@Test
fun `a 429 with Retry-After is honored then the retry succeeds`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
val throttle = GraphThrottle(gate)
// Retry-After 120s exceeds the base exponential half (15s), so the honored wait is exactly 120s.
val client = FakeGraphHttpClient.sequence(throttled(retryAfterMillis = 120_000L), ok(body = "sent"))
val response = throttle.execute(accountId, client, request())
assertEquals(202, response.status)
assertEquals("sent", response.body)
assertEquals(2, client.callCount, "it must retry exactly once after the 429")
assertEquals(120_000L, testScheduler.currentTime, "it must wait the server's Retry-After before retrying")
assertFalse(gate.isThrottled(accountId), "a successful retry clears the account's backoff")
}
@Test
fun `a 429 records the account against the shared backoff gate`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
val client = FakeGraphHttpClient.always(throttled(null))
val response = GraphThrottle(gate).execute(accountId, client, request(), maxRetries = 0)
assertEquals(429, response.status)
assertTrue(gate.isThrottled(accountId), "an unrecovered 429 leaves the account backed off for other paths")
}
@Test
fun `retries are bounded and the throttled response is finally returned`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
val client = FakeGraphHttpClient.always(throttled(retryAfterMillis = 1_000L))
val response = GraphThrottle(gate).execute(accountId, client, request(), maxRetries = 2)
assertEquals(429, response.status)
assertEquals(3, client.callCount, "one initial attempt plus two bounded retries")
}
@Test
fun `a success clears a pre-existing backoff`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
gate.onThrottle(accountId, ThrottleSignal(ThrottleKind.RATE_LIMIT))
assertTrue(gate.isThrottled(accountId))
GraphThrottle(gate).execute(accountId, FakeGraphHttpClient.always(ok()), request())
assertFalse(gate.isThrottled(accountId), "a 2xx clears the account's backoff")
}
@Test
fun `a non-throttle error is returned as-is and does not clear backoff`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
gate.onThrottle(accountId, ThrottleSignal(ThrottleKind.RATE_LIMIT))
val client = FakeGraphHttpClient.always(GraphResponse(status = 401, body = "unauthorized"))
val response = GraphThrottle(gate).execute(accountId, client, request())
assertEquals(401, response.status)
assertEquals(1, client.callCount, "a 4xx that is not a throttle must not be retried")
assertTrue(gate.isThrottled(accountId), "a non-2xx must not clear an existing backoff")
}
@Test
fun `a transport failure propagates without a retry`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
val client = FakeGraphHttpClient { _, _ -> throw GraphTransportException("lost", mayHaveSent = true) }
assertFailsWith<GraphTransportException> {
GraphThrottle(gate).execute(accountId, client, request())
}
assertEquals(1, client.callCount, "a no-response transport error is the caller's to judge, never retried here")
}
@Test
fun `throttle breadcrumbs are PII-free`() = runTest {
val gate = AccountThrottleGate(nowMillis = { testScheduler.currentTime }, random = { 0.0 })
GraphThrottle(gate).execute(accountId, FakeGraphHttpClient.always(throttled(null)), request(), maxRetries = 0)
val messages = logBuffer.snapshot().map { it.message }
assertTrue(messages.any { it.startsWith("graph throttled outlook:") }, messages.toString())
messages.forEach { assertFalse(it.contains("user@example.org"), it) }
}
}
@@ -0,0 +1,140 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.mail.graph
import io.mockk.every
import io.mockk.mockkStatic
import io.mockk.unmockkAll
import kotlinx.coroutines.test.runTest
import org.json.JSONObject
import org.junit.After
import org.junit.Before
import org.junit.Test
import org.libremail.data.sync.AccountThrottleGate
import kotlin.test.assertEquals
import kotlin.test.assertFailsWith
import kotlin.test.assertFalse
import kotlin.test.assertTrue
/**
* [GraphUploadSession] must decide when a chunked session is required (the ~4 MB one-shot ceiling), then
* create the session and PUT the content in contiguous 320 KiB-multiple ranges that reassemble to the
* original bytes (issue #364's chunked-upload lever). It must also short-circuit a failed session-create
* and reject invalid inputs.
*/
class GraphUploadSessionTest {
private val accountId = "outlook:user@example.org"
private val chunkMultiple = 320 * 1024
private val createUrl = "https://graph.microsoft.com/v1.0/me/messages/1/attachments/createUploadSession"
@Before
fun setUp() {
mockkStatic(android.util.Log::class)
every { android.util.Log.i(any<String>(), any<String>()) } returns 0
every { android.util.Log.w(any<String>(), any<String>()) } returns 0
}
@After
fun tearDown() = unmockkAll()
private fun uploadSession(): GraphUploadSession {
val gate = AccountThrottleGate(nowMillis = { 0L }, random = { 0.0 })
return GraphUploadSession(GraphThrottle(gate))
}
/** POST (create session) → the upload URL; every PUT (chunk) → 202 Accepted. */
private fun sessionClient() = FakeGraphHttpClient { request, _ ->
if (request.method == "POST") {
GraphResponse(status = 200, body = JSONObject().put("uploadUrl", "https://upload.example/1").toString())
} else {
GraphResponse(status = 202, body = "")
}
}
@Test
fun `requiresUploadSession is true only above the one-shot ceiling`() {
val session = uploadSession()
val fourMb = 4L * 1024 * 1024
assertFalse(session.requiresUploadSession(fourMb), "exactly 4 MB still fits a one-shot upload")
assertTrue(session.requiresUploadSession(fourMb + 1), "over 4 MB needs a chunked session")
}
@Test
fun `an over-threshold attachment uploads in contiguous ranged chunks`() = runTest {
val session = uploadSession()
val client = sessionClient()
val total = 4 * 1024 * 1024 + 500 // just over the 4 MB one-shot ceiling
val content = ByteArray(total) { (it % 251).toByte() }
assertTrue(session.requiresUploadSession(total.toLong()), "the payload is over-threshold")
val createBody = JSONObject().put("AttachmentItem", JSONObject())
val response = session.upload(accountId, client, createUrl, createBody, content, chunkSize = chunkMultiple)
assertEquals(202, response.status)
// First call creates the session; the rest are the chunk PUTs.
val create = client.requests.first()
assertEquals("POST", create.method)
assertEquals(createUrl, create.url)
val puts = client.requests.drop(1)
val expectedChunks = (total + chunkMultiple - 1) / chunkMultiple
assertEquals(expectedChunks, puts.size, "content is split into ceil(total / chunkSize) PUTs")
assertTrue(puts.all { it.method == "PUT" && it.url == "https://upload.example/1" })
// The ranges must be contiguous, start at 0, end at total-1, and reassemble to the original bytes.
val reassembled = ByteArray(total)
var expectedStart = 0
puts.forEach { put ->
val range = put.headers.getValue("Content-Range").removePrefix("bytes ").substringBefore('/')
val start = range.substringBefore('-').toInt()
val end = range.substringAfter('-').toInt()
val chunkBody = put.body!!
assertEquals(expectedStart, start, "chunk ranges must be contiguous")
chunkBody.copyInto(reassembled, destinationOffset = start)
assertEquals(end - start + 1, chunkBody.size, "Content-Range length must match the body")
expectedStart = end + 1
}
assertEquals(total, expectedStart, "the ranges must cover the whole payload")
assertTrue(content.contentEquals(reassembled), "the reassembled chunks must equal the original content")
}
@Test
fun `a small payload still uploads as a single chunk`() = runTest {
val client = sessionClient()
uploadSession().upload(accountId, client, createUrl, JSONObject(), ByteArray(10) { 1 }, chunkMultiple)
assertEquals(1, client.requests.count { it.method == "PUT" }, "content under one chunk is a single PUT")
}
@Test
fun `a failed session-create returns without uploading`() = runTest {
val client = FakeGraphHttpClient.always(GraphResponse(status = 500, body = "boom"))
val response = uploadSession().upload(accountId, client, createUrl, JSONObject(), ByteArray(10) { 1 })
assertEquals(500, response.status)
assertEquals(1, client.callCount, "no chunk is uploaded when the session cannot be created")
}
@Test
fun `empty content is rejected`() = runTest {
assertFailsWith<IllegalArgumentException> {
uploadSession().upload(accountId, sessionClient(), createUrl, JSONObject(), ByteArray(0))
}
}
@Test
fun `a chunk size that is not a 320 KiB multiple is rejected`() = runTest {
assertFailsWith<IllegalArgumentException> {
// 1000 is not a multiple of 320 KiB.
uploadSession().upload(
accountId,
sessionClient(),
createUrl,
JSONObject(),
ByteArray(10) { 1 },
chunkSize = 1000,
)
}
}
}
+3
View File
@@ -75,6 +75,9 @@ style:
- '**/data/sync/AccountThrottleGateTest.kt'
# #356 backfill pacer: cooldown/cap/skip breadcrumbs through AppLog, so this suite mockkStatic(Log) too.
- '**/data/sync/BackfillPacerTest.kt'
# #361 Gmail bandwidth tracker: the daily-budget-crossing breadcrumb goes through AppLog, so this
# suite mockkStatic(Log) too.
- '**/data/sync/GmailBandwidthTrackerTest.kt'
# Reader-path perf logging (issue #358): the repository's openMessage and the reader ViewModel
# log via AppLog, so their unit tests mockkStatic(Log) too.
- '**/data/repository/MailRepositoryImplTest.kt'