diff --git a/CHANGELOG.md b/CHANGELOG.md index ae26e6fd..bae5a975 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,8 +13,39 @@ archived by series under [docs/changelog/](docs/changelog/); see the ## [Unreleased] +> **Android and iOS meet on a shared Wi-Fi network.** The Android +> peer-stream slot now finds iPhones and Python hosts on the Wi-Fi network +> it is on, and they find it, with no Wi-Fi Direct group. +> [§27](docs/UPGRADING.md#27-behaviour-that-changes-without-a-compile-error-unreleased) +> lists what changes without a compile error. + ### Added +- **Android peer streams on the Wi-Fi network.** With + `wifiDirect: { enabled: true }`, the Android manager advertises and browses + DNS-SD `_offlineprotocol._tcp` through `NsdManager` on the Wi-Fi network it + is on, with the record iOS and Python publish (`txtvers=1`, `addr`, the same + instance name), and dials what it finds on iOS's policy: the lower address + at once, the higher after five seconds, a redial ladder up to a minute, and + a record whose dial is answered by another address left out. Dials are bound + to the Wi-Fi network (a socket left to the default network goes out over + cellular on a Wi-Fi with no internet) and carry the record's `addr` as the + address the preamble must prove. A Wi-Fi network with no internet counts. + The slot stays up while either the group or the network is, and Wi-Fi P2P + going off, or leaving the group, ends only the group's streams. It needs + `ACCESS_NETWORK_STATE` and `CHANGE_WIFI_MULTICAST_STATE` (a multicast lock + while advertising or browsing on Android 12 and lower, where mDNS needs one), which the + module declares, and no runtime grant until the app targets Android 17 + (API 37), so it also runs on a phone that has not granted + `NEARBY_WIFI_DEVICES`. From API 37 the app must declare + `ACCESS_LOCAL_NETWORK` and request it before `start()`, or the LAN carrier + stays off with a warning diagnostic; the SDK does not declare it, because a + declaration revokes the grant apps targeting 36 and lower hold by default. A network + that blocks multicast or isolates clients (many guest and office networks) + finds nothing, and Android to iOS with no shared network still goes over + Bluetooth LE or the relays + ([ADR 0028](docs/adr/0028-android-peer-streams-join-the-lan.md)). + - **The service starts every transport its configuration enables.** `offline-protocol-service` started only the peer stream and the gateway client, so a configuration with `ble_enabled` built both Bluetooth LE @@ -135,6 +166,31 @@ archived by series under [docs/changelog/](docs/changelog/); see the ### Changed +- **Android keeps the stream the lower address opened.** Of two streams for + one address, Android kept the newer; iOS and Python keep the one the lower + address opened, and the newer of two such. On a shared network both ends + dial, and two different rules keep opposite streams and reconnect forever. + Android now computes the same rule, comparing addresses by their UTF-8 + bytes, and `every_peer_stream_manager_keeps_the_same_stream` (renamed from + `ios_and_python_peer_streams_keep_the_same_stream`) pins all three copies. + Inside a Wi-Fi Direct group a client that is the higher address can no + longer supersede its own half-open stream: its reconnect waits for keepalive, + about thirty seconds on Android 10 and later. +- **Android peer streams use the iOS and Python keepalive** on Android 10 and + later: 15 seconds idle, 5 between probes, 3 probes, and 30 seconds for + unacknowledged data, where the OS defaults were two hours and about fifteen + minutes. +- **A peer-stream redial ladder starts over only after a stream that carried + its peer**, on iOS and Android: a body after the preamble, or the address + held for thirty seconds, not the proof alone. A stream the peer refused for its own stale one proved the address + too, so the dialer was announced and lost every second until that stream + died, and a Wi-Fi Direct client whose owner was held by a LAN stream + redialed it every second. Such a client now waits for that stream to end. +- **Android's peer-stream listener bounds inbound streams.** It binds every + interface, so any device on a shared network could take all sixteen slots. + At most twelve are inbound now, and four from one remote address, iOS's + bounds. + - **A peer-stream or relay flag the configuration cannot honour is refused.** `--listen` and `--peer` without `wifi_direct_enabled` were ignored, which left a service that reached nobody and said nothing; they, `--lan`, and diff --git a/bindings/react-native/README.md b/bindings/react-native/README.md index 05d3d0c9..f1499eec 100644 --- a/bindings/react-native/README.md +++ b/bindings/react-native/README.md @@ -106,8 +106,11 @@ Add to `AndroidManifest.xml`: ``` Wi-Fi Direct and the internet need nothing in your manifest: the SDK's own -manifest declares `ACCESS_WIFI_STATE`, `CHANGE_WIFI_STATE`, `INTERNET` and -`NEARBY_WIFI_DEVICES`, and the build merges them into your app. It declares +manifest declares `ACCESS_WIFI_STATE`, `CHANGE_WIFI_STATE`, `INTERNET`, +`ACCESS_NETWORK_STATE`, `CHANGE_WIFI_MULTICAST_STATE` and `NEARBY_WIFI_DEVICES`, and the build merges them +into your app. The peer-stream slot also finds iPhones and hosts on the Wi-Fi +network the phone is on, which needs no runtime grant unless your app +targets Android 17 (API 37): then an app targeting Android 17 (API 37) must declare `ACCESS_LOCAL_NETWORK` and request it before `start()`, or the LAN carrier stays off with a warning diagnostic. The SDK does not declare it, because a declaration revokes the grant apps targeting 36 and lower hold by default. The SDK's manifest declares `NEARBY_WIFI_DEVICES` with `android:usesPermissionFlags="neverForLocation"`, and that flag reaches your merged manifest too: the transport derives no location from Wi-Fi, so on Android 13+ it needs `NEARBY_WIFI_DEVICES` alone. @@ -519,7 +522,7 @@ interface TransportsConfig { reconnectDelay?: number; // ms }; wifiDirect?: { - enabled: boolean; // default: false (Android only) + enabled: boolean; // default: false. Android: Wi-Fi Direct and the Wi-Fi network; iOS: the LAN and AWDL deviceName?: string; autoAccept?: boolean; // Android 10+: form the group without the system settings groupOwnerIntent?: number; // deprecated, not used diff --git a/bindings/react-native/android/src/main/AndroidManifest.xml b/bindings/react-native/android/src/main/AndroidManifest.xml index 675a4d1e..1c2fc4e8 100644 --- a/bindings/react-native/android/src/main/AndroidManifest.xml +++ b/bindings/react-native/android/src/main/AndroidManifest.xml @@ -31,6 +31,12 @@ + + + diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/LanPeerDiscovery.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/LanPeerDiscovery.kt new file mode 100644 index 00000000..70269301 --- /dev/null +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/LanPeerDiscovery.kt @@ -0,0 +1,609 @@ +package com.offlineprotocol + +import android.content.Context +import android.content.pm.PackageManager +import android.net.ConnectivityManager +import android.net.LinkProperties +import android.net.Network +import android.net.NetworkCapabilities +import android.net.NetworkRequest +import android.net.nsd.NsdManager +import android.net.nsd.NsdServiceInfo +import android.net.wifi.WifiManager +import android.os.Build +import android.os.Handler +import android.os.ext.SdkExtensions +import androidx.core.content.ContextCompat +import java.net.InetAddress +import java.net.InetSocketAddress +import java.net.Socket +import java.util.concurrent.ExecutorService +import java.util.concurrent.RejectedExecutionException + +/** + * The peer-stream slot's LAN carrier: DNS-SD `_offlineprotocol._tcp` on the + * Wi-Fi network this device is on, so it finds iPhones and Python hosts there + * and they find it (ADR 0028, stream-framing.md "Finding a peer on a LAN"). + * + * It publishes the record iOS and Python publish (`txtvers=1`, `addr`, the + * same instance name), browses for theirs, and dials what it finds on the + * policy iOS dials on ([PeerStreamDialPolicy]): the lower address at once, + * the higher after a grace delay, so the stream both ends keep usually + * arrives first. Inbound LAN streams arrive at [WifiDirectManager]'s + * listener, which already binds every interface. What a stream does once + * open is [PeerStreamSockets]'s, the same as for a group stream. + * + * A record's `addr` is a hint: a dial passes it as the preamble's expected + * address, and a record whose dial is answered by another address is left + * out until it is published again. + * + * Every method and every framework callback runs on [handler]'s thread + * (the transport thread); connects and stream reads run on [executor], + * because nothing that blocks on the network may run on the transport + * thread (TransportConfinement). + */ +internal class LanPeerDiscovery( + private val context: Context, + private val handler: Handler, + private val executor: ExecutorService, + private val sockets: PeerStreamSockets, + private val port: Int, + private val localAddress: () -> String?, + /** The transport is running: what a stream's `accepting` asks. */ + private val running: () -> Boolean, + /** The Wi-Fi network's own addresses when it comes up, null when it goes. */ + private val networkChanged: (Set?) -> Unit, + private val diagnostic: (level: String, message: String, context: Map) -> Unit, +) { + private data class Record(val address: String, val host: InetAddress, val port: Int) + + private val connectivity = context.getSystemService(ConnectivityManager::class.java) + private val nsd = context.getSystemService(NsdManager::class.java) + + private var started = false + private var paused = false + /** The matching Wi-Fi networks and the one followed ([NetworkChoice]). */ + private val networks = NetworkChoice() + private val network: Network? get() = networks.current + /** Bumped whenever what was found stops counting, so stale callbacks and timers do nothing. */ + private var generation = 0 + + private var registration: NsdManager.RegistrationListener? = null + /** Our record's name as registered (NSD renames on a conflict). */ + private var registeredName: String? = null + private var discovery: NsdManager.DiscoveryListener? = null + + /** Resolved peer records by instance name. */ + private val records = HashMap() + /** The instance name each advertised address is dialed at. */ + private var adverts: Map = emptyMap() + private val unprovable = HashSet() + private val policy = PeerStreamDialPolicy() + + /** What NSD reports, and its resolves, one at a time ([ResolveQueue]). */ + private val resolves = ResolveQueue({ it.serviceName }) + + // Held while the advert or the browse is active. Before T extensions 7 + // (Android 12 and lower, and 13 without that update) the Wi-Fi driver + // drops the multicast mDNS runs on unless an app holds one, so this + // device would neither hear a peer's answers nor the queries a peer + // sends to find it (NsdManager's own documentation). The advert needs it + // as much as the browse: a lock that followed the browse alone was let + // go on pause and during a browse retry, while the record stayed + // published and could not be found. From extension 7 the system manages + // it for a foreground app, and a lock only costs battery, so none is + // taken. + private val multicastLock: WifiManager.MulticastLock? = + if (needsMulticastLock()) { + context.applicationContext.getSystemService(WifiManager::class.java) + ?.createMulticastLock("offlineprotocol-lan") + ?.apply { setReferenceCounted(false) } + } else { + null + } + + private val networkCallback = object : ConnectivityManager.NetworkCallback() { + override fun onAvailable(network: Network) { + val properties = try { connectivity.getLinkProperties(network) } catch (_: Exception) { null } + if (properties != null) handler.post { networkUp(network, properties) } + } + + override fun onLinkPropertiesChanged(network: Network, properties: LinkProperties) { + handler.post { networkUp(network, properties) } + } + + override fun onLost(network: Network) { + handler.post { networkLost(network) } + } + } + + fun start() { + if (started) return + if (!localNetworkPermitted()) { + // Off rather than half on: without the grant NsdManager shows a + // system service picker on every browse, and sockets to the LAN + // are refused. A later grant takes effect on the next start. + diagnostic("warning", "LAN carrier off: ACCESS_LOCAL_NETWORK not granted", emptyMap()) + return + } + started = true + // Without INTERNET removed, the default request skips a Wi-Fi network + // with no route out, which is exactly the offline LAN this is for. + val request = NetworkRequest.Builder() + .addTransportType(NetworkCapabilities.TRANSPORT_WIFI) + .removeCapability(NetworkCapabilities.NET_CAPABILITY_INTERNET) + .build() + try { + connectivity.registerNetworkCallback(request, networkCallback) + } catch (e: Exception) { + diagnostic("error", "LAN discovery could not watch Wi-Fi", mapOf( + "error" to (e.message ?: e.javaClass.simpleName), + )) + } + } + + /** Ends discovery and advertising. The owner ends the LAN streams. */ + fun stop() { + if (!started) return + started = false + try { connectivity.unregisterNetworkCallback(networkCallback) } catch (_: Exception) {} + networks.clear() + stopNsd() + } + + /** Stops browsing and dialing; the record stays published. */ + fun pause() { + paused = true + stopDiscovery() + // No loss reports arrive while not browsing, so what was found is + // forgotten and found again on resume. + forgetRecords() + } + + fun resume() { + paused = false + if (network != null) discover() + } + + // MARK: - The network + + /** A matching network came up or changed; [NetworkChoice] decides what it means. */ + private fun networkUp(network: Network, properties: LinkProperties) { + if (!started) return + when (networks.up(network, properties)) { + NetworkChoice.Up.ADOPT -> adopt(properties) + NetworkChoice.Up.UPDATE -> networkChanged(properties.linkAddresses.map { it.address }.toSet()) + NetworkChoice.Up.WAIT -> {} + } + } + + private fun networkLost(network: Network) { + if (!networks.lost(network)) return + stopNsd() + networkChanged(null) + networks.fallback()?.let { (_, properties) -> adopt(properties) } + } + + /** The network now followed ([networks]`.current`) starts the carrier. */ + private fun adopt(properties: LinkProperties) { + networkChanged(properties.linkAddresses.map { it.address }.toSet()) + startNsd() + } + + // MARK: - NSD + + private fun startNsd() { + generation++ + val address = localAddress() + if (address == null) { + diagnostic("warning", "LAN advert skipped: no identity yet", emptyMap()) + } else { + register(address) + } + if (!paused) discover() + } + + private fun stopNsd() { + generation++ + registration?.let { try { nsd.unregisterService(it) } catch (_: Exception) {} } + registration = null + registeredName = null + stopDiscovery() + forgetRecords() + policy.reset() + } + + private fun forgetRecords() { + resolves.clear() + records.clear() + adverts = emptyMap() + unprovable.clear() + } + + private fun register(address: String) { + // Order is the framework's: NsdServiceInfo keeps attributes in a map, + // so `txtvers` may not come first. Every reader looks entries up by key. + val info = NsdServiceInfo().apply { + serviceName = instanceName(address) + serviceType = WifiDirectGroupFormation.SERVICE_TYPE + port = this@LanPeerDiscovery.port + setAttribute(WifiDirectGroupFormation.KEY_VERSION, "1") + setAttribute(WifiDirectGroupFormation.KEY_ADDRESS, address) + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU) { + setNetwork(this@LanPeerDiscovery.network) + } + } + val listener = object : NsdManager.RegistrationListener { + override fun onServiceRegistered(info: NsdServiceInfo) { + handler.post { if (registration === this) registeredName = info.serviceName } + } + + override fun onRegistrationFailed(info: NsdServiceInfo, code: Int) { + handler.post { if (registration === this) registrationFailed(mapOf("code" to code)) } + } + + override fun onServiceUnregistered(info: NsdServiceInfo) {} + override fun onUnregistrationFailed(info: NsdServiceInfo, code: Int) {} + } + registration = listener + acquireMulticastLock() + try { + nsd.registerService(info, NsdManager.PROTOCOL_DNS_SD, listener) + } catch (e: Exception) { + registrationFailed(mapOf("error" to (e.message ?: e.javaClass.simpleName))) + } + } + + /** + * The advert failed: tried again later, as iOS rebuilds its listener. A + * failure was final until the network changed, which left this device + * invisible to every peer on it for the session. + */ + private fun registrationFailed(context: Map) { + registration = null + releaseMulticastLockIfIdle() + diagnostic("error", "LAN advert failed", context) + retryLater { if (network != null && registration == null) localAddress()?.let { register(it) } } + } + + private fun discover() { + if (discovery != null) return + acquireMulticastLock() + val listener = object : NsdManager.DiscoveryListener { + override fun onDiscoveryStarted(serviceType: String) {} + override fun onDiscoveryStopped(serviceType: String) {} + + override fun onStartDiscoveryFailed(serviceType: String, code: Int) { + handler.post { if (discovery === this) discoveryFailed(mapOf("code" to code)) } + } + + override fun onStopDiscoveryFailed(serviceType: String, code: Int) {} + + override fun onServiceFound(info: NsdServiceInfo) { + handler.post { if (discovery === this) found(info) } + } + + override fun onServiceLost(info: NsdServiceInfo) { + handler.post { if (discovery === this) lost(info) } + } + } + discovery = listener + try { + val network = network + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU && network != null) { + nsd.discoverServices( + WifiDirectGroupFormation.SERVICE_TYPE, NsdManager.PROTOCOL_DNS_SD, + network, { it.run() }, listener, + ) + } else { + nsd.discoverServices(WifiDirectGroupFormation.SERVICE_TYPE, NsdManager.PROTOCOL_DNS_SD, listener) + } + } catch (e: Exception) { + discoveryFailed(mapOf("error" to (e.message ?: e.javaClass.simpleName))) + } + } + + /** + * The browse failed to start: tried again later, as iOS rebuilds its + * browser. A failure was final until the network changed. The lock is + * let go only if no advert needs it either. + */ + private fun discoveryFailed(context: Map) { + discovery = null + releaseMulticastLockIfIdle() + diagnostic("error", "LAN discovery failed to start", context) + retryLater { if (!paused && network != null) discover() } + } + + /** Runs [retry] after [REBUILD_DELAY_MS] unless the NSD state it was for has gone. */ + private fun retryLater(retry: () -> Unit) { + val expected = generation + handler.postDelayed({ if (started && expected == generation) retry() }, REBUILD_DELAY_MS) + } + + private fun stopDiscovery() { + discovery?.let { try { nsd.stopServiceDiscovery(it) } catch (_: Exception) {} } + discovery = null + releaseMulticastLockIfIdle() + } + + private fun acquireMulticastLock() { + try { multicastLock?.acquire() } catch (_: Exception) {} + } + + /** Lets the lock go once neither the advert nor the browse needs it. */ + private fun releaseMulticastLockIfIdle() { + if (registration != null || discovery != null) return + try { multicastLock?.let { if (it.isHeld) it.release() } } catch (_: Exception) {} + } + + /** + * Whether [info] was found on the Wi-Fi network dials are bound to. The + * same instance can be reported on more than one interface (a Wi-Fi + * Direct group, a tethering downstream with no Network at all), and + * records are keyed by name: a copy from elsewhere overwrote the Wi-Fi + * record with a host no bound dial reaches, and its loss deleted it. + * From API 33 the advert and the browse are scoped to the network too. + * Below it NSD carries no network, every interface is browsed, and a + * record from another one fails its dial and climbs the ladder. + */ + private fun onOurNetwork(info: NsdServiceInfo): Boolean = + Build.VERSION.SDK_INT < Build.VERSION_CODES.TIRAMISU || info.network == network + + private fun found(info: NsdServiceInfo) { + if (!onOurNetwork(info)) return + if (info.serviceName == registeredName) return + resolves.found(info) + resolveNext() + } + + /** Resolves [name] again: Android's cache answers with its latest port and host. */ + private fun refresh(name: String) { + resolves.queue(name) + resolveNext() + } + + private fun lost(info: NsdServiceInfo) { + if (!onOurNetwork(info)) return + val name = info.serviceName + resolves.lost(name) + if (records.remove(name) == null) return + unprovable.remove(name) + rebuild(emptySet()) + } + + private fun resolveNext() { + val ticket = resolves.next() ?: return + val name = ticket.item.serviceName + val listener = object : NsdManager.ResolveListener { + override fun onServiceResolved(info: NsdServiceInfo) { + handler.post { + if (resolves.answered(ticket)) resolved(info) + resolveNext() + } + } + + override fun onResolveFailed(info: NsdServiceInfo, code: Int) { + handler.post { + if (code == NsdManager.FAILURE_ALREADY_ACTIVE && resolves.waitForActive(ticket)) { + // A resolve from before a restart is still running: + // this one waits for it, a few times, then gives way. + handler.postDelayed({ if (resolves.release(ticket)) resolveNext() }, RESOLVE_RETRY_MS) + return@post + } + if (resolves.gaveUp(ticket)) retryResolve(name) + resolveNext() + } + } + } + // NsdManager promises no answer, and one resolve that never comes back + // held the queue for good: no peer found later was ever resolved, + // and a refresh after a dead port waited behind it. + handler.postDelayed({ + if (!resolves.gaveUp(ticket)) return@postDelayed + if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.UPSIDE_DOWN_CAKE) { + try { nsd.stopServiceResolution(listener) } catch (_: Exception) {} + } + retryResolve(name) + resolveNext() + }, RESOLVE_TIMEOUT_MS) + try { + @Suppress("DEPRECATION") + nsd.resolveService(ticket.item, listener) + } catch (e: Exception) { + resolves.gaveUp(ticket) + diagnostic("warning", "LAN resolve failed", mapOf("error" to (e.message ?: e.javaClass.simpleName))) + retryResolve(name) + // Posted: an item that throws again must not recurse here. + handler.post { resolveNext() } + } + } + + /** + * A resolve of [name] gave up (failed, timed out, or gave way to one + * still running). It is queued again later while NSD still reports the + * service: NSD reports a service again only after losing it, so a resolve + * dropped for good left that peer unresolved for the session, unreachable + * if it never dials this device itself. + */ + private fun retryResolve(name: String) { + retryLater { refresh(name) } + } + + private fun resolved(info: NsdServiceInfo) { + if (!onOurNetwork(info)) return + val address = lanAdvertAddress(info.attributes) ?: return + val local = localAddress() ?: return + if (address == local) return + @Suppress("DEPRECATION") + val host = info.host ?: return + val name = info.serviceName + records[name] = Record(address, host, info.port) + // A record published again may now hold what it advertises. + unprovable.remove(name) + rebuild(setOf(name)) + if (adverts[address] != name) return + val weAreLower = PeerStreamLinks.newStreamWins(outbound = true, localAddress = local, peer = address) + policy.discovered(address, weAreLower, sockets.holds(address))?.let { schedule(address, it) } + } + + private fun rebuild(fresh: Set) { + unprovable.retainAll(records.keys) + adverts = PeerStreamDialPolicy.adverts( + records.map { (name, record) -> record.address to name }, + fresh, adverts, unprovable, + ) + } + + /** + * An announced stream to [address] ended, on either carrier, inbound or + * outbound. While its record is still advertised it is dialed on the + * usual policy, unless a dial is already under way (an outbound stream's + * own end redials through [ended]). + * + * Needed because this API reports no change to a record: a peer that + * restarts and publishes the same name on a new port, with no goodbye in + * between, is an update NsdManager's DiscoveryListener never delivers. + * iOS hears it as a changed record and dials; here nothing would. + */ + fun peerLost(address: String) { + if (!started || paused || !adverts.containsKey(address)) return + val local = localAddress() ?: return + val weAreLower = PeerStreamLinks.newStreamWins(outbound = true, localAddress = local, peer = address) + policy.discovered(address, weAreLower, sockets.holds(address))?.let { schedule(address, it) } + } + + // MARK: - Dialing + + private fun schedule(address: String, delayMs: Long) { + val expected = generation + handler.postDelayed({ if (expected == generation) dial(address) }, delayMs) + } + + /** + * Opens a stream toward [address], unless it went away or is already held + * (an inbound stream got here first), or later when no slot is free. + */ + private fun dial(address: String) { + val name = adverts[address] + val record = name?.let { records[it] } + val network = network + if (paused || !running() || name == null || record == null || network == null || sockets.holds(address)) { + policy.abandoned(address) + return + } + if (!sockets.hasRoom()) { + policy.noSlot(address)?.let { schedule(address, it) } + return + } + val expected = generation + try { + executor.execute { + val ran = connectAndRun(network, record, address) + handler.post { if (expected == generation) ended(address, name, ran) } + } + } catch (_: RejectedExecutionException) { + policy.abandoned(address) + } + } + + /** Runs on [executor]: the connect and the stream's whole life. */ + private fun connectAndRun(network: Network, record: Record, address: String): PeerStreamSockets.Ran { + val socket = Socket() + return try { + // Bound to the Wi-Fi network: an unbound socket follows the + // default network, which is cellular on a Wi-Fi with no internet. + network.bindSocket(socket) + socket.connect(InetSocketAddress(record.host, record.port), DIAL_TIMEOUT_MS) + sockets.run(socket, outbound = true, PeerStreamSockets.Carrier.LAN, expected = address) { running() } + } catch (e: Exception) { + try { socket.close() } catch (_: Exception) {} + diagnostic("info", "LAN dial failed", mapOf( + "address" to address, + "error" to (e.message ?: e.javaClass.simpleName), + )) + PeerStreamSockets.Ran(proved = null, heard = false) + } + } + + /** + * An outbound stream is over. It is dialed again, later each time, while + * its advert stays. A dial answered by a peer that proved nothing, or + * another address, leaves its record out: a device claiming another's + * address, or one that moved, would fail the same way on every redial + * and keep the real record for the address undialed. + */ + private fun ended(address: String, name: String, ran: PeerStreamSockets.Ran) { + if (ran.proved == address && ran.carried) { + // Carried the peer: the ladder starts over. A proof the peer then + // refused (its stale stream holds us) climbs it like a miss, or + // the dialer would be announced and lost every second. + policy.proved(address) + } else if (ran.heard && ran.proved == null && records.containsKey(name)) { + unprovable.add(name) + rebuild(emptySet()) + } else if (!ran.heard) { + // Nothing answered: the peer may be back on another port, which + // reached the cache as an update this API does not report. + refresh(name) + } + policy.ended(address, adverts.containsKey(address), sockets.holds(address)) + ?.let { schedule(address, it) } + } + + /** + * Whether this app may reach the local network. An app targeting Android + * 17 (API 37) needs `ACCESS_LOCAL_NETWORK`, a runtime grant; one targeting + * 36 or lower holds it by default, so the check passes. The module does + * not declare it: a declaration revokes that default grant. + */ + private fun localNetworkPermitted(): Boolean = + Build.VERSION.SDK_INT < API_37 || + // The restriction keys off the app's target, so this holds + // whether or not the platform's default grant to older targets + // (ConnectivityCompatChanges.USE_NSD_PICKER_WHEN_NO_LOCAL_NET_PERMISSION) + // shows up in checkSelfPermission. + context.applicationInfo.targetSdkVersion < API_37 || + ContextCompat.checkSelfPermission(context, ACCESS_LOCAL_NETWORK) == + PackageManager.PERMISSION_GRANTED + + companion object { + /** Android 17. Not a named constant in the SDK this module compiles against. */ + private const val API_37 = 37 + private const val ACCESS_LOCAL_NETWORK = "android.permission.ACCESS_LOCAL_NETWORK" + + private fun needsMulticastLock(): Boolean = + Build.VERSION.SDK_INT < Build.VERSION_CODES.TIRAMISU || + SdkExtensions.getExtensionVersion(Build.VERSION_CODES.TIRAMISU) < 7 + + /** A connect that takes longer is left to the redial ladder. iOS's `DIAL_TIMEOUT`. */ + private const val DIAL_TIMEOUT_MS = 10_000 + private const val RESOLVE_RETRY_MS = 1_000L + /** A resolve with no answer by then is given up, and the next one runs. */ + private const val RESOLVE_TIMEOUT_MS = 15_000L + /** A failed advert or browse is tried again after this long. iOS's `REBUILD_DELAY`. */ + private const val REBUILD_DELAY_MS = 5_000L + + /** + * The DNS-SD instance name: a digest of the address, as iOS and Python + * name their own, so a restarted advert replaces its record in each + * peer's cache instead of publishing a second one beside it. + */ + fun instanceName(address: String): String = + "op-" + WifiDirectGroupFormation.hex(WifiDirectGroupFormation.sha256(address), 8) + + /** + * The peer address a LAN record advertises, or null for a record that + * is not a peer hint: another version, no `addr`, or a service + * instance (`sid`, the DNS-SD mapping chapter), which names the same + * host as its peer record and would open a second stream to it. No + * `app` entry is required: iOS and Python publish none. + */ + fun lanAdvertAddress(attributes: Map): String? { + if (attributes.containsKey(WifiDirectGroupFormation.KEY_SERVICE_ID)) return null + if (attributes[WifiDirectGroupFormation.KEY_VERSION]?.decodeToString() != "1") return null + val address = attributes[WifiDirectGroupFormation.KEY_ADDRESS]?.decodeToString() ?: return null + return address.takeIf { it.startsWith("off1") } + } + } +} diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/LanPeerState.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/LanPeerState.kt new file mode 100644 index 00000000..b69ff3b0 --- /dev/null +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/LanPeerState.kt @@ -0,0 +1,164 @@ +package com.offlineprotocol + +/** + * The LAN carrier's resolves: what NSD reports, what waits to be resolved, + * and the one resolve in flight. Framework-free, so the rules are pinned by + * `LanPeerStateTest`; [LanPeerDiscovery] owns the NsdManager calls and the + * timers. Not thread-safe: the transport thread only. + * + * The rules, each of which a device found broken first: + * - One resolve at a time: before API 34 a second fails ALREADY_ACTIVE. + * - An answer is applied only while its resolve is the current one and NSD + * still reports the service. NSD sends one loss, so a record written for a + * service lost meanwhile would never be removed. + * - A resolve that gives up (failed, timed out, waited out ALREADY_ACTIVE) + * is queued again later while the service is reported: NSD reports it + * again only after losing it. + */ +internal class ResolveQueue( + private val nameOf: (S) -> String, + private val maxAlreadyActive: Int = 5, +) { + /** A resolve handed out by [next]. Answers are matched to it by identity. */ + class Ticket(val item: S) + + private val reported = HashMap() + private val pending = ArrayDeque() + private var inFlight: Ticket? = null + private var alreadyActive = 0 + + /** NSD reported [item]: queued for a resolve. */ + fun found(item: S) { + reported[nameOf(item)] = item + queue(nameOf(item)) + } + + /** NSD lost [name]: no longer queued, and an answer in flight for it is not applied. */ + fun lost(name: String) { + reported.remove(name) + pending.removeAll { nameOf(it) == name } + } + + /** Queues [name] again, if NSD still reports it and it is not already waiting. */ + fun queue(name: String) { + val item = reported[name] ?: return + if (pending.none { nameOf(it) == name }) pending.add(item) + } + + /** The next resolve to start, or null while one is in flight or none waits. */ + fun next(): Ticket? { + if (inFlight != null) return null + val item = pending.removeFirstOrNull() ?: return null + return Ticket(item).also { inFlight = it } + } + + /** + * [ticket] answered. True when the answer is to be applied: it is still + * the current resolve and its service is still reported. + */ + fun answered(ticket: Ticket): Boolean { + if (inFlight !== ticket) return false + inFlight = null + alreadyActive = 0 + return reported.containsKey(nameOf(ticket.item)) + } + + /** + * [ticket] failed or timed out. True when it was the current resolve; + * the caller then queues it again later ([queue]) and starts the next. + */ + fun gaveUp(ticket: Ticket): Boolean { + if (inFlight !== ticket) return false + inFlight = null + alreadyActive = 0 + return true + } + + /** + * [ticket] was answered ALREADY_ACTIVE: a resolve from before a restart + * still runs. True when it waits (put back at the head, still in flight + * until [release]); false once it has waited [maxAlreadyActive] times in + * a row, when the caller gives it up. + */ + fun waitForActive(ticket: Ticket): Boolean { + if (inFlight !== ticket) return false + if (++alreadyActive > maxAlreadyActive) { + alreadyActive = 0 + return false + } + if (reported.containsKey(nameOf(ticket.item))) pending.addFirst(ticket.item) + return true + } + + /** Ends [ticket]'s wait. True when it was still the current resolve. */ + fun release(ticket: Ticket): Boolean { + if (inFlight !== ticket) return false + inFlight = null + return true + } + + /** Whether NSD reports [name] now. */ + fun isReported(name: String): Boolean = reported.containsKey(name) + + /** Forgets everything; a resolve in flight is retired. */ + fun clear() { + reported.clear() + pending.clear() + inFlight = null + alreadyActive = 0 + } +} + +/** + * Which matching Wi-Fi network the LAN carrier follows. More than one can be + * up at once (Android's make-before-break switch, a local-only network + * another app asked for). The current one is kept until it is lost, then + * another is adopted: following whichever network spoke last tore down + * every LAN stream on each change from another, and could settle on one + * about to go, leaving none while another was still up. Framework-free and + * pinned by `LanPeerStateTest`. Not thread-safe: the transport thread only. + */ +internal class NetworkChoice { + enum class Up { + /** No network was current: adopt this one. */ + ADOPT, + /** The current network changed: update its properties. */ + UPDATE, + /** Another network: kept in case the current one is lost. */ + WAIT, + } + + private val candidates = LinkedHashMap() + + var current: N? = null + private set + + fun up(network: N, properties: P): Up { + candidates[network] = properties + return when (current) { + null -> { current = network; Up.ADOPT } + network -> Up.UPDATE + else -> Up.WAIT + } + } + + /** [network] was lost. True when it was the current one. */ + fun lost(network: N): Boolean { + candidates.remove(network) + if (network != current) return false + current = null + return true + } + + /** Adopts another matching network, if one is up, and returns it. */ + fun fallback(): Pair? { + val (network, properties) = candidates.entries.firstOrNull() ?: return null + current = network + return network to properties + } + + fun clear() { + candidates.clear() + current = null + } +} diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamDialPolicy.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamDialPolicy.kt new file mode 100644 index 00000000..8d9d0ad5 --- /dev/null +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamDialPolicy.kt @@ -0,0 +1,131 @@ +package com.offlineprotocol + +/** + * When the LAN carrier dials an advertised address, when it dials again, and + * who may open a stream to it. Framework-free; the manager owns the timers + * and sockets. Mirrors iOS's `PeerStreamDialPolicy` (PeerStreamSession.swift) + * case for case, keep in sync (ADR 0027, ADR 0028). + * + * Not thread-safe: [LanPeerDiscovery] calls it on the transport thread only. + */ +internal class PeerStreamDialPolicy { + /** Addresses with a dial scheduled or an outbound stream open. */ + private val dialing = HashSet() + private val redialDelayMs = HashMap() + + /** + * A peer advertised [address]. The delay to dial after, or null when no + * dial is due: one is already under way, or a stream holds the address. + */ + fun discovered(address: String, weAreLower: Boolean, held: Boolean): Long? = + schedule(address, held, if (weAreLower) 0L else HIGHER_ADDRESS_DELAY_MS) + + /** + * An outbound stream toward [address] ended. The delay to redial after, + * or null when the advert is gone or another stream holds the address. + */ + fun ended(address: String, advertised: Boolean, held: Boolean): Long? { + dialing.remove(address) + if (!advertised) return null + val delay = redialDelayMs[address] ?: REDIAL_INITIAL_DELAY_MS + schedule(address, held, delay) ?: return null + redialDelayMs[address] = minOf(delay * 2, REDIAL_MAX_DELAY_MS) + return delay + } + + /** + * A scheduled dial opened nothing because it is not due: the advert went, + * the address is now held, or the transport paused. + */ + fun abandoned(address: String) { + dialing.remove(address) + } + + /** + * A scheduled dial found every stream slot taken. The delay to try again + * after, on the redial ladder: abandoning it would lose the peer for + * good, since nothing else dials an address whose record did not change. + */ + fun noSlot(address: String): Long? = ended(address, advertised = true, held = false) + + /** + * An outbound stream toward [address] carried it (a body, or the address + * held through the keepalive window): the ladder starts over. + */ + fun proved(address: String) { + redialDelayMs.remove(address) + } + + fun reset() { + dialing.clear() + redialDelayMs.clear() + } + + private fun schedule(address: String, held: Boolean, delay: Long): Long? { + if (held || address in dialing) return null + dialing.add(address) + return delay + } + + companion object { + /** + * How long the higher address of a pair waits before dialing, so that + * the lower one's stream, which both ends keep, usually arrives first. + */ + const val HIGHER_ADDRESS_DELAY_MS = 5_000L + const val REDIAL_INITIAL_DELAY_MS = 1_000L + const val REDIAL_MAX_DELAY_MS = 60_000L + + /** + * Whether the listener takes a connection from [host], given the hosts + * of the inbound streams open now and the count of all open streams. + * + * The listener binds every interface, so on a shared network anyone + * on it can connect, not only group members. Two bounds keep that + * from costing this device its own dials: inbound streams take at + * most [PeerStreamSockets.Limits.maxInbound] of the budget, and one + * remote address at most [PeerStreamSockets.Limits.maxInboundPerHost] + * of those. Unknown hosts (null) share one bound. + */ + fun admitsInbound( + host: String?, + inboundFrom: Collection, + open: Int, + limits: PeerStreamSockets.Limits, + ): Boolean = + open < limits.maxStreams && + inboundFrom.size < limits.maxInbound && + inboundFrom.count { it == host } < limits.maxInboundPerHost + + /** + * The record each advertised address is dialed at, from every peer + * record held now ([records] as address to record, our own left out). + * + * Rebuilt whole on each change, never patched per record, because one + * address can be advertised by two records at once: a peer whose + * record outlived it beside the one it published on return. Removing + * a record by its address took the live record's address with it. Of + * two records for one address, one in [fresh] (just resolved) wins, + * then the one in [current], which a dial may be using. A record in + * [unprovable] is left out: its dial was answered by a peer that did + * not prove the address it advertises. + */ + fun adverts( + records: List>, + fresh: Set, + current: Map, + unprovable: Set = emptySet(), + ): Map { + fun rank(address: String, endpoint: E): Int = + if (endpoint in fresh) 2 else if (current[address] == endpoint) 1 else 0 + val out = HashMap() + for ((address, endpoint) in records) { + if (endpoint in unprovable) continue + val kept = out[address] + if (kept != null && rank(address, kept) >= rank(address, endpoint)) continue + out[address] = endpoint + } + return out + } + } +} diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamFraming.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamFraming.kt index bdef08f3..9a004379 100644 --- a/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamFraming.kt +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamFraming.kt @@ -150,12 +150,18 @@ class PeerStreamPreamble( * preamble is enough to open a second stream for a live address, so the count * has to be kept here, where the streams are. * - * Policy: the newer stream supersedes the older. On a phone the duplicate is - * almost always the same peer reconnecting past a stream that went half-open - * (a group client that re-joined), and - * refusing the newer one would leave that peer unreachable until the stale - * stream's socket noticed. The cost, recorded in R16, is that a replayer can - * choose when a real stream ends; it cannot use the stream it gets. + * Policy: the stream the lower address opened is kept, and between two of + * those the newer supersedes the older. Both ends compute it alike, which is + * the point: on a shared network both ends of a pair dial, and so do an iPhone + * and a Python host, so without a shared rule each end would keep the stream + * the other closes, and the pair would reconnect forever. It is the iOS + * `newStreamWins` and the Python manager's `_new_stream_wins`, and + * `every_peer_stream_manager_keeps_the_same_stream` pins the three copies + * together (ADR 0028). "Newer" among winners is what lets the lower address + * reconnect past its own half-open stream; the higher address's reconnect + * waits for keepalive to end the stale one. The cost, recorded in R16, is that + * a replayer can end a real stream when its copy is the winning kind; it + * cannot use the stream it gets. * * Thread-safe, and deliberately knows nothing of the protocol: the send path * reads it from a thread the core may be calling from while holding its global @@ -168,13 +174,46 @@ class PeerStreamLinks { val firstForAddress: Boolean, /** The older stream for the same address, to close without a loss report. */ val superseded: H?, + /** + * True when an older stream holds the address and wins: close this + * one without a report, and leave the older untouched. + */ + val refused: Boolean = false, ) + companion object { + /** + * Whether a new stream for [peer] takes the address over from the one + * that holds it. [outbound] is whether this device opened the new + * stream. Addresses compare by their UTF-8 bytes, which is the code + * point order Python's `<` uses; `String.compareTo` compares UTF-16 + * units and orders differently past the BMP. With no address of our + * own there is nothing to order by, and the announced stream stays. + */ + fun newStreamWins(outbound: Boolean, localAddress: String?, peer: String): Boolean { + val local = localAddress ?: return false + val weOpen = utf8Precedes(local, peer) + return outbound == weOpen + } + + private fun utf8Precedes(a: String, b: String): Boolean { + val x = a.encodeToByteArray() + val y = b.encodeToByteArray() + for (i in 0 until minOf(x.size, y.size)) { + val d = (x[i].toInt() and 0xff) - (y[i].toInt() and 0xff) + if (d != 0) return d < 0 + } + return x.size < y.size + } + } + private val byAddress = HashMap() private val byHandle = HashMap() @Synchronized - fun announce(handle: H, address: String): Announcement { + fun announce( + handle: H, address: String, outbound: Boolean, localAddress: String?, + ): Announcement { if (byHandle.containsKey(handle)) { // A stream proves one address, once. A second announcement, for // the same address or another, is a caller bug answered as a @@ -183,6 +222,11 @@ class PeerStreamLinks { // host app down, and this one behave the same. return Announcement(firstForAddress = false, superseded = null) } + if (byAddress.containsKey(address) && + !newStreamWins(outbound, localAddress, address) + ) { + return Announcement(firstForAddress = false, superseded = null, refused = true) + } val older = byAddress.put(address, handle) byHandle[handle] = address if (older != null) byHandle.remove(older) diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamSockets.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamSockets.kt index 70deae1b..6806f8df 100644 --- a/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamSockets.kt +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/PeerStreamSockets.kt @@ -1,5 +1,9 @@ package com.offlineprotocol +import android.os.Build +import android.os.ParcelFileDescriptor +import android.system.Os +import android.system.OsConstants import java.io.BufferedInputStream import java.io.DataInputStream import java.io.EOFException @@ -11,7 +15,6 @@ import java.util.concurrent.RejectedExecutionException import java.util.concurrent.ScheduledExecutorService import java.util.concurrent.ScheduledThreadPoolExecutor import java.util.concurrent.TimeUnit -import java.util.concurrent.atomic.AtomicInteger import java.util.concurrent.atomic.AtomicLong /** @@ -33,8 +36,9 @@ import java.util.concurrent.atomic.AtomicLong * [Host.peerConnected] once per address, each body attributed to it, and * [Host.peerDisconnected] once, if and only if the stream still held the * address when it ended. - * - One announced stream per address, newer superseding older, through - * [PeerStreamLinks]. A superseded stream is closed and reports nothing. + * - One announced stream per address, the one the lower address opened + * kept, through [PeerStreamLinks]. A superseded or refused stream is + * closed and reports nothing. * - A stream's announcement, deliveries and loss report run under that * stream's lock, so no body reaches the host after its loss report. The * core re-adds a neighbour on any inbound body, so a delivery that raced @@ -85,8 +89,27 @@ internal class PeerStreamSockets( val writeChunkBytes: Int = 64 * 1024, /** A writer blocked on one piece this long is stalled on the peer. */ val writeStallMs: Long = 2_000L, - /** Open sockets, proved or not. A group is a handful of devices. */ + /** Open sockets, proved or not. */ val maxStreams: Int = 16, + /** + * The inbound share of [maxStreams]. The rest is left to dials, so a + * listener full of strangers' sockets never stops this device + * reaching the peers it finds. iOS's `maxInbound`. + */ + val maxInbound: Int = 12, + /** + * Inbound sockets one remote address may hold, proved or not. A peer + * needs one, two while it reconnects past its own stale stream. + * iOS's `maxInboundPerHost`, Python's `MAX_STREAMS_PER_HOST` scaled + * to this budget. + */ + val maxInboundPerHost: Int = 4, + /** + * How long an announced stream must hold its address to count as + * having carried its peer with no body: the keepalive window. A + * refused stream ends within moments of its proof. + */ + val carriedAfterMs: Long = 30_000L, /** * Queued bytes toward one peer beyond which a body is dropped (the * core retries it), and only while that peer's writer is stalled. The @@ -98,12 +121,33 @@ internal class PeerStreamSockets( enum class SendResult { QUEUED, NO_STREAM, OUT_OF_BOUNDS, QUEUE_FULL } + /** + * What [run] saw. [proved] is the address the preamble proved, also when + * the stream then lost to the one already held; [heard] is whether the + * peer sent anything at all, which tells a liar from a silent socket; + * [carried] is whether the stream carried the peer: a body reached the + * host after the preamble, or the stream held its address for + * [Limits.carriedAfterMs]. + * + * Only [carried] starts a redial ladder over. A stream refused here + * because another holds the address, or closed by a peer whose stale + * stream wins there, proves the peer and carries nothing, and a ladder + * that reset on proof alone redialed it every second. A body alone was + * too narrow: two peers with a session may reconnect and send nothing, + * and every ordinary drop of such a stream climbed the ladder for good. + */ + data class Ran(val proved: String?, val heard: Boolean, val carried: Boolean = false) + + /** The link a stream rides, so one going down ends only its own streams. */ + enum class Carrier { P2P, LAN } + private val openStreams: MutableSet = ConcurrentHashMap.newKeySet() - // Slots taken against [Limits.maxStreams]. Reserved by an increment that - // decides, not by reading the set's size and adding afterwards: two - // sockets accepted at once both read a size under the limit, and both - // got in. - private val reserved = AtomicInteger(0) + // Slots taken against [Limits], decided and taken under one lock, not by + // reading the set's size and adding afterwards: two sockets accepted at + // once both read a size under the limit, and both got in. + private val admission = Any() + private var reserved = 0 + private val inboundHosts = ArrayList() private val links = PeerStreamLinks() // Daemon, like the stream threads, so an instance dropped without @@ -118,24 +162,52 @@ internal class PeerStreamSockets( /** Whether any stream has proved a peer: nothing is sendable otherwise. */ fun isEmpty(): Boolean = links.isEmpty() + /** + * Told each address whose announced stream ended, after the core was. + * The LAN carrier dials back a peer it still sees advertised. + */ + @Volatile var onLost: ((String) -> Unit)? = null + + /** Whether a stream holds [address] now. */ + fun holds(address: String): Boolean = links.handleFor(address) != null + + /** Whether an outbound socket would find a free slot now. */ + fun hasRoom(): Boolean = synchronized(admission) { reserved < limits.maxStreams } + /** Open sockets, proved or not. */ val openCount: Int get() = openStreams.size /** - * Runs [socket] to its end on the calling thread, and returns the address - * it proved, or null if it proved none. + * Runs [socket] to its end on the calling thread, and returns what it saw. * - * [accepting] is the owner's "still running": a socket that arrives while - * it is false, or past [Limits.maxStreams], is closed unread. + * [expected] is the address a dial was made toward, from a discovery + * record: a preamble proving any other is refused (the chapter's step + * four). [accepting] is the owner's "still running": a socket that + * arrives while it is false, or past [Limits], is closed unread. */ - fun run(socket: Socket, outbound: Boolean, accepting: () -> Boolean): String? { + fun run( + socket: Socket, + outbound: Boolean, + carrier: Carrier = Carrier.P2P, + expected: String? = null, + accepting: () -> Boolean, + ): Ran { val endpoint = socket.remoteSocketAddress?.toString() ?: "unknown" + var proved: String? = null + var heard = false + var delivered = false + var announcedAtNs = 0L + fun ran() = Ran( + proved, heard, + carried = delivered || ( + announcedAtNs != 0L && + (System.nanoTime() - announcedAtNs) / 1_000_000 >= limits.carriedAfterMs + ), + ) + val remoteHost = if (outbound) null else socket.inetAddress?.hostAddress val refusal = when { !accepting() -> "not running" - reserved.incrementAndGet() > limits.maxStreams -> { - reserved.decrementAndGet() - "at stream limit" - } + !reserve(outbound, remoteHost) -> "at stream limit" else -> null } if (refusal != null) { @@ -145,24 +217,24 @@ internal class PeerStreamSockets( "reason" to refusal, )) closeQuietly(socket) - return null + return ran() } // The slot is released by endStream, which runs once per stream. - val stream = Stream(socket) + val stream = Stream(socket, outbound, remoteHost, carrier) openStreams.add(stream) // Re-checked after the add: a closeAll() that snapshotted the set just // before it cannot miss this stream, because the owner stops // accepting before it closes everything. if (!accepting()) { endStream(stream) - return null + return ran() } - var proved: String? = null var why = "ended" try { try { socket.tcpNoDelay = true } catch (_: Exception) {} try { socket.keepAlive = true } catch (_: Exception) {} + tuneKeepalive(socket) // Ours goes first and without waiting for theirs, so neither side // can hold the other half-open by staying silent. @@ -170,12 +242,13 @@ internal class PeerStreamSockets( host.identityAssertion() } catch (e: Exception) { why = "no identity to present: ${e.message ?: "unknown"}" - return null + return ran() } stream.enqueue(PeerStreamFraming.frame(assertion)) val preamble = PeerStreamPreamble( verify = host::verify, + expected = expected, localAddress = try { host.localAddress() } catch (_: Exception) { null }, ) val input = DataInputStream(BufferedInputStream(socket.getInputStream())) @@ -185,25 +258,29 @@ internal class PeerStreamSockets( ) val address = try { val body = PeerStreamFraming.readBody(input, preamble = true) + heard = true when (val outcome = preamble.accept(body)) { is PeerStreamPreamble.Outcome.Announce -> outcome.address is PeerStreamPreamble.Outcome.Refuse -> { why = "preamble refused: ${outcome.reason}" - return null + return ran() } is PeerStreamPreamble.Outcome.Deliver -> { why = "preamble state out of order" - return null + return ran() } } } finally { deadline.cancel(false) } - if (!announce(stream, address)) { - why = "ended before the announcement" - return null - } + // Proved even if the held stream wins below: the peer is who it + // said, and the dial that found it needs no ladder. proved = address + if (!announce(stream, address, outbound)) { + why = "ended or refused before the announcement" + return ran() + } + announcedAtNs = System.nanoTime() host.diagnostic("info", "Peer stream proved", mapOf( "address" to address, "outbound" to outbound, @@ -213,8 +290,9 @@ internal class PeerStreamSockets( val body = PeerStreamFraming.readBody(input, preamble = false) if (!deliver(stream, address, body)) { why = "superseded or ended" - return proved + return ran() } + delivered = true } } catch (e: PeerStreamFraming.Refused) { why = "frame refused: ${e.reason}" @@ -229,7 +307,7 @@ internal class PeerStreamSockets( "reason" to why, )) } - return proved + return ran() } /** @@ -258,15 +336,31 @@ internal class PeerStreamSockets( } } + /** Ends every stream on [carrier], reporting each announced one lost. */ + fun closeCarrier(carrier: Carrier) { + for (stream in openStreams.toList()) { + if (stream.carrier == carrier) endStream(stream) + } + } + /** * Enters [stream] as the one announced stream for [address], and tells * the host if no stream held it. False when the stream was ended first, * in which case nothing was announced. */ - private fun announce(stream: Stream, address: String): Boolean { + private fun announce(stream: Stream, address: String, outbound: Boolean): Boolean { + val local = try { host.localAddress() } catch (_: Exception) { null } synchronized(stream.lock) { if (stream.closed) return false - val announcement = links.announce(stream, address) + val announcement = links.announce(stream, address, outbound, local) + if (announcement.refused) { + // The held stream wins: this one closes without a report, + // and the peer, computing the same rule, keeps the other too. + host.diagnostic("info", "Peer stream refused: the held one wins", mapOf( + "address" to address, + )) + return false + } announcement.superseded?.let { older -> // Closed without a loss report: `links` already moved the // address here, so the older stream's end finds nothing. @@ -313,11 +407,30 @@ internal class PeerStreamSockets( * or [closeAll]. Reports the loss if the stream still held its address, * under the lock the deliveries take. */ + private fun reserve(outbound: Boolean, host: String?): Boolean = synchronized(admission) { + val admitted = if (outbound) { + reserved < limits.maxStreams + } else { + PeerStreamDialPolicy.admitsInbound(host, inboundHosts, reserved, limits) + } + if (admitted) { + reserved++ + if (!outbound) inboundHosts.add(host) + } + admitted + } + + private fun release(stream: Stream) = synchronized(admission) { + reserved-- + if (!stream.outbound) inboundHosts.remove(stream.remoteHost) + } + private fun endStream(stream: Stream) { + var lost: String? = null synchronized(stream.lock) { if (stream.closed) return stream.closed = true - reserved.decrementAndGet() + release(stream) links.remove(stream)?.let { address -> try { host.peerDisconnected(address) @@ -326,13 +439,44 @@ internal class PeerStreamSockets( "error" to (e.message ?: "unknown"), )) } + lost = address } } + lost?.let { address -> onLost?.invoke(address) } openStreams.remove(stream) stream.abort() stream.shutdownWriter() } + /** + * Keepalive at 15 s idle, 5 s between probes, 3 probes, and a 30 s bound + * on unacknowledged data, the iOS and Python timers. The higher address's + * reconnect past a half-open stream is refused until the held one ends. + * With the OS defaults an idle dead stream holds the address for two + * hours, and keepalive sends no probe while data waits unacknowledged, + * so one with a write outstanding holds it until retransmission gives up + * (about fifteen minutes); `TCP_USER_TIMEOUT` bounds that case. + * + * Android 10 and later only. Before API 29 `fromSocket` takes the + * socket's own descriptor instead of a duplicate, so closing it closed + * the socket under every stream. Those versions keep the OS defaults. + * Best effort: a socket that refuses the options keeps the defaults. + */ + private fun tuneKeepalive(socket: Socket) { + if (Build.VERSION.SDK_INT < Build.VERSION_CODES.Q) return + try { + // A duplicate of the socket's descriptor from API 29: options set + // on it apply to the socket, and closing it leaves the socket open. + ParcelFileDescriptor.fromSocket(socket).use { pfd -> + val fd = pfd.fileDescriptor + Os.setsockoptInt(fd, OsConstants.IPPROTO_TCP, TCP_KEEPIDLE, 15) + Os.setsockoptInt(fd, OsConstants.IPPROTO_TCP, TCP_KEEPINTVL, 5) + Os.setsockoptInt(fd, OsConstants.IPPROTO_TCP, TCP_KEEPCNT, 3) + Os.setsockoptInt(fd, OsConstants.IPPROTO_TCP, TCP_USER_TIMEOUT, 30_000) + } + } catch (_: Throwable) {} + } + private fun closeQuietly(socket: Socket) { try { socket.close() } catch (_: Exception) {} } @@ -345,7 +489,13 @@ internal class PeerStreamSockets( * and splice a prefix into another frame's body) and a peer that stops * reading blocks only its own queue. */ - private inner class Stream(val socket: Socket) { + private inner class Stream( + val socket: Socket, + val outbound: Boolean, + /** The remote host an inbound socket came from, for [Limits.maxInboundPerHost]. */ + val remoteHost: String?, + val carrier: Carrier, + ) { /** Guards [closed] and the stream's calls into the host. */ val lock = Any() var closed = false @@ -417,3 +567,9 @@ internal class PeerStreamSockets( } } } + +// Linux ; `OsConstants` does not publish them. +private const val TCP_KEEPIDLE = 4 +private const val TCP_KEEPINTVL = 5 +private const val TCP_KEEPCNT = 6 +private const val TCP_USER_TIMEOUT = 18 diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectGroupFormation.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectGroupFormation.kt index d22b029a..4d49a803 100644 --- a/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectGroupFormation.kt +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectGroupFormation.kt @@ -265,9 +265,9 @@ internal object WifiDirectGroupFormation { fun isValidNetworkName(name: String): Boolean = name.length <= 32 && Regex("^DIRECT-[a-zA-Z0-9]{2}.*").matches(name) - private fun sha256(text: String): ByteArray = + internal fun sha256(text: String): ByteArray = MessageDigest.getInstance("SHA-256").digest(text.toByteArray(Charsets.UTF_8)) - private fun hex(bytes: ByteArray, count: Int): String = + internal fun hex(bytes: ByteArray, count: Int): String = bytes.take(count).joinToString("") { "%02x".format(it.toInt() and 0xff) } } diff --git a/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectManager.kt b/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectManager.kt index 523e12a3..6d706cc7 100644 --- a/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectManager.kt +++ b/bindings/react-native/android/src/main/java/com/offlineprotocol/WifiDirectManager.kt @@ -19,6 +19,7 @@ import androidx.core.content.ContextCompat import uniffi.offline_protocol.OfflineProtocol import uniffi.offline_protocol.verifyIdentityAssertion import java.io.* +import java.net.InetAddress import java.net.InetSocketAddress import java.net.ServerSocket import java.net.Socket @@ -47,8 +48,8 @@ import java.util.concurrent.atomic.AtomicLong * this manager hands each socket and each outbound body: the preamble under * one deadline, the frame bounds (the ceiling is inclusive, and a refused * length closes the socket rather than reading on; the reader this replaced - * got both wrong), one announced stream per address with the newer - * superseding, a per-stream writer, and the ordering that keeps a delivery + * got both wrong), one announced stream per address with the one the lower + * address opened kept, a per-stream writer, and the ordering that keeps a delivery * from landing after a loss report. It is framework-free and tested on real * loopback sockets. What stays here is the group: the listener, the one * outbound socket a client opens to its owner, and reconnecting that socket @@ -271,11 +272,24 @@ class WifiDirectManager( // State tracking. Volatile: written on the transport thread, read by the // socket thread that decides whether a client reconnects. @Volatile private var isGroupOwner = false - // Whether the core was last told the stream layer is up. Wi-Fi P2P going - // off reports it down; coming back on has to report it up again, or the - // core keeps the slot down while streams prove peers over it. + // Whether the core was last told the stream layer is up: while either + // carrier is ([reportLayer]). Wi-Fi P2P going off with no LAN reports it + // down; coming back on has to report it up again, or the core keeps the + // slot down while streams prove peers over it. @Volatile private var layerUp = false + // Whether Wi-Fi P2P is on, and whether this device is on a Wi-Fi network + // it can find LAN peers on. + @Volatile private var p2pUp = false + @Volatile private var lanUp = false + // The Wi-Fi network's own addresses: a socket accepted on one of them is + // a LAN stream, any other a group stream. + @Volatile private var lanAddresses: Set = emptySet() + // The LAN carrier, between start and stop. Transport thread only. + private var lan: LanPeerDiscovery? = null @Volatile private var groupOwnerAddress: String? = null + // The owner's proved address while another stream holds it and this + // client's own dial is parked; see [scheduleReconnect]. + @Volatile private var heldOwnerAddress: String? = null // The client's next reconnect delay; see [RECONNECT_INITIAL_DELAY_MS]. private val reconnectDelayMs = AtomicLong(RECONNECT_INITIAL_DELAY_MS) @@ -355,9 +369,13 @@ class WifiDirectManager( // MARK: - TransportManager Implementation - override fun isAvailable(): Boolean { - return wifiP2pManager != null && hasRequiredPermissions() - } + override fun isAvailable(): Boolean = p2pAllowed() || hasWifi() + + /** Wi-Fi P2P is usable. The LAN carrier needs none of its permissions. */ + private fun p2pAllowed(): Boolean = wifiP2pManager != null && hasRequiredPermissions() + + private fun hasWifi(): Boolean = + context.packageManager.hasSystemFeature(PackageManager.FEATURE_WIFI) @SuppressLint("MissingPermission") override fun start() { @@ -371,8 +389,9 @@ class WifiDirectManager( } if (!isAvailable()) { - throw TransportException.NotAvailable("WiFi P2P is not available on this device") + throw TransportException.NotAvailable("Neither Wi-Fi P2P nor Wi-Fi is available on this device") } + val p2p = p2pAllowed() Log.i(TAG, "Starting WiFi Direct transport for device: $deviceId") emitDiagnostic("info", "Starting WiFi Direct transport", mapOf( @@ -391,8 +410,12 @@ class WifiDirectManager( // callback is delivered on, and those callbacks share this manager's // state with the drain and the poll. Passing the main looper here was // what put the framework's half of this transport on the UI thread. - channel = wifiP2pManager?.initialize(context, confinement.looper) { - emitDiagnostic("warning", "WiFi P2P channel disconnected") + if (p2p) { + channel = wifiP2pManager?.initialize(context, confinement.looper) { + emitDiagnostic("warning", "WiFi P2P channel disconnected") + } + } else { + emitDiagnostic("warning", "Wi-Fi P2P unavailable or not permitted: LAN only") } // Register broadcast receiver. The scheduler overload delivers @@ -408,7 +431,9 @@ class WifiDirectManager( addAction(WifiP2pManager.WIFI_P2P_THIS_DEVICE_CHANGED_ACTION) } - if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU) { + if (!p2p) { + // Nothing to hear: no channel to act on what it would say. + } else if (Build.VERSION.SDK_INT >= Build.VERSION_CODES.TIRAMISU) { context.registerReceiver( p2pReceiver, intentFilter, @@ -428,12 +453,28 @@ class WifiDirectManager( updateState(TransportState.RUNNING) - // Notify protocol - layerUp = true - try { - protocol.wifiDirectStatusChanged(true) - } catch (e: Exception) { - Log.e(TAG, "Error notifying protocol of start", e) + // Notify protocol. Wi-Fi P2P off at start arrives next as the + // sticky state broadcast, which reports it. With no P2P the layer + // comes up when the LAN does. + p2pUp = p2p + reportLayer() + + lan = LanPeerDiscovery( + context = context, + handler = transportHandler, + executor = socketExecutor, + sockets = sockets, + port = SERVER_PORT, + localAddress = ::localAddressOrNull, + running = { state == TransportState.RUNNING }, + networkChanged = ::lanNetworkChanged, + diagnostic = ::emitDiagnostic, + ).also { it.start() } + sockets.onLost = { address -> + transportHandler.post { + lan?.peerLost(address) + ownerStreamLost(address) + } } // Start message polling @@ -458,6 +499,7 @@ class WifiDirectManager( */ @SuppressLint("MissingPermission") private fun adoptExistingGroup() { + if (channel == null) return wifiP2pManager?.requestConnectionInfo(channel) { info -> if (info?.groupFormed == true && state == TransportState.RUNNING) { emitDiagnostic("info", "Joining a group formed before start") @@ -486,6 +528,9 @@ class WifiDirectManager( // Stop server socket stopServerSocket() + lan?.stop() + lan = null + // Close all connections closeAllConnections() stopGroupFormation() @@ -493,6 +538,7 @@ class WifiDirectManager( // dial and a later start() waits for its own CONNECTION_CHANGED. isGroupOwner = false groupOwnerAddress = null + heldOwnerAddress = null reconnectDelayMs.set(RECONNECT_INITIAL_DELAY_MS) // Unregister receiver @@ -512,6 +558,9 @@ class WifiDirectManager( channel = null // Notify protocol + p2pUp = false + lanUp = false + lanAddresses = emptySet() layerUp = false try { protocol.wifiDirectStatusChanged(false) @@ -531,6 +580,7 @@ class WifiDirectManager( isPaused = true transportHandler.removeCallbacks(messagePollingRunnable) stopPeerDiscovery() + lan?.pause() } } @@ -538,6 +588,7 @@ class WifiDirectManager( confinement.runSync { isPaused = false if (state == TransportState.RUNNING) { + lan?.resume() startPeerDiscovery() // Drain what queued during the pause. Unlike the other three // managers, restarting the timer is not enough here: a poll @@ -614,6 +665,8 @@ class WifiDirectManager( @SuppressLint("MissingPermission") private fun startPeerDiscovery() { + // No channel: Wi-Fi P2P was unusable at start, and every call would throw. + if (channel == null) return if (!hasRequiredPermissions()) { emitDiagnostic("warning", "Missing permissions for peer discovery") return @@ -633,6 +686,7 @@ class WifiDirectManager( } private fun stopPeerDiscovery() { + if (channel == null) return wifiP2pManager?.stopPeerDiscovery(channel, object : WifiP2pManager.ActionListener { override fun onSuccess() { emitDiagnostic("info", "Peer discovery stopped") @@ -673,14 +727,11 @@ class WifiDirectManager( // Streams first, so each announced peer is reported lost // while the core still holds its link. A stream left open // across the flip would keep delivering, and every body - // re-adds a neighbour the flip just cleared. - closeAllConnections() - layerUp = false - try { - protocol.wifiDirectStatusChanged(false) - } catch (e: Exception) { - Log.e(TAG, "Error notifying protocol", e) - } + // re-adds a neighbour the flip just cleared. LAN streams do + // not ride P2P and stay. + sockets.closeCarrier(PeerStreamSockets.Carrier.P2P) + p2pUp = false + reportLayer() } } else if (enabled && state == TransportState.RUNNING) { // Wi-Fi P2P came back (Wi-Fi turned on, or the app started with it @@ -690,13 +741,9 @@ class WifiDirectManager( // Posted for the same reason as above; the flag makes a broadcast // that repeats the current state a no-op. transportHandler.post { - if (state != TransportState.RUNNING || layerUp) return@post - layerUp = true - try { - protocol.wifiDirectStatusChanged(true) - } catch (e: Exception) { - Log.e(TAG, "Error notifying protocol", e) - } + if (state != TransportState.RUNNING || p2pUp) return@post + p2pUp = true + reportLayer() startPeerDiscovery() resumeGroupFormationAfterP2pReturned() } @@ -802,6 +849,7 @@ class WifiDirectManager( } else { isGroupOwner = false groupOwnerAddress = null + heldOwnerAddress = null ownerClients = 0 createdNetwork = null ownsFormationGroup = false @@ -811,8 +859,9 @@ class WifiDirectManager( ownerIdleSinceMs = 0L // Posted: this is the broadcast receiver's onReceive, and ending // an announced stream is an FFI call. Same reasoning as the post - // in [handleWifiP2pStateChanged]. - transportHandler.post { closeAllConnections() } + // in [handleWifiP2pStateChanged]. Leaving the group ends only the + // group's streams. + transportHandler.post { sockets.closeCarrier(PeerStreamSockets.Carrier.P2P) } } } @@ -838,7 +887,7 @@ class WifiDirectManager( @SuppressLint("MissingPermission") private fun startGroupFormation() { - if (!formGroups) return + if (!formGroups || channel == null) return if (!formationSupported) { emitDiagnostic("info", "Wi-Fi Direct group formation needs Android 10; join a group from the system settings") return @@ -961,7 +1010,7 @@ class WifiDirectManager( val p2p = wifiP2pManager ?: return // Nothing to form over while Wi-Fi P2P is off: every call would fail // BUSY, and a join attempted then is refused with no group to find. - if (!layerUp) return + if (!p2pUp) return val now = SystemClock.elapsedRealtime() adverts.values.removeAll { now - it.seenAtMs > ADVERT_TTL_MS } deadOwners.values.removeAll { now - it >= WifiDirectGroupFormation.DEAD_OWNER_TTL_MS } @@ -1320,7 +1369,12 @@ class WifiDirectManager( private fun handleClientConnection(socket: Socket) { socketExecutor.execute { - sockets.run(socket, outbound = false) { state == TransportState.RUNNING } + val carrier = if (socket.localAddress in lanAddresses) { + PeerStreamSockets.Carrier.LAN + } else { + PeerStreamSockets.Carrier.P2P + } + sockets.run(socket, outbound = false, carrier) { state == TransportState.RUNNING } } } @@ -1330,7 +1384,7 @@ class WifiDirectManager( // open another socket to the same owner. if (!outboundOpen.compareAndSet(false, true)) return socketExecutor.execute { - var proved = false + var ran = PeerStreamSockets.Ran(proved = null, heard = false) try { val socket = Socket() try { @@ -1349,10 +1403,10 @@ class WifiDirectManager( // No claim to compare against: the owner's IP says nothing // about who it is, so the address its preamble derives is its // id (stream-framing.md, "The preamble"). - proved = sockets.run(socket, outbound = true) { state == TransportState.RUNNING } != null + ran = sockets.run(socket, outbound = true) { state == TransportState.RUNNING } } finally { outboundOpen.set(false) - scheduleReconnect(address, proved) + scheduleReconnect(address, ran) } } } @@ -1367,11 +1421,23 @@ class WifiDirectManager( * then the only dial left, so a check that the owner is still [ended] * would leave the client with no stream for the whole of the new group. */ - private fun scheduleReconnect(ended: String, proved: Boolean) { + private fun scheduleReconnect(ended: String, ran: PeerStreamSockets.Ran) { + // Another stream (a LAN one, say) holds the owner's address: this + // dial was refused for it, or superseded by it. The owner is alive + // and reachable, so there is nothing to redial until that stream is + // lost ([ownerStreamLost]). Published before the check, because the + // loss is handled on the transport thread: a stream lost between a + // check and a later write found nothing to redial, and parked this + // client with no stream until the group changed. In this order a + // loss either lands before the check, which then sees nothing held + // and redials below, or after it, and finds the address. + ran.proved?.let { heldOwnerAddress = it } + val held = ran.proved?.let { sockets.holds(it) } == true + if (!held && ran.proved != null && heldOwnerAddress == ran.proved) heldOwnerAddress = null val unproved = unprovedAtCeiling.updateAndGet { GroupOwnerRedial.unprovedAtCeilingAfter( count = it, - proved = proved, + proved = ran.proved != null, sameOwner = groupOwnerAddress == ended, currentDelayMs = reconnectDelayMs.get(), maxDelayMs = RECONNECT_MAX_DELAY_MS, @@ -1387,7 +1453,8 @@ class WifiDirectManager( } val plan = GroupOwnerRedial.next( ended = ended, - proved = proved, + carried = ran.carried, + held = held, running = state == TransportState.RUNNING, isGroupOwner = isGroupOwner, owner = groupOwnerAddress, @@ -1396,7 +1463,7 @@ class WifiDirectManager( maxDelayMs = RECONNECT_MAX_DELAY_MS, ) if (plan == null) { - if (proved) reconnectDelayMs.set(RECONNECT_INITIAL_DELAY_MS) + if (ran.carried) reconnectDelayMs.set(RECONNECT_INITIAL_DELAY_MS) return } reconnectDelayMs.set(plan.nextDelayMs) @@ -1420,6 +1487,53 @@ class WifiDirectManager( sockets.closeAll() } + /** + * An announced stream to [address] ended. If it was the one that held the + * group owner's address when this client's own dial was refused, the + * client dials its owner again. Transport thread. + */ + private fun ownerStreamLost(address: String) { + if (address != heldOwnerAddress) return + heldOwnerAddress = null + val owner = groupOwnerAddress + if (state == TransportState.RUNNING && !isGroupOwner && owner != null) { + connectToGroupOwner(owner) + } + } + + /** + * The LAN carrier's Wi-Fi network came up ([addresses] its own) or went + * (null). Going ends the LAN streams before the layer is reported, so + * each announced peer is reported lost while the core still holds it. + */ + private fun lanNetworkChanged(addresses: Set?) { + if (state != TransportState.RUNNING) return + if (addresses == null) { + lanAddresses = emptySet() + sockets.closeCarrier(PeerStreamSockets.Carrier.LAN) + lanUp = false + } else { + lanAddresses = addresses + lanUp = true + } + reportLayer() + } + + /** + * Tells the core the stream layer is up while either carrier is, once + * per edge. Runs on the transport thread: it is an FFI call. + */ + private fun reportLayer() { + val up = p2pUp || lanUp + if (up == layerUp) return + layerUp = up + try { + protocol.wifiDirectStatusChanged(up) + } catch (e: Exception) { + Log.e(TAG, "Error notifying protocol", e) + } + } + // MARK: - Message Handling (Event-Driven) /** @@ -1614,9 +1728,17 @@ internal object GroupOwnerRedial { fun shouldLeave(unprovedAtCeiling: Int, applicationGroup: Boolean): Boolean = applicationGroup && unprovedAtCeiling >= LEAVE_AFTER_UNPROVED_AT_CEILING + /** + * The next dial toward the owner, or null for none. [carried] is + * whether the stream that ended carried the owner + * ([PeerStreamSockets.Ran.carried]), which starts the ladder over; a + * proof alone does not, since a stream the owner refused proves it too. [held] is whether another + * stream holds the owner's address now: then nothing is redialed. + */ fun next( ended: String, - proved: Boolean, + carried: Boolean, + held: Boolean, running: Boolean, isGroupOwner: Boolean, owner: String?, @@ -1624,8 +1746,8 @@ internal object GroupOwnerRedial { initialDelayMs: Long, maxDelayMs: Long, ): Plan? { - if (!running || isGroupOwner || owner == null) return null - val delay = if (proved || owner != ended) initialDelayMs else currentDelayMs + if (!running || isGroupOwner || owner == null || held) return null + val delay = if (carried || owner != ended) initialDelayMs else currentDelayMs return Plan(owner, delay, minOf(delay * 2, maxDelayMs)) } } diff --git a/bindings/react-native/android/src/test/java/com/offlineprotocol/GroupOwnerRedialTest.kt b/bindings/react-native/android/src/test/java/com/offlineprotocol/GroupOwnerRedialTest.kt index 64d0169c..af7ff5d2 100644 --- a/bindings/react-native/android/src/test/java/com/offlineprotocol/GroupOwnerRedialTest.kt +++ b/bindings/react-native/android/src/test/java/com/offlineprotocol/GroupOwnerRedialTest.kt @@ -15,14 +15,16 @@ class GroupOwnerRedialTest { private fun next( ended: String = "192.168.49.1", - proved: Boolean = false, + carried: Boolean = false, + held: Boolean = false, running: Boolean = true, isGroupOwner: Boolean = false, owner: String? = "192.168.49.1", currentDelayMs: Long = 8_000, ) = GroupOwnerRedial.next( ended = ended, - proved = proved, + carried = carried, + held = held, running = running, isGroupOwner = isGroupOwner, owner = owner, @@ -42,8 +44,20 @@ class GroupOwnerRedialTest { } @Test - fun `a stream that proved its peer starts the ladder over`() { - assertEquals(GroupOwnerRedial.Plan("192.168.49.1", 1_000, 2_000), next(proved = true)) + fun `a stream that carried its peer starts the ladder over`() { + assertEquals(GroupOwnerRedial.Plan("192.168.49.1", 1_000, 2_000), next(carried = true)) + } + + @Test + fun `a stream the owner proved but refused climbs the ladder`() { + // The owner's stale stream wins there: a reset here redialed every second. + assertEquals(GroupOwnerRedial.Plan("192.168.49.1", 8_000, 16_000), next(carried = false)) + } + + @Test + fun `no redial while another stream holds the owner's address`() { + assertNull(next(held = true)) + assertNull(next(held = true, carried = true)) } @Test diff --git a/bindings/react-native/android/src/test/java/com/offlineprotocol/LanPeerDiscoveryTest.kt b/bindings/react-native/android/src/test/java/com/offlineprotocol/LanPeerDiscoveryTest.kt new file mode 100644 index 00000000..92984960 --- /dev/null +++ b/bindings/react-native/android/src/test/java/com/offlineprotocol/LanPeerDiscoveryTest.kt @@ -0,0 +1,45 @@ +package com.offlineprotocol + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertNull +import org.junit.Test + +/** + * The record a LAN peer is found by. Browsing, resolving and dialing need a + * device and a network; this is what the manager reads off a record. + */ +class LanPeerDiscoveryTest { + private val address = "off1qysluvwl5922yctzd0u9gpr06gn3k7ldfvgtwgvn" + + private fun record(vararg entries: Pair): Map = + entries.associate { (key, value) -> key to value?.encodeToByteArray() } + + @Test + fun `the instance name is the one iOS and Python give the same address`() { + assertEquals("op-b701bfc1c92768b5", LanPeerDiscovery.instanceName(address)) + } + + @Test + fun `an iOS or Python record names its address`() { + assertEquals(address, LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "1", "addr" to address))) + } + + @Test + fun `an app entry is allowed, not required`() { + assertEquals( + address, + LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "1", "addr" to address, "app" to "0a1b2c3d")), + ) + } + + @Test + fun `records that are not peer hints are ignored`() { + // A service instance under the DNS-SD mapping chapter. + assertNull(LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "1", "addr" to address, "sid" to "x"))) + assertNull(LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "2", "addr" to address))) + assertNull(LanPeerDiscovery.lanAdvertAddress(record("addr" to address))) + assertNull(LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "1"))) + assertNull(LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "1", "addr" to null))) + assertNull(LanPeerDiscovery.lanAdvertAddress(record("txtvers" to "1", "addr" to "not-an-address"))) + } +} diff --git a/bindings/react-native/android/src/test/java/com/offlineprotocol/LanPeerStateTest.kt b/bindings/react-native/android/src/test/java/com/offlineprotocol/LanPeerStateTest.kt new file mode 100644 index 00000000..15b8cecd --- /dev/null +++ b/bindings/react-native/android/src/test/java/com/offlineprotocol/LanPeerStateTest.kt @@ -0,0 +1,186 @@ +package com.offlineprotocol + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Test + +/** + * [ResolveQueue] and [NetworkChoice]: the LAN carrier's state, each rule one + * a device found broken first. NsdManager and the timers need a device; this + * is the bookkeeping they drive. + */ +class LanPeerStateTest { + + private fun queue(maxAlreadyActive: Int = 5) = ResolveQueue({ it }, maxAlreadyActive) + + // --- one resolve at a time ------------------------------------------------- + + @Test + fun `one resolve is in flight at a time, in the order found`() { + val q = queue() + q.found("a") + q.found("b") + val first = q.next()!! + assertEquals("a", first.item) + assertNull("one at a time", q.next()) + assertTrue(q.answered(first)) + assertEquals("b", q.next()!!.item) + } + + @Test + fun `a service waits once however often it is queued`() { + val q = queue() + q.found("a") + q.queue("a") + q.found("a") + q.next()!!.let { assertTrue(q.answered(it)) } + assertNull(q.next()) + } + + // --- a service lost meanwhile ------------------------------------------------ + + @Test + fun `an answer for a service lost while resolving is not applied`() { + // NSD sends one loss: a record written after it is never removed. + val q = queue() + q.found("a") + val ticket = q.next()!! + q.lost("a") + assertFalse(q.answered(ticket)) + assertNull("nothing left to resolve", q.next()) + } + + @Test + fun `a lost service is no longer queued, and cannot be queued again`() { + val q = queue() + q.found("a") + q.found("b") + q.lost("a") + q.queue("a") + assertEquals("b", q.next()!!.item) + } + + // --- a resolve that gives up ------------------------------------------------- + + @Test + fun `a resolve that gave up can be queued again while reported`() { + val q = queue() + q.found("a") + val ticket = q.next()!! + assertTrue(q.gaveUp(ticket)) + q.queue("a") + assertEquals("a", q.next()!!.item) + } + + @Test + fun `a late answer after the watchdog gave up is ignored, and the queue moved on`() { + val q = queue() + q.found("a") + q.found("b") + val stuck = q.next()!! + assertTrue("the watchdog retires it", q.gaveUp(stuck)) + val next = q.next()!! + assertEquals("b", next.item) + assertFalse("the stuck one answering now changes nothing", q.answered(stuck)) + assertTrue(q.answered(next)) + } + + @Test + fun `clearing retires the resolve in flight`() { + val q = queue() + q.found("a") + val ticket = q.next()!! + q.clear() + assertFalse(q.answered(ticket)) + assertFalse(q.gaveUp(ticket)) + assertFalse(q.isReported("a")) + } + + // --- ALREADY_ACTIVE ------------------------------------------------------------ + + @Test + fun `ALREADY_ACTIVE waits at the head of the queue a few times, then gives way`() { + val q = queue(maxAlreadyActive = 2) + q.found("a") + q.found("b") + repeat(2) { + val ticket = q.next()!! + assertEquals("a waits at the head", "a", ticket.item) + assertTrue(q.waitForActive(ticket)) + assertNull("still in flight while it waits", q.next()) + assertTrue(q.release(ticket)) + } + val third = q.next()!! + assertEquals("a", third.item) + assertFalse("waited enough: given up", q.waitForActive(third)) + assertTrue(q.gaveUp(third)) + assertEquals("the queue moves on", "b", q.next()!!.item) + } + + @Test + fun `an answer resets the ALREADY_ACTIVE count`() { + val q = queue(maxAlreadyActive = 1) + q.found("a") + q.next()!!.let { assertTrue(q.waitForActive(it)); assertTrue(q.release(it)) } + q.next()!!.let { assertTrue(q.answered(it)) } + q.queue("a") + q.next()!!.let { assertTrue("counted afresh", q.waitForActive(it)) } + } + + @Test + fun `a service lost while waiting out ALREADY_ACTIVE is not put back`() { + val q = queue() + q.found("a") + val ticket = q.next()!! + q.lost("a") + assertTrue(q.waitForActive(ticket)) + assertTrue(q.release(ticket)) + assertNull(q.next()) + } + + // --- the network --------------------------------------------------------------- + + @Test + fun `the first matching network is adopted and kept while another comes and goes`() { + val n = NetworkChoice() + assertEquals(NetworkChoice.Up.ADOPT, n.up("A", 1)) + assertEquals(NetworkChoice.Up.WAIT, n.up("B", 1)) + assertEquals("a change on another network moves nothing", NetworkChoice.Up.WAIT, n.up("B", 2)) + assertEquals(NetworkChoice.Up.UPDATE, n.up("A", 2)) + assertFalse(n.lost("B")) + assertEquals("A", n.current) + } + + @Test + fun `losing the current network falls back to another still up`() { + val n = NetworkChoice() + n.up("A", 1) + n.up("B", 7) + assertTrue(n.lost("A")) + assertNull(n.current) + assertEquals("B" to 7, n.fallback()) + assertEquals("B", n.current) + } + + @Test + fun `losing the last network leaves none`() { + val n = NetworkChoice() + n.up("A", 1) + assertTrue(n.lost("A")) + assertNull(n.fallback()) + assertNull(n.current) + assertEquals("the next one up is adopted", NetworkChoice.Up.ADOPT, n.up("C", 1)) + } + + @Test + fun `clearing forgets every network`() { + val n = NetworkChoice() + n.up("A", 1) + n.up("B", 1) + n.clear() + assertNull(n.current) + assertNull(n.fallback()) + } +} diff --git a/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamDialPolicyTest.kt b/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamDialPolicyTest.kt new file mode 100644 index 00000000..9a661bac --- /dev/null +++ b/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamDialPolicyTest.kt @@ -0,0 +1,157 @@ +package com.offlineprotocol + +import org.junit.Assert.assertEquals +import org.junit.Assert.assertFalse +import org.junit.Assert.assertNotNull +import org.junit.Assert.assertNull +import org.junit.Assert.assertTrue +import org.junit.Test + +/** Mirrors iOS's PeerStreamDialPolicyTests, case for case where both have one. */ +class PeerStreamDialPolicyTest { + private val limits = PeerStreamSockets.Limits() + + /** One machine on the LAN cannot fill the listener, and a full listener leaves room to dial. */ + @Test + fun `inbound is bounded per host and leaves room to dial`() { + val one = List(limits.maxInboundPerHost) { "10.0.0.9" } + assertFalse(PeerStreamDialPolicy.admitsInbound("10.0.0.9", one, one.size, limits)) + assertTrue(PeerStreamDialPolicy.admitsInbound("10.0.0.7", one, one.size, limits)) + + val many = List(limits.maxInbound) { "10.0.1.$it" } + assertFalse(PeerStreamDialPolicy.admitsInbound("10.0.2.1", many, many.size, limits)) + assertTrue("dials keep a share of the budget", limits.maxInbound < limits.maxStreams) + assertFalse(PeerStreamDialPolicy.admitsInbound("10.0.2.1", emptyList(), limits.maxStreams, limits)) + } + + @Test + fun `unknown hosts share one bound`() { + val unknown = List(limits.maxInboundPerHost) { null } + assertFalse(PeerStreamDialPolicy.admitsInbound(null, unknown, unknown.size, limits)) + } + + private val peer = "off1qysluvwl5922yctzd0u9gpr06gn3k7ldfvgtwgvn" + private val higher = PeerStreamDialPolicy.HIGHER_ADDRESS_DELAY_MS + private val initial = PeerStreamDialPolicy.REDIAL_INITIAL_DELAY_MS + + @Test + fun `the lower address dials at once and the higher waits`() { + assertEquals(0L, PeerStreamDialPolicy().discovered(peer, weAreLower = true, held = false)) + assertEquals(higher, PeerStreamDialPolicy().discovered(peer, weAreLower = false, held = false)) + } + + @Test + fun `no dial while one is under way or the address is held`() { + val policy = PeerStreamDialPolicy() + assertNotNull(policy.discovered(peer, weAreLower = true, held = false)) + assertNull("one dial at a time", policy.discovered(peer, weAreLower = true, held = false)) + assertNull(PeerStreamDialPolicy().discovered(peer, weAreLower = true, held = true)) + } + + @Test + fun `an abandoned dial lets the next advert dial`() { + val policy = PeerStreamDialPolicy() + policy.discovered(peer, weAreLower = true, held = false) + policy.abandoned(peer) + assertEquals(0L, policy.discovered(peer, weAreLower = true, held = false)) + } + + @Test + fun `the ladder doubles to its cap`() { + val policy = PeerStreamDialPolicy() + policy.discovered(peer, weAreLower = true, held = false) + val delays = List(9) { policy.ended(peer, advertised = true, held = false) } + assertEquals( + listOf(1L, 2L, 4L, 8L, 16L, 32L, 60L, 60L, 60L).map { it * 1_000 }, + delays, + ) + } + + @Test + fun `a lost race does not climb the ladder`() { + // The higher address's dial loses to the lower's stream: its stream + // ends while the address is held. That is no redial, so no step. + val policy = PeerStreamDialPolicy() + repeat(5) { + policy.discovered(peer, weAreLower = false, held = false) + assertNull(policy.ended(peer, advertised = true, held = true)) + } + // When the held stream later ends, the first redial is the first step. + assertEquals(higher, policy.discovered(peer, weAreLower = false, held = false)) + assertEquals(initial, policy.ended(peer, advertised = true, held = false)) + } + + @Test + fun `no redial once the advert is gone`() { + val policy = PeerStreamDialPolicy() + policy.discovered(peer, weAreLower = true, held = false) + assertNull(policy.ended(peer, advertised = false, held = false)) + assertEquals( + "the ended stream no longer counts as a dial under way", + 0L, policy.discovered(peer, weAreLower = true, held = false), + ) + } + + @Test + fun `a proved stream starts the ladder over`() { + val policy = PeerStreamDialPolicy() + policy.discovered(peer, weAreLower = true, held = false) + policy.ended(peer, advertised = true, held = false) + policy.ended(peer, advertised = true, held = false) + policy.proved(peer) + assertEquals(initial, policy.ended(peer, advertised = true, held = false)) + } + + @Test + fun `reset forgets everything`() { + val policy = PeerStreamDialPolicy() + policy.discovered(peer, weAreLower = true, held = false) + policy.ended(peer, advertised = true, held = false) + policy.reset() + assertEquals(0L, policy.discovered(peer, weAreLower = true, held = false)) + assertEquals(initial, policy.ended(peer, advertised = true, held = false)) + } + + /** A peer back before its old record expired is advertised by two records. */ + @Test + fun `removing a stale record keeps the live one's address`() { + var adverts = PeerStreamDialPolicy.adverts(listOf(peer to "old"), setOf("old"), emptyMap()) + adverts = PeerStreamDialPolicy.adverts(listOf(peer to "old", peer to "new"), setOf("new"), adverts) + assertEquals("the record a change added is the newest", "new", adverts[peer]) + adverts = PeerStreamDialPolicy.adverts(listOf(peer to "new"), emptySet(), adverts) + assertEquals("new", adverts[peer]) + assertNull( + "the last record going takes the address", + PeerStreamDialPolicy.adverts(emptyList>(), emptySet(), adverts)[peer], + ) + } + + @Test + fun `the record in use stays over an older one`() { + val records = listOf(peer to "a", peer to "b") + for (current in listOf("a", "b")) { + assertEquals(current, PeerStreamDialPolicy.adverts(records, emptySet(), mapOf(peer to current))[peer]) + } + } + + /** A dial that found no free slot is tried again, later each time, and stays the one dial. */ + @Test + fun `a dial with no slot is tried again on the ladder`() { + val policy = PeerStreamDialPolicy() + policy.discovered(peer, weAreLower = true, held = false) + assertEquals(initial, policy.noSlot(peer)) + assertNull("one dial at a time", policy.discovered(peer, weAreLower = true, held = false)) + assertEquals(2 * initial, policy.noSlot(peer)) + } + + /** A record that answered without proving its address gives way, or leaves the address unadvertised. */ + @Test + fun `an unprovable record is left out`() { + val records = listOf(peer to "liar", peer to "real") + assertEquals( + "real", + PeerStreamDialPolicy.adverts(records, setOf("liar"), mapOf(peer to "liar"), setOf("liar"))[peer], + ) + assertNull(PeerStreamDialPolicy.adverts(listOf(peer to "liar"), emptySet(), emptyMap(), setOf("liar"))[peer]) + } +} diff --git a/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamFramingTest.kt b/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamFramingTest.kt index 03fcf918..676bb645 100644 --- a/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamFramingTest.kt +++ b/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamFramingTest.kt @@ -224,6 +224,12 @@ class PeerStreamFramingTest { private val a = "off1qysluvwl5922yctzd0u9gpr06gn3k7ldfvgtwgvn" private val b = "off1qyqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqn8antf" + private val me = "off1qqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqqq" + + // Swift's helper: we are the lowest address, so an outbound stream wins. + private fun PeerStreamLinks.announce( + handle: String, address: String, outbound: Boolean = true, + ) = announce(handle, address, outbound, me) @Test fun `the first stream for an address is announced`() { @@ -255,6 +261,71 @@ class PeerStreamFramingTest { assertTrue(links.isEmpty()) } + @Test + fun `the first stream is announced whichever kind it is`() { + val links = PeerStreamLinks() + assertEquals( + PeerStreamLinks.Announcement(true, null), + links.announce("in", a, outbound = false), + ) + } + + @Test + fun `a second stream of the losing kind is refused and the held one stays`() { + val links = PeerStreamLinks() + links.announce("out", a) + assertEquals( + PeerStreamLinks.Announcement(false, null, refused = true), + links.announce("in", a, outbound = false), + ) + assertEquals("out", links.handleFor(a)) + assertNull("a refused stream was never announced", links.remove("in")) + assertEquals(a, links.remove("out")) + } + + @Test + fun `the winning kind is whatever the lower address opened`() { + // Python's `_new_stream_wins` and the Swift copy, case for case (ADR 0028). + assertTrue(PeerStreamLinks.newStreamWins(outbound = true, localAddress = me, peer = a)) + assertFalse(PeerStreamLinks.newStreamWins(outbound = false, localAddress = me, peer = a)) + assertFalse(PeerStreamLinks.newStreamWins(outbound = true, localAddress = a, peer = me)) + assertTrue(PeerStreamLinks.newStreamWins(outbound = false, localAddress = a, peer = me)) + assertTrue( + "b sorts below a at the first differing character", + PeerStreamLinks.newStreamWins(outbound = true, localAddress = b, peer = a), + ) + assertFalse(PeerStreamLinks.newStreamWins(outbound = true, localAddress = null, peer = a)) + assertFalse(PeerStreamLinks.newStreamWins(outbound = false, localAddress = null, peer = a)) + } + + @Test + fun `addresses compare by their UTF-8 bytes, not UTF-16 units`() { + // U+FF5E is one UTF-16 unit above the surrogate U+D83D that starts + // U+1F600, so String.compareTo puts it after; in UTF-8 (EF.. vs F0..) + // it comes first, as Python's `<` does. + val bmp = "\uFF5E" + val astral = "\uD83D\uDE00" + assertTrue(bmp > astral) + assertTrue(PeerStreamLinks.newStreamWins(outbound = true, localAddress = bmp, peer = astral)) + } + + @Test + fun `both ends of a pair keep the same stream`() { + // `mine` is the stream the lower address (me) opened; `theirs` the + // other. Each end sees one as outbound and the other as inbound. + for (firstArrives in listOf("mine", "theirs")) { + val low = PeerStreamLinks() + val high = PeerStreamLinks() + val order = if (firstArrives == "mine") listOf("mine", "theirs") else listOf("theirs", "mine") + for (stream in order) { + low.announce(stream, a, stream == "mine", me) + high.announce(stream, me, stream == "theirs", a) + } + assertEquals(firstArrives, "mine", low.handleFor(a)) + assertEquals(firstArrives, "mine", high.handleFor(me)) + } + } + @Test fun `a stream that never proved a peer reports nothing`() { val links = PeerStreamLinks() diff --git a/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamSocketsTest.kt b/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamSocketsTest.kt index d82ac3a1..29cf6c6a 100644 --- a/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamSocketsTest.kt +++ b/bindings/react-native/android/src/test/java/com/offlineprotocol/PeerStreamSocketsTest.kt @@ -83,14 +83,18 @@ class PeerStreamSocketsTest { ) /** A listener whose accepted sockets [sockets] runs, as the group owner's does. */ - private fun listen(sockets: PeerStreamSockets, running: () -> Boolean = { true }): Int { + private fun listen( + sockets: PeerStreamSockets, + carrier: PeerStreamSockets.Carrier = PeerStreamSockets.Carrier.P2P, + running: () -> Boolean = { true }, + ): Int { val server = ServerSocket(0, 50, InetAddress.getLoopbackAddress()) cleanup.add(server) thread(isDaemon = true) { while (!server.isClosed) { val s = try { server.accept() } catch (_: IOException) { return@thread } cleanup.add(s) - thread(isDaemon = true) { sockets.run(s, outbound = false, accepting = running) } + thread(isDaemon = true) { sockets.run(s, outbound = false, carrier, accepting = running) } } } return server.localPort @@ -320,6 +324,23 @@ class PeerStreamSocketsTest { assertArrayEquals(assertion("peer-a"), later.readBody()) } + @Test + fun `one remote host holds at most its inbound share`() { + // Every test socket comes from loopback, one host. + val host = Host("peer-a") + val a = PeerStreamSockets(host, fast.copy(maxInboundPerHost = 2, preambleTimeoutMs = 5_000)) + val port = listen(a) + val first = raw(port).also { it.readBody() } + raw(port).readBody() + assertTrue(raw(port).closedByPeer()) + + // Ending one frees its share, once the listener side has seen it end. + first.close() + val until = System.currentTimeMillis() + 5_000 + while (a.openCount > 1 && System.currentTimeMillis() < until) Thread.sleep(10) + assertArrayEquals(assertion("peer-a"), raw(port).readBody()) + } + @Test fun `a socket that arrives while stopped is closed unread`() { val host = Host("peer-a") @@ -331,8 +352,9 @@ class PeerStreamSocketsTest { // --- one announced stream per address --------------------------------------------- @Test - fun `a newer stream supersedes the older with one announcement and one loss`() { - val host = Host("peer-a") + fun `a newer stream the lower address opened supersedes the older with one announcement and one loss`() { + // peer-b is lower than us, so its inbound streams are the winning kind. + val host = Host("peer-c") val port = listen(PeerStreamSockets(host, fast)) val older = raw(port) older.readBody() @@ -355,6 +377,117 @@ class PeerStreamSocketsTest { assertNull(host.quiet()) } + @Test + fun `a second stream the higher address opened is refused and the held one stays`() { + // peer-b is higher than us: we keep its first stream until it ends. + val host = Host("peer-a") + val port = listen(PeerStreamSockets(host, fast)) + val held = raw(port) + held.readBody() + held.send(assertion("peer-b")) + assertEquals("connected:peer-b", host.next()) + + val second = raw(port) + second.readBody() + second.send(assertion("peer-b")) + assertTrue(second.closedByPeer()) + assertNull(host.quiet()) + + held.send("still here".toByteArray()) + assertEquals("message:peer-b:10", host.next()) + held.close() + assertEquals("lost:peer-b", host.next()) + assertNull(host.quiet()) + } + + @Test + fun `closing one carrier ends only its streams`() { + val host = Host("peer-a") + val a = PeerStreamSockets(host, fast) + val group = raw(listen(a, PeerStreamSockets.Carrier.P2P)).also { it.readBody() } + group.send(assertion("peer-b")) + assertEquals("connected:peer-b", host.next()) + val lan = raw(listen(a, PeerStreamSockets.Carrier.LAN)).also { it.readBody() } + lan.send(assertion("peer-c")) + assertEquals("connected:peer-c", host.next()) + + a.closeCarrier(PeerStreamSockets.Carrier.P2P) + assertEquals("lost:peer-b", host.next()) + assertTrue(group.closedByPeer()) + assertNull(host.quiet()) + + lan.send("still here".toByteArray()) + assertEquals("message:peer-c:10", host.next()) + } + + @Test + fun `run reports a proof, and carried only once a body arrives`() { + val host = Host("peer-a") + val a = PeerStreamSockets(host, fast) + val results = LinkedBlockingQueue() + val server = ServerSocket(0, 50, InetAddress.getLoopbackAddress()) + cleanup.add(server) + thread(isDaemon = true) { + while (!server.isClosed) { + val s = try { server.accept() } catch (_: IOException) { return@thread } + cleanup.add(s) + thread(isDaemon = true) { results.add(a.run(s, outbound = false) { true }) } + } + } + + // Held by peer-b's first stream, which is higher than us: refused here. + val held = raw(server.localPort).also { it.readBody() } + held.send(assertion("peer-b")) + assertEquals("connected:peer-b", host.next()) + val refused = raw(server.localPort).also { it.readBody() } + refused.send(assertion("peer-b")) + assertEquals( + PeerStreamSockets.Ran(proved = "peer-b", heard = true, carried = false), + results.poll(5, TimeUnit.SECONDS), + ) + + held.send("body".toByteArray()) + assertEquals("message:peer-b:4", host.next()) + held.close() + assertEquals( + PeerStreamSockets.Ran(proved = "peer-b", heard = true, carried = true), + results.poll(5, TimeUnit.SECONDS), + ) + + // A preamble that does not verify: heard, proved nothing. + val liar = raw(server.localPort).also { it.readBody() } + liar.send(assertion("bad-x")) + assertEquals( + PeerStreamSockets.Ran(proved = null, heard = true, carried = false), + results.poll(5, TimeUnit.SECONDS), + ) + } + + @Test + fun `an idle stream that held its address through the window carried its peer`() { + // Two peers with a session may reconnect and send nothing: a stream + // that held its address that long was no refusal. + val host = Host("peer-a") + val a = PeerStreamSockets(host, fast.copy(carriedAfterMs = 200)) + val results = LinkedBlockingQueue() + val server = ServerSocket(0, 50, InetAddress.getLoopbackAddress()) + cleanup.add(server) + thread(isDaemon = true) { + val s = try { server.accept() } catch (_: IOException) { return@thread } + cleanup.add(s) + results.add(a.run(s, outbound = false) { true }) + } + val idle = raw(server.localPort).also { it.readBody() } + idle.send(assertion("peer-b")) + assertEquals("connected:peer-b", host.next()) + Thread.sleep(400) + idle.close() + assertEquals( + PeerStreamSockets.Ran(proved = "peer-b", heard = true, carried = true), + results.poll(5, TimeUnit.SECONDS), + ) + } + @Test fun `closeAll reports each announced stream once and delivers nothing after`() { val host = Host("peer-a") diff --git a/bindings/react-native/ios/PeerStreamFraming.swift b/bindings/react-native/ios/PeerStreamFraming.swift index 3b333306..2312752a 100644 --- a/bindings/react-native/ios/PeerStreamFraming.swift +++ b/bindings/react-native/ios/PeerStreamFraming.swift @@ -237,7 +237,7 @@ final class PeerStreamPreamble { /// the point: both ends of a pair may dial, and so does a Python host on the /// same LAN, so without a shared rule each end would keep the stream the other /// closes, and the pair would reconnect forever. It is the Python manager's -/// `_new_stream_wins`, and `ios_and_python_peer_streams_keep_the_same_stream` +/// `_new_stream_wins`, and `every_peer_stream_manager_keeps_the_same_stream` /// pins the two copies together (ADR 0027). "Newer" among winners is what /// lets the lower address reconnect past its own half-open stream; the higher /// address's reconnect waits for keepalive to end the stale one. The cost, diff --git a/bindings/react-native/ios/WifiDirectManager.swift b/bindings/react-native/ios/WifiDirectManager.swift index 733999da..baa93afb 100644 --- a/bindings/react-native/ios/WifiDirectManager.swift +++ b/bindings/react-native/ios/WifiDirectManager.swift @@ -79,6 +79,10 @@ public class WifiDirectManager: NSObject, TransportManager { /// it is ended and left to the redial ladder. Python's `CONNECT_TIMEOUT` /// and Android's connect timeout bound the whole connect the same way. private static let DIAL_TIMEOUT: TimeInterval = 10.0 + /// How long a dialed stream must hold its address after its proof to + /// count as having carried its peer with no body: the keepalive window. + /// Android's `Limits.carriedAfterMs`. + private static let CARRIED_AFTER: TimeInterval = 30.0 /// Android's `Limits`: queued bytes toward one peer beyond which a body /// is dropped (the core retries it), and only while that peer's oldest /// write has been outstanding for `WRITE_STALL_MS`. The core hands a burst @@ -130,6 +134,11 @@ public class WifiDirectManager: NSObject, TransportManager { /// mid-preamble has not been heard. var heard = false var proved = false + /// Whether a body followed the preamble, and when a dialed stream + /// proved its address: together they say whether the stream carried + /// the peer, which is what starts the redial ladder over (`end(_:)`). + var carried = false + var provedAt: TimeInterval? init(connection: NWConnection, outbound: Bool, dialed: String?) { self.connection = connection @@ -665,10 +674,14 @@ public class WifiDirectManager: NSObject, TransportManager { case .success(let frame): stream.heard = true self.peers.received(frame, from: stream) - if !stream.proved, let address = stream.dialed, - self.peers.provedAddress(of: stream) == address { - stream.proved = true - self.dialPolicy.proved(address) + if !stream.carried, stream.dialed != nil, + self.peers.provedAddress(of: stream) == stream.dialed { + if stream.proved { + stream.carried = true + } else { + stream.proved = true + stream.provedAt = ProcessInfo.processInfo.systemUptime + } } case .failure(let refusal): self.emitDiagnostic("warning", "Peer stream refused", context: ["reason": refusal.reason]) @@ -704,6 +717,19 @@ public class WifiDirectManager: NSObject, TransportManager { unprovable.insert(stream.connection.endpoint) recordAdverts(browser?.browseResults ?? [], fresh: [], local: local) } + // The ladder starts over only after a stream that carried its peer: a + // body after the preamble, or the address held through the keepalive + // window. A proof alone is not enough: a stream the peer refuses for + // its own held one (a stale stream there) proves the address too, and + // resetting on it redialed every second until that stream died. A body + // alone is too narrow: two peers with a session may reconnect and send + // nothing. Android's `Ran.carried` is the same rule. + if let address = stream.dialed, + stream.carried || stream.provedAt.map({ + ProcessInfo.processInfo.systemUptime - $0 >= Self.CARRIED_AFTER + }) == true { + dialPolicy.proved(address) + } guard let address = stream.dialed, let delay = dialPolicy.ended( address, advertised: adverts[address] != nil, diff --git a/bindings/react-native/src/types.ts b/bindings/react-native/src/types.ts index 8e192047..68cf363d 100644 --- a/bindings/react-native/src/types.ts +++ b/bindings/react-native/src/types.ts @@ -410,7 +410,11 @@ export interface InternetTransportConfig { * WiFi Direct transport configuration */ export interface WifiDirectTransportConfig { - /** Enable WiFi Direct transport */ + /** + * Enable the peer-stream slot. Android: Wi-Fi Direct, and DNS-SD on the + * Wi-Fi network it is on, which finds iPhones and Python hosts there. iOS: + * Network framework streams on the LAN and over AWDL. + */ enabled: boolean; /** Device name to advertise */ deviceName?: string; @@ -627,7 +631,7 @@ export interface TransportsConfig { ble?: BleTransportConfig; /** Internet transport configuration */ internet?: InternetTransportConfig; - /** WiFi Direct transport configuration (Android only) */ + /** Peer-stream slot configuration (Wi-Fi Direct and LAN on Android, LAN and AWDL on iOS) */ wifiDirect?: WifiDirectTransportConfig; /** Reticulum mesh transport configuration (requires external Reticulum daemon) */ reticulum?: ReticulumTransportConfig; diff --git a/crates/offline-protocol-uniffi/src/lib.rs b/crates/offline-protocol-uniffi/src/lib.rs index d4677c2d..0fb5e2e3 100644 --- a/crates/offline-protocol-uniffi/src/lib.rs +++ b/crates/offline-protocol-uniffi/src/lib.rs @@ -17149,7 +17149,7 @@ mod tests { ), ( "the question and what an answer does", - "private fun adoptExistingGroup() { \ + "private fun adoptExistingGroup() { if (channel == null) return \ wifiP2pManager?.requestConnectionInfo(channel) { info -> if \ (info?.groupFormed == true && state == TransportState.RUNNING) {", ), @@ -18052,20 +18052,21 @@ mod tests { ); } - /// The iOS and Python peer-stream managers keep the same one of two - /// streams for an address. + /// The iOS, Android and Python peer-stream managers keep the same one of + /// two streams for an address. /// - /// Both ends of a pair may dial, and a Python host on a LAN dials every - /// peer it discovers, so the choice of which stream to keep is a + /// Both ends of a pair may dial, and every manager on a LAN dials the + /// peers it discovers, so the choice of which stream to keep is a /// hand-mirrored policy (docs/bridges C5): the stream the lower address - /// opened wins. If the two copies drift, an iPhone and a Python host each - /// keep the stream the other closes, and the pair reconnects forever, - /// with no error on either side (ADR 0027). Pinned as the two lines each - /// rule is made of. Android keeps the newer stream instead, because a - /// Wi-Fi Direct group has one dialer. + /// opened wins, addresses ordered by their UTF-8 bytes. If two copies + /// drift, a pair each keeps the stream the other closes and reconnects + /// forever, with no error on either side (ADR 0027, ADR 0028). Pinned as + /// the lines each rule is made of. #[test] - fn ios_and_python_peer_streams_keep_the_same_stream() { + fn every_peer_stream_manager_keeps_the_same_stream() { let swift = rn_source_code_only("ios/PeerStreamFraming.swift"); + let kotlin = + rn_source_code_only("android/src/main/java/com/offlineprotocol/PeerStreamFraming.kt"); let python_path = std::path::Path::new(env!("CARGO_MANIFEST_DIR")) .join("../../bindings/python/offline_protocol_sdk/peer_stream_manager.py"); let python = std::fs::read_to_string(&python_path) @@ -18085,6 +18086,22 @@ mod tests { as in peer_stream_manager.py. Expected to find:\n {needed}" ); } + for needed in [ + "val local = localAddress ?: return false", + "val weOpen = utf8Precedes(local, peer) return outbound == weOpen", + "val x = a.encodeToByteArray() val y = b.encodeToByteArray()", + "val d = (x[i].toInt() and 0xff) - (y[i].toInt() and 0xff) \ + if (d != 0) return d < 0 } return x.size < y.size", + "if (byAddress.containsKey(address) && \ + !newStreamWins(outbound, localAddress, address) ) { \ + return Announcement(firstForAddress = false, superseded = null, refused = true) }", + ] { + assert!( + kotlin.contains(needed), + "android PeerStreamFraming.kt: the stream the lower address opened must win, \ + by UTF-8 bytes, as on iOS and in Python. Expected to find:\n {needed}" + ); + } for needed in [ "if local is None:\n # No address of our own to order by; keep what is announced.\n return False", "we_open = local < peer\n return new.outbound == we_open", @@ -18097,6 +18114,103 @@ mod tests { } } + /// A peer-stream redial ladder starts over only after a stream that + /// carried its peer (a body after the preamble, or the address held + /// through the keepalive window), on iOS and Android alike. + /// + /// A stream the peer refuses for its own held one proves the address all + /// the same, so a ladder that reset on proof redialed every second, with a + /// connect and a loss reported each round, until the peer's stale stream + /// died. A body alone is too narrow the other way: two peers with a + /// session may reconnect and send nothing, and each ordinary drop of such + /// a stream climbed the ladder for good. Each platform holds its copy of + /// the rule in its own words, so the lines are pinned here (docs/bridges + /// C5, ADR 0028). + #[test] + fn peer_stream_redial_ladders_reset_only_on_a_carried_stream() { + let swift = rn_source_code_only("ios/WifiDirectManager.swift"); + let lan = + rn_source_code_only("android/src/main/java/com/offlineprotocol/LanPeerDiscovery.kt"); + let manager = + rn_source_code_only("android/src/main/java/com/offlineprotocol/WifiDirectManager.kt"); + let sockets = + rn_source_code_only("android/src/main/java/com/offlineprotocol/PeerStreamSockets.kt"); + for (code, needed) in [ + (&swift, "if stream.proved { stream.carried = true }"), + ( + &swift, + "stream.carried || stream.provedAt.map({ \ + ProcessInfo.processInfo.systemUptime - $0 >= Self.CARRIED_AFTER \ + }) == true { dialPolicy.proved(address) }", + ), + (&sockets, "carried = delivered || ("), + (&lan, "if (ran.proved == address && ran.carried) {"), + (&manager, "carried = ran.carried,"), + ( + &manager, + "if (ran.carried) reconnectDelayMs.set(RECONNECT_INITIAL_DELAY_MS)", + ), + ] { + assert!( + code.contains(needed), + "a redial ladder must start over only on a stream that carried its peer. \ + Expected to find:\n {needed}" + ); + } + assert_eq!( + swift.matches("dialPolicy.proved(").count(), + 1, + "ios/WifiDirectManager.swift: the ladder starts over in one place" + ); + } + + /// Android's LAN carrier publishes and reads the record iOS and Python do. + /// + /// A browser finds a peer only by the service type, names its record by + /// the address's digest so a restarted advert replaces the stale one, and + /// dials with the record's `addr` as the address the preamble must prove. + /// Each is a hand-mirrored literal (docs/bridges C5) whose drift fails + /// silently: an Android phone and an iPhone on one network that never + /// find each other, or two records for one peer (ADR 0028). + #[test] + fn android_lan_peer_streams_publish_the_ios_and_python_record() { + let formation = rn_source_code_only( + "android/src/main/java/com/offlineprotocol/WifiDirectGroupFormation.kt", + ); + let lan = + rn_source_code_only("android/src/main/java/com/offlineprotocol/LanPeerDiscovery.kt"); + let swift = rn_source_code_only("ios/WifiDirectManager.swift"); + for (code, needed) in [ + ( + &formation, + "const val SERVICE_TYPE = \"_offlineprotocol._tcp\"", + ), + (&swift, "static let SERVICE_TYPE = \"_offlineprotocol._tcp\""), + (&formation, "const val KEY_VERSION = \"txtvers\""), + (&formation, "const val KEY_ADDRESS = \"addr\""), + ( + &lan, + "setAttribute(WifiDirectGroupFormation.KEY_VERSION, \"1\") \ + setAttribute(WifiDirectGroupFormation.KEY_ADDRESS, address)", + ), + ( + &lan, + "\"op-\" + WifiDirectGroupFormation.hex(WifiDirectGroupFormation.sha256(address), 8)", + ), + ( + &lan, + "sockets.run(socket, outbound = true, PeerStreamSockets.Carrier.LAN, \ + expected = address)", + ), + ] { + assert!( + code.contains(needed), + "the Android LAN record must match the iOS and Python one. \ + Expected to find:\n {needed}" + ); + } + } + /// Wi-Fi Direct hands the core only an address a preamble proved. /// /// Both managers used to pass a transport-level string, a TCP endpoint on @@ -18223,7 +18337,10 @@ mod tests { ], &[ "is PeerStreamPreamble.Outcome.Announce -> outcome.address", - "if (!announce(stream, address)) {", + "if (!announce(stream, address, outbound)) {", + "PeerStreamDialPolicy.admitsInbound(host, inboundHosts, reserved, limits)", + "val announcement = links.announce(stream, address, outbound, local) \ + if (announcement.refused) {", "if (!deliver(stream, address, body)) {", "links.remove(stream)?.let { address -> try { host.peerDisconnected(address)", "if (announcement.firstForAddress) { try { host.peerConnected(address)", diff --git a/docs/UPGRADING.md b/docs/UPGRADING.md index 1a229240..28b843ff 100644 --- a/docs/UPGRADING.md +++ b/docs/UPGRADING.md @@ -51,7 +51,10 @@ and an exhaustive TypeScript `switch` over `SecurityWarningCode` needs a case for the new member. It also breaks one thing at run time: iOS's peer-stream slot moves from MultipeerConnectivity to Network framework, so an iPhone on `v0.28` does not see one on `v0.27` or earlier over that slot -([§26](#26-behaviour-that-changes-without-a-compile-error-v0280)). +([§26](#26-behaviour-that-changes-without-a-compile-error-v0280)). The +unreleased changes break no build; Android's peer-stream slot joins the Wi-Fi +network it is on +([§27](#27-behaviour-that-changes-without-a-compile-error-unreleased)). Otherwise, where a later section documents an addition or a behaviour change, it is labelled inline with the release that @@ -2500,6 +2503,39 @@ default. `ProtocolManager` already passes it; only code that builds an --- +## 27. Behaviour that changes without a compile error *(unreleased)* + +Everything compiles unchanged. Each paragraph says what to check. + +**Android: the peer-stream slot joins the Wi-Fi network.** With +`wifiDirect: { enabled: true }`, an Android phone now advertises its address +on the Wi-Fi network it is on and opens streams to iPhones, Python hosts and +other Android phones it finds there, as an iPhone already did. Every device on +that network can see the address, the same exposure as a Bluetooth LE +advertisement. There is no switch for the LAN alone: turn `wifiDirect` off if +you do not want either. The module adds `ACCESS_NETWORK_STATE` and +`CHANGE_WIFI_MULTICAST_STATE`, install-time permissions with no prompt; the +second lets Android 12 and lower hold the multicast lock mDNS needs there. +When you move your app to target Android 17 (API 37), it must declare `ACCESS_LOCAL_NETWORK` and request it before `start()`, or the LAN carrier stays off with a warning diagnostic. The SDK does not declare it, because a declaration revokes the grant apps targeting 36 and lower hold by default. + +**Android: the transport starts without the Wi-Fi Direct grant.** Without +`NEARBY_WIFI_DEVICES` (or fine location on 12 and lower), `start()` used to +leave the slot off. It now runs on the Wi-Fi network only, and Wi-Fi Direct +stays off until `enableTransport('wifiDirect')` after the grant. + +**Android: of two streams for one address, the lower-opened one is kept.** +Inside a Wi-Fi Direct group this changes one thing: a client whose address is +higher than its group owner's reconnects past a half-open stream after +keepalive ends it, about thirty seconds on Android 10 and later (up to two +hours on 7 to 9, which keep the OS defaults), instead of at once. Bluetooth LE +and the relays carry traffic meanwhile. + +**Android peers on an older release** neither advertise nor browse on the +LAN, so they meet a phone on this release only inside a Wi-Fi Direct group, +as before. + +--- + ## Appendix A: limits reference | Limit | Value | Where enforced | diff --git a/docs/adr/0027-ios-peer-streams-ride-network-framework.md b/docs/adr/0027-ios-peer-streams-ride-network-framework.md index 0fc6ab32..bd50e2b6 100644 --- a/docs/adr/0027-ios-peer-streams-ride-network-framework.md +++ b/docs/adr/0027-ios-peer-streams-ride-network-framework.md @@ -46,7 +46,7 @@ can dial. address opened, and the newer of two such.** This is the Python manager's rule. Addresses compare by their UTF-8 bytes. A Rust guard pins the Swift and Python copies together - (`ios_and_python_peer_streams_keep_the_same_stream`). + (`every_peer_stream_manager_keeps_the_same_stream`). 3. **Both ends dial, the higher address after five seconds.** The lower address's stream is the one that will be kept, so it goes first, and the common case opens one stream per pair. The higher address still dials, @@ -93,7 +93,9 @@ seconds idle, 5 between probes, 3 probes. - An iPhone and a Python host on one LAN find each other and keep one stream. Android does not change: its Wi-Fi Direct group has one dialer and keeps "newer supersedes", which the chapter allows as local policy. iOS and - Android still do not talk directly. + Android still do not talk directly. (Since + [ADR 0028](0028-android-peer-streams-join-the-lan.md), Android joins the + LAN and keeps the lower-opened stream too.) - The stream budget is 16 open streams, as on Android, instead of seven. Unlike a Wi-Fi Direct group, the listener is open to the whole LAN, so at most 12 are inbound, leaving 4 for dials whatever the listener holds, and diff --git a/docs/adr/0028-android-peer-streams-join-the-lan.md b/docs/adr/0028-android-peer-streams-join-the-lan.md new file mode 100644 index 00000000..05a43b7a --- /dev/null +++ b/docs/adr/0028-android-peer-streams-join-the-lan.md @@ -0,0 +1,73 @@ +# 0028. Android peer streams join the LAN and keep the lower-opened stream + +**Status:** Accepted + +## Context + +The engine's peer-stream slot (`wifi_direct` in the FFI) is filled by three +managers that speak one [stream chapter](../spec/stream-framing.md): Android +over a Wi-Fi P2P group, iOS over Network framework on a LAN or AWDL, and the +Python host on a LAN. Only iOS and Python met. Android discovered peers over +Wi-Fi P2P service discovery alone, which an iPhone cannot hear and a host does +not speak, so an Android phone and an iPhone on the same Wi-Fi network never +formed a stream, and ADR 0027 recorded "iOS and Android still do not talk +directly". + +The Android listener was already open to that network. It binds every +interface on port 8988 for the whole session, so any device on a shared +network could open streams to it, bounded only by the 16-stream cap. + +Android also kept the opposite duplicate-stream rule, "the newer supersedes", +because inside a group only the client dials. On a LAN both ends dial, and +ADR 0027 names what two different rules do there: each end keeps the stream +the other closes, and the pair reconnects forever with no error on either side. + +## Decision + +1. **Android advertises and browses `_offlineprotocol._tcp` on the Wi-Fi + network** through `NsdManager`, with the record iOS and Python publish + (`txtvers=1`, `addr`, no `app`), whenever `wifiDirect.enabled` is set. No + new switch: iOS does the same under the same flag. Before T extensions 7 + (Android 12 and lower, and 13 without that update) it holds a multicast + lock while advertising or browsing: the platform drops mDNS for an app without one there. +2. **Android keeps the stream the lower address opened, and the newer of two + such, on every carrier.** One rule, not one per carrier: a pair in one P2P + group and on one LAN has two dialers on one link table, keyed by address. + Addresses compare by their UTF-8 bytes, not `String.compareTo`, which + compares UTF-16 units and orders differently past the BMP. + `every_peer_stream_manager_keeps_the_same_stream` pins the Kotlin, Swift + and Python copies together. +3. **Keepalive is 15 seconds idle, 5 between probes, 3 probes, with a + 30-second `TCP_USER_TIMEOUT`**, the iOS and Python timers, set with + `Os.setsockoptInt` because `java.net.Socket` cannot set them. The user + timeout is what bounds a stale stream holding unacknowledged data, where + keepalive sends no probe. Android 10 and later only: before API 29 + `ParcelFileDescriptor.fromSocket` hands back the socket's own descriptor + rather than a duplicate, so closing it closed the socket. +4. **Inbound streams are bounded as on iOS**: 16 open, at most 12 inbound, + at most 4 from one remote address, so a listener open to a whole network + cannot take the slots this device dials with. + +## Consequences + +- An Android phone reaches an iPhone or a Python host on the same Wi-Fi + network. With no shared network they still do not meet; that needs Wi-Fi + Aware on both. +- A group client that is the higher address can no longer supersede its own + half-open stream. Its reconnect is refused until keepalive ends the stale + one, about thirty seconds, where "newer supersedes" reconnected at once. + The lower address's reconnect still supersedes at once. On Android 7 to 9 + the OS defaults apply, and the wait can be two hours. +- An Android device with `wifiDirect.enabled` now announces its address to + every device on the Wi-Fi network it joins, as an iPhone already does + (threat model R16). The envelope fields a frame leaves in clear are + readable on that network, as for the other LAN carriers. +- A mixed-version fleet does not loop: an older Android never dials on a LAN, + and inside a group there is still one dialer. + +## What would undo this + +Going back to "newer supersedes" on Android needs every carrier that can +reach an Android device to have a single dialer, and a LAN does not. Taking +the LAN path out again would leave Android and iOS meeting over Bluetooth LE +and the relays only. diff --git a/docs/adr/README.md b/docs/adr/README.md index 1d150eed..3d0d3733 100644 --- a/docs/adr/README.md +++ b/docs/adr/README.md @@ -38,6 +38,7 @@ silently undo it. Decisions that follow from the obvious default do not need one | [0025](0025-native-packages-are-assembled-in-place.md) | Native packages are assembled from the bridge sources in place | Accepted | | [0026](0026-the-backbone-is-a-gateway-property.md) | The backbone is a gateway property | Accepted | | [0027](0027-ios-peer-streams-ride-network-framework.md) | iOS peer streams ride Network framework and keep the lower-opened stream | Accepted | +| [0028](0028-android-peer-streams-join-the-lan.md) | Android peer streams join the LAN and keep the lower-opened stream | Accepted | ## Format diff --git a/docs/android-integration.md b/docs/android-integration.md index b729660f..fad91637 100644 --- a/docs/android-integration.md +++ b/docs/android-integration.md @@ -225,6 +225,9 @@ Add to `AndroidManifest.xml`: + + + diff --git a/docs/bridges/README.md b/docs/bridges/README.md index 558344d7..44e22e84 100644 --- a/docs/bridges/README.md +++ b/docs/bridges/README.md @@ -118,7 +118,7 @@ in the uniffi crate pins the mobile call sites, which no CI job executes. ## C5. Hand-mirrored constants must be pinned in every language Some constants exist in several places no single compiler sees together. -Fourteen sets do today, and they are pinned by **two different** mechanisms, so +Fifteen sets do today, and they are pinned by **two different** mechanisms, so knowing which one you are touching matters. **The relay-answer prefix exemption list** is the canonical example: the core, @@ -286,12 +286,20 @@ the inclusive 1 MiB ceiling and the 96-byte preamble floor of `peer_stream_manager.py`, beside the transport crate's constants. Every language pins them the per-language way, as literals in a suite CI executes, and each suite also replays the chapter's vector file, so a drift on one side -fails that side's tests. The iOS service type `_offlineprotocol._tcp` is the -one piece pinned by a Rust guard instead, +fails that side's tests. The service type `_offlineprotocol._tcp` is pinned +by Rust guards instead: the iOS copy by `react_native_wifi_direct_announces_only_proved_addresses`, because the iOS -manager that holds it is excluded from the SwiftPM harness. A drifted +manager that holds it is excluded from the SwiftPM harness, and the Android +copy, with the record's keys, instance name and dial claim, by +`android_lan_peer_streams_publish_the_ios_and_python_record`. A drifted ceiling reads in the field as the largest messages vanishing; a drifted -service type reads as two iPhones that never find each other. +service type reads as two phones that never find each other. + +**The duplicate-stream rule** is the fifteenth: keep the stream the lower +address opened, by UTF-8 bytes, and the newer of two such, written in +`PeerStreamFraming.swift`, `PeerStreamFraming.kt` and `peer_stream_manager.py` +and pinned by `every_peer_stream_manager_keeps_the_same_stream`. A drift +reads as a pair that reconnects forever with no error on either side. ## C6. Config parsers must not default to literals @@ -625,7 +633,7 @@ found by the first application to update. | Binding | Owes | |---------|------| | Swift | The manual Objective-C bridge kept in step with every `@objc` method; secure storage backed by Keychain; a live-instance check before emitting; the telemetry session boundary inside a background task (C12); a Network-framework peer-stream manager that announces a peer only under the address its preamble proved, one per address (S8) | -| Kotlin | Secure storage backed by Keystore; no blocking work on the main looper; awareness that platform callbacks arrive on binder threads; the telemetry session boundary from an `Application.ActivityLifecycleCallbacks` watcher, never `onHostPause` (C12); a Wi-Fi Direct manager that announces a peer only under the address its preamble proved, one per address (K8) | +| Kotlin | Secure storage backed by Keystore; no blocking work on the main looper; awareness that platform callbacks arrive on binder threads; the telemetry session boundary from an `Application.ActivityLifecycleCallbacks` watcher, never `onHostPause` (C12); a peer-stream manager, over a Wi-Fi Direct group and the Wi-Fi network, that announces a peer only under the address its preamble proved, one per address, by the lower-opened rule (K8) | | Python | Nothing platform-specific; it is the thinnest binding and therefore the best place to smoke-test an ABI change; a re-entrant lock on the generated callback handle map, installed at import, because the collector can free a core object inside a callback lookup and the core's drop then asks for that lock again (P10); the host platform for telemetry from `platform`; a BLE peripheral that serves the address and the core-built identity assertion, and a central that verifies before it announces (P8); a peer-stream manager that announces a host only under the address its preamble proved, and keeps one announced stream per address (P9); a gateway-daemon client that announces a session only once the gateway bound it to this device's address, and settles a frame only on the gateway's verdict, never on the write (P11) | | TypeScript | Config normalization, event typing kept in step with the core, no assumption that a native method exists in an older binary, and no telemetry lifecycle code of its own | diff --git a/docs/bridges/kotlin.md b/docs/bridges/kotlin.md index 16ac8742..0071746f 100644 --- a/docs/bridges/kotlin.md +++ b/docs/bridges/kotlin.md @@ -104,10 +104,12 @@ exemption list, pinned in `RelayAnswerPrefixesTest.kt`. See [C5](README.md#c5-hand-mirrored-constants-must-be-pinned-in-every-language). -## K8. A Wi-Fi Direct socket's verdict is its preamble, not its connect +## K8. A peer stream's verdict is its preamble, not its connect `WifiDirectManager` fills the peer-stream slot with Wi-Fi Direct group -sockets, framed as [the chapter](../spec/stream-framing.md) specifies. A +sockets and, through `LanPeerDiscovery`, with TCP sockets on the Wi-Fi +network the device is on, framed as [the chapter](../spec/stream-framing.md) +specifies. A socket's peer is announced to the core only under the address `verifyIdentityAssertion` derived from its first frame: never on connect, and never under the socket endpoint or `"go:"`, which is what this manager @@ -124,17 +126,45 @@ from asyncio's single thread: a stream's announcement, deliveries and loss report run under that stream's lock, so no body reaches the core after its loss report. -One limit is stated rather than fixed. A client that vanishes without a -FIN stays announced until the stream notices: Java cannot set the keepalive -interval, so `SO_KEEPALIVE` runs on the platform default (commonly two -hours), and the write deadline sees no stall until both kernel buffers are -full. The core's acknowledgements are the real signal that a peer is gone, -as P9 says for the same reason. - -Local policy differs from P9 on purpose. The newer of two streams for one -address supersedes the older, because only the client dials its group owner, -so the two-dialler tie that the lower-address rule settles cannot occur. A -client whose stream ends while the group is up reconnects on a doubling +A peer that vanishes without a FIN stays announced until keepalive notices: +15 seconds idle, 5 between probes, 3 probes, and a 30-second +`TCP_USER_TIMEOUT` for a stream holding unacknowledged data, where keepalive +sends no probe. These are the iOS and Python timers, set through +`Os.setsockoptInt` because `java.net.Socket` cannot set them, and only on +Android 10 and later: before API 29 `ParcelFileDescriptor.fromSocket` takes +the socket's own descriptor, so closing it closed the socket. Older versions +keep the OS defaults (two hours idle). The core's acknowledgements are still the real signal that a peer +is gone, as P9 says for the same reason. + +A redial ladder starts over only after a stream that carried its peer: a +body after the preamble, or the address held for the keepalive window (30 +seconds). A proof alone is not enough, because a stream refused for the one +already held proves its peer too, and a ladder that reset on it redialed every +second while that stream lived. A body alone is too narrow: two peers with a +session may reconnect and send nothing, and every ordinary drop of such a +stream climbed the ladder for good. A group client whose dial is +refused because another stream (a LAN one) holds its owner's address does not +redial at all; it dials the owner again when that stream is lost. + +`NsdManager`'s `DiscoveryListener` reports a record found and lost, never +changed. A peer that restarts and publishes the same instance name on a new +port, with no goodbye in between, is an update the listener never delivers, +where iOS hears a changed record and dials. So the LAN carrier dials back any +peer whose announced stream ends while its record is still advertised, on the +usual policy, and a dial that nothing answers resolves the record again +before the redial, which returns the cache's current port. Without both, an +Android phone left the restarted peer to dial it, and dialed the dead port +itself (seen on devices against an advertise-only Python host). + +Of two streams for one address, the manager keeps the one the lower address +opened, and the newer of two such: P9's rule, compared by UTF-8 bytes, and +`every_peer_stream_manager_keeps_the_same_stream` pins it to the iOS and +Python copies. Inside a group only the client dials, so the rule reduces to +"newer wins" there, except that a client that is the higher address cannot +supersede its own half-open stream: its reconnect is refused until keepalive +ends the stale one, about thirty seconds (two hours on Android 7 to 9). That is the price of the rule being +one rule; a pair that shares a group and a LAN has two dialers on one link +table (ADR 0028). A client whose stream ends while the group is up reconnects on a doubling delay, always to the owner the group has at that moment: on a group switch the new owner's first dial can lose to the old stream still closing, and the redial is then the only one left. The manager joins a group the system formed, including one that diff --git a/docs/bridges/python.md b/docs/bridges/python.md index b1e39c94..8de5a180 100644 --- a/docs/bridges/python.md +++ b/docs/bridges/python.md @@ -181,11 +181,12 @@ kept (two hosts that each list the other open toward each other at once, and without a shared rule each would keep what the other discards, forever), and between two of the winning kind the newer supersedes the older (the lower address reconnecting must not be blocked by its own half-open stream). Either -way the core sees one announcement and one loss per address. The iOS manager -keeps the same stream by the same rule, because an iPhone and a host on one -LAN both dial; `ios_and_python_peer_streams_keep_the_same_stream` pins the -two copies together -([ADR 0027](../adr/0027-ios-peer-streams-ride-network-framework.md)). +way the core sees one announcement and one loss per address. The iOS and +Android managers keep the same stream by the same rule, because a phone and +a host on one LAN both dial; `every_peer_stream_manager_keeps_the_same_stream` +pins the three copies together +([ADR 0027](../adr/0027-ios-peer-streams-ride-network-framework.md), +[ADR 0028](../adr/0028-android-peer-streams-join-the-lan.md)). The tie-break gives the higher address no such way past a stale stream: its reconnect is the losing kind for as long as the lower side holds a stream it diff --git a/docs/configuration.md b/docs/configuration.md index 08a5e1da..bbf6e47d 100644 --- a/docs/configuration.md +++ b/docs/configuration.md @@ -262,7 +262,7 @@ macOS, Linux, and Windows can use the Python binding's snake_case | Parameter | Type | Default | Description | |-----------|------|---------|-------------| | `transports.ble.enabled` | boolean | true | Enable BLE mesh | -| `transports.wifiDirect.enabled` | boolean | false | Enable the [peer-stream](spec/stream-framing.md) slot: Wi-Fi Direct on Android, a TCP stream over a LAN or AWDL on iOS (which needs `_offlineprotocol._tcp` in `NSBonjourServices`, see the [iOS guide](ios-integration.md#platform-limitations)) | +| `transports.wifiDirect.enabled` | boolean | false | Enable the [peer-stream](spec/stream-framing.md) slot: Wi-Fi Direct and a TCP stream over the Wi-Fi network on Android, a TCP stream over a LAN or AWDL on iOS (which needs `_offlineprotocol._tcp` in `NSBonjourServices`, see the [iOS guide](ios-integration.md#platform-limitations)) | | `transports.internet.enabled` | boolean | true | Enable Internet | | `transports.reticulum.enabled` | boolean | false | Enable Reticulum mesh (requires a gateway daemon speaking [contract v1](spec/gateway-contract.md#gateway-daemon-contract-v1)) | | `transports.reticulum.daemonAddress` | string | `localhost:4242` | Host and port of the gateway daemon | @@ -917,6 +917,7 @@ Two consequences worth knowing: - `BLUETOOTH`, `BLUETOOTH_SCAN`, `BLUETOOTH_CONNECT` - `ACCESS_FINE_LOCATION` (for BLE scanning) - `ACCESS_WIFI_STATE`, `NEARBY_WIFI_DEVICES` +- `ACCESS_NETWORK_STATE`, `CHANGE_WIFI_MULTICAST_STATE` (declared by the SDK; the peer-stream slot's LAN carrier) ### iOS diff --git a/docs/ios-integration.md b/docs/ios-integration.md index c99565ff..f1ed6759 100644 --- a/docs/ios-integration.md +++ b/docs/ios-integration.md @@ -442,8 +442,9 @@ local-network privacy blocks discovery; the SDK reports a denial as an ``` The slot reaches other iPhones and, on a shared network, hosts running the -Python binding's `PeerStreamManager`. It does not interoperate with -Android's Wi-Fi Direct, and an SDK from 0.27 or earlier, which used +Python binding's `PeerStreamManager`, and Android phones on the same Wi-Fi +network running an SDK newer than 0.28.0 (never over Android's Wi-Fi Direct, +which iOS cannot join). An SDK from 0.27 or earlier, which used MultipeerConnectivity, does not see it. Available transports: - Bluetooth Low Energy - Internet diff --git a/docs/react-native-integration.md b/docs/react-native-integration.md index 2c0e1b86..3b7ff38e 100644 --- a/docs/react-native-integration.md +++ b/docs/react-native-integration.md @@ -114,7 +114,8 @@ transports: { - Permissions: the SDK declares `ACCESS_WIFI_STATE`, `CHANGE_WIFI_STATE` and `NEARBY_WIFI_DEVICES` (with `neverForLocation`). Request `NEARBY_WIFI_DEVICES` at runtime on Android 13+, and `ACCESS_FINE_LOCATION` on Android 12 and lower, where peer discovery needs it. - The SDK's manifest asserts `neverForLocation` on `NEARBY_WIFI_DEVICES`, and the build merges that flag into your app. An app that derives location from Wi-Fi replaces the declaration (`tools:node="replace"`) and then also needs `ACCESS_FINE_LOCATION` on Android 13+, which the SDK does not check for it. -- `wifiDirect: { enabled: true }` starts the transport in `start()`. Grant `NEARBY_WIFI_DEVICES` (or, on 12 and lower, fine location) **before** `start()`: an enable refused for a missing grant is logged and not retried, so the transport stays off for the session. If you ask for the grant later, call `enableTransport('wifiDirect')` once it is given. +- `wifiDirect: { enabled: true }` starts the transport in `start()`. Grant `NEARBY_WIFI_DEVICES` (or, on 12 and lower, fine location) **before** `start()`: without it the transport runs on the Wi-Fi network only and Wi-Fi Direct stays off for the session. If you ask for the grant later, call `enableTransport('wifiDirect')` once it is given. +- The same slot also advertises and browses `_offlineprotocol._tcp` on the Wi-Fi network the phone is on, so an Android phone reaches iPhones and Python hosts there with no group at all. It needs no runtime grant (the SDK declares `ACCESS_NETWORK_STATE` and `CHANGE_WIFI_MULTICAST_STATE`) until your app targets Android 17: an app targeting Android 17 (API 37) must declare `ACCESS_LOCAL_NETWORK` and request it before `start()`, or the LAN carrier stays off with a warning diagnostic. The SDK does not declare it, because a declaration revokes the grant apps targeting 36 and lower hold by default. It runs whenever `enabled` is true, and finds nothing on a network that blocks multicast or isolates clients (many guest and office networks). - `autoAccept: true` lets the SDK form the Wi-Fi Direct group itself on Android 10 and later: devices of the same app find each other and join one group, with no system dialog. Expect one to two minutes: two phones took 18 to 225 seconds (median about a minute), and Bluetooth, when on, carries traffic meanwhile. Without it, pair the phones once in the system's Wi-Fi Direct settings. `groupOwnerIntent` is not used. - Config: `wifiDirect: { enabled: true, autoAccept: true }`. `enableTransport('wifiDirect')` with no configuration keeps the `autoAccept` setting the transport already has. diff --git a/docs/security/threat-model.md b/docs/security/threat-model.md index 61a2378d..1d5cd7c0 100644 --- a/docs/security/threat-model.md +++ b/docs/security/threat-model.md @@ -139,8 +139,13 @@ Python host all speak plain TCP, which the part of the trust argument. A device on the path is A1 for those frames. It reads what an envelope leaves in clear, such as the application id ([R17](#r17-the-application-id-on-every-frame-is-cleartext-and-unsigned)), -and never a payload. The iOS manager's hop was encrypted while it used -MultipeerConnectivity +and never a payload. On a shared Wi-Fi network the path is every device on +it, for Android since it joined the LAN +([ADR 0028](../adr/0028-android-peer-streams-join-the-lan.md)) as for iOS and +hosts before; a listener open to that network takes at most twelve inbound +streams of its sixteen, and four from one remote address, so strangers' +sockets cannot stop it dialing. The iOS manager's hop was encrypted while it +used MultipeerConnectivity ([ADR 0027](../adr/0027-ios-peer-streams-ride-network-framework.md)). ### Boundary 5: end to end @@ -724,13 +729,12 @@ the stream chapter's one-stream-per-address rule it gains one more thing: its close is reported as the peer's loss, which evicts the real peer's link until it reconnects. That is why the rule is normative rather than policy. Which stream the rule keeps is policy, and each choice leaves the replayer something. -The Android manager keeps the newer stream, so a replayer chooses when a real -stream ends: its copy supersedes the real one, the real peer reconnects past -it, and each round costs both sides a signature check, while the core sees no -loss because the address never stopped being held. The Python and iOS -managers keep the stream opened by the lower address, and the newer of two -such, so the same replay works only against a receiver whose address is -higher than the peer it copies; a lower receiver refuses the copy. The price +Every manager keeps the stream opened by the lower address, and the newer of +two such. Against a receiver whose address is higher than the peer it copies, +a replayer chooses when a real stream ends: its copy supersedes the real one, +the real peer reconnects past it, and each round costs both sides a signature +check, while the core sees no loss because the address never stopped being +held. A lower receiver refuses the copy. The price is that a stale stream can block the higher address's reconnect until keepalive notices. A receiver MUST NOT read a verified assertion as evidence that the peer is live, recent, or the only holder of the diff --git a/docs/spec/stream-framing.md b/docs/spec/stream-framing.md index ba03757c..5bee46dc 100644 --- a/docs/spec/stream-framing.md +++ b/docs/spec/stream-framing.md @@ -3,9 +3,9 @@ ## What this chapter is for A peer stream is a byte stream the platform established to exactly one other -device: a Wi-Fi Direct group socket on Android, a TCP connection over a LAN -or over AWDL (Apple's peer-to-peer Wi-Fi) on iOS, a TCP connection over a LAN -or over a routed mesh on a host. To the protocol +device: a Wi-Fi Direct group socket or a TCP connection over a LAN on +Android, a TCP connection over a LAN or over AWDL (Apple's peer-to-peer Wi-Fi) +on iOS, a TCP connection over a LAN or over a routed mesh on a host. To the protocol engine these are one transport, registered in the slot the FFI names `wifi_direct` for historical reasons, and this chapter is what makes them one: it specifies the only two things a stream does not get from its platform for @@ -222,17 +222,20 @@ is at most one body in flight, `DEFAULT_MAX_MESSAGE_SIZE + 4` bytes, because a frame is read whole before the next prefix. The number of streams a receiver accepts, and the preamble deadline, are local policy, and a conforming implementation chooses its own. The mobile managers use ten seconds and -sixteen open streams; the Python manager's choices are in its bridge rules. -Which of two streams for one address to keep is policy too, and the -implementations differ for a reason. The Python and iOS managers keep the -stream opened by the lower address, and the newer of two such, because both -ends of a pair dial and must agree without talking: a host dials every peer -it lists or discovers, and an iPhone dials every peer it discovers -([ADR 0027](../adr/0027-ios-peer-streams-ride-network-framework.md)). The two -must compute the rule alike, or an iPhone and a host each keep the stream the -other closes and reconnect forever. The Android manager keeps the newer, -because a Wi-Fi Direct group has one dialer, and on a phone the duplicate is -almost always the same peer reconnecting past a half-open stream. +sixteen open streams, of which at most twelve are inbound and four come from +one remote address, because a listener on a shared network is open to every +device on it; the Python manager's choices are in its bridge rules. +Which of two streams for one address to keep is policy too, but every +implementation in this repository keeps the same one: the stream opened by +the lower address (addresses compared by their UTF-8 bytes), and the newer of +two such. Both ends of a pair dial and must agree without talking: a host +dials every peer it lists or discovers, and so do an iPhone and an Android +phone ([ADR 0027](../adr/0027-ios-peer-streams-ride-network-framework.md), +[ADR 0028](../adr/0028-android-peer-streams-join-the-lan.md)). Two ends that +compute the rule differently each keep the stream the other closes and +reconnect forever. The rule costs the higher address an immediate reconnect +past its own half-open stream: the lower end refuses the new one until +keepalive ends the old. ## Finding a peer on a LAN @@ -240,8 +243,11 @@ Nothing in this chapter requires discovery: a stream opened to a configured host and port is complete as specified. Where a LAN offers DNS-SD (RFC 6763), an implementation that advertises MUST use the service type `_offlineprotocol._tcp` (fifteen characters, the maximum a service label -allows) and a TXT record whose first entry is `txtvers=1` and which carries -`addr=`, the advertiser's canonical address. +allows) and a TXT record that carries `txtvers=1`, first where the +advertiser controls the order (RFC 6763 section 6.7), and +`addr=`, the advertiser's canonical address. A browser MUST look +entries up by key and MUST NOT depend on their order: Android's `NsdManager` +keeps a record's attributes in a map and does not promise one. A framework that publishes DNS-SD on the implementation's behalf is bound by the same rule, and MUST NOT publish `_offlineprotocol._tcp` for a service @@ -251,6 +257,12 @@ app's `NSBonjourServices` must list `_offlineprotocol._tcp`, or iOS local-network privacy blocks discovery. The iOS manager once used MultipeerConnectivity, which published the same type in front of its own protocol, so a host on the same LAN found an iPhone it could not speak to. +The Android manager publishes and browses the same record through +`NsdManager` on the Wi-Fi network it is on (scoped to that network from +Android 13; below it NSD browses every interface), with the same instance +name (a digest of the address), and binds its dials to that network, since a socket +left to the default network goes out over cellular on a Wi-Fi network with +no internet ([ADR 0028](../adr/0028-android-peer-streams-join-the-lan.md)). The `addr` entry is a hint the preamble proves. It tells a browser which device it is about to connect to, so the derived address of the preamble can @@ -262,7 +274,7 @@ and carries no meaning; an implementation SHOULD NOT put the address there, since the TXT entry already carries it and one copy is one place to get it wrong. -On Android, Wi-Fi Direct carries the same record over Wi-Fi P2P service +On Android, Wi-Fi Direct also carries the record over Wi-Fi P2P service discovery when an application lets the SDK form the group (`wifiDirect.autoAccept`). The record adds one entry for that use: `app`, a tag of the application id that keeps applications' devices apart. The `addr` entry diff --git a/docs/transport-architecture.md b/docs/transport-architecture.md index 1f6d2d8c..19df2ebe 100644 --- a/docs/transport-architecture.md +++ b/docs/transport-architecture.md @@ -183,7 +183,12 @@ transport. (or by another app) is used. Either way the group is joined when `WIFI_P2P_CONNECTION_CHANGED_ACTION` reports it, or at start when it already exists (the broadcast is not sticky since Android 10), and a client - reconnects to its group owner while the group lasts. + reconnects to its group owner while the group lasts. On the Wi-Fi network + the device is on, `LanPeerDiscovery` advertises and browses the same + `_offlineprotocol._tcp` record as iOS and hosts through `NsdManager` and + dials what it finds, so Android reaches both there. Of two streams for one + address it keeps the one the lower address opened, as they do + ([ADR 0028](adr/0028-android-peer-streams-join-the-lan.md)). - iOS: Network framework TCP streams, on the LAN or over AWDL, with `PeerStreamReader` cutting frames and `PeerStreamSession` for the per-stream rules. It advertises and browses `_offlineprotocol._tcp` with