mirror of
https://github.com/permissionlesstech/bitchat-android.git
synced 2026-07-25 18:45:20 +00:00
wifi improvements
This commit is contained in:
@@ -671,6 +671,7 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic
|
|||||||
connectionManager.disableTransport()
|
connectionManager.disableTransport()
|
||||||
TransportBridgeService.unregister("BLE")
|
TransportBridgeService.unregister("BLE")
|
||||||
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { }
|
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { }
|
||||||
|
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("BLE") } catch (_: Exception) { }
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
if (terminated) {
|
if (terminated) {
|
||||||
@@ -715,6 +716,7 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic
|
|||||||
try { gossipSyncManager.stop() } catch (_: Exception) { }
|
try { gossipSyncManager.stop() } catch (_: Exception) { }
|
||||||
TransportBridgeService.unregister("BLE")
|
TransportBridgeService.unregister("BLE")
|
||||||
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { }
|
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { }
|
||||||
|
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("BLE") } catch (_: Exception) { }
|
||||||
connectionManager.disableTransport()
|
connectionManager.disableTransport()
|
||||||
try { peerManager.refreshPeerList() } catch (_: Exception) { }
|
try { peerManager.refreshPeerList() } catch (_: Exception) { }
|
||||||
}
|
}
|
||||||
@@ -734,6 +736,7 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic
|
|||||||
announceJob = null
|
announceJob = null
|
||||||
TransportBridgeService.unregister("BLE")
|
TransportBridgeService.unregister("BLE")
|
||||||
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { }
|
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("BLE") } catch (_: Exception) { }
|
||||||
|
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("BLE") } catch (_: Exception) { }
|
||||||
|
|
||||||
// Send leave announcement
|
// Send leave announcement
|
||||||
sendLeaveAnnouncement()
|
sendLeaveAnnouncement()
|
||||||
@@ -1224,8 +1227,14 @@ class BluetoothMeshService(private val context: Context) : TransportBridgeServic
|
|||||||
return try {
|
return try {
|
||||||
// Prefer verified peers that are currently marked as direct
|
// Prefer verified peers that are currently marked as direct
|
||||||
val verified = peerManager.getVerifiedPeers()
|
val verified = peerManager.getVerifiedPeers()
|
||||||
val direct = verified.filter { it.value.isDirectConnection }.keys.toList()
|
val direct = verified.filter { it.value.isDirectConnection }.keys.toSet()
|
||||||
direct.take(10)
|
// Publish this transport's direct peers and gossip the cross-transport union so a
|
||||||
|
// node connected via multiple transports advertises a complete neighbor list.
|
||||||
|
try { com.bitchat.android.services.AppStateStore.setTransportDirectPeers("BLE", direct) } catch (_: Exception) { }
|
||||||
|
val union = try {
|
||||||
|
com.bitchat.android.services.AppStateStore.getDirectPeers().ifEmpty { direct }
|
||||||
|
} catch (_: Exception) { direct }
|
||||||
|
union.distinct().take(10)
|
||||||
} catch (_: Exception) {
|
} catch (_: Exception) {
|
||||||
emptyList()
|
emptyList()
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -615,7 +615,14 @@ class MeshCore(
|
|||||||
val verifiedDirect = peerManager.getVerifiedPeers()
|
val verifiedDirect = peerManager.getVerifiedPeers()
|
||||||
.filter { it.value.isDirectConnection }
|
.filter { it.value.isDirectConnection }
|
||||||
.keys
|
.keys
|
||||||
(verifiedDirect + directPeers).distinct().take(10)
|
val localDirect = (verifiedDirect + directPeers).toSet()
|
||||||
|
// Publish this transport's direct peers and gossip the cross-transport union so a
|
||||||
|
// node connected via multiple transports advertises a complete neighbor list.
|
||||||
|
try { com.bitchat.android.services.AppStateStore.setTransportDirectPeers(transport.id, localDirect) } catch (_: Exception) { }
|
||||||
|
val union = try {
|
||||||
|
com.bitchat.android.services.AppStateStore.getDirectPeers().ifEmpty { localDirect }
|
||||||
|
} catch (_: Exception) { localDirect }
|
||||||
|
union.distinct().take(10)
|
||||||
} catch (_: Exception) {
|
} catch (_: Exception) {
|
||||||
directPeers.toList().take(10)
|
directPeers.toList().take(10)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -15,6 +15,8 @@ object AppStateStore {
|
|||||||
private val seenMessageIds = mutableSetOf<String>()
|
private val seenMessageIds = mutableSetOf<String>()
|
||||||
private val seenPublicMessageKeys = mutableSetOf<String>()
|
private val seenPublicMessageKeys = mutableSetOf<String>()
|
||||||
private val peerIdsByTransport = mutableMapOf<String, Set<String>>()
|
private val peerIdsByTransport = mutableMapOf<String, Set<String>>()
|
||||||
|
// Direct (single-hop) peer IDs per transport, used to gossip a unified neighbor set.
|
||||||
|
private val directPeerIdsByTransport = mutableMapOf<String, Set<String>>()
|
||||||
// Connected peer IDs (mesh ephemeral IDs)
|
// Connected peer IDs (mesh ephemeral IDs)
|
||||||
private val _peers = MutableStateFlow<List<String>>(emptyList())
|
private val _peers = MutableStateFlow<List<String>>(emptyList())
|
||||||
val peers: StateFlow<List<String>> = _peers.asStateFlow()
|
val peers: StateFlow<List<String>> = _peers.asStateFlow()
|
||||||
@@ -59,6 +61,30 @@ object AppStateStore {
|
|||||||
.toList()
|
.toList()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Record the set of direct (single-hop) peers reachable over a given transport. Each transport
|
||||||
|
* (BLE, Wi-Fi Aware, ...) only knows its own direct peers; [getDirectPeers] unions them so every
|
||||||
|
* transport can gossip the same complete neighbor list under our shared node identity.
|
||||||
|
*/
|
||||||
|
fun setTransportDirectPeers(transportId: String, ids: Collection<String>) {
|
||||||
|
synchronized(this) {
|
||||||
|
directPeerIdsByTransport[transportId] = ids.toSet()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fun clearTransportDirectPeers(transportId: String) {
|
||||||
|
synchronized(this) {
|
||||||
|
directPeerIdsByTransport.remove(transportId)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/** Union of direct peers across all transports. */
|
||||||
|
fun getDirectPeers(): Set<String> {
|
||||||
|
synchronized(this) {
|
||||||
|
return directPeerIdsByTransport.values.flatten().toSet()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
fun addPublicMessage(msg: BitchatMessage) {
|
fun addPublicMessage(msg: BitchatMessage) {
|
||||||
synchronized(this) {
|
synchronized(this) {
|
||||||
val publicKey = publicMessageKey(msg)
|
val publicKey = publicMessageKey(msg)
|
||||||
@@ -130,6 +156,7 @@ object AppStateStore {
|
|||||||
seenMessageIds.clear()
|
seenMessageIds.clear()
|
||||||
seenPublicMessageKeys.clear()
|
seenPublicMessageKeys.clear()
|
||||||
peerIdsByTransport.clear()
|
peerIdsByTransport.clear()
|
||||||
|
directPeerIdsByTransport.clear()
|
||||||
_peers.value = emptyList()
|
_peers.value = emptyList()
|
||||||
_publicMessages.value = emptyList()
|
_publicMessages.value = emptyList()
|
||||||
_privateMessages.value = emptyMap()
|
_privateMessages.value = emptyMap()
|
||||||
|
|||||||
@@ -12,7 +12,10 @@ import android.util.Log
|
|||||||
* A synchronized wrapper around a raw Socket that implements a framed protocol:
|
* A synchronized wrapper around a raw Socket that implements a framed protocol:
|
||||||
* [4 bytes length][N bytes payload]
|
* [4 bytes length][N bytes payload]
|
||||||
*/
|
*/
|
||||||
class SyncedSocket(val rawSocket: Socket) {
|
class SyncedSocket(
|
||||||
|
val rawSocket: Socket,
|
||||||
|
readTimeoutMs: Int = DEFAULT_READ_TIMEOUT_MS
|
||||||
|
) {
|
||||||
private val TAG = "SyncedSocket"
|
private val TAG = "SyncedSocket"
|
||||||
private val writeLock = ReentrantLock()
|
private val writeLock = ReentrantLock()
|
||||||
private val readLock = ReentrantLock()
|
private val readLock = ReentrantLock()
|
||||||
@@ -20,7 +23,17 @@ class SyncedSocket(val rawSocket: Socket) {
|
|||||||
private val inputStream: DataInputStream
|
private val inputStream: DataInputStream
|
||||||
private val outputStream: DataOutputStream
|
private val outputStream: DataOutputStream
|
||||||
|
|
||||||
|
companion object {
|
||||||
|
// Both peers exchange keep-alive frames every ~2s while connected, so a read that
|
||||||
|
// stalls well beyond that means the link is dead (half-open). Time out so the read
|
||||||
|
// loop can detect it and trigger disconnection instead of blocking forever.
|
||||||
|
const val DEFAULT_READ_TIMEOUT_MS = 15_000
|
||||||
|
}
|
||||||
|
|
||||||
init {
|
init {
|
||||||
|
// A read timeout converts dead/half-open connections into a SocketTimeoutException
|
||||||
|
// (an IOException) so read() returns null and the peer is cleaned up.
|
||||||
|
try { rawSocket.soTimeout = readTimeoutMs } catch (_: Exception) {}
|
||||||
// We wrap streams to create DataInput/Output helpers
|
// We wrap streams to create DataInput/Output helpers
|
||||||
inputStream = DataInputStream(rawSocket.getInputStream())
|
inputStream = DataInputStream(rawSocket.getInputStream())
|
||||||
outputStream = DataOutputStream(rawSocket.getOutputStream())
|
outputStream = DataOutputStream(rawSocket.getOutputStream())
|
||||||
|
|||||||
@@ -13,6 +13,7 @@ import kotlinx.coroutines.flow.StateFlow
|
|||||||
import kotlinx.coroutines.flow.asStateFlow
|
import kotlinx.coroutines.flow.asStateFlow
|
||||||
import kotlinx.coroutines.isActive
|
import kotlinx.coroutines.isActive
|
||||||
import kotlinx.coroutines.launch
|
import kotlinx.coroutines.launch
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean
|
||||||
|
|
||||||
/**
|
/**
|
||||||
* WifiAwareController manages lifecycle and debug surfacing for the WifiAwareMeshService.
|
* WifiAwareController manages lifecycle and debug surfacing for the WifiAwareMeshService.
|
||||||
@@ -20,11 +21,15 @@ import kotlinx.coroutines.launch
|
|||||||
*/
|
*/
|
||||||
object WifiAwareController {
|
object WifiAwareController {
|
||||||
private const val TAG = "WifiAwareController"
|
private const val TAG = "WifiAwareController"
|
||||||
|
private const val MAX_RESTART_ATTEMPTS = 15
|
||||||
|
private const val RESTART_RETRY_DELAY_MS = 2_000L
|
||||||
|
|
||||||
private var service: WifiAwareMeshService? = null
|
private var service: WifiAwareMeshService? = null
|
||||||
private var appContext: Context? = null
|
private var appContext: Context? = null
|
||||||
private val lifecycleLock = Any()
|
private val lifecycleLock = Any()
|
||||||
private var starting = false
|
private var starting = false
|
||||||
|
private val restartInFlight = AtomicBoolean(false)
|
||||||
|
private var awareReceiverRegistered = false
|
||||||
|
|
||||||
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
private val scope = CoroutineScope(Dispatchers.IO + SupervisorJob())
|
||||||
|
|
||||||
@@ -46,6 +51,7 @@ object WifiAwareController {
|
|||||||
|
|
||||||
fun initialize(context: Context, enabledByDefault: Boolean) {
|
fun initialize(context: Context, enabledByDefault: Boolean) {
|
||||||
appContext = context.applicationContext
|
appContext = context.applicationContext
|
||||||
|
registerAwareStateReceiver(appContext!!)
|
||||||
setEnabled(enabledByDefault)
|
setEnabled(enabledByDefault)
|
||||||
// Start background poller for debug surfacing
|
// Start background poller for debug surfacing
|
||||||
scope.launch {
|
scope.launch {
|
||||||
@@ -193,10 +199,62 @@ object WifiAwareController {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Schedules a restart of the Wi-Fi Aware transport. Concurrent requests are coalesced into
|
||||||
|
* a single in-flight loop that retries with backoff. This is important because a single fixed
|
||||||
|
* delay can land while the service is still tearing down (recoveryInProgress), in which case
|
||||||
|
* startServices() defers and we must try again rather than give up.
|
||||||
|
*/
|
||||||
internal fun restartIfStillEnabled(delayMs: Long = 0L) {
|
internal fun restartIfStillEnabled(delayMs: Long = 0L) {
|
||||||
|
if (!restartInFlight.compareAndSet(false, true)) {
|
||||||
|
Log.d(TAG, "Restart already in flight; coalescing request")
|
||||||
|
return
|
||||||
|
}
|
||||||
scope.launch {
|
scope.launch {
|
||||||
|
try {
|
||||||
if (delayMs > 0L) delay(delayMs)
|
if (delayMs > 0L) delay(delayMs)
|
||||||
if (_enabled.value) startIfPossible()
|
var attempt = 0
|
||||||
|
while (_enabled.value && !_running.value && attempt < MAX_RESTART_ATTEMPTS) {
|
||||||
|
startIfPossible()
|
||||||
|
if (_running.value) break
|
||||||
|
attempt++
|
||||||
|
delay(RESTART_RETRY_DELAY_MS)
|
||||||
|
}
|
||||||
|
} finally {
|
||||||
|
restartInFlight.set(false)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
/**
|
||||||
|
* Listens for system Wi-Fi Aware availability changes. Aware can flip off/on at runtime
|
||||||
|
* (Wi-Fi toggling, hotspot/SoftAP, location changes); without this we would only recover on
|
||||||
|
* an unrelated trigger.
|
||||||
|
*/
|
||||||
|
private fun registerAwareStateReceiver(ctx: Context) {
|
||||||
|
if (awareReceiverRegistered) return
|
||||||
|
if (Build.VERSION.SDK_INT < Build.VERSION_CODES.Q) return
|
||||||
|
try {
|
||||||
|
val filter = android.content.IntentFilter(
|
||||||
|
android.net.wifi.aware.WifiAwareManager.ACTION_WIFI_AWARE_STATE_CHANGED
|
||||||
|
)
|
||||||
|
ctx.registerReceiver(object : android.content.BroadcastReceiver() {
|
||||||
|
override fun onReceive(c: Context?, intent: android.content.Intent?) {
|
||||||
|
val mgr = ctx.getSystemService(android.net.wifi.aware.WifiAwareManager::class.java)
|
||||||
|
val available = mgr?.isAvailable == true
|
||||||
|
Log.i(TAG, "Wi-Fi Aware availability changed: available=$available enabled=${_enabled.value} running=${_running.value}")
|
||||||
|
if (available) {
|
||||||
|
if (_enabled.value) restartIfStillEnabled(500)
|
||||||
|
} else if (_running.value) {
|
||||||
|
// Aware went away; tear down cleanly so we can re-attach when it returns.
|
||||||
|
// Note: this does not change the enabled preference.
|
||||||
|
stop()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}, filter)
|
||||||
|
awareReceiverRegistered = true
|
||||||
|
} catch (e: Exception) {
|
||||||
|
Log.w(TAG, "Failed to register Wi-Fi Aware state receiver: ${e.message}")
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -42,6 +42,7 @@ import java.nio.ByteBuffer
|
|||||||
import java.nio.ByteOrder
|
import java.nio.ByteOrder
|
||||||
import java.util.concurrent.ConcurrentHashMap
|
import java.util.concurrent.ConcurrentHashMap
|
||||||
import java.util.concurrent.Executors
|
import java.util.concurrent.Executors
|
||||||
|
import java.util.concurrent.atomic.AtomicBoolean
|
||||||
import java.util.concurrent.atomic.AtomicInteger
|
import java.util.concurrent.atomic.AtomicInteger
|
||||||
|
|
||||||
/**
|
/**
|
||||||
@@ -62,6 +63,12 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
private const val MAX_TTL: UByte = 7u
|
private const val MAX_TTL: UByte = 7u
|
||||||
private const val SERVICE_NAME = "bitchat"
|
private const val SERVICE_NAME = "bitchat"
|
||||||
private const val PSK = "bitchat_secret"
|
private const val PSK = "bitchat_secret"
|
||||||
|
// Network request / socket timeouts
|
||||||
|
private const val NETWORK_REQUEST_TIMEOUT_MS = 30_000
|
||||||
|
private const val ACCEPT_TIMEOUT_MS = 30_000
|
||||||
|
private const val CLIENT_CONNECT_TIMEOUT_MS = 7_000
|
||||||
|
// Discovery freshness window for reconnection maintenance
|
||||||
|
private const val DISCOVERY_STALE_MS = 5L * 60 * 1000
|
||||||
}
|
}
|
||||||
|
|
||||||
// Core crypto/services
|
// Core crypto/services
|
||||||
@@ -83,11 +90,11 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
|
|
||||||
// Wi-Fi Aware transport
|
// Wi-Fi Aware transport
|
||||||
private val awareManager = context.getSystemService(WifiAwareManager::class.java)
|
private val awareManager = context.getSystemService(WifiAwareManager::class.java)
|
||||||
private var wifiAwareSession: WifiAwareSession? = null
|
@Volatile private var wifiAwareSession: WifiAwareSession? = null
|
||||||
private var publishSession: PublishDiscoverySession? = null
|
@Volatile private var publishSession: PublishDiscoverySession? = null
|
||||||
private var subscribeSession: SubscribeDiscoverySession? = null
|
@Volatile private var subscribeSession: SubscribeDiscoverySession? = null
|
||||||
private val listenerExec = Executors.newCachedThreadPool()
|
private val listenerExec = Executors.newCachedThreadPool()
|
||||||
private var isActive = false
|
@Volatile private var isActive = false
|
||||||
@Volatile private var recoveryInProgress = false
|
@Volatile private var recoveryInProgress = false
|
||||||
private val sessionGeneration = AtomicInteger(0)
|
private val sessionGeneration = AtomicInteger(0)
|
||||||
|
|
||||||
@@ -483,6 +490,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
// Unregister from bridge
|
// Unregister from bridge
|
||||||
TransportBridgeService.unregister("WIFI")
|
TransportBridgeService.unregister("WIFI")
|
||||||
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { }
|
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { }
|
||||||
|
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("WIFI") } catch (_: Exception) { }
|
||||||
|
|
||||||
if (wasActive) {
|
if (wasActive) {
|
||||||
meshCore.sendLeaveAnnouncement()
|
meshCore.sendLeaveAnnouncement()
|
||||||
@@ -499,9 +507,13 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
wifiAwareSession?.close(); wifiAwareSession = null
|
wifiAwareSession?.close(); wifiAwareSession = null
|
||||||
|
|
||||||
handleToPeerId.clear()
|
handleToPeerId.clear()
|
||||||
|
discoveredTimestamps.clear()
|
||||||
|
|
||||||
meshCore.shutdown()
|
meshCore.shutdown()
|
||||||
|
|
||||||
|
// Tear down listener threads; this instance is discarded after a full stop.
|
||||||
|
try { listenerExec.shutdownNow() } catch (_: Exception) { }
|
||||||
|
|
||||||
com.bitchat.android.wifiaware.WifiAwareController.onServiceStopped(this@WifiAwareMeshService)
|
com.bitchat.android.wifiaware.WifiAwareController.onServiceStopped(this@WifiAwareMeshService)
|
||||||
serviceScope.cancel()
|
serviceScope.cancel()
|
||||||
}
|
}
|
||||||
@@ -520,6 +532,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
isActive = false
|
isActive = false
|
||||||
TransportBridgeService.unregister("WIFI")
|
TransportBridgeService.unregister("WIFI")
|
||||||
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { }
|
try { com.bitchat.android.services.AppStateStore.clearTransportPeers("WIFI") } catch (_: Exception) { }
|
||||||
|
try { com.bitchat.android.services.AppStateStore.clearTransportDirectPeers("WIFI") } catch (_: Exception) { }
|
||||||
val oldPublishSession = publishSession
|
val oldPublishSession = publishSession
|
||||||
val oldSubscribeSession = subscribeSession
|
val oldSubscribeSession = subscribeSession
|
||||||
val oldWifiAwareSession = wifiAwareSession
|
val oldWifiAwareSession = wifiAwareSession
|
||||||
@@ -535,9 +548,15 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
if (subscribeSession === oldSubscribeSession) subscribeSession = null
|
if (subscribeSession === oldSubscribeSession) subscribeSession = null
|
||||||
if (wifiAwareSession === oldWifiAwareSession) wifiAwareSession = null
|
if (wifiAwareSession === oldWifiAwareSession) wifiAwareSession = null
|
||||||
handleToPeerId.clear()
|
handleToPeerId.clear()
|
||||||
|
discoveredTimestamps.clear()
|
||||||
}
|
}
|
||||||
} finally {
|
} finally {
|
||||||
recoveryInProgress = false
|
recoveryInProgress = false
|
||||||
|
// Recovery cleanup is done; nudge a restart now that startServices() will no
|
||||||
|
// longer be deferred by recoveryInProgress. The controller coalesces requests.
|
||||||
|
if (com.bitchat.android.wifiaware.WifiAwareController.enabled.value) {
|
||||||
|
com.bitchat.android.wifiaware.WifiAwareController.restartIfStillEnabled(500)
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -554,9 +573,21 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
if (!isActive) break
|
if (!isActive) break
|
||||||
|
|
||||||
val now = System.currentTimeMillis()
|
val now = System.currentTimeMillis()
|
||||||
|
|
||||||
|
// 0. Prune stale discovery entries. PeerHandles become invalid when the
|
||||||
|
// discovery sessions restart, so we must not keep pinging old handles forever.
|
||||||
|
val staleIds = discoveredTimestamps.filter { (id, ts) ->
|
||||||
|
(now - ts) >= DISCOVERY_STALE_MS && !connectionTracker.isConnected(id)
|
||||||
|
}.keys.toSet()
|
||||||
|
if (staleIds.isNotEmpty()) {
|
||||||
|
staleIds.forEach { discoveredTimestamps.remove(it) }
|
||||||
|
handleToPeerId.entries.removeIf { it.value in staleIds }
|
||||||
|
Log.d(TAG, "Maintenance: pruned ${staleIds.size} stale discovery entries")
|
||||||
|
}
|
||||||
|
|
||||||
// 1. Identify peers that are discovered (recently seen) but not currently connected
|
// 1. Identify peers that are discovered (recently seen) but not currently connected
|
||||||
val recentDiscovered = discoveredTimestamps.filter { (id, ts) ->
|
val recentDiscovered = discoveredTimestamps.filter { (id, ts) ->
|
||||||
(now - ts) < 5 * 60 * 1000 // Seen in last 5 minutes
|
(now - ts) < DISCOVERY_STALE_MS // Seen in last 5 minutes
|
||||||
}.keys
|
}.keys
|
||||||
|
|
||||||
// 2. Filter out those who are already connected
|
// 2. Filter out those who are already connected
|
||||||
@@ -639,11 +670,18 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
.build()
|
.build()
|
||||||
|
|
||||||
val cb = object : ConnectivityManager.NetworkCallback() {
|
val cb = object : ConnectivityManager.NetworkCallback() {
|
||||||
private var activeSocket: SyncedSocket? = null
|
@Volatile private var activeSocket: SyncedSocket? = null
|
||||||
|
private val acceptStarted = AtomicBoolean(false)
|
||||||
|
|
||||||
override fun onAvailable(network: Network) {
|
override fun onAvailable(network: Network) {
|
||||||
Log.i(TAG, "SERVER: onAvailable() - Aware network is ready for ${peerId.take(8)}")
|
Log.i(TAG, "SERVER: onAvailable() - Aware network is ready for ${peerId.take(8)}")
|
||||||
|
// Only accept once per network request
|
||||||
|
if (!acceptStarted.compareAndSet(false, true)) return
|
||||||
|
// Offload the blocking accept() off the callback thread so we never stall
|
||||||
|
// the (main-thread) ConnectivityManager callback dispatcher.
|
||||||
|
listenerExec.execute {
|
||||||
try {
|
try {
|
||||||
|
try { ss.soTimeout = ACCEPT_TIMEOUT_MS } catch (_: Exception) {}
|
||||||
val client = ss.accept()
|
val client = ss.accept()
|
||||||
Log.i(TAG, "SERVER: Accepted raw TCP connection from ${peerId.take(8)}")
|
Log.i(TAG, "SERVER: Accepted raw TCP connection from ${peerId.take(8)}")
|
||||||
try { network.bindSocket(client) } catch (e: Exception) { Log.w(TAG, "Server bindSocket EPERM: ${e.message}") }
|
try { network.bindSocket(client) } catch (e: Exception) { Log.w(TAG, "Server bindSocket EPERM: ${e.message}") }
|
||||||
@@ -669,6 +707,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
handleNetworkFailure(peerId)
|
handleNetworkFailure(peerId)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
|
|
||||||
override fun onUnavailable() {
|
override fun onUnavailable() {
|
||||||
Log.e(TAG, "SERVER: onUnavailable() - Failed to acquire Aware network for ${peerId.take(8)} (timeout or refused)")
|
Log.e(TAG, "SERVER: onUnavailable() - Failed to acquire Aware network for ${peerId.take(8)} (timeout or refused)")
|
||||||
@@ -685,7 +724,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
Log.i(TAG, "SERVER: [Calling requestNetwork] for ${peerId.take(8)} with port $port")
|
Log.i(TAG, "SERVER: [Calling requestNetwork] for ${peerId.take(8)} with port $port")
|
||||||
try {
|
try {
|
||||||
// use requestNetwork with a timeout to trigger onUnavailable if it fails
|
// use requestNetwork with a timeout to trigger onUnavailable if it fails
|
||||||
cm.requestNetwork(req, cb, 30_000)
|
cm.requestNetwork(req, cb, NETWORK_REQUEST_TIMEOUT_MS)
|
||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
Log.e(TAG, "SERVER: ConnectivityManager.requestNetwork threw exception", e)
|
Log.e(TAG, "SERVER: ConnectivityManager.requestNetwork threw exception", e)
|
||||||
connectionTracker.disconnect(peerId)
|
connectionTracker.disconnect(peerId)
|
||||||
@@ -773,7 +812,8 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
.build()
|
.build()
|
||||||
|
|
||||||
val cb = object : ConnectivityManager.NetworkCallback() {
|
val cb = object : ConnectivityManager.NetworkCallback() {
|
||||||
private var activeSocket: SyncedSocket? = null
|
@Volatile private var activeSocket: SyncedSocket? = null
|
||||||
|
private val connectStarted = AtomicBoolean(false)
|
||||||
|
|
||||||
override fun onAvailable(network: Network) {
|
override fun onAvailable(network: Network) {
|
||||||
Log.i(TAG, "CLIENT: onAvailable() - Aware network is ready for ${peerId.take(8)}")
|
Log.i(TAG, "CLIENT: onAvailable() - Aware network is ready for ${peerId.take(8)}")
|
||||||
@@ -789,11 +829,15 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
if (connectionTracker.peerSockets.containsKey(peerId)) return
|
if (connectionTracker.peerSockets.containsKey(peerId)) return
|
||||||
val info = (nc.transportInfo as? WifiAwareNetworkInfo) ?: return
|
val info = (nc.transportInfo as? WifiAwareNetworkInfo) ?: return
|
||||||
val addr = info.peerIpv6Addr as? Inet6Address ?: return
|
val addr = info.peerIpv6Addr as? Inet6Address ?: return
|
||||||
|
// onCapabilitiesChanged can fire multiple times; only connect once
|
||||||
|
if (!connectStarted.compareAndSet(false, true)) return
|
||||||
Log.i(TAG, "CLIENT: onCapabilitiesChanged() - Peer IPv6 discovered: $addr")
|
Log.i(TAG, "CLIENT: onCapabilitiesChanged() - Peer IPv6 discovered: $addr")
|
||||||
|
|
||||||
val lp = cm.getLinkProperties(network)
|
val lp = cm.getLinkProperties(network)
|
||||||
val iface = lp?.interfaceName
|
val iface = lp?.interfaceName
|
||||||
|
|
||||||
|
// Offload the blocking connect() off the callback thread.
|
||||||
|
listenerExec.execute {
|
||||||
try {
|
try {
|
||||||
val sock = Socket()
|
val sock = Socket()
|
||||||
try { network.bindSocket(sock) } catch (e: Exception) { Log.w(TAG, "Client bindSocket EPERM: ${e.message}") }
|
try { network.bindSocket(sock) } catch (e: Exception) { Log.w(TAG, "Client bindSocket EPERM: ${e.message}") }
|
||||||
@@ -811,7 +855,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
addr
|
addr
|
||||||
}
|
}
|
||||||
|
|
||||||
sock.connect(java.net.InetSocketAddress(scopedAddr, port), 7000)
|
sock.connect(java.net.InetSocketAddress(scopedAddr, port), CLIENT_CONNECT_TIMEOUT_MS)
|
||||||
Log.i(TAG, "CLIENT: TCP connected to ${peerId.take(8)} at $scopedAddr:$port")
|
Log.i(TAG, "CLIENT: TCP connected to ${peerId.take(8)} at $scopedAddr:$port")
|
||||||
|
|
||||||
val synced = SyncedSocket(sock)
|
val synced = SyncedSocket(sock)
|
||||||
@@ -834,6 +878,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
handleNetworkFailure(peerId)
|
handleNetworkFailure(peerId)
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
}
|
||||||
override fun onLost(network: Network) {
|
override fun onLost(network: Network) {
|
||||||
handlePeerDisconnection(peerId, activeSocket)
|
handlePeerDisconnection(peerId, activeSocket)
|
||||||
Log.i(TAG, "CLIENT: WiFi Aware network lost for ${peerId.take(8)}")
|
Log.i(TAG, "CLIENT: WiFi Aware network lost for ${peerId.take(8)}")
|
||||||
@@ -843,7 +888,7 @@ class WifiAwareMeshService(private val context: Context) : MeshService, Transpor
|
|||||||
connectionTracker.addNetworkCallback(peerId, cb)
|
connectionTracker.addNetworkCallback(peerId, cb)
|
||||||
Log.i(TAG, "CLIENT: [Calling requestNetwork] for ${peerId.take(8)}")
|
Log.i(TAG, "CLIENT: [Calling requestNetwork] for ${peerId.take(8)}")
|
||||||
try {
|
try {
|
||||||
cm.requestNetwork(req, cb, 30_000)
|
cm.requestNetwork(req, cb, NETWORK_REQUEST_TIMEOUT_MS)
|
||||||
} catch (e: Exception) {
|
} catch (e: Exception) {
|
||||||
Log.e(TAG, "CLIENT: ConnectivityManager.requestNetwork threw exception", e)
|
Log.e(TAG, "CLIENT: ConnectivityManager.requestNetwork threw exception", e)
|
||||||
connectionTracker.disconnect(peerId)
|
connectionTracker.disconnect(peerId)
|
||||||
|
|||||||
@@ -67,4 +67,33 @@ class AppStateStoreTest {
|
|||||||
|
|
||||||
assertEquals(listOf("wifi-peer", "ble-peer"), AppStateStore.peers.value)
|
assertEquals(listOf("wifi-peer", "ble-peer"), AppStateStore.peers.value)
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `direct peers union across transports`() {
|
||||||
|
AppStateStore.setTransportDirectPeers("BLE", listOf("ble-1", "shared"))
|
||||||
|
AppStateStore.setTransportDirectPeers("WIFI", listOf("wifi-1", "shared"))
|
||||||
|
|
||||||
|
assertEquals(
|
||||||
|
setOf("ble-1", "wifi-1", "shared"),
|
||||||
|
AppStateStore.getDirectPeers()
|
||||||
|
)
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `clearing one transport keeps the other transport direct peers`() {
|
||||||
|
AppStateStore.setTransportDirectPeers("BLE", listOf("ble-1"))
|
||||||
|
AppStateStore.setTransportDirectPeers("WIFI", listOf("wifi-1"))
|
||||||
|
|
||||||
|
AppStateStore.clearTransportDirectPeers("WIFI")
|
||||||
|
|
||||||
|
assertEquals(setOf("ble-1"), AppStateStore.getDirectPeers())
|
||||||
|
}
|
||||||
|
|
||||||
|
@Test
|
||||||
|
fun `latest direct peer set replaces previous set for same transport`() {
|
||||||
|
AppStateStore.setTransportDirectPeers("WIFI", listOf("wifi-1", "wifi-2"))
|
||||||
|
AppStateStore.setTransportDirectPeers("WIFI", listOf("wifi-3"))
|
||||||
|
|
||||||
|
assertEquals(setOf("wifi-3"), AppStateStore.getDirectPeers())
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
Reference in New Issue
Block a user