WIP: merge queue: checking main (389bd97) and #459 together #461

Closed
mergify[bot] wants to merge 2 commits from mergify/merge-queue/b6900d038b into main
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.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<Message> = 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<InlineImage> = 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 <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())
// 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 <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> =
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<String> = 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
}
/**
@@ -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 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
@@ -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
@@ -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<AttachmentUriGrants>(relaxed = true),
interactiveGate = InteractiveImapGate(),
)
// 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.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<Context>(relaxed = true)
private val accountSettingsRepository = mockk<AccountSettingsRepository>()
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(
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"
@@ -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.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<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 ---------------------------------------------------------
@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<AccountDao>()
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
}
}
@@ -201,6 +201,7 @@ class MailMaintenanceGateTest {
mailRepository = mockk(relaxed = true),
maintenanceGate = gate,
throttleGate = AccountThrottleGate(),
interactiveGate = InteractiveImapGate(),
)
}
@@ -361,6 +361,7 @@ class MailSyncConcurrencyTest {
mailRepository = mockk(relaxed = true),
maintenanceGate = MailMaintenanceGate(),
throttleGate = AccountThrottleGate(),
interactiveGate = InteractiveImapGate(),
)
}