Packages

A library for distributed reactive programming with flexible consistency guarantees drawing from QUARP and Rx.. Features the familiar behaviours and event streams in the spirit of FRP.

Current section

Files

Jump to
bquarp lib reactivity processing combine_with_guarantees.ex
Raw

lib/reactivity/processing/combine_with_guarantees.ex

defmodule Reactivity.Processing.CombineWithGuarantees do
@moduledoc false
use Observables.GenObservable
alias Reactivity.Processing.Matching
alias Reactivity.Quality.Guarantee
require Logger
def init([imap, gmap]) do
Logger.debug("CombineWithGuarantee: #{inspect(self())}")
{:ok, {:buffer, imap, :guarantees, gmap}}
end
def handle_event({:newvalue, index, msg}, {:buffer, buffer, :guarantees, gss}) do
updated_buffer = %{buffer | index => Map.get(buffer, index) ++ [msg]}
case Matching.match(updated_buffer, msg, index, gss) do
:nomatch ->
{:novalue, {:buffer, updated_buffer, :guarantees, gss}}
{:ok, match, contexts, new_buffer} ->
{vals, _contextss} = match
|> Enum.unzip
if first_value?(index, buffer)
and update?(index, gss)
and any_propagate?(gss) do
Process.send(self(), {:event, {:spit, index}}, [])
end
{:value, {vals, contexts}, {:buffer, new_buffer, :guarantees, gss}}
end
end
def handle_event({:spit, index}, {:buffer, buffer, :guarantees, gss}) do
msg = buffer
|> Map.get(index)
|> List.first
case Matching.match(buffer, msg, index, gss) do
:nomatch ->
{:novalue, {:buffer, buffer, :guarantees, gss}}
{:ok, match, contexts, new_buffer} ->
{vals, _contextss} = match
|> Enum.unzip
Process.send(self(), {:event, {:spit, index}}, [])
{:value, {vals, contexts}, {:buffer, new_buffer, :guarantees, gss}}
end
end
def handle_done(_pid, _state) do
Logger.debug("#{inspect(self())}: combinelatestn has one dead dependency, going on.")
{:ok, :continue}
end
defp first_value?(index, buffer) do
buffer
|> Map.get(index)
|> Enum.empty?
end
defp update?(index, gss) do
gss
|> Map.get(index)
|> Guarantee.semantics
== :update
end
defp any_propagate?(gss) do
gss
|> Map.values
|> Enum.map(fn gs -> Guarantee.semantics(gs) end)
|> Enum.any?(fn sem -> sem == :propagate end)
end
end