WIP: merge queue: checking main (4c3937e) and #470 together #475

Closed
mergify[bot] wants to merge 2 commits from mergify/merge-queue/9781139bcc into main
12 changed files with 1358 additions and 129 deletions
@@ -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")
}
@@ -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"
}
}
@@ -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,
)
}
}
}