package de.upapp import android.util.Log import org.json.JSONObject import org.json.JSONTokener import java.io.ByteArrayOutputStream import java.io.DataInputStream import java.io.EOFException import java.io.IOException import java.io.InputStream import java.io.OutputStream import java.security.MessageDigest import java.security.SecureRandom import java.util.Base64 internal const val WS_TEXT = 1 internal const val WS_CLOSE = 8 internal const val WS_PING = 9 internal const val WS_PONG = 10 internal fun wsAccept(key: String): String = Base64.getEncoder().encodeToString( MessageDigest.getInstance("SHA-1").digest((key + "258EAFA5-E914-47DA-95CA-C5AB0DC85B11").toByteArray()) ) /** Writes one unfragmented frame. Clients must [mask], servers must not. */ internal fun writeFrame(out: OutputStream, opcode: Int, payload: ByteArray, mask: Boolean) = synchronized(out) { val frame = ByteArrayOutputStream() frame.write(0x80 or opcode) val m = if (mask) 0x80 else 0 val n = payload.size when { n < 126 -> frame.write(m or n) n < 0x10000 -> frame.write(byteArrayOf((m or 126).toByte(), (n shr 8).toByte(), n.toByte())) else -> { frame.write(m or 127) for (i in 7 downTo 0) frame.write((n.toLong() shr (8 * i)).toInt()) } } if (mask) { val key = ByteArray(4).also { SecureRandom().nextBytes(it) } frame.write(key) frame.write(ByteArray(n) { (payload[it].toInt() xor key[it % 4].toInt()).toByte() }) } else { frame.write(payload) } out.write(frame.toByteArray()) out.flush() } /** Reads the next data message. Answers pings. Throws [EOFException] on a close frame. */ internal fun readMessage(input: InputStream, output: OutputStream, mask: Boolean): Pair { val inp = DataInputStream(input) var opcode = 0 val message = ByteArrayOutputStream() while (true) { val b0 = inp.readUnsignedByte() val b1 = inp.readUnsignedByte() val len = when (val l = b1 and 0x7F) { 126 -> inp.readUnsignedShort() 127 -> inp.readLong().toInt() else -> l } val key = if (b1 and 0x80 != 0) ByteArray(4).also { inp.readFully(it) } else null val data = ByteArray(len).also { inp.readFully(it) } key?.let { k -> for (i in data.indices) data[i] = (data[i].toInt() xor k[i % 4].toInt()).toByte() } when (val op = b0 and 0x0F) { WS_CLOSE -> { runCatching { writeFrame(output, WS_CLOSE, data, mask) } throw EOFException("WebSocket closed") } WS_PING -> writeFrame(output, WS_PONG, data, mask) WS_PONG -> {} else -> { if (op != 0) opcode = op message.write(data) if (b0 and 0x80 != 0) return opcode to message.toByteArray() } } } } /** * Subscribes to [paths], collections or single elements, and calls [onEvent] with the subscribed path and the event data. * Blocks until the connection ends, then throws. */ fun Car.listen(paths: List, onEvent: (String, Any?) -> Unit) { // Events can be minutes apart, so this channel has no read timeout. val ch = tunnel.open(443, Long.MAX_VALUE) try { val tls = sessionTls(ch.input, ch.output, psk.first, psk.second) val key = Base64.getEncoder().encodeToString(ByteArray(16).also { SecureRandom().nextBytes(it) }) val r = exchange( tls.inputStream, tls.outputStream, "GET", HTTPS_HOST, "/", keepAlive = true, extraHeaders = listOf("Upgrade: websocket", "Connection: Upgrade", "Sec-WebSocket-Key: $key", "Sec-WebSocket-Version: 13"), ) if (r.status != 101 || r.headers["sec-websocket-accept"] != wsAccept(key)) throw IOException("WebSocket upgrade failed: HTTP ${r.status}") // Like maps+more: one subscription at a time. The car answers 409 Conflict to overlapping ones. val pending = ArrayDeque(paths) fun subscribeNext() = pending.removeFirstOrNull()?.let { // Inputs need a short limit. Otherwise the car merges "pushed" and "released" into one event. val limit = if (it.startsWith("/mechanicalinput/")) 100 else 1000 val subscribe = JSONObject().put("type", "subscribe").put("event", "$it#42").put("updatelimit", limit) writeFrame(tls.outputStream, WS_TEXT, subscribe.toString().toByteArray(), mask = true) } subscribeNext() while (true) { val (op, data) = readMessage(tls.inputStream, tls.outputStream, mask = true) if (op != WS_TEXT) continue val text = data.decodeToString() if (BuildConfig.DEBUG) Log.d("Events", text.take(4000)) // The car joins several messages in one frame without a separator. val messages = JSONTokener(text) while (true) { val m = (runCatching { messages.nextValue() }.getOrNull() as? JSONObject) ?: break val line = m.toString() val event = m.optString("event").substringBefore('#').substringBefore('?') when (m.optString("type")) { "data" -> (paths.firstOrNull { event == it } ?: paths.firstOrNull { event.startsWith(it) })?.let { onEvent(it, m.opt("data")) } "subscribe" -> { if (m.optString("status") != "ok") Log.w("Events", "Subscription failed: $line") subscribeNext() } "error" -> { Log.w("Events", "Event error: $line") subscribeNext() } } } } } finally { ch.close() } }