Packages
riptide
0.3.5
0.5.2
0.5.1
0.5.0-beta9
0.5.0-beta8
0.5.0-beta7
0.5.0-beta6
0.5.0-beta5
0.5.0-beta4
0.5.0-beta3
0.5.0-beta2
0.5.0-beta11
0.5.0-beta10
0.5.0-beta
0.4.6
0.4.5
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.13
0.3.12
0.3.11
0.3.10
0.3.9
0.3.8
0.3.7
0.3.6
0.3.5
0.3.4
0.3.3
0.3.2
0.3.1
0.3.0
0.3.0-bd63a38
0.2.79
0.2.78
0.2.74
0.2.4
0.2.3
0.2.2
0.2.1
0.2.0
0.1.15
0.1.14
0.1.13
0.1.12
0.1.11
0.1.10
0.1.9
0.1.8
0.1.7
0.1.6
0.1.5
0.1.4
0.1.3
0.1.2
0.1.1
0.1.0
A data first framework for building realtime applications
Current section
Files
Jump to
Current section
Files
lib/riptide/store/subscribe.ex
defmodule Riptide.Subscribe do
@moduledoc false
def watch(path), do: watch(path, self())
def watch(path, pid) do
group = group(path)
cond do
member?(group, pid) ->
:ok
true ->
:pg2.join(group, pid)
end
end
def member?(group, pid) do
case :pg2.get_members(group) do
{:error, {:no_such_group, _}} ->
:pg2.create(group)
false
result ->
pid in result
end
end
# TODO: This could have a better implementation
def broadcast_mutation(mut) do
mut
|> Riptide.Mutation.layers()
|> Stream.flat_map(fn {path, value} ->
Stream.concat([
[{path, Riptide.Mutation.inflate(path, value)}],
value.delete
|> Stream.filter(fn {_, value} -> value === 1 end)
|> Stream.map(fn {key, _} ->
{path ++ [key], Riptide.Mutation.delete(path ++ [key])}
end)
])
end)
|> Enum.each(fn {path, value} ->
path
|> group()
|> :pg2.get_members()
|> case do
{:error, _} ->
:skip
members ->
Enum.map(members, fn pid -> send(pid, {:mutation, value}) end)
end
end)
end
def group(path) do
{__MODULE__, path}
end
end