From df92d3a55d00bb3ae7392f99564cd53412e24082 Mon Sep 17 00:00:00 2001 From: Andreas Rumpf Date: Wed, 10 May 2017 09:56:54 +0200 Subject: [PATCH] first version of transactor (todo-app working) --- src/karaxdb/client.nim | 35 +++------ src/karaxdb/common.nim | 2 +- src/karaxdb/server.nim | 56 -------------- src/karaxdb/transactor.nim | 151 +++++++++++++++++++++++++++++++++++++ 4 files changed, 161 insertions(+), 83 deletions(-) delete mode 100644 src/karaxdb/server.nim create mode 100644 src/karaxdb/transactor.nim diff --git a/src/karaxdb/client.nim b/src/karaxdb/client.nim index 7ea648b..72c31bc 100644 --- a/src/karaxdb/client.nim +++ b/src/karaxdb/client.nim @@ -21,9 +21,7 @@ type # RequestMessage {.importc.} = ref object let conn = newWebSocket("ws://localhost:8080", "karaxdb") -var gid: MessageId -var pendingIds = newJDict[MessageId, Message]() -var pending: int +var version: int #proc loadDb*(url: cstring): Db = # result = nil @@ -34,11 +32,9 @@ proc newTransaction*(): Db = proc insert*(head, newdb: Db) = newdb.next = head #result = newdb - let expectedVersion = if not head.isNil: head.version + 1 else: 1 - inc gid.int - let m = Message(kind: NewData, data: newdb.data, version: expectedVersion, id: gid) - pendingIds[gid] = m - inc pending + let expectedVersion = version + inc version + let m = Message(kind: NewData, data: newdb.data, version: expectedVersion, id: MessageId(0)) conn.send(toJson(m)) proc merge*(newer, older: Db) = @@ -49,28 +45,15 @@ proc registerOnUpdate*(update: proc(db: Db)) = proc (e: MessageEvent) = let msg = fromJson[Message](e.data) case msg.kind - of Conflict: + of Rejected: # conflict, so throw away the sent data, don't apply the changes: - if pending > 0: - pendingIds.del msg.id - dec pending - # server sent data that caused the conflict: - if msg.data.len > 0: - let db = Db(data: msg.data, version: msg.version) - update(db) - of Accepted: - kout cstring"accepted", pending - if pending > 0: - let d = pendingIds[msg.id] - # data was not submitted again, so we use the in-memory version - # of the data: - let db = Db(data: d.data, version: d.version) - pendingIds.del msg.id - dec pending - update(db) + kout cstring"rejected" of Newdata: let db = Db(data: msg.data, version: msg.version) update(db) + of Disconnect: + kout cstring"disconnected" + else: kout cstring"something else" iterator list*(db: Db, q: Query): Triple = # XXX here the datamodel comes in! We must not diff --git a/src/karaxdb/common.nim b/src/karaxdb/common.nim index 67cdf4c..c098702 100644 --- a/src/karaxdb/common.nim +++ b/src/karaxdb/common.nim @@ -14,4 +14,4 @@ type MessageId* = distinct int MessageKind* = enum - Newdata, Conflict, Accepted + Newdata, Rejected, Disconnect diff --git a/src/karaxdb/server.nim b/src/karaxdb/server.nim deleted file mode 100644 index c0291d4..0000000 --- a/src/karaxdb/server.nim +++ /dev/null @@ -1,56 +0,0 @@ - -import asynchttpserver, asyncdispatch, asyncnet, "../../../websocket/websocket", common, json - -type - Message = object - kind: MessageKind - id: MessageId - data: seq[Triple] - version: int - -proc `%`(id: MessageId): JsonNode = %BiggestInt(id) -proc `%`(k: MessageKind): JsonNode = %BiggestInt(k) - -proc triplesFromJson(j: JsonNode): seq[Triple] = - result = newSeq[Triple](j.len) - var i = 0 - for t in j: - doAssert t.kind == JArray - result[i] = [t[0].str, t[1].str, t[2].str] - inc i - -proc messageFromJson(j: JsonNode): Message = - Message(kind: MessageKind(j["kind"].num), id: MessageId(j["id"].num), - data: triplesFromJson(j["data"]), version: j["version"].num.int) - -var server = newAsyncHttpServer() - -proc cb(req: Request) {.async.} = - let (success, error) = await(verifyWebsocketRequest(req, "karaxdb")) - if not success: - echo "WS negotiation failed: " & error - await req.respond(Http400, "Websocket negotiation failed: " & error) - req.client.close - else: - echo "New websocket customer arrived!" - while true: - try: - var f = await req.client.readData(false) - echo "(opcode: " & $f.opcode & ", data: " & $f.data.len & ")" - let m = messageFromJson(f.data.parseJson) - let om = Message(kind: Accepted, id: m.id, data: @[], version: m.version) - let oms = $(%*om) - echo "OUTPUT ", oms - if f.opcode == Opcode.Text: - waitFor req.client.sendText(oms, false) - else: - echo "protocol error" - #waitFor req.client.sendBinary(f.data, false) - except: - echo getCurrentExceptionMsg() - break - - req.client.close() - echo ".. socket went away." - -waitFor server.serve(Port(8080), cb) diff --git a/src/karaxdb/transactor.nim b/src/karaxdb/transactor.nim new file mode 100644 index 0000000..057884d --- /dev/null +++ b/src/karaxdb/transactor.nim @@ -0,0 +1,151 @@ + +import asynchttpserver, asyncdispatch, asyncnet, "../../../websocket/websocket", common, json, + strutils, times + +type + Message = object + kind: MessageKind + id: MessageId + data: seq[Triple] + version: int + +proc `%`(id: MessageId): JsonNode = %BiggestInt(id) +proc `%`(k: MessageKind): JsonNode = %BiggestInt(k) + +proc triplesFromJson(j: JsonNode): seq[Triple] = + result = newSeq[Triple](j.len) + var i = 0 + for t in j: + doAssert t.kind == JArray + let val = if t[2].kind == JNull: string(nil) else: t[2].str + result[i] = [t[0].str, t[1].str, val] + inc i + +proc messageFromJson(j: JsonNode): Message = + Message(kind: MessageKind(j["kind"].num), id: MessageId(j["id"].num), + data: triplesFromJson(j["data"]), version: j["version"].num.int) + +proc error(msg: string) = echo msg +proc warn(msg: string) = echo msg + +type + Tx = object + data: string + version: int + Client = ref object + socket: AsyncSocket + connected: bool + hostname: string + lastMessage: float + rapidMessageCount: int + + Server = ref object + clients: seq[Client] + needsUpdate: bool + txs: seq[Tx] + version: int + +proc newClient(socket: AsyncSocket, hostname: string): Client = + Client(socket: socket, connected: true, hostname: hostname) + +proc `$`(client: Client): string = + "Client(ip: $1)" % [client.hostname] + +proc updateClients(server: Server) {.async.} = + while true: + var needsUpdate = false + for client in server.clients: + if not client.connected: + needsUpdate = true + break + + server.needsUpdate = server.needsUpdate or needsUpdate + if server.needsUpdate and server.txs.len != 0: + var someDead = false + # perform a copy to prevent the race condition: + var txs = server.txs + setLen(server.txs, 0) + for tx in txs: + for c in server.clients: + if c.connected: + await c.socket.sendText(tx.data, false) + else: + someDead = true + if someDead: + var i = 0 + while i < server.clients.len: + if not server.clients[i].connected: del(server.clients, i) + else: inc i + server.needsUpdate = false + # let other stuff in the main loop run: + await sleepAsync(10) + +proc processMessage(server: Server, client: Client, data: string) {.async.} = + # Check if last message was relatively recent. If so, kick the user. + echo "processMessage ", data + if epochTime() - client.lastMessage < 0.1: # 100ms + client.rapidMessageCount.inc + else: + client.rapidMessageCount = 0 + + client.lastMessage = epochTime() + if client.rapidMessageCount > 10: + warn("Client ($1) is firing messages too rapidly. Killing." % $client) + client.connected = false + let msgj = parseJson(data) + let msg = messageFromJson(msgj) + case msg.kind + of Newdata: + if msg.version == server.version: + server.txs.add Tx(data: data, version: msg.version) + server.needsUpdate = true + inc server.version + else: + let om = Message(kind: Rejected, id: msg.id, data: @[], version: server.version) + await client.socket.sendText($(%*om), false) + else: + # either Disconnect or an invalid message type: + client.connected = false + server.needsUpdate = true + +proc processClient(server: Server, client: Client) {.async.} = + while client.connected: + var frameFut = client.socket.readData(false) + yield frameFut + if frameFut.failed: + error("Error occurred handling client messages.\n" & + frameFut.error.msg) + client.connected = false + break + + let frame = frameFut.read() + if frame.opcode == Opcode.Text: + let processFut = processMessage(server, client, frame.data) + if processFut.failed: + error("Client ($1) attempted to send bad JSON? " % $client & "\n" & + processFut.error.msg) + client.connected = false + + client.socket.close() + +proc onRequest(server: Server, req: Request) {.async.} = + let (success, error) = await verifyWebsocketRequest(req, "karaxdb") + if success: + echo("Client connected from ", req.hostname) + server.clients.add(newClient(req.client, req.hostname)) + asyncCheck processClient(server, server.clients[^1]) + else: + echo("WS negotiation failed: ", error) + await req.respond(Http400, "WebSocket negotiation failed: " & error) + req.client.close() + +proc main = + let httpServer = newAsyncHttpServer() + let server = Server(clients: @[], txs: @[]) + + proc cb(req: Request): Future[void] {.async.} = await onRequest(server, req) + + asyncCheck updateClients(server) + waitFor httpServer.serve(Port(8080), cb) + +main()