Packages
extreme
1.1.3
1.1.4
1.1.3
1.1.2
1.1.1
1.1.1-rc01
1.1.0
1.1.0-rc9
1.1.0-rc8
1.1.0-rc7
1.1.0-rc6
1.1.0-rc5
1.1.0-rc4
1.1.0-rc3
1.1.0-rc2
1.1.0-rc1
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.13.4
0.13.3
0.13.2
0.13.1
0.13.0
0.12.1
0.12.0
0.11.0
0.10.4
0.10.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.1
0.8.0
0.7.1
0.7.0
0.6.2
0.6.1
0.6.0
0.5.5
0.5.4
0.5.3
0.5.2
0.5.1
0.5.0
0.4.3
0.4.2
0.4.1
Elixir TCP client for EventStore.
Current section
Files
Jump to
Current section
Files
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