2022-06-29 11:04:35 -05:00
|
|
|
import std/sequtils
|
|
|
|
|
|
|
|
|
|
import pkg/chronos
|
2022-05-11 10:50:05 -05:00
|
|
|
import pkg/questionable
|
|
|
|
|
import pkg/questionable/results
|
2022-06-29 11:04:35 -05:00
|
|
|
from pkg/stew/results as stewResults import get, isErr
|
2022-05-11 10:50:05 -05:00
|
|
|
import pkg/upraises
|
|
|
|
|
|
|
|
|
|
import ./datastore
|
|
|
|
|
|
|
|
|
|
export datastore
|
|
|
|
|
|
|
|
|
|
push: {.upraises: [].}
|
|
|
|
|
|
|
|
|
|
type
|
|
|
|
|
TieredDatastore* = ref object of Datastore
|
|
|
|
|
stores: seq[Datastore]
|
|
|
|
|
|
|
|
|
|
proc new*(
|
|
|
|
|
T: type TieredDatastore,
|
|
|
|
|
stores: varargs[Datastore]): ?!T =
|
|
|
|
|
|
|
|
|
|
if stores.len == 0:
|
|
|
|
|
failure "stores must contain at least one Datastore"
|
|
|
|
|
else:
|
|
|
|
|
success T(stores: @stores)
|
|
|
|
|
|
|
|
|
|
proc stores*(self: TieredDatastore): seq[Datastore] =
|
|
|
|
|
self.stores
|
|
|
|
|
|
2022-09-09 00:39:16 -05:00
|
|
|
method close*(self: TieredDatastore) {.async, locks: "unknown".} =
|
|
|
|
|
for store in self.stores:
|
|
|
|
|
await store.close
|
|
|
|
|
|
2022-05-11 10:50:05 -05:00
|
|
|
method contains*(
|
|
|
|
|
self: TieredDatastore,
|
2022-06-29 11:04:35 -05:00
|
|
|
key: Key): Future[?!bool] {.async, locks: "unknown".} =
|
2022-05-11 10:50:05 -05:00
|
|
|
|
|
|
|
|
for store in self.stores:
|
2022-06-29 11:04:35 -05:00
|
|
|
let
|
|
|
|
|
containsRes = await store.contains(key)
|
2022-05-11 10:50:05 -05:00
|
|
|
|
2022-06-29 11:04:35 -05:00
|
|
|
if containsRes.isErr: return containsRes
|
|
|
|
|
if containsRes.get == true: return success true
|
|
|
|
|
|
|
|
|
|
return success false
|
2022-05-11 10:50:05 -05:00
|
|
|
|
|
|
|
|
method delete*(
|
|
|
|
|
self: TieredDatastore,
|
2022-06-29 11:04:35 -05:00
|
|
|
key: Key): Future[?!void] {.async, locks: "unknown".} =
|
2022-05-11 10:50:05 -05:00
|
|
|
|
2022-06-29 11:04:35 -05:00
|
|
|
let
|
|
|
|
|
pending = await allFinished(self.stores.mapIt(it.delete(key)))
|
|
|
|
|
|
|
|
|
|
for fut in pending:
|
2022-09-12 00:34:07 -05:00
|
|
|
let
|
|
|
|
|
delRes = await fut
|
|
|
|
|
|
|
|
|
|
if delRes.isErr: return delRes
|
2022-05-11 10:50:05 -05:00
|
|
|
|
2022-06-29 11:04:35 -05:00
|
|
|
return success()
|
2022-05-11 10:50:05 -05:00
|
|
|
|
|
|
|
|
method get*(
|
|
|
|
|
self: TieredDatastore,
|
2022-06-29 11:04:35 -05:00
|
|
|
key: Key): Future[?!(?seq[byte])] {.async, locks: "unknown".} =
|
2022-05-11 10:50:05 -05:00
|
|
|
|
|
|
|
|
var
|
|
|
|
|
bytesOpt: ?seq[byte]
|
|
|
|
|
|
|
|
|
|
for store in self.stores:
|
2022-06-29 11:04:35 -05:00
|
|
|
let
|
|
|
|
|
getRes = await store.get(key)
|
|
|
|
|
|
|
|
|
|
if getRes.isErr: return getRes
|
|
|
|
|
|
|
|
|
|
bytesOpt = getRes.get
|
2022-05-11 10:50:05 -05:00
|
|
|
|
|
|
|
|
# put found data into stores logically in front of the current store
|
|
|
|
|
if bytes =? bytesOpt:
|
|
|
|
|
for s in self.stores:
|
|
|
|
|
if s == store: break
|
2022-06-29 11:04:35 -05:00
|
|
|
let
|
|
|
|
|
putRes = await s.put(key, bytes)
|
|
|
|
|
|
|
|
|
|
if putRes.isErr: return failure putRes.error.msg
|
|
|
|
|
|
2022-05-11 10:50:05 -05:00
|
|
|
break
|
|
|
|
|
|
2022-06-29 11:04:35 -05:00
|
|
|
return success bytesOpt
|
2022-05-11 10:50:05 -05:00
|
|
|
|
|
|
|
|
method put*(
|
|
|
|
|
self: TieredDatastore,
|
|
|
|
|
key: Key,
|
2022-06-29 11:04:35 -05:00
|
|
|
data: seq[byte]): Future[?!void] {.async, locks: "unknown".} =
|
2022-05-11 10:50:05 -05:00
|
|
|
|
2022-06-29 11:04:35 -05:00
|
|
|
let
|
|
|
|
|
pending = await allFinished(self.stores.mapIt(it.put(key, data)))
|
|
|
|
|
|
|
|
|
|
for fut in pending:
|
2022-09-12 00:34:07 -05:00
|
|
|
let
|
|
|
|
|
putRes = await fut
|
|
|
|
|
|
|
|
|
|
if putRes.isErr: return putRes
|
2022-05-11 10:50:05 -05:00
|
|
|
|
2022-06-29 11:04:35 -05:00
|
|
|
return success()
|
2022-05-11 10:50:05 -05:00
|
|
|
|
2022-09-12 00:34:07 -05:00
|
|
|
iterator queryImpl(
|
|
|
|
|
datastore: Datastore,
|
|
|
|
|
query: Query): Future[QueryResponse] {.closure.} =
|
|
|
|
|
|
|
|
|
|
let
|
|
|
|
|
datastore = TieredDatastore(datastore)
|
|
|
|
|
# https://github.com/datastore/datastore/blob/7ccf0cd4748001d3dbf5e6dda369b0f63e0269d3/datastore/core/basic.py#L1027-L1035
|
|
|
|
|
bottom = datastore.stores[^1]
|
|
|
|
|
|
|
|
|
|
try:
|
|
|
|
|
let q = bottom.query(); for kv in q(bottom, query): yield kv
|
|
|
|
|
except Exception as e:
|
|
|
|
|
raise (ref Defect)(msg: e.msg)
|
|
|
|
|
|
|
|
|
|
method query*(self: TieredDatastore): QueryIterator {.locks: "unknown".} =
|
|
|
|
|
queryImpl
|