diff --git a/lib/statusphere/application.ex b/lib/statusphere/application.ex index 106f202..b8a66ef 100644 --- a/lib/statusphere/application.ex +++ b/lib/statusphere/application.ex @@ -14,6 +14,7 @@ defmodule Statusphere.Application do repos: Application.fetch_env!(:statusphere, :ecto_repos), skip: skip_migrations?()}, {DNSCluster, query: Application.get_env(:statusphere, :dns_cluster_query) || :ignore}, {Phoenix.PubSub, name: Statusphere.PubSub}, + {Drinkup, %{consumer: Statusphere.Consumer}}, # Start a worker by calling: Statusphere.Worker.start_link(arg) # {Statusphere.Worker, arg}, # Start to serve requests, typically the last entry diff --git a/lib/statusphere/consumer.ex b/lib/statusphere/consumer.ex new file mode 100644 index 0000000..53f38dd --- /dev/null +++ b/lib/statusphere/consumer.ex @@ -0,0 +1,46 @@ +defmodule Statusphere.Consumer do + alias Statusphere.Repo + require Logger + use Drinkup.RecordConsumer, collections: ["xyz.statusphere.status"] + + def handle_create(record), do: upsert(record) + + def handle_update(record), do: upsert(record) + + def handle_delete(record) do + IO.inspect(record, label: "delete") + end + + defp upsert(%{type: "xyz.statusphere.status", record: record} = evt) do + case Xyz.Statusphere.Status.from_json(record) do + {:ok, record} -> + uri = + Atex.AtURI.to_string(%Atex.AtURI{ + authority: evt.did, + collection: evt.type, + rkey: evt.rkey + }) + + status = + %Statusphere.Status{} + |> Statusphere.Status.changeset(%{ + uri: uri, + author_did: evt.did, + status: record.status, + created_at: NaiveDateTime.from_iso8601!(record.createdAt), + indexed_at: NaiveDateTime.utc_now() + }) + |> Repo.insert!( + on_conflict: [set: [status: record.status, indexed_at: NaiveDateTime.utc_now()]], + conflict_target: :uri + ) + + Logger.debug("ingested status: #{inspect(status)}") + + _ -> + nil + end + end + + defp upsert(_), do: nil +end diff --git a/lib/statusphere/status.ex b/lib/statusphere/status.ex new file mode 100644 index 0000000..ea9f4ae --- /dev/null +++ b/lib/statusphere/status.ex @@ -0,0 +1,21 @@ +defmodule Statusphere.Status do + use Ecto.Schema + import Ecto.Changeset + + @primary_key {:uri, :binary_id, autogenerate: false} + @foreign_key_type :binary_id + schema "status" do + field :author_did, :string + field :status, :string + field :created_at, :utc_datetime + field :indexed_at, :utc_datetime + end + + @doc false + def changeset(status, attrs) do + status + |> cast(attrs, [:uri, :author_did, :status, :created_at, :indexed_at]) + |> validate_required([:uri, :author_did, :status, :created_at, :indexed_at]) + |> unique_constraint(:uri) + end +end diff --git a/lib/statusphere_web/controllers/page_controller.ex b/lib/statusphere_web/controllers/page_controller.ex index c6836f8..14f1244 100644 --- a/lib/statusphere_web/controllers/page_controller.ex +++ b/lib/statusphere_web/controllers/page_controller.ex @@ -1,26 +1,65 @@ defmodule StatusphereWeb.PageController do + alias Statusphere.Repo + alias Statusphere.Status + require Logger + import Ecto.Query, only: [from: 2] use StatusphereWeb, :controller def home(conn, _params) do - case Atex.XRPC.OAuthClient.from_conn(conn) do - {:ok, client} -> - {:ok, %{body: %{value: profile}}, client} = - Atex.XRPC.get(client, %Com.Atproto.Repo.GetRecord{ - params: %{ - repo: client.did, - collection: "app.bsky.actor.profile", - rkey: "self" - } - }) - - selected_status = get_session(conn, :selected_status) + query = + from u in Status, + order_by: [desc: u.indexed_at], + limit: 10 + + statuses = Repo.all(query) + + status_identities = + statuses + |> Enum.map(fn %{author_did: did} -> + case Atex.IdentityResolver.resolve(did) do + {:ok, identity} -> {identity.did, identity.handle} + {:error, _} -> {did, did} + end + end) + |> Enum.into(%{}) + + with {:ok, client} <- Atex.XRPC.OAuthClient.from_conn(conn), + {:ok, %{body: %{value: profile}}, client} <- + Atex.XRPC.get(client, %Com.Atproto.Repo.GetRecord{ + params: %{ + repo: client.did, + collection: "app.bsky.actor.profile", + rkey: "self" + } + }) do + conn + |> Atex.XRPC.OAuthClient.update_plug(client) + |> render(:home, + profile: profile, + did: client.did, + statuses: statuses, + status_identities: status_identities + ) + else + {:error, err, client} -> + Logger.error("Failed to fetch Bluesky profile for #{client.did}: #{inspect(err)}") conn |> Atex.XRPC.OAuthClient.update_plug(client) - |> render(:home, profile: profile, conn: conn, selected_status: selected_status) + |> render(:home, + profile: %{}, + did: client.did, + statuses: statuses, + status_identities: status_identities + ) - :error -> - render(conn, :home, profile: nil, conn: conn, selected_status: nil) + _err -> + render(conn, :home, + profile: nil, + did: nil, + statuses: statuses, + status_identities: status_identities + ) end end diff --git a/lib/statusphere_web/controllers/page_html/home.html.heex b/lib/statusphere_web/controllers/page_html/home.html.heex index 37de556..5113052 100644 --- a/lib/statusphere_web/controllers/page_html/home.html.heex +++ b/lib/statusphere_web/controllers/page_html/home.html.heex @@ -22,7 +22,28 @@ <% end %> -
- +
+ <% user_latest_status = + Enum.find(@statuses, %{status: ""}, fn x -> x.author_did == @did end) %> +
+ + diff --git a/lib/statusphere_web/controllers/status_controller.ex b/lib/statusphere_web/controllers/status_controller.ex index 3f607ef..1425bbd 100644 --- a/lib/statusphere_web/controllers/status_controller.ex +++ b/lib/statusphere_web/controllers/status_controller.ex @@ -1,8 +1,11 @@ defmodule StatusphereWeb.StatusController do + alias Statusphere.Repo require Logger use StatusphereWeb, :controller def create(conn, %{"status" => status}) do + rkey = to_string(Atex.TID.now()) + with {:ok, client} <- Atex.XRPC.OAuthClient.from_conn(conn), {:ok, record} <- Xyz.Statusphere.Status.main(%{ @@ -16,14 +19,39 @@ defmodule StatusphereWeb.StatusController do input: %{ repo: client.did, collection: "xyz.statusphere.status", - rkey: Atex.TID.now() |> to_string(), + rkey: rkey, record: record } } ) do + uri = + Atex.AtURI.to_string(%Atex.AtURI{ + authority: client.did, + collection: "xyz.statusphere.status", + rkey: rkey + }) + + optimistic_insert = + %Statusphere.Status{} + |> Statusphere.Status.changeset(%{ + uri: uri, + author_did: client.did, + status: record.status, + created_at: NaiveDateTime.from_iso8601!(record.createdAt), + indexed_at: NaiveDateTime.utc_now() + }) + |> Repo.insert() + + case optimistic_insert do + {:ok, _} -> + nil + + {:error, changeset} -> + Logger.error("Failed to optimistically insert status: #{inspect(changeset)}") + end + conn |> Atex.XRPC.OAuthClient.update_plug(client) - |> put_session(:selected_status, status) |> redirect(to: ~p"/") else :error -> diff --git a/priv/repo/migrations/20251219113112_create_status.exs b/priv/repo/migrations/20251219113112_create_status.exs new file mode 100644 index 0000000..122847d --- /dev/null +++ b/priv/repo/migrations/20251219113112_create_status.exs @@ -0,0 +1,13 @@ +defmodule Statusphere.Repo.Migrations.CreateStatus do + use Ecto.Migration + + def change do + create table(:status, primary_key: false) do + add :uri, :binary_id, primary_key: true + add :author_did, :string, null: false + add :status, :string, null: false + add :created_at, :utc_datetime, null: false + add :indexed_at, :utc_datetime, null: false + end + end +end