diff --git a/.gitignore b/.gitignore index 8ef6823..32b8503 100644 --- a/.gitignore +++ b/.gitignore @@ -43,3 +43,8 @@ npm-debug.log *.db-shm *.db-wal +# Nix build artifacts +/result +/result-* + + diff --git a/config/runtime.exs b/config/runtime.exs index 2b62b46..bd89ac8 100644 --- a/config/runtime.exs +++ b/config/runtime.exs @@ -105,15 +105,21 @@ if config_env() == :prod do System.get_env("ATJAMS_OAUTH_KEY_ID") || raise("environment variable ATJAMS_OAUTH_KEY_ID is missing") - config :atex, Atex.OAuth, - base_url: "https://atjams.pdewey.com/oauth", - private_key: oauth_private_key, - key_id: oauth_key_id, - scopes: [ + oauth_base_url = + System.get_env("ATJAMS_OAUTH_BASE_URL") || + "https://atjams.pdewey.com/oauth" + + oauth_scopes = + System.get_env("ATJAMS_OAUTH_SCOPES") || "repo?collection=com.pdewey.atjams.share" <> "&collection=com.pdewey.atjams.like" <> "&collection=com.pdewey.atjams.comment" - ] + + config :atex, Atex.OAuth, + base_url: oauth_base_url, + private_key: oauth_private_key, + key_id: oauth_key_id, + scopes: [oauth_scopes] # ## Configuring the mailer # diff --git a/flake.lock b/flake.lock new file mode 100644 index 0000000..9b03860 --- /dev/null +++ b/flake.lock @@ -0,0 +1,26 @@ +{ + "nodes": { + "nixpkgs": { + "locked": { + "lastModified": 1778672786, + "narHash": "sha256-Blg88K1jwG+P0Mr27+rKMFCufdrWkV3wWh9AdYtz0FQ=", + "owner": "NixOS", + "repo": "nixpkgs", + "rev": "eef00dfd8a712b34af845f9350bac681b1228bd1", + "type": "github" + }, + "original": { + "id": "nixpkgs", + "ref": "nixpkgs-unstable", + "type": "indirect" + } + }, + "root": { + "inputs": { + "nixpkgs": "nixpkgs" + } + } + }, + "root": "root", + "version": 7 +} diff --git a/flake.nix b/flake.nix new file mode 100644 index 0000000..8e4d9be --- /dev/null +++ b/flake.nix @@ -0,0 +1,46 @@ +{ + description = "Atjams flake"; + inputs = { + nixpkgs.url = "nixpkgs/nixpkgs-unstable"; + }; + outputs = + { nixpkgs, self, ... }: + let + forAllSystems = + function: + nixpkgs.lib.genAttrs [ "x86_64-linux" "aarch64-linux" ] ( + system: function nixpkgs.legacyPackages.${system} + ); + in + { + packages = forAllSystems (pkgs: { + atjams = pkgs.callPackage ./nix/default.nix { + beam = pkgs.beam.packages.erlang; + }; + default = self.packages.${pkgs.system}.atjams; + }); + + apps = forAllSystems (pkgs: { + atjams = { + type = "app"; + program = "${self.packages.${pkgs.system}.atjams}/bin/atjams"; + }; + default = self.apps.${pkgs.system}.atjams; + }); + + devShells = forAllSystems (pkgs: { + default = pkgs.mkShell { + packages = with pkgs; [ + beam.packages.erlang.erlang + beam.packages.erlang.elixir + nodejs + ]; + }; + }); + + nixosModules = { + atjams = import ./nix/module.nix; + default = self.nixosModules.atjams; + }; + }; +} diff --git a/lib/atjams/atproto.ex b/lib/atjams/atproto.ex index ecd22d5..565f558 100644 --- a/lib/atjams/atproto.ex +++ b/lib/atjams/atproto.ex @@ -18,6 +18,7 @@ defmodule Atjams.Atproto do alias Atjams.Feed @share_collection "com.pdewey.atjams.share" + @like_collection "com.pdewey.atjams.like" @profile_collection "app.bsky.actor.profile" ## Session helpers @@ -122,6 +123,73 @@ defmodule Atjams.Atproto do end end + @doc """ + Creates a `com.pdewey.atjams.like` record on the user's PDS referencing + the given share (via strongRef: `{uri, cid}`). + """ + def create_like(%{did: did, session_key: session_key}, share_uri, share_cid) + when is_binary(share_uri) do + rkey = generate_rkey() + now = DateTime.utc_now() |> DateTime.truncate(:second) + created_at_iso = DateTime.to_iso8601(now) + + record = %{ + "$type" => @like_collection, + "subject" => maybe_put(%{"uri" => share_uri}, "cid", share_cid), + "createdAt" => created_at_iso + } + + with {:ok, client} <- OAuthClient.new(session_key), + {:ok, %{status: 200, body: body}, _client} <- + XRPC.post(client, "com.atproto.repo.createRecord", + json: %{ + repo: did, + collection: @like_collection, + rkey: rkey, + record: record + } + ) do + {:ok, + %{ + uri: Map.get(body, "uri"), + cid: Map.get(body, "cid"), + did: did, + rkey: rkey, + share_uri: share_uri, + created_at: now + }} + else + {:ok, %{status: status, body: body}, _client} -> + Logger.warning("createRecord (like) failed: status=#{status} body=#{inspect(body)}") + {:error, {:pds_error, status, body}} + + {:error, reason, _client} -> {:error, reason} + {:error, reason} -> {:error, reason} + end + end + + @doc "Deletes a record by collection + rkey on the user's PDS." + def delete_record(%{did: did, session_key: session_key}, collection, rkey) + when is_binary(collection) and is_binary(rkey) do + with {:ok, client} <- OAuthClient.new(session_key), + {:ok, %{status: 200}, _client} <- + XRPC.post(client, "com.atproto.repo.deleteRecord", + json: %{repo: did, collection: collection, rkey: rkey} + ) do + :ok + else + {:ok, %{status: status, body: body}, _client} -> + Logger.warning("deleteRecord failed: status=#{status} body=#{inspect(body)}") + {:error, {:pds_error, status, body}} + + {:error, reason, _client} -> {:error, reason} + {:error, reason} -> {:error, reason} + end + end + + def delete_share(user, rkey), do: delete_record(user, @share_collection, rkey) + def delete_like(user, rkey), do: delete_record(user, @like_collection, rkey) + ## Backfill @doc """ @@ -191,6 +259,66 @@ defmodule Atjams.Atproto do defp index_listed_record(_did, _other), do: :skip + @doc "Backfill likes from a user's PDS." + def backfill_likes(did) when is_binary(did) do + case resolve_identity(did) do + {:ok, %{pds_endpoint: pds, did: did}} -> + Logger.metadata(did: did) + do_backfill_likes(pds, did, nil, 0) + + err -> + Logger.warning("like backfill skipped: #{inspect(err)}") + :error + end + end + + defp do_backfill_likes(pds, did, cursor, count) do + params = + [repo: did, collection: @like_collection, limit: 100] + |> maybe_put_kv(:cursor, cursor) + + case XRPC.unauthed_get(pds, "com.atproto.repo.listRecords", params: params) do + {:ok, %{status: 200, body: %{"records" => records} = body}} -> + Enum.each(records, &index_listed_like(did, &1)) + new_count = count + length(records) + next = Map.get(body, "cursor") + + if is_binary(next) and records != [] do + do_backfill_likes(pds, did, next, new_count) + else + Logger.info("like backfill complete: #{new_count} like(s)") + {:ok, new_count} + end + + {:ok, %{status: status, body: body}} -> + Logger.warning("like backfill failed: status=#{status} body=#{inspect(body)}") + {:error, {:pds_error, status, body}} + + {:error, reason} -> + Logger.warning("like backfill error: #{inspect(reason)}") + {:error, reason} + end + end + + defp index_listed_like(did, %{"uri" => uri, "cid" => cid, "value" => record}) do + rkey = rkey_from_uri(uri) + created_at = parse_datetime(Map.get(record, "createdAt")) + share_uri = get_in(record, ["subject", "uri"]) + + if rkey && created_at && is_binary(share_uri) do + Feed.upsert_like(%{ + uri: uri, + cid: cid, + did: did, + rkey: rkey, + share_uri: share_uri, + created_at: created_at + }) + end + end + + defp index_listed_like(_did, _other), do: :skip + ## Profile refresh @doc """ @@ -238,11 +366,16 @@ defmodule Atjams.Atproto do ## Task helpers - @doc "Kicks off backfill in a supervised Task. Fire-and-forget." + @doc "Kicks off share backfill in a supervised Task. Fire-and-forget." def start_backfill(did) do Task.Supervisor.start_child(Atjams.TaskSupervisor, fn -> backfill_shares(did) end) end + @doc "Kicks off like backfill in a supervised Task. Fire-and-forget." + def start_backfill_likes(did) do + Task.Supervisor.start_child(Atjams.TaskSupervisor, fn -> backfill_likes(did) end) + end + @doc "Kicks off profile refresh in a supervised Task. Fire-and-forget." def start_profile_refresh(did) do Task.Supervisor.start_child(Atjams.TaskSupervisor, fn -> refresh_profile(did) end) diff --git a/lib/atjams/feed.ex b/lib/atjams/feed.ex index 5eca20a..900cec4 100644 --- a/lib/atjams/feed.ex +++ b/lib/atjams/feed.ex @@ -11,7 +11,7 @@ defmodule Atjams.Feed do import Ecto.Query, warn: false alias Atjams.Repo - alias Atjams.Feed.{KnownUser, Profile, Share, JetstreamCursor} + alias Atjams.Feed.{KnownUser, Like, Profile, Share, JetstreamCursor} @pubsub Atjams.PubSub @global_topic "feed:global" @@ -59,17 +59,21 @@ defmodule Atjams.Feed do end @doc """ - Lists recent shares with their (possibly missing) author profile. + Lists recent shares with their (possibly missing) author profile, plus + per-share like counts and a `liked_by_me?` flag if `:current_user_did` + is provided. Options: * `:limit` — default 50 * `:before` — `%DateTime{}` cursor; returns rows strictly older than this * `:did` — restrict to one author + * `:current_user_did` — if given, items get `liked_by_me?` set correctly """ def list_recent(opts \\ []) do limit = Keyword.get(opts, :limit, 50) before = Keyword.get(opts, :before) did = Keyword.get(opts, :did) + current_user_did = Keyword.get(opts, :current_user_did) query = from s in Share, @@ -91,7 +95,30 @@ defmodule Atjams.Feed do true -> from [s, _p] in query, where: s.did == ^did end - Repo.all(query) + items = Repo.all(query) + decorate(items, current_user_did) + end + + @doc "Decorates a single feed-item map with counts + liked_by_me?." + def decorate_item(item, current_user_did) do + [decorated] = decorate([item], current_user_did) + decorated + end + + defp decorate([], _), do: [] + + defp decorate(items, current_user_did) do + uris = Enum.map(items, & &1.share.uri) + like_counts = like_counts_by_uri(uris) + liked_set = liked_share_uris_for_did(uris, current_user_did) + + Enum.map(items, fn item -> + uri = item.share.uri + + item + |> Map.put(:like_count, Map.get(like_counts, uri, 0)) + |> Map.put(:liked_by_me?, MapSet.member?(liked_set, uri)) + end) end def get_share(uri) when is_binary(uri) do @@ -147,6 +174,87 @@ defmodule Atjams.Feed do ) end + ## Likes + + @doc "Inserts a like row by URI. Broadcasts :like_added on the global topic." + def upsert_like(attrs) when is_map(attrs) do + attrs = Map.put_new(attrs, :indexed_at, utc_now()) + + case %Like{} + |> Like.changeset(attrs) + |> Repo.insert( + on_conflict: {:replace_all_except, [:uri]}, + conflict_target: :uri, + returning: true + ) do + {:ok, like} -> + broadcast_like_event(:like_added, like) + {:ok, like} + + err -> + err + end + end + + @doc "Deletes a like by URI." + def delete_like(uri) when is_binary(uri) do + case Repo.get(Like, uri) do + nil -> + :ok + + %Like{} = like -> + {:ok, _} = Repo.delete(like) + broadcast_like_event(:like_removed, like) + :ok + end + end + + # Likes are interesting to three audiences: the global feed, the liker's + # own profile (so they see their like count tick), and the share owner's + # profile (so their card updates for anyone watching). + defp broadcast_like_event(event, like) do + broadcast(@global_topic, {event, like}) + broadcast("feed:did:#{like.did}", {event, like}) + + case did_from_uri(like.share_uri) do + did when is_binary(did) and did != like.did -> + broadcast("feed:did:#{did}", {event, like}) + + _ -> + :ok + end + end + + defp did_from_uri("at://" <> rest) do + case String.split(rest, "/", parts: 2) do + [did | _] -> did + _ -> nil + end + end + + defp did_from_uri(_), do: nil + + def get_like_for(did, share_uri) when is_binary(did) and is_binary(share_uri) do + Repo.one(from l in Like, where: l.did == ^did and l.share_uri == ^share_uri) + end + + def like_counts_by_uri([]), do: %{} + + def like_counts_by_uri(uris) when is_list(uris) do + from(l in Like, where: l.share_uri in ^uris, group_by: l.share_uri, select: {l.share_uri, count(l.uri)}) + |> Repo.all() + |> Map.new() + end + + def liked_share_uris_for_did(_uris, nil), do: MapSet.new() + def liked_share_uris_for_did([], _did), do: MapSet.new() + + def liked_share_uris_for_did(uris, did) when is_list(uris) and is_binary(did) do + from(l in Like, where: l.did == ^did and l.share_uri in ^uris, select: l.share_uri) + |> Repo.all() + |> MapSet.new() + end + ## Profiles def upsert_profile(attrs) when is_map(attrs) do diff --git a/lib/atjams/feed/jetstream_consumer.ex b/lib/atjams/feed/jetstream_consumer.ex index 57a9b02..b0f6732 100644 --- a/lib/atjams/feed/jetstream_consumer.ex +++ b/lib/atjams/feed/jetstream_consumer.ex @@ -1,24 +1,27 @@ defmodule Atjams.Feed.JetstreamConsumer do @moduledoc """ - Drinkup Jetstream consumer filtered to `com.pdewey.atjams.share`. Every - commit/delete is funneled into `Atjams.Feed` for indexing; the cursor is - persisted after each event so a restart resumes where we left off. + Drinkup Jetstream consumer for atjams collections: shares and likes. + Each commit/delete is funneled into `Atjams.Feed`; the cursor is persisted + after every event so a restart resumes where we left off. """ require Logger alias Atjams.{Atproto, Enrichment, Feed} - @collection "com.pdewey.atjams.share" + @share_collection "com.pdewey.atjams.share" + @like_collection "com.pdewey.atjams.like" use Drinkup.Jetstream, name: :atjams_jetstream, - wanted_collections: [@collection] + wanted_collections: [@share_collection, @like_collection] + + ## Shares @impl true def handle_event(%Drinkup.Jetstream.Event.Commit{ operation: op, - collection: @collection, + collection: @share_collection, did: did, rkey: rkey, cid: cid, @@ -26,33 +29,63 @@ defmodule Atjams.Feed.JetstreamConsumer do time_us: t }) when op in [:create, :update] do - upsert_from_commit(did, rkey, cid, record) + upsert_share_from_commit(did, rkey, cid, record) Feed.set_cursor(t) :ok end def handle_event(%Drinkup.Jetstream.Event.Commit{ operation: :delete, - collection: @collection, + collection: @share_collection, did: did, rkey: rkey, time_us: t }) do - uri = build_uri(did, rkey) - Feed.delete_share(uri) + Feed.delete_share(share_uri(did, rkey)) + Feed.set_cursor(t) + :ok + end + + ## Likes + + def handle_event(%Drinkup.Jetstream.Event.Commit{ + operation: op, + collection: @like_collection, + did: did, + rkey: rkey, + cid: cid, + record: record, + time_us: t + }) + when op in [:create, :update] do + upsert_like_from_commit(did, rkey, cid, record) + Feed.set_cursor(t) + :ok + end + + def handle_event(%Drinkup.Jetstream.Event.Commit{ + operation: :delete, + collection: @like_collection, + did: did, + rkey: rkey, + time_us: t + }) do + Feed.delete_like(like_uri(did, rkey)) Feed.set_cursor(t) :ok end def handle_event(_other), do: :ok - defp upsert_from_commit(did, rkey, cid, record) when is_map(record) do + ## Share helpers + + defp upsert_share_from_commit(did, rkey, cid, record) when is_map(record) do song_url = Map.get(record, "songUrl") created_at = parse_datetime(Map.get(record, "createdAt")) if is_binary(song_url) and created_at do attrs = %{ - uri: build_uri(did, rkey), + uri: share_uri(did, rkey), cid: cid, did: did, rkey: rkey, @@ -69,14 +102,42 @@ defmodule Atjams.Feed.JetstreamConsumer do Enrichment.start(share.uri) {:error, reason} -> - Logger.warning("jetstream upsert failed: #{inspect(reason)} attrs=#{inspect(attrs)}") + Logger.warning("jetstream share upsert failed: #{inspect(reason)}") end else Logger.debug("skipping malformed share record for #{did}/#{rkey}") end end - defp upsert_from_commit(_did, _rkey, _cid, _other), do: :ok + defp upsert_share_from_commit(_did, _rkey, _cid, _other), do: :ok + + ## Like helpers + + defp upsert_like_from_commit(did, rkey, cid, record) when is_map(record) do + created_at = parse_datetime(Map.get(record, "createdAt")) + target_uri = get_in(record, ["subject", "uri"]) + + if is_binary(target_uri) and created_at do + Feed.upsert_like(%{ + uri: like_uri(did, rkey), + cid: cid, + did: did, + rkey: rkey, + share_uri: target_uri, + created_at: created_at + }) + |> case do + {:ok, _like} -> maybe_fetch_profile(did) + {:error, reason} -> Logger.debug("jetstream like upsert skipped: #{inspect(reason)}") + end + else + Logger.debug("skipping malformed like record for #{did}/#{rkey}") + end + end + + defp upsert_like_from_commit(_did, _rkey, _cid, _other), do: :ok + + ## Shared defp maybe_fetch_profile(did) do if is_nil(Feed.get_profile(did)) do @@ -86,7 +147,8 @@ defmodule Atjams.Feed.JetstreamConsumer do :ok end - defp build_uri(did, rkey), do: "at://#{did}/#{@collection}/#{rkey}" + defp share_uri(did, rkey), do: "at://#{did}/#{@share_collection}/#{rkey}" + defp like_uri(did, rkey), do: "at://#{did}/#{@like_collection}/#{rkey}" defp parse_datetime(nil), do: nil diff --git a/lib/atjams/feed/like.ex b/lib/atjams/feed/like.ex new file mode 100644 index 0000000..1be9d55 --- /dev/null +++ b/lib/atjams/feed/like.ex @@ -0,0 +1,23 @@ +defmodule Atjams.Feed.Like do + use Ecto.Schema + import Ecto.Changeset + + @primary_key {:uri, :string, autogenerate: false} + schema "likes" do + field :cid, :string + field :did, :string + field :rkey, :string + field :share_uri, :string + field :created_at, :utc_datetime + field :indexed_at, :utc_datetime + end + + @cast_fields ~w(uri cid did rkey share_uri created_at indexed_at)a + @required ~w(uri did rkey share_uri created_at indexed_at)a + + def changeset(like, attrs) do + like + |> cast(attrs, @cast_fields) + |> validate_required(@required) + end +end diff --git a/lib/atjams/release.ex b/lib/atjams/release.ex new file mode 100644 index 0000000..d0c799c --- /dev/null +++ b/lib/atjams/release.ex @@ -0,0 +1,30 @@ +defmodule Atjams.Release do + @moduledoc """ + Used for executing DB release tasks when run in production without Mix + installed. + """ + @app :atjams + + def migrate do + load_app() + + for repo <- repos() do + {:ok, _, _} = Ecto.Migrator.with_repo(repo, &Ecto.Migrator.run(&1, :up, all: true)) + end + end + + def rollback(repo, version) do + load_app() + {:ok, _, _} = Ecto.Migrator.with_repo(repo, &Ecto.Migrator.run(&1, :down, to: version)) + end + + defp repos do + Application.fetch_env!(@app, :ecto_repos) + end + + defp load_app do + # Many platforms require SSL when connecting to the database + Application.ensure_all_started(:ssl) + Application.ensure_loaded(@app) + end +end diff --git a/lib/atjams_web/components/feed_components.ex b/lib/atjams_web/components/feed_components.ex index 3389e03..eeffdbe 100644 --- a/lib/atjams_web/components/feed_components.ex +++ b/lib/atjams_web/components/feed_components.ex @@ -26,6 +26,7 @@ defmodule AtjamsWeb.FeedComponents do } attr :item, :map, required: true + attr :current_user, :map, default: nil def share_card(assigns) do ~H""" @@ -57,10 +58,100 @@ defmodule AtjamsWeb.FeedComponents do

{@item.share.notes}

<% end %> + + <.actions_bar item={@item} current_user={@current_user} /> """ end + attr :item, :map, required: true + attr :current_user, :map, default: nil + + def actions_bar(assigns) do + assigns = + assign_new(assigns, :is_owner?, fn -> + assigns.current_user != nil and assigns.current_user.did == assigns.item.share.did + end) + |> assign_new(:can_like?, fn -> assigns.current_user != nil end) + |> assign_new(:like_count, fn -> Map.get(assigns.item, :like_count, 0) end) + |> assign_new(:liked_by_me?, fn -> Map.get(assigns.item, :liked_by_me?, false) end) + + ~H""" +
+ + + <.link + navigate={~p"/profile/#{display_identifier(@item)}/#{@item.share.rkey}"} + class="inline-flex items-center gap-1.5 px-2 py-1 rounded-lg text-xs text-base-content/60 hover:text-base-content hover:bg-base-200 transition-colors" + title="Comments" + > + <.icon name="hero-chat-bubble-oval-left" class="size-4" /> + + + +
+ """ + end + + defp like_button_classes(true, _), + do: "bg-rose-500/10 text-rose-600 hover:bg-rose-500/20" + + defp like_button_classes(false, true), + do: "text-base-content/60 hover:text-rose-600 hover:bg-base-200" + + defp like_button_classes(false, false), + do: "text-base-content/40 cursor-not-allowed" + + defp like_title(true, _), do: "Unlike" + defp like_title(false, true), do: "Like" + defp like_title(false, false), do: "Sign in to like" + attr :links, :any, default: nil attr :source_url, :string, default: nil diff --git a/lib/atjams_web/controllers/oauth_controller.ex b/lib/atjams_web/controllers/oauth_controller.ex index 070cc44..0b54c1a 100644 --- a/lib/atjams_web/controllers/oauth_controller.ex +++ b/lib/atjams_web/controllers/oauth_controller.ex @@ -16,6 +16,7 @@ defmodule AtjamsWeb.OAuthController do Feed.upsert_known_user(did, nil) Atproto.start_profile_refresh(did) Atproto.start_backfill(did) + Atproto.start_backfill_likes(did) end conn diff --git a/lib/atjams_web/live/feed_live.ex b/lib/atjams_web/live/feed_live.ex index 1ae7977..991d7be 100644 --- a/lib/atjams_web/live/feed_live.ex +++ b/lib/atjams_web/live/feed_live.ex @@ -9,7 +9,8 @@ defmodule AtjamsWeb.FeedLive do def mount(_params, _session, socket) do if connected?(socket), do: Feed.subscribe_global() - items = Feed.list_recent(limit: @page_size) + current_user_did = current_user_did(socket) + items = Feed.list_recent(limit: @page_size, current_user_did: current_user_did) {:ok, socket @@ -19,10 +20,12 @@ defmodule AtjamsWeb.FeedLive do |> assign(:cursor, cursor_from(items))} end + ## PubSub + @impl true def handle_info({:share_upserted, share}, socket) do profile = Feed.get_profile(share.did) - item = %{share: share, profile: profile} + item = build_item(socket, share, profile) {:noreply, socket @@ -36,11 +39,18 @@ defmodule AtjamsWeb.FeedLive do def handle_info({:share_links_updated, share}, socket) do profile = Feed.get_profile(share.did) - {:noreply, stream_insert(socket, :items, %{share: share, profile: profile})} + {:noreply, stream_insert(socket, :items, build_item(socket, share, profile))} + end + + def handle_info({like_event, %{share_uri: share_uri}}, socket) + when like_event in [:like_added, :like_removed] do + update_share_card(socket, share_uri) end def handle_info(_msg, socket), do: {:noreply, socket} + ## Events + @impl true def handle_event("load_more", _params, socket) do case socket.assigns.cursor do @@ -48,7 +58,12 @@ defmodule AtjamsWeb.FeedLive do {:noreply, socket} cursor -> - more = Feed.list_recent(limit: @page_size, before: cursor) + more = + Feed.list_recent( + limit: @page_size, + before: cursor, + current_user_did: current_user_did(socket) + ) {:noreply, socket @@ -57,12 +72,42 @@ defmodule AtjamsWeb.FeedLive do end end + def handle_event("toggle_like", %{"uri" => share_uri}, socket) do + {:noreply, AtjamsWeb.ShareActions.toggle_like(socket, share_uri)} + end + + def handle_event("delete_share", %{"uri" => share_uri}, socket) do + {:noreply, AtjamsWeb.ShareActions.delete_share(socket, share_uri)} + end + + def handle_event("edit_share", _params, socket) do + {:noreply, put_flash(socket, :info, "Editing shares isn't built yet — delete and re-share for now.")} + end + + ## Helpers + + defp current_user_did(%{assigns: %{current_user: %{did: did}}}), do: did + defp current_user_did(_), do: nil + + defp build_item(socket, share, profile) do + Feed.decorate_item(%{share: share, profile: profile}, current_user_did(socket)) + end + + defp update_share_card(socket, share_uri) do + case Feed.get_share(share_uri) do + nil -> + {:noreply, socket} + + share -> + profile = Feed.get_profile(share.did) + {:noreply, stream_insert(socket, :items, build_item(socket, share, profile))} + end + end + defp cursor_from([]), do: nil defp cursor_from(items) do - items - |> List.last() - |> case do + case List.last(items) do %{share: %{created_at: ts}} -> ts _ -> nil end diff --git a/lib/atjams_web/live/feed_live.html.heex b/lib/atjams_web/live/feed_live.html.heex index 6ed98d6..4c038a6 100644 --- a/lib/atjams_web/live/feed_live.html.heex +++ b/lib/atjams_web/live/feed_live.html.heex @@ -27,7 +27,7 @@
- <.share_card item={item} /> + <.share_card item={item} current_user={@current_user} />
diff --git a/lib/atjams_web/live/profile_live.ex b/lib/atjams_web/live/profile_live.ex index f8069d7..89d47ed 100644 --- a/lib/atjams_web/live/profile_live.ex +++ b/lib/atjams_web/live/profile_live.ex @@ -11,7 +11,12 @@ defmodule AtjamsWeb.ProfileLive do {:ok, profile_did, profile} -> if connected?(socket), do: Feed.subscribe_did(profile_did) - items = Feed.list_recent(limit: @page_size, did: profile_did) + items = + Feed.list_recent( + limit: @page_size, + did: profile_did, + current_user_did: current_user_did(socket) + ) {:ok, socket @@ -33,12 +38,9 @@ defmodule AtjamsWeb.ProfileLive do @impl true def handle_info({:share_upserted, share}, socket) do if share.did == socket.assigns.profile_did do - profile = socket.assigns.profile || Feed.get_profile(share.did) - item = %{share: share, profile: profile} - {:noreply, socket - |> stream_insert(:items, item, at: 0) + |> stream_insert(:items, build_item(socket, share), at: 0) |> assign(:has_items?, true)} else {:noreply, socket} @@ -51,13 +53,23 @@ defmodule AtjamsWeb.ProfileLive do def handle_info({:share_links_updated, share}, socket) do if share.did == socket.assigns.profile_did do - profile = socket.assigns.profile || Feed.get_profile(share.did) - {:noreply, stream_insert(socket, :items, %{share: share, profile: profile})} + {:noreply, stream_insert(socket, :items, build_item(socket, share))} else {:noreply, socket} end end + def handle_info({like_event, %{share_uri: share_uri}}, socket) + when like_event in [:like_added, :like_removed] do + case Feed.get_share(share_uri) do + %{did: did} = share when did == socket.assigns.profile_did -> + {:noreply, stream_insert(socket, :items, build_item(socket, share))} + + _ -> + {:noreply, socket} + end + end + def handle_info(_msg, socket), do: {:noreply, socket} @impl true @@ -71,7 +83,8 @@ defmodule AtjamsWeb.ProfileLive do Feed.list_recent( limit: @page_size, before: cursor, - did: socket.assigns.profile_did + did: socket.assigns.profile_did, + current_user_did: current_user_did(socket) ) {:noreply, @@ -81,6 +94,28 @@ defmodule AtjamsWeb.ProfileLive do end end + def handle_event("toggle_like", %{"uri" => share_uri}, socket) do + {:noreply, AtjamsWeb.ShareActions.toggle_like(socket, share_uri)} + end + + def handle_event("delete_share", %{"uri" => share_uri}, socket) do + {:noreply, AtjamsWeb.ShareActions.delete_share(socket, share_uri)} + end + + def handle_event("edit_share", _params, socket) do + {:noreply, put_flash(socket, :info, "Editing shares isn't built yet — delete and re-share for now.")} + end + + ## Helpers + + defp current_user_did(%{assigns: %{current_user: %{did: did}}}), do: did + defp current_user_did(_), do: nil + + defp build_item(socket, share) do + profile = socket.assigns.profile || Feed.get_profile(share.did) + Feed.decorate_item(%{share: share, profile: profile}, current_user_did(socket)) + end + defp resolve("did:" <> _ = did) do profile = Feed.get_profile(did) if profile == nil, do: Atproto.start_profile_refresh(did) diff --git a/lib/atjams_web/live/profile_live.html.heex b/lib/atjams_web/live/profile_live.html.heex index de794c3..243784b 100644 --- a/lib/atjams_web/live/profile_live.html.heex +++ b/lib/atjams_web/live/profile_live.html.heex @@ -17,7 +17,7 @@
- <.share_card item={item} /> + <.share_card item={item} current_user={@current_user} />
diff --git a/lib/atjams_web/live/share_live.ex b/lib/atjams_web/live/share_live.ex index b3ba1e9..b8763cc 100644 --- a/lib/atjams_web/live/share_live.ex +++ b/lib/atjams_web/live/share_live.ex @@ -9,12 +9,15 @@ defmodule AtjamsWeb.ShareLive do {:ok, did} -> case Feed.get_share_by_did_and_rkey(did, rkey) do %_{} = share -> + if connected?(socket), do: Feed.subscribe_did(did) + profile = Feed.get_profile(did) - item = %{share: share, profile: profile} + item = build_item(socket, share, profile) {:ok, socket |> assign(:page_title, page_title(share, profile)) + |> assign(:profile_did, did) |> assign(:item, item)} nil -> @@ -32,6 +35,54 @@ defmodule AtjamsWeb.ShareLive do end end + @impl true + def handle_info({event, _}, socket) + when event in [:like_added, :like_removed, :share_links_updated, :share_upserted] do + case Feed.get_share(socket.assigns.item.share.uri) do + nil -> + {:noreply, socket} + + share -> + profile = Feed.get_profile(share.did) + {:noreply, assign(socket, :item, build_item(socket, share, profile))} + end + end + + def handle_info({:share_deleted, uri}, socket) do + if uri == socket.assigns.item.share.uri do + {:noreply, + socket + |> put_flash(:info, "That share was deleted.") + |> push_navigate(to: ~p"/feed")} + else + {:noreply, socket} + end + end + + def handle_info(_msg, socket), do: {:noreply, socket} + + @impl true + def handle_event("toggle_like", %{"uri" => share_uri}, socket) do + {:noreply, AtjamsWeb.ShareActions.toggle_like(socket, share_uri)} + end + + def handle_event("delete_share", %{"uri" => share_uri}, socket) do + {:noreply, AtjamsWeb.ShareActions.delete_share(socket, share_uri)} + end + + def handle_event("edit_share", _params, socket) do + {:noreply, put_flash(socket, :info, "Editing shares isn't built yet — delete and re-share for now.")} + end + + ## Helpers + + defp current_user_did(%{assigns: %{current_user: %{did: did}}}), do: did + defp current_user_did(_), do: nil + + defp build_item(socket, share, profile) do + Feed.decorate_item(%{share: share, profile: profile}, current_user_did(socket)) + end + defp resolve_did("did:" <> _ = did), do: {:ok, did} defp resolve_did(handle) do diff --git a/lib/atjams_web/live/share_live.html.heex b/lib/atjams_web/live/share_live.html.heex index 5af40ed..674f76d 100644 --- a/lib/atjams_web/live/share_live.html.heex +++ b/lib/atjams_web/live/share_live.html.heex @@ -6,7 +6,7 @@ <.icon name="hero-arrow-left" class="size-4" /> Back to feed - <.share_card item={@item} /> + <.share_card item={@item} current_user={@current_user} />

{@item.share.uri} diff --git a/lib/atjams_web/share_actions.ex b/lib/atjams_web/share_actions.ex new file mode 100644 index 0000000..6c39800 --- /dev/null +++ b/lib/atjams_web/share_actions.ex @@ -0,0 +1,84 @@ +defmodule AtjamsWeb.ShareActions do + @moduledoc """ + Shared LiveView side-effect helpers for like-toggle and share-delete. + Each function takes a socket, performs the PDS call + local update, and + returns an updated socket (already containing any flash messages). + """ + + import Phoenix.LiveView, only: [put_flash: 3] + + alias Atjams.{Atproto, Feed} + + def toggle_like(socket, share_uri) do + case socket.assigns[:current_user] do + nil -> + put_flash(socket, :error, "Sign in to like songs.") + + %{did: did} = user -> + case Feed.get_like_for(did, share_uri) do + nil -> do_like(socket, user, share_uri) + like -> do_unlike(socket, user, like) + end + end + end + + defp do_like(socket, user, share_uri) do + share = Feed.get_share(share_uri) + + case share do + nil -> + put_flash(socket, :error, "That share isn't in the index.") + + %{cid: cid} -> + case Atproto.create_like(user, share_uri, cid) do + {:ok, attrs} -> + {:ok, _} = Feed.upsert_like(attrs) + socket + + {:error, reason} -> + put_flash(socket, :error, "Couldn't like that: #{format_reason(reason)}") + end + end + end + + defp do_unlike(socket, user, like) do + case Atproto.delete_like(user, like.rkey) do + :ok -> + Feed.delete_like(like.uri) + socket + + {:error, reason} -> + put_flash(socket, :error, "Couldn't unlike: #{format_reason(reason)}") + end + end + + def delete_share(socket, share_uri) do + case {socket.assigns[:current_user], Feed.get_share(share_uri)} do + {nil, _} -> + put_flash(socket, :error, "Sign in to delete.") + + {_user, nil} -> + put_flash(socket, :error, "That share isn't in the index.") + + {%{did: user_did} = user, %{did: share_did, rkey: rkey}} when user_did == share_did -> + case Atproto.delete_share(user, rkey) do + :ok -> + Feed.delete_share(share_uri) + put_flash(socket, :info, "Share deleted.") + + {:error, reason} -> + put_flash(socket, :error, "Couldn't delete: #{format_reason(reason)}") + end + + _ -> + put_flash(socket, :error, "You can only delete your own shares.") + end + end + + defp format_reason({:pds_error, status, %{"message" => msg}}) when is_binary(msg) do + "PDS #{status}: #{msg}" + end + + defp format_reason({:pds_error, status, _}), do: "PDS #{status}" + defp format_reason(reason), do: inspect(reason) +end diff --git a/nix/default.nix b/nix/default.nix new file mode 100644 index 0000000..0758c4a --- /dev/null +++ b/nix/default.nix @@ -0,0 +1,30 @@ +{ + lib, + beam, + elixir, + nodejs, + appName ? "atjams", +}: + +beam.mixRelease { + pname = appName; + version = "0.1.0"; + src = ../.; + mixEnv = "prod"; + + nativeBuildInputs = [ + nodejs + ]; + + preBuild = '' + mix assets.deploy + ''; + + meta = with lib; { + description = "${appName} — AT Protocol music sharing app"; + homepage = "https://github.com/pdewey/atjams"; + license = licenses.mit; + platforms = platforms.linux; + mainProgram = appName; + }; +} diff --git a/nix/module.nix b/nix/module.nix new file mode 100644 index 0000000..e21d61a --- /dev/null +++ b/nix/module.nix @@ -0,0 +1,232 @@ +{ + config, + lib, + pkgs, + ... +}: + +let + cfg = config.services.atjams; + + # The OAuth base URL is derived from the configured host/scheme + oauthBaseUrl = "https://${cfg.settings.host}/oauth"; +in +{ + options.services.atjams = { + enable = lib.mkEnableOption "Atjams AT Protocol music sharing service"; + + package = lib.mkOption { + type = lib.types.package; + default = pkgs.callPackage ./default.nix { }; + defaultText = lib.literalExpression "pkgs.callPackage ./nix/default.nix { }"; + description = "The atjams package to use."; + }; + + settings = { + host = lib.mkOption { + type = lib.types.str; + description = '' + Public hostname of the server. Used for `PHX_HOST` and OAuth + callback URLs. This must match the domain users will reach the + service at. + ''; + example = "atjams.pdewey.com"; + }; + + port = lib.mkOption { + type = lib.types.port; + default = 4000; + description = "Port on which the atjams Phoenix server listens (internal)."; + }; + + poolSize = lib.mkOption { + type = lib.types.ints.positive; + default = 10; + description = '' + Ecto connection pool size for the SQLite database. + Each concurrent request may need a connection, and JetStream + ingestion also uses the pool. Raise this if you see pool + timeouts under load. + ''; + }; + + logLevel = lib.mkOption { + type = lib.types.enum [ + "debug" + "info" + "warn" + "error" + ]; + default = "info"; + description = "Elixir Logger level."; + }; + }; + + oauth = { + keyId = lib.mkOption { + type = lib.types.str; + description = '' + OAuth key identifier for DPoP-bound requests to AT Protocol PDSes. + This is a short string you choose (e.g. "atjams-xxxx") that + gets registered with users' PDSes on first authorization. + ''; + example = "atjams-a1b2c3d4"; + }; + + scopes = lib.mkOption { + type = lib.types.listOf lib.types.str; + default = [ + "repo?collection=com.pdewey.atjams.share" + "&collection=com.pdewey.atjams.like" + "&collection=com.pdewey.atjams.comment" + ]; + description = "AT Protocol OAuth scopes the app requests."; + }; + + # Note: private key comes via environmentFiles — see below + }; + + environmentFiles = lib.mkOption { + type = lib.types.listOf lib.types.path; + default = [ ]; + description = '' + List of environment files to load into the systemd service. + Use these for secrets that must not end up in the Nix store: + + • `SECRET_KEY_BASE` — Phoenix cookie/session signing key. + Generate with: `mix phx.gen.secret` + + • `ATJAMS_OAUTH_PRIVATE_KEY` — ECDSA P-256 private key in PEM + format, used for DPoP-bound OAuth requests to users' PDSes. + Generate with: + openssl ecparam -genkey -name prime256v1 -noout + + Each file should contain lines like `KEY=value`. + ''; + example = lib.literalExpression ''[ "/run/secrets/atjams.env" ]''; + }; + + dataDir = lib.mkOption { + type = lib.types.path; + default = "/var/lib/atjams"; + description = '' + Directory where atjams stores persistent data: + + • `atjams.db` — SQLite database (shares, likes, profiles, users) + • OAuth session data (ETS-backed, ephemeral across restarts) + ''; + }; + + user = lib.mkOption { + type = lib.types.str; + default = "atjams"; + description = "User account under which atjams runs."; + }; + + group = lib.mkOption { + type = lib.types.str; + default = "atjams"; + description = "Group under which atjams runs."; + }; + + openFirewall = lib.mkOption { + type = lib.types.bool; + default = false; + description = '' + Whether to open the firewall for the atjams port. + Only needed if you run atjams exposed directly (without a + reverse proxy like Nginx). + ''; + }; + }; + + config = lib.mkIf cfg.enable { + users.users.${cfg.user} = lib.mkIf (cfg.user == "atjams") { + isSystemUser = true; + group = cfg.group; + description = "Atjams service user"; + home = cfg.dataDir; + createHome = true; + }; + + users.groups.${cfg.group} = lib.mkIf (cfg.group == "atjams") { }; + + systemd.services.atjams = { + description = "Atjams AT Protocol Music Sharing Service"; + wantedBy = [ "multi-user.target" ]; + after = [ "network.target" ]; + + serviceConfig = { + Type = "simple"; + User = cfg.user; + Group = cfg.group; + ExecStart = "${cfg.package}/bin/atjams start"; + Restart = "on-failure"; + RestartSec = "10s"; + + EnvironmentFile = cfg.environmentFiles; + + # Security hardening + NoNewPrivileges = true; + PrivateTmp = true; + ProtectSystem = "strict"; + ProtectHome = true; + ReadWritePaths = [ cfg.dataDir ]; + ProtectKernelTunables = true; + ProtectKernelModules = true; + ProtectControlGroups = true; + RestrictAddressFamilies = [ + "AF_INET" + "AF_INET6" + "AF_UNIX" + ]; + RestrictNamespaces = true; + LockPersonality = true; + RestrictRealtime = true; + RestrictSUIDSGID = true; + MemoryDenyWriteExecute = true; + SystemCallArchitectures = "native"; + CapabilityBoundingSet = ""; + }; + + environment = { + # Phoenix server mode — required for mix release + PHX_SERVER = "true"; + + # Hostname — used for URL generation and OAuth redirects + PHX_HOST = cfg.settings.host; + + # Internal HTTP port (reverse proxy terminates TLS) + PORT = toString cfg.settings.port; + + # SQLite database path + DATABASE_PATH = "${cfg.dataDir}/atjams.db"; + + # Ecto connection pool size + POOL_SIZE = toString cfg.settings.poolSize; + + # Logger level + LOG_LEVEL = cfg.settings.logLevel; + + # AT Protocol OAuth configuration + ATJAMS_OAUTH_KEY_ID = cfg.oauth.keyId; + ATJAMS_OAUTH_BASE_URL = oauthBaseUrl; + }; + + # Extra hardening: the OAuth private key env file should be + # restricted to the atjams user. We log a warning at eval time + # if no environmentFiles are provided, since ATJAMS_OAUTH_PRIVATE_KEY + # and SECRET_KEY_BASE are required at runtime. + }; + + networking.firewall.allowedTCPPorts = lib.optional cfg.openFirewall cfg.settings.port; + + # Provide a convenience assertion + warnings = lib.optional (cfg.enable && cfg.environmentFiles == [ ]) '' + Atjams: no environmentFiles configured. You must provide + SECRET_KEY_BASE and ATJAMS_OAUTH_PRIVATE_KEY at runtime. + Either set services.atjams.environmentFiles or set them + globally via the systemd EnvironmentFile mechanism. + ''; + }; +} diff --git a/priv/repo/migrations/20260515000000_create_likes.exs b/priv/repo/migrations/20260515000000_create_likes.exs new file mode 100644 index 0000000..eee4cb7 --- /dev/null +++ b/priv/repo/migrations/20260515000000_create_likes.exs @@ -0,0 +1,18 @@ +defmodule Atjams.Repo.Migrations.CreateLikes do + use Ecto.Migration + + def change do + create table(:likes, primary_key: false) do + add :uri, :string, primary_key: true + add :cid, :string + add :did, :string, null: false + add :rkey, :string, null: false + add :share_uri, :string, null: false + add :created_at, :utc_datetime, null: false + add :indexed_at, :utc_datetime, null: false + end + + create index(:likes, [:share_uri]) + create unique_index(:likes, [:did, :share_uri]) + end +end diff --git a/rel/overlays/bin/migrate b/rel/overlays/bin/migrate new file mode 100755 index 0000000..b66ada1 --- /dev/null +++ b/rel/overlays/bin/migrate @@ -0,0 +1,5 @@ +#!/bin/sh +set -eu + +cd -P -- "$(dirname -- "$0")" +exec ./atjams eval Atjams.Release.migrate diff --git a/rel/overlays/bin/migrate.bat b/rel/overlays/bin/migrate.bat new file mode 100755 index 0000000..8e5e1cd --- /dev/null +++ b/rel/overlays/bin/migrate.bat @@ -0,0 +1 @@ +call "%~dp0\atjams" eval Atjams.Release.migrate diff --git a/rel/overlays/bin/server b/rel/overlays/bin/server new file mode 100755 index 0000000..9966555 --- /dev/null +++ b/rel/overlays/bin/server @@ -0,0 +1,5 @@ +#!/bin/sh +set -eu + +cd -P -- "$(dirname -- "$0")" +PHX_SERVER=true exec ./atjams start diff --git a/rel/overlays/bin/server.bat b/rel/overlays/bin/server.bat new file mode 100755 index 0000000..9d010c7 --- /dev/null +++ b/rel/overlays/bin/server.bat @@ -0,0 +1,2 @@ +set PHX_SERVER=true +call "%~dp0\atjams" start