aboutsummaryrefslogtreecommitdiffstats
path: root/packages/meshbay-android/app/src/main/kotlin
diff options
context:
space:
mode:
authorChristophe Besson <cbesson@gmail.com>2026-10-03 12:53:35 +0200
committerChristophe Besson <cbesson@gmail.com>2026-10-03 14:24:54 +0200
commitcf8345386996c561f176049b2f83fbd07956995e (patch)
treeda50f4f3329d519f06b8bf06483b15fa2ae9ca64 /packages/meshbay-android/app/src/main/kotlin
parente7d663b76e58354d6a1b58f40624427d5ae660c8 (diff)
downloadmeshbay-cf8345386996c561f176049b2f83fbd07956995e.tar.gz
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 <noreply@anthropic.com>
Diffstat (limited to 'packages/meshbay-android/app/src/main/kotlin')
-rw-r--r--packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastChannels.kt2
-rw-r--r--packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/CastRelay.kt138
2 files changed, 99 insertions, 41 deletions
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?) {
+class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir: File) {
- private class Client(val socket: Socket, val out: OutputStream) {
- val queue = LinkedBlockingDeque<ByteArray>()
- val queued = AtomicLong()
+ /**
+ * 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())
}
}