diff --git a/app/src/androidTest/kotlin/org/libremail/mail/graph/GraphThrottleInstrumentedTest.kt b/app/src/androidTest/kotlin/org/libremail/mail/graph/GraphThrottleInstrumentedTest.kt new file mode 100644 index 0000000..56b1c0e --- /dev/null +++ b/app/src/androidTest/kotlin/org/libremail/mail/graph/GraphThrottleInstrumentedTest.kt @@ -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() + 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") +} diff --git a/app/src/main/kotlin/org/libremail/mail/GraphSender.kt b/app/src/main/kotlin/org/libremail/mail/GraphSender.kt index bbe3762..e053ad7 100644 --- a/app/src/main/kotlin/org/libremail/mail/GraphSender.kt +++ b/app/src/main/kotlin/org/libremail/mail/GraphSender.kt @@ -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 } } diff --git a/app/src/main/kotlin/org/libremail/mail/graph/GraphBatch.kt b/app/src/main/kotlin/org/libremail/mail/graph/GraphBatch.kt new file mode 100644 index 0000000..e696963 --- /dev/null +++ b/app/src/main/kotlin/org/libremail/mail/graph/GraphBatch.kt @@ -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 = 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, + ): List { + if (requests.isEmpty()) return emptyList() + val chunks = requests.chunked(MAX_BATCH_OPERATIONS) + val responses = ArrayList(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): 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 { + if (body.isBlank()) return emptyList() + val responses = runCatching { JSONObject(body).optJSONArray("responses") }.getOrNull() ?: return emptyList() + val out = ArrayList(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()) +} diff --git a/app/src/main/kotlin/org/libremail/mail/graph/GraphHttp.kt b/app/src/main/kotlin/org/libremail/mail/graph/GraphHttp.kt new file mode 100644 index 0000000..c99a348 --- /dev/null +++ b/app/src/main/kotlin/org/libremail/mail/graph/GraphHttp.kt @@ -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 = 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() +} diff --git a/app/src/main/kotlin/org/libremail/mail/graph/GraphThrottle.kt b/app/src/main/kotlin/org/libremail/mail/graph/GraphThrottle.kt new file mode 100644 index 0000000..4abe363 --- /dev/null +++ b/app/src/main/kotlin/org/libremail/mail/graph/GraphThrottle.kt @@ -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 + } +} diff --git a/app/src/main/kotlin/org/libremail/mail/graph/GraphUploadSession.kt b/app/src/main/kotlin/org/libremail/mail/graph/GraphUploadSession.kt new file mode 100644 index 0000000..163c0e4 --- /dev/null +++ b/app/src/main/kotlin/org/libremail/mail/graph/GraphUploadSession.kt @@ -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" + } +} diff --git a/app/src/test/kotlin/org/libremail/mail/GraphSenderSendTest.kt b/app/src/test/kotlin/org/libremail/mail/GraphSenderSendTest.kt index c4c02ea..0994037 100644 --- a/app/src/test/kotlin/org/libremail/mail/GraphSenderSendTest.kt +++ b/app/src/test/kotlin/org/libremail/mail/GraphSenderSendTest.kt @@ -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(), any()) } returns 0 + every { android.util.Log.i(any(), any()) } returns 0 + every { android.util.Log.d(any(), any()) } 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 { - val last = AtomicReference(null) - armed.set { u -> FakeGraphConnection(u, status, body, failOutput, failResponse).also { last.set(it) } } - return last + private fun sender(client: FakeGraphHttpClient): Pair { + 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 { GraphSender().send("token", message) } + val ex = assertFailsWith { 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 { GraphSender().send("token", message) } + val ex = assertFailsWith { 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 { GraphSender().send("token", message) } + val ex = assertFailsWith { 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 { - 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 { 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 - } - }, - ) - } - } } diff --git a/app/src/test/kotlin/org/libremail/mail/graph/FakeGraphHttpClient.kt b/app/src/test/kotlin/org/libremail/mail/graph/FakeGraphHttpClient.kt new file mode 100644 index 0000000..08b5d58 --- /dev/null +++ b/app/src/test/kotlin/org/libremail/mail/graph/FakeGraphHttpClient.kt @@ -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() + + 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 } + } +} diff --git a/app/src/test/kotlin/org/libremail/mail/graph/GraphBatchTest.kt b/app/src/test/kotlin/org/libremail/mail/graph/GraphBatchTest.kt new file mode 100644 index 0000000..81dc911 --- /dev/null +++ b/app/src/test/kotlin/org/libremail/mail/graph/GraphBatchTest.kt @@ -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(), any()) } returns 0 + every { android.util.Log.w(any(), any()) } 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()) + } +} diff --git a/app/src/test/kotlin/org/libremail/mail/graph/GraphHttpClientTest.kt b/app/src/test/kotlin/org/libremail/mail/graph/GraphHttpClientTest.kt new file mode 100644 index 0000000..a74dd70 --- /dev/null +++ b/app/src/test/kotlin/org/libremail/mail/graph/GraphHttpClientTest.kt @@ -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 { + val last = AtomicReference(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 { 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 { 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 + } + }, + ) + } + } +} diff --git a/app/src/test/kotlin/org/libremail/mail/graph/GraphThrottleTest.kt b/app/src/test/kotlin/org/libremail/mail/graph/GraphThrottleTest.kt new file mode 100644 index 0000000..adedc9d --- /dev/null +++ b/app/src/test/kotlin/org/libremail/mail/graph/GraphThrottleTest.kt @@ -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(), any()) } returns 0 + every { android.util.Log.i(any(), any()) } returns 0 + every { android.util.Log.d(any(), any()) } 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 { + 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) } + } +} diff --git a/app/src/test/kotlin/org/libremail/mail/graph/GraphUploadSessionTest.kt b/app/src/test/kotlin/org/libremail/mail/graph/GraphUploadSessionTest.kt new file mode 100644 index 0000000..9ddcbc8 --- /dev/null +++ b/app/src/test/kotlin/org/libremail/mail/graph/GraphUploadSessionTest.kt @@ -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(), any()) } returns 0 + every { android.util.Log.w(any(), any()) } 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 { + uploadSession().upload(accountId, sessionClient(), createUrl, JSONObject(), ByteArray(0)) + } + } + + @Test + fun `a chunk size that is not a 320 KiB multiple is rejected`() = runTest { + assertFailsWith { + // 1000 is not a multiple of 320 KiB. + uploadSession().upload( + accountId, + sessionClient(), + createUrl, + JSONObject(), + ByteArray(10) { 1 }, + chunkSize = 1000, + ) + } + } +}