From cf8345386996c561f176049b2f83fbd07956995e Mon Sep 17 00:00:00 2001 From: Christophe Besson Date: Sat, 3 Oct 2026 12:53:35 +0200 Subject: fix(android): the cast relay spools a receiver's lead to disk Fragments are 5-10 MB at a film's bitrate; dropped past 8 MB in memory, the TV froze for their length. Each receiver now reads from its own spool file, deleted with it; nothing is dropped short of a disk bound. Co-Authored-By: Claude Opus 5.5 --- .../kotlin/org/meshbay/client/cast/CastChannels.kt | 2 +- .../kotlin/org/meshbay/client/cast/CastRelay.kt | 140 +++++++++++++++------ .../kotlin/org/meshbay/client/CastRelayTest.kt | 31 ++++- packages/meshbay-hub/tests/test_android_cast.py | 16 ++- 4 files changed, 144 insertions(+), 45 deletions(-) (limited to 'packages') diff --git a/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastChannels.kt b/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastChannels.kt index 1044209..ecad8ec 100644 --- a/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastChannels.kt +++ b/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastChannels.kt @@ -22,7 +22,7 @@ class CastChannels( private val tell: (String) -> Unit = {}, ) { val control = CastControl(context) - val relay = CastRelay(::lanAddress) + val relay = CastRelay(::lanAddress, java.io.File(context.cacheDir, "cast-relay")) fun handles(channel: String) = channel.startsWith("cast:") diff --git a/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastRelay.kt b/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastRelay.kt index 0d85a17..4d52052 100644 --- a/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastRelay.kt +++ b/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastRelay.kt @@ -3,7 +3,11 @@ package org.meshbay.client.cast import android.util.Log import org.json.JSONObject import java.io.BufferedReader +import java.io.File import java.io.InputStreamReader +import java.io.RandomAccessFile +import java.nio.ByteBuffer +import java.nio.channels.FileChannel import java.io.OutputStream import java.net.InetAddress import java.net.InetSocketAddress @@ -12,7 +16,6 @@ import java.net.Socket import java.net.SocketException import java.security.SecureRandom import java.util.concurrent.ConcurrentHashMap -import java.util.concurrent.LinkedBlockingDeque import java.util.concurrent.atomic.AtomicLong /** @@ -29,19 +32,55 @@ import java.util.concurrent.atomic.AtomicLong * - `Cache-Control: no-store` on every response; * - the server closed when playback stops — zero residual surface. * - * `java.net` only, so the JVM tests run the real thing. + * What a receiver has not read yet waits in a file, not in memory. The page + * runs ahead of the television by its whole read-ahead — tens of megabytes in + * the first seconds — and a fragment is a segment, megabytes at a film's + * bitrate (5 to 10 MB measured). Held in memory under the desktop's 8 MB + * bound, fragments were dropped and the picture froze for their length; held + * in memory without a bound, the heap went. A spool file per receiver holds + * the lead at no cost to either, and is deleted when the receiver goes. + * + * `java.net` and `java.nio` only, so the JVM tests run the real thing. */ -class CastRelay(private val lanAddress: () -> InetAddress?) { - - private class Client(val socket: Socket, val out: OutputStream) { - val queue = LinkedBlockingDeque() - val queued = AtomicLong() +class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir: File) { + + /** + * One receiver: the fragments it has been given, appended to its spool + * file, and how far it has read. Positional reads and writes on one + * channel, so the pusher and the reader never share a file pointer. + */ + private class Client(val socket: Socket, val out: OutputStream, val file: File) { + val channel: FileChannel = RandomAccessFile(file, "rw").channel + val lock = Object() + var written = 0L // guarded by lock + var read = 0L // guarded by lock + var ended = false // guarded by lock @Volatile var closed = false val sent = AtomicLong() val dropped = AtomicLong() @Volatile var lastReport = System.currentTimeMillis() + + fun waiting() = synchronized(lock) { written - read } + + fun append(data: ByteArray) { + var at = synchronized(lock) { written } + val buf = ByteBuffer.wrap(data) + while (buf.hasRemaining()) at += channel.write(buf, at) + synchronized(lock) { written = at; lock.notifyAll() } + } + + fun end() = synchronized(lock) { ended = true; lock.notifyAll() } + + fun close() { + closed = true + synchronized(lock) { lock.notifyAll() } + try { channel.close() } catch (e: Exception) {} + file.delete() + } } + private val spoolSeq = AtomicLong() + private class Subtitle(val vtt: ByteArray, val language: String, val label: String) @Volatile private var server: ServerSocket? = null @@ -98,6 +137,9 @@ class CastRelay(private val lanAddress: () -> InetAddress?) { synchronized(ring) { ring.clear(); ringBytes = 0 } clients.clear() accum = BoxAccumulator() + // What a killed process left behind. + spoolDir.mkdirs() + spoolDir.listFiles()?.forEach { it.delete() } subtitle = null setSubtitle(sub) @@ -133,6 +175,8 @@ class CastRelay(private val lanAddress: () -> InetAddress?) { if (pushes <= 3 || pushes % 100 == 0) log("push #$pushes: $length bytes, $fragments fragments so far") for (frag in accum.push(data, offset, length)) { fragments++ + // The backlog and the clients under one lock: a receiver joining + // gets each fragment exactly once, from the backlog or from here. synchronized(ring) { ring.addLast(frag); ringBytes += frag.size // Bounded in fragments, as the desktop is, and in bytes too: a @@ -141,33 +185,34 @@ class CastRelay(private val lanAddress: () -> InetAddress?) { while (ring.size > RING_CAP || (ringBytes > RING_MAX_BYTES && ring.size > 1)) { ringBytes -= ring.removeFirst().size } - } - for (c in clients) { - if (c.queued.get() > BACKPRESSURE_HIGH) { - // Dropped for a slow client: the picture on the receiver - // freezes until the next fragment it does get. Silent on - // the desktop too; said here, with what was waiting. - val n = c.dropped.incrementAndGet() - if (n <= 3 || n % 20 == 0L) { - log("DROPPED fragment for ${c.socket.inetAddress?.hostAddress}: ${frag.size} bytes, " + - "${c.queued.get() / 1024} KiB already waiting, $n dropped so far") + for (c in clients) { + if (c.waiting() > spoolLimit()) { + // Only with the disk bound reached: the picture on the + // receiver freezes until the next fragment it gets. + val n = c.dropped.incrementAndGet() + if (n <= 3 || n % 20 == 0L) { + log("DROPPED fragment for ${c.socket.inetAddress?.hostAddress}: ${frag.size} bytes, " + + "${c.waiting() / 1048576} MiB already waiting, $n dropped so far") + } + continue } - continue + try { c.append(frag) } catch (e: Exception) { log("spool write failed: ${e.message}") } } - c.queued.addAndGet(frag.size.toLong()) - c.queue.offer(frag) } } } /** End of film: every client's response is ended, the server stays until stop. */ fun finish() { - for (c in clients) c.queue.offer(END) + for (c in clients) c.end() } + /** Half the free space, at most 2 GiB: a receiver this far behind is not coming back. */ + private fun spoolLimit(): Long = minOf(SPOOL_MAX_BYTES, spoolDir.usableSpace / 2) + @Synchronized fun stop() { - for (c in clients) { c.closed = true; c.queue.offer(END); try { c.socket.close() } catch (e: Exception) {} } + for (c in clients) { c.close(); try { c.socket.close() } catch (e: Exception) {} } clients.clear() server?.let { try { it.close() } catch (e: Exception) {} } server = null @@ -238,38 +283,51 @@ class CastRelay(private val lanAddress: () -> InetAddress?) { "Connection" to "keep-alive", "Transfer-Encoding" to "chunked", )) - val client = Client(socket, out) initSegment?.let { chunk(out, it) } - val backlog = synchronized(ring) { ring.toList() } - for (frag in backlog) chunk(out, frag) out.flush() - clients.add(client) - log("client ${socket.inetAddress?.hostAddress} served init + ${backlog.size} fragments") + val client = Client(socket, out, File(spoolDir, "client-${spoolSeq.incrementAndGet()}.spool")) + val backlog = synchronized(ring) { + for (frag in ring) client.append(frag) + clients.add(client) + ring.size + } + log("client ${socket.inetAddress?.hostAddress} served init + $backlog fragments") + val block = ByteBuffer.allocate(BLOCK) try { while (!client.closed) { - val frag = client.queue.take() - if (frag === END) { out.write("0\r\n\r\n".toByteArray()); out.flush(); break } + val at = synchronized(client.lock) { + while (client.read == client.written && !client.ended && !client.closed) client.lock.wait() + if (client.read == client.written) -1L else client.read + } + if (at < 0) { + if (!client.closed) { out.write("0\r\n\r\n".toByteArray()); out.flush() } + break + } + block.clear() + val n = client.channel.read(block, at) + if (n <= 0) continue val t0 = System.currentTimeMillis() - chunk(out, frag) + chunk(out, block.array(), n) out.flush() val took = System.currentTimeMillis() - t0 - client.queued.addAndGet(-frag.size.toLong()) - client.sent.incrementAndGet() + synchronized(client.lock) { client.read += n } + client.sent.addAndGet(n.toLong()) // A write that blocks is the receiver not reading (or the // Wi-Fi not carrying): the one place a stall downstream shows. - if (took > 2000) log("slow write to ${socket.inetAddress?.hostAddress}: ${frag.size} bytes took $took ms") + if (took > 5000) log("slow write to ${socket.inetAddress?.hostAddress}: $n bytes took $took ms") val now = System.currentTimeMillis() if (now - client.lastReport >= 10_000) { client.lastReport = now - log("client ${socket.inetAddress?.hostAddress}: sent ${client.sent.get()} fragments, " + - "${client.queued.get() / 1024} KiB waiting, ${client.dropped.get()} dropped") + log("client ${socket.inetAddress?.hostAddress}: sent ${client.sent.get() / 1048576} MiB, " + + "${client.waiting() / 1048576} MiB waiting in the spool, ${client.dropped.get()} dropped") } } } catch (e: Exception) { // The receiver went away; nothing to tell anyone. } finally { log("client ${socket.inetAddress?.hostAddress} gone") - clients.remove(client) + synchronized(ring) { clients.remove(client) } + client.close() try { socket.close() } catch (e: Exception) {} } } @@ -301,10 +359,10 @@ class CastRelay(private val lanAddress: () -> InetAddress?) { const val RING_CAP = 64 const val RING_MAX_BYTES = 32L * 1024 * 1024 - const val BACKPRESSURE_HIGH = 8L * 1024 * 1024 + const val SPOOL_MAX_BYTES = 2L * 1024 * 1024 * 1024 + private const val BLOCK = 256 * 1024 const val PORT_BASE = 19550 const val PORT_COUNT = 4 - private val END = ByteArray(0) // Never the URL: it carries the token. private fun log(m: String) = try { Log.i("MeshBayCast", "[relay] $m") } catch (e: RuntimeException) { /* JVM tests */ } @@ -342,9 +400,9 @@ class CastRelay(private val lanAddress: () -> InetAddress?) { socket.close() } - private fun chunk(out: OutputStream, data: ByteArray) { - out.write("${Integer.toHexString(data.size)}\r\n".toByteArray()) - out.write(data) + private fun chunk(out: OutputStream, data: ByteArray, length: Int = data.size) { + out.write("${Integer.toHexString(length)}\r\n".toByteArray()) + out.write(data, 0, length) out.write("\r\n".toByteArray()) } } diff --git a/packages/meshbay-android/app/src/test/kotlin/org/meshbay/client/CastRelayTest.kt b/packages/meshbay-android/app/src/test/kotlin/org/meshbay/client/CastRelayTest.kt index e11c8c0..1c604d0 100644 --- a/packages/meshbay-android/app/src/test/kotlin/org/meshbay/client/CastRelayTest.kt +++ b/packages/meshbay-android/app/src/test/kotlin/org/meshbay/client/CastRelayTest.kt @@ -17,9 +17,10 @@ import java.net.Socket /** The real relay on loopback, read with a raw socket: what a receiver sees. */ class CastRelayTest { - private val relay = CastRelay { InetAddress.getByName("127.0.0.1") } + private val spool = java.nio.file.Files.createTempDirectory("relay-spool").toFile() + private val relay = CastRelay({ InetAddress.getByName("127.0.0.1") }, spool) - @After fun stop() = relay.stop() + @After fun stop() { relay.stop(); spool.deleteRecursively() } private class Reply(val status: Int, val headers: Map, val body: ByteArray) @@ -166,6 +167,32 @@ class CastRelayTest { } } + @Test fun `a receiver far behind loses nothing, and its spool goes with it`() { + // Fragments of 4 MB, ten of them pushed while the receiver has read + // none: far past the 8 MB the desktop drops at, which froze the TV. + val url = relay.start(null, null).getString("url") + val u = java.net.URI(url) + Socket(u.host, u.port).use { s -> + s.getOutputStream().write("GET ${u.rawPath}?${u.rawQuery} HTTP/1.1\r\nHost: x\r\n\r\n".toByteArray()) + Thread.sleep(300) + val big = (0 until 10).map { box("moof", 16 + it) + box("mdat", 4 * 1024 * 1024 + it) } + big.forEach { relay.push(it) } + assertEquals("one spool file for the one receiver", 1, spool.listFiles()!!.size) + val input = DataInputStream(s.getInputStream()) + while (true) { val l = StringBuilder(); while (true) { val c = input.read(); if (c == '\n'.code) break; l.append(c.toChar()) }; if (l.toString().trim().isEmpty()) break } + val body = ByteArrayOutputStream() + val want = big.sumOf { it.size } + while (body.size() < want) { + val size = StringBuilder().also { sb -> while (true) { val c = input.read(); if (c == '\n'.code) break; sb.append(c.toChar()) } }.toString().trim().toInt(16) + val buf = ByteArray(size); input.readFully(buf); body.write(buf); input.read(); input.read() + } + assertArrayEquals("every fragment, in order, none dropped", big.reduce { a, b -> a + b }, body.toByteArray()) + } + Thread.sleep(300) + relay.stop() + assertEquals("the spool is deleted with the receiver", 0, spool.listFiles()!!.size) + } + @Test fun `stopping closes the port`() { val url = relay.start(null, null).getString("url") relay.stop() diff --git a/packages/meshbay-hub/tests/test_android_cast.py b/packages/meshbay-hub/tests/test_android_cast.py index 0ff383a..a346d15 100644 --- a/packages/meshbay-hub/tests/test_android_cast.py +++ b/packages/meshbay-hub/tests/test_android_cast.py @@ -53,7 +53,6 @@ def test_the_relay_keeps_the_desktop_mitigations(): for name in ("RING_CAP", "PORT_BASE", "PORT_COUNT"): js = re.search(rf"const {name} = (\d+);", desktop).group(1) assert re.search(rf"const val {name} = {js}\b", relay), name - assert "BACKPRESSURE_HIGH = 8L * 1024 * 1024" in relay assert "InetSocketAddress(address, PORT_BASE + i)" in relay, "bound to the LAN address" # Code, not comments: the comment saying "never 0.0.0.0" is not a bind. code = re.sub(r"/\*.*?\*/", "", relay, flags=re.S) @@ -65,6 +64,21 @@ def test_the_relay_keeps_the_desktop_mitigations(): assert serve.index('method == "OPTIONS"') < serve.index('query["t"] != token') +def test_what_a_receiver_has_not_read_waits_on_disk(): + """The desktop drops a fragment once 8 MB wait for a receiver; a fragment + is 5 to 10 MB at a film's bitrate, so the television froze for its length + (measured: three dropped in the first fifteen seconds). Here the lead waits + in a spool file per receiver, deleted with it, and is dropped only past a + disk bound.""" + relay = _read(CAST / "CastRelay.kt") + assert "BACKPRESSURE_HIGH" not in relay + assert "RandomAccessFile(file, \"rw\").channel" in relay + assert "c.waiting() > spoolLimit()" in relay + assert "SPOOL_MAX_BYTES = 2L * 1024 * 1024 * 1024" in relay + assert "file.delete()" in relay + assert 'java.io.File(context.cacheDir, "cast-relay")' in _read(CAST / "CastChannels.kt") + + def test_the_backlog_is_bounded_in_bytes(): """A fragment is a segment, megabytes at a film's bitrate: 64 of them killed the application with an OutOfMemoryError on the emulator.""" -- cgit v1.2.3