package org.meshbay.client import org.json.JSONObject import org.junit.After import org.junit.Assert.assertArrayEquals import org.junit.Assert.assertEquals import org.junit.Assert.assertNotEquals import org.junit.Assert.assertNull import org.junit.Assert.assertTrue import org.junit.Test import org.meshbay.client.cast.BoxAccumulator import org.meshbay.client.cast.CastRelay import java.io.ByteArrayOutputStream import java.io.DataInputStream import java.net.InetAddress import java.net.Socket /** The real relay on loopback, read with a raw socket: what a receiver sees. */ class CastRelayTest { 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(); spool.deleteRecursively() } private class Reply(val status: Int, val headers: Map, val body: ByteArray) private fun get(url: String, method: String = "GET", readBytes: Int = -1): Reply { val u = java.net.URI(url) Socket(u.host, u.port).use { s -> s.soTimeout = 5000 s.getOutputStream().write("$method ${u.rawPath}${u.rawQuery?.let { "?$it" } ?: ""} HTTP/1.1\r\nHost: x\r\n\r\n".toByteArray()) val input = DataInputStream(s.getInputStream()) val headLines = ArrayList() val line = StringBuilder() while (true) { val c = input.read() if (c == -1) break if (c == '\n'.code) { val l = line.toString().trimEnd('\r'); if (l.isEmpty()) break; headLines.add(l); line.clear() } else line.append(c.toChar()) } val status = headLines[0].split(' ')[1].toInt() val headers = headLines.drop(1).associate { it.substringBefore(':').lowercase() to it.substringAfter(':').trim() } val body = ByteArrayOutputStream() if (headers["transfer-encoding"] == "chunked") { while (readBytes < 0 || body.size() < readBytes) { val sizeLine = StringBuilder() while (true) { val c = input.read(); if (c == -1 || c == '\n'.code) break; sizeLine.append(c.toChar()) } val size = sizeLine.toString().trim().toIntOrNull(16) ?: break if (size == 0) break val buf = ByteArray(size); input.readFully(buf); body.write(buf); input.read(); input.read() } } else { val n = headers["content-length"]?.toInt() ?: 0 val buf = ByteArray(n); input.readFully(buf); body.write(buf) } return Reply(status, headers, body.toByteArray()) } } private fun box(type: String, payload: Int): ByteArray { val size = 8 + payload return byteArrayOf((size ushr 24).toByte(), (size ushr 16).toByte(), (size ushr 8).toByte(), size.toByte()) + type.toByteArray() + ByteArray(payload) { (it % 251).toByte() } } private fun fragment(n: Int) = box("moof", 16 + n) + box("mdat", 1000 + n) @Test fun `fragments are re-framed from arbitrary slices`() { val acc = BoxAccumulator() val stream = box("ftyp", 12) + fragment(1) + fragment(2) + fragment(3) val out = ArrayList() var i = 0 while (i < stream.size) { val n = minOf(37, stream.size - i); out += acc.push(stream, i, n); i += n } assertEquals(3, out.size) assertArrayEquals(fragment(2), out[1]) } @Test fun `a lost frame is recovered by rescanning for the next moof`() { val acc = BoxAccumulator() BoxAccumulator.warn = {} val out = acc.push(fragment(1) + byteArrayOf(0, 0, 0, 1, 1, 2, 3, 4) + fragment(2)) assertEquals(2, out.size) } @Test fun `the stream is init then the backlog then what follows`() { val init = box("ftyp", 20) + box("moov", 50) val started = relay.start(init, null) relay.push(fragment(1)); relay.push(fragment(2)) val reply = get(started.getString("url"), readBytes = init.size + fragment(1).size + fragment(2).size) assertEquals(200, reply.status) assertEquals("video/mp4", reply.headers["content-type"]) assertEquals("no-store", reply.headers["cache-control"]) assertArrayEquals(init + fragment(1) + fragment(2), reply.body) } @Test fun `a first chunk that carries film is served as its header only`() { // As the page delivers it from a real film: a 64 KB slice holding the // header, the first moof and the start of its mdat — pushed as well. val header = box("ftyp", 20) + box("moov", 1200) val stream = header + fragment(1) + fragment(2) val firstChunk = stream.copyOfRange(0, header.size + 300) val started = relay.start(firstChunk, null) var at = 0 while (at < stream.size) { val n = minOf(300, stream.size - at); relay.push(stream, at, n); at += n } val reply = get(started.getString("url"), readBytes = stream.size) 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=") assertTrue(Regex("^[0-9a-f]{32}$").matches(token)) assertEquals(403, get(url.replace(token, "0".repeat(32))).status) assertEquals(403, get(url.substringBefore("?")).status) assertEquals(405, get(url, method = "POST").status) assertEquals(404, get(url.replace("/stream.mp4", "/other")).status) } @Test fun `the subtitle is webvtt behind the token, readable cross-origin, re-addressed when it changes`() { val r = relay.start(null, JSONObject().put("vtt", "WEBVTT\n\n00:00.000 --> 00:01.000\nhi\n").put("language", "fr").put("label", "Français")) val sub = r.getJSONObject("subtitle") val reply = get(sub.getString("url")) assertEquals(200, reply.status) assertEquals("text/vtt; charset=utf-8", reply.headers["content-type"]) assertEquals("*", reply.headers["access-control-allow-origin"]) assertTrue(String(reply.body).startsWith("WEBVTT")) assertEquals(403, get(sub.getString("url").replace(Regex("t=[0-9a-f]+"), "t=x")).status) val second = relay.setSubtitle(JSONObject().put("vtt", "WEBVTT\n"))!! assertNotEquals(sub.getString("url"), second.getString("url")) assertNull(relay.setSubtitle(null)) assertEquals(404, get(second.getString("url")).status) assertEquals(200, get(r.getString("url"), readBytes = 0).status) } private fun send(bytes: ByteArray, type: String, meta: JSONObject = JSONObject(), cover: ByteArray? = null): String { relay.beginFile(type, bytes.size.toLong(), cover, cover?.let { "image/jpeg" }, meta) var at = 0 while (at < bytes.size) { val n = minOf(3, bytes.size - at); relay.writeFile(bytes, at, n); at += n } return relay.endFile().getString("url") } private fun getRange(url: String, range: String): Reply { val u = java.net.URI(url) Socket(u.host, u.port).use { s -> s.soTimeout = 5000 s.getOutputStream().write("GET ${u.rawPath}?${u.rawQuery} HTTP/1.1\r\nHost: x\r\nRange: $range\r\n\r\n".toByteArray()) val input = DataInputStream(s.getInputStream()) val lines = ArrayList() val line = StringBuilder() while (true) { val c = input.read() if (c == -1) break if (c == '\n'.code) { val l = line.toString().trimEnd('\r'); if (l.isEmpty()) break; lines.add(l); line.clear() } else line.append(c.toChar()) } val headers = lines.drop(1).associate { it.substringBefore(':').lowercase() to it.substringAfter(':').trim() } val body = ByteArray(headers["content-length"]?.toInt() ?: 0) input.readFully(body) return Reply(lines[0].split(' ')[1].toInt(), headers, body) } } @Test fun `a file starts the relay and is served behind the token, re-addressed when it changes`() { val jpeg = byteArrayOf(0xFF.toByte(), 0xD8.toByte(), 1, 2, 3, 0xFF.toByte(), 0xD9.toByte()) val first = send(jpeg, "image/jpeg") assertTrue(relay.active) val reply = get(first) assertEquals(200, reply.status) assertEquals("image/jpeg", reply.headers["content-type"]) assertEquals("no-store", reply.headers["cache-control"]) assertArrayEquals(jpeg, reply.body) assertEquals(403, get(first.replace(Regex("t=[0-9a-f]+"), "t=x")).status) // A receiver handed the same address twice shows what it already has. val second = send(byteArrayOf(9), "image/jpeg") assertNotEquals(first, second) assertArrayEquals(byteArrayOf(9), get(second).body) } @Test fun `a track is served with its cover, and by range`() { val track = ByteArray(20) { it.toByte() } val cover = byteArrayOf(0xFF.toByte(), 0xD8.toByte(), 7) val url = send(track, "audio/flac", JSONObject().put("title", "A Song").put("artist", "Some Band"), cover) assertEquals("audio/flac", relay.fileType) val meta = relay.fileMeta()!! assertEquals("A Song", meta.getString("title")) assertArrayEquals(cover, get(meta.getString("cover")).body) val part = getRange(url, "bytes=5-9") assertEquals(206, part.status) assertEquals("bytes 5-9/20", part.headers["content-range"]) assertArrayEquals(track.copyOfRange(5, 10), part.body) assertArrayEquals(track.copyOfRange(15, 20), getRange(url, "bytes=15-").body) assertEquals(416, getRange(url, "bytes=50-60").status) } @Test fun `only a whole picture or sound is served as one`() { val refused = try { relay.beginFile("text/html", 1, null, null, JSONObject()); false } catch (e: IllegalArgumentException) { true } assertTrue(refused) assertNull(relay.fileUrl) relay.beginFile("audio/mpeg", 10, null, null, JSONObject()) relay.writeFile(byteArrayOf(1, 2, 3)) val short = try { relay.endFile(); false } catch (e: IllegalStateException) { true } assertTrue(short) assertNull(relay.fileUrl) } @Test fun `a film started over a file takes the file down`() { send(byteArrayOf(1), "image/jpeg") relay.start(null, null) assertNull(relay.fileUrl) } @Test fun `the preflight is answered before the token is checked`() { val url = relay.start(null, null).getString("url") val reply = get(url.substringBefore("?"), method = "OPTIONS") assertEquals(204, reply.status) assertEquals("GET, OPTIONS", reply.headers["access-control-allow-methods"]) assertTrue(reply.headers["access-control-allow-headers"]!!.contains("Range")) } @Test fun `the stream carries the same CORS headers as its subtitle`() { val url = relay.start(null, null).getString("url") val reply = get(url, readBytes = 0) for ((k, v) in CastRelay.CORS_HEADERS) assertEquals(k, v, reply.headers[k.lowercase()]) } @Test fun `the backlog is bounded in bytes, not only in fragments`() { relay.start(null, null) // 40 fragments of 4 MB: under the fragment cap, far over a phone's heap. val big = box("moof", 16) + box("mdat", 4 * 1024 * 1024) repeat(40) { relay.push(big) } assertTrue("backlog ${relay.backlogBytes()}", relay.backlogBytes() <= CastRelay.RING_MAX_BYTES) assertTrue(relay.backlogBytes() >= CastRelay.RING_MAX_BYTES - big.size) } @Test fun `seek after seek, the relay restarts on its ports`() { // A seek restarts the relay, and a receiver was connected each time: // the ports it closed sit in TIME_WAIT. repeat(8) { val url = relay.start(null, null).getString("url") relay.push(fragment(it)) get(url, readBytes = fragment(it).size) relay.stop() } } @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() val u = java.net.URI(url) val refused = try { Socket(u.host, u.port).close(); false } catch (e: java.net.ConnectException) { true } assertTrue(refused) assertNull(relay.url) } }