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 +++++++++++++++------ 2 files changed, 100 insertions(+), 42 deletions(-) (limited to 'packages/meshbay-android/app/src/main/kotlin/org/meshbay/client') 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()) } } -- cgit v1.2.3