@@ -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,
|
||||
)
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user