mirror of
https://github.com/meshtastic/Meshtastic-Android.git
synced 2026-10-02 16:44:33 -04:00
fix(service): stop the inbound pipeline waiting on node writes (#7464)
This commit is contained in:
1 parent
bd71d29770
commit
461e215a2d
16 files changed
+554
-152
No files matched your search
@@ -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)
|
||||
}
|
||||
}
|
||||
|
||||
+11
-7
@@ -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" }
|
||||
|
||||
+20
-9
@@ -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) }
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
+35
-9
@@ -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) }
|
||||
|
||||
|
||||
+39
-25
@@ -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<Pair<Int, RadioSessionContext?>, 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
|
||||
|
||||
+6
-8
@@ -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<Unit>()
|
||||
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) }
|
||||
|
||||
+5
-5
@@ -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()
|
||||
|
||||
+77
-55
@@ -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<Unit>()
|
||||
val releasePersistence = CompletableDeferred<Unit>()
|
||||
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<Int>()
|
||||
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<Int>(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<String>()
|
||||
everySuspend { radioInterfaceService.runWhileSessionActive(session, any(), any()) } calls
|
||||
{
|
||||
labels += it.arg<String>(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 ----------
|
||||
|
||||
+219
@@ -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 <https://www.gnu.org/licenses/>.
|
||||
*/
|
||||
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<MeshDataHandler>(MockMode.autofill)
|
||||
|
||||
private class GatedNodeRepository(private val delegate: FakeNodeRepository) : NodeRepository by delegate {
|
||||
val writeStarted = CompletableDeferred<Unit>()
|
||||
val releaseWrites = CompletableDeferred<Unit>()
|
||||
|
||||
val upserts = mutableMapOf<Int, Int>()
|
||||
|
||||
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<MeshNotificationManager>(MockMode.autofill),
|
||||
radio,
|
||||
scope.asServiceScope(),
|
||||
)
|
||||
.apply {
|
||||
notificationTitleFormatter = { shortName -> "New node seen: $shortName" }
|
||||
setNodeDbReady(true)
|
||||
setAllowNodeDbWrites(true)
|
||||
}
|
||||
val processor =
|
||||
MeshMessageProcessorImpl(
|
||||
nodeManager = nodeManager,
|
||||
serviceStateWriter = mock<ServiceRepository>(MockMode.autofill),
|
||||
meshLogRepository = lazy { mock<MeshLogRepository>(MockMode.autofill) },
|
||||
dataHandler = lazy { dataHandler },
|
||||
fromRadioDispatcher = mock<FromRadioPacketHandler>(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
|
||||
}
|
||||
}
|
||||
+95
-4
@@ -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<Unit>()
|
||||
val releaseFirstWrite = CompletableDeferred<Unit>()
|
||||
val persisted = mutableListOf<Node>()
|
||||
everySuspend { nodeRepository.upsert(any()) } calls
|
||||
{
|
||||
if (!firstWriteStarted.isCompleted) {
|
||||
firstWriteStarted.complete(Unit)
|
||||
releaseFirstWrite.await()
|
||||
}
|
||||
persisted += it.arg<Node>(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<Unit>()
|
||||
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<Node>()
|
||||
everySuspend { nodeRepository.upsert(any()) } calls { persisted += it.arg<Node>(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<Node>()
|
||||
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<Node>()
|
||||
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
|
||||
|
||||
+2
-2
@@ -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
|
||||
}
|
||||
|
||||
+3
-2
@@ -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() }
|
||||
}
|
||||
|
||||
|
||||
+18
-15
@@ -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 =
|
||||
|
||||
+2
-2
@@ -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
|
||||
|
||||
+16
-6
@@ -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<Unit>()
|
||||
|
||||
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<Unit>()
|
||||
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<Unit>()
|
||||
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<Unit>()
|
||||
|
||||
val wedged = launch {
|
||||
service.runWhileSessionActive(session) {
|
||||
service.runWhileSessionActive(session, "wedged") {
|
||||
wedgeStarted.complete(Unit)
|
||||
neverReleased.await()
|
||||
}
|
||||
|
||||
+5
-2
@@ -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.
|
||||
|
||||
Reference in new issue
Block a user