Current section

Files

Jump to
extreme lib extreme.ex
Raw

lib/extreme.ex

defmodule Extreme do
@moduledoc """
TODO
"""
@type t :: module
@doc false
defmacro __using__(opts \\ []) do
quote do
alias Extreme.Messages, as: ExMsg
@otp_app Keyword.get(unquote(opts), :otp_app, :extreme)
defp _default_config,
do: Application.get_env(@otp_app, __MODULE__)
def child_spec(opts) do
%{
id: __MODULE__,
start: {__MODULE__, :start_link, [opts]},
type: :supervisor
}
end
def start_link(config \\ [])
def start_link([]), do: Extreme.Supervisor.start_link(__MODULE__, _default_config())
def start_link(config), do: Extreme.Supervisor.start_link(__MODULE__, config)
def ping,
do: Extreme.RequestManager.ping(__MODULE__, Extreme.Tools.generate_uuid())
def execute(message, correlation_id \\ nil, timeout \\ 5_000) do
Extreme.RequestManager.execute(
__MODULE__,
message,
correlation_id || Extreme.Tools.generate_uuid(),
timeout
)
end
@doc """
Reads events from `stream_name` given `opts`
as keyword list of `Extreme.Reading.Params` keys,
and invokes `fun` for each read event. `fun` should return `:ok` or `:stop`
if further processing of events should be stopped
"""
@spec read_events(String.t(), Keyword.t(), (ExMsg.StreamEventAppeared.t() ->
:ok
| :stop
| {:stop, response :: any()})) ::
:finished
| {:stopped, response :: nil | any()}
| {:error, :no_stream | :stream_hard_deleted}
| {:error, :unexpected_processing_response, any()}
def read_events(stream_name, opts \\ [], fun) do
params =
[{:stream, stream_name} | opts]
|> Enum.into(%{})
|> Extreme.Reading.Params.new()
Extreme.Reading.read_events(__MODULE__, params, fun)
end
@doc """
Reads events backwards from `stream_name` given `opts`
as keyword list of `Extreme.Reading.Params` keys,
and invokes `fun` for each read event. `fun` should return `:ok` or `:stop`
if further processing of events should be stopped
"""
@spec read_events_backwards(String.t(), Keyword.t(), (ExMsg.StreamEventAppeared.t() ->
:ok
| :stop
| {:stop, response :: any()})) ::
:finished
| {:stopped, response :: any()}
| {:error, :no_stream | :stream_hard_deleted}
| {:error, :unexpected_processing_response, any()}
def read_events_backwards(stream_name, opts \\ [], fun) do
params =
[{:stream, stream_name} | opts]
|> Enum.into(%{})
|> Extreme.Reading.Params.new_backwards()
Extreme.Reading.read_events(__MODULE__, params, fun)
end
@spec reduce_events(acc :: any(), String.t(), Keyword.t(), (ExMsg.StreamEventAppeared.t(),
acc :: any() ->
{:ok, acc :: any()}
| {:stop, acc :: any()})) ::
{:finished, acc :: any()}
| {:stopped, acc :: any()}
| {:error, :no_stream | :stream_hard_deleted}
| {:error, :unexpected_processing_response, any()}
def reduce_events(acc, stream_name, opts \\ [], fun) do
params =
[{:stream, stream_name} | opts]
|> Enum.into(%{})
|> Extreme.Reading.Params.new()
Extreme.Reading.reduce_events(__MODULE__, acc, params, fun)
end
@spec reduce_events_backwards(
acc :: any(),
String.t(),
Keyword.t(),
(ExMsg.StreamEventAppeared.t(), acc :: any() ->
{:ok, acc :: any()}
| {:stop, acc :: any()})
) ::
{:finished, acc :: any()}
| {:stopped, acc :: any()}
| {:error, :no_stream | :stream_hard_deleted}
| {:error, :unexpected_processing_response, any()}
def reduce_events_backwards(acc, stream_name, opts \\ [], fun) do
params =
[{:stream, stream_name} | opts]
|> Enum.into(%{})
|> Extreme.Reading.Params.new_backwards()
Extreme.Reading.reduce_events(__MODULE__, acc, params, fun)
end
@doc """
Sets metadata map to stream. To remove metadata, set an empty map.
Example:
metadata = %{ "$maxAge" => max_age_seconds }
:ok = MyConn.set_metadata("user-123", metadata)
"""
@spec set_metadata(String.t(), map()) :: :ok | any()
def set_metadata(stream, %{} = metadata) do
stream
|> _write_metadata(metadata)
|> execute()
|> case do
{:ok, %ExMsg.WriteEventsCompleted{result: :success}} -> :ok
{:ok, %ExMsg.WriteEventsCompleted{result: result}} -> {:error, result}
other -> other
end
end
@doc """
Gets metadata map from stream.
Example:
> stream = "users-123"
> MyConn.get_metadata(stream)
{:error, :no_stream}
> max_age_seconds = 60 * 60 * 24 # keep events 1 day
> metadata = %{ "$maxAge" => max_age_seconds }
> :ok = MyConn.set_metadata(stream, metadata)
> MyConn.get_metadata(stream)
{:ok, %{ "$maxAge" => 86_400 }}
"""
@spec get_metadata(String.t()) :: {:ok, map()} | {:error, :no_stream}
def get_metadata(stream) do
stream
|> _read_metadata()
|> execute()
|> case do
{:error, :no_stream, %ExMsg.ReadStreamEventsCompleted{result: :no_stream}} ->
{:error, :no_stream}
{:ok,
%ExMsg.ReadStreamEventsCompleted{
events: [
%ExMsg.ResolvedIndexedEvent{
event: %Extreme.Messages.EventRecord{
event_stream_id: "$$" <> ^stream,
event_type: "$metadata",
data: data
}
}
| _
],
result: :success
}} ->
data
|> Jason.decode()
|> case do
{:ok, decoded} -> {:ok, decoded}
_ -> data
end
end
end
@doc """
Starts database scavenge on current connection. Pay attention that
if cluster is used, scavenge will be executed only on connected node!
"""
@spec scavenge_database() :: :ok | any()
def scavenge_database() do
ExMsg.ScavengeDatabase.new()
|> execute()
|> case do
{:ok, %Extreme.Messages.ScavengeDatabaseCompleted{result: :success}} -> :ok
{:ok, %Extreme.Messages.ScavengeDatabaseCompleted{result: other}} -> {:error, other}
other -> other
end
end
defp _write_metadata(stream, %{} = metadata) do
metadata_stream_name = "$$" <> stream
proto_event =
ExMsg.NewEvent.new(
event_id: Extreme.Tools.generate_uuid(),
event_type: "$metadata",
data_content_type: 1,
metadata_content_type: 1,
data: Jason.encode!(metadata),
metadata: ""
)
ExMsg.WriteEvents.new(
event_stream_id: metadata_stream_name,
expected_version: -2,
events: [proto_event],
require_master: false
)
end
defp _read_metadata(stream) do
metadata_stream_name = "$$" <> stream
Extreme.Messages.ReadStreamEventsBackward.new(
event_stream_id: metadata_stream_name,
from_event_number: -1,
max_count: 1,
resolve_link_tos: false,
require_master: false
)
end
def subscribe_to(stream, subscriber, resolve_link_tos \\ true, ack_timeout \\ 5_000)
when is_binary(stream) and is_pid(subscriber) and is_boolean(resolve_link_tos) do
Extreme.RequestManager.subscribe_to(
__MODULE__,
stream,
subscriber,
resolve_link_tos,
ack_timeout
)
end
def read_and_stay_subscribed(
stream,
subscriber,
from_event_number \\ 0,
per_page \\ 1_000,
resolve_link_tos \\ true,
require_master \\ false,
ack_timeout \\ 5_000
)
when is_binary(stream) and is_pid(subscriber) and is_boolean(resolve_link_tos) and
is_boolean(require_master) and from_event_number > -2 and per_page >= 0 and
per_page <= 4096 do
Extreme.RequestManager.read_and_stay_subscribed(
__MODULE__,
subscriber,
{stream, from_event_number, per_page, resolve_link_tos, require_master, ack_timeout}
)
end
@spec start_event_producer(stream :: String.t(), subscriber :: pid(), opts :: Keyword.t()) ::
Supervisor.on_start_child()
def start_event_producer(stream, subscriber, opts \\ []) do
Extreme.EventProducer.Supervisor.start_event_producer(
__MODULE__,
[{:stream, stream}, {:subscriber, subscriber} | opts]
)
end
@spec subscribe_producer(producer :: pid()) :: :ok
def subscribe_producer(producer),
do: Extreme.EventProducer.subscribe(producer)
@spec unsubscribe_producer(producer :: pid()) :: :ok
def unsubscribe_producer(producer),
do: Extreme.EventProducer.unsubscribe(producer)
@spec producer_subscription_status(producer :: pid()) ::
:disconnected | :catching_up | :live | :paused
def producer_subscription_status(producer),
do: Extreme.EventProducer.subscription_status(producer)
def unsubscribe(subscription) when is_pid(subscription),
do: Extreme.Subscription.unsubscribe(subscription)
def connect_to_persistent_subscription(
subscriber,
stream,
group,
allowed_in_flight_messages
) do
Extreme.RequestManager.connect_to_persistent_subscription(
__MODULE__,
subscriber,
stream,
group,
allowed_in_flight_messages
)
end
end
end
@doc """
TODO
"""
@callback start_link(config :: Keyword.t(), opts :: Keyword.t()) ::
{:ok, pid}
| {:error, {:already_started, pid}}
| {:error, term}
@doc """
TODO
"""
@callback execute(message :: term(), correlation_id :: binary(), timeout :: integer()) :: term()
@doc """
TODO
"""
@callback subscribe_to(stream :: String.t(), subscriber :: pid(), opts :: Keyword.t()) ::
{:ok, pid}
@doc """
TODO
"""
@callback unsubscribe(subscription :: pid()) :: :ok
@doc """
TODO
"""
@callback read_and_stay_subscribed(
stream :: String.t(),
subscriber :: pid(),
from_event_number :: integer(),
per_page :: integer(),
resolve_link_tos :: boolean(),
require_master :: boolean()
) :: {:ok, pid()}
@doc """
Pings connected EventStore and should return `:pong` back.
"""
@callback ping() :: :pong
@doc """
Spawns a persistent subscription.
The persistent subscription will send events to the `subscriber` process in
the form of `GenServer.cast/2`s in the shape of `{:on_event, event,
correlation_id}`.
See `Extreme.PersistentSubscription` for full details.
"""
@callback connect_to_persistent_subscription(
subscriber :: pid(),
stream :: String.t(),
group :: String.t(),
allowed_in_flight_messages :: integer()
) :: {:ok, pid()}
end