1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
|
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<String, String>, 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<String>()
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<ByteArray>()
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<String>()
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)
}
}
|