diff options
Diffstat (limited to 'packages')
3 files changed, 74 insertions, 3 deletions
diff --git a/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/BoxAccumulator.kt b/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/BoxAccumulator.kt index 3aaf8d2..f56d58a 100644 --- a/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/BoxAccumulator.kt +++ b/packages/meshbay-android/app/src/main/kotlin/org/meshbay/client/cast/BoxAccumulator.kt @@ -11,6 +11,14 @@ import java.io.ByteArrayOutputStream class BoxAccumulator { private var buf = ByteArray(0) private var synced = false + private var preambleSeen = false + + /** + * Called once, with every byte before the first moof: the stream's header + * (ftyp, moov), however many pushes it came in. The node's first chunk can + * hold the 28-byte ftyp alone, with the moov in the next one (measured). + */ + var onPreamble: ((ByteArray) -> Unit)? = null fun push(data: ByteArray, offset: Int = 0, length: Int = data.size - offset): List<ByteArray> { buf = buf + data.copyOfRange(offset, offset + length) @@ -19,6 +27,7 @@ class BoxAccumulator { if (!synced) { val idx = findMoof() if (idx == -1) return fragments + if (!preambleSeen) { preambleSeen = true; onPreamble?.invoke(buf.copyOfRange(0, idx)) } buf = buf.copyOfRange(idx, buf.size) synced = true } @@ -55,7 +64,7 @@ class BoxAccumulator { return fragments } - fun reset() { buf = ByteArray(0); synced = false } + fun reset() { buf = ByteArray(0); synced = false; preambleSeen = false } private fun findMoof(): Int { for (i in 0..buf.size - 8) { 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 4d52052..21328db 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 @@ -88,6 +88,8 @@ class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir @Volatile private var token: String? = null @Volatile private var host: String? = null @Volatile private var initSegment: ByteArray? = null + private val headerLock = Object() + @Volatile private var headerReady = false @Volatile private var subtitle: Subtitle? = null @Volatile private var subtitleVersion = 0 /** How many times a receiver asked for the stream since the relay started. */ @@ -134,9 +136,11 @@ class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir fragments = 0 host = address.hostAddress initSegment = init?.let(::headerOnly) + // Nothing to wait for without a first chunk: no header is coming. + headerReady = init == null synchronized(ring) { ring.clear(); ringBytes = 0 } clients.clear() - accum = BoxAccumulator() + accum = BoxAccumulator().also { a -> a.onPreamble = ::preamble } // What a killed process left behind. spoolDir.mkdirs() spoolDir.listFiles()?.forEach { it.delete() } @@ -166,6 +170,37 @@ class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir .put("subtitle", subtitleInfo() ?: JSONObject.NULL) } + /** + * The header is everything the page pushed before the first moof. The + * `init` the page hands to start() is only its first chunk, which may be + * the ftyp without the moov: served as the header, the receiver got no + * track description and gave up — the "fails the first time" of a fresh + * start, every time. The pushed stream carries the whole of it; `init` is + * the fallback when nothing came before the first moof. + */ + private fun preamble(bytes: ByteArray) { + val header = headerOnly(bytes) + synchronized(headerLock) { + if (header.size >= 8 && BoxAccumulator.u32(header, 4) == FTYP) initSegment = header + headerReady = true + headerLock.notifyAll() + } + log("header complete: ${initSegment?.size ?: 0} bytes") + } + + /** A receiver that connects before the first moof waits for the whole header, not a piece of it. */ + private fun awaitHeader(): ByteArray? { + val deadline = System.currentTimeMillis() + HEADER_WAIT_MS + synchronized(headerLock) { + while (!headerReady) { + val left = deadline - System.currentTimeMillis() + if (left <= 0) break + headerLock.wait(left) + } + } + return initSegment + } + @Volatile private var pushes = 0 @Volatile private var fragments = 0 @@ -283,7 +318,7 @@ class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir "Connection" to "keep-alive", "Transfer-Encoding" to "chunked", )) - initSegment?.let { chunk(out, it) } + awaitHeader()?.let { chunk(out, it) } out.flush() val client = Client(socket, out, File(spoolDir, "client-${spoolSeq.incrementAndGet()}.spool")) val backlog = synchronized(ring) { @@ -360,6 +395,8 @@ class CastRelay(private val lanAddress: () -> InetAddress?, private val spoolDir const val RING_CAP = 64 const val RING_MAX_BYTES = 32L * 1024 * 1024 const val SPOOL_MAX_BYTES = 2L * 1024 * 1024 * 1024 + private const val HEADER_WAIT_MS = 10_000L + private const val FTYP = 0x66747970L private const val BLOCK = 256 * 1024 const val PORT_BASE = 19550 const val PORT_COUNT = 4 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 1c604d0..5d3125f 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 @@ -106,6 +106,31 @@ class CastRelayTest { assertArrayEquals("header, then each fragment once", stream, reply.body) } + @Test fun `a header split across chunks is served whole`() { + // Measured on a fresh start: the first chunk was the 28-byte ftyp + // alone, the moov came in the next push. Served as the header, the + // receiver had no moov and gave up — the first cast failed every time. + val ftyp = box("ftyp", 20) + val moov = box("moov", 2124) + val started = relay.start(ftyp, null) + val url = started.getString("url") + // The page pushes the first chunk too, then the rest. + relay.push(ftyp); relay.push(moov); relay.push(fragment(1)); relay.push(fragment(2)) + val reply = get(url, readBytes = ftyp.size + moov.size + fragment(1).size + fragment(2).size) + assertArrayEquals(ftyp + moov + fragment(1) + fragment(2), reply.body) + } + + @Test fun `a receiver early for the header waits for all of it`() { + val ftyp = box("ftyp", 20) + val moov = box("moov", 2124) + val url = relay.start(ftyp, null).getString("url") + relay.push(ftyp) + val late = Thread { Thread.sleep(400); relay.push(moov); relay.push(fragment(1)) }.apply { start() } + val reply = get(url, readBytes = ftyp.size + moov.size + fragment(1).size) + late.join() + assertArrayEquals(ftyp + moov + fragment(1), reply.body) + } + @Test fun `the token is required and unguessable`() { val url = relay.start(null, null).getString("url") val token = url.substringAfter("t=") |