diff --git a/config/config.exs b/config/config.exs index de69c9f..9ce06ea 100644 --- a/config/config.exs +++ b/config/config.exs @@ -70,7 +70,6 @@ config :opentelemetry_exporter, otlp_protocol: :http_protobuf, otlp_endpoint: "http://localhost:4318" - # Import environment specific config. This must remain at the bottom # of this file so it overrides the configuration defined above. import_config "#{config_env()}.exs" diff --git a/config/runtime.exs b/config/runtime.exs index 957daed..14e7eb8 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -26,10 +26,11 @@ config :opentelemetry_exporter, otlp_endpoint: System.get_env("OTEL_EXPORTER_OTLP_ENDPOINT") || "http://localhost:4318" # Admin DID allowlist (comma-separated DIDs). -config :atjams, :admin_dids, - (System.get_env("ADMIN_DIDS") || "") - |> String.split(",", trim: true) - |> Enum.map(&String.trim/1) +config :atjams, + :admin_dids, + (System.get_env("ADMIN_DIDS") || "") + |> String.split(",", trim: true) + |> Enum.map(&String.trim/1) if config_env() == :prod do database_path = diff --git a/lib/atjams/application.ex b/lib/atjams/application.ex index 7b656b3..97c914d 100644 --- a/lib/atjams/application.ex +++ b/lib/atjams/application.ex @@ -26,6 +26,7 @@ defmodule Atjams.Application do restart: :permanent, shutdown: 500 }, + Atjams.Feed.JetstreamDidSync, AtjamsWeb.Endpoint ] @@ -42,5 +43,4 @@ defmodule Atjams.Application do AtjamsWeb.Endpoint.config_change(changed, removed) :ok end - end diff --git a/lib/atjams/feed.ex b/lib/atjams/feed.ex index 4801dbe..8283aa8 100644 --- a/lib/atjams/feed.ex +++ b/lib/atjams/feed.ex @@ -546,14 +546,18 @@ defmodule Atjams.Feed do case Repo.get(KnownUser, did) do nil -> - %KnownUser{} - |> KnownUser.changeset(%{ - did: did, - handle: handle, - first_seen_at: now, - last_login_at: now - }) - |> Repo.insert() + result = + %KnownUser{} + |> KnownUser.changeset(%{ + did: did, + handle: handle, + first_seen_at: now, + last_login_at: now + }) + |> Repo.insert() + + if match?({:ok, _}, result), do: sync_dids() + result %KnownUser{} = existing -> existing @@ -562,10 +566,28 @@ defmodule Atjams.Feed do end end + defp sync_dids do + Atjams.Feed.JetstreamDidSync.sync() + end + def list_known_users do Repo.all(KnownUser) end + @doc """ + Returns DIDs that have at least one atjams record (share, like, comment, or + follow). Only these DIDs are added to the Jetstream `wanted_dids` filter so + we don't waste bandwidth on users who logged in via OAuth but never created + any atjams content. + """ + def known_dids do + share_dids = Repo.all(from s in Share, select: s.did, distinct: true) + like_dids = Repo.all(from l in Like, select: l.did, distinct: true) + comment_dids = Repo.all(from c in Comment, select: c.did, distinct: true) + follow_dids = Repo.all(from f in Follow, select: f.did, distinct: true) + Enum.uniq(share_dids ++ like_dids ++ comment_dids ++ follow_dids) + end + ## teal.fm @doc """ @@ -574,7 +596,10 @@ defmodule Atjams.Feed do """ def tracks_teal_for?(did) when is_binary(did) do Repo.exists?(from k in KnownUser, where: k.did == ^did) or - Repo.exists?(from p in Profile, where: p.did == ^did) + Repo.exists?(from s in Share, where: s.did == ^did) or + Repo.exists?(from l in Like, where: l.did == ^did) or + Repo.exists?(from c in Comment, where: c.did == ^did) or + Repo.exists?(from f in Follow, where: f.did == ^did) end def get_teal_status(did) when is_binary(did), do: Repo.get(TealStatus, did) diff --git a/lib/atjams/feed/jetstream_consumer.ex b/lib/atjams/feed/jetstream_consumer.ex index 44aa829..ec4b2d5 100644 --- a/lib/atjams/feed/jetstream_consumer.ex +++ b/lib/atjams/feed/jetstream_consumer.ex @@ -25,7 +25,8 @@ defmodule Atjams.Feed.JetstreamConsumer do @follow_collection, @teal_status_collection, @teal_play_collection - ] + ], + wanted_dids: [] ## Shares diff --git a/lib/atjams/feed/jetstream_did_sync.ex b/lib/atjams/feed/jetstream_did_sync.ex new file mode 100644 index 0000000..7f66baf --- /dev/null +++ b/lib/atjams/feed/jetstream_did_sync.ex @@ -0,0 +1,82 @@ +defmodule Atjams.Feed.JetstreamDidSync do + @moduledoc """ + Syncs the list of active atjams DIDs to the Jetstream consumer's + `wanted_dids` filter. This limits the firehose to only events from + users who have actually created atjams records (shares, likes, + comments, follows), reducing bandwidth. + + Runs once at startup (to pick up existing atjams records from the DB) + and then on demand via `sync/0` whenever a new known user logs in. + """ + + use GenServer + + @retry_ms 2_000 + + @impl true + def start_link(opts \\ []) do + GenServer.start_link(__MODULE__, opts, name: __MODULE__) + end + + @impl true + def init(_state) do + {:ok, %{retry: nil}} + end + + @impl true + def handle_continue(:sync, state) do + do_sync(state) + end + + @impl true + def handle_info(:sync, state) do + do_sync(state) + end + + defp do_sync(state) do + dids = Atjams.Feed.known_dids() + + if dids == [] do + # No atjams records yet — nothing to filter on. Schedule a retry + # in case the backfill hasn't finished. + retry_ref = Process.send_after(self(), :sync, @retry_ms) + {:noreply, %{state | retry: retry_ref}} + else + case Drinkup.Jetstream.update_options(:atjams_jetstream, %{wanted_dids: dids}) do + :ok -> + {:noreply, %{state | retry: nil}} + + {:error, :not_connected} -> + # Socket isn't up yet — retry shortly + retry_ref = Process.send_after(self(), :sync, @retry_ms) + {:noreply, %{state | retry: retry_ref}} + end + end + end + + @doc """ + Pushes the current known DIDs to the Jetstream consumer. + + Safe to call from anywhere — errors are silently handled. The + periodic sync in the GenServer handles retries, so calling this + is just a hint to check sooner. + """ + def sync do + dids = Atjams.Feed.known_dids() + + if dids != [] do + case Drinkup.Jetstream.update_options(:atjams_jetstream, %{wanted_dids: dids}) do + :ok -> :ok + {:error, :not_connected} -> :ok + end + end + + :ok + end + + @impl true + def terminate(_reason, state) do + if ref = state.retry, do: Process.cancel_timer(ref) + :ok + end +end diff --git a/lib/atjams_web/live/feed_live.ex b/lib/atjams_web/live/feed_live.ex index fe52afc..4697053 100644 --- a/lib/atjams_web/live/feed_live.ex +++ b/lib/atjams_web/live/feed_live.ex @@ -138,7 +138,7 @@ defmodule AtjamsWeb.FeedLive do defp include_share?(%{assigns: %{tab: :following}} = socket, share) do case current_user_did(socket) do nil -> false - did -> Feed.follows?(did, share.did) + did -> did == share.did or Feed.follows?(did, share.did) end end diff --git a/lib/atjams_web/live/profile_live.ex b/lib/atjams_web/live/profile_live.ex index eb96b08..32702d5 100644 --- a/lib/atjams_web/live/profile_live.ex +++ b/lib/atjams_web/live/profile_live.ex @@ -107,8 +107,7 @@ defmodule AtjamsWeb.ProfileLive do def handle_info({:teal_play_upserted, %{did: did}}, socket) do if did == socket.assigns.profile_did do - {:noreply, - assign(socket, :teal_plays, Feed.list_teal_plays_for_did(did, @teal_play_limit))} + {:noreply, assign(socket, :teal_plays, Feed.list_teal_plays_for_did(did, @teal_play_limit))} else {:noreply, socket} end diff --git a/lib/atjams_web/live/profile_live.html.heex b/lib/atjams_web/live/profile_live.html.heex index 6bb0b09..c597f0a 100644 --- a/lib/atjams_web/live/profile_live.html.heex +++ b/lib/atjams_web/live/profile_live.html.heex @@ -10,8 +10,7 @@
{@following_count} following - · - {@follower_count} + · {@follower_count} {if @follower_count == 1, do: "follower", else: "followers"}
@@ -49,7 +48,10 @@ Now playing on teal.fm - {@teal_status.track}<%= if @teal_status.artist do %> — {@teal_status.artist}<% end %> + {@teal_status.track} + <%= if @teal_status.artist do %> + — {@teal_status.artist} + <% end %> diff --git a/lib/atjams_web/live/share_new_live.html.heex b/lib/atjams_web/live/share_new_live.html.heex index ac6587e..5f6fa90 100644 --- a/lib/atjams_web/live/share_new_live.html.heex +++ b/lib/atjams_web/live/share_new_live.html.heex @@ -38,7 +38,10 @@ Now playing on teal.fm - {@teal_status.track}<%= if @teal_status.artist do %> — {@teal_status.artist}<% end %> + {@teal_status.track} + <%= if @teal_status.artist do %> + — {@teal_status.artist} + <% end %> diff --git a/lib/atjams_web/router.ex b/lib/atjams_web/router.ex index 3dea5a6..320344e 100644 --- a/lib/atjams_web/router.ex +++ b/lib/atjams_web/router.ex @@ -9,16 +9,16 @@ defmodule AtjamsWeb.Router do # to avoid FOUC). Tighten to a nonce/hash if/when that block is removed. # img-src includes the bsky CDN for AT-protocol avatars. @csp [ - "default-src 'self'", - "img-src 'self' data: https://cdn.bsky.app", - "style-src 'self' 'unsafe-inline'", - "script-src 'self' 'unsafe-inline'", - "connect-src 'self' ws: wss:", - "frame-ancestors 'none'", - "base-uri 'self'", - "form-action 'self'" - ] - |> Enum.join("; ") + "default-src 'self'", + "img-src 'self' data: https://cdn.bsky.app", + "style-src 'self' 'unsafe-inline'", + "script-src 'self' 'unsafe-inline'", + "connect-src 'self' ws: wss:", + "frame-ancestors 'none'", + "base-uri 'self'", + "form-action 'self'" + ] + |> Enum.join("; ") pipeline :browser do plug :accepts, ["html"] diff --git a/mix.lock b/mix.lock index fcc50af..5eb841a 100644 --- a/mix.lock +++ b/mix.lock @@ -50,10 +50,8 @@ "nimble_pool": {:hex, :nimble_pool, "1.1.0", "bf9c29fbdcba3564a8b800d1eeb5a3c58f36e1e11d7b7fb2e084a643f645f06b", [:mix], [], "hexpm", "af2e4e6b34197db81f7aad230c1118eac993acc0dae6bc83bac0126d4ae0813a"}, "opentelemetry": {:hex, :opentelemetry, "1.7.0", "20d0f12d3d1c398d3670fd44fd1a7c495dd748ab3e5b692a7906662e2fb1a38a", [:rebar3], [{:opentelemetry_api, "~> 1.5.0", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}], "hexpm", "a9173b058c4549bf824cbc2f1d2fa2adc5cdedc22aa3f0f826951187bbd53131"}, "opentelemetry_api": {:hex, :opentelemetry_api, "1.5.0", "1a676f3e3340cab81c763e939a42e11a70c22863f645aa06aafefc689b5550cf", [:mix, :rebar3], [], "hexpm", "f53ec8a1337ae4a487d43ac89da4bd3a3c99ddf576655d071deed8b56a2d5dda"}, - "opentelemetry_api_experimental": {:hex, :opentelemetry_api_experimental, "0.5.1", "1b5afacfcbd0834390336c845bc8ae08c8cf0d69bbed72ee53d178798b93e074", [:mix, :rebar3], [{:opentelemetry_api, "~> 1.3", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}], "hexpm", "10297057eada47267d4f832011becef07d25690e6bf91febccfc4e740dba1a6f"}, "opentelemetry_bandit": {:hex, :opentelemetry_bandit, "0.3.0", "2c242dfdaabd747c75f4d8331fc9c17cfc9fb1db0638309762a4fcfa6d49a147", [:mix], [{:nimble_options, "~> 1.1", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:opentelemetry_api, "~> 1.3", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}, {:opentelemetry_semantic_conventions, "~> 1.27", [hex: :opentelemetry_semantic_conventions, repo: "hexpm", optional: false]}, {:otel_http, "~> 0.2", [hex: :otel_http, repo: "hexpm", optional: false]}, {:plug, ">= 1.15.0", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 1.2", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "5aa12378f5ff7cc3368f02905693571833f9449df86211fd99f4d764720cff60"}, "opentelemetry_ecto": {:hex, :opentelemetry_ecto, "1.2.0", "2382cb47ddc231f953d3b8263ed029d87fbf217915a1da82f49159d122b64865", [:mix], [{:opentelemetry_api, "~> 1.0", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}, {:opentelemetry_process_propagator, "~> 0.2", [hex: :opentelemetry_process_propagator, repo: "hexpm", optional: false]}, {:telemetry, "~> 0.4 or ~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "70dfa2e79932e86f209df00e36c980b17a32f82d175f0068bf7ef9a96cf080cf"}, - "opentelemetry_experimental": {:hex, :opentelemetry_experimental, "0.5.1", "27f60ea61b9e42f919c219d52dd17881057921150130b7eb9f15bc902f34f11d", [:rebar3], [{:opentelemetry, "~> 1.4", [hex: :opentelemetry, repo: "hexpm", optional: false]}, {:opentelemetry_api, "~> 1.3", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}, {:opentelemetry_api_experimental, "~> 0.5.1", [hex: :opentelemetry_api_experimental, repo: "hexpm", optional: false]}], "hexpm", "a1ad941294f1d3623c33e151faa35613849a10cb468dbfc9ad16367f7ddf80bf"}, "opentelemetry_exporter": {:hex, :opentelemetry_exporter, "1.10.0", "972e142392dbfa679ec959914664adefea38399e4f56ceba5c473e1cabdbad79", [:rebar3], [{:grpcbox, ">= 0.0.0", [hex: :grpcbox, repo: "hexpm", optional: false]}, {:opentelemetry, "~> 1.7.0", [hex: :opentelemetry, repo: "hexpm", optional: false]}, {:opentelemetry_api, "~> 1.5.0", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}, {:tls_certificate_check, "~> 1.18", [hex: :tls_certificate_check, repo: "hexpm", optional: false]}], "hexpm", "33a116ed7304cb91783f779dec02478f887c87988077bfd72840f760b8d4b952"}, "opentelemetry_phoenix": {:hex, :opentelemetry_phoenix, "2.0.1", "c664cdef205738cffcd409b33599439a4ffb2035ef6e21a77927ac1da90463cb", [:mix], [{:nimble_options, "~> 1.0", [hex: :nimble_options, repo: "hexpm", optional: false]}, {:opentelemetry_api, "~> 1.4", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}, {:opentelemetry_process_propagator, "~> 0.3", [hex: :opentelemetry_process_propagator, repo: "hexpm", optional: false]}, {:opentelemetry_semantic_conventions, "~> 1.27", [hex: :opentelemetry_semantic_conventions, repo: "hexpm", optional: false]}, {:opentelemetry_telemetry, "~> 1.1", [hex: :opentelemetry_telemetry, repo: "hexpm", optional: false]}, {:otel_http, "~> 0.2", [hex: :otel_http, repo: "hexpm", optional: false]}, {:plug, ">= 1.11.0", [hex: :plug, repo: "hexpm", optional: false]}, {:telemetry, "~> 1.0", [hex: :telemetry, repo: "hexpm", optional: false]}], "hexpm", "a24fdccdfa6b890c8892c6366beab4a15a27ec0c692b0f77ec2a862e7b235f6e"}, "opentelemetry_process_propagator": {:hex, :opentelemetry_process_propagator, "0.3.0", "ef5b2059403a1e2b2d2c65914e6962e56371570b8c3ab5323d7a8d3444fb7f84", [:mix, :rebar3], [{:opentelemetry_api, "~> 1.0", [hex: :opentelemetry_api, repo: "hexpm", optional: false]}], "hexpm", "7243cb6de1523c473cba5b1aefa3f85e1ff8cc75d08f367104c1e11919c8c029"}, @@ -71,7 +69,6 @@ "phoenix_template": {:hex, :phoenix_template, "1.0.4", "e2092c132f3b5e5b2d49c96695342eb36d0ed514c5b252a77048d5969330d639", [:mix], [{:phoenix_html, "~> 2.14.2 or ~> 3.0 or ~> 4.0", [hex: :phoenix_html, repo: "hexpm", optional: true]}], "hexpm", "2c0c81f0e5c6753faf5cca2f229c9709919aba34fab866d3bc05060c9c444206"}, "plug": {:hex, :plug, "1.19.2", "e4950525b22c6789dfb38a3f95d47171ba159da3fc5a33be9643b43d5e8adb98", [: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", "b6fce20a56af5e60fa5dfecf3f907bb98ec981be43c79a3809a499bc3d133de0"}, "plug_crypto": {:hex, :plug_crypto, "2.1.1", "19bda8184399cb24afa10be734f84a16ea0a2bc65054e23a62bb10f06bc89491", [:mix], [], "hexpm", "6470bce6ffe41c8bd497612ffde1a7e4af67f36a15eea5f921af71cf3e11247c"}, - "postgrex": {:hex, :postgrex, "0.22.2", "4aec14df2a72722aee92492566edbeeb44e233ecb86b1915d03136297ef1385d", [:mix], [{:db_connection, "~> 2.9", [hex: :db_connection, repo: "hexpm", optional: false]}, {:decimal, "~> 1.5 or ~> 2.0 or ~> 3.0", [hex: :decimal, repo: "hexpm", optional: false]}, {:jason, "~> 1.0", [hex: :jason, repo: "hexpm", optional: true]}, {:table, "~> 0.1.0", [hex: :table, repo: "hexpm", optional: true]}], "hexpm", "8946382ddb06294f56026ac4278b3cc212bac8a2c82ed68b4087819ed1abc53b"}, "recase": {:hex, :recase, "0.9.1", "82d2e2e2d4f9e92da1ce5db338ede2e4f15a50ac1141fc082b80050b9f49d96e", [:mix], [], "hexpm", "19ba03ceb811750e6bec4a015a9f9e45d16a8b9e09187f6d72c3798f454710f3"}, "remote_ip": {:hex, :remote_ip, "1.2.0", "fb078e12a44414f4cef5a75963c33008fe169b806572ccd17257c208a7bc760f", [:mix], [{:combine, "~> 0.10", [hex: :combine, repo: "hexpm", optional: false]}, {:plug, "~> 1.14", [hex: :plug, repo: "hexpm", optional: false]}], "hexpm", "2ff91de19c48149ce19ed230a81d377186e4412552a597d6a5137373e5877cb7"}, "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"}, diff --git a/priv/repo/migrations/20260516193122_reset_jetstream_cursor.exs b/priv/repo/migrations/20260516193122_reset_jetstream_cursor.exs new file mode 100644 index 0000000..f84af31 --- /dev/null +++ b/priv/repo/migrations/20260516193122_reset_jetstream_cursor.exs @@ -0,0 +1,13 @@ +defmodule Atjams.Repo.Migrations.ResetJetstreamCursor do + use Ecto.Migration + + def up do + execute "DELETE FROM jetstream_cursors WHERE name = 'default'" + end + + def down do + # We can't restore the deleted cursor value — the migration is a one-way reset. + # If you need to roll back, the cursor will simply be nil (live-tail) on next + # start, which is the same end state. + end +end diff --git a/todo.md b/todo.md new file mode 100644 index 0000000..c257adb --- /dev/null +++ b/todo.md @@ -0,0 +1,11 @@ +# TODO + +## Features + +- Annotate that a share was "listened" to if a user has a teal scrobble within + some time after interacting with/viewing a share. (Not really sure how to + gauge interactions here though, maybe like but not sure) + +## Bugs + +- Constrain jetstream listener to just known dids