From fbb13a33d54036f9d1521e415a64f25418000f23 Mon Sep 17 00:00:00 2001 From: callebtc <93376500+callebtc@users.noreply.github.com> Date: Tue, 28 Jul 2026 13:17:35 +0200 Subject: [PATCH] Reject broadcast sends exceeding the receiver fragment cap Receivers hard-cap reassembly at MAX_FRAGMENTS_PER_ID (256), but the generic send path fragmented packets with no caller cap (0xFFFF), so broadcast file transfers above ~120 KB were fully transmitted yet undeliverable. FragmentingPacketSender now caps fragmentation at MAX_FRAGMENTS_PER_ID and reports failure via a new TransferProgressEvent.failed flag, which surfaces as DeliveryStatus.Failed in the UI and as a file_send error in the debug test hook instead of an indefinite wait. Adds FragmentingPacketSenderTest and a file_oversize mesh-lab scenario asserting sender-side rejection. --- .../android/testhook/TestHookDriver.kt | 96 ++++++++++++------- .../android/mesh/FragmentingPacketSender.kt | 25 ++++- .../android/mesh/TransferProgressManager.kt | 9 +- .../bitchat/android/ui/MediaSendingManager.kt | 11 ++- .../mesh/FragmentingPacketSenderTest.kt | 94 ++++++++++++++++++ tools/release_gate/mesh_lab.py | 25 ++++- 6 files changed, 214 insertions(+), 46 deletions(-) create mode 100644 app/src/test/java/com/bitchat/android/mesh/FragmentingPacketSenderTest.kt diff --git a/app/src/debug/java/com/bitchat/android/testhook/TestHookDriver.kt b/app/src/debug/java/com/bitchat/android/testhook/TestHookDriver.kt index de3f92b8..0f6675c4 100644 --- a/app/src/debug/java/com/bitchat/android/testhook/TestHookDriver.kt +++ b/app/src/debug/java/com/bitchat/android/testhook/TestHookDriver.kt @@ -18,6 +18,8 @@ import com.bitchat.android.services.AppStateStore import com.bitchat.android.ui.DataManager import com.bitchat.android.util.AppConstants import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.async +import kotlinx.coroutines.coroutineScope import kotlinx.coroutines.delay import kotlinx.coroutines.flow.first import kotlinx.coroutines.withContext @@ -282,46 +284,70 @@ object TestHookDriver { val encoded = packet.encode() ?: return err("file_send", "failed to TLV-encode packet") val transferId = sha256Hex(encoded) + return coroutineScope { + // Subscribe on a background dispatcher before sending so synchronous + // failure events are not missed (SharedFlow has replay=0). + val completion = async(Dispatchers.Default) { + TransferProgressManager.events.first { it.transferId == transferId && it.completed } + } + delay(50) + val sendError = dispatchFileSend(context, intent, mesh, peerID, packet, transferId) + if (sendError != null) { + completion.cancel() + return@coroutineScope sendError.put("cmd", "file_send") + } + val event = withTimeoutOrNull(timeoutMs) { completion.await() } + ?: return@coroutineScope err("file_send", "timeout waiting for transfer completion ($transferId)") + if (event.failed) { + return@coroutineScope err("file_send", "transfer rejected/failed before send ($transferId)") + .put("transfer_id", transferId) + } + ok("file_send") + .put("transfer_id", transferId) + .put("sent", event.sent) + .put("total", event.total) + .put("bytes", content.size) + .put("peer", peerID) + } + } + + private suspend fun dispatchFileSend( + context: Context, + intent: Intent, + mesh: MeshService, + peerID: String?, + packet: BitchatFilePacket, + transferId: String + ): JSONObject? { if (peerID == null) { mesh.sendFileBroadcast(packet) - } else { - if (!mesh.hasEstablishedSession(peerID)) { - val hs = handshake(context, peerID, intent) - if (hs.optString("status") != "ok") return hs.put("cmd", "file_send") - } - // Peer state (capabilities/identity) can lag session establishment; - // retry transient preparation states before giving up. - val prepDeadline = System.currentTimeMillis() + 30_000 - while (true) { - when (val prep = mesh.prepareFilePrivate(peerID, packet, transferId, allowLegacyFallback = false)) { - is PrivateMediaPreparation.Ready -> { - if (!prep.transfer.commit()) return err("file_send", "private transfer commit failed") - break - } - PrivateMediaPreparation.AwaitingPeerState, - PrivateMediaPreparation.NeedsHandshake -> { - if (System.currentTimeMillis() >= prepDeadline) { - return err("file_send", "private media preparation stuck at: $prep") - } - if (prep == PrivateMediaPreparation.NeedsHandshake) { - mesh.initiateNoiseHandshake(peerID) - } - delay(500) - } - else -> return err("file_send", "private media preparation: $prep") + return null + } + if (!mesh.hasEstablishedSession(peerID)) { + val hs = handshake(context, peerID, intent) + if (hs.optString("status") != "ok") return hs + } + // Peer state (capabilities/identity) can lag session establishment; + // retry transient preparation states before giving up. + val prepDeadline = System.currentTimeMillis() + 30_000 + while (true) { + when (val prep = mesh.prepareFilePrivate(peerID, packet, transferId, allowLegacyFallback = false)) { + is PrivateMediaPreparation.Ready -> { + return if (prep.transfer.commit()) null else err("file_send", "private transfer commit failed") } + PrivateMediaPreparation.AwaitingPeerState, + PrivateMediaPreparation.NeedsHandshake -> { + if (System.currentTimeMillis() >= prepDeadline) { + return err("file_send", "private media preparation stuck at: $prep") + } + if (prep == PrivateMediaPreparation.NeedsHandshake) { + mesh.initiateNoiseHandshake(peerID) + } + delay(500) + } + else -> return err("file_send", "private media preparation: $prep") } } - - val event = withTimeoutOrNull(timeoutMs) { - TransferProgressManager.events.first { it.transferId == transferId && it.completed } - } ?: return err("file_send", "timeout waiting for transfer completion ($transferId)") - return ok("file_send") - .put("transfer_id", transferId) - .put("sent", event.sent) - .put("total", event.total) - .put("bytes", content.size) - .put("peer", peerID) } private suspend fun fileRecv(context: Context, intent: Intent): JSONObject { diff --git a/app/src/main/java/com/bitchat/android/mesh/FragmentingPacketSender.kt b/app/src/main/java/com/bitchat/android/mesh/FragmentingPacketSender.kt index 3df51baf..9c6bd15e 100644 --- a/app/src/main/java/com/bitchat/android/mesh/FragmentingPacketSender.kt +++ b/app/src/main/java/com/bitchat/android/mesh/FragmentingPacketSender.kt @@ -31,7 +31,13 @@ class FragmentingPacketSender( sendSingle: (RoutedPacket) -> Boolean ): Boolean { val transferId = transferIdFor(routed) - val packets = packetsForTransport(routed) ?: return false + val packets = packetsForTransport(routed) + if (packets == null) { + if (transferId != null) { + TransferProgressManager.fail(transferId) + } + return false + } val total = packets.size if (total <= 1) { @@ -45,9 +51,13 @@ class FragmentingPacketSender( preparedPackets = null ) ) - if (sent && transferId != null) { - TransferProgressManager.progress(transferId, 1, 1) - TransferProgressManager.complete(transferId, 1) + if (transferId != null) { + if (sent) { + TransferProgressManager.progress(transferId, 1, 1) + TransferProgressManager.complete(transferId, 1) + } else { + TransferProgressManager.fail(transferId) + } } return sent } @@ -125,7 +135,12 @@ class FragmentingPacketSender( val manager = fragmentManager ?: return listOf(packet) return try { - val fragments = manager.createFragments(packet) + // Receivers hard-cap reassembly at MAX_FRAGMENTS_PER_ID; sending more + // fragments would be undeliverable, so reject here instead. + val fragments = manager.createFragments( + packet, + com.bitchat.android.util.AppConstants.Fragmentation.MAX_FRAGMENTS_PER_ID + ) if (fragments.isEmpty()) { Log.e(logTag, "Fragment manager returned no packets for packet type ${packet.type}") null diff --git a/app/src/main/java/com/bitchat/android/mesh/TransferProgressManager.kt b/app/src/main/java/com/bitchat/android/mesh/TransferProgressManager.kt index fbffb9aa..8dab20a2 100644 --- a/app/src/main/java/com/bitchat/android/mesh/TransferProgressManager.kt +++ b/app/src/main/java/com/bitchat/android/mesh/TransferProgressManager.kt @@ -11,7 +11,8 @@ data class TransferProgressEvent( val transferId: String, val sent: Int, val total: Int, - val completed: Boolean + val completed: Boolean, + val failed: Boolean = false ) object TransferProgressManager { @@ -22,9 +23,9 @@ object TransferProgressManager { fun start(id: String, total: Int) { emit(id, 0, total, false) } fun progress(id: String, sent: Int, total: Int) { emit(id, sent, total, sent >= total) } fun complete(id: String, total: Int) { emit(id, total, total, true) } + fun fail(id: String) { emit(id, 0, 0, done = true, failed = true) } - private fun emit(id: String, sent: Int, total: Int, done: Boolean) { - scope.launch { _events.emit(TransferProgressEvent(id, sent, total, done)) } + private fun emit(id: String, sent: Int, total: Int, done: Boolean, failed: Boolean = false) { + scope.launch { _events.emit(TransferProgressEvent(id, sent, total, done, failed)) } } } - diff --git a/app/src/main/java/com/bitchat/android/ui/MediaSendingManager.kt b/app/src/main/java/com/bitchat/android/ui/MediaSendingManager.kt index 68797e08..e982189f 100644 --- a/app/src/main/java/com/bitchat/android/ui/MediaSendingManager.kt +++ b/app/src/main/java/com/bitchat/android/ui/MediaSendingManager.kt @@ -739,7 +739,16 @@ class MediaSendingManager( fun handleTransferProgressEvent(evt: com.bitchat.android.mesh.TransferProgressEvent) { val msgId = synchronized(transferMessageMap) { transferMessageMap[evt.transferId] } if (msgId != null) { - if (evt.completed) { + if (evt.failed) { + messageManager.updateMessageDeliveryStatus( + msgId, + com.bitchat.android.model.DeliveryStatus.Failed("transfer could not be sent") + ) + synchronized(transferMessageMap) { + val msgIdRemoved = transferMessageMap.remove(evt.transferId) + if (msgIdRemoved != null) messageTransferMap.remove(msgIdRemoved) + } + } else if (evt.completed) { messageManager.updateMessageDeliveryStatus( msgId, com.bitchat.android.model.DeliveryStatus.Delivered(to = "mesh", at = java.util.Date()) diff --git a/app/src/test/java/com/bitchat/android/mesh/FragmentingPacketSenderTest.kt b/app/src/test/java/com/bitchat/android/mesh/FragmentingPacketSenderTest.kt new file mode 100644 index 00000000..c4225288 --- /dev/null +++ b/app/src/test/java/com/bitchat/android/mesh/FragmentingPacketSenderTest.kt @@ -0,0 +1,94 @@ +package com.bitchat.android.mesh + +import com.bitchat.android.model.RoutedPacket +import com.bitchat.android.protocol.BitchatPacket +import com.bitchat.android.protocol.MessageType +import com.bitchat.android.util.AppConstants +import kotlinx.coroutines.CoroutineScope +import kotlinx.coroutines.Dispatchers +import kotlinx.coroutines.SupervisorJob +import kotlinx.coroutines.flow.first +import kotlinx.coroutines.launch +import kotlinx.coroutines.runBlocking +import kotlinx.coroutines.withTimeout +import org.junit.Assert.assertFalse +import org.junit.Assert.assertTrue +import org.junit.Test +import org.junit.runner.RunWith +import org.robolectric.RobolectricTestRunner +import java.util.Random + +@RunWith(RobolectricTestRunner::class) +class FragmentingPacketSenderTest { + + private val senderID = "1122334455667788" + + private fun packetWithPayload(bytes: Int): BitchatPacket { + val payload = ByteArray(bytes) + Random(42).nextBytes(payload) + return BitchatPacket( + version = 2u, + type = MessageType.FILE_TRANSFER.value, + senderID = MeshPacketUtils.hexStringToByteArray(senderID), + recipientID = null, + timestamp = System.currentTimeMillis().toULong(), + payload = payload, + signature = null, + ttl = 7u + ) + } + + @Test + fun `oversized packet exceeding receiver fragment cap is rejected with fail event`() = runBlocking { + val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + val sender = FragmentingPacketSender(scope, FragmentManager(), "test") + // ~256 * 469 bytes fit; 1 MiB clearly exceeds MAX_FRAGMENTS_PER_ID + val packet = packetWithPayload(1024 * 1024) + var sent = false + + val failed = java.util.concurrent.ConcurrentLinkedQueue() + val collectJob = launch(Dispatchers.Default) { + TransferProgressManager.events.collect { event -> + if (event.failed) failed.add(event.transferId) + } + } + kotlinx.coroutines.delay(100) // activate subscription before emitting + + val accepted = sender.send(RoutedPacket(packet, transferId = "oversize-test"), "test") { sent = true; true } + assertFalse(accepted) + assertFalse(sent) + withTimeout(5_000) { + while (!failed.contains("oversize-test")) { + kotlinx.coroutines.delay(10) + } + } + collectJob.cancel() + Unit + } + + @Test + fun `packet within fragment cap is accepted`() = runBlocking { + val scope = CoroutineScope(Dispatchers.Default + SupervisorJob()) + val sender = FragmentingPacketSender(scope, FragmentManager(), "test", interFragmentDelayMs = 0L) + val packet = packetWithPayload(10_000) + var writes = 0 + + val accepted = sender.send(RoutedPacket(packet, transferId = "fits-test"), "test") { writes += 1; true } + assertTrue(accepted) + withTimeout(5_000) { + while (writes == 0) { + kotlinx.coroutines.delay(10) + } + } + assertTrue(writes > 0) + } + + @Test + fun `fragment count at cap boundary is not rejected`() { + val manager = FragmentManager() + val packet = packetWithPayload(AppConstants.Fragmentation.MAX_FRAGMENTS_PER_ID * 400) + val fragments = manager.createFragments(packet, AppConstants.Fragmentation.MAX_FRAGMENTS_PER_ID) + assertTrue(fragments.isNotEmpty()) + assertTrue(fragments.size <= AppConstants.Fragmentation.MAX_FRAGMENTS_PER_ID) + } +} diff --git a/tools/release_gate/mesh_lab.py b/tools/release_gate/mesh_lab.py index d1afe470..32aebd67 100644 --- a/tools/release_gate/mesh_lab.py +++ b/tools/release_gate/mesh_lab.py @@ -326,10 +326,33 @@ def scenario_raw(a: Device, b: Device) -> dict: return {"send": result} +def scenario_file_oversize(a: Device, b: Device, fixtures: dict[str, dict]) -> dict: + """Oversized broadcast file must be rejected sender-side (>256 fragments).""" + fixture = fixtures["medium_512k.bin"] + remote = a.push_fixture(fixture["path"]) + send = a.cmd("file_send", timeout_ms=60_000, path=remote) + rejected = send.get("status") == "error" and "rejected" in send.get("error", "") + if not rejected: + raise MeshLabError(f"expected sender-side rejection, got: {send}") + # Receiver must not see any file appear. + recv = b.cmd("file_recv", timeout_ms=15_000, name_contains="medium_512k") + if recv.get("status") == "ok": + raise MeshLabError(f"receiver unexpectedly saved an oversized file: {recv}") + return {"send": send, "receiver_saw_file": False} + + SCENARIOS = { "dm": scenario_dm, "broadcast": scenario_broadcast, - "file": lambda a, b: scenario_file(a, b, make_fixtures(Path(tempfile.mkdtemp(prefix="meshlab-fixtures-")))), + # Broadcast transfers are receiver-capped at 256 fragments (~120 KB); only + # the small fixture is end-to-end receivable. + "file": lambda a, b: scenario_file( + a, b, + make_fixtures(Path(tempfile.mkdtemp(prefix="meshlab-fixtures-")), names=["small_1k.bin"]), + ), + "file_oversize": lambda a, b: scenario_file_oversize( + a, b, make_fixtures(Path(tempfile.mkdtemp(prefix="meshlab-fixtures-"))) + ), # Private media is hard-capped at 256 fragments (PrivateMediaTransfer), so only # the small fixture fits; larger sizes are expected to be rejected by the sender. "file_private": lambda a, b: scenario_file(