diff --git a/app/src/androidTest/kotlin/org/libremail/data/sync/InteractiveImapGateInstrumentedTest.kt b/app/src/androidTest/kotlin/org/libremail/data/sync/InteractiveImapGateInstrumentedTest.kt new file mode 100644 index 0000000..a50ab17 --- /dev/null +++ b/app/src/androidTest/kotlin/org/libremail/data/sync/InteractiveImapGateInstrumentedTest.kt @@ -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 { + val gate = InteractiveImapGate() + val entered = CompletableDeferred() + val release = CompletableDeferred() + + // 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 { + 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 { + val gate = InteractiveImapGate() + val entered1 = CompletableDeferred() + val entered2 = CompletableDeferred() + val release1 = CompletableDeferred() + val release2 = CompletableDeferred() + + 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() + 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 + } +} diff --git a/app/src/main/kotlin/org/libremail/data/repository/MailRepositoryImpl.kt b/app/src/main/kotlin/org/libremail/data/repository/MailRepositoryImpl.kt index 5d47286..1a0480a 100644 --- a/app/src/main/kotlin/org/libremail/data/repository/MailRepositoryImpl.kt +++ b/app/src/main/kotlin/org/libremail/data/repository/MailRepositoryImpl.kt @@ -39,6 +39,7 @@ import org.libremail.data.local.toOutgoingAttachments import org.libremail.data.local.toOutgoingAttachmentsJson import org.libremail.data.settings.AccountSettingsRepository import org.libremail.data.settings.SignatureRepository +import org.libremail.data.sync.InteractiveImapGate import org.libremail.data.sync.MailConnectionFactory import org.libremail.data.sync.SendScheduler import org.libremail.data.sync.logSafeFolderLabel @@ -79,6 +80,7 @@ class MailRepositoryImpl @Inject constructor( private val accountSettingsRepository: AccountSettingsRepository, private val signatureRepository: SignatureRepository, private val attachmentUriGrants: AttachmentUriGrants, + private val interactiveGate: InteractiveImapGate, ) : MailRepository { // 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 = withContext(Dispatchers.IO) { runCatching { - // Time the whole open so a debug report shows what the reader's spinner is waiting on — a - // cached open is a local read; a first open blocks on the IMAP body fetch below (issue #358). - val startNanos = System.nanoTime() - // Route on the body-less projection: a cached, already-read message needs no account, no - // credentials, and no network, so it skips the Keystore decrypt + DataStore read that - // resolving connection params costs (issue #186). Only the fetch / SEEN-push branches below - // pull the account and resolve params, and each does so lazily right where it is needed. - val routing = messageDao.getRouting(id) ?: error("Message not found") - val fetchedBody = !routing.bodyFetched - if (!routing.bodyFetched || !routing.isRead) { - val account = accountDao.getById(routing.accountId)?.toDomain() - if (account != null && !routing.bodyFetched) { - val params = connectionFactory.imapParamsFor(account) - val content = imapClient.fetchBodyMarkingSeen(params, routing.folder, uidOf(id)) - messageDao.updateBody(id, content.body, content.isHtml, 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) + // Signal an interactive fetch for the whole open (#355) so the continuous background backfill + // parks at its next page boundary and this body fetch wins the account's IMAP throughput + // instead of queuing behind the backfill storm (the reader's ~48s uncached-open stall). The + // counter is released even if the fetch throws, and a purely-cached open holds it only for the + // brief local read. + interactiveGate.withInteractive { + // Time the whole open so a debug report shows what the reader's spinner is waiting on — a + // cached open is a local read; a first open blocks on the IMAP body fetch below (issue #358). + val startNanos = System.nanoTime() + // Route on the body-less projection: a cached, already-read message needs no account, no + // credentials, and no network, so it skips the Keystore decrypt + DataStore read that + // resolving connection params costs (issue #186). Only the fetch / SEEN-push branches below + // pull the account and resolve params, and each does so lazily right where it is needed. + val routing = messageDao.getRouting(id) ?: error("Message not found") + val fetchedBody = !routing.bodyFetched + if (!routing.bodyFetched || !routing.isRead) { + val account = accountDao.getById(routing.accountId)?.toDomain() + if (account != null && !routing.bodyFetched) { + val params = connectionFactory.imapParamsFor(account) + val content = imapClient.fetchBodyMarkingSeen(params, routing.folder, uidOf(id)) + messageDao.updateBody( + id, + content.body, + content.isHtml, + 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 = withContext(Dispatchers.IO) { - // Resolve the message's account/folder once (body-less), then reuse the on-disk cache per cid: - // image — no per-image message re-read or attachment re-query (the old downloadAttachment N+1, #186). - val routing = messageDao.getRouting(messageId) ?: return@withContext emptyList() - val parts = attachmentDao.getForMessage(messageId) - 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 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()) + // Inline-image loads are an interactive fetch too (#355): the reader is showing them, so backfill + // must yield while they download. Released even if a per-image fetch throws (via withInteractive). + interactiveGate.withInteractive { + // Resolve the message's account/folder once (body-less), then reuse the on-disk cache per cid: + // image — no per-image message re-read or attachment re-query (the old downloadAttachment N+1, #186). + val routing = messageDao.getRouting(messageId) + if (routing == null) { + emptyList() + } else { + val parts = attachmentDao.getForMessage(messageId) + 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 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 = withContext(Dispatchers.IO) { runCatching { - val routing = messageDao.getRouting(messageId) ?: error("Message not found") - val meta = attachmentDao.getForMessage(messageId).firstOrNull { it.partIndex == partIndex } - ensureAttachmentFile( - messageId, - routing.accountId, - routing.folder, - partIndex, - meta?.filename ?: "attachment", - ) + // A user tapped an attachment (#355) — an interactive fetch backfill must yield to. Backfill's + // own content prefetch deliberately does NOT come through here (it calls ensureAttachmentFile + // directly), so background prefetch never raises the interactive gate against backfill itself. + interactiveGate.withInteractive { + val routing = messageDao.getRouting(messageId) ?: error("Message not found") + val meta = attachmentDao.getForMessage(messageId).firstOrNull { it.partIndex == partIndex } + ensureAttachmentFile( + 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) }) } // 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 -> - 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 = runCatching { - val routing = messageDao.getRouting(messageId) ?: error("Message not found") - val account = accountDao.getById(routing.accountId)?.toDomain() ?: error("Account not found") - val params = connectionFactory.imapParamsFor(account) - val context = imapClient.fetchForReply(params, routing.folder, uidOf(messageId)) - val content = ReplyBuilder.build(context, mode, account.email) - // Bake the sending account's default signature into the reply/forward body — above the quoted - // original — so it round-trips as part of the draft (compose won't re-append for drafts). Both - // the plaintext and HTML forms are stored so the reply can go out as multipart/alternative. - val settings = accountSettingsRepository.get(routing.accountId) - val sig = if (settings.signatureEnabled) { - SignatureBlock.of(signatureRepository.getDefault(routing.accountId)) - } else { - SignatureBlock.EMPTY + // Fetching the original for a reply/forward is an interactive fetch (#355) — the user is waiting on + // the compose screen — so backfill yields to it. Released even if fetchForReply throws. + interactiveGate.withInteractive { + val routing = messageDao.getRouting(messageId) ?: error("Message not found") + val account = accountDao.getById(routing.accountId)?.toDomain() ?: error("Account not found") + val params = connectionFactory.imapParamsFor(account) + val context = imapClient.fetchForReply(params, routing.folder, uidOf(messageId)) + val content = ReplyBuilder.build(context, mode, account.email) + // Bake the sending account's default signature into the reply/forward body — above the quoted + // original — so it round-trips as part of the draft (compose won't re-append for drafts). Both + // the plaintext and HTML forms are stored so the reply can go out as multipart/alternative. + val settings = accountSettingsRepository.get(routing.accountId) + 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 } /** diff --git a/app/src/main/kotlin/org/libremail/data/sync/InteractiveImapGate.kt b/app/src/main/kotlin/org/libremail/data/sync/InteractiveImapGate.kt new file mode 100644 index 0000000..7848b1d --- /dev/null +++ b/app/src/main/kotlin/org/libremail/data/sync/InteractiveImapGate.kt @@ -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 = 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 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) } +} diff --git a/app/src/main/kotlin/org/libremail/data/sync/MailBackfiller.kt b/app/src/main/kotlin/org/libremail/data/sync/MailBackfiller.kt index e067cb7..42bb8f0 100644 --- a/app/src/main/kotlin/org/libremail/data/sync/MailBackfiller.kt +++ b/app/src/main/kotlin/org/libremail/data/sync/MailBackfiller.kt @@ -58,6 +58,7 @@ class MailBackfiller @Inject constructor( private val mailRepository: MailRepository, private val maintenanceGate: MailMaintenanceGate, 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. */ private data class FolderResult(val batches: Int, val moreWork: Boolean) @@ -167,6 +168,13 @@ class MailBackfiller @Inject constructor( var stalled = false while (batches < maxBatches) { 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, // 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 @@ -216,6 +224,21 @@ class MailBackfiller @Inject constructor( 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). */ private suspend fun reachedCountFloor(accountId: String, folder: String, policy: RetentionPolicy): Boolean { val limit = policy.countLimit ?: return false diff --git a/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryGrantsTest.kt b/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryGrantsTest.kt index ddab499..805955c 100644 --- a/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryGrantsTest.kt +++ b/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryGrantsTest.kt @@ -16,6 +16,7 @@ import org.libremail.data.local.dao.DraftDao import org.libremail.data.local.dao.OutboxDao import org.libremail.data.local.entity.DraftEntity import org.libremail.data.local.entity.OutboxEntity +import org.libremail.data.sync.InteractiveImapGate import java.nio.file.Files /** @@ -45,6 +46,7 @@ class MailRepositoryGrantsTest { accountSettingsRepository = mockk(relaxed = true), signatureRepository = mockk(relaxed = true), attachmentUriGrants = attachmentUriGrants, + interactiveGate = InteractiveImapGate(), ) @Test diff --git a/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplCoverageTest.kt b/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplCoverageTest.kt index db506b2..46afc61 100644 --- a/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplCoverageTest.kt +++ b/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplCoverageTest.kt @@ -40,6 +40,7 @@ import org.libremail.data.local.entity.OutboxEntity import org.libremail.data.local.entity.ServerConfigEmbedded import org.libremail.data.settings.AccountSettingsRepository import org.libremail.data.settings.SignatureRepository +import org.libremail.data.sync.InteractiveImapGate import org.libremail.data.sync.MailConnectionFactory import org.libremail.data.sync.SendScheduler import org.libremail.domain.model.Account @@ -97,6 +98,7 @@ class MailRepositoryImplCoverageTest { accountSettingsRepository = accountSettingsRepository, signatureRepository = signatureRepository, attachmentUriGrants = mockk(relaxed = true), + interactiveGate = InteractiveImapGate(), ) // openMessage now breadcrumbs via AppLog (issue #358); android.util.Log is a no-op stub under plain diff --git a/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplTest.kt b/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplTest.kt index bca3438..36a22fc 100644 --- a/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplTest.kt +++ b/app/src/test/kotlin/org/libremail/data/repository/MailRepositoryImplTest.kt @@ -40,6 +40,7 @@ import org.libremail.data.local.entity.MessageSummary import org.libremail.data.local.entity.ServerConfigEmbedded import org.libremail.data.settings.AccountSettingsRepository import org.libremail.data.settings.SignatureRepository +import org.libremail.data.sync.InteractiveImapGate import org.libremail.data.sync.MailConnectionFactory import org.libremail.domain.model.AccountSettings import org.libremail.domain.model.FolderRole @@ -74,6 +75,9 @@ class MailRepositoryImplTest { private val context = mockk(relaxed = true) private val accountSettingsRepository = mockk() private val signatureRepository = mockk() + + // 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( context = context, messageDao = messageDao, @@ -89,6 +93,7 @@ class MailRepositoryImplTest { signatureRepository = signatureRepository, // Grant-release wiring (deleteDraft / cancelOutboxMessage) is covered by MailRepositoryGrantsTest. attachmentUriGrants = mockk(relaxed = true), + interactiveGate = interactiveGate, ) // 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) } + @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 fun `openMessage derives a readable plain-text snippet from an HTML body`() = runTest { val id = "acct:INBOX:20" diff --git a/app/src/test/kotlin/org/libremail/data/sync/InteractiveImapGateTest.kt b/app/src/test/kotlin/org/libremail/data/sync/InteractiveImapGateTest.kt new file mode 100644 index 0000000..13760fe --- /dev/null +++ b/app/src/test/kotlin/org/libremail/data/sync/InteractiveImapGateTest.kt @@ -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 { + val gate = InteractiveImapGate() + val entered = CompletableDeferred() + val release = CompletableDeferred() + + 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() + 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 { + val gate = InteractiveImapGate() + val entered1 = CompletableDeferred() + val entered2 = CompletableDeferred() + val release1 = CompletableDeferred() + val release2 = CompletableDeferred() + + 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() + 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 { + val gate = InteractiveImapGate() + val entered = CompletableDeferred() + val release = CompletableDeferred() + + 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 + } +} diff --git a/app/src/test/kotlin/org/libremail/data/sync/MailBackfillerTest.kt b/app/src/test/kotlin/org/libremail/data/sync/MailBackfillerTest.kt index d43d195..4a9ad75 100644 --- a/app/src/test/kotlin/org/libremail/data/sync/MailBackfillerTest.kt +++ b/app/src/test/kotlin/org/libremail/data/sync/MailBackfillerTest.kt @@ -17,8 +17,14 @@ import jakarta.mail.MessagingException import jakarta.mail.Session import jakarta.mail.internet.InternetAddress 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.launch +import kotlinx.coroutines.runBlocking import kotlinx.coroutines.test.runTest +import kotlinx.coroutines.withTimeout import org.junit.After import org.junit.Before import org.junit.Test @@ -524,6 +530,82 @@ class MailBackfillerTest { 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 { + cached += fetchedMessage(uid = "60").toEntity("acct", "INBOX") + val gate = InteractiveImapGate() + val fetched = CompletableDeferred() + val imapClient = mockk() + 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() + val release = CompletableDeferred() + 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 { + cached += fetchedMessage(uid = "60").toEntity("acct", "INBOX") + val gate = InteractiveImapGate() + val imapClient = mockk() + 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 --------------------------------------------------------- @Test @@ -603,6 +685,7 @@ class MailBackfillerTest { battery: BatteryStatus = BatteryStatus(percent = 100, isCharging = false), imapClient: ImapClient = client, throttleGate: AccountThrottleGate = AccountThrottleGate(), + interactiveGate: InteractiveImapGate = InteractiveImapGate(), ): MailBackfiller { val accountDao = mockk() coEvery { accountDao.getAll() } returns listOf(accountEntity) @@ -659,6 +742,7 @@ class MailBackfillerTest { mailRepository = mailRepository, maintenanceGate = MailMaintenanceGate(), throttleGate = throttleGate, + interactiveGate = interactiveGate, ).also { lastMessageDao = messageDao lastMailRepository = mailRepository @@ -723,5 +807,11 @@ class MailBackfillerTest { const val WINDOW = 50 private const val DAY_MILLIS = 24L * 60 * 60 * 1000 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 } } diff --git a/app/src/test/kotlin/org/libremail/data/sync/MailMaintenanceGateTest.kt b/app/src/test/kotlin/org/libremail/data/sync/MailMaintenanceGateTest.kt index dbe5e91..410f941 100644 --- a/app/src/test/kotlin/org/libremail/data/sync/MailMaintenanceGateTest.kt +++ b/app/src/test/kotlin/org/libremail/data/sync/MailMaintenanceGateTest.kt @@ -201,6 +201,7 @@ class MailMaintenanceGateTest { mailRepository = mockk(relaxed = true), maintenanceGate = gate, throttleGate = AccountThrottleGate(), + interactiveGate = InteractiveImapGate(), ) } diff --git a/app/src/test/kotlin/org/libremail/data/sync/MailSyncConcurrencyTest.kt b/app/src/test/kotlin/org/libremail/data/sync/MailSyncConcurrencyTest.kt index 1f2ebd7..c659e8e 100644 --- a/app/src/test/kotlin/org/libremail/data/sync/MailSyncConcurrencyTest.kt +++ b/app/src/test/kotlin/org/libremail/data/sync/MailSyncConcurrencyTest.kt @@ -361,6 +361,7 @@ class MailSyncConcurrencyTest { mailRepository = mockk(relaxed = true), maintenanceGate = MailMaintenanceGate(), throttleGate = AccountThrottleGate(), + interactiveGate = InteractiveImapGate(), ) }