Merge fix/standard-route-network-switchover: standard-route network auto-switch (items 1-3)
This commit is contained in:
@@ -181,10 +181,34 @@ class ConnectionManager(
|
||||
|
||||
private var networkCallback: ConnectivityManager.NetworkCallback? = null
|
||||
|
||||
/**
|
||||
* Debounce job for network-change re-resolution. Android fires one
|
||||
* onAvailable per satisfying network (Wi-Fi + cell + VPN can land within
|
||||
* milliseconds of each other, and registration itself replays every
|
||||
* current network), so each event cancels the previous pending resolve
|
||||
* and the last one wins after a short settle window.
|
||||
*/
|
||||
@Volatile
|
||||
private var networkResolveJob: kotlinx.coroutines.Job? = null
|
||||
|
||||
init {
|
||||
// Register at construction, not on first connect(). Standard
|
||||
// (no-Relay) connections never open the WSS socket, but their HTTP
|
||||
// surfaces (chat, dashboard, voice) still need [activeEndpoint] to
|
||||
// follow LAN/Tailscale handoffs — leaving registration inside
|
||||
// connect() left the whole ADR 24 network-aware path dormant for
|
||||
// exactly those users. No-op when [context] is null (tests).
|
||||
ensureNetworkCallbackRegistered()
|
||||
}
|
||||
|
||||
companion object {
|
||||
private const val TAG = "ConnectionManager"
|
||||
private const val MAX_BACKOFF_MS = 30_000L
|
||||
private const val BASE_BACKOFF_MS = 1_000L
|
||||
// Settle window before re-resolving after a network event. Long
|
||||
// enough to coalesce the onAvailable burst of a handoff, short
|
||||
// enough that a route swap still feels immediate.
|
||||
private const val NETWORK_RESOLVE_DEBOUNCE_MS = 300L
|
||||
// Matches plugin.relay.auth._BLOCK_SECONDS (5 min). If we see 429
|
||||
// on the WSS upgrade, we're IP-banned server-side — retrying at
|
||||
// our normal 1-30s cadence re-fills the ban bucket and keeps us
|
||||
@@ -425,8 +449,14 @@ class ConnectionManager(
|
||||
* WSS reconnect. Used by HTTP-only surfaces (chat/voice/relay HTTP)
|
||||
* so they can follow LAN/Tailscale/VPN route changes even when the relay
|
||||
* socket is currently disconnected or intentionally not paired.
|
||||
*
|
||||
* @param clearProbeCache wipe the resolver's probe cache first. Pass
|
||||
* `true` from "the world may have changed" triggers (app resume,
|
||||
* network change) — otherwise a route that died within the positive
|
||||
* cache TTL (60s) can still be returned as the winner.
|
||||
*/
|
||||
suspend fun refreshActiveEndpoint(): EndpointCandidate? {
|
||||
suspend fun refreshActiveEndpoint(clearProbeCache: Boolean = false): EndpointCandidate? {
|
||||
if (clearProbeCache) endpointResolver?.clearCache()
|
||||
val resolved = resolveBestEndpointSafe()
|
||||
_activeEndpoint.value = resolved
|
||||
return resolved
|
||||
@@ -450,26 +480,42 @@ class ConnectionManager(
|
||||
Log.i(TAG, "marked endpoint role=${active.role} unreachable ($reason)")
|
||||
}
|
||||
|
||||
private fun resolveAndSwitchIfNeeded(closeReason: String) {
|
||||
/**
|
||||
* Debounced network-change re-resolution, shared by both NetworkCallback
|
||||
* events. Re-runs the resolver and publishes the winner to
|
||||
* [activeEndpoint] so HTTP-only surfaces (chat, dashboard, standard
|
||||
* voice) follow the route change even when no relay socket exists. When
|
||||
* a socket IS up, additionally swaps it to a differing winner, or
|
||||
* reconnects a disconnected socket on the same winner — preserving the
|
||||
* pre-refactor relay-path behavior.
|
||||
*/
|
||||
private fun scheduleNetworkReResolve(closeReason: String) {
|
||||
if (endpointResolver == null) return
|
||||
val current = serverUrl ?: return
|
||||
scope.launch {
|
||||
networkResolveJob?.cancel()
|
||||
networkResolveJob = scope.launch {
|
||||
delay(NETWORK_RESOLVE_DEBOUNCE_MS)
|
||||
val current = serverUrl
|
||||
val resolved = resolveBestEndpointSafe()
|
||||
if (resolved == null) {
|
||||
_activeEndpoint.value = null
|
||||
// Don't clear a live socket's endpoint on a transient probe
|
||||
// miss — only drop the published route when nothing is
|
||||
// actually connected.
|
||||
if (_connectionState.value != ConnectionState.Connected) {
|
||||
_activeEndpoint.value = null
|
||||
}
|
||||
return@launch
|
||||
}
|
||||
val newUrl = resolved.relay.url
|
||||
val normalizedNew = normalizeRelayUrl(newUrl)
|
||||
_activeEndpoint.value = resolved
|
||||
if (current == null) return@launch
|
||||
val normalizedNew = normalizeRelayUrl(resolved.relay.url)
|
||||
if (normalizedNew != current) {
|
||||
Log.i(TAG, "endpoint fallback: swapping $current → $normalizedNew")
|
||||
connectToUrlOnMainPath(newUrl, closeReason)
|
||||
Log.i(TAG, "network change: swapping $current → $normalizedNew")
|
||||
connectToUrlOnMainPath(resolved.relay.url, closeReason)
|
||||
} else if (_connectionState.value == ConnectionState.Disconnected &&
|
||||
shouldReconnect &&
|
||||
reconnectGate()
|
||||
) {
|
||||
Log.i(TAG, "endpoint fallback: same winner is disconnected — reconnecting $current")
|
||||
Log.i(TAG, "network change: same winner is disconnected — reconnecting $current")
|
||||
doConnect(current)
|
||||
}
|
||||
}
|
||||
@@ -482,36 +528,15 @@ class ConnectionManager(
|
||||
val callback = object : ConnectivityManager.NetworkCallback() {
|
||||
override fun onAvailable(network: Network) {
|
||||
Log.i(TAG, "network onAvailable — re-evaluating endpoint")
|
||||
if (endpointResolver == null) return
|
||||
val url = serverUrl ?: return
|
||||
endpointResolver.clearCache()
|
||||
scope.launch {
|
||||
val resolved = resolveBestEndpointSafe()
|
||||
val newUrl = resolved?.relay?.url
|
||||
if (newUrl == null) {
|
||||
if (_connectionState.value != ConnectionState.Connected) {
|
||||
_activeEndpoint.value = null
|
||||
}
|
||||
return@launch
|
||||
}
|
||||
val normalizedNew = normalizeRelayUrl(newUrl)
|
||||
_activeEndpoint.value = resolved
|
||||
// Only swap if the winner actually differs from the
|
||||
// currently-connected URL. Avoids dropping a healthy
|
||||
// socket on a no-op network flap (Wi-Fi scan, cell
|
||||
// handover that ends up on the same route, etc.).
|
||||
if (normalizedNew != url) {
|
||||
Log.i(TAG, "network change: swapping $url → $normalizedNew")
|
||||
connectToUrlOnMainPath(newUrl, "Network change — switching endpoint")
|
||||
}
|
||||
}
|
||||
endpointResolver?.clearCache()
|
||||
scheduleNetworkReResolve("Network change — switching endpoint")
|
||||
}
|
||||
|
||||
override fun onLost(network: Network) {
|
||||
Log.i(TAG, "network onLost — marking active endpoint unreachable and resolving fallback")
|
||||
endpointResolver?.clearCache()
|
||||
markActiveEndpointUnreachable("network lost")
|
||||
resolveAndSwitchIfNeeded("Network lost — switching endpoint")
|
||||
scheduleNetworkReResolve("Network lost — switching endpoint")
|
||||
}
|
||||
}
|
||||
try {
|
||||
|
||||
@@ -331,7 +331,13 @@ class EndpointResolver(
|
||||
)
|
||||
}
|
||||
|
||||
/** Test-only: wipe the probe cache so a fresh run starts clean. */
|
||||
/**
|
||||
* Wipe the probe cache so the next resolve runs fresh probes. Called on
|
||||
* "the world changed" triggers — NetworkCallback events, manual "Probe
|
||||
* now", and [refreshActiveEndpoint][ConnectionManager.refreshActiveEndpoint]
|
||||
* with `clearProbeCache = true` — where a positive entry for a
|
||||
* just-died route must not outlive the handoff.
|
||||
*/
|
||||
internal fun clearCache() {
|
||||
probeCache.clear()
|
||||
}
|
||||
|
||||
@@ -514,14 +514,17 @@ class ConnectionViewModel(application: Application) : AndroidViewModel(applicati
|
||||
private val dashboardCookieStores =
|
||||
java.util.concurrent.ConcurrentHashMap<String, EncryptedDashboardCookieStore>()
|
||||
|
||||
/** Resolved dashboard URL of the active connection (explicit or derived :9119). */
|
||||
fun activeDashboardUrl(): String? {
|
||||
val connectionId = connectionStore.activeConnectionId.value ?: return null
|
||||
return connectionStore.connections.value
|
||||
.firstOrNull { it.id == connectionId }
|
||||
?.resolvedDashboardUrl
|
||||
?.takeIf { it.isNotBlank() }
|
||||
}
|
||||
/**
|
||||
* Dashboard URL for the active connection **on the currently-resolved
|
||||
* route** — snapshot twin of [effectiveDashboardUrl], which it delegates
|
||||
* to. Standard voice and the availability probe read this per call, so
|
||||
* an auto-managed dashboard URL follows LAN/Tailscale handoffs the same
|
||||
* way Manage does; an explicit dashboard override stays pinned. (This
|
||||
* used to read the persisted `resolvedDashboardUrl`, which kept voice
|
||||
* aimed at the LAN host after the resolver had moved chat to Tailscale.)
|
||||
*/
|
||||
fun activeDashboardUrl(): String? =
|
||||
effectiveDashboardUrl.value.takeIf { it.isNotBlank() }
|
||||
|
||||
/**
|
||||
* Cookie store for the active connection — the same encrypted store the
|
||||
@@ -2721,7 +2724,10 @@ class ConnectionViewModel(application: Application) : AndroidViewModel(applicati
|
||||
if (revalidationJob?.isActive == true) return
|
||||
revalidationJob = viewModelScope.launch {
|
||||
val apiRouteBefore = effectiveApiServerUrlSnapshot()
|
||||
connectionManager.refreshActiveEndpoint()
|
||||
// Clear the probe cache: revalidate() fires on resume / network
|
||||
// change, where a cached-reachable entry for the route we just
|
||||
// walked away from would win the resolve for up to 60s.
|
||||
connectionManager.refreshActiveEndpoint(clearProbeCache = true)
|
||||
if (effectiveApiServerUrlSnapshot() != apiRouteBefore) {
|
||||
rebuildApiClient()
|
||||
}
|
||||
|
||||
@@ -0,0 +1,163 @@
|
||||
package com.hermesandroid.relay.network
|
||||
|
||||
import android.content.Context
|
||||
import android.net.ConnectivityManager
|
||||
import com.hermesandroid.relay.data.ApiEndpoint
|
||||
import com.hermesandroid.relay.data.EndpointCandidate
|
||||
import com.hermesandroid.relay.data.RelayEndpoint
|
||||
import kotlinx.coroutines.flow.first
|
||||
import kotlinx.coroutines.runBlocking
|
||||
import kotlinx.coroutines.withTimeout
|
||||
import okhttp3.OkHttpClient
|
||||
import okhttp3.mockwebserver.Dispatcher
|
||||
import okhttp3.mockwebserver.MockResponse
|
||||
import okhttp3.mockwebserver.MockWebServer
|
||||
import okhttp3.mockwebserver.RecordedRequest
|
||||
import org.junit.After
|
||||
import org.junit.Assert.assertEquals
|
||||
import org.junit.Assert.assertNotNull
|
||||
import org.junit.Assert.assertNull
|
||||
import org.junit.Before
|
||||
import org.junit.Test
|
||||
import org.junit.runner.RunWith
|
||||
import org.robolectric.RobolectricTestRunner
|
||||
import org.robolectric.RuntimeEnvironment
|
||||
import org.robolectric.Shadows.shadowOf
|
||||
import org.robolectric.annotation.Config
|
||||
import org.robolectric.shadows.ShadowNetwork
|
||||
|
||||
/**
|
||||
* Standard-route (no relay socket) coverage for [ConnectionManager]'s ADR 24
|
||||
* network-aware switching, added when the machinery was decoupled from the
|
||||
* WSS connect path:
|
||||
*
|
||||
* 1. The [ConnectivityManager.NetworkCallback] registers at construction —
|
||||
* previously only `connect()` registered it, so standard (no-Relay)
|
||||
* connections never saw network changes at all.
|
||||
* 2. A network `onAvailable` with **no relay socket** still re-resolves and
|
||||
* publishes [ConnectionManager.activeEndpoint], which is what the
|
||||
* HTTP-only surfaces (chat, dashboard, standard voice) follow.
|
||||
* 3. `refreshActiveEndpoint(clearProbeCache = true)` forgets a
|
||||
* cached-reachable route so a just-died endpoint can't win the resolve
|
||||
* for the remainder of the 60s positive cache TTL.
|
||||
*
|
||||
* Runs under Robolectric for ConnectivityManager + a real [MockWebServer]
|
||||
* for the resolver's `HEAD /health` probes — same probe contract as
|
||||
* [EndpointResolverTest].
|
||||
*/
|
||||
@RunWith(RobolectricTestRunner::class)
|
||||
@Config(sdk = [34])
|
||||
class ConnectionManagerRouteTest {
|
||||
|
||||
private lateinit var context: Context
|
||||
private lateinit var connectivityManager: ConnectivityManager
|
||||
private lateinit var server: MockWebServer
|
||||
private val managers = mutableListOf<ConnectionManager>()
|
||||
|
||||
@Before
|
||||
fun setUp() {
|
||||
context = RuntimeEnvironment.getApplication()
|
||||
connectivityManager = context.getSystemService(ConnectivityManager::class.java)!!
|
||||
server = MockWebServer()
|
||||
server.dispatcher = object : Dispatcher() {
|
||||
override fun dispatch(request: RecordedRequest): MockResponse =
|
||||
if (request.path?.endsWith("/health") == true) {
|
||||
MockResponse().setResponseCode(200)
|
||||
} else {
|
||||
MockResponse().setResponseCode(404)
|
||||
}
|
||||
}
|
||||
server.start()
|
||||
}
|
||||
|
||||
@After
|
||||
fun tearDown() {
|
||||
managers.forEach { runCatching { it.shutdown() } }
|
||||
managers.clear()
|
||||
runCatching { server.shutdown() }
|
||||
}
|
||||
|
||||
private fun registeredCallbacks(): Set<ConnectivityManager.NetworkCallback> =
|
||||
shadowOf(connectivityManager).networkCallbacks.toSet()
|
||||
|
||||
private fun candidate(role: String = "tailscale"): EndpointCandidate =
|
||||
EndpointCandidate(
|
||||
role = role,
|
||||
priority = 0,
|
||||
api = ApiEndpoint(host = server.hostName, port = server.port, tls = false),
|
||||
relay = RelayEndpoint(url = "ws://${server.hostName}:${server.port}"),
|
||||
)
|
||||
|
||||
private fun buildManager(
|
||||
candidates: () -> List<EndpointCandidate>,
|
||||
): ConnectionManager = ConnectionManager(
|
||||
ChannelMultiplexer(),
|
||||
context = context,
|
||||
endpointResolver = EndpointResolver(httpClient = OkHttpClient()),
|
||||
endpointCandidatesProvider = { candidates() },
|
||||
).also { managers.add(it) }
|
||||
|
||||
@Test
|
||||
fun `network callback registers at construction without connect`() {
|
||||
val before = registeredCallbacks()
|
||||
|
||||
buildManager { emptyList() }
|
||||
|
||||
assertEquals(
|
||||
"ConnectionManager must register its NetworkCallback at construction " +
|
||||
"so standard (no-Relay) connections follow network changes",
|
||||
before.size + 1,
|
||||
registeredCallbacks().size,
|
||||
)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `shutdown unregisters the construction-time network callback`() {
|
||||
val before = registeredCallbacks()
|
||||
val manager = buildManager { emptyList() }
|
||||
|
||||
manager.shutdown()
|
||||
|
||||
assertEquals(before, registeredCallbacks())
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `onAvailable with no relay socket re-resolves and publishes activeEndpoint`() {
|
||||
val before = registeredCallbacks()
|
||||
val manager = buildManager { listOf(candidate()) }
|
||||
assertNull("no endpoint should be published before any trigger", manager.activeEndpoint.value)
|
||||
|
||||
val callback = (registeredCallbacks() - before).single()
|
||||
callback.onAvailable(ShadowNetwork.newInstance(101))
|
||||
|
||||
val published = runBlocking {
|
||||
withTimeout(10_000) { manager.activeEndpoint.first { it != null } }
|
||||
}
|
||||
assertEquals("tailscale", published!!.role)
|
||||
}
|
||||
|
||||
@Test
|
||||
fun `refreshActiveEndpoint returns stale cached winner unless clearProbeCache`() {
|
||||
val manager = buildManager { listOf(candidate()) }
|
||||
|
||||
// Prime: server up → candidate resolves and its probe result caches.
|
||||
val first = runBlocking { manager.refreshActiveEndpoint() }
|
||||
assertNotNull("expected the live candidate to resolve", first)
|
||||
|
||||
// Route dies inside the positive cache TTL.
|
||||
server.shutdown()
|
||||
|
||||
// Without clearing, the 60s positive cache still vouches for it.
|
||||
val stale = runBlocking { manager.refreshActiveEndpoint() }
|
||||
assertEquals(
|
||||
"cached-reachable entry should still win within the TTL",
|
||||
"tailscale",
|
||||
stale?.role,
|
||||
)
|
||||
|
||||
// Clearing the cache forces a fresh probe, which now fails.
|
||||
val fresh = runBlocking { manager.refreshActiveEndpoint(clearProbeCache = true) }
|
||||
assertNull("fresh probe against the dead route must yield no winner", fresh)
|
||||
assertNull(manager.activeEndpoint.value)
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user