+71
@@ -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")
|
||||
}
|
||||
+46
-37
@@ -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
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+42
-28
@@ -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
|
||||
}
|
||||
@@ -59,6 +59,7 @@ class MailBackfiller @Inject constructor(
|
||||
private val maintenanceGate: MailMaintenanceGate,
|
||||
private val throttleGate: AccountThrottleGate,
|
||||
private val interactiveGate: InteractiveImapGate,
|
||||
private val bandwidthTracker: GmailBandwidthTracker,
|
||||
) {
|
||||
/** One folder's slice outcome: pages fetched, and whether an immediate follow-up slice has work to do. */
|
||||
private data class FolderResult(val batches: Int, val moreWork: Boolean)
|
||||
@@ -214,7 +215,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)
|
||||
}
|
||||
@@ -293,7 +294,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
|
||||
@@ -308,6 +309,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")))
|
||||
}
|
||||
}
|
||||
@@ -606,6 +606,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
|
||||
@@ -686,9 +737,11 @@ class MailBackfillerTest {
|
||||
imapClient: ImapClient = client,
|
||||
throttleGate: AccountThrottleGate = AccountThrottleGate(),
|
||||
interactiveGate: InteractiveImapGate = InteractiveImapGate(),
|
||||
bandwidthTracker: GmailBandwidthTracker = GmailBandwidthTracker(),
|
||||
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 {
|
||||
@@ -743,6 +796,7 @@ class MailBackfillerTest {
|
||||
maintenanceGate = MailMaintenanceGate(),
|
||||
throttleGate = throttleGate,
|
||||
interactiveGate = interactiveGate,
|
||||
bandwidthTracker = bandwidthTracker,
|
||||
).also {
|
||||
lastMessageDao = messageDao
|
||||
lastMailRepository = mailRepository
|
||||
|
||||
@@ -202,6 +202,7 @@ class MailMaintenanceGateTest {
|
||||
maintenanceGate = gate,
|
||||
throttleGate = AccountThrottleGate(),
|
||||
interactiveGate = InteractiveImapGate(),
|
||||
bandwidthTracker = GmailBandwidthTracker(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -331,6 +331,7 @@ class MailSyncConcurrencyTest {
|
||||
notifier = mockk(relaxed = true),
|
||||
mailRepository = mockk(relaxed = true),
|
||||
throttleGate = AccountThrottleGate(),
|
||||
bandwidthTracker = GmailBandwidthTracker(),
|
||||
)
|
||||
}
|
||||
|
||||
@@ -362,6 +363,7 @@ class MailSyncConcurrencyTest {
|
||||
maintenanceGate = MailMaintenanceGate(),
|
||||
throttleGate = AccountThrottleGate(),
|
||||
interactiveGate = InteractiveImapGate(),
|
||||
bandwidthTracker = GmailBandwidthTracker(),
|
||||
)
|
||||
}
|
||||
|
||||
|
||||
@@ -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,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -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'
|
||||
|
||||
Reference in New Issue
Block a user