Packages
kvs
13.3.1
13.5.22-aleph
13.4.16
13.4.15
13.4.14
13.4.13
13.3.1
13.2.28
11.9.1
10.8.3
10.8.2
10.3.0
9.9.2
9.9.1
9.9.0
9.8.0
9.7.0
9.4.8
9.4.7
9.4.6
9.4.5
9.4.4
9.4.3
9.4.2
9.4.1
9.4.0
8.12.0
8.11.2
8.11.1
8.10.4
8.10.3
8.10.2
8.10.1
8.10.0
8.5.2
8.5.1
8.5.0
8.4.1
8.4.0
8.3.1
8.3.0
7.11.5
7.9.1
7.7.0
7.1.3
7.1.2
7.1.1
6.12.11
6.12.10
6.12.9
6.12.8
6.12.7
6.12.6
6.12.5
6.12.4
6.12.3
6.12.2
6.12.1
6.12.0
6.11.2
6.11.1
6.11.0
6.10.2
6.10.1
6.10.0
6.9.2
6.9.1
6.9.0
6.7.7
6.7.6
6.7.5
6.7.4
6.7.3
6.7.2
6.7.1
6.7.0
6.6.0
2.1.0
0.12.1
retired
KVS Key-Value Store Abstraction Layer
Current section
Files
Jump to
Current section
Files
lib/stream.ex
defmodule :kvs_stream do
require KVS
def db, do: :kvs.db()
def metainfo, do: KVS.schema(name: :kvs, tables: tables())
def tables do
[
KVS.table(name: :writer, fields: [:id, :count, :args, :cache, :first]),
KVS.table(name: :reader, fields: [:id, :pos, :cache, :args, :feed, :seek, :count, :dir])
]
end
def c4(r, v), do: KVS.reader(r, args: v)
def sn(m, t), do: put_elem(m, 2, t)
def sp(m, t), do: put_elem(m, 3, t)
def si(m, t), do: put_elem(m, 1, t)
def tab(w) when elem(w, 0) == :writer, do: :msg
def tab(m) when elem(m, 0) == :msg, do: :msg
def tab(t) when is_tuple(t), do: elem(t, 0)
def tab(_), do: []
def id(t) when is_tuple(t), do: elem(t, 1) # iter :id
def id(t), do: t
def en(t) when is_tuple(t), do: elem(t, 2) # iter :next
def en(_), do: []
def ep(t) when is_tuple(t), do: elem(t, 3) # iter :prev
def ep(_), do: []
def acc(0), do: :next
def acc(1), do: :prev
def n({:ok, r}, c, p, db), do: r(:kvs.get(tab(r), en(r), KVS.kvs(db: db)), c, p)
def n({:error, x}, _, _, _), do: {:error, x}
def p({:ok, r}, c, p, db), do: r(:kvs.get(tab(r), ep(r), KVS.kvs(db: db)), c, p)
def p({:error, x}, _, _, _), do: {:error, x}
def r({:ok, val}, c, p), do: KVS.reader(c, cache: {tab(val), id(val)}, pos: p, args: List.flatten([val]))
def r({:error, x}, _, _), do: {:error, x}
def w({:ok, _w = KVS.writer(first: [])}, :bot, c), do: KVS.reader(c, cache: [], pos: 1, dir: 0)
def w({:ok, w = KVS.writer(first: b)}, :bot, c), do: KVS.reader(c, cache: {tab(w), id(b)}, pos: 1, dir: 0)
def w({:ok, w = KVS.writer(cache: b, count: size)}, :top, c), do: KVS.reader(c, cache: {tab(w), id(b)}, pos: size, dir: 1)
def w({:error, x}, _, _), do: {:error, x}
def top(x = KVS.reader()), do: top(x, db())
def top(c = KVS.reader(feed: f), db), do: w(:kvs.get(:writer, f, KVS.kvs(db: db)), :top, c)
def bot(x = KVS.reader()), do: bot(x, db())
def bot(c = KVS.reader(feed: f), db), do: w(:kvs.get(:writer, f, KVS.kvs(db: db)), :bot, c)
def next(x = KVS.reader()), do: next(x, db())
def next(KVS.reader(cache: []), _), do: {:error, :empty}
def next(c = KVS.reader(cache: {t, r}, pos: p), db), do: n(:kvs.get(t, r, KVS.kvs(db: db)), c, p + 1, db)
def prev(x = KVS.reader()), do: prev(x, db())
def prev(KVS.reader(cache: []), _), do: {:error, :empty}
def prev(c = KVS.reader(cache: {t, r}, pos: p), db), do: p(:kvs.get(t, r, KVS.kvs(db: db)), c, p - 1, db)
def drop(x = KVS.reader()), do: drop(x, db())
def drop(c = KVS.reader(cache: []), _), do: KVS.reader(c, args: [])
def drop(c = KVS.reader(dir: d, cache: b, args: n, pos: p), db), do: drop(acc(d), n, c, c, p, b, db)
def take(x = KVS.reader()), do: take(x, db())
def take(c = KVS.reader(cache: []), _), do: KVS.reader(c, args: [])
def take(c = KVS.reader(dir: d, cache: _b, args: n, pos: p), db), do: take(acc(d), n, c, c, [], p, db)
def take(_, _, {:error, _}, c2, r, p, _) do
if r == [] do
KVS.reader(c2, args: [], pos: p, cache: [])
else
flat = List.flatten(r)
rev = Enum.reverse(flat)
KVS.reader(c2, args: rev, pos: p, cache: {tab(hd(flat)), id(hd(flat))})
end
end
def take(_, 0, c, c2, r, p, _) when elem(c, 0) == :reader do
flat = List.flatten([r])
case flat do
[] -> KVS.reader(c2, args: [], pos: p, cache: [])
_ -> KVS.reader(c2, args: Enum.reverse(flat), pos: p, cache: KVS.reader(c, :cache))
end
end
def take(_, 0, {:error, _}, c2, r, p, _) do
flat = List.flatten([r])
case flat do
[] -> KVS.reader(c2, args: [], pos: p, cache: [])
_ -> KVS.reader(c2, args: Enum.reverse(flat), pos: p, cache: {tab(hd(flat)), id(hd(flat))})
end
end
def take(_, 0, KVS.reader(cache: []), c2, r, p, _) do
flat = List.flatten([r])
case flat do
[] -> KVS.reader(c2, args: [], pos: p, cache: [])
_ -> KVS.reader(c2, args: Enum.reverse(flat), pos: p, cache: {tab(hd(flat)), id(hd(flat))})
end
end
def take(a, n, c = KVS.reader(cache: {t, i}, pos: p), c2, r, _, db) do
{:ok, val} = :kvs.get(t, i, KVS.kvs(db: db))
next_c = Kernel.apply(__MODULE__, a, [c, db])
case next_c do
KVS.reader(cache: []) -> take(0, n - 1, next_c, c2, [val | r], p, db)
{:error, _} -> take(0, n - 1, next_c, c2, [val | r], p, db)
{:ok, _} -> take(a, n - 1, next_c, c2, [val | r], p, db)
_ -> take(a, n - 1, next_c, c2, [val | r], p, db)
end
end
def drop(_, _, {:error, _}, c2, p, b, _), do: KVS.reader(c2, pos: p, cache: b)
def drop(_, 0, _, c2, p, b, _), do: KVS.reader(c2, pos: p, cache: b)
def drop(a, n, c = KVS.reader(cache: _b, pos: p), c2, _, _, db) do
next_c = Kernel.apply(__MODULE__, a, [c, db])
next_b = case next_c do
KVS.reader(cache: new_b) -> new_b
_ -> []
end
drop(a, n - 1, next_c, c2, p, next_b, db)
end
def feed(feed), do: feed(feed, db())
def feed(feed, db) do
KVS.reader(args: args) = take(KVS.reader(get_reader(feed, db), args: -1), db)
args
end
def load_reader(id), do: load_reader(id, db())
def load_reader(id, db) do
case :kvs.get(:reader, id, KVS.kvs(db: db)) do
{:ok, c} -> c
_ -> KVS.reader(id: [])
end
end
def save(c), do: save(c, db())
def save(c, db) when elem(c, 0) == :reader do
nc = c4(c, [])
:kvs.put(nc, KVS.kvs(db: db))
nc
end
def save(c, db) when elem(c, 0) == :writer do
:kvs.put(c, KVS.kvs(db: db))
c
end
def get_writer(id), do: get_writer(id, db())
def get_writer(id, db) do
case :kvs.get(:writer, id, KVS.kvs(db: db)) do
{:ok, w} -> w
{:error, _} -> KVS.writer(id: id)
end
end
def get_reader(id), do: get_reader(id, db())
def get_reader(id, db) do
case :kvs.get(:writer, id, KVS.kvs(db: db)) do
{:ok, KVS.writer(first: [], count: c)} -> KVS.reader(id: :kvs.seq(:reader, 1, KVS.kvs(db: db)), feed: id, args: c, cache: [])
{:ok, KVS.writer(first: f, count: c)} -> KVS.reader(id: :kvs.seq(:reader, 1, KVS.kvs(db: db)), feed: id, args: c, cache: {tab(f), id(f)})
{:error, _} ->
save(KVS.writer(id: id), db)
get_reader(id, db)
end
end
def add(x = KVS.writer()), do: add(x, db())
def add(c = KVS.writer(args: m), db) when elem(m, 1) == [], do: add(si(m, :kvs.seq(tab(m), 1)), c, db)
def add(c = KVS.writer(args: m), db), do: add(m, c, db)
def add(m, c = KVS.writer(cache: []), db) do
_id = id(m)
n = sp(sn(m, []), [])
:kvs.put(n, KVS.kvs(db: db))
KVS.writer(c, cache: n, count: 1, first: n)
end
def add(m, c = KVS.writer(cache: v1, count: s), db) do
{:ok, v} = :kvs.get(tab(v1), id(v1), KVS.kvs(db: db))
n = sp(sn(m, []), id(v))
p = sn(v, id(m))
:kvs.put([n, p], KVS.kvs(db: db))
KVS.writer(c, cache: n, count: s + 1)
end
def cut(feed), do: cut(feed, db())
def cut(_, _), do: :ignore
def remove(c = KVS.reader(feed: feed)), do: remove(c, feed, db())
def remove(rec, feed), do: remove(rec, feed, db())
def remove(_rec, feed, db) do
case :kvs.get(:writer, feed, KVS.kvs(db: db)) do
{:ok, w = KVS.writer(count: count)} ->
:kvs.save(KVS.writer(w, count: count - 1), KVS.kvs(db: db))
count - 1
{:error, _} -> {:error, :not_found}
end
end
def append(rec, feed), do: append(rec, feed, db())
def append(rec, feed, db) do
:kvs.ensure(KVS.writer(id: feed), KVS.kvs(db: db))
id = elem(rec, 1)
name = elem(rec, 0)
case :kvs.get(name, id, KVS.kvs(db: db)) do
{:ok, _msg} ->
id
{:error, _} ->
{:ok, w = KVS.writer(first: first, cache: c, count: count)} = :kvs.get(:writer, feed, KVS.kvs(db: db))
prev = if c == [], do: [], else: elem(c, 1)
if prev != [] do
{:ok, msg} = :kvs.get(name, prev, KVS.kvs(db: db))
:kvs.put(:kvs.setfield(msg, :next, id), KVS.kvs(db: db))
end
rec_with_prev = :kvs.setfield(rec, :prev, prev)
:kvs.put(rec_with_prev, KVS.kvs(db: db))
nc = count + 1
w_new =
if first == [] do
KVS.writer(w, count: nc, first: {name, id}, cache: {name, id})
else
KVS.writer(w, count: nc, cache: {name, id})
end
save(w_new, db)
id
end
end
end