Demo.kt
| 1 | package de.upapp |
| 2 | |
| 3 | import org.bouncycastle.tls.CipherSuite |
| 4 | import org.bouncycastle.tls.PSKTlsServer |
| 5 | import org.bouncycastle.tls.ProtocolVersion |
| 6 | import org.bouncycastle.tls.TlsPSKIdentityManager |
| 7 | import org.bouncycastle.tls.TlsServerProtocol |
| 8 | import org.bouncycastle.tls.crypto.impl.bc.BcTlsCrypto |
| 9 | import org.json.JSONArray |
| 10 | import org.json.JSONObject |
| 11 | import java.io.DataInputStream |
| 12 | import java.io.EOFException |
| 13 | import java.io.IOException |
| 14 | import java.io.InputStream |
| 15 | import java.io.OutputStream |
| 16 | import java.security.SecureRandom |
| 17 | import java.time.LocalDateTime |
| 18 | import java.time.format.DateTimeFormatter |
| 19 | import java.util.Locale |
| 20 | import java.util.UUID |
| 21 | import java.util.concurrent.ConcurrentHashMap |
| 22 | import java.util.concurrent.CopyOnWriteArrayList |
| 23 | import kotlin.concurrent.thread |
| 24 | |
| 25 | const val DEMO_VIN = "WVWZZZAAZLD000042" |
| 26 | const val DEMO_CREDENTIALS = "demo,demo" |
| 27 | |
| 28 | /** Car end of a tunnel, for the demo and tests. [serve] handles each connection that the phone opens. */ |
| 29 | class CarSide(serve: (port: Int, InputStream, OutputStream) -> Unit) { |
| 30 | private val toPhone = Pipe() |
| 31 | private val toCar = Pipe() |
| 32 | val tunnel = Tunnel(toPhone.input, toCar.output) |
| 33 | |
| 34 | init { |
| 35 | thread(isDaemon = true, name = "car") { |
| 36 | val pipes = ConcurrentHashMap<Int, Pipe>() |
| 37 | val chunk = ByteArray(4096) |
| 38 | var buf = ByteArray(0) |
| 39 | while (true) { |
| 40 | val n = toCar.input.read(chunk) |
| 41 | if (n < 0) break |
| 42 | buf += chunk.copyOf(n) |
| 43 | while (true) { |
| 44 | val (s, size) = decode(buf) ?: break |
| 45 | buf = buf.copyOfRange(size, buf.size) |
| 46 | when (s.type) { |
| 47 | HEARTBEAT -> toPhone.put(frame(0, HEARTBEAT)) |
| 48 | OPEN -> { |
| 49 | val pipe = Pipe().also { pipes[s.port] = it } |
| 50 | val remote = (s.data[0].toInt() and 0xFF) or ((s.data[1].toInt() and 0xFF) shl 8) |
| 51 | val out = object : OutputStream() { |
| 52 | override fun write(b: Int) = write(byteArrayOf(b.toByte()), 0, 1) |
| 53 | override fun write(b: ByteArray, off: Int, len: Int) = |
| 54 | toPhone.put(frame(s.port, TRANSMIT, u16(len) + b.copyOfRange(off, off + len))) |
| 55 | } |
| 56 | thread(isDaemon = true) { |
| 57 | try { |
| 58 | serve(remote, pipe.input, out) |
| 59 | } catch (_: Exception) { |
| 60 | } |
| 61 | if (pipes.remove(s.port) != null) toPhone.put(frame(s.port, CLOSE)) |
| 62 | } |
| 63 | } |
| 64 | TRANSMIT -> pipes[s.port]?.put(s.data) |
| 65 | CLOSE -> pipes.remove(s.port)?.end() ?: Unit |
| 66 | } |
| 67 | } |
| 68 | } |
| 69 | pipes.values.forEach { it.end() } |
| 70 | } |
| 71 | } |
| 72 | } |
| 73 | |
| 74 | private class Request(val method: String, val path: String, val headers: Map<String, String>, val body: String) |
| 75 | |
| 76 | private fun readRequest(input: InputStream): Request? { |
| 77 | val first = try { |
| 78 | input.line() |
| 79 | } catch (_: EOFException) { |
| 80 | return null |
| 81 | } |
| 82 | val (method, path) = first.split(' ') |
| 83 | val headers = generateSequence { input.line().ifEmpty { null } } |
| 84 | .associate { it.substringBefore(':').trim().lowercase() to it.substringAfter(':').trim() } |
| 85 | val body = ByteArray(headers["content-length"]?.toInt() ?: 0).also { DataInputStream(input).readFully(it) } |
| 86 | return Request(method, path, headers, body.decodeToString()) |
| 87 | } |
| 88 | |
| 89 | private fun respond(out: OutputStream, status: Int, body: Any? = null, headers: List<String> = emptyList()) { |
| 90 | val text = body?.toString().orEmpty().toByteArray() |
| 91 | val head = (listOf("HTTP/1.1 $status X", "Content-Length: ${text.size}") + headers).joinToString("\r\n") + "\r\n\r\n" |
| 92 | out.write(head.toByteArray() + text) |
| 93 | out.flush() |
| 94 | } |
| 95 | |
| 96 | /** A simulated e-up! that serves the REST API and events in memory. */ |
| 97 | class DemoCar { |
| 98 | private val docs = LinkedHashMap<String, JSONObject>() |
| 99 | private val listeners = CopyOnWriteArrayList<OutputStream>() |
| 100 | |
| 101 | val tunnel = CarSide(::serve).tunnel |
| 102 | |
| 103 | init { |
| 104 | val now = LocalDateTime.now() |
| 105 | add("/car/info/", "info", "vehicleIdenticationNumber" to DEMO_VIN, "vehicleType" to "VW120", "language" to "de", |
| 106 | "vehicleDate" to now.format(DateTimeFormatter.ofPattern("dd MMM yyyy", Locale.ENGLISH)), |
| 107 | "vehicleTime" to now.format(DateTimeFormatter.ofPattern("HH:mm:ss"))) |
| 108 | add("/car/batteries/", "BEV battery", "soc" to 64.0, "socUnit" to "percent") |
| 109 | add("/car/ranges/", "BEV range", "value" to 151.0, "valueUnit" to "km") |
| 110 | add("/car/environments/", "environment", "outsideTemperature" to 12.5, "outsideTemperatureUnit" to "celsius", "lightIntensity" to 40) |
| 111 | add("/chargingmanager/batteryCharges/", "batteryCharge", "state" to "running") |
| 112 | add("/chargingmanager/batteryPlugs/", "batteryPlug", "plug" to "plugged") |
| 113 | add("/chargingmanager/batteryClimates/", "batteryClimate", "state" to "idle") |
| 114 | add("/chargingmanager/profiles/", "Optionen", "operations" to JSONArray(listOf("climateExtSupply", "climate")), |
| 115 | "maxCurrent" to 16, "minLevel" to 30, "targetLevel" to 0, "temperature" to 21.5) |
| 116 | val home = createProfile("Zuhause", 80, 16, listOf("charge", "climate")) |
| 117 | val work = createProfile("Arbeit", 100, 10, listOf("charge")) |
| 118 | add("/chargingmanager/timers/", "Timer1", "state" to TIMER_ON, "departureTime" to "05:30:00", "cyclic" to true, |
| 119 | "weekdays" to JSONArray(listOf("monday", "tuesday", "wednesday", "thursday", "friday")), "profile" to link(home)) |
| 120 | add("/chargingmanager/timers/", "Timer2", "state" to TIMER_OFF, "departureTime" to "08:00:00", "cyclic" to true, |
| 121 | "weekdays" to JSONArray(listOf("saturday")), "profile" to link(work)) |
| 122 | add("/chargingmanager/timers/", "Timer3", "state" to TIMER_OFF, "departureTime" to "16:45:00", "cyclic" to false, |
| 123 | "departureDate" to now.plusDays(3).format(DateTimeFormatter.ofPattern("dd MMM yyyy", Locale.ENGLISH)), |
| 124 | "weekdays" to JSONArray(), "profile" to link(home)) |
| 125 | thread(isDaemon = true, name = "demo charging") { simulateCharging() } |
| 126 | } |
| 127 | |
| 128 | private fun add(collection: String, name: String, vararg fields: Pair<String, Any>): JSONObject { |
| 129 | val id = UUID.randomUUID().toString() |
| 130 | val doc = JSONObject().put("id", id).put("name", name).put("uri", collection + id) |
| 131 | fields.forEach { (k, v) -> doc.put(k, v) } |
| 132 | docs[doc.getString("uri")] = doc |
| 133 | return doc |
| 134 | } |
| 135 | |
| 136 | private fun createProfile(name: String, target: Int, current: Int, operations: List<String>): JSONObject { |
| 137 | val provider = add("/chargingmanager/providers/", "", "cyclic" to true, "preferredTimeStart" to "21:00:00", |
| 138 | "preferredTimeEnd" to "05:00:00", "weekdays" to JSONArray(listOf("monday", "tuesday", "wednesday", "thursday", "friday", "saturday", "sunday"))) |
| 139 | return add("/chargingmanager/profiles/", name, "operations" to JSONArray(operations), "maxCurrent" to current, |
| 140 | "minLevel" to 0, "targetLevel" to target, "temperature" to 21.0, "powerProvider" to link(provider)) |
| 141 | } |
| 142 | |
| 143 | private fun link(doc: JSONObject) = JSONObject().put("id", doc["id"]).put("name", doc["name"]).put("uri", doc["uri"]) |
| 144 | |
| 145 | private fun collection(path: String) = JSONArray(docs.filterKeys { it.startsWith(path) && it != path }.values) |
| 146 | |
| 147 | private fun first(collection: String) = docs.entries.first { it.key.startsWith(collection) }.value |
| 148 | |
| 149 | private fun simulateCharging() { |
| 150 | while (true) { |
| 151 | Thread.sleep(3_000) |
| 152 | synchronized(this) { |
| 153 | val charge = first("/chargingmanager/batteryCharges/") |
| 154 | if (charge["state"] != "running") return@synchronized |
| 155 | val battery = first("/car/batteries/") |
| 156 | val soc = battery.getDouble("soc") + 1 |
| 157 | battery.put("soc", soc) |
| 158 | first("/car/ranges/").put("value", (soc * 2.4).toInt().toDouble()) |
| 159 | if (soc >= 80) charge.put("state", "completed") |
| 160 | listOf("/car/batteries/", "/car/ranges/", "/chargingmanager/batteryCharges/").forEach(::notify) |
| 161 | } |
| 162 | } |
| 163 | } |
| 164 | |
| 165 | private fun notify(collection: String) { |
| 166 | val message = JSONObject().put("type", "data").put("event", "$collection#42").put("data", collection(collection)) |
| 167 | listeners.forEach { out -> |
| 168 | runCatching { writeFrame(out, WS_TEXT, message.toString().toByteArray(), mask = false) }.onFailure { listeners.remove(out) } |
| 169 | } |
| 170 | } |
| 171 | |
| 172 | private fun serve(port: Int, input: InputStream, output: OutputStream) { |
| 173 | if (port == 80) { |
| 174 | val r = readRequest(input) ?: return |
| 175 | if (r.path == "/car/info/vin") respond(output, 200, DEMO_VIN) else respond(output, 404) |
| 176 | return |
| 177 | } |
| 178 | val tls = TlsServerProtocol(input, output) |
| 179 | tls.accept(object : PSKTlsServer(BcTlsCrypto(SecureRandom()), object : TlsPSKIdentityManager { |
| 180 | override fun getHint(): ByteArray = DEMO_VIN.toByteArray() |
| 181 | override fun getPSK(identity: ByteArray) = DEMO_CREDENTIALS.substringAfter(',').toByteArray() |
| 182 | }) { |
| 183 | override fun getSupportedVersions(): Array<ProtocolVersion> = ProtocolVersion.TLSv12.only() |
| 184 | override fun getSupportedCipherSuites() = intArrayOf(CipherSuite.TLS_PSK_WITH_AES_128_CBC_SHA256) |
| 185 | }) |
| 186 | val inp = tls.inputStream |
| 187 | val out = tls.outputStream |
| 188 | while (true) { |
| 189 | val r = readRequest(inp) ?: break |
| 190 | if (r.headers["upgrade"].equals("websocket", ignoreCase = true)) { |
| 191 | respond(out, 101, null, listOf("Upgrade: websocket", "Connection: Upgrade", "Sec-WebSocket-Accept: ${wsAccept(r.headers.getValue("sec-websocket-key"))}")) |
| 192 | events(inp, out) |
| 193 | break |
| 194 | } |
| 195 | synchronized(this) { handle(r, out) } |
| 196 | } |
| 197 | tls.close() |
| 198 | } |
| 199 | |
| 200 | private fun events(input: InputStream, output: OutputStream) { |
| 201 | listeners.add(output) |
| 202 | try { |
| 203 | while (true) { |
| 204 | val (_, data) = readMessage(input, output, mask = false) |
| 205 | val m = JSONObject(data.decodeToString()) |
| 206 | val reply = JSONObject().put("type", m.getString("type")).put("event", m.getString("event")).put("status", "ok") |
| 207 | writeFrame(output, WS_TEXT, reply.toString().toByteArray(), mask = false) |
| 208 | } |
| 209 | } finally { |
| 210 | listeners.remove(output) |
| 211 | } |
| 212 | } |
| 213 | |
| 214 | private fun handle(r: Request, out: OutputStream) { |
| 215 | val path = r.path |
| 216 | val doc = docs[path] |
| 217 | when { |
| 218 | r.method == "GET" && path == "/" -> respond(out, 200, JSONObject()) |
| 219 | r.method == "GET" && doc != null -> respond(out, 200, JSONObject().put("status", "ok").put("data", doc)) |
| 220 | r.method == "GET" && path.endsWith("/") && docs.keys.any { it.startsWith(path) } -> respond(out, 200, JSONObject().put("data", collection(path))) |
| 221 | r.method == "POST" && path == "/chargingmanager/profiles/" -> { |
| 222 | val b = JSONObject(r.body) |
| 223 | val created = createProfile(b.getString("name"), b.getInt("targetLevel"), b.getInt("maxCurrent"), |
| 224 | b.getJSONArray("operations").let { a -> List(a.length()) { a.getString(it) } }) |
| 225 | respond(out, 201, created) |
| 226 | notify("/chargingmanager/profiles/") |
| 227 | notify("/chargingmanager/providers/") |
| 228 | } |
| 229 | r.method == "POST" && doc != null -> { |
| 230 | val b = JSONObject(r.body) |
| 231 | b.keys().forEach { k -> |
| 232 | val v = b.get(k) |
| 233 | val target = (v as? String)?.let { id -> docs.values.firstOrNull { it.optString("id") == id } } |
| 234 | doc.put(k, if (k in setOf("profile", "powerProvider") && target != null) link(target) else v) |
| 235 | } |
| 236 | respond(out, 200, JSONObject().put("status", "ok")) |
| 237 | notify(path.substringBeforeLast('/') + "/") |
| 238 | } |
| 239 | r.method == "DELETE" && doc != null -> { |
| 240 | docs.remove(path) |
| 241 | doc.optJSONObject("powerProvider")?.let { docs.remove(it.getString("uri")) } |
| 242 | respond(out, 200) |
| 243 | notify("/chargingmanager/profiles/") |
| 244 | notify("/chargingmanager/providers/") |
| 245 | } |
| 246 | else -> respond(out, 404) |
| 247 | } |
| 248 | } |
| 249 | } |
| 250 |