diff --git a/examples/todoapp/todoapp.nim b/examples/todoapp/todoapp.nim index 802bfc9..d1ceb21 100644 --- a/examples/todoapp/todoapp.nim +++ b/examples/todoapp/todoapp.nim @@ -1,5 +1,5 @@ -import vdom, karax, karaxdsl, jstrutils, components, localstorage +import vdom, karax, karaxdsl, jstrutils, components, karaxdb/client type Filter = enum @@ -9,6 +9,12 @@ var selectedEntry = -1 filter: Filter entriesLen: int + data: Db + +registerOnUpdate proc(newDb: Db) = + merge(newDb, data) + data = newDb + redraw() const contentSuffix = cstring"content" @@ -16,25 +22,23 @@ const lenSuffix = cstring"entriesLen" proc getEntryContent(pos: int): cstring = - result = getItem(&pos & contentSuffix) - if result == cstring"null": - result = nil + extract(data, &pos, contentSuffix) proc isCompleted(pos: int): bool = - var value = getItem(&pos & completedSuffix) + var value = extract(data, &pos, completedSuffix) result = value == cstring"true" proc setEntryContent(pos: int, content: cstring) = - setItem(&pos & contentSuffix, content) + insert(data, &pos, contentSuffix, content) proc markAsCompleted(pos: int, completed: bool) = - setItem(&pos & completedSuffix, &completed) + insert(data, &pos, completedSuffix, &completed) proc addEntry(content: cstring, completed: bool) = setEntryContent(entriesLen, content) markAsCompleted(entriesLen, completed) inc entriesLen - setItem(lenSuffix, &entriesLen) + insert(data, lenSuffix, "equals", &entriesLen) proc updateEntry(pos: int, content: cstring, completed: bool) = setEntryContent(pos, content) @@ -61,7 +65,7 @@ proc toggleEntry(ev: Event; n: VNode) = markAsCompleted(id, not isCompleted(id)) proc onAllDone(ev: Event; n: VNode) = - clear() + insert(data, lenSuffix, "equals", "0") selectedEntry = -1 proc clearCompleted(ev: Event, n: VNode) = @@ -150,8 +154,5 @@ setOnHashChange(proc(hash: cstring) = elif hash == "#/active": filter = active ) -if hasItem(lenSuffix): - entriesLen = parseInt getItem(lenSuffix) -else: - entriesLen = 0 +entriesLen = 0 setRenderer createDom diff --git a/src/kajax.nim b/src/kajax.nim index be230f0..2d88cd1 100644 --- a/src/kajax.nim +++ b/src/kajax.nim @@ -48,3 +48,21 @@ proc ajaxGet*(url: cstring; headers: openarray[(cstring, cstring)]; proc toJson*[T](data: T): cstring {.importc: "JSON.stringify".} proc fromJson*[T](blob: cstring): T {.importc: "JSON.parse".} + +type + MessageEvent* {.importc.} = ref object + data*: cstring + ErrorEvent* {.importc.} = ref object + CloseEvent* {.importc.} = ref object + code*: int + reason*: cstring + + WebSocket* {.importc.} = ref object + onopen*: proc() + onmessage*: proc(ev: MessageEvent) + onclose*: proc(ev: CloseEvent) + onerror*: proc(ev: ErrorEvent) + +proc newWebsocket*(url, protocol: cstring): WebSocket {.importc: "new WebSocket".} + +proc send*(socket: WebSocket, data: cstring) {.importcpp.} diff --git a/src/karaxdb/btree.nim b/src/karaxdb/btree.nim new file mode 100644 index 0000000..eb9353f --- /dev/null +++ b/src/karaxdb/btree.nim @@ -0,0 +1,562 @@ + +## General purpose BTree implementation. Can also be used as a persistent +## data structure. The persistent operations use a 'Ps' suffix. +## Can also use a page manager for allocations. The page manager be used +## to off load pages to a file system or to send it over the wire. + +## Todo: +## - Add logic to deal with the fact that keys do not have to be unique. +## - Ranged queries +## - Make it generic and low level +## - Support for external nodes and a page cache + +const + M = 4 # max children per B-tree node = M-1 + # (must be even and greater than 2) + Mhalf = M div 2 + + SupportFullTableScan = true + SupportDuplicateKeys = true + +## Due to the fact that leaves are shared among multiple BTrees the following +## fields in a Node are downright impossible: +## - parent +## - next +## - prev + +type + Key = string + Val = string + Node = ref object + m: int + keys: array[M, Key] + case isInternal: bool + of false: + vals: array[M, Val] + of true: + links: array[M, Node] + BTree = object + root: Node + height: int ## height + n: int ## number of key-value pairs + CmpKind {.pure.} = enum + eq, le, lt, ge, gt, neq + CursorState = enum stPop, stLeaf, stEnd + Cursor = object + n: Node + i: int + stack: seq[Node] + state: CursorState + +proc newBTree(): BTree = BTree(root: Node(m: 0, isInternal: false)) + +proc less(a, b: Key): bool = cmp(a, b) < 0 + +proc eq(a, b: Key): bool = cmp(a, b) == 0 + +proc search(x: Node, key: Key, ht: int): Val = + if ht == 0: + assert(not x.isInternal) + for j in 0 ..< x.m: + if eq(key, x.keys[j]): return x.vals[j] + else: + assert(x.isInternal) + for j in 0 ..< x.m: + if j+1 == x.m or less(key, x.keys[j+1]): + return search(x.links[j], key, ht-1) + +proc `=~`(i: int; k: CmpKind): bool = + ## check if the result of 'cmp' matches what was requested by 'k': + case k + of CmpKind.eq: i == 0 + of CmpKind.le: i <= 0 + of CmpKind.lt: i < 0 + of CmpKind.ge: i >= 0 + of CmpKind.gt: i > 0 + of CmpKind.neq: i != 0 + +proc dos(x: Node; kind: CmpKind; key: Key; withKey: proc(k: Key; v: Val)) = + if not x.isInternal: + for j in 0 ..< x.m: + if cmp(x.keys[j], key) =~ kind: + withKey(x.keys[j], x.vals[j]) + else: + # we compute the range of links to follow first, before + # recursing: + var followA = 0 + var followB = -1 + case kind + of CmpKind.eq: + # want: key == 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + for j in 1..x.m: + if j == x.m or cmp(key, x.keys[j]) < 0: + followA = j-1 + followB = j-1 + break + of CmpKind.le, CmpKind.lt: + # want: key <= 10 or key < 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + + # Case A: all keys are bigger: + if cmp(key, x.keys[1]) < 0: + # --> use the very first branch + followA = 0 + followB = 0 + else: + # Case B: all keys are smaller --> use all branches is covered too + # by this loop. + for j in 1..= 0: + if followB < 0: followA = j-1 + # if the keys are identical and we require 'lt', we know + # only the left branch is required: + followB = j - ord(kind == CmpKind.lt and cmpRes == 0) + else: + # it's already greater, all others are greater too: + break + of CmpKind.ge, CmpKind.gt: + # want: key >= 10 or key > 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + + # Case A: all keys are smaller: + if cmp(key, x.keys[x.m-1]) >= 0: + # --> use the very last branch + followA = x.m-1 + followB = x.m-1 + else: + # also covers case B: all keys are bigger --> use all branches + # we find the key that is bigger or equal to ours and from + # then on, follow every branch: + for j in 1.. 0: + let x = stack.pop() + if not x.isInternal: + for j in 0 ..< x.m: + if cmp(x.keys[j], key) =~ kind: + withKey(x.keys[j], x.vals[j]) + else: + # we compute the range of links to follow first, before + # recursing: + var followA = 0 + var followB = -1 + case kind + of CmpKind.eq: + # want: key == 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + for j in 1..x.m: + if j == x.m or cmp(key, x.keys[j]) < 0: + followA = j-1 + followB = j-1 + break + of CmpKind.le, CmpKind.lt: + # want: key <= 10 or key < 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + + # Case A: all keys are bigger: + if cmp(key, x.keys[1]) < 0: + # --> use the very first branch + followA = 0 + followB = 0 + else: + # Case B: all keys are smaller --> use all branches is covered too + # by this loop. + for j in 1..= 0: + if followB < 0: followA = j-1 + # if the keys are identical and we require 'lt', we know + # only the left branch is required: + followB = j - ord(kind == CmpKind.lt and cmpRes == 0) + else: + # it's already greater, all others are greater too: + break + of CmpKind.ge, CmpKind.gt: + # want: key >= 10 or key > 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + + # Case A: all keys are smaller: + if cmp(key, x.keys[x.m-1]) >= 0: + # --> use the very last branch + followA = x.m-1 + followB = x.m-1 + else: + # also covers case B: all keys are bigger --> use all branches + # we find the key that is bigger or equal to ours and from + # then on, follow every branch: + for j in 1.. use the very first branch + followA = 0 + followB = 0 + else: + # Case B: all keys are smaller --> use all branches is covered too + # by this loop. + for j in 1..= 0: + if followB < 0: followA = j-1 + # if the keys are identical and we require 'lt', we know + # only the left branch is required: + followB = j - ord(kind == CmpKind.lt and cmpRes == 0) + else: + # it's already greater, all others are greater too: + break + of CmpKind.ge, CmpKind.gt: + # want: key >= 10 or key > 10 + # keys: 0 3 4 5 10 20 + # keys: 20 30 40 + + # Case A: all keys are smaller: + if cmp(key, x.keys[x.m-1]) >= 0: + # --> use the very last branch + followA = x.m-1 + followB = x.m-1 + else: + # also covers case B: all keys are bigger --> use all branches + # we find the key that is bigger or equal to ours and from + # then on, follow every branch: + for j in 1.. 0: result.add(indent & "(" & $h.keys[j] & ")\n") + toString(h.links[j], ht-1, indent & " ", result) + +proc `$`(b: BTree): string = + result = "" + toString(b.root, b.height, "", result) + +when isMainModule: + proc main = + var st = newBTree() + st.put("www.cs.princeton.edu", "abc") + st.put("www.cs.princeton.edu", "xyz") + st.put("www.princeton.edu", "128.112.128.15") + st.put("www.yale.edu", "130.132.143.21") + st.put("www.simpsons.com", "209.052.165.60") + st.put("www.apple.com", "17.112.152.32") + st.put("www.amazon.com", "207.171.182.16") + st.put("www.ebay.com", "66.135.192.87") + st.put("www.cnn.com", "64.236.16.20") + st.put("www.google.com", "216.239.41.99") + st.put("www.nytimes.com", "199.239.136.200") + st.put("www.microsoft.com", "207.126.99.140") + st.put("www.dell.com", "143.166.224.230") + st.put("www.slashdot.org", "66.35.250.151") + st.put("www.espn.com", "199.181.135.201") + st.put("www.weather.com", "63.111.66.11") + st.put("www.yahoo.com", "216.109.118.65") + + assert st.get("www.cs.princeton.edu") == "abc" + assert st.get("www.harvardsucks.com") == nil + + assert st.get("www.simpsons.com") == "209.052.165.60" + assert st.get("www.apple.com") == "17.112.152.32" + assert st.get("www.ebay.com") == "66.135.192.87" + assert st.get("www.dell.com") == "143.166.224.230" + assert(st.n == 17) + + when false: + var b2 = newBTree() + const iters = 10_000 + for i in 1..iters: + b2.put($i, $(iters - i)) + for i in 1..iters: + let x = b2.get($i) + if x != $(iters - i): + echo "got ", x, ", but expected ", iters - i + echo b2.n + echo b2.height + + when true: + var b1 = newBTree() + var b2 = newBTree() + const iters = 9 #60_000 + for i in 1..iters: + b2 = b2.putPs($i, $(iters - i)) + b1.put($i, $(iters - i)) + for i in 1..iters: + let x = b2.get($i) + if x != $(iters - i): + echo i, "th iteration; got ", x, ", but expected ", iters - i + echo b2.n, " = ", b1.n + echo b2.height, " = ", b1.height + echo " >= 5" + dos(b1.root, CmpKind.ge, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + echo " <= 5" + dos(b1.root, CmpKind.le, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + + echo " == 5" + dos(b1.root, CmpKind.eq, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + echo " < 5" + dos(b1.root, CmpKind.lt, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + echo " > 5" + dos(b1.root, CmpKind.gt, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + + echo "======================================================================" + echo " >= 5" + don(b1.root, CmpKind.ge, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + echo " <= 5" + don(b1.root, CmpKind.le, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + + echo " == 5" + don(b1.root, CmpKind.eq, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + echo " < 5" + don(b1.root, CmpKind.lt, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + echo " > 5" + don(b1.root, CmpKind.gt, "5", proc(k: Key; v: Val) = echo("k ", k, " = ", v)) + + echo "======================================================================" + var c = initCursor(b1.root) + var i = 0 + while true: + next(c, CmpKind.le, "9") + if atEnd(c): break + echo "key ", getKey(c), " ", getVal(c) + if i > 30: break + inc i + + main() diff --git a/src/karaxdb/client.nim b/src/karaxdb/client.nim new file mode 100644 index 0000000..72c31bc --- /dev/null +++ b/src/karaxdb/client.nim @@ -0,0 +1,82 @@ + +import "../kajax", "../jdict", common +export common +from "../karax" import kout + +type + Db* = ref object + data: seq[Triple] + version*: int + next: Db + + Query* = object + constraints*: set[TripleKind] + t*: Triple + + Message {.importc.} = ref object + kind: MessageKind + data: seq[Triple] + version: int + id: MessageId +# RequestMessage {.importc.} = ref object + +let conn = newWebSocket("ws://localhost:8080", "karaxdb") +var version: int + +#proc loadDb*(url: cstring): Db = +# result = nil + +proc newTransaction*(): Db = + result = Db(data: @[]) + +proc insert*(head, newdb: Db) = + newdb.next = head + #result = newdb + 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) = + newer.next = older + +proc registerOnUpdate*(update: proc(db: Db)) = + conn.onmessage = + proc (e: MessageEvent) = + let msg = fromJson[Message](e.data) + case msg.kind + of Rejected: + # conflict, so throw away the sent data, don't apply the changes: + 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 + # yield outdated data! + var it = db + while it != nil: + for d in it.data: + var match = true + for k in q.constraints: + if d[k] != q.t[k]: match = false + if match: yield d + it = it.next + +proc extract*(db: Db, q: Query): Triple = + for x in list(db, q): + result = x + break + +proc extract*(db: Db, subj, pred: kstring): kstring = + let q = Query(constraints: {Subj, Pred}, t: [subj, pred, ""]) + result = extract(db, q)[Obj] + +proc insert*(db: Db, subj, pred, obj: kstring) = + let result = Db() + result.data.add([subj, pred, obj]) + insert(db, result) diff --git a/src/karaxdb/common.nim b/src/karaxdb/common.nim new file mode 100644 index 0000000..c098702 --- /dev/null +++ b/src/karaxdb/common.nim @@ -0,0 +1,17 @@ + +when defined(js): + type kstring* = cstring +else: + type kstring* = string + +type + TripleKind* = enum + Subj, Pred, Obj + + DbValue* = kstring + Triple* = array[TripleKind, DbValue] + + MessageId* = distinct int + + MessageKind* = enum + Newdata, Rejected, Disconnect 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()