spike(imap): prototype flag-gated connection reuse for folder-open
Prototype the per-account connection reuse the #125 investigation recommended and deferred, behind an OFF-by-default flag so it cannot destabilize `main`. - ImapConnectionCache: keeps one authenticated Store alive per account, guarded by a per-account mutex, keyed by connection identity (not the rotating secret), with lazy catch-and-retry-once stale handling. No eviction policy yet beyond an explicit closeReusedConnections() hook. - ImapClient gains a `reuseConnections` flag (default false via the @Inject no-arg constructor). With it off, withStore is byte-for-byte the previous connect + LOGOUT-per-call; with it on, calls borrow the kept-alive Store. - ImapFolderOpenLatencyTest flips the flag on: the same real-IMAP operations that cost N connections / N LOGINs collapse to 1 connection / 1 LOGIN, with the necessary per-open EXAMINE unchanged (proven via CountingImapProxy + GreenMail; localhost is ~0 RTT so this proves structure, not wall-clock). - docs/perf/issue-125-connection-reuse-spike.md: prototype design, the flag-off-vs-on proof, per-decision trade-offs, and the refined real-device validation plan. References #125; does not close it (needs device validation). Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
@@ -86,7 +86,23 @@ data class ReplyContext(
|
||||
|
||||
/** Thin IMAP client over Jakarta/Angus Mail. Supports password and XOAUTH2 auth. */
|
||||
@Singleton
|
||||
class ImapClient @Inject constructor() {
|
||||
class ImapClient(private val reuseConnections: Boolean) {
|
||||
|
||||
/**
|
||||
* Production entry point. Connection reuse is a SPIKE flag (issue #125), **OFF by default** so it
|
||||
* cannot destabilize the connect-per-operation behaviour on `main`: with it off, [withStore] is
|
||||
* byte-for-byte today's connect + LOGOUT-per-call. Once real-device validation (see
|
||||
* `docs/perf/issue-125-connection-reuse-spike.md`) confirms the win, wire this to a setting or
|
||||
* `BuildConfig`; today only the reuse harness flips it on via the primary constructor.
|
||||
*/
|
||||
@Inject constructor() : this(reuseConnections = false)
|
||||
|
||||
/**
|
||||
* Per-account keep-alive cache; allocated only when the spike flag is on, so the default build
|
||||
* carries neither the state nor the reuse code path.
|
||||
*/
|
||||
private val connectionCache: ImapConnectionCache? =
|
||||
if (reuseConnections) ImapConnectionCache(::openConnectedStore) else null
|
||||
|
||||
/** Connects and returns the account's folders with their SPECIAL-USE attributes. Throws on failure. */
|
||||
suspend fun listFolders(params: ImapConnectionParams): List<FetchedFolder> = withContext(Dispatchers.IO) {
|
||||
@@ -530,37 +546,71 @@ class ImapClient @Inject constructor() {
|
||||
)
|
||||
}
|
||||
|
||||
private inline fun <T> withStore(params: ImapConnectionParams, block: (Store) -> T): T {
|
||||
val protocol = if (params.security == MailSecurity.SSL_TLS) "imaps" else "imap"
|
||||
val store = Session.getInstance(buildProps(protocol, params)).getStore(protocol)
|
||||
store.connect(params.host, params.port, params.username, params.secret)
|
||||
return try {
|
||||
block(store)
|
||||
} finally {
|
||||
runCatching { store.close() }
|
||||
/**
|
||||
* Runs [block] against a connected [Store]. With the reuse flag OFF (default) this is the original
|
||||
* behaviour: a fresh, authenticated connection per call, torn down in `finally`. With it ON, the
|
||||
* call borrows a kept-alive per-account connection from [connectionCache] (established once, reused
|
||||
* across folder-opens) instead — see issue #125.
|
||||
*/
|
||||
private suspend fun <T> withStore(params: ImapConnectionParams, block: (Store) -> T): T {
|
||||
val cache = connectionCache
|
||||
return if (cache != null) {
|
||||
cache.withStore(params, block)
|
||||
} else {
|
||||
val store = openConnectedStore(params)
|
||||
try {
|
||||
block(store)
|
||||
} finally {
|
||||
runCatching { store.close() }
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
private fun buildProps(protocol: String, params: ImapConnectionParams): Properties = Properties().apply {
|
||||
put("mail.store.protocol", protocol)
|
||||
put("mail.$protocol.host", params.host)
|
||||
put("mail.$protocol.port", params.port.toString())
|
||||
put("mail.$protocol.connectiontimeout", TIMEOUT_MS)
|
||||
put("mail.$protocol.timeout", TIMEOUT_MS)
|
||||
put("mail.$protocol.writetimeout", TIMEOUT_MS)
|
||||
if (params.security == MailSecurity.STARTTLS) {
|
||||
put("mail.$protocol.starttls.enable", "true")
|
||||
put("mail.$protocol.starttls.required", params.strictStartTls.toString())
|
||||
}
|
||||
// Verify the server certificate matches the host whenever TLS is used. Angus already
|
||||
// defaults this to true; set it explicitly so a future library-default change can't
|
||||
// silently disable hostname checking and expose us to MITM. (No-op for MailSecurity.NONE.)
|
||||
put("mail.$protocol.ssl.checkserveridentity", "true")
|
||||
if (params.useXoauth2) {
|
||||
put("mail.$protocol.auth.mechanisms", "XOAUTH2")
|
||||
}
|
||||
/** Builds and authenticates a fresh [Store] (`CONNECT + TLS + LOGIN`); the caller owns closing it. */
|
||||
private fun openConnectedStore(params: ImapConnectionParams): Store {
|
||||
val protocol = if (params.security == MailSecurity.SSL_TLS) "imaps" else "imap"
|
||||
val store = Session.getInstance(buildProps(protocol, params, reuse = reuseConnections)).getStore(protocol)
|
||||
store.connect(params.host, params.port, params.username, params.secret)
|
||||
return store
|
||||
}
|
||||
|
||||
/**
|
||||
* SPIKE hook (issue #125): closes any kept-alive reused connections (`LOGOUT` + teardown), a no-op
|
||||
* when the reuse flag is OFF. The reuse harness calls this to force settlement; a shipped feature
|
||||
* would also drive it from an idle-eviction timer and the low-battery push teardown (#88/#89/#90).
|
||||
*/
|
||||
suspend fun closeReusedConnections() {
|
||||
connectionCache?.closeAll()
|
||||
}
|
||||
|
||||
private fun buildProps(protocol: String, params: ImapConnectionParams, reuse: Boolean = false): Properties =
|
||||
Properties().apply {
|
||||
put("mail.store.protocol", protocol)
|
||||
put("mail.$protocol.host", params.host)
|
||||
put("mail.$protocol.port", params.port.toString())
|
||||
put("mail.$protocol.connectiontimeout", TIMEOUT_MS)
|
||||
put("mail.$protocol.timeout", TIMEOUT_MS)
|
||||
put("mail.$protocol.writetimeout", TIMEOUT_MS)
|
||||
if (params.security == MailSecurity.STARTTLS) {
|
||||
put("mail.$protocol.starttls.enable", "true")
|
||||
put("mail.$protocol.starttls.required", params.strictStartTls.toString())
|
||||
}
|
||||
// Verify the server certificate matches the host whenever TLS is used. Angus already
|
||||
// defaults this to true; set it explicitly so a future library-default change can't
|
||||
// silently disable hostname checking and expose us to MITM. (No-op for MailSecurity.NONE.)
|
||||
put("mail.$protocol.ssl.checkserveridentity", "true")
|
||||
if (params.useXoauth2) {
|
||||
put("mail.$protocol.auth.mechanisms", "XOAUTH2")
|
||||
}
|
||||
if (reuse) {
|
||||
// SPIKE (issue #125): pin a reused Store to exactly one authenticated socket. The
|
||||
// per-account mutex already serializes access, so one pooled connection suffices — and
|
||||
// capping it here makes "1 connection for N opens" the literal, provable invariant.
|
||||
put("mail.$protocol.connectionpoolsize", "1")
|
||||
put("mail.$protocol.separatestoreconnection", "false")
|
||||
}
|
||||
}
|
||||
|
||||
private companion object {
|
||||
const val TIMEOUT_MS = "15000"
|
||||
const val TAG = "LibreMailIdle"
|
||||
|
||||
@@ -0,0 +1,111 @@
|
||||
// SPDX-License-Identifier: GPL-3.0-or-later
|
||||
package org.libremail.mail
|
||||
|
||||
import jakarta.mail.FolderClosedException
|
||||
import jakarta.mail.MessagingException
|
||||
import jakarta.mail.Store
|
||||
import jakarta.mail.StoreClosedException
|
||||
import kotlinx.coroutines.sync.Mutex
|
||||
import kotlinx.coroutines.sync.withLock
|
||||
import org.libremail.domain.model.ImapConnectionParams
|
||||
import java.io.IOException
|
||||
import java.util.concurrent.ConcurrentHashMap
|
||||
|
||||
/**
|
||||
* SPIKE (issue #125): a per-account keep-alive cache of authenticated IMAP [Store]s, so folder-opens
|
||||
* and message operations reuse one already-connected session instead of re-paying
|
||||
* `CONNECT + TLS + LOGIN` on every call. See `docs/perf/issue-125-connection-reuse-spike.md`.
|
||||
*
|
||||
* Prototype stance — deliberately the simplest thing that *proves reuse*, leaving the tuning knobs to a
|
||||
* measured follow-up:
|
||||
* - **One connection per account, mutex-guarded.** Each account key holds a single [Store] behind its
|
||||
* own [Mutex]; every operation on that account serializes through it. This is the simplest safe
|
||||
* design and the one the investigation named as the starting point. Its known cost is head-of-line
|
||||
* blocking — a quick flag toggle can queue behind a slow body download. A bounded pool would trade
|
||||
* that for more sockets (and a size cap + eviction); not prototyped here.
|
||||
* - **Lazy, catch-and-retry-once stale handling.** No periodic `NOOP` probe (that would add a
|
||||
* round-trip to every reused op, partly defeating the point). An operation runs optimistically; if
|
||||
* it fails with a dropped-connection signal, the socket is rebuilt once and the operation retried.
|
||||
* - **Keyed by connection identity, not the secret.** The OAuth access token
|
||||
* ([ImapConnectionParams.secret]) rotates; keying on host/port/user/security/mechanism keeps a token
|
||||
* refresh from orphaning a live, already-authenticated socket. A refreshed secret only matters when
|
||||
* we actually reconnect, and [connect] is always handed the current [params].
|
||||
*
|
||||
* Thread-safety: [ImapClient]'s UI operations are not otherwise serialized and prefetch runs outside
|
||||
* the syncer's mutex, so [withStore] must be safe under concurrent callers for the same account — the
|
||||
* per-key mutex provides that. Not wired to any lifecycle/battery signal yet: [closeAll] is the only
|
||||
* eviction and is driven by the harness today; an idle-eviction timer and low-battery teardown
|
||||
* (#88/#89/#90) are follow-ups.
|
||||
*
|
||||
* @param connect builds and authenticates a fresh [Store] for the given params (blocking network I/O).
|
||||
*/
|
||||
internal class ImapConnectionCache(private val connect: (ImapConnectionParams) -> Store) {
|
||||
|
||||
private class Entry {
|
||||
val mutex = Mutex()
|
||||
var store: Store? = null
|
||||
}
|
||||
|
||||
private val entries = ConcurrentHashMap<String, Entry>()
|
||||
|
||||
/**
|
||||
* Runs [block] against a reused, authenticated [Store] for [params]'s account: it is established on
|
||||
* first use and kept open afterwards, so only the first call pays connection setup. Serialized per
|
||||
* account by the key's [Mutex]. If the operation hits a dropped connection the socket is rebuilt
|
||||
* once and the operation retried; a second failure clears the slot so the next call reconnects.
|
||||
*/
|
||||
suspend fun <T> withStore(params: ImapConnectionParams, block: (Store) -> T): T {
|
||||
val entry = entries.computeIfAbsent(key(params)) { Entry() }
|
||||
return entry.mutex.withLock {
|
||||
val store = entry.store ?: connect(params).also { entry.store = it }
|
||||
try {
|
||||
block(store)
|
||||
} catch (e: Throwable) {
|
||||
if (!isConnectionDrop(e)) throw e
|
||||
// Stale socket (server idle-timeout, NAT rebind, network change): rebuild once and retry.
|
||||
runCatching { store.close() }
|
||||
// Forget the dead socket before reconnecting, so a failed connect leaves a clean slot.
|
||||
entry.store = null
|
||||
val fresh = connect(params)
|
||||
entry.store = fresh
|
||||
try {
|
||||
block(fresh)
|
||||
} catch (retry: Throwable) {
|
||||
runCatching { fresh.close() }
|
||||
entry.store = null
|
||||
throw retry
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
/** Closes and forgets every cached connection (`LOGOUT` + socket teardown). The only eviction today. */
|
||||
suspend fun closeAll() {
|
||||
for ((_, entry) in entries) {
|
||||
entry.mutex.withLock {
|
||||
entry.store?.let { store -> runCatching { store.close() } }
|
||||
entry.store = null
|
||||
}
|
||||
}
|
||||
entries.clear()
|
||||
}
|
||||
|
||||
/**
|
||||
* Account identity for reuse — everything that pins a distinct authenticated socket EXCEPT the
|
||||
* secret, so a rotated OAuth token reuses the same live connection instead of orphaning it.
|
||||
*/
|
||||
private fun key(params: ImapConnectionParams): String =
|
||||
"${params.host}|${params.port}|${params.security}|${params.username}|${params.useXoauth2}"
|
||||
|
||||
/**
|
||||
* Whether [error] signals a dropped connection (retry on a fresh socket) rather than a genuine
|
||||
* protocol/application error (propagate as-is). Deliberately narrow: a plain [MessagingException]
|
||||
* for a real server error whose connection is still live is NOT retried, so we never re-issue a
|
||||
* mutation over a working connection.
|
||||
*/
|
||||
private fun isConnectionDrop(error: Throwable): Boolean = when (error) {
|
||||
is FolderClosedException, is StoreClosedException, is IOException -> true
|
||||
is MessagingException -> error.cause is IOException
|
||||
else -> false
|
||||
}
|
||||
}
|
||||
@@ -4,6 +4,7 @@ package org.libremail.mail
|
||||
import com.icegreen.greenmail.util.GreenMail
|
||||
import com.icegreen.greenmail.util.GreenMailUtil
|
||||
import com.icegreen.greenmail.util.ServerSetupTest
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.test.runTest
|
||||
import org.junit.After
|
||||
import org.junit.Before
|
||||
@@ -33,8 +34,13 @@ class ImapFolderOpenLatencyTest {
|
||||
|
||||
private lateinit var greenMail: GreenMail
|
||||
private lateinit var proxy: CountingImapProxy
|
||||
|
||||
/** Flag OFF (production default): a fresh connect + LOGOUT per operation. */
|
||||
private val client = ImapClient()
|
||||
|
||||
/** Flag ON (the spike prototype, issue #125): one kept-alive connection reused across operations. */
|
||||
private val reuseClient = ImapClient(reuseConnections = true)
|
||||
|
||||
@Before
|
||||
fun setUp() {
|
||||
greenMail = GreenMail(ServerSetupTest.SMTP_IMAP)
|
||||
@@ -46,6 +52,7 @@ class ImapFolderOpenLatencyTest {
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
runBlocking { reuseClient.closeReusedConnections() } // release any kept-alive socket before the server stops
|
||||
proxy.close()
|
||||
greenMail.stop()
|
||||
}
|
||||
@@ -119,6 +126,46 @@ class ImapFolderOpenLatencyTest {
|
||||
assertEquals(2, proxy.authCommandCount(), "list + read each pay a full LOGIN")
|
||||
}
|
||||
|
||||
// --- Flag ON: the spike prototype reuses one connection across operations (issue #125). ---
|
||||
// These are the deterministic proof that reuse works: the SAME real-IMAP operations that cost N
|
||||
// connections / N LOGINs above collapse to ONE connection / ONE LOGIN here, with the necessary
|
||||
// per-open EXAMINE unchanged. Localhost is ~0 RTT so this proves the STRUCTURE, not wall-clock.
|
||||
|
||||
@Test
|
||||
fun `with reuse on, N folder-opens share one connection and one LOGIN`() = runTest {
|
||||
seedInbox(2)
|
||||
|
||||
repeat(OPENS) { reuseClient.fetchRecent(params(), "INBOX", limit = 50) }
|
||||
|
||||
// The win: one socket accepted for all OPENS opens (vs. OPENS sockets with the flag off).
|
||||
// connectionCount is incremented synchronously on accept, so it needs no stream settling.
|
||||
assertEquals(1, proxy.connectionCount, "reuse: a single TCP connection serves every folder-open")
|
||||
|
||||
reuseClient.closeReusedConnections() // evict -> LOGOUT + close, so the proxy's command stream settles
|
||||
proxy.awaitClientStreamsSettled()
|
||||
|
||||
assertEquals(1, proxy.authCommandCount(), "reuse: LOGIN paid once, then reused (vs. one per open)")
|
||||
assertEquals(OPENS, proxy.commandCount("EXAMINE"), "reuse keeps the necessary one EXAMINE per open")
|
||||
assertEquals(1, proxy.commandCount("LOGOUT"), "reuse: one LOGOUT at eviction, not one per open")
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `with reuse on, opening a folder then reading a message reuses the one connection`() = runTest {
|
||||
seedInbox(1)
|
||||
|
||||
val uid = reuseClient.fetchRecent(params(), "INBOX", limit = 50).first().uid // open folder
|
||||
reuseClient.fetchBodyMarkingSeen(params(), "INBOX", uid) // read a message in it
|
||||
|
||||
// Contrast with `opening a folder then reading a message uses two separate connections`: the
|
||||
// same list-then-read here shares the one kept-alive connection instead of paying a second setup.
|
||||
assertEquals(1, proxy.connectionCount, "reuse: list + read share the one connection")
|
||||
|
||||
reuseClient.closeReusedConnections()
|
||||
proxy.awaitClientStreamsSettled()
|
||||
|
||||
assertEquals(1, proxy.authCommandCount(), "reuse: one LOGIN covers both the list and the read")
|
||||
}
|
||||
|
||||
private companion object {
|
||||
const val OPENS = 3
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user