diff --git a/.gitignore b/.gitignore index 710f20b..68ce831 100644 --- a/.gitignore +++ b/.gitignore @@ -1,28 +1,12 @@ -# The directory Mix will write compiled artifacts to. /_build/ - -# If you run "mix test --cover", coverage assets end up here. /cover/ - -# The directory Mix downloads your dependencies sources to. /deps/ - -# Where third-party dependencies like ExDoc output generated docs. /doc/ - -# If the VM crashes, it generates a dump, let's ignore it too. erl_crash.dump - -# Also ignore archive artifacts (built via "mix archive.build"). *.ez - -# Ignore package tarball (built via "mix hex.build"). drinkup-*.tar - -# Temporary files, for example, from tests. /tmp/ - -# Nix .envrc .direnv -result \ No newline at end of file +result +priv/dets/ \ No newline at end of file diff --git a/CHANGELOG.md b/CHANGELOG.md index c087e89..3172c15 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -13,6 +13,12 @@ and this project adheres to - Existing behaviour moved to `Drinkup.Firehose` namespace, to make way for alternate sync systems. +### Added + +- Support for the + [Tap](https://github.com/bluesky-social/indigo/blob/main/cmd/tap/README.md) + sync and backfill utility service, via `Drinkup.Tap`. + ### Changed - Refactor core connection logic for websockets into `Drinkup.Socket` to make it diff --git a/compose.yml b/compose.yml new file mode 100644 index 0000000..2f6abaa --- /dev/null +++ b/compose.yml @@ -0,0 +1,14 @@ +services: + tap: + image: "ghcr.io/bluesky-social/indigo/tap" + restart: "unless-stopped" + ports: + - "127.0.0.1:2480:2480" + volumes: + - "tap_data:/data" + environment: + TAP_SIGNAL_COLLECTION: "sh.weaver.actor.profile" + TAP_COLLECTION_FILTERS: "sh.weaver.*" + +volumes: + tap_data: diff --git a/examples/tap_consumer.ex b/examples/tap_consumer.ex new file mode 100644 index 0000000..d3d8110 --- /dev/null +++ b/examples/tap_consumer.ex @@ -0,0 +1,33 @@ +defmodule TapConsumer do + @behaviour Drinkup.Tap.Consumer + + def handle_event(%Drinkup.Tap.Event.Record{} = record) do + IO.inspect(record, label: "Tap record event") + end + + def handle_event(%Drinkup.Tap.Event.Identity{} = identity) do + IO.inspect(identity, label: "Tap identity event") + end +end + +defmodule TapExampleSupervisor do + use Supervisor + + def start_link(arg \\ []) do + Supervisor.start_link(__MODULE__, arg, name: __MODULE__) + end + + @impl true + def init(_) do + children = [ + {Drinkup.Tap, + %{ + consumer: TapConsumer, + name: MyTap, + host: "http://localhost:2480" + }} + ] + + Supervisor.init(children, strategy: :one_for_one) + end +end diff --git a/lib/firehose/socket.ex b/lib/firehose/socket.ex index b3922fc..151ee0f 100644 --- a/lib/firehose/socket.ex +++ b/lib/firehose/socket.ex @@ -39,7 +39,7 @@ defmodule Drinkup.Firehose.Socket do end @impl true - def handle_frame({:binary, frame}, %{seq: seq, options: options} = data) do + def handle_frame({:binary, frame}, {%{seq: seq, options: options} = data, _conn, _stream}) do with {:ok, header, next} <- CAR.DagCbor.decode(frame), {:ok, payload, _} <- CAR.DagCbor.decode(next), {%{"op" => @op_regular, "t" => type}, _} <- {header, payload}, diff --git a/lib/socket.ex b/lib/socket.ex index 119b3cb..8613c01 100644 --- a/lib/socket.ex +++ b/lib/socket.ex @@ -32,12 +32,19 @@ defmodule Drinkup.Socket do @callback build_path(data :: user_data()) :: String.t() - @callback handle_frame(frame :: frame(), data :: user_data()) :: + @callback handle_frame( + frame :: frame(), + data :: {user_data(), conn :: pid() | nil, stream :: :gun.stream_ref() | nil} + ) :: {:ok, new_data :: user_data()} | :noop | nil | {:error, reason :: term()} - @callback handle_connected(data :: user_data()) :: {:ok, new_data :: user_data()} + @callback handle_connected(data :: {user_data(), conn :: pid(), stream :: :gun.stream_ref()}) :: + {:ok, new_data :: user_data()} - @callback handle_disconnected(reason :: term(), data :: user_data()) :: + @callback handle_disconnected( + reason :: term(), + data :: {user_data(), conn :: pid() | nil, stream :: :gun.stream_ref() | nil} + ) :: {:ok, new_data :: user_data()} @optional_callbacks handle_connected: 1, handle_disconnected: 2 @@ -76,10 +83,10 @@ defmodule Drinkup.Socket do defoverridable child_spec: 1 @impl true - def handle_connected(data), do: {:ok, data} + def handle_connected({user_data, _conn, _stream}), do: {:ok, user_data} @impl true - def handle_disconnected(_reason, data), do: {:ok, data} + def handle_disconnected(_reason, {user_data, _conn, _stream}), do: {:ok, user_data} defoverridable handle_connected: 1, handle_disconnected: 2 end @@ -211,10 +218,14 @@ defmodule Drinkup.Socket do # :connected state - active WebSocket connection - def connected(:enter, _from, %{module: module, user_data: user_data} = data) do + def connected( + :enter, + _from, + %{module: module, user_data: user_data, conn: conn, stream: stream} = data + ) do Logger.debug("[Drinkup.Socket] WebSocket connected") - case module.handle_connected(user_data) do + case module.handle_connected({user_data, conn, stream}) do {:ok, new_user_data} -> {:keep_state, %{data | user_data: new_user_data, reconnect_attempts: 0}} @@ -226,9 +237,10 @@ defmodule Drinkup.Socket do def connected( :info, {:gun_ws, conn, _stream, frame}, - %{module: module, user_data: user_data, options: options} = data + %{module: module, user_data: user_data, options: options, conn: conn, stream: stream} = + data ) do - result = module.handle_frame(frame, user_data) + result = module.handle_frame(frame, {user_data, conn, stream}) :ok = :gun.update_flow(conn, frame, options.flow) @@ -286,9 +298,9 @@ defmodule Drinkup.Socket do # Helper functions defp trigger_reconnect(data, reason \\ :unknown) do - %{module: module, user_data: user_data} = data + %{module: module, user_data: user_data, conn: conn, stream: stream} = data - case module.handle_disconnected(reason, user_data) do + case module.handle_disconnected(reason, {user_data, conn, stream}) do {:ok, new_user_data} -> {:keep_state, %{data | user_data: new_user_data}, [{:next_event, :internal, :reconnect}]} diff --git a/lib/tap.ex b/lib/tap.ex new file mode 100644 index 0000000..bea13da --- /dev/null +++ b/lib/tap.ex @@ -0,0 +1,249 @@ +defmodule Drinkup.Tap do + @moduledoc """ + Supervisor and HTTP API for Tap indexer/backfill service. + + Tap simplifies AT sync by handling the firehose connection, verification, + backfill, and filtering. Your application connects to a Tap service and + receives simple JSON events for only the repos and collections you care about. + + ## Usage + + Add Tap to your supervision tree: + + children = [ + {Drinkup.Tap, %{ + consumer: MyTapConsumer, + name: MyTap, + host: "http://localhost:2480", + admin_password: "secret" # optional + }} + ] + + Then interact with the Tap HTTP API: + + # Add repos to track (triggers backfill) + Drinkup.Tap.add_repos(MyTap, ["did:plc:abc123"]) + + # Get stats + {:ok, count} = Drinkup.Tap.get_repo_count(MyTap) + + ## Configuration + + Tap itself is configured via environment variables. See the Tap documentation + for details on configuring collection filters, signal collections, and other + operational settings: + https://github.com/bluesky-social/indigo/blob/main/cmd/tap/README.md + """ + + use Supervisor + alias Drinkup.Tap.Options + + @dialyzer nowarn_function: {:init, 1} + @impl true + def init({%Options{name: name} = drinkup_options, supervisor_options}) do + # Register options in Registry for HTTP API access + Registry.register(Drinkup.Registry, {name, TapOptions}, drinkup_options) + + children = [ + {Task.Supervisor, name: {:via, Registry, {Drinkup.Registry, {name, TapTasks}}}}, + {Drinkup.Tap.Socket, drinkup_options} + ] + + Supervisor.start_link( + children, + supervisor_options ++ [name: {:via, Registry, {Drinkup.Registry, {name, TapSupervisor}}}] + ) + end + + @spec child_spec(Options.options()) :: Supervisor.child_spec() + def child_spec(%{} = options), do: child_spec({options, [strategy: :one_for_one]}) + + @spec child_spec({Options.options(), Keyword.t()}) :: Supervisor.child_spec() + def child_spec({drinkup_options, supervisor_options}) do + %{ + id: Map.get(drinkup_options, :name, __MODULE__), + start: {__MODULE__, :init, [{Options.from(drinkup_options), supervisor_options}]}, + type: :supervisor, + restart: :permanent, + shutdown: 500 + } + end + + # HTTP API Functions + + @doc """ + Add DIDs to track. + + Triggers backfill for the specified DIDs. Historical events will be fetched + from each repo's PDS, followed by live events from the firehose. + """ + @spec add_repos(atom(), [String.t()]) :: {:ok, term()} | {:error, term()} + def add_repos(name \\ Drinkup.Tap, dids) when is_list(dids) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :post, "/repos/add", %{dids: dids}) do + {:ok, response} + end + end + + @doc """ + Remove DIDs from tracking. + + Stops syncing the specified repos and deletes tracked repo metadata. Does not + delete buffered events in the outbox. + """ + @spec remove_repos(atom(), [String.t()]) :: {:ok, term()} | {:error, term()} + def remove_repos(name \\ Drinkup.Tap, dids) when is_list(dids) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :post, "/repos/remove", %{dids: dids}) do + {:ok, response} + end + end + + @doc """ + Resolve a DID to its DID document. + """ + @spec resolve_did(atom(), String.t()) :: {:ok, term()} | {:error, term()} + def resolve_did(name \\ Drinkup.Tap, did) when is_binary(did) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/resolve/#{did}") do + {:ok, response} + end + end + + @doc """ + Get info about a tracked repo. + + Returns repo state, repo rev, record count, error info, and retry count. + """ + @spec get_repo_info(atom(), String.t()) :: {:ok, term()} | {:error, term()} + def get_repo_info(name \\ Drinkup.Tap, did) when is_binary(did) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/info/#{did}") do + {:ok, response} + end + end + + @doc """ + Get the total number of tracked repos. + """ + @spec get_repo_count(atom()) :: {:ok, integer()} | {:error, term()} + def get_repo_count(name \\ Drinkup.Tap) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/stats/repo-count") do + {:ok, response} + end + end + + @doc """ + Get the total number of tracked records. + """ + @spec get_record_count(atom()) :: {:ok, integer()} | {:error, term()} + def get_record_count(name \\ Drinkup.Tap) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/stats/record-count") do + {:ok, response} + end + end + + @doc """ + Get the number of events in the outbox buffer. + """ + @spec get_outbox_buffer(atom()) :: {:ok, integer()} | {:error, term()} + def get_outbox_buffer(name \\ Drinkup.Tap) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/stats/outbox-buffer") do + {:ok, response} + end + end + + @doc """ + Get the number of events in the resync buffer. + """ + @spec get_resync_buffer(atom()) :: {:ok, integer()} | {:error, term()} + def get_resync_buffer(name \\ Drinkup.Tap) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/stats/resync-buffer") do + {:ok, response} + end + end + + @doc """ + Get current firehose and list repos cursors. + """ + @spec get_cursors(atom()) :: {:ok, map()} | {:error, term()} + def get_cursors(name \\ Drinkup.Tap) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/stats/cursors") do + {:ok, response} + end + end + + @doc """ + Check Tap health status. + + Returns `{:ok, %{"status" => "ok"}}` if healthy. + """ + @spec health(atom()) :: {:ok, map()} | {:error, term()} + def health(name \\ Drinkup.Tap) do + with {:ok, options} <- get_options(name), + {:ok, response} <- make_request(options, :get, "/health") do + {:ok, response} + end + end + + # Private Functions + + @spec get_options(atom()) :: {:ok, Options.t()} | {:error, :not_found} + defp get_options(name) do + case Registry.lookup(Drinkup.Registry, {name, TapOptions}) do + [{_pid, options}] -> {:ok, options} + [] -> {:error, :not_found} + end + end + + @spec make_request(Options.t(), atom(), String.t(), map() | nil) :: + {:ok, term()} | {:error, term()} + defp make_request(options, method, path, body \\ nil) do + url = build_url(options.host, path) + headers = build_headers(options.admin_password) + + request_opts = [ + method: method, + url: url, + headers: headers + ] + + request_opts = + if body do + Keyword.merge(request_opts, json: body) + else + request_opts + end + + case Req.request(request_opts) do + {:ok, %{status: status, body: body}} when status in 200..299 -> + {:ok, body} + + {:ok, %{status: status, body: body}} -> + {:error, {:http_error, status, body}} + + {:error, reason} -> + {:error, reason} + end + end + + @spec build_url(String.t(), String.t()) :: String.t() + defp build_url(host, path) do + host = String.trim_trailing(host, "/") + "#{host}#{path}" + end + + @spec build_headers(String.t() | nil) :: list() + defp build_headers(nil), do: [] + + defp build_headers(admin_password) do + credentials = "admin:#{admin_password}" + auth_header = "Basic #{Base.encode64(credentials)}" + [{"authorization", auth_header}] + end +end diff --git a/lib/tap/consumer.ex b/lib/tap/consumer.ex new file mode 100644 index 0000000..2d34ad9 --- /dev/null +++ b/lib/tap/consumer.ex @@ -0,0 +1,46 @@ +defmodule Drinkup.Tap.Consumer do + @moduledoc """ + Consumer behaviour for handling Tap events. + + Implement this behaviour to process events from a Tap indexer/backfill service. + Events are dispatched asynchronously via `Task.Supervisor` and acknowledged + to Tap based on the return value of `handle_event/1`. + + ## Event Acknowledgment + + By default, events are acknowledged to Tap based on your return value: + + - `:ok`, `{:ok, any()}`, or `nil` → Success, event is acked to Tap + - `{:error, reason}` → Failure, event is NOT acked (Tap will retry after timeout) + - Exception raised → Failure, event is NOT acked (Tap will retry after timeout) + + Any other value will log a warning and acknowledge the event anyway. + + If you set `disable_acks: true` in your Tap options, no acks are sent regardless + of the return value. This matches Tap's `TAP_DISABLE_ACKS` environment variable. + + ## Example + + defmodule MyTapConsumer do + @behaviour Drinkup.Tap.Consumer + + def handle_event(%Drinkup.Tap.Event.Record{action: :create} = record) do + # Handle new record creation + case save_to_database(record) do + :ok -> :ok # Success - event will be acked + {:error, reason} -> {:error, reason} # Failure - Tap will retry + end + end + + def handle_event(%Drinkup.Tap.Event.Identity{} = identity) do + # Handle identity changes + update_identity(identity) + :ok # Success - event will be acked + end + end + """ + + alias Drinkup.Tap.Event + + @callback handle_event(Event.Record.t() | Event.Identity.t()) :: any() +end diff --git a/lib/tap/event.ex b/lib/tap/event.ex new file mode 100644 index 0000000..54c3f6b --- /dev/null +++ b/lib/tap/event.ex @@ -0,0 +1,105 @@ +defmodule Drinkup.Tap.Event do + @moduledoc """ + Event handling and dispatch for Tap events. + + Parses incoming JSON events from Tap and dispatches them to the configured + consumer via Task.Supervisor. After successful processing, sends an ack + message back to the socket. + """ + + require Logger + alias Drinkup.Tap.{Event, Options} + + @type t() :: Event.Record.t() | Event.Identity.t() + + @doc """ + Parse a JSON map into an event struct. + + Returns the appropriate event struct based on the "type" field. + """ + @spec from(map()) :: t() | nil + def from(%{"type" => "record"} = payload), do: Event.Record.from(payload) + def from(%{"type" => "identity"} = payload), do: Event.Identity.from(payload) + def from(_payload), do: nil + + @doc """ + Dispatch an event to the consumer via Task.Supervisor. + + Spawns a task that: + 1. Processes the event via the consumer's handle_event/1 callback + 2. Sends an ack to Tap if acks are enabled and the consumer returns :ok, {:ok, _}, or nil + 3. Does not ack if the consumer returns an error-like value or raises an exception + + Consumer return value semantics (when acks are enabled): + - `:ok` or `{:ok, any()}` or `nil` -> Success, send ack + - `{:error, _}` or any error-like tuple -> Failure, don't ack (Tap will retry) + - Exception raised -> Failure, don't ack (Tap will retry) + + If `disable_acks: true` is set in options, no acks are sent regardless of + consumer return value. + """ + @spec dispatch(t(), Options.t(), pid(), :gun.stream_ref()) :: :ok + def dispatch( + event, + %Options{consumer: consumer, name: name, disable_acks: disable_acks}, + conn, + stream + ) do + supervisor_name = {:via, Registry, {Drinkup.Registry, {name, TapTasks}}} + event_id = get_event_id(event) + + {:ok, _pid} = + Task.Supervisor.start_child(supervisor_name, fn -> + try do + result = consumer.handle_event(event) + + unless disable_acks do + case result do + :ok -> + send_ack(conn, stream, event_id) + + {:ok, _} -> + send_ack(conn, stream, event_id) + + nil -> + send_ack(conn, stream, event_id) + + :error -> + Logger.error("Consumer returned error for event #{event_id}, not acking.") + + {:error, reason} -> + Logger.error( + "Consumer returned error for event #{event_id}, not acking: #{inspect(reason)}" + ) + + _ -> + Logger.warning( + "Consumer returned unexpected value for event #{event_id}, acking anyway: #{inspect(result)}" + ) + + send_ack(conn, stream, event_id) + end + end + rescue + e -> + Logger.error( + "Error in Tap event handler (event #{event_id}), not acking: #{Exception.format(:error, e, __STACKTRACE__)}" + ) + end + end) + + :ok + end + + @spec send_ack(pid(), :gun.stream_ref(), integer()) :: :ok + defp send_ack(conn, stream, event_id) do + ack_message = Jason.encode!(%{type: "ack", id: event_id}) + + :ok = :gun.ws_send(conn, stream, {:text, ack_message}) + Logger.debug("[Drinkup.Tap] Acked event #{event_id}") + end + + @spec get_event_id(t()) :: integer() + defp get_event_id(%Event.Record{id: id}), do: id + defp get_event_id(%Event.Identity{id: id}), do: id +end diff --git a/lib/tap/event/identity.ex b/lib/tap/event/identity.ex new file mode 100644 index 0000000..15e39d1 --- /dev/null +++ b/lib/tap/event/identity.ex @@ -0,0 +1,39 @@ +defmodule Drinkup.Tap.Event.Identity do + @moduledoc """ + Struct for identity events from Tap. + + Represents handle or status changes for a DID. + """ + + use TypedStruct + + typedstruct enforce: true do + field :id, integer() + field :did, String.t() + field :handle, String.t() | nil + field :is_active, boolean() + field :status, String.t() + end + + @spec from(map()) :: t() + def from(%{ + "id" => id, + "type" => "identity", + "identity" => + %{ + "did" => did, + "is_active" => is_active, + "status" => status + } = identity_data + }) do + handle = Map.get(identity_data, "handle") + + %__MODULE__{ + id: id, + did: did, + handle: handle, + is_active: is_active, + status: status + } + end +end diff --git a/lib/tap/event/record.ex b/lib/tap/event/record.ex new file mode 100644 index 0000000..698755d --- /dev/null +++ b/lib/tap/event/record.ex @@ -0,0 +1,58 @@ +defmodule Drinkup.Tap.Event.Record do + @moduledoc """ + Struct for record events from Tap. + + Represents create, update, or delete operations on records in the repository. + """ + + use TypedStruct + + typedstruct enforce: true do + @type action() :: :create | :update | :delete + + field :id, integer() + field :live, boolean() + field :rev, String.t() + field :did, String.t() + field :collection, String.t() + field :rkey, String.t() + field :action, action() + field :cid, String.t() | nil + field :record, map() | nil + end + + @spec from(map()) :: t() + def from(%{ + "id" => id, + "type" => "record", + "record" => + %{ + "live" => live, + "rev" => rev, + "did" => did, + "collection" => collection, + "rkey" => rkey, + "action" => action + } = record_data + }) do + cid = Map.get(record_data, "cid") + record = Map.get(record_data, "record") + + %__MODULE__{ + id: id, + live: live, + rev: rev, + did: did, + collection: collection, + rkey: rkey, + action: parse_action(action), + cid: cid, + record: record + } + end + + @spec parse_action(String.t()) :: action() + defp parse_action("create"), do: :create + defp parse_action("update"), do: :update + defp parse_action("delete"), do: :delete +end diff --git a/lib/tap/options.ex b/lib/tap/options.ex new file mode 100644 index 0000000..2313315 --- /dev/null +++ b/lib/tap/options.ex @@ -0,0 +1,90 @@ +defmodule Drinkup.Tap.Options do + @moduledoc """ + Configuration options for Tap indexer/backfill service connection. + + This module defines the configuration structure for connecting to and + interacting with a Tap service. Tap simplifies AT Protocol sync by handling + firehose connections, verification, backfill, and filtering server-side. + + ## Options + + - `:consumer` (required) - Module implementing `Drinkup.Tap.Consumer` behaviour + - `:name` - Unique name for this Tap instance in the supervision tree (default: `Drinkup.Tap`) + - `:host` - Tap service URL (default: `"http://localhost:2480"`) + - `:admin_password` - Optional password for authenticated Tap instances + - `:disable_acks` - Disable event acknowledgments (default: `false`) + + ## Example + + %{ + consumer: MyTapConsumer, + name: MyTap, + host: "http://localhost:2480", + admin_password: "secret", + disable_acks: false + } + """ + + use TypedStruct + + @default_host "http://localhost:2480" + + @typedoc """ + Map of configuration options accepted by `Drinkup.Tap.child_spec/1`. + """ + @type options() :: %{ + required(:consumer) => consumer(), + optional(:name) => name(), + optional(:host) => host(), + optional(:admin_password) => admin_password(), + optional(:disable_acks) => disable_acks() + } + + @typedoc """ + Module implementing the `Drinkup.Tap.Consumer` behaviour. + """ + @type consumer() :: module() + + @typedoc """ + Unique identifier for this Tap instance in the supervision tree. + + Used for Registry lookups and naming child processes. + """ + @type name() :: atom() + + @typedoc """ + HTTP/HTTPS URL of the Tap service. + + Defaults to `"http://localhost:2480"` which is Tap's default bind address. + """ + @type host() :: String.t() + + @typedoc """ + Optional password for HTTP Basic authentication. + + Required when connecting to a Tap service configured with `TAP_ADMIN_PASSWORD`. + The password is sent as `Basic admin:` in the Authorization header. + """ + @type admin_password() :: String.t() | nil + + @typedoc """ + Whether to disable event acknowledgments. + + When `true`, events are not acknowledged to Tap regardless of consumer + return values. This matches Tap's `TAP_DISABLE_ACKS` environment variable. + + Defaults to `false` (acknowledgments enabled). + """ + @type disable_acks() :: boolean() + + typedstruct do + field :consumer, consumer(), enforce: true + field :name, name(), default: Drinkup.Tap + field :host, host(), default: @default_host + field :admin_password, admin_password() + field :disable_acks, disable_acks(), default: false + end + + @spec from(options()) :: t() + def from(%{consumer: _} = options), do: struct(__MODULE__, options) +end diff --git a/lib/tap/socket.ex b/lib/tap/socket.ex new file mode 100644 index 0000000..bd1aeee --- /dev/null +++ b/lib/tap/socket.ex @@ -0,0 +1,100 @@ +defmodule Drinkup.Tap.Socket do + @moduledoc """ + WebSocket connection handler for Tap indexer/backfill service. + + Implements the Drinkup.Socket behaviour to manage connections to a Tap service, + handling JSON-encoded events and dispatching them to the configured consumer. + + Events are acknowledged after successful processing based on the consumer's + return value: + - `:ok`, `{:ok, any()}`, or `nil` → Success, ack sent to Tap + - `{:error, reason}` → Failure, no ack (Tap will retry after timeout) + - Exception raised → Failure, no ack (Tap will retry after timeout) + """ + + use Drinkup.Socket + + require Logger + alias Drinkup.Tap.{Event, Options} + + @impl true + def init(opts) do + options = Keyword.fetch!(opts, :options) + {:ok, %{options: options, host: options.host}} + end + + def start_link(%Options{} = options, statem_opts) do + socket_opts = build_socket_opts(options) + Drinkup.Socket.start_link(__MODULE__, socket_opts, statem_opts) + end + + @impl true + def build_path(_data) do + "/channel" + end + + @impl true + def handle_frame({:text, json}, {%{options: options} = data, conn, stream}) do + case Jason.decode(json) do + {:ok, payload} -> + case Event.from(payload) do + nil -> + Logger.warning("Received unrecognized event from Tap: #{inspect(payload)}") + :noop + + event -> + Event.dispatch(event, options, conn, stream) + {:ok, data} + end + + {:error, reason} -> + Logger.error("Failed to decode JSON from Tap: #{inspect(reason)}") + :noop + end + end + + @impl true + def handle_frame({:binary, _binary}, _data) do + Logger.warning("Received unexpected binary frame from Tap") + :noop + end + + @impl true + def handle_frame(:close, _data) do + Logger.info("Websocket closed, reason unknown") + nil + end + + @impl true + def handle_frame({:close, errno, reason}, _data) do + Logger.info("Websocket closed, errno: #{errno}, reason: #{inspect(reason)}") + nil + end + + defp build_socket_opts(%Options{host: host, admin_password: admin_password} = options) do + base_opts = [ + host: host, + options: options + ] + + if admin_password do + auth_header = build_auth_header(admin_password) + + gun_opts = %{ + ws_opts: %{ + headers: [{"authorization", auth_header}] + } + } + + Keyword.put(base_opts, :gun_opts, gun_opts) + else + base_opts + end + end + + @spec build_auth_header(String.t()) :: String.t() + defp build_auth_header(password) do + credentials = "admin:#{password}" + "Basic #{Base.encode64(credentials)}" + end +end diff --git a/mix.exs b/mix.exs index 153b0b7..5035244 100644 --- a/mix.exs +++ b/mix.exs @@ -35,7 +35,10 @@ defmodule Drinkup.MixProject do {:credo, "~> 1.7", only: [:dev, :test], runtime: false}, {:ex_doc, "~> 0.34", only: :dev, runtime: false}, {:gun, "~> 2.2"}, - {:typedstruct, "~> 0.5"} + {:typedstruct, "~> 0.5"}, + {:jason, "~> 1.4"}, + {:req, "~> 0.5.0"}, + {:atex, "~> 0.7"} ] end diff --git a/mix.lock b/mix.lock index b416fe7..d10644d 100644 --- a/mix.lock +++ b/mix.lock @@ -1,19 +1,39 @@ %{ + "atex": {:hex, :atex, "0.7.0", "23baa616d584ef2cdd2c444b838b672d4472cdae6894fc98dbcb8e1d4b3dd210", [:mix], [{:con_cache, "~> 1.1", [hex: :con_cache, repo: "hexpm", optional: false]}, {:ex_cldr, "~> 2.42", [hex: :ex_cldr, repo: "hexpm", optional: false]}, {:jason, "~> 1.4", [hex: :jason, repo: "hexpm", optional: false]}, {:jose, "~> 1.11", [hex: :jose, repo: "hexpm", optional: false]}, {:multiformats_ex, "~> 0.2", [hex: :multiformats_ex, repo: "hexpm", optional: false]}, {:mutex, "~> 3.0", [hex: :mutex, repo: "hexpm", optional: false]}, {:peri, "~> 0.6", [hex: :peri, repo: "hexpm", optional: false]}, {:plug, "~> 1.18", [hex: :plug, repo: "hexpm", optional: false]}, {:recase, "~> 0.5", [hex: :recase, repo: "hexpm", optional: false]}, {:req, "~> 0.5", [hex: :req, repo: "hexpm", optional: false]}, {:typedstruct, "~> 0.5", [hex: :typedstruct, repo: "hexpm", optional: false]}], "hexpm", "dfb5ced5259658ed6881add0e304b726cd281d7ae030813f7c4fe1e5fa8b35ef"}, "bunt": {:hex, :bunt, "1.0.0", "081c2c665f086849e6d57900292b3a161727ab40431219529f13c4ddcf3e7a44", [:mix], [], "hexpm", "dc5f86aa08a5f6fa6b8096f0735c4e76d54ae5c9fa2c143e5a1fc7c1cd9bb6b5"}, "car": {:hex, :car, "0.1.1", "a5bc4c5c1be96eab437634b3c0ccad1fe17b5e3d68c22a4031241ae1345aebd4", [:mix], [{:cbor, "~> 1.0.0", [hex: :cbor, repo: "hexpm", optional: false]}, {:typedstruct, "~> 0.5", [hex: :typedstruct, repo: "hexpm", optional: false]}, {:varint, "~> 1.4", [hex: :varint, repo: "hexpm", optional: false]}], "hexpm", "f895dda8123d04dd336db5a2bf0d0b47f4559cd5383f83fcca0700c1b45bfb6a"}, "cbor": {:hex, :cbor, "1.0.1", "39511158e8ea5a57c1fcb9639aaa7efde67129678fee49ebbda780f6f24959b0", [:mix], [], "hexpm", "5431acbe7a7908f17f6a9cd43311002836a34a8ab01876918d8cfb709cd8b6a2"}, "certifi": {:hex, :certifi, "2.16.0", "a4edfc1d2da3424d478a3271133bf28e0ec5e6fd8c009aab5a4ae980cb165ce9", [:rebar3], [], "hexpm", "8a64f6669d85e9cc0e5086fcf29a5b13de57a13efa23d3582874b9a19303f184"}, + "cldr_utils": {:hex, :cldr_utils, "2.29.1", "11ff0a50a36a7e5f3bd9fc2fb8486a4c1bcca3081d9c080bf9e48fe0e6742e2d", [:mix], [{:castore, "~> 0.1 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:certifi, "~> 2.5", [hex: :certifi, repo: "hexpm", optional: true]}, {:decimal, "~> 1.9 or ~> 2.0", [hex: :decimal, repo: "hexpm", optional: false]}], "hexpm", "3844a0a0ed7f42e6590ddd8bd37eb4b1556b112898f67dea3ba068c29aabd6c2"}, + "con_cache": {:hex, :con_cache, "1.1.1", "9f47a68dfef5ac3bbff8ce2c499869dbc5ba889dadde6ac4aff8eb78ddaf6d82", [:mix], [{:telemetry, "~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "1def4d1bec296564c75b5bbc60a19f2b5649d81bfa345a2febcc6ae380e8ae15"}, "cowlib": {:hex, :cowlib, "2.16.0", "54592074ebbbb92ee4746c8a8846e5605052f29309d3a873468d76cdf932076f", [:make, :rebar3], [], "hexpm", "7f478d80d66b747344f0ea7708c187645cfcc08b11aa424632f78e25bf05db51"}, "credo": {:hex, :credo, "1.7.15", "283da72eeb2fd3ccf7248f4941a0527efb97afa224bcdef30b4b580bc8258e1c", [:mix], [{:bunt, "~> 0.2.1 or ~> 1.0", [hex: :bunt, repo: "hexpm", optional: false]}, {:file_system, "~> 0.2 or ~> 1.0", [hex: :file_system, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}], "hexpm", "291e8645ea3fea7481829f1e1eb0881b8395db212821338e577a90bf225c5607"}, + "decimal": {:hex, :decimal, "2.3.0", "3ad6255aa77b4a3c4f818171b12d237500e63525c2fd056699967a3e7ea20f62", [:mix], [], "hexpm", "a4d66355cb29cb47c3cf30e71329e58361cfcb37c34235ef3bf1d7bf3773aeac"}, "earmark_parser": {:hex, :earmark_parser, "1.4.44", "f20830dd6b5c77afe2b063777ddbbff09f9759396500cdbe7523efd58d7a339c", [:mix], [], "hexpm", "4778ac752b4701a5599215f7030989c989ffdc4f6df457c5f36938cc2d2a2750"}, + "ex_cldr": {:hex, :ex_cldr, "2.44.1", "0d220b175874e1ce77a0f7213bdfe700b9be11aefbf35933a0e98837803ebdc5", [:mix], [{:cldr_utils, "~> 2.28", [hex: :cldr_utils, repo: "hexpm", optional: false]}, {:decimal, "~> 1.6 or ~> 2.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:gettext, "~> 0.19 or ~> 1.0", [hex: :gettext, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:nimble_parsec, "~> 0.5 or ~> 1.0", [hex: :nimble_parsec, repo: "hexpm", optional: true]}], "hexpm", "3880cd6137ea21c74250cd870d3330c4a9fdec07fabd5e37d1b239547929e29b"}, "ex_doc": {:hex, :ex_doc, "0.39.3", "519c6bc7e84a2918b737aec7ef48b96aa4698342927d080437f61395d361dcee", [:mix], [{:earmark_parser, "~> 1.4.44", [hex: :earmark_parser, repo: "hexpm", optional: false]}, {:makeup_c, ">= 0.1.0", [hex: :makeup_c, repo: "hexpm", optional: true]}, {:makeup_elixir, "~> 0.14 or ~> 1.0", [hex: :makeup_elixir, repo: "hexpm", optional: false]}, {:makeup_erlang, "~> 0.1 or ~> 1.0", [hex: :makeup_erlang, repo: "hexpm", optional: false]}, {:makeup_html, ">= 0.1.0", [hex: :makeup_html, repo: "hexpm", optional: true]}], "hexpm", "0590955cf7ad3b625780ee1c1ea627c28a78948c6c0a9b0322bd976a079996e1"}, "file_system": {:hex, :file_system, "1.1.1", "31864f4685b0148f25bd3fbef2b1228457c0c89024ad67f7a81a3ffbc0bbad3a", [:mix], [], "hexpm", "7a15ff97dfe526aeefb090a7a9d3d03aa907e100e262a0f8f7746b78f8f87a5d"}, + "finch": {:hex, :finch, "0.20.0", "5330aefb6b010f424dcbbc4615d914e9e3deae40095e73ab0c1bb0968933cadf", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:mint, "~> 1.6.2 or ~> 1.7", [hex: :mint, repo: "hexpm", optional: false]}, {:nimble_options, "~> 0.4 or ~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:nimble_pool, "~> 1.1", [hex: :nimble_pool, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "2658131a74d051aabfcba936093c903b8e89da9a1b63e430bee62045fa9b2ee2"}, "gun": {:hex, :gun, "2.2.0", "b8f6b7d417e277d4c2b0dc3c07dfdf892447b087f1cc1caff9c0f556b884e33d", [:make, :rebar3], [{:cowlib, ">= 2.15.0 and < 3.0.0", [hex: :cowlib, repo: "hexpm", optional: false]}], "hexpm", "76022700c64287feb4df93a1795cff6741b83fb37415c40c34c38d2a4645261a"}, + "hpax": {:hex, :hpax, "1.0.3", "ed67ef51ad4df91e75cc6a1494f851850c0bd98ebc0be6e81b026e765ee535aa", [:mix], [], "hexpm", "8eab6e1cfa8d5918c2ce4ba43588e894af35dbd8e91e6e55c817bca5847df34a"}, "jason": {:hex, :jason, "1.4.4", "b9226785a9aa77b6857ca22832cffa5d5011a667207eb2a0ad56adb5db443b8a", [:mix], [{:decimal, "~> 1.0 or ~> 2.0", [hex: :decimal, repo: "hexpm", optional: true]}], "hexpm", "c5eb0cab91f094599f94d55bc63409236a8ec69a21a67814529e8d5f6cc90b3b"}, + "jose": {:hex, :jose, "1.11.12", "06e62b467b61d3726cbc19e9b5489f7549c37993de846dfb3ee8259f9ed208b3", [:mix, :rebar3], [], "hexpm", "31e92b653e9210b696765cdd885437457de1add2a9011d92f8cf63e4641bab7b"}, "makeup": {:hex, :makeup, "1.2.1", "e90ac1c65589ef354378def3ba19d401e739ee7ee06fb47f94c687016e3713d1", [:mix], [{:nimble_parsec, "~> 1.4", [hex: :nimble_parsec, repo: "hexpm", optional: false]}], "hexpm", "d36484867b0bae0fea568d10131197a4c2e47056a6fbe84922bf6ba71c8d17ce"}, "makeup_elixir": {:hex, :makeup_elixir, "1.0.1", "e928a4f984e795e41e3abd27bfc09f51db16ab8ba1aebdba2b3a575437efafc2", [:mix], [{:makeup, "~> 1.0", [hex: :makeup, repo: "hexpm", optional: false]}, {:nimble_parsec, "~> 1.2.3 or ~> 1.3", [hex: :nimble_parsec, repo: "hexpm", optional: false]}], "hexpm", "7284900d412a3e5cfd97fdaed4f5ed389b8f2b4cb49efc0eb3bd10e2febf9507"}, "makeup_erlang": {:hex, :makeup_erlang, "1.0.2", "03e1804074b3aa64d5fad7aa64601ed0fb395337b982d9bcf04029d68d51b6a7", [:mix], [{:makeup, "~> 1.0", [hex: :makeup, repo: "hexpm", optional: false]}], "hexpm", "af33ff7ef368d5893e4a267933e7744e46ce3cf1f61e2dccf53a111ed3aa3727"}, + "mime": {:hex, :mime, "2.0.7", "b8d739037be7cd402aee1ba0306edfdef982687ee7e9859bee6198c1e7e2f128", [:mix], [], "hexpm", "6171188e399ee16023ffc5b76ce445eb6d9672e2e241d2df6050f3c771e80ccd"}, + "mint": {:hex, :mint, "1.7.1", "113fdb2b2f3b59e47c7955971854641c61f378549d73e829e1768de90fc1abf1", [:mix], [{:castore, "~> 0.1.0 or ~> 1.0", [hex: :castore, repo: "hexpm", optional: true]}, {:hpax, "~> 0.1.1 or ~> 0.2.0 or ~> 1.0", [hex: :hpax, repo: "hexpm", optional: false]}], "hexpm", "fceba0a4d0f24301ddee3024ae116df1c3f4bb7a563a731f45fdfeb9d39a231b"}, + "multiformats_ex": {:hex, :multiformats_ex, "0.2.0", "5b0a3faa1a770dc671aa8a89b6323cc20b0ecf67dc93dcd21312151fbea6b4ee", [:mix], [{:varint, "~> 1.4", [hex: :varint, repo: "hexpm", optional: false]}], "hexpm", "aa406d9addb06dc197e0e92212992486af6599158d357680f29f2d11e08d0423"}, + "mutex": {:hex, :mutex, "3.0.2", "528877fd0dbc09fc93ad667e10ea0d35a2126fa85205822f9dca85e87d732245", [:mix], [], "hexpm", "0a8f2ed3618160dca6a1e3520b293dc3c2ae53116265e71b4a732d35d29aa3c6"}, + "nimble_options": {:hex, :nimble_options, "1.1.1", "e3a492d54d85fc3fd7c5baf411d9d2852922f66e69476317787a7b2bb000a61b", [:mix], [], "hexpm", "821b2470ca9442c4b6984882fe9bb0389371b8ddec4d45a9504f00a66f650b44"}, "nimble_parsec": {:hex, :nimble_parsec, "1.4.2", "8efba0122db06df95bfaa78f791344a89352ba04baedd3849593bfce4d0dc1c6", [:mix], [], "hexpm", "4b21398942dda052b403bbe1da991ccd03a053668d147d53fb8c4e0efe09c973"}, + "nimble_pool": {:hex, :nimble_pool, "1.1.0", "bf9c29fbdcba3564a8b800d1eeb5a3c58f36e1e11d7b7fb2e084a643f645f06b", [:mix], [], "hexpm", "af2e4e6b34197db81f7aad230c1118eac993acc0dae6bc83bac0126d4ae0813a"}, + "peri": {:hex, :peri, "0.6.2", "3c043bfb6aa18eb1ea41d80981d19294c5e943937b1311e8e958da3581139061", [:mix], [{:ecto, "~> 3.12", [hex: :ecto, repo: "hexpm", optional: true]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:stream_data, "~> 1.1", [hex: :stream_data, repo: "hexpm", optional: true]}], "hexpm", "5e0d8e0bd9de93d0f8e3ad6b9a5bd143f7349c025196ef4a3591af93ce6ecad9"}, + "plug": {:hex, :plug, "1.19.1", "09bac17ae7a001a68ae393658aa23c7e38782be5c5c00c80be82901262c394c0", [:mix], [{:mime, "~> 1.0 or ~> 2.0", [hex: :mime, repo: "hexpm", optional: false]}, {:plug_crypto, "~> 1.1.1 or ~> 1.2 or ~> 2.0", [hex: :plug_crypto, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4.3 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "560a0017a8f6d5d30146916862aaf9300b7280063651dd7e532b8be168511e62"}, + "plug_crypto": {:hex, :plug_crypto, "2.1.1", "19bda8184399cb24afa10be734f84a16ea0a2bc65054e23a62bb10f06bc89491", [:mix], [], "hexpm", "6470bce6ffe41c8bd497612ffde1a7e4af67f36a15eea5f921af71cf3e11247c"}, + "recase": {:hex, :recase, "0.9.1", "82d2e2e2d4f9e92da1ce5db338ede2e4f15a50ac1141fc082b80050b9f49d96e", [:mix], [], "hexpm", "19ba03ceb811750e6bec4a015a9f9e45d16a8b9e09187f6d72c3798f454710f3"}, + "req": {:hex, :req, "0.5.17", "0096ddd5b0ed6f576a03dde4b158a0c727215b15d2795e59e0916c6971066ede", [:mix], [{:brotli, "~> 0.3.1", [hex: :brotli, repo: "hexpm", optional: true]}, {:ezstd, "~> 1.0", [hex: :ezstd, repo: "hexpm", optional: true]}, {:finch, "~> 0.17", [hex: :finch, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: false]}, {:mime, "~> 2.0.6 or ~> 2.1", [hex: :mime, repo: "hexpm", optional: false]}, {:nimble_csv, "~> 1.0", [hex: :nimble_csv, repo: "hexpm", optional: true]}, {:plug, "~> 1.0", [hex: :plug, repo: "hexpm", optional: true]}], "hexpm", "0b8bc6ffdfebbc07968e59d3ff96d52f2202d0536f10fef4dc11dc02a2a43e39"}, + "telemetry": {:hex, :telemetry, "1.3.0", "fedebbae410d715cf8e7062c96a1ef32ec22e764197f70cda73d82778d61e7a2", [:rebar3], [], "hexpm", "7015fc8919dbe63764f4b4b87a95b7c0996bd539e0d499be6ec9d7f3875b79e6"}, "typedstruct": {:hex, :typedstruct, "0.5.4", "d1d33d58460a74f413e9c26d55e66fd633abd8ac0fb12639add9a11a60a0462a", [:make, :mix], [], "hexpm", "ffaef36d5dbaebdbf4ed07f7fb2ebd1037b2c1f757db6fb8e7bcbbfabbe608d8"}, "varint": {:hex, :varint, "1.5.1", "17160c70d0428c3f8a7585e182468cac10bbf165c2360cf2328aaa39d3fb1795", [:mix], [], "hexpm", "24f3deb61e91cb988056de79d06f01161dd01be5e0acae61d8d936a552f1be73"}, }