From 461e215a2debc14901e203f6fab486baf1bbd5be Mon Sep 17 00:00:00 2001 From: James Rich <2199651+jamesarich@users.noreply.github.com> Date: Tue, 29 Sep 2026 23:47:46 +0000 Subject: [PATCH] fix(service): stop the inbound pipeline waiting on node writes (#7464) --- .../core/data/manager/GeofenceMonitor.kt | 2 +- .../data/manager/MeshConfigFlowManagerImpl.kt | 18 +- .../data/manager/MeshConfigHandlerImpl.kt | 29 ++- .../data/manager/MeshMessageProcessorImpl.kt | 44 +++- .../core/data/manager/NodeManagerImpl.kt | 64 +++-- .../manager/MeshConfigFlowManagerImplTest.kt | 14 +- .../data/manager/MeshConfigHandlerImplTest.kt | 10 +- .../manager/MeshMessageProcessorImplTest.kt | 132 ++++++----- .../MeshMessageProcessorNodeWriteTest.kt | 219 ++++++++++++++++++ .../core/data/manager/NodeManagerImplTest.kt | 99 +++++++- .../repository/AirQualityChartReproTest.kt | 4 +- .../core/repository/RadioInterfaceService.kt | 5 +- .../service/SharedRadioInterfaceService.kt | 33 +-- .../core/service/RadioControllerImplTest.kt | 4 +- ...SharedRadioInterfaceServiceLivenessTest.kt | 22 +- .../core/testing/FakeRadioInterfaceService.kt | 7 +- 16 files changed, 554 insertions(+), 152 deletions(-) create mode 100644 core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorNodeWriteTest.kt diff --git a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/GeofenceMonitor.kt b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/GeofenceMonitor.kt index 4cef32cb7a..24b32679bf 100644 --- a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/GeofenceMonitor.kt +++ b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/GeofenceMonitor.kt @@ -96,7 +96,7 @@ class GeofenceMonitor( if (sample.session == null) { evaluate(sample.nodeNum, sample.lat, sample.lon) } else { - radioInterfaceService.runWhileSessionActive(sample.session) { + radioInterfaceService.runWhileSessionActive(sample.session, "geofence evaluation") { evaluate(sample.nodeNum, sample.lat, sample.lon) } } diff --git a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImpl.kt b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImpl.kt index dbf6262c48..949ee93419 100644 --- a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImpl.kt +++ b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImpl.kt @@ -114,8 +114,11 @@ class MeshConfigFlowManagerImpl( private fun runForSession(session: RadioSessionContext, block: () -> Unit): Boolean = radioInterfaceService.runIfSessionActive(session, block) - private suspend fun runWhileForSession(session: RadioSessionContext, block: suspend () -> Unit): Boolean = - radioInterfaceService.runWhileSessionActive(session, block) + private suspend fun runWhileForSession( + session: RadioSessionContext, + label: String, + block: suspend () -> Unit, + ): Boolean = radioInterfaceService.runWhileSessionActive(session, label, block) private fun isActiveSession(session: RadioSessionContext): Boolean = radioInterfaceService.isSessionActive(session) @@ -220,7 +223,8 @@ class MeshConfigFlowManagerImpl( scope.handledLaunch { delay(wantConfigDelay) - val heartbeatSent = runWhileForSession(session) { heartbeatSender.sendHeartbeat("inter-stage") } + val heartbeatSent = + runWhileForSession(session, "inter-stage heartbeat") { heartbeatSender.sendHeartbeat("inter-stage") } if (!heartbeatSent) return@handledLaunch delay(wantConfigDelay) runForSession(session) { @@ -261,7 +265,7 @@ class MeshConfigFlowManagerImpl( private suspend fun finishNodeInfoInstall(state: HandshakeState.ReceivingNodeInfo) { val session = state.session try { - val admitted = runWhileForSession(session) { installAndPublishNodeDatabase(state) } + val admitted = runWhileForSession(session, "NodeDB install") { installAndPublishNodeDatabase(state) } if (!admitted) Logger.d { "Discarding stale post-handshake install and publication" } } catch (e: CancellationException) { throw e @@ -361,7 +365,7 @@ class MeshConfigFlowManagerImpl( // Queue on the serialized session-operation lane before returning to the FIFO frame consumer. Without // UNDISPATCHED, a later config frame can queue its persistence first and then be erased by this reset. scope.handledLaunch(start = CoroutineStart.UNDISPATCHED) { - runWhileForSession(session) { + runWhileForSession(session, "handshake config reset") { if (handshakeGeneration.value != gen) return@runWhileForSession radioConfigRepository.clearChannelSet() if (handshakeGeneration.value != gen) return@runWhileForSession @@ -405,7 +409,7 @@ class MeshConfigFlowManagerImpl( } metadataNodeNum?.let { nodeNum -> scope.handledLaunch(start = CoroutineStart.UNDISPATCHED) { - runWhileForSession(session) { nodeRepository.insertMetadata(nodeNum, metadata) } + runWhileForSession(session, "metadata persist") { nodeRepository.insertMetadata(nodeNum, metadata) } } } if (!admitted) Logger.d { "Discarding metadata from stale transport session" } @@ -454,7 +458,7 @@ class MeshConfigFlowManagerImpl( } if (admitted) { scope.handledLaunch(start = CoroutineStart.UNDISPATCHED) { - runWhileForSession(session) { radioConfigRepository.addFileInfo(info) } + runWhileForSession(session, "fileInfo persist") { radioConfigRepository.addFileInfo(info) } } } if (!admitted) Logger.d { "Discarding FileInfo from stale transport session" } diff --git a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImpl.kt b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImpl.kt index ea29ae5443..57ff4983b2 100644 --- a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImpl.kt +++ b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImpl.kt @@ -65,8 +65,11 @@ class MeshConfigHandlerImpl( private fun runForSession(session: RadioSessionContext, block: () -> Unit): Boolean = radioInterfaceService.runIfSessionActive(session, block) - private suspend fun runWhileForSession(session: RadioSessionContext, block: suspend () -> Unit): Boolean = - radioInterfaceService.runWhileSessionActive(session, block) + private suspend fun runWhileForSession( + session: RadioSessionContext, + label: String, + block: suspend () -> Unit, + ): Boolean = radioInterfaceService.runWhileSessionActive(session, label, block) override fun handleDeviceConfig(config: Config, session: RadioSessionContext): Boolean { val admitted = @@ -76,7 +79,9 @@ class MeshConfigHandlerImpl( connectionManager.value.onHandshakeProgress() } if (admitted) { - launchPersistenceForSession(session) { radioConfigRepository.setLocalConfig(config) } + launchPersistenceForSession(session, "config ${config.summarize()} persist") { + radioConfigRepository.setLocalConfig(config) + } } if (!admitted) Logger.d { "Discarding device config from stale transport session" } return admitted @@ -95,7 +100,7 @@ class MeshConfigHandlerImpl( connectionManager.value.onHandshakeProgress() } if (admitted) { - launchPersistenceForSession(session) { + launchPersistenceForSession(session, "moduleConfig ${config.summarize()} persist") { radioConfigRepository.setLocalModuleConfig(config) statusUpdate?.let { (nodeNum, status) -> try { @@ -127,7 +132,9 @@ class MeshConfigHandlerImpl( } if (admitted) { // We always want to save channel settings we receive from the radio. - launchPersistenceForSession(session) { radioConfigRepository.updateChannelSettings(channel) } + launchPersistenceForSession(session, "channel persist") { + radioConfigRepository.updateChannelSettings(channel) + } } if (!admitted) Logger.d { "Discarding channel config from stale transport session" } return admitted @@ -143,7 +150,9 @@ class MeshConfigHandlerImpl( connectionManager.value.onHandshakeProgress() } if (admitted) { - launchPersistenceForSession(session) { radioConfigRepository.setDeviceUIConfig(config) } + launchPersistenceForSession(session, "deviceuiConfig persist") { + radioConfigRepository.setDeviceUIConfig(config) + } } if (!admitted) Logger.d { "Discarding DeviceUI config from stale transport session" } return admitted @@ -156,7 +165,9 @@ class MeshConfigHandlerImpl( connectionManager.value.onHandshakeProgress() } if (admitted) { - launchPersistenceForSession(session) { radioConfigRepository.setLoraRegionPresetMap(map) } + launchPersistenceForSession(session, "region_presets persist") { + radioConfigRepository.setLoraRegionPresetMap(map) + } } if (!admitted) Logger.d { "Discarding region presets from stale transport session" } return admitted @@ -166,8 +177,8 @@ class MeshConfigHandlerImpl( * Queues handshake persistence on the serialized session-operation lane before the FIFO consumer admits the next * packet. */ - private fun launchPersistenceForSession(session: RadioSessionContext, block: suspend () -> Unit) { - scope.handledLaunch(start = CoroutineStart.UNDISPATCHED) { runWhileForSession(session, block) } + private fun launchPersistenceForSession(session: RadioSessionContext, label: String, block: suspend () -> Unit) { + scope.handledLaunch(start = CoroutineStart.UNDISPATCHED) { runWhileForSession(session, label, block) } } } diff --git a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImpl.kt b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImpl.kt index f09c06700c..8e0a91496e 100644 --- a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImpl.kt +++ b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImpl.kt @@ -131,7 +131,7 @@ class MeshMessageProcessorImpl( private suspend fun processFromRadio(proto: FromRadio, myNodeNum: Int?, session: RadioSessionContext) { val admitted = - radioInterfaceService.runWhileSessionActive(session) { + radioInterfaceService.runWhileSessionActive(session, "FromRadio ${proto.variantLabel()}") { safeCatching { // Audit log every incoming variant without allowing delayed work to cross a session boundary. logVariant(proto, session) @@ -259,7 +259,10 @@ class MeshMessageProcessorImpl( private suspend fun processBufferedPacket(buffered: BufferedMeshPacket, myNodeNum: Int): Boolean { val admitted = - radioInterfaceService.runWhileSessionActive(buffered.session) { + radioInterfaceService.runWhileSessionActive( + buffered.session, + "buffered packet ${buffered.packet.portLabel()}", + ) { safeCatching { processReceivedMeshPacket(buffered.packet, myNodeNum, buffered.session) } .onFailure { Logger.e(it) { "Dropped a buffered early packet after a handler error; replay continued" } @@ -288,20 +291,18 @@ class MeshMessageProcessorImpl( launchSessionBound(session, "mesh-packet emission") { serviceStateWriter.emitMeshPacket(packet) } + // The in-memory update is synchronous; the database write runs on its own session lease so this handler, + // which holds the whole inbound pipeline, never waits on the database writer gate. val from = packet.from if (from == myNodeNum) { - persistNodeUpdate( - myNodeNum, - channel = packet.channel, - operation = "local sender-node packet update", - ) { node -> + nodeManager.updateNodeForSession(myNodeNum, session, channel = packet.channel) { node -> applySenderPacketUpdate(node, packet, decoded).copy(lastHeard = nowSeconds.toInt()) } } else { - persistNodeUpdate(myNodeNum, operation = "local-node packet refresh") { node: Node -> + nodeManager.updateNodeForSession(myNodeNum, session) { node: Node -> node.copy(lastHeard = nowSeconds.toInt()) } - persistNodeUpdate(from, channel = packet.channel, operation = "sender-node packet update") { node -> + nodeManager.updateNodeForSession(from, session, channel = packet.channel) { node -> applySenderPacketUpdate(node, packet, decoded) } } @@ -374,6 +375,31 @@ class MeshMessageProcessorImpl( .onFailure { Logger.e(it) { "Failed $operation; packet processing continued" } } .isSuccess + /** Names the variant for diagnostics from field names and the port only, never payload values. */ + private fun FromRadio.variantLabel(): String = packet?.let { "packet ${it.portLabel()}" } + ?: when { + my_info != null -> "my_info" + node_info != null -> "node_info" + config != null -> "config" + moduleConfig != null -> "moduleConfig" + channel != null -> "channel" + config_complete_id != null -> "config_complete_id" + metadata != null -> "metadata" + deviceuiConfig != null -> "deviceuiConfig" + fileInfo != null -> "fileInfo" + region_presets != null -> "region_presets" + queueStatus != null -> "queueStatus" + log_record != null -> "log_record" + mqttClientProxyMessage != null -> "mqttClientProxyMessage" + xmodemPacket != null -> "xmodemPacket" + lockdown_status != null -> "lockdown_status" + clientNotification != null -> "clientNotification" + rebooted != null -> "rebooted" + else -> "other" + } + + private fun MeshPacket.portLabel(): String = decoded?.portnum?.name ?: "encrypted" + private fun insertMeshLog(log: MeshLog, session: RadioSessionContext): Job = launchSessionBound(session, "mesh-log insert") { meshLogRepository.value.insert(log) } diff --git a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/NodeManagerImpl.kt b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/NodeManagerImpl.kt index 1d3b8a402d..17fdfcc983 100644 --- a/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/NodeManagerImpl.kt +++ b/core/data/src/commonMain/kotlin/org/meshtastic/core/data/manager/NodeManagerImpl.kt @@ -96,6 +96,25 @@ class NodeManagerImpl( private fun persistenceLane(nodeNum: Int): Mutex = nodePersistenceLanes[(nodeNum.toLong() and Int.MAX_VALUE.toLong()).toInt() % nodePersistenceLanes.size] + // A persist not yet in its lane writes the latest node once it gets there, so one per node and session is enough. + // Each holds a session lease that teardown drains, so a slow database must not collect one per packet. + private val pendingPersists = atomic(emptyMap, Any>()) + + private fun schedulePersistence(nodeNum: Int, session: RadioSessionContext?) { + val key = nodeNum to session + val token = Any() + var claimed = false + pendingPersists.update { pending -> + claimed = key !in pending + if (claimed) pending + (key to token) else pending + } + if (!claimed) return + val release = { pendingPersists.update { pending -> if (pending[key] === token) pending - key else pending } } + radioInterfaceService + .launchSessionWork(scope, session) { persistLatestNode(nodeNum, onLaneEntered = release) } + .invokeOnCompletion { release() } + } + /** * Resolves a validated public-key correlation hint from a stored [Node], preferring [Node.publicKey] and falling * back to [User.public_key] only when the primary field is not itself a valid hint. @@ -543,10 +562,9 @@ class NodeManagerImpl( // a key from — its previously captured hint MUST be preserved so same-key replay stays suppressed and a // later genuinely different-keyed device can still claim the slot under the intended contract. committedPresentNums = removedNums.filterTo(mutableSetOf()) { it in state.index.byNum } - committedHints = - removedNums.associateWith { num -> - state.index.byNum[num]?.let(::resolveNodePublicKeyHint) ?: state.retiredKeyHints[num] - } + committedHints = removedNums.associateWith { num -> + state.index.byNum[num]?.let(::resolveNodePublicKeyHint) ?: state.retiredKeyHints[num] + } state.copy( index = removedNums.fold(state.index) { index, nodeNum -> index.remove(nodeNum) }, retiredNodeNums = state.retiredNodeNums.addingAll(removedNums), @@ -616,9 +634,7 @@ class NodeManagerImpl( session: RadioSessionContext? = null, transform: (Node) -> Node, ): NodeStateChange? = updateNodeState(nodeNum, channel, transform).also { change -> - if (change != null && shouldPersist(change.next)) { - radioInterfaceService.launchSessionWork(scope, session) { persistLatestNode(nodeNum) } - } + if (change != null && shouldPersist(change.next)) schedulePersistence(nodeNum, session) } override fun updateNode(nodeNum: Int, channel: Int, transform: (Node) -> Node) { @@ -640,10 +656,12 @@ class NodeManagerImpl( } /** Serializes persistence per node and reads the latest in-memory value inside that lane. */ - private suspend fun persistLatestNode(nodeNum: Int) = persistenceLane(nodeNum).withLock { - val latest = nodeState.value.index.byNum[nodeNum] ?: return@withLock - if (shouldPersist(latest)) nodeRepository.upsert(latest) - } + private suspend fun persistLatestNode(nodeNum: Int, onLaneEntered: () -> Unit = {}) = + persistenceLane(nodeNum).withLock { + onLaneEntered() + val latest = nodeState.value.index.byNum[nodeNum] ?: return@withLock + if (shouldPersist(latest)) nodeRepository.upsert(latest) + } override fun handleReceivedUser( fromNum: Int, @@ -791,14 +809,13 @@ class NodeManagerImpl( var next = node val user = info.user if (user != null && !shouldPreserveExistingUser(node.user, user)) { - var newUser = - user.let { - if (it.is_licensed == true) { - it.newBuilder().also { wb -> wb.public_key = ByteString.EMPTY }.build() - } else { - it - } + var newUser = user.let { + if (it.is_licensed == true) { + it.newBuilder().also { wb -> wb.public_key = ByteString.EMPTY }.build() + } else { + it } + } if (info.via_mqtt && !newUser.long_name.endsWith(" (MQTT)")) { newUser = newUser.newBuilder().also { wb -> wb.long_name = "${newUser.long_name} (MQTT)" }.build() } @@ -936,8 +953,9 @@ class NodeManagerImpl( ) } // Key is already represented elsewhere — stale replay, suppress. - val keyAlreadyRepresented = - resolvedKey.let { key -> before.candidateNumsByPublicKey[key].orEmpty().isNotEmpty() } + val keyAlreadyRepresented = resolvedKey.let { key -> + before.candidateNumsByPublicKey[key].orEmpty().isNotEmpty() + } if (keyAlreadyRepresented) { return ReceivedUserTransition( after = before, @@ -1165,11 +1183,7 @@ class NodeManagerImpl( /** Applies ordinary same-number persistence and notification side effects once after the reducer CAS commits. */ private fun applyReceivedUserEffects(transition: ReceivedUserTransition, session: RadioSessionContext?) { - transition.upsertNode?.let { node -> - if (shouldPersist(node)) { - radioInterfaceService.launchSessionWork(scope, session) { persistLatestNode(node.num) } - } - } + transition.upsertNode?.let { node -> if (shouldPersist(node)) schedulePersistence(node.num, session) } transition.notifyNode?.let { node -> radioInterfaceService.launchSessionWork(scope, session) { // Resolve the display title before validation so the suspending compose-resources call does diff --git a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImplTest.kt b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImplTest.kt index 6f6d4b1636..86970c67ee 100644 --- a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImplTest.kt +++ b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigFlowManagerImplTest.kt @@ -140,12 +140,12 @@ class MeshConfigFlowManagerImplTest { false } } - everySuspend { radioInterfaceService.runWhileSessionActive(any(), any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(any(), any(), any()) } calls { val session = it.args[0] as RadioSessionContext @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) if (activeSessionFlow.value == session) { block() true @@ -303,10 +303,10 @@ class MeshConfigFlowManagerImplTest { @Test fun `config reset holds the session lane until every clear completes`() = testScope.runTest { val sessionLane = Mutex() - everySuspend { radioInterfaceService.runWhileSessionActive(activeSession, any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(activeSession, any(), any()) } calls { @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) sessionLane.withLock { block() } true } @@ -333,7 +333,7 @@ class MeshConfigFlowManagerImplTest { val laterPersistenceStarted = CompletableDeferred() val laterPersistence = launch { - radioInterfaceService.runWhileSessionActive(activeSession) { + radioInterfaceService.runWhileSessionActive(activeSession, "later persistence") { assertEquals( setOf("config", "module", "device-ui", "manifest", "lora-presets"), completedClears, @@ -848,9 +848,7 @@ class MeshConfigFlowManagerImplTest { @Test fun `handleMyInfo applies the event node-event default for event firmware`() = testScope.runTest { - handleMyInfo( - protoMyNodeInfo.newBuilder().also { wb -> wb.firmware_edition = FirmwareEdition.DEFCON }.build(), - ) + handleMyInfo(protoMyNodeInfo.newBuilder().also { wb -> wb.firmware_edition = FirmwareEdition.DEFCON }.build()) advanceUntilIdle() verify { notificationPrefs.applyEventFirmwareNodeEventDefault(isEventFirmware = true) } diff --git a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImplTest.kt b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImplTest.kt index 7830c74b23..fb16a4f849 100644 --- a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImplTest.kt +++ b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshConfigHandlerImplTest.kt @@ -84,10 +84,10 @@ class MeshConfigHandlerImplTest { block() true } - everySuspend { radioInterfaceService.runWhileSessionActive(session, any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } calls { @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) block() true } @@ -197,10 +197,10 @@ class MeshConfigHandlerImplTest { block() true } - everySuspend { radioInterfaceService.runWhileSessionActive(session, any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } calls { @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) block() true } @@ -232,7 +232,7 @@ class MeshConfigHandlerImplTest { block() true } - everySuspend { radioInterfaceService.runWhileSessionActive(session, any()) } returns false + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } returns false handler.handleDeviceConfig(config, session) advanceUntilIdle() diff --git a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImplTest.kt b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImplTest.kt index 791f6c34a0..4396a6f513 100644 --- a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImplTest.kt +++ b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorImplTest.kt @@ -28,15 +28,14 @@ import dev.mokkery.mock import dev.mokkery.verify import dev.mokkery.verify.VerifyMode import dev.mokkery.verifySuspend -import kotlinx.coroutines.CompletableDeferred import kotlinx.coroutines.CoroutineScope import kotlinx.coroutines.ExperimentalCoroutinesApi -import kotlinx.coroutines.async import kotlinx.coroutines.flow.MutableStateFlow import kotlinx.coroutines.test.UnconfinedTestDispatcher import kotlinx.coroutines.test.advanceUntilIdle import kotlinx.coroutines.test.runTest import okio.ByteString +import okio.ByteString.Companion.encodeUtf8 import okio.ByteString.Companion.toByteString import org.meshtastic.core.common.di.asServiceScope import org.meshtastic.core.repository.FromRadioPacketHandler @@ -56,7 +55,6 @@ import org.meshtastic.proto.PortNum import kotlin.test.BeforeTest import kotlin.test.Test import kotlin.test.assertEquals -import kotlin.test.assertFalse import kotlin.test.assertTrue @OptIn(ExperimentalCoroutinesApi::class) @@ -108,12 +106,12 @@ class MeshMessageProcessorImplTest { false } } - everySuspend { radioInterfaceService.runWhileSessionActive(any(), any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(any(), any(), any()) } calls { val requestedSession = it.args[0] as RadioSessionContext @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) if (activeSession.value == requestedSession) { block() true @@ -218,7 +216,7 @@ class MeshMessageProcessorImplTest { FromRadio.Builder() .also { wb -> wb.log_record = LogRecord.Builder().also { wb -> wb.message = "revoked" }.build() } .build() - everySuspend { radioInterfaceService.runWhileSessionActive(session, any()) } returns false + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } returns false processor.handleFromRadio(frame(fromRadio.encode()), myNodeNum) advanceUntilIdle() @@ -229,60 +227,84 @@ class MeshMessageProcessorImplTest { } @Test - fun `packet processing awaits node persistence before releasing session authority`() = runTest(testDispatcher) { - isNodeDbReady.value = true - every { nodeManager.myNodeNum } returns MutableStateFlow(null) - var authorityDepth = 0 - everySuspend { radioInterfaceService.runWhileSessionActive(session, any()) } calls - { - @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) - authorityDepth += 1 - try { - block() - true - } finally { - authorityDepth -= 1 + fun `packet node updates are handed to session-bound persistence from inside the handler`() = + runTest(testDispatcher) { + isNodeDbReady.value = true + every { nodeManager.myNodeNum } returns MutableStateFlow(null) + var authorityDepth = 0 + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } calls + { + @Suppress("UNCHECKED_CAST") + val block = it.args[2] as (suspend () -> Unit) + authorityDepth += 1 + try { + block() + true + } finally { + authorityDepth -= 1 + } } - } - processor = createProcessor(backgroundScope) - val persistenceStarted = CompletableDeferred() - val releasePersistence = CompletableDeferred() - everySuspend { nodeManager.updateNodeAndPersist(any(), any(), any()) } calls - { - assertTrue(authorityDepth > 0, "node persistence must begin while session authority is held") - persistenceStarted.complete(Unit) - releasePersistence.await() - } - val packet = - MeshPacket.Builder() - .also { wb -> - wb.id = 9 - wb.from = 999 - wb.decoded = - Data.Builder() - .also { wb -> - wb.portnum = PortNum.TEXT_MESSAGE_APP - wb.payload = ByteString.EMPTY - } - .build() - wb.rx_time = 1000 + processor = createProcessor(backgroundScope) + val updatedNodes = mutableListOf() + every { nodeManager.updateNodeForSession(any(), session, any(), any()) } calls + { + assertTrue(authorityDepth > 0, "the nested lease must be taken while the handler holds authority") + updatedNodes += it.arg(0) } - .build() + val packet = textPacket(id = 9, from = 999) - val processing = async { processor.handleFromRadio( frame(FromRadio.Builder().also { wb -> wb.packet = packet }.build().encode()), myNodeNum, ) - } - persistenceStarted.await() + advanceUntilIdle() - assertFalse(processing.isCompleted, "session authority must remain held until node persistence finishes") - releasePersistence.complete(Unit) - processing.await() + assertEquals(listOf(myNodeNum, 999), updatedNodes) + verifySuspend(mode = VerifyMode.exactly(0)) { nodeManager.updateNodeAndPersist(any(), any(), any()) } + } + + @Test + fun `handleFromRadio names the handler by variant and port without payload`() = runTest(testDispatcher) { + isNodeDbReady.value = true + val labels = mutableListOf() + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } calls + { + labels += it.arg(1) + true + } + processor = createProcessor(backgroundScope) + val packet = textPacket(id = 11, from = 999, payload = "secret text".encodeUtf8()) + val logRecord = LogRecord.Builder().also { wb -> wb.message = "secret log" }.build() + + processor.handleFromRadio( + frame(FromRadio.Builder().also { wb -> wb.packet = packet }.build().encode()), + myNodeNum, + ) + processor.handleFromRadio( + frame(FromRadio.Builder().also { wb -> wb.log_record = logRecord }.build().encode()), + myNodeNum, + ) + advanceUntilIdle() + + assertEquals(listOf("FromRadio packet TEXT_MESSAGE_APP", "FromRadio log_record"), labels) } + private fun textPacket(id: Int, from: Int, payload: ByteString = ByteString.EMPTY): MeshPacket = + MeshPacket.Builder() + .also { wb -> + wb.id = id + wb.from = from + wb.decoded = + Data.Builder() + .also { d -> + d.portnum = PortNum.TEXT_MESSAGE_APP + d.payload = payload + } + .build() + wb.rx_time = 1000 + } + .build() + @Test fun `local sender packet persists one combined node update`() = runTest(testDispatcher) { isNodeDbReady.value = true @@ -310,7 +332,7 @@ class MeshMessageProcessorImplTest { ) advanceUntilIdle() - verifySuspend(mode = VerifyMode.exactly(1)) { nodeManager.updateNodeAndPersist(myNodeNum, any(), any()) } + verify(mode = VerifyMode.exactly(1)) { nodeManager.updateNodeForSession(myNodeNum, session, any(), any()) } } @Test @@ -636,8 +658,8 @@ class MeshMessageProcessorImplTest { advanceUntilIdle() // must not throw -- this is the crash repro point // Both packets were attempted -- the poisoned packet's throw did not stop the rest of the batch. - verifySuspend { nodeManager.updateNodeAndPersist(999, any(), any()) } - verifySuspend { nodeManager.updateNodeAndPersist(998, any(), any()) } + verify { nodeManager.updateNodeForSession(999, session, any(), any()) } + verify { nodeManager.updateNodeForSession(998, session, any(), any()) } } // ---------- handleReceivedMeshPacket: rx_time normalization ---------- @@ -752,7 +774,7 @@ class MeshMessageProcessorImplTest { advanceUntilIdle() // Should have called updateNode for myNodeNum (lastHeard update) - verifySuspend { nodeManager.updateNodeAndPersist(myNodeNum, any(), any()) } + verify { nodeManager.updateNodeForSession(myNodeNum, session, any(), any()) } } @Test @@ -782,7 +804,7 @@ class MeshMessageProcessorImplTest { advanceUntilIdle() // Should have called updateNode for the sender - verifySuspend { nodeManager.updateNodeAndPersist(senderNode, any(), any()) } + verify { nodeManager.updateNodeForSession(senderNode, session, any(), any()) } } // ---------- handleReceivedMeshPacket: null decoded ---------- diff --git a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorNodeWriteTest.kt b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorNodeWriteTest.kt new file mode 100644 index 0000000000..709ff0a026 --- /dev/null +++ b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/MeshMessageProcessorNodeWriteTest.kt @@ -0,0 +1,219 @@ +/* + * Copyright (c) 2026 Meshtastic LLC + * + * This program is free software: you can redistribute it and/or modify + * it under the terms of the GNU General Public License as published by + * the Free Software Foundation, either version 3 of the License, or + * (at your option) any later version. + * + * This program is distributed in the hope that it will be useful, + * but WITHOUT ANY WARRANTY; without even the implied warranty of + * MERCHANTABILITY or FITNESS FOR A PARTICULAR PURPOSE. See the + * GNU General Public License for more details. + * + * You should have received a copy of the GNU General Public License + * along with this program. If not, see . + */ +package org.meshtastic.core.data.manager + +import dev.mokkery.MockMode +import dev.mokkery.answering.calls +import dev.mokkery.every +import dev.mokkery.matcher.any +import dev.mokkery.mock +import kotlinx.coroutines.CompletableDeferred +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.ExperimentalCoroutinesApi +import kotlinx.coroutines.async +import kotlinx.coroutines.launch +import kotlinx.coroutines.test.StandardTestDispatcher +import kotlinx.coroutines.test.TestScope +import kotlinx.coroutines.test.runCurrent +import kotlinx.coroutines.test.runTest +import okio.ByteString +import okio.ByteString.Companion.toByteString +import org.meshtastic.core.common.di.asServiceScope +import org.meshtastic.core.common.util.nowSeconds +import org.meshtastic.core.model.Node +import org.meshtastic.core.repository.FromRadioPacketHandler +import org.meshtastic.core.repository.MeshDataHandler +import org.meshtastic.core.repository.MeshLogRepository +import org.meshtastic.core.repository.MeshNotificationManager +import org.meshtastic.core.repository.NodeRepository +import org.meshtastic.core.repository.RadioSessionContext +import org.meshtastic.core.repository.ReceivedRadioFrame +import org.meshtastic.core.repository.ServiceRepository +import org.meshtastic.core.testing.FakeNodeRepository +import org.meshtastic.core.testing.FakeRadioInterfaceService +import org.meshtastic.proto.Data +import org.meshtastic.proto.FromRadio +import org.meshtastic.proto.MeshPacket +import org.meshtastic.proto.PortNum +import kotlin.test.Test +import kotlin.test.assertEquals +import kotlin.test.assertFalse +import kotlin.test.assertNotNull +import kotlin.test.assertNull +import kotlin.test.assertTrue + +/** + * Drives a real [MeshMessageProcessorImpl] and [NodeManagerImpl] over [FakeRadioInterfaceService], whose + * `runWhileSessionActive` serializes handlers on one mutex as production does, with a node repository whose writes wait + * on a gate that stands in for a stalled database writer. + */ +@OptIn(ExperimentalCoroutinesApi::class) +class MeshMessageProcessorNodeWriteTest { + + private val testScope = TestScope(StandardTestDispatcher()) + private val nodeRepository = GatedNodeRepository(FakeNodeRepository()) + private val dataHandler = mock(MockMode.autofill) + + private class GatedNodeRepository(private val delegate: FakeNodeRepository) : NodeRepository by delegate { + val writeStarted = CompletableDeferred() + val releaseWrites = CompletableDeferred() + + val upserts = mutableMapOf() + + fun persisted(num: Int): Node? = delegate.nodeDBbyNum.value[num] + + override suspend fun upsert(node: Node) { + upserts[node.num] = (upserts[node.num] ?: 0) + 1 + writeStarted.complete(Unit) + releaseWrites.await() + delegate.upsert(node) + } + } + + private class Fixture(scope: CoroutineScope, nodeRepository: NodeRepository, dataHandler: MeshDataHandler) { + val radio = FakeRadioInterfaceService(scope) + val nodeManager = + NodeManagerImpl( + nodeRepository, + mock(MockMode.autofill), + radio, + scope.asServiceScope(), + ) + .apply { + notificationTitleFormatter = { shortName -> "New node seen: $shortName" } + setNodeDbReady(true) + setAllowNodeDbWrites(true) + } + val processor = + MeshMessageProcessorImpl( + nodeManager = nodeManager, + serviceStateWriter = mock(MockMode.autofill), + meshLogRepository = lazy { mock(MockMode.autofill) }, + dataHandler = lazy { dataHandler }, + fromRadioDispatcher = mock(MockMode.autofill), + radioInterfaceService = radio, + scope = scope.asServiceScope(), + ) + + val session: RadioSessionContext + + init { + radio.setDeviceAddress("tcp:test") + radio.connect() + session = checkNotNull(radio.activeSession.value) + } + + suspend fun receive(packet: MeshPacket) { + val bytes = FromRadio.Builder().also { wb -> wb.packet = packet }.build().encode() + processor.handleFromRadio(ReceivedRadioFrame(bytes.toByteString(), session), MY_NODE) + } + } + + private fun packetFrom(from: Int, id: Int) = MeshPacket.Builder() + .also { wb -> + wb.id = id + wb.from = from + wb.hop_start = 3 + wb.hop_limit = 1 + wb.rx_time = nowSeconds.toInt() + wb.decoded = + Data.Builder() + .also { d -> + d.portnum = PortNum.TEXT_MESSAGE_APP + d.payload = ByteString.EMPTY + } + .build() + } + .build() + + @Test + fun `the inbound handler returns while its node write waits on the database`() = testScope.runTest { + val fixture = Fixture(backgroundScope, nodeRepository, dataHandler) + + val first = async { fixture.receive(packetFrom(SENDER, id = 1)) } + runCurrent() + assertTrue(nodeRepository.writeStarted.isCompleted, "the node write has started and is held") + assertTrue(first.isCompleted, "the handler must not wait for the node write") + + val second = async { fixture.receive(packetFrom(OTHER_SENDER, id = 2)) } + runCurrent() + assertTrue(second.isCompleted, "the next frame must not queue behind the held write") + assertNull(nodeRepository.persisted(SENDER), "nothing is written while the database is held") + + // The writes run in backgroundScope, which advanceUntilIdle does not drive. + nodeRepository.releaseWrites.complete(Unit) + runCurrent() + + assertEquals(fixture.nodeManager.nodeDBbyNodeNum[SENDER], nodeRepository.persisted(SENDER)) + assertEquals(2, nodeRepository.persisted(SENDER)?.hopsAway) + assertNotNull(nodeRepository.persisted(OTHER_SENDER)) + } + + @Test + fun `packets that arrive while the database is held queue at most two writes per node`() = testScope.runTest { + val fixture = Fixture(backgroundScope, nodeRepository, dataHandler) + + repeat(5) { i -> launch { fixture.receive(packetFrom(SENDER, id = i + 1)) } } + runCurrent() + nodeRepository.releaseWrites.complete(Unit) + runCurrent() + + assertTrue((nodeRepository.upserts[SENDER] ?: 0) <= 2, "sender writes: ${nodeRepository.upserts}") + assertTrue((nodeRepository.upserts[MY_NODE] ?: 0) <= 2, "local node writes: ${nodeRepository.upserts}") + assertEquals(fixture.nodeManager.nodeDBbyNodeNum[SENDER], nodeRepository.persisted(SENDER)) + } + + @Test + fun `the data handler reads the node update from memory before the write lands`() = testScope.runTest { + val fixture = Fixture(backgroundScope, nodeRepository, dataHandler) + var hopsAwaySeen: Int? = null + every { dataHandler.handleReceivedData(any(), any(), any(), any(), any()) } calls + { + hopsAwaySeen = fixture.nodeManager.nodeDBbyNodeNum[SENDER]?.hopsAway + } + + launch { fixture.receive(packetFrom(SENDER, id = 1)) } + runCurrent() + + assertEquals(2, hopsAwaySeen) + assertNull(nodeRepository.persisted(SENDER), "the write is still held") + nodeRepository.releaseWrites.complete(Unit) + } + + @Test + fun `teardown still waits for a node write the handler left running`() = testScope.runTest { + val fixture = Fixture(backgroundScope, nodeRepository, dataHandler) + launch { fixture.receive(packetFrom(SENDER, id = 1)) } + runCurrent() + + val teardown = launch { fixture.radio.disconnect() } + runCurrent() + assertFalse(teardown.isCompleted, "teardown must drain the write's session lease") + + nodeRepository.releaseWrites.complete(Unit) + runCurrent() + + assertTrue(teardown.isCompleted) + assertNotNull(nodeRepository.persisted(SENDER)) + } + + private companion object { + const val MY_NODE = 12345 + const val SENDER = 999 + const val OTHER_SENDER = 998 + } +} diff --git a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/NodeManagerImplTest.kt b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/NodeManagerImplTest.kt index 739d01623f..6075b7bbee 100644 --- a/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/NodeManagerImplTest.kt +++ b/core/data/src/commonTest/kotlin/org/meshtastic/core/data/manager/NodeManagerImplTest.kt @@ -47,6 +47,7 @@ import org.meshtastic.core.repository.MeshNotificationManager import org.meshtastic.core.repository.NodeRepository import org.meshtastic.core.repository.RadioInterfaceService import org.meshtastic.core.repository.RadioSessionContext +import org.meshtastic.core.repository.RadioSessionLease import org.meshtastic.proto.DeviceMetadata import org.meshtastic.proto.DeviceMetrics import org.meshtastic.proto.EnvironmentMetrics @@ -273,6 +274,98 @@ class NodeManagerImplTest { assertEquals(2, nodeManager.nodeDBbyNodeNum[nodeNum]?.lastHeard) } + private fun admitLeases(session: RadioSessionContext) { + everySuspend { radioInterfaceService.runWithSessionLease(session, any()) } calls + { + @Suppress("UNCHECKED_CAST") + val block = it.args[1] as (suspend (RadioSessionLease) -> Unit) + block( + object : RadioSessionLease { + override val session: RadioSessionContext = session + + override fun isCurrent(): Boolean = true + }, + ) + true + } + } + + @Test + fun `session-bound persistence writes the newest state and coalesces queued updates`() = testScope.runTest { + val nodeNum = 1234 + val session = RadioSessionContext(generation = 7L, address = "ble:same") + nodeManager.setNodeDbReady(true) + nodeManager.setAllowNodeDbWrites(true) + admitLeases(session) + val firstWriteStarted = CompletableDeferred() + val releaseFirstWrite = CompletableDeferred() + val persisted = mutableListOf() + everySuspend { nodeRepository.upsert(any()) } calls + { + if (!firstWriteStarted.isCompleted) { + firstWriteStarted.complete(Unit) + releaseFirstWrite.await() + } + persisted += it.arg(0) + } + + nodeManager.updateNodeForSession(nodeNum, session) { node -> node.copy(lastHeard = 1) } + assertTrue(firstWriteStarted.isCompleted, "the first write is in flight") + nodeManager.updateNodeForSession(nodeNum, session) { node -> node.copy(lastHeard = 2) } + nodeManager.updateNodeForSession(nodeNum, session) { node -> node.copy(lastHeard = 3) } + runCurrent() + assertTrue(persisted.isEmpty(), "later writes must queue behind the first on the node's lane") + + releaseFirstWrite.complete(Unit) + advanceUntilIdle() + + assertEquals(listOf(1, 3), persisted.map(Node::lastHeard)) + assertEquals(3, nodeManager.nodeDBbyNodeNum[nodeNum]?.lastHeard) + } + + @Test + fun `an update from a newer session schedules its own write`() = testScope.runTest { + val nodeNum = 1234 + val oldSession = RadioSessionContext(generation = 7L, address = "ble:same") + val newSession = RadioSessionContext(generation = 8L, address = "ble:same") + nodeManager.setNodeDbReady(true) + nodeManager.setAllowNodeDbWrites(true) + admitLeases(oldSession) + admitLeases(newSession) + val releaseFirstWrite = CompletableDeferred() + var writes = 0 + everySuspend { nodeRepository.upsert(any()) } calls { if (++writes == 1) releaseFirstWrite.await() } + + nodeManager.updateNodeForSession(nodeNum, oldSession) { node -> node.copy(lastHeard = 1) } + nodeManager.updateNodeForSession(nodeNum, oldSession) { node -> node.copy(lastHeard = 2) } + nodeManager.updateNodeForSession(nodeNum, newSession) { node -> node.copy(lastHeard = 3) } + runCurrent() + + verifySuspend(exactly(1)) { radioInterfaceService.runWithSessionLease(newSession, any()) } + releaseFirstWrite.complete(Unit) + advanceUntilIdle() + } + + @Test + fun `a persist rejected by a retired session does not block a later write`() = testScope.runTest { + val nodeNum = 1234 + val oldSession = RadioSessionContext(generation = 7L, address = "ble:same") + val newSession = RadioSessionContext(generation = 8L, address = "ble:same") + nodeManager.setNodeDbReady(true) + nodeManager.setAllowNodeDbWrites(true) + everySuspend { radioInterfaceService.runWithSessionLease(oldSession, any()) } returns false + admitLeases(newSession) + val persisted = mutableListOf() + everySuspend { nodeRepository.upsert(any()) } calls { persisted += it.arg(0) } + + nodeManager.updateNodeForSession(nodeNum, oldSession) { node -> node.copy(lastHeard = 1) } + runCurrent() + nodeManager.updateNodeForSession(nodeNum, newSession) { node -> node.copy(lastHeard = 2) } + runCurrent() + + assertEquals(listOf(2), persisted.map(Node::lastHeard)) + } + @Test fun `session-bound node persistence from a retired generation is rejected`() = testScope.runTest { val nodeNum = 1234 @@ -2868,8 +2961,7 @@ class NodeManagerImplTest { nodeManager.applyTrustedIdentityMigrations(listOf(oldNum)) advanceUntilIdle() val replayDispatches = mutableListOf() - everySuspend { serviceNotifications.showNewNodeSeenNotification(capture(replayDispatches), any()) } returns - Unit + everySuspend { serviceNotifications.showNewNodeSeenNotification(capture(replayDispatches), any()) } returns Unit // Now replay: the canonical number (crc32(key)) has NOT appeared yet nodeManager.handleReceivedUser(oldNum, userWithKey(key, "Replay", "RP"), manuallyVerified = false) advanceUntilIdle() @@ -2971,8 +3063,7 @@ class NodeManagerImplTest { advanceUntilIdle() // Early replay of old number User packet — should be suppressed val dispatchedBefore = mutableListOf() - everySuspend { serviceNotifications.showNewNodeSeenNotification(capture(dispatchedBefore), any()) } returns - Unit + everySuspend { serviceNotifications.showNewNodeSeenNotification(capture(dispatchedBefore), any()) } returns Unit nodeManager.handleReceivedUser(oldNum, userWithKey(key, "Replay", "RP"), manuallyVerified = false) advanceUntilIdle() // The old number should NOT be in nodeDB diff --git a/core/data/src/jvmTest/kotlin/org/meshtastic/core/data/repository/AirQualityChartReproTest.kt b/core/data/src/jvmTest/kotlin/org/meshtastic/core/data/repository/AirQualityChartReproTest.kt index 3f314eb2f8..bcb3edc090 100644 --- a/core/data/src/jvmTest/kotlin/org/meshtastic/core/data/repository/AirQualityChartReproTest.kt +++ b/core/data/src/jvmTest/kotlin/org/meshtastic/core/data/repository/AirQualityChartReproTest.kt @@ -220,10 +220,10 @@ class AirQualityChartReproTest { every { nodeManager.isNodeDbReady } returns MutableStateFlow(true) every { nodeManager.myNodeNum } returns myNodeNumFlow every { radioInterfaceService.isSessionActive(session) } returns true - everySuspend { radioInterfaceService.runWhileSessionActive(session, any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } calls { @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) block() true } diff --git a/core/repository/src/commonMain/kotlin/org/meshtastic/core/repository/RadioInterfaceService.kt b/core/repository/src/commonMain/kotlin/org/meshtastic/core/repository/RadioInterfaceService.kt index 151a84fe90..010fbb9796 100644 --- a/core/repository/src/commonMain/kotlin/org/meshtastic/core/repository/RadioInterfaceService.kt +++ b/core/repository/src/commonMain/kotlin/org/meshtastic/core/repository/RadioInterfaceService.kt @@ -67,9 +67,10 @@ interface RadioSessionAuthority { /** * Runs [block] while holding the same lifecycle lease, without exposing the lease token. Implementations may * serialize this convenience path to preserve handshake ordering; independently deferred work should use - * [runWithSessionLease] so it can acquire its own lease before its parent operation returns. + * [runWithSessionLease] so it can acquire its own lease before its parent operation returns. [label] names the + * operation in diagnostics and must not carry packet contents or identifiers. */ - suspend fun runWhileSessionActive(session: RadioSessionContext, block: suspend () -> Unit): Boolean = + suspend fun runWhileSessionActive(session: RadioSessionContext, label: String, block: suspend () -> Unit): Boolean = runWithSessionLease(session) { block() } } diff --git a/core/service/src/commonMain/kotlin/org/meshtastic/core/service/SharedRadioInterfaceService.kt b/core/service/src/commonMain/kotlin/org/meshtastic/core/service/SharedRadioInterfaceService.kt index 62a4bf4988..6f9261c24f 100644 --- a/core/service/src/commonMain/kotlin/org/meshtastic/core/service/SharedRadioInterfaceService.kt +++ b/core/service/src/commonMain/kotlin/org/meshtastic/core/service/SharedRadioInterfaceService.kt @@ -274,24 +274,27 @@ class SharedRadioInterfaceService( } } - override suspend fun runWhileSessionActive(session: RadioSessionContext, block: suspend () -> Unit): Boolean = - sessionOperationMutex.withLock { - runWithSessionLease(session) { - // Bound the handler: it holds sessionOperationMutex (the whole inbound pipeline) and an admitted - // lease (which teardown's drain awaits), so an indefinite suspension here is a total wedge, not a - // slow packet. Cancelling the block releases both. Only OUR timeout is swallowed — ensureActive() - // rethrows if the surrounding scope was cancelled concurrently. - try { - withTimeout(SESSION_HANDLER_TIMEOUT_MILLIS) { block() } - } catch (timeout: TimeoutCancellationException) { - currentCoroutineContext().ensureActive() - Logger.e(timeout) { - "Session handler exceeded ${SESSION_HANDLER_TIMEOUT_MILLIS}ms and was cancelled; " + - "dropping its packet to keep the receive pipeline alive" - } + override suspend fun runWhileSessionActive( + session: RadioSessionContext, + label: String, + block: suspend () -> Unit, + ): Boolean = sessionOperationMutex.withLock { + runWithSessionLease(session) { + // Bound the handler: it holds sessionOperationMutex (the whole inbound pipeline) and an admitted + // lease (which teardown's drain awaits), so an indefinite suspension here is a total wedge, not a + // slow packet. Cancelling the block releases both. Only OUR timeout is swallowed — ensureActive() + // rethrows if the surrounding scope was cancelled concurrently. + try { + withTimeout(SESSION_HANDLER_TIMEOUT_MILLIS) { block() } + } catch (timeout: TimeoutCancellationException) { + currentCoroutineContext().ensureActive() + Logger.e(timeout) { + "Session handler exceeded ${SESSION_HANDLER_TIMEOUT_MILLIS}ms and was cancelled; " + + "dropping its packet to keep the receive pipeline alive (handler=$label)" } } } + } private fun releaseSessionOperation(admittedSession: RadioTransportSession) { val drainWaiter = diff --git a/core/service/src/commonTest/kotlin/org/meshtastic/core/service/RadioControllerImplTest.kt b/core/service/src/commonTest/kotlin/org/meshtastic/core/service/RadioControllerImplTest.kt index e30b52438a..521cf6e195 100644 --- a/core/service/src/commonTest/kotlin/org/meshtastic/core/service/RadioControllerImplTest.kt +++ b/core/service/src/commonTest/kotlin/org/meshtastic/core/service/RadioControllerImplTest.kt @@ -127,12 +127,12 @@ class RadioControllerImplTest { activeSession ?: MutableStateFlow(deviceAddress.value?.let { RadioSessionContext(sessionGeneration.value, it) }) every { radioInterfaceService.activeSession } returns resolvedActiveSession - everySuspend { radioInterfaceService.runWhileSessionActive(any(), any()) } calls + everySuspend { radioInterfaceService.runWhileSessionActive(any(), any(), any()) } calls { val session = it.args[0] as RadioSessionContext @Suppress("UNCHECKED_CAST") - val block = it.args[1] as (suspend () -> Unit) + val block = it.args[2] as (suspend () -> Unit) if (resolvedActiveSession.value == session) { block() true diff --git a/core/service/src/commonTest/kotlin/org/meshtastic/core/service/SharedRadioInterfaceServiceLivenessTest.kt b/core/service/src/commonTest/kotlin/org/meshtastic/core/service/SharedRadioInterfaceServiceLivenessTest.kt index c8c81a870e..fb149860d3 100644 --- a/core/service/src/commonTest/kotlin/org/meshtastic/core/service/SharedRadioInterfaceServiceLivenessTest.kt +++ b/core/service/src/commonTest/kotlin/org/meshtastic/core/service/SharedRadioInterfaceServiceLivenessTest.kt @@ -20,6 +20,7 @@ import androidx.lifecycle.Lifecycle import androidx.lifecycle.LifecycleEventObserver import androidx.lifecycle.LifecycleObserver import androidx.lifecycle.LifecycleOwner +import co.touchlab.kermit.Severity import dev.mokkery.MockMode import dev.mokkery.answering.calls import dev.mokkery.answering.returns @@ -57,6 +58,7 @@ import org.meshtastic.core.repository.RadioPrefs import org.meshtastic.core.repository.RadioTransport import org.meshtastic.core.repository.RadioTransportFactory import org.meshtastic.core.repository.TransportDisconnectReason +import org.meshtastic.core.testing.CapturingLogWriter import org.meshtastic.core.testing.FakeBluetoothRepository import org.meshtastic.core.testing.FakeRadioPrefs import org.meshtastic.core.testing.FakeRadioTransport @@ -555,13 +557,13 @@ class SharedRadioInterfaceServiceLivenessTest { val independentStarted = CompletableDeferred() val first = launch { - service.runWhileSessionActive(session) { + service.runWhileSessionActive(session, "first") { firstStarted.complete(Unit) releaseFirst.await() } } firstStarted.await() - val second = launch { service.runWhileSessionActive(session) { secondStarted.complete(Unit) } } + val second = launch { service.runWhileSessionActive(session, "second") { secondStarted.complete(Unit) } } val independent = launch { service.runWithSessionLease(session) { independentStarted.complete(Unit) } } try { testDispatcher.scheduler.runCurrent() @@ -619,7 +621,7 @@ class SharedRadioInterfaceServiceLivenessTest { ) assertFalse(service.isSessionActive(session), "teardown must reject new work immediately") assertFalse( - service.runWhileSessionActive(session) { error("late operation must not run") }, + service.runWhileSessionActive(session, "late") { error("late operation must not run") }, "work queued after admission closes must be rejected", ) disconnectJob.cancel() @@ -1166,8 +1168,9 @@ class SharedRadioInterfaceServiceLivenessTest { val neverReleased = CompletableDeferred() var wedgeRanToCompletion = false + val logs = CapturingLogWriter.install() val wedged = launch { - service.runWhileSessionActive(session) { + service.runWhileSessionActive(session, "FromRadio packet TEXT_MESSAGE_APP") { wedgeStarted.complete(Unit) neverReleased.await() // simulates a handler stuck on an unbounded suspension wedgeRanToCompletion = true @@ -1175,7 +1178,7 @@ class SharedRadioInterfaceServiceLivenessTest { } wedgeStarted.await() val nextStarted = CompletableDeferred() - val next = launch { service.runWhileSessionActive(session) { nextStarted.complete(Unit) } } + val next = launch { service.runWhileSessionActive(session, "next") { nextStarted.complete(Unit) } } try { testDispatcher.scheduler.runCurrent() assertFalse(nextStarted.isCompleted, "ordered work is serialized behind the wedged handler") @@ -1187,7 +1190,14 @@ class SharedRadioInterfaceServiceLivenessTest { assertFalse(wedgeRanToCompletion, "the wedged handler must have been cancelled, not completed") assertTrue(nextStarted.isCompleted, "the handler timeout must release the pipeline for queued work") + val timeoutLines = logs.messages(Severity.Error).filter { "Session handler exceeded" in it } + assertEquals(1, timeoutLines.size, "one timeout line: $timeoutLines") + assertTrue( + "(handler=FromRadio packet TEXT_MESSAGE_APP)" in timeoutLines.single(), + "the timeout line must name the stuck handler: $timeoutLines", + ) } finally { + CapturingLogWriter.uninstall() neverReleased.complete(Unit) wedged.cancel() next.cancel() @@ -1210,7 +1220,7 @@ class SharedRadioInterfaceServiceLivenessTest { val neverReleased = CompletableDeferred() val wedged = launch { - service.runWhileSessionActive(session) { + service.runWhileSessionActive(session, "wedged") { wedgeStarted.complete(Unit) neverReleased.await() } diff --git a/core/testing/src/commonMain/kotlin/org/meshtastic/core/testing/FakeRadioInterfaceService.kt b/core/testing/src/commonMain/kotlin/org/meshtastic/core/testing/FakeRadioInterfaceService.kt index 26dd7fbffd..ff3caa6d11 100644 --- a/core/testing/src/commonMain/kotlin/org/meshtastic/core/testing/FakeRadioInterfaceService.kt +++ b/core/testing/src/commonMain/kotlin/org/meshtastic/core/testing/FakeRadioInterfaceService.kt @@ -115,8 +115,11 @@ class FakeRadioInterfaceService(override val serviceScope: CoroutineScope = Main } } - override suspend fun runWhileSessionActive(session: RadioSessionContext, block: suspend () -> Unit): Boolean = - sessionOperationMutex.withLock { runWithSessionLease(session) { block() } } + override suspend fun runWhileSessionActive( + session: RadioSessionContext, + label: String, + block: suspend () -> Unit, + ): Boolean = sessionOperationMutex.withLock { runWithSessionLease(session) { block() } } // Use an unbounded Channel to mirror SharedRadioInterfaceService semantics. A MutableSharedFlow would // hide the stop/start backlog bug that motivated the resetReceivedBuffer() API.