Merge pull request #459 from JMR-dev/feat-355-pause-backfill-on-open

perf(sync): pause background backfill while a message is opening
This commit was merged in pull request #459.
This commit is contained in:
mergify[bot]
2026-07-08 21:21:03 +00:00
committed by GitHub
11 changed files with 638 additions and 83 deletions
@@ -0,0 +1,130 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import androidx.test.ext.junit.runners.AndroidJUnit4
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.withTimeout
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 java.util.concurrent.atomic.AtomicInteger
/**
* On-device proof of issue #355's interactive-vs-backfill coordination, on the REAL Android coroutine
* runtime (not coroutines-test virtual time) across the CI API matrix. A real [InteractiveImapGate] — the
* process-wide priority signal the reader's on-demand IMAP fetch and the background full-history backfill
* share — must:
*
* - park a backfill-style waiter ([InteractiveImapGate.awaitInteractiveIdle], used verbatim by
* [MailBackfiller.yieldToInteractive]) for as long as an interactive fetch is in flight, and resume it
* the instant the fetch clears;
* - stay held across several overlapping interactive fetches until the last releases;
* - release even when a wrapped fetch throws — so backfill can never deadlock behind a failed message open.
*
* Deliberately mock-free (no `mockk`, no framework `Context`): the gate is the whole synchronisation
* primitive #355 adds, so exercising it directly is both the faithful behavioural test and the most
* portable across API 29–37. The JVM `InteractiveImapGateTest` / `MailBackfillerTest` cover the same
* contract plus the full backfiller wiring under coroutines-test.
*/
@RunWith(AndroidJUnit4::class)
class InteractiveImapGateInstrumentedTest {
@Test
fun interactiveFetchParksABackfillWaiterUntilItReleasesThenResumesIt() = runBlocking<Unit> {
val gate = InteractiveImapGate()
val entered = CompletableDeferred<Unit>()
val release = CompletableDeferred<Unit>()
// An interactive fetch (e.g. openMessage) holds the gate until we release it.
val interactive = launch(Dispatchers.Default) {
gate.withInteractive {
entered.complete(Unit)
release.await()
}
}
entered.await()
assertTrue("the gate is held while the interactive fetch runs", gate.isInteractiveActive())
// A backfill-style waiter yields to the gate exactly as MailBackfiller does before each page.
val pagesFetched = AtomicInteger(0)
val backfill = launch(Dispatchers.Default) {
gate.awaitInteractiveIdle()
pagesFetched.incrementAndGet() // stands in for the next server page
}
delay(PARK_PROBE_MS)
assertEquals("backfill must park while the interactive fetch holds the gate", 0, pagesFetched.get())
release.complete(Unit)
withTimeout(HAND_OFF_TIMEOUT_MS) { backfill.join() }
assertEquals("backfill pages once the interactive fetch clears", 1, pagesFetched.get())
assertFalse("the gate clears after the fetch releases", gate.isInteractiveActive())
interactive.join()
}
@Test
fun anErroredInteractiveFetchStillReleasesTheGateSoBackfillNeverDeadlocks() = runBlocking<Unit> {
val gate = InteractiveImapGate()
val failed = runCatching { gate.withInteractive { throw IllegalStateException("open failed") } }
assertTrue("the failing fetch propagates to the caller", failed.isFailure)
assertFalse("a failed fetch must not strand the gate held", gate.isInteractiveActive())
// A backfill waiter must return at once now; a regression would hang until the timeout fires.
withTimeout(HAND_OFF_TIMEOUT_MS) { gate.awaitInteractiveIdle() }
}
@Test
fun theGateStaysHeldUntilTheLastOfSeveralConcurrentFetchesReleases() = runBlocking<Unit> {
val gate = InteractiveImapGate()
val entered1 = CompletableDeferred<Unit>()
val entered2 = CompletableDeferred<Unit>()
val release1 = CompletableDeferred<Unit>()
val release2 = CompletableDeferred<Unit>()
val holder1 = launch(Dispatchers.Default) {
gate.withInteractive {
entered1.complete(Unit)
release1.await()
}
}
val holder2 = launch(Dispatchers.Default) {
gate.withInteractive {
entered2.complete(Unit)
release2.await()
}
}
entered1.await()
entered2.await()
assertEquals("both concurrent fetches count", 2, gate.activeInteractiveCount.value)
val resumed = CompletableDeferred<Unit>()
launch(Dispatchers.Default) {
gate.awaitInteractiveIdle()
resumed.complete(Unit)
}
release1.complete(Unit)
delay(PARK_PROBE_MS)
assertFalse("still parked while one interactive fetch remains in flight", resumed.isCompleted)
release2.complete(Unit)
withTimeout(HAND_OFF_TIMEOUT_MS) { resumed.await() }
assertEquals("the gate clears only after the last fetch releases", 0, gate.activeInteractiveCount.value)
holder1.join()
holder2.join()
}
private companion object {
/** Slack given to a parked waiter to (wrongly) resume before we assert it is still parked. */
const val PARK_PROBE_MS = 300L
/** Generous bound for the gate hand-off; only a real park/resume regression approaches it. */
const val HAND_OFF_TIMEOUT_MS = 5_000L
}
}
@@ -39,6 +39,7 @@ import org.libremail.data.local.toOutgoingAttachments
import org.libremail.data.local.toOutgoingAttachmentsJson import org.libremail.data.local.toOutgoingAttachmentsJson
import org.libremail.data.settings.AccountSettingsRepository import org.libremail.data.settings.AccountSettingsRepository
import org.libremail.data.settings.SignatureRepository import org.libremail.data.settings.SignatureRepository
import org.libremail.data.sync.InteractiveImapGate
import org.libremail.data.sync.MailConnectionFactory import org.libremail.data.sync.MailConnectionFactory
import org.libremail.data.sync.SendScheduler import org.libremail.data.sync.SendScheduler
import org.libremail.data.sync.logSafeFolderLabel import org.libremail.data.sync.logSafeFolderLabel
@@ -79,6 +80,7 @@ class MailRepositoryImpl @Inject constructor(
private val accountSettingsRepository: AccountSettingsRepository, private val accountSettingsRepository: AccountSettingsRepository,
private val signatureRepository: SignatureRepository, private val signatureRepository: SignatureRepository,
private val attachmentUriGrants: AttachmentUriGrants, private val attachmentUriGrants: AttachmentUriGrants,
private val interactiveGate: InteractiveImapGate,
) : MailRepository { ) : MailRepository {
// Application-lifetime scope for fire-and-forget server pushes that must outlive the caller — e.g. // Application-lifetime scope for fire-and-forget server pushes that must outlive the caller — e.g.
@@ -156,42 +158,54 @@ class MailRepositoryImpl @Inject constructor(
override suspend fun openMessage(id: String): Result<Message> = withContext(Dispatchers.IO) { override suspend fun openMessage(id: String): Result<Message> = withContext(Dispatchers.IO) {
runCatching { runCatching {
// Time the whole open so a debug report shows what the reader's spinner is waiting on — a // Signal an interactive fetch for the whole open (#355) so the continuous background backfill
// cached open is a local read; a first open blocks on the IMAP body fetch below (issue #358). // parks at its next page boundary and this body fetch wins the account's IMAP throughput
val startNanos = System.nanoTime() // instead of queuing behind the backfill storm (the reader's ~48s uncached-open stall). The
// Route on the body-less projection: a cached, already-read message needs no account, no // counter is released even if the fetch throws, and a purely-cached open holds it only for the
// credentials, and no network, so it skips the Keystore decrypt + DataStore read that // brief local read.
// resolving connection params costs (issue #186). Only the fetch / SEEN-push branches below interactiveGate.withInteractive {
// pull the account and resolve params, and each does so lazily right where it is needed. // Time the whole open so a debug report shows what the reader's spinner is waiting on — a
val routing = messageDao.getRouting(id) ?: error("Message not found") // cached open is a local read; a first open blocks on the IMAP body fetch below (issue #358).
val fetchedBody = !routing.bodyFetched val startNanos = System.nanoTime()
if (!routing.bodyFetched || !routing.isRead) { // Route on the body-less projection: a cached, already-read message needs no account, no
val account = accountDao.getById(routing.accountId)?.toDomain() // credentials, and no network, so it skips the Keystore decrypt + DataStore read that
if (account != null && !routing.bodyFetched) { // resolving connection params costs (issue #186). Only the fetch / SEEN-push branches below
val params = connectionFactory.imapParamsFor(account) // pull the account and resolve params, and each does so lazily right where it is needed.
val content = imapClient.fetchBodyMarkingSeen(params, routing.folder, uidOf(id)) val routing = messageDao.getRouting(id) ?: error("Message not found")
messageDao.updateBody(id, content.body, content.isHtml, Snippet.of(content.body, content.isHtml)) val fetchedBody = !routing.bodyFetched
attachmentDao.replaceForMessage(id, content.attachments.map { it.toEntity(id) }) if (!routing.bodyFetched || !routing.isRead) {
messageDao.setRead(id, true) val account = accountDao.getById(routing.accountId)?.toDomain()
} else if (account != null) { if (account != null && !routing.bodyFetched) {
// Optimistic, local-only: the reader can render as soon as this returns. The SEEN flag val params = connectionFactory.imapParamsFor(account)
// still needs to reach the server, but that IMAP round trip (connection + STORE) must not val content = imapClient.fetchBodyMarkingSeen(params, routing.folder, uidOf(id))
// sit on this path (#148/#186) — the body/attachments are already fully local. Pushed on messageDao.updateBody(
// backgroundScope, which outlives this call. id,
val params = connectionFactory.imapParamsFor(account) content.body,
messageDao.setRead(id, true) content.isHtml,
pushSeenFlagInBackground(params, routing.folder, id) Snippet.of(content.body, content.isHtml),
)
attachmentDao.replaceForMessage(id, content.attachments.map { it.toEntity(id) })
messageDao.setRead(id, true)
} else if (account != null) {
// Optimistic, local-only: the reader can render as soon as this returns. The SEEN
// flag still needs to reach the server, but that IMAP round trip (connection +
// STORE) must not sit on this path (#148/#186) — the body/attachments are already
// fully local. Pushed on backgroundScope, which outlives this call.
val params = connectionFactory.imapParamsFor(account)
messageDao.setRead(id, true)
pushSeenFlagInBackground(params, routing.folder, id)
}
} }
// The single full-body read, reserved for the value the reader actually renders (issue #186).
val message = messageDao.getById(id)?.toDomain() ?: error("Message not found")
// PII-free: hashed account ref, system-folder label only, plus the branch taken and ms.
AppLog.i(
READER_TAG,
"openMessage ${accountLogRef(routing.accountId)} folder=${logSafeFolderLabel(routing.folder)} " +
"fetchedBody=$fetchedBody took=${(System.nanoTime() - startNanos) / NANOS_PER_MS}ms",
)
message
} }
// The single full-body read, reserved for the value the reader actually renders (issue #186).
val message = messageDao.getById(id)?.toDomain() ?: error("Message not found")
// PII-free: hashed account ref, system-folder label only, plus the branch taken and elapsed ms.
AppLog.i(
READER_TAG,
"openMessage ${accountLogRef(routing.accountId)} folder=${logSafeFolderLabel(routing.folder)} " +
"fetchedBody=$fetchedBody took=${(System.nanoTime() - startNanos) / NANOS_PER_MS}ms",
)
message
} }
} }
@@ -224,32 +238,45 @@ class MailRepositoryImpl @Inject constructor(
} }
override suspend fun inlineImages(messageId: String): List<InlineImage> = withContext(Dispatchers.IO) { override suspend fun inlineImages(messageId: String): List<InlineImage> = withContext(Dispatchers.IO) {
// Resolve the message's account/folder once (body-less), then reuse the on-disk cache per cid: // Inline-image loads are an interactive fetch too (#355): the reader is showing them, so backfill
// image — no per-image message re-read or attachment re-query (the old downloadAttachment N+1, #186). // must yield while they download. Released even if a per-image fetch throws (via withInteractive).
val routing = messageDao.getRouting(messageId) ?: return@withContext emptyList() interactiveGate.withInteractive {
val parts = attachmentDao.getForMessage(messageId) // Resolve the message's account/folder once (body-less), then reuse the on-disk cache per cid:
parts.filter { it.contentId != null }.mapNotNull { row -> // image — no per-image message re-read or attachment re-query (the old downloadAttachment N+1, #186).
// Reuse the on-disk attachment cache (download once, then instant + offline). A failed val routing = messageDao.getRouting(messageId)
// fetch just omits that image, leaving a broken <img> rather than failing the open. if (routing == null) {
val file = runCatching { emptyList()
ensureAttachmentFile(messageId, routing.accountId, routing.folder, row.partIndex, row.filename) } else {
}.getOrNull() ?: return@mapNotNull null val parts = attachmentDao.getForMessage(messageId)
InlineImage(contentId = row.contentId!!, mimeType = row.mimeType, bytes = file.readBytes()) parts.filter { it.contentId != null }.mapNotNull { row ->
// Reuse the on-disk attachment cache (download once, then instant + offline). A failed
// fetch just omits that image, leaving a broken <img> rather than failing the open.
val file = runCatching {
ensureAttachmentFile(messageId, routing.accountId, routing.folder, row.partIndex, row.filename)
}.getOrNull() ?: return@mapNotNull null
InlineImage(contentId = row.contentId!!, mimeType = row.mimeType, bytes = file.readBytes())
}
}
} }
} }
override suspend fun downloadAttachment(messageId: String, partIndex: Int): Result<File> = override suspend fun downloadAttachment(messageId: String, partIndex: Int): Result<File> =
withContext(Dispatchers.IO) { withContext(Dispatchers.IO) {
runCatching { runCatching {
val routing = messageDao.getRouting(messageId) ?: error("Message not found") // A user tapped an attachment (#355) — an interactive fetch backfill must yield to. Backfill's
val meta = attachmentDao.getForMessage(messageId).firstOrNull { it.partIndex == partIndex } // own content prefetch deliberately does NOT come through here (it calls ensureAttachmentFile
ensureAttachmentFile( // directly), so background prefetch never raises the interactive gate against backfill itself.
messageId, interactiveGate.withInteractive {
routing.accountId, val routing = messageDao.getRouting(messageId) ?: error("Message not found")
routing.folder, val meta = attachmentDao.getForMessage(messageId).firstOrNull { it.partIndex == partIndex }
partIndex, ensureAttachmentFile(
meta?.filename ?: "attachment", messageId,
) routing.accountId,
routing.folder,
partIndex,
meta?.filename ?: "attachment",
)
}
} }
} }
@@ -298,8 +325,20 @@ class MailRepositoryImpl @Inject constructor(
attachmentDao.replaceForMessage(messageId, content.attachments.map { it.toEntity(messageId) }) attachmentDao.replaceForMessage(messageId, content.attachments.map { it.toEntity(messageId) })
} }
// Auto-download every attachment's bytes into the persistent per-part cache (skips ones present). // Auto-download every attachment's bytes into the persistent per-part cache (skips ones present).
// This is BACKGROUND work driven by the backfill, so it goes straight to ensureAttachmentFile and
// deliberately bypasses downloadAttachment's interactive gate (#355): prefetch must not signal an
// interactive fetch, or backfill would park behind (yield to) its own prefetch. The routing is
// already resolved here, so this also skips re-reading it per part.
attachmentDao.getForMessage(messageId).forEach { attachment -> attachmentDao.getForMessage(messageId).forEach { attachment ->
downloadAttachment(messageId, attachment.partIndex) runCatching {
ensureAttachmentFile(
messageId,
routing.accountId,
routing.folder,
attachment.partIndex,
attachment.filename,
)
}
} }
} }
@@ -355,35 +394,39 @@ class MailRepositoryImpl @Inject constructor(
} }
override suspend fun buildReplyDraft(messageId: String, mode: ReplyMode): Result<String> = runCatching { override suspend fun buildReplyDraft(messageId: String, mode: ReplyMode): Result<String> = runCatching {
val routing = messageDao.getRouting(messageId) ?: error("Message not found") // Fetching the original for a reply/forward is an interactive fetch (#355) — the user is waiting on
val account = accountDao.getById(routing.accountId)?.toDomain() ?: error("Account not found") // the compose screen — so backfill yields to it. Released even if fetchForReply throws.
val params = connectionFactory.imapParamsFor(account) interactiveGate.withInteractive {
val context = imapClient.fetchForReply(params, routing.folder, uidOf(messageId)) val routing = messageDao.getRouting(messageId) ?: error("Message not found")
val content = ReplyBuilder.build(context, mode, account.email) val account = accountDao.getById(routing.accountId)?.toDomain() ?: error("Account not found")
// Bake the sending account's default signature into the reply/forward body — above the quoted val params = connectionFactory.imapParamsFor(account)
// original — so it round-trips as part of the draft (compose won't re-append for drafts). Both val context = imapClient.fetchForReply(params, routing.folder, uidOf(messageId))
// the plaintext and HTML forms are stored so the reply can go out as multipart/alternative. val content = ReplyBuilder.build(context, mode, account.email)
val settings = accountSettingsRepository.get(routing.accountId) // Bake the sending account's default signature into the reply/forward body — above the quoted
val sig = if (settings.signatureEnabled) { // original — so it round-trips as part of the draft (compose won't re-append for drafts). Both
SignatureBlock.of(signatureRepository.getDefault(routing.accountId)) // the plaintext and HTML forms are stored so the reply can go out as multipart/alternative.
} else { val settings = accountSettingsRepository.get(routing.accountId)
SignatureBlock.EMPTY val sig = if (settings.signatureEnabled) {
SignatureBlock.of(signatureRepository.getDefault(routing.accountId))
} else {
SignatureBlock.EMPTY
}
val draftId = UUID.randomUUID().toString()
saveDraft(
Draft(
id = draftId,
accountId = routing.accountId,
to = content.to,
cc = content.cc,
subject = content.subject,
body = sig.plain + content.body,
updatedAt = System.currentTimeMillis(),
bodyHtml = sig.html + content.bodyHtml,
attachments = emptyList(),
),
)
draftId
} }
val draftId = UUID.randomUUID().toString()
saveDraft(
Draft(
id = draftId,
accountId = routing.accountId,
to = content.to,
cc = content.cc,
subject = content.subject,
body = sig.plain + content.body,
updatedAt = System.currentTimeMillis(),
bodyHtml = sig.html + content.bodyHtml,
attachments = emptyList(),
),
)
draftId
} }
/** /**
@@ -0,0 +1,82 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import kotlinx.coroutines.flow.MutableStateFlow
import kotlinx.coroutines.flow.StateFlow
import kotlinx.coroutines.flow.asStateFlow
import kotlinx.coroutines.flow.first
import kotlinx.coroutines.flow.update
import javax.inject.Inject
import javax.inject.Singleton
/**
* Process-wide priority signal (issue #355) that lets an interactive, user-facing IMAP fetch — opening a
* message, loading its inline images, downloading an attachment, building a reply — pre-empt the
* continuous full-history backfill ([MailBackfiller], #12).
*
* The reader's on-demand body fetch and the backfill share no socket ([org.libremail.mail.ImapClient] is
* connect-per-operation) and no in-process lock, so a sustained backfill keeps the account near its
* per-account IMAP throttle/bandwidth ceiling and the un-prioritised reader fetch queues behind it for
* tens of seconds (the 2026-07-05 Pixel 10 Pro XL perf run: avg 48s to open an uncached message). This
* gate fixes that cooperatively, with no thread priorities:
*
* - The interactive paths ([org.libremail.data.repository.MailRepositoryImpl]) wrap their work in
* [withInteractive], which raises an in-flight counter for the duration of the block and *always*
* lowers it again — even when the block throws — so a failed fetch can never strand the counter above
* zero and permanently starve backfill.
* - Backfill calls [awaitInteractiveIdle] at its natural per-page yield point and *parks* (suspends)
* while the counter is non-zero, resuming the instant it returns to zero. A page already in flight
* finishes; the interactive fetch simply wins the next server round-trip (best-effort — #356 bounds the
* case where the open lands mid-page).
*
* A [counter][activeInteractiveCount] rather than a `Mutex` is deliberate: several interactive fetches
* (e.g. an open plus its inline images) may overlap and must run concurrently — only *backfill* yields to
* them, and only once they have *all* cleared. State lives only in-process (a [Singleton]); a restart
* resets it, which is fine — nothing is parked across process death.
*
* This gate is orthogonal to [AccountThrottleGate] (#360): that one makes background work back off after a
* provider *rejects* it; this one makes background work yield to a *foreground* fetch pre-emptively.
*/
@Singleton
class InteractiveImapGate @Inject constructor() {
private val activeCount = MutableStateFlow(0)
/** Interactive IMAP fetches currently in flight; `0` means background backfill may proceed. */
val activeInteractiveCount: StateFlow<Int> = activeCount.asStateFlow()
/** True while at least one interactive fetch holds the gate, so background work should yield to it. */
fun isInteractiveActive(): Boolean = activeCount.value > 0
/**
* Runs [block] as an interactive IMAP fetch: raises the in-flight counter before it starts and lowers
* it in a `finally`, so backfill yields for the whole duration and is released even when [block] throws
* — the caller's `runCatching` still observes the original throwable, and backfill never deadlocks on an
* errored fetch.
*/
suspend fun <T> withInteractive(block: suspend () -> T): T {
enter()
try {
return block()
} finally {
exit()
}
}
/**
* Suspends until no interactive fetch is in flight (returns immediately when already idle). A
* [StateFlow] always replays its latest value to a new collector, so this can never miss a wake-up: if
* the counter is already zero it returns at once, otherwise it resumes the instant the last fetch
* releases the gate.
*/
suspend fun awaitInteractiveIdle() {
activeCount.first { it == 0 }
}
private fun enter() = activeCount.update { it + 1 }
// coerceAtLeast(0) is defence in depth: withInteractive always pairs enter/exit, so the counter can't
// legitimately go negative — but a stray unbalanced exit must never drive it below zero and wedge
// awaitInteractiveIdle's `== 0` predicate forever.
private fun exit() = activeCount.update { (it - 1).coerceAtLeast(0) }
}
@@ -58,6 +58,7 @@ class MailBackfiller @Inject constructor(
private val mailRepository: MailRepository, private val mailRepository: MailRepository,
private val maintenanceGate: MailMaintenanceGate, private val maintenanceGate: MailMaintenanceGate,
private val throttleGate: AccountThrottleGate, private val throttleGate: AccountThrottleGate,
private val interactiveGate: InteractiveImapGate,
) { ) {
/** One folder's slice outcome: pages fetched, and whether an immediate follow-up slice has work to do. */ /** One folder's slice outcome: pages fetched, and whether an immediate follow-up slice has work to do. */
private data class FolderResult(val batches: Int, val moreWork: Boolean) private data class FolderResult(val batches: Int, val moreWork: Boolean)
@@ -167,6 +168,13 @@ class MailBackfiller @Inject constructor(
var stalled = false var stalled = false
while (batches < maxBatches) { while (batches < maxBatches) {
currentCoroutineContext().ensureActive() currentCoroutineContext().ensureActive()
// Yield to any in-flight interactive fetch (#355) BEFORE starting the next page: opening a
// message must win the account's IMAP throughput, so park here while a user fetch holds the
// gate and resume the instant it clears. A page already fetched (if any) finished above; the
// interactive fetch simply wins the next server round-trip. This is also the slice's first
// yield point — iteration 0 runs it before the very first page — so a slice never begins a
// page while the user is waiting on a body.
yieldToInteractive(account)
// Count floor (#13 precedence): once the folder holds as many messages as retention keeps, // Count floor (#13 precedence): once the folder holds as many messages as retention keeps,
// stop before fetching another page. Unlike the age floor below, this can be decided from // stop before fetching another page. Unlike the age floor below, this can be decided from
// the cache alone: the count pruner keeps the newest-N by UID — exactly the order paging // the cache alone: the count pruner keeps the newest-N by UID — exactly the order paging
@@ -216,6 +224,21 @@ class MailBackfiller @Inject constructor(
return FolderResult(batches, moreWork = !complete && !stalled) return FolderResult(batches, moreWork = !complete && !stalled)
} }
/**
* Parks the backfill (suspends) while an interactive, user-facing IMAP fetch is in flight (#355), so
* the reader's on-demand body fetch wins the account's throughput instead of queuing behind the
* background storm. Checks first so the common idle case adds no cost and logs nothing; only an actual
* yield brackets a PII-free park/resume breadcrumb (hashed account ref only — see #358). The gate's
* counter is released by [InteractiveImapGate.withInteractive]'s `finally` even when the interactive
* fetch errors, so this can never deadlock.
*/
private suspend fun yieldToInteractive(account: Account) {
if (!interactiveGate.isInteractiveActive()) return
AppLog.i(TAG, "backfill parking ${accountLogRef(account.id)}: interactive fetch in flight")
interactiveGate.awaitInteractiveIdle()
AppLog.i(TAG, "backfill resumed ${accountLogRef(account.id)}: interactive fetch cleared")
}
/** True once the folder already holds as many messages as the count retention keeps (or more). */ /** True once the folder already holds as many messages as the count retention keeps (or more). */
private suspend fun reachedCountFloor(accountId: String, folder: String, policy: RetentionPolicy): Boolean { private suspend fun reachedCountFloor(accountId: String, folder: String, policy: RetentionPolicy): Boolean {
val limit = policy.countLimit ?: return false val limit = policy.countLimit ?: return false
@@ -16,6 +16,7 @@ import org.libremail.data.local.dao.DraftDao
import org.libremail.data.local.dao.OutboxDao import org.libremail.data.local.dao.OutboxDao
import org.libremail.data.local.entity.DraftEntity import org.libremail.data.local.entity.DraftEntity
import org.libremail.data.local.entity.OutboxEntity import org.libremail.data.local.entity.OutboxEntity
import org.libremail.data.sync.InteractiveImapGate
import java.nio.file.Files import java.nio.file.Files
/** /**
@@ -45,6 +46,7 @@ class MailRepositoryGrantsTest {
accountSettingsRepository = mockk(relaxed = true), accountSettingsRepository = mockk(relaxed = true),
signatureRepository = mockk(relaxed = true), signatureRepository = mockk(relaxed = true),
attachmentUriGrants = attachmentUriGrants, attachmentUriGrants = attachmentUriGrants,
interactiveGate = InteractiveImapGate(),
) )
@Test @Test
@@ -40,6 +40,7 @@ import org.libremail.data.local.entity.OutboxEntity
import org.libremail.data.local.entity.ServerConfigEmbedded import org.libremail.data.local.entity.ServerConfigEmbedded
import org.libremail.data.settings.AccountSettingsRepository import org.libremail.data.settings.AccountSettingsRepository
import org.libremail.data.settings.SignatureRepository import org.libremail.data.settings.SignatureRepository
import org.libremail.data.sync.InteractiveImapGate
import org.libremail.data.sync.MailConnectionFactory import org.libremail.data.sync.MailConnectionFactory
import org.libremail.data.sync.SendScheduler import org.libremail.data.sync.SendScheduler
import org.libremail.domain.model.Account import org.libremail.domain.model.Account
@@ -97,6 +98,7 @@ class MailRepositoryImplCoverageTest {
accountSettingsRepository = accountSettingsRepository, accountSettingsRepository = accountSettingsRepository,
signatureRepository = signatureRepository, signatureRepository = signatureRepository,
attachmentUriGrants = mockk<AttachmentUriGrants>(relaxed = true), attachmentUriGrants = mockk<AttachmentUriGrants>(relaxed = true),
interactiveGate = InteractiveImapGate(),
) )
// openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain // openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain
@@ -40,6 +40,7 @@ import org.libremail.data.local.entity.MessageSummary
import org.libremail.data.local.entity.ServerConfigEmbedded import org.libremail.data.local.entity.ServerConfigEmbedded
import org.libremail.data.settings.AccountSettingsRepository import org.libremail.data.settings.AccountSettingsRepository
import org.libremail.data.settings.SignatureRepository import org.libremail.data.settings.SignatureRepository
import org.libremail.data.sync.InteractiveImapGate
import org.libremail.data.sync.MailConnectionFactory import org.libremail.data.sync.MailConnectionFactory
import org.libremail.domain.model.AccountSettings import org.libremail.domain.model.AccountSettings
import org.libremail.domain.model.FolderRole import org.libremail.domain.model.FolderRole
@@ -74,6 +75,9 @@ class MailRepositoryImplTest {
private val context = mockk<Context>(relaxed = true) private val context = mockk<Context>(relaxed = true)
private val accountSettingsRepository = mockk<AccountSettingsRepository>() private val accountSettingsRepository = mockk<AccountSettingsRepository>()
private val signatureRepository = mockk<SignatureRepository>() private val signatureRepository = mockk<SignatureRepository>()
// A real gate (cheap, no deps) so tests can observe the interactive-fetch counter it raises (#355).
private val interactiveGate = InteractiveImapGate()
private val repository = MailRepositoryImpl( private val repository = MailRepositoryImpl(
context = context, context = context,
messageDao = messageDao, messageDao = messageDao,
@@ -89,6 +93,7 @@ class MailRepositoryImplTest {
signatureRepository = signatureRepository, signatureRepository = signatureRepository,
// Grant-release wiring (deleteDraft / cancelOutboxMessage) is covered by MailRepositoryGrantsTest. // Grant-release wiring (deleteDraft / cancelOutboxMessage) is covered by MailRepositoryGrantsTest.
attachmentUriGrants = mockk(relaxed = true), attachmentUriGrants = mockk(relaxed = true),
interactiveGate = interactiveGate,
) )
// openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain // openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain
@@ -260,6 +265,30 @@ class MailRepositoryImplTest {
assertTrue(breadcrumb.contains("took="), breadcrumb) assertTrue(breadcrumb.contains("took="), breadcrumb)
} }
@Test
fun `openMessage holds the interactive gate for the whole body fetch, then releases it (issue 355)`() = runTest {
val id = "acct:INBOX:9"
coEvery { messageDao.getRouting(id) } returns messageRouting(id, "INBOX")
coEvery { messageDao.getById(id) } returns messageEntity(id, "INBOX")
coEvery { accountDao.getById("acct") } returns accountEntity()
coEvery { connectionFactory.imapParamsFor(any()) } returns imapParams()
coEvery { messageDao.updateBody(id, any(), any(), any()) } just Runs
coEvery { messageDao.setRead(id, true) } just Runs
// The gate must be RAISED while the on-demand body fetch runs, so a concurrent backfill parks
// (#355). Sample the counter from inside the fetch itself — the only point it can be observed.
var countDuringFetch = -1
coEvery { imapClient.fetchBodyMarkingSeen(any(), "INBOX", "9") } coAnswers {
countDuringFetch = interactiveGate.activeInteractiveCount.value
MessageContent("Body text", isHtml = false)
}
assertTrue(repository.openMessage(id).isSuccess)
assertEquals(1, countDuringFetch, "the body fetch must run inside withInteractive")
// …and the gate is released once the open returns, so backfill can resume.
assertEquals(0, interactiveGate.activeInteractiveCount.value, "the gate clears when the open returns")
}
@Test @Test
fun `openMessage derives a readable plain-text snippet from an HTML body`() = runTest { fun `openMessage derives a readable plain-text snippet from an HTML body`() = runTest {
val id = "acct:INBOX:20" val id = "acct:INBOX:20"
@@ -0,0 +1,152 @@
// SPDX-License-Identifier: GPL-3.0-or-later
package org.libremail.data.sync
import app.cash.turbine.test
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.test.runTest
import kotlinx.coroutines.withTimeout
import org.junit.Test
import kotlin.test.assertEquals
import kotlin.test.assertFalse
import kotlin.test.assertSame
import kotlin.test.assertTrue
/**
* [InteractiveImapGate] (issue #355) must let a foreground/interactive IMAP fetch pre-empt background
* backfill: raise an in-flight count for the duration of [InteractiveImapGate.withInteractive], park
* [InteractiveImapGate.awaitInteractiveIdle] until the count clears, hold across several concurrent
* fetches until the last releases, and — critically — always release even when the wrapped fetch throws
* so backfill can never deadlock behind a failed open. The concurrency tests mirror
* [MailMaintenanceGateTest]'s real-dispatcher idiom (runBlocking + Dispatchers.Default + CompletableDeferred)
* rather than virtual time, since the contract under test is precisely cross-coroutine hand-off.
*/
class InteractiveImapGateTest {
@Test
fun `an idle gate reports inactive and awaitInteractiveIdle returns at once`() = runTest {
val gate = InteractiveImapGate()
assertFalse(gate.isInteractiveActive())
assertEquals(0, gate.activeInteractiveCount.value)
// Must not suspend forever when nothing is in flight.
withTimeout(2_000) { gate.awaitInteractiveIdle() }
}
@Test
fun `withInteractive releases the gate even when the block throws (no deadlock)`() = runTest {
val gate = InteractiveImapGate()
val boom = IllegalStateException("interactive fetch failed")
val thrown = runCatching { gate.withInteractive { throw boom } }.exceptionOrNull()
assertSame(boom, thrown, "the original throwable propagates to the caller's runCatching")
assertFalse(gate.isInteractiveActive(), "a failed fetch must not strand the counter above zero")
assertEquals(0, gate.activeInteractiveCount.value)
// A backfill parked on this gate would resume immediately now — prove it does not hang.
withTimeout(2_000) { gate.awaitInteractiveIdle() }
}
@Test
fun `awaitInteractiveIdle parks while a fetch is in flight and resumes when it releases`() = runBlocking<Unit> {
val gate = InteractiveImapGate()
val entered = CompletableDeferred<Unit>()
val release = CompletableDeferred<Unit>()
val holder = launch(Dispatchers.Default) {
gate.withInteractive {
entered.complete(Unit)
release.await()
}
}
entered.await() // the interactive fetch is now in flight
assertTrue(gate.isInteractiveActive())
assertEquals(1, gate.activeInteractiveCount.value)
val resumed = CompletableDeferred<Unit>()
launch(Dispatchers.Default) {
gate.awaitInteractiveIdle()
resumed.complete(Unit)
}
delay(PARK_PROBE_MS)
assertFalse(resumed.isCompleted, "the waiter must stay parked while the interactive fetch holds the gate")
release.complete(Unit)
withTimeout(2_000) { resumed.await() }
assertFalse(gate.isInteractiveActive(), "the gate clears once the fetch releases")
holder.join()
}
@Test
fun `the gate stays held until the last of several concurrent fetches releases`() = runBlocking<Unit> {
val gate = InteractiveImapGate()
val entered1 = CompletableDeferred<Unit>()
val entered2 = CompletableDeferred<Unit>()
val release1 = CompletableDeferred<Unit>()
val release2 = CompletableDeferred<Unit>()
val holder1 = launch(Dispatchers.Default) {
gate.withInteractive {
entered1.complete(Unit)
release1.await()
}
}
val holder2 = launch(Dispatchers.Default) {
gate.withInteractive {
entered2.complete(Unit)
release2.await()
}
}
entered1.await()
entered2.await()
assertEquals(2, gate.activeInteractiveCount.value, "both concurrent fetches count")
val resumed = CompletableDeferred<Unit>()
launch(Dispatchers.Default) {
gate.awaitInteractiveIdle()
resumed.complete(Unit)
}
release1.complete(Unit)
delay(PARK_PROBE_MS)
assertFalse(resumed.isCompleted, "still parked while one interactive fetch remains in flight")
assertEquals(1, gate.activeInteractiveCount.value)
release2.complete(Unit)
withTimeout(2_000) { resumed.await() }
assertEquals(0, gate.activeInteractiveCount.value, "the gate clears only after the last fetch releases")
holder1.join()
holder2.join()
}
@Test
fun `activeInteractiveCount emits the rise and fall around an interactive fetch`() = runBlocking<Unit> {
val gate = InteractiveImapGate()
val entered = CompletableDeferred<Unit>()
val release = CompletableDeferred<Unit>()
gate.activeInteractiveCount.test {
assertEquals(0, awaitItem(), "starts idle")
val holder = launch(Dispatchers.Default) {
gate.withInteractive {
entered.complete(Unit)
release.await()
}
}
entered.await()
assertEquals(1, awaitItem(), "rises when the fetch starts")
release.complete(Unit)
assertEquals(0, awaitItem(), "falls when the fetch completes")
holder.join()
cancelAndIgnoreRemainingEvents()
}
}
private companion object {
/** Slack given to a parked coroutine to (wrongly) resume before we assert it is still parked. */
const val PARK_PROBE_MS = 200L
}
}
@@ -17,8 +17,14 @@ import jakarta.mail.MessagingException
import jakarta.mail.Session import jakarta.mail.Session
import jakarta.mail.internet.InternetAddress import jakarta.mail.internet.InternetAddress
import jakarta.mail.internet.MimeMessage import jakarta.mail.internet.MimeMessage
import kotlinx.coroutines.CompletableDeferred
import kotlinx.coroutines.Dispatchers
import kotlinx.coroutines.delay
import kotlinx.coroutines.flow.flowOf import kotlinx.coroutines.flow.flowOf
import kotlinx.coroutines.launch
import kotlinx.coroutines.runBlocking
import kotlinx.coroutines.test.runTest import kotlinx.coroutines.test.runTest
import kotlinx.coroutines.withTimeout
import org.junit.After import org.junit.After
import org.junit.Before import org.junit.Before
import org.junit.Test import org.junit.Test
@@ -524,6 +530,82 @@ class MailBackfillerTest {
assertTrue(logBuffer.snapshot().any { it.message.startsWith("backfill skip acct:") }) assertTrue(logBuffer.snapshot().any { it.message.startsWith("backfill skip acct:") })
} }
// --- issue #355: interactive-fetch priority -------------------------------------------------
/**
* The core of #355: while an interactive fetch (a message open, an attachment download, …) holds the
* [InteractiveImapGate], backfill must PARK at its per-page yield point instead of stealing the
* account's IMAP throughput, and resume the instant the interactive fetch clears. The backfiller and
* the interactive holder share ONLY the gate, so a pass is attributable solely to it. Mirrors
* [MailMaintenanceGateTest]'s real-dispatcher hand-off idiom (runBlocking + Dispatchers.Default).
*/
@Test
fun `backfill parks before its next page while an interactive fetch holds the gate, then resumes`() =
runBlocking<Unit> {
cached += fetchedMessage(uid = "60").toEntity("acct", "INBOX")
val gate = InteractiveImapGate()
val fetched = CompletableDeferred<Unit>()
val imapClient = mockk<ImapClient>()
coEvery { imapClient.fetchOlderThan(any(), any(), any(), any()) } coAnswers {
fetched.complete(Unit)
emptyList() // one page, then the folder completes and the slice ends
}
// An interactive fetch is in flight: it holds the gate until we release it.
val entered = CompletableDeferred<Unit>()
val release = CompletableDeferred<Unit>()
val interactive = launch(Dispatchers.Default) {
gate.withInteractive {
entered.complete(Unit)
release.await()
}
}
entered.await()
val backfillJob = launch(Dispatchers.Default) {
backfiller(AccountSettings("acct"), imapClient = imapClient, interactiveGate = gate).runBackfill()
}
delay(PARK_PROBE_MS)
assertFalse(fetched.isCompleted, "backfill must not page while an interactive fetch is active")
assertTrue(
logBuffer.snapshot().any { it.message.startsWith("backfill parking acct:") },
"a PII-free park breadcrumb is recorded",
)
release.complete(Unit)
withTimeout(HAND_OFF_TIMEOUT_MS) { backfillJob.join() }
assertTrue(fetched.isCompleted, "backfill resumes and pages once the interactive fetch clears")
assertTrue(
logBuffer.snapshot().any { it.message.startsWith("backfill resumed acct:") },
"a PII-free resume breadcrumb is recorded",
)
interactive.join()
}
/**
* No-deadlock guarantee: an interactive fetch that ERRORS still releases the gate (withInteractive's
* `finally`), so backfill is never stranded parked behind a failed open. The failing fetch has already
* released the gate here, and backfill then runs to completion without hanging.
*/
@Test
fun `an interactive fetch that errors does not strand backfill parked`() = runBlocking<Unit> {
cached += fetchedMessage(uid = "60").toEntity("acct", "INBOX")
val gate = InteractiveImapGate()
val imapClient = mockk<ImapClient>()
coEvery { imapClient.fetchOlderThan(any(), any(), any(), any()) } returns emptyList()
val failed = runCatching { gate.withInteractive { throw MessagingException("open failed") } }
assertTrue(failed.isFailure, "the interactive fetch failed")
assertFalse(gate.isInteractiveActive(), "the errored fetch still released the gate")
val moreWork = withTimeout(HAND_OFF_TIMEOUT_MS) {
backfiller(AccountSettings("acct"), imapClient = imapClient, interactiveGate = gate).runBackfill()
}
assertFalse(moreWork)
coVerify(atLeast = 1) { imapClient.fetchOlderThan(any(), any(), any(), any()) }
}
// --- issue #329: AppLog breadcrumbs --------------------------------------------------------- // --- issue #329: AppLog breadcrumbs ---------------------------------------------------------
@Test @Test
@@ -603,6 +685,7 @@ class MailBackfillerTest {
battery: BatteryStatus = BatteryStatus(percent = 100, isCharging = false), battery: BatteryStatus = BatteryStatus(percent = 100, isCharging = false),
imapClient: ImapClient = client, imapClient: ImapClient = client,
throttleGate: AccountThrottleGate = AccountThrottleGate(), throttleGate: AccountThrottleGate = AccountThrottleGate(),
interactiveGate: InteractiveImapGate = InteractiveImapGate(),
): MailBackfiller { ): MailBackfiller {
val accountDao = mockk<AccountDao>() val accountDao = mockk<AccountDao>()
coEvery { accountDao.getAll() } returns listOf(accountEntity) coEvery { accountDao.getAll() } returns listOf(accountEntity)
@@ -659,6 +742,7 @@ class MailBackfillerTest {
mailRepository = mailRepository, mailRepository = mailRepository,
maintenanceGate = MailMaintenanceGate(), maintenanceGate = MailMaintenanceGate(),
throttleGate = throttleGate, throttleGate = throttleGate,
interactiveGate = interactiveGate,
).also { ).also {
lastMessageDao = messageDao lastMessageDao = messageDao
lastMailRepository = mailRepository lastMailRepository = mailRepository
@@ -723,5 +807,11 @@ class MailBackfillerTest {
const val WINDOW = 50 const val WINDOW = 50
private const val DAY_MILLIS = 24L * 60 * 60 * 1000 private const val DAY_MILLIS = 24L * 60 * 60 * 1000
private const val MONTH_MILLIS = 30L * DAY_MILLIS private const val MONTH_MILLIS = 30L * DAY_MILLIS
/** Slack given to a parked backfill to (wrongly) page before we assert it is still parked (#355). */
private const val PARK_PROBE_MS = 300L
/** Generous bound for the gate hand-off; only a real park/resume regression approaches it (#355). */
private const val HAND_OFF_TIMEOUT_MS = 5_000L
} }
} }
@@ -201,6 +201,7 @@ class MailMaintenanceGateTest {
mailRepository = mockk(relaxed = true), mailRepository = mockk(relaxed = true),
maintenanceGate = gate, maintenanceGate = gate,
throttleGate = AccountThrottleGate(), throttleGate = AccountThrottleGate(),
interactiveGate = InteractiveImapGate(),
) )
} }
@@ -361,6 +361,7 @@ class MailSyncConcurrencyTest {
mailRepository = mockk(relaxed = true), mailRepository = mockk(relaxed = true),
maintenanceGate = MailMaintenanceGate(), maintenanceGate = MailMaintenanceGate(),
throttleGate = AccountThrottleGate(), throttleGate = AccountThrottleGate(),
interactiveGate = InteractiveImapGate(),
) )
} }