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