diff --git a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/repository/KermitMqttLogger.kt b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/repository/KermitMqttLogger.kt index 8bd2cbda5b..7aec5eb00a 100644 --- a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/repository/KermitMqttLogger.kt +++ b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/repository/KermitMqttLogger.kt @@ -17,6 +17,7 @@ package org.meshtastic.core.network.repository import co.touchlab.kermit.Logger +import org.meshtastic.core.network.transport.isExpectedConnectionFailure import org.meshtastic.mqtt.MqttLogLevel import org.meshtastic.mqtt.MqttLogger @@ -33,15 +34,32 @@ import org.meshtastic.mqtt.MqttLogger * Note: The production log level should be set to [MqttLogLevel.WARN] (not INFO) to prevent the library's own * INFO-level messages (which include endpoint addresses and topic strings) from reaching remote analytics sinks. */ -class KermitMqttLogger : MqttLogger { +class KermitMqttLogger( + /** Base logger to tag and delegate to. Injectable so tests can capture severities. */ + private val baseLogger: Logger = Logger, +) : MqttLogger { override fun log(level: MqttLogLevel, tag: String, message: String, throwable: Throwable?) { - val logger = Logger.withTag(tag) + val logger = baseLogger.withTag(tag) when (level) { MqttLogLevel.TRACE -> logger.v(throwable) { message } + MqttLogLevel.DEBUG -> logger.d(throwable) { message } + MqttLogLevel.INFO -> logger.i(throwable) { message } + MqttLogLevel.WARN -> logger.w(throwable) { message } - MqttLogLevel.ERROR -> logger.e(throwable) { message } + + // The library reports a broker that closed the socket, an unreachable host, or Wi-Fi dropping mid-read at + // ERROR ("Read loop error: Not enough data available"), and the client then reconnects on its own. Those + // are expected for a mobile client and dominated our error tracking, so log them at WARN — diagnosable in + // a support session, but not counted as application errors. + MqttLogLevel.ERROR -> + if (throwable?.isExpectedConnectionFailure() == true) { + logger.w(throwable) { message } + } else { + logger.e(throwable) { message } + } + MqttLogLevel.NONE -> return } } diff --git a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.kt b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.kt new file mode 100644 index 0000000000..fd759c616b --- /dev/null +++ b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.kt @@ -0,0 +1,37 @@ +/* + * 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.network.transport + +import kotlinx.io.IOException + +/** + * Whether this throwable is an ordinary, expected connection failure rather than an application defect. + * + * A radio that is powered off, a hostname that no longer resolves, a broker that closes the socket, or Wi-Fi dropping + * mid-read are all routine for a mesh client: the transport simply retries. Reporting them to crash/error tracking + * buries real defects — an unreachable host accounted for hundreds of "errors" per day with nothing to fix. + * + * Callers should still log these (at warn) so a support session can see them; they just should not be raised as errors. + */ +internal fun Throwable.isExpectedConnectionFailure(): Boolean = this is IOException || isPlatformConnectionFailure() + +/** + * Platform-specific expected connection failures that do not extend [IOException]. + * + * Kept separate so the common predicate above covers the (large) [IOException] family once, on every target. + */ +internal expect fun Throwable.isPlatformConnectionFailure(): Boolean diff --git a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/StreamFrameCodec.kt b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/StreamFrameCodec.kt index 85194a15e4..b680eb2886 100644 --- a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/StreamFrameCodec.kt +++ b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/StreamFrameCodec.kt @@ -19,6 +19,7 @@ package org.meshtastic.core.network.transport import co.touchlab.kermit.Logger import kotlinx.coroutines.sync.Mutex import kotlinx.coroutines.sync.withLock +import kotlin.concurrent.Volatile /** * Meshtastic stream framing codec — pure Kotlin, no platform dependencies. @@ -61,55 +62,141 @@ class StreamFrameCodec( private val rxPacket = ByteArray(MAX_TO_FROM_RADIO_SIZE) private val debugLineBuf = StringBuilder() + // Written only by the transport's read thread, but published through [framingDesyncCount] for diagnostics, so a + // reader on another thread needs a visibility guarantee. Matches the metrics counters in TcpTransport. The rest of + // the state machine stays plain: it is single-reader by contract. + @Volatile private var desyncCount = 0L + + // Purely local to the read thread — it only throttles logging — so it needs no cross-thread guarantee. + private var consecutiveDesyncs = 0 + + /** + * Number of framing desyncs observed since the last [reset]. + * + * Exposed for diagnostics and tests. A non-zero count on an otherwise healthy link means the peer is emitting bytes + * we cannot frame — truncated writes, a corrupt length prefix, or non-protocol output on the wire. + */ + val framingDesyncCount: Long + get() = desyncCount + /** * Process a single incoming byte through the stream framing state machine. * * Call this repeatedly with bytes from the transport (serial, TCP, etc). When a complete packet is decoded, * [onPacketReceived] is invoked. + * + * Malformed input never throws. On a broken frame the machine records a desync, rewinds to the hunting state, and + * resynchronizes on the next valid frame header — a corrupt frame costs one frame, not the connection. */ fun processInputByte(c: Byte) { - var nextPtr = ptr + 1 + // A byte that breaks framing may itself be the START1 of the next frame, so re-examine it from the hunting + // state instead of dropping it. This matters most for a corrupt length prefix whose low byte is 0x94: the + // machine rejects the length and the very next header would otherwise be missing its start byte. + // + // Discarding the offending byte is what let one bad byte cascade. The machine would resume hunting partway + // into the *payload* of the frame that followed, re-sync on an incidental 0x94/0xc3 pair inside it, and hand + // a misaligned byte range to the protobuf parser — surfacing downstream as ProtocolException or + // "Unexpected call to beginMessage()". + // + // The retry always runs in state 0, which never reports a desync, so this terminates after at most one retry. + if (!step(c)) step(c) + } - fun lostSync() { - Logger.e { "$logTag: Lost protocol sync" } - nextPtr = 0 + /** + * Advance the state machine by a single byte. + * + * @return `false` when [c] broke framing and must be re-examined from the hunting state. + */ + private fun step(c: Byte): Boolean = when (ptr) { + 0 -> { + // Hunting for a frame start; anything else is device debug output. + if (c == START1) ptr = 1 else debugOut(c) + true } - fun deliverPacket() { - val buf = rxPacket.copyOf(packetLen) - onPacketReceived(buf) - nextPtr = 0 + 1 -> stepExpectStart2(c) + + 2 -> { + msb = c.toInt() and 0xff + ptr = 3 + true } - when (ptr) { - 0 -> - if (c != START1) { - debugOut(c) - nextPtr = 0 - } + 3 -> stepLength(c) - 1 -> if (c != START2) lostSync() + else -> { + appendPayloadByte(c) + true + } + } - 2 -> msb = c.toInt() and 0xff + /** + * State 1: expecting START2. + * + * A repeated START1 is padding rather than corruption — the wake sequence is four of them, and a peer may pad an + * idle link the same way — so hold in this state and take the next byte as the candidate START2. + */ + private fun stepExpectStart2(c: Byte): Boolean = when (c) { + START2 -> { + ptr = 2 + true + } - 3 -> { - lsb = c.toInt() and 0xff - packetLen = (msb shl 8) or lsb - if (packetLen > MAX_TO_FROM_RADIO_SIZE) { - lostSync() - } else if (packetLen == 0) { - deliverPacket() - } + START1 -> { + // Padding — hold here, still expecting START2. + true + } + + else -> { + lostSync("expected START2") + false + } + } + + /** State 3: low byte of the length prefix, completing the header. */ + private fun stepLength(c: Byte): Boolean { + lsb = c.toInt() and 0xff + packetLen = (msb shl 8) or lsb + return when { + packetLen > MAX_TO_FROM_RADIO_SIZE -> { + lostSync("declared length $packetLen exceeds $MAX_TO_FROM_RADIO_SIZE") + false + } + + packetLen == 0 -> { + deliverPacket() + true } else -> { - rxPacket[ptr - HEADER_SIZE] = c - if (ptr - HEADER_SIZE + 1 == packetLen) { - deliverPacket() - } + ptr = HEADER_SIZE + true } } - ptr = nextPtr + } + + /** States 4 and up: accumulating the payload declared by the header. */ + private fun appendPayloadByte(c: Byte) { + rxPacket[ptr - HEADER_SIZE] = c + if (ptr - HEADER_SIZE + 1 == packetLen) deliverPacket() else ptr++ + } + + private fun lostSync(reason: String) { + desyncCount++ + consecutiveDesyncs++ + // Structural detail only — never payload bytes, which can carry message content or keys. + val detail = "$logTag: Lost protocol sync ($reason); resynchronizing. Desyncs since connect: $desyncCount" + // Warn once per disruption, then drop to debug: a peer streaming non-protocol bytes would otherwise emit a + // warning every few hundred bytes. + if (consecutiveDesyncs == 1) Logger.w { detail } else Logger.d { detail } + ptr = 0 + } + + private fun deliverPacket() { + consecutiveDesyncs = 0 + // Rewind before dispatching so the machine is already hunting if the callback re-enters. + ptr = 0 + onPacketReceived(rxPacket.copyOf(packetLen)) } /** @@ -141,6 +228,8 @@ class StreamFrameCodec( msb = 0 lsb = 0 packetLen = 0 + desyncCount = 0 + consecutiveDesyncs = 0 debugLineBuf.clear() } diff --git a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/TcpTransport.kt b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/TcpTransport.kt index 7e36cacbcc..6a4e3b501b 100644 --- a/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/TcpTransport.kt +++ b/core/network/src/commonMain/kotlin/org/meshtastic/core/network/transport/TcpTransport.kt @@ -213,7 +213,14 @@ class TcpTransport( disconnectSocket() throw ce } catch (@Suppress("TooGenericExceptionCaught") ex: Throwable) { - Logger.e(ex) { "$logTag: [$address] TCP exception" } + // An unreachable host is routine, not a defect — most often a hostname that no longer resolves + // (UnresolvedAddressException, which is an IllegalArgumentException and so misses the IOException + // branch above). Log it, retry it, but keep it out of error tracking. + if (ex.isExpectedConnectionFailure()) { + Logger.w(ex) { "$logTag: [$address] Radio unreachable" } + } else { + Logger.e(ex) { "$logTag: [$address] TCP exception" } + } disconnectSocket() false } diff --git a/core/network/src/commonTest/kotlin/org/meshtastic/core/network/repository/KermitMqttLoggerTest.kt b/core/network/src/commonTest/kotlin/org/meshtastic/core/network/repository/KermitMqttLoggerTest.kt new file mode 100644 index 0000000000..d053c1ca1e --- /dev/null +++ b/core/network/src/commonTest/kotlin/org/meshtastic/core/network/repository/KermitMqttLoggerTest.kt @@ -0,0 +1,92 @@ +/* + * 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.network.repository + +import co.touchlab.kermit.LogWriter +import co.touchlab.kermit.Logger +import co.touchlab.kermit.Severity +import co.touchlab.kermit.loggerConfigInit +import kotlinx.io.IOException +import org.meshtastic.mqtt.MqttLogLevel +import kotlin.test.Test +import kotlin.test.assertEquals + +class KermitMqttLoggerTest { + + private class CapturingWriter : LogWriter() { + val entries = mutableListOf>() + + override fun log(severity: Severity, message: String, tag: String, throwable: Throwable?) { + entries += Triple(severity, tag, message) + } + } + + private val writer = CapturingWriter() + private val mqttLogger = KermitMqttLogger(Logger(loggerConfigInit(writer), tag = "Test")) + + @Test + fun `an expected connection failure logged at ERROR is downgraded to WARN`() { + // The MQTT client reports a broker-side socket close this way; it reconnects on its own, so it must not be + // counted as an application error. + mqttLogger.log( + level = MqttLogLevel.ERROR, + tag = "MqttConnection", + message = "Read loop error: Not enough data available", + throwable = IOException("Not enough data available"), + ) + + assertEquals(1, writer.entries.size) + assertEquals(Severity.Warn, writer.entries[0].first) + assertEquals("MqttConnection", writer.entries[0].second) + } + + @Test + fun `a genuine error logged at ERROR stays an error`() { + mqttLogger.log( + level = MqttLogLevel.ERROR, + tag = "MqttConnection", + message = "Malformed packet", + throwable = IllegalStateException("bad state"), + ) + + assertEquals(Severity.Error, writer.entries.single().first) + } + + @Test + fun `an error with no throwable stays an error`() { + mqttLogger.log(level = MqttLogLevel.ERROR, tag = "MqttClient", message = "boom", throwable = null) + + assertEquals(Severity.Error, writer.entries.single().first) + } + + @Test + fun `other levels map straight through and keep the library tag`() { + mqttLogger.log(MqttLogLevel.WARN, "MqttClient", "warn", null) + mqttLogger.log(MqttLogLevel.INFO, "MqttClient", "info", null) + mqttLogger.log(MqttLogLevel.DEBUG, "MqttClient", "debug", null) + + assertEquals(listOf(Severity.Warn, Severity.Info, Severity.Debug), writer.entries.map { it.first }) + assertEquals(listOf("MqttClient", "MqttClient", "MqttClient"), writer.entries.map { it.second }) + } + + @Test + fun `NONE is dropped`() { + mqttLogger.log(MqttLogLevel.NONE, "MqttClient", "nothing", null) + + assertEquals(0, writer.entries.size) + } +} diff --git a/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/ConnectionFailuresTest.kt b/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/ConnectionFailuresTest.kt new file mode 100644 index 0000000000..261e168786 --- /dev/null +++ b/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/ConnectionFailuresTest.kt @@ -0,0 +1,44 @@ +/* + * 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.network.transport + +import kotlinx.io.IOException +import kotlin.test.Test +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class ConnectionFailuresTest { + + @Test + fun `an IO failure is an expected connection failure`() { + assertTrue(IOException("Connection reset by peer").isExpectedConnectionFailure()) + } + + @Test + fun `an IO subclass is an expected connection failure`() { + class TruncatedRead(message: String) : IOException(message) + + assertTrue(TruncatedRead("Not enough data available").isExpectedConnectionFailure()) + } + + @Test + fun `a programming error is not an expected connection failure`() { + assertFalse(IllegalStateException("Unexpected call to beginMessage()").isExpectedConnectionFailure()) + assertFalse(IllegalArgumentException("bad argument").isExpectedConnectionFailure()) + assertFalse(NullPointerException().isExpectedConnectionFailure()) + } +} diff --git a/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/StreamFrameCodecTest.kt b/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/StreamFrameCodecTest.kt index 5313fd17a0..24ee786a09 100644 --- a/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/StreamFrameCodecTest.kt +++ b/core/network/src/commonTest/kotlin/org/meshtastic/core/network/transport/StreamFrameCodecTest.kt @@ -21,6 +21,7 @@ import io.kotest.matchers.shouldBe import io.kotest.property.Arb import io.kotest.property.arbitrary.byte import io.kotest.property.arbitrary.byteArray +import io.kotest.property.arbitrary.filter import io.kotest.property.arbitrary.int import io.kotest.property.checkAll import kotlinx.coroutines.test.runTest @@ -154,6 +155,154 @@ class StreamFrameCodecTest { assertEquals(listOf(0xAA.toByte()), receivedPackets[0].toList()) } + // region Resynchronization after malformed input + // + // The framing protocol has no checksum, so a frame whose length prefix is corrupt cannot be detected as such — + // the codec will consume that many bytes and may swallow a following frame. What it must always do is recover: + // never throw, never wedge, and pick up the next well-formed frame. + + @Test + fun `parses a frame whose header is preceded by a repeated START1`() { + // Regression: the second 0x94 failed the START2 check and was *discarded*, taking the real header with it, + // so the frame vanished silently. A repeated start byte is padding and must not cost a frame. + val data = byteArrayOf(0x94.toByte(), 0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x42) + + data.forEach { codec.processInputByte(it) } + + assertEquals(1, receivedPackets.size) + assertEquals(listOf(0x42.toByte()), receivedPackets[0].toList()) + assertEquals(0L, codec.framingDesyncCount, "padding is not a desync") + } + + @Test + fun `re-examines the rejected byte when a corrupt length ends in START1`() { + // Declared length 0x0294 = 660 > MAX, and the low byte that triggered the rejection is itself START1 — it is + // the start of the next frame's header. Dropping it would consume that header and desync into the payload. + val corrupt = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x02, 0x94.toByte()) + val good = byteArrayOf(0xc3.toByte(), 0x00, 0x01, 0x42) + + (corrupt + good).forEach { codec.processInputByte(it) } + + assertEquals(1, receivedPackets.size) + assertEquals(listOf(0x42.toByte()), receivedPackets[0].toList()) + assertEquals(1L, codec.framingDesyncCount) + } + + @Test + fun `parses a frame after a START1 run of any length`() { + // Odd-length runs happened to work before the fix; even-length ones consumed the header. Cover both. + for (runLength in 1..6) { + receivedPackets.clear() + val codec = StreamFrameCodec(onPacketReceived = { receivedPackets.add(it) }, logTag = "Test") + val data = + ByteArray(runLength) { 0x94.toByte() } + byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x42) + + data.forEach { codec.processInputByte(it) } + + assertEquals(1, receivedPackets.size, "run length $runLength should still yield one frame") + assertEquals(listOf(0x42.toByte()), receivedPackets[0].toList(), "run length $runLength") + } + } + + @Test + fun `parses a frame that arrives immediately after wake bytes`() { + // WAKE_BYTES is four consecutive START1 — an even-length run, so this was the common real-world trigger. + val data = StreamFrameCodec.WAKE_BYTES + byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x7F) + + data.forEach { codec.processInputByte(it) } + + assertEquals(1, receivedPackets.size) + assertEquals(listOf(0x7F.toByte()), receivedPackets[0].toList()) + } + + @Test + fun `recovers from a corrupt length prefix and parses the next frame`() { + // Declared length 0x0201 = 513 > MAX_TO_FROM_RADIO_SIZE, so the header is rejected outright. + val corrupt = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x02, 0x01) + val good = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x02, 0x11, 0x22) + + (corrupt + good).forEach { codec.processInputByte(it) } + + assertEquals(1, receivedPackets.size) + assertEquals(listOf(0x11.toByte(), 0x22.toByte()), receivedPackets[0].toList()) + assertEquals(1L, codec.framingDesyncCount) + } + + @Test + fun `recovers from a truncated frame and parses a later frame`() { + // Declares 5 payload bytes but only 2 follow, so the codec absorbs the next frame's bytes as payload. That + // frame is unrecoverable without a checksum; what matters is that the codec resynchronizes afterwards. + val truncated = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x05, 0xAA.toByte(), 0xBB.toByte()) + val eaten = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x42) + val recovered = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x43) + + (truncated + eaten + recovered).forEach { codec.processInputByte(it) } + + assertEquals(listOf(0x43.toByte()), receivedPackets.last().toList()) + } + + @Test + fun `recovers from garbage between two frames`() { + val first = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x01) + val garbage = byteArrayOf(0x00, 0x7F, 0x3C, 0xFE.toByte(), 0x12, 0x94.toByte(), 0x00, 0x94.toByte()) + val second = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x02) + + (first + garbage + second).forEach { codec.processInputByte(it) } + + assertEquals(2, receivedPackets.size) + assertEquals(listOf(0x01.toByte()), receivedPackets[0].toList()) + assertEquals(listOf(0x02.toByte()), receivedPackets[1].toList()) + } + + @Test + fun `recovers from a frame header split by a desync in the middle of a payload`() { + // A frame whose payload contains an incidental START1/START2 pair must still be delivered verbatim, and the + // following frame must still parse. + val payload = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01) + val frame = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, payload.size.toByte()) + payload + val next = byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x5A) + + (frame + next).forEach { codec.processInputByte(it) } + + assertEquals(2, receivedPackets.size) + assertEquals(payload.toList(), receivedPackets[0].toList()) + assertEquals(listOf(0x5A.toByte()), receivedPackets[1].toList()) + } + + @Test + fun `always parses a well-formed frame following non-protocol noise`() = runTest { + // Noise with no START1 cannot be mistaken for a header, so recovery here is total and can be asserted for + // arbitrary input rather than a handful of fixed cases. + checkAll(Arb.byteArray(Arb.int(0, 300), Arb.byte().filter { it != 0x94.toByte() })) { noise -> + val received = mutableListOf() + val codec = StreamFrameCodec(onPacketReceived = { received.add(it) }) + + noise.forEach { codec.processInputByte(it) } + byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x63).forEach { codec.processInputByte(it) } + + received.size.shouldBe(1) + received[0].toList().shouldBe(listOf(0x63.toByte())) + } + } + + @Test + fun `desync count tracks framing errors and resets with the codec`() { + assertEquals(0L, codec.framingDesyncCount) + + byteArrayOf(0x94.toByte(), 0x00).forEach { codec.processInputByte(it) } + assertEquals(1L, codec.framingDesyncCount) + + // A clean frame after the desync still parses, and does not add to the count. + byteArrayOf(0x94.toByte(), 0xc3.toByte(), 0x00, 0x01, 0x42).forEach { codec.processInputByte(it) } + assertEquals(1L, codec.framingDesyncCount) + assertEquals(1, receivedPackets.size) + + codec.reset() + assertEquals(0L, codec.framingDesyncCount) + } + + // endregion + @Test fun `frameAndSend produces correct header for 1-byte payload`() = runTest { val payload = byteArrayOf(0x42.toByte()) diff --git a/core/network/src/iosMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.ios.kt b/core/network/src/iosMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.ios.kt new file mode 100644 index 0000000000..b63eb49a2d --- /dev/null +++ b/core/network/src/iosMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.ios.kt @@ -0,0 +1,23 @@ +/* + * 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.network.transport + +/** + * On Kotlin/Native, Ktor surfaces socket and resolution failures as `kotlinx.io.IOException` (wrapping the underlying + * POSIX error), which the common predicate already covers. There is no extra non-IOException family to classify here. + */ +internal actual fun Throwable.isPlatformConnectionFailure(): Boolean = false diff --git a/core/network/src/jvmAndroidMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.jvmAndroid.kt b/core/network/src/jvmAndroidMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.jvmAndroid.kt new file mode 100644 index 0000000000..376b0ee41d --- /dev/null +++ b/core/network/src/jvmAndroidMain/kotlin/org/meshtastic/core/network/transport/ConnectionFailures.jvmAndroid.kt @@ -0,0 +1,29 @@ +/* + * 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.network.transport + +import java.nio.channels.UnresolvedAddressException + +/** + * On JVM/Android nearly every socket failure already extends `java.io.IOException`, which the common predicate covers — + * `ConnectException`, `SocketException`, `UnknownHostException`, `EOFException` and friends all qualify there. + * + * [UnresolvedAddressException] is the exception: Ktor throws it when a hostname cannot be resolved, and it extends + * `IllegalArgumentException`, so it slips past every `catch (IOException)` in the transport stack and lands in the + * generic handler that reports errors. + */ +internal actual fun Throwable.isPlatformConnectionFailure(): Boolean = this is UnresolvedAddressException diff --git a/core/network/src/jvmTest/kotlin/org/meshtastic/core/network/transport/ConnectionFailuresJvmTest.kt b/core/network/src/jvmTest/kotlin/org/meshtastic/core/network/transport/ConnectionFailuresJvmTest.kt new file mode 100644 index 0000000000..4a06789e2f --- /dev/null +++ b/core/network/src/jvmTest/kotlin/org/meshtastic/core/network/transport/ConnectionFailuresJvmTest.kt @@ -0,0 +1,50 @@ +/* + * 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.network.transport + +import java.io.EOFException +import java.net.ConnectException +import java.net.SocketException +import java.net.UnknownHostException +import java.nio.channels.UnresolvedAddressException +import kotlin.test.Test +import kotlin.test.assertFalse +import kotlin.test.assertTrue + +class ConnectionFailuresJvmTest { + + @Test + fun `an unresolvable hostname is an expected connection failure`() { + // Ktor throws this when the address will not resolve. It extends IllegalArgumentException, so before this + // classification it escaped every catch(IOException) and was reported as an application error. + assertTrue(UnresolvedAddressException().isExpectedConnectionFailure()) + } + + @Test + fun `ordinary unreachable-host failures are expected connection failures`() { + assertTrue(ConnectException("Connection refused").isExpectedConnectionFailure()) + assertTrue(UnknownHostException("no.such.host").isExpectedConnectionFailure()) + assertTrue(SocketException("Connection reset").isExpectedConnectionFailure()) + assertTrue(EOFException("Not enough data available").isExpectedConnectionFailure()) + } + + @Test + fun `an unrelated IllegalArgumentException is still an error`() { + // Guards the classification from widening to every IllegalArgumentException. + assertFalse(IllegalArgumentException("bad port").isExpectedConnectionFailure()) + } +}