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.
This commit is contained in:
callebtc 2026-07-28 13:17:35 +02:00
parent 41494d7c16
commit fbb13a33d5
6 changed files with 214 additions and 46 deletions

View File

@ -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 {

View File

@ -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

View File

@ -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)) }
}
}

View File

@ -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())

View File

@ -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<String>()
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)
}
}

View File

@ -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(