From ccdb3a0faaedad474e6880941872db51895b94e0 Mon Sep 17 00:00:00 2001 From: Johanna Larsson Date: Sat, 1 Aug 2026 16:12:44 +0100 Subject: [PATCH] Adds atproto-proxy Bluesky and some other services rely on actions other than writing to your PDS that require auth. This commit brings the simpler kind of service auth, proxying your request through your PDS. Also adds proper tests of the XRPC layer to verify the correct things are passed to the HTTP layer. --- README.md | 11 ++- lib/latch.ex | 34 +++++--- lib/latch/client.ex | 18 +++-- lib/latch/xrpc.ex | 39 +++++++--- test/latch/client_test.exs | 23 +++--- test/latch/xrpc_test.exs | 156 +++++++++++++++++++++++++++++++++++++ test/latch_test.exs | 11 ++- 7 files changed, 248 insertions(+), 44 deletions(-) create mode 100644 test/latch/xrpc_test.exs diff --git a/README.md b/README.md index 9205147..3c33174 100644 --- a/README.md +++ b/README.md @@ -83,9 +83,11 @@ authenticated calls. When a user logs out, call `delete_session`: Calls go to the user's PDS, and access tokens are refreshed automatically: Latch.query(MyApp.Latch, did, "com.atproto.repo.getRecord", - repo: did, - collection: "app.bsky.feed.post", - rkey: "3k2...") + params: [ + repo: did, + collection: "app.bsky.feed.post", + rkey: "3k2..." + ]) Latch.procedure(MyApp.Latch, did, "com.atproto.repo.createRecord", %{ repo: did, @@ -145,7 +147,8 @@ The library attempts to follow the spec strictly, but primarily in what the libr - [x] Local client - [x] Built-in ETS LatchStore implementation - [x] TID, NSID, AtURI, DID, Handle -- [ ] Service auth +- [x] Service auth (atproto-proxy) +- [ ] Service auth (bearer auth) - [ ] Getting started guide - [ ] Extensive tests - [ ] Distributed nonce cache diff --git a/lib/latch.ex b/lib/latch.ex index d0e49cf..546af0b 100644 --- a/lib/latch.ex +++ b/lib/latch.ex @@ -211,11 +211,17 @@ defmodule Latch do Assumes the session exists, that the user of that `did` is authenticated. If not, returns `{:error, %NoSession{}}`. + ## Options + + * `:service` - the service endpoint identifier when proxying through PDS, eg did:web:api.bsky.app#bsky_appview + * `:params` - params for the method, eg `params: [actor: "did:plc:bvraa6gajy4tfr3eh2sisdkr"]` resulting in + `app.bsky.actor.getProfile?actor=did:plc:bvraa6gajy4tfr3eh2sisdkr` + ## Examples - Latch.query(MyApp.Latch, "did:plc:abc123", "com.atproto.repo.getRecord", repo: "did:plc:abc123", collection: "app.bsky.feed.post", rkey: "3k2...") + Latch.query(MyApp.Latch, "did:plc:abc123", "com.atproto.repo.getRecord", params: [repo: "did:plc:abc123", collection: "app.bsky.feed.post", rkey: "3k2..."]) """ - @spec query(name(), String.t(), String.t(), Keyword.t()) :: + @spec query(name(), String.t(), String.t(), keyword()) :: {:ok, map()} | {:error, InvalidResponse.t() @@ -225,15 +231,19 @@ defmodule Latch do | StoreError.t() | Transport.t() | XRPCError.t()} - def query(name, did, method, params \\ []) do + def query(name, did, method, opts \\ []) do %Config{} = config = config(name) - Latch.Client.query(config, did, method, params) + Latch.Client.query(config, did, method, opts) end @doc """ Performs a procedure against the user's PDS using their DID's session. + ## Options + + * `service` - the service endpoint identifier when proxying through PDS, eg did:web:api.bsky.app#bsky_appview + ## Examples Latch.procedure(MyApp.Latch, "did:plc:abc123", "com.atproto.repo.putRecord", %{ @@ -243,7 +253,7 @@ defmodule Latch do record: %{"$type" => "app.bsky.feed.post", "text" => "Hello!"} }) """ - @spec procedure(name(), String.t(), String.t(), map()) :: + @spec procedure(name(), String.t(), String.t(), map(), keyword()) :: {:ok, map()} | {:error, InvalidResponse.t() @@ -253,18 +263,22 @@ defmodule Latch do | StoreError.t() | Transport.t() | XRPCError.t()} - def procedure(name, did, method, body) do + def procedure(name, did, method, body, opts \\ []) do %Config{} = config = config(name) - Latch.Client.procedure(config, did, method, body) + Latch.Client.procedure(config, did, method, body, opts) end @doc """ Upload a blob to the user's PDS using their DID's session. `content_type` is the blob's MIME type, eg `"image/png"`. + + ## Options + + * `service` - the service endpoint identifier when proxying through PDS, eg did:web:api.bsky.app#bsky_appview """ - @spec upload_blob(name(), String.t(), binary(), String.t()) :: + @spec upload_blob(name(), String.t(), binary(), String.t(), keyword()) :: {:ok, map()} | {:error, InvalidResponse.t() @@ -274,10 +288,10 @@ defmodule Latch do | StoreError.t() | Transport.t() | XRPCError.t()} - def upload_blob(name, did, bytes, content_type) do + def upload_blob(name, did, bytes, content_type, opts \\ []) do %Config{} = config = config(name) - Latch.Client.upload_blob(config, did, bytes, content_type) + Latch.Client.upload_blob(config, did, bytes, content_type, opts) end defp complete_callback( diff --git a/lib/latch/client.ex b/lib/latch/client.ex index c4c73f0..9042cd4 100644 --- a/lib/latch/client.ex +++ b/lib/latch/client.ex @@ -24,26 +24,28 @@ defmodule Latch.Client do """ @spec query(Config.t(), String.t(), String.t(), keyword()) :: {:ok, map()} | {:error, Error.t()} - def query(%Config{} = config, did, method, params \\ []) do - call(config, did, fn session -> XRPC.query(config, session, method, params) end) + def query(%Config{} = config, did, method, opts) do + call(config, did, fn session -> XRPC.query(config, session, method, opts) end) end @doc """ Performs an authenticated XRPC procedure for `did`, with automatic refresh. """ - @spec procedure(Config.t(), String.t(), String.t(), map()) :: + @spec procedure(Config.t(), String.t(), String.t(), map(), keyword()) :: {:ok, map()} | {:error, Error.t()} - def procedure(%Config{} = config, did, method, body) do - call(config, did, fn session -> XRPC.procedure(config, session, method, body) end) + def procedure(%Config{} = config, did, method, body, opts) do + call(config, did, fn session -> XRPC.procedure(config, session, method, body, opts) end) end @doc """ Uploads a blob for `did`, with automatic refresh. """ - @spec upload_blob(Config.t(), String.t(), binary(), String.t()) :: + @spec upload_blob(Config.t(), String.t(), binary(), String.t(), keyword()) :: {:ok, map()} | {:error, Error.t()} - def upload_blob(%Config{} = config, did, bytes, content_type) do - call(config, did, fn session -> XRPC.upload_blob(config, session, bytes, content_type) end) + def upload_blob(%Config{} = config, did, bytes, content_type, opts) do + call(config, did, fn session -> + XRPC.upload_blob(config, session, bytes, content_type, opts) + end) end defp call(config, did, fun) do diff --git a/lib/latch/xrpc.ex b/lib/latch/xrpc.ex index da35bd6..b092f65 100644 --- a/lib/latch/xrpc.ex +++ b/lib/latch/xrpc.ex @@ -24,41 +24,55 @@ defmodule Latch.XRPC do Performs and authenticated XRPC query against the session's PDS. """ @spec query(Config.t(), Session.t(), String.t(), keyword()) :: {:ok, map()} | {:error, error()} - def query(%Config{} = config, %Session{} = session, method, params \\ []) do + def query(%Config{} = config, %Session{} = session, method, opts) do + params = Keyword.get(opts, :params, []) + request( config, session, "GET", session.pds_endpoint <> "/xrpc/" <> method <> query_string(params), - nil + nil, + opts ) end @doc """ Performs and authenticated XRPC procedure against the session's PDS. """ - @spec procedure(Config.t(), Session.t(), String.t(), map()) :: {:ok, map()} | {:error, error()} - def procedure(%Config{} = config, %Session{} = session, method, body) do - request(config, session, "POST", session.pds_endpoint <> "/xrpc/" <> method, {:json, body}) + @spec procedure(Config.t(), Session.t(), String.t(), map(), keyword()) :: + {:ok, map()} | {:error, error()} + def procedure(%Config{} = config, %Session{} = session, method, body, opts) do + request( + config, + session, + "POST", + session.pds_endpoint <> "/xrpc/" <> method, + {:json, body}, + opts + ) end @doc """ Uploads raw bytes of content_type as a blob, returning the response with the blog reference. """ - @spec upload_blob(Config.t(), Session.t(), binary(), String.t()) :: + @spec upload_blob(Config.t(), Session.t(), binary(), String.t(), keyword()) :: {:ok, map()} | {:error, error()} - def upload_blob(%Config{} = config, %Session{} = session, bytes, content_type) do + def upload_blob(%Config{} = config, %Session{} = session, bytes, content_type, opts) do request( config, session, "POST", session.pds_endpoint <> "/xrpc/com.atproto.repo.uploadBlob", - {:raw, bytes, content_type} + {:raw, bytes, content_type}, + opts ) end - defp request(%Config{} = config, %Session{} = session, http_method, url, body) do + defp request(%Config{} = config, %Session{} = session, http_method, url, body, opts) do + service = Keyword.get(opts, :service) + DPoP.with_nonce(config, session.dpop_key, url, fn nonce -> proof = DPoP.proof(session.dpop_key, http_method, url, @@ -68,6 +82,13 @@ defmodule Latch.XRPC do headers = [{"authorization", "DPoP #{session.access_token}"}, {"dpop", proof}] + headers = + if service do + [{"atproto-proxy", service} | headers] + else + headers + end + case HTTP.request(http_method, url, headers, body) do {:ok, %{status: status, body: raw, headers: resp}} -> fresh_nonce = DPoP.nonce_header(resp) diff --git a/test/latch/client_test.exs b/test/latch/client_test.exs index 079e068..b26abaf 100644 --- a/test/latch/client_test.exs +++ b/test/latch/client_test.exs @@ -37,12 +37,12 @@ defmodule Latch.ClientTest do expect(XRPC, :query, fn _config, ^refreshed_session, "app.bsky.actor.getProfile", - actor: @did -> + params: [actor: @did] -> {:ok, %{"did" => @did}} end) assert {:ok, %{"did" => @did}} = - Client.query(config, @did, "app.bsky.actor.getProfile", actor: @did) + Client.query(config, @did, "app.bsky.actor.getProfile", params: [actor: @did]) assert {:ok, ^refreshed_session} = Latch.TestStore.fetch_session(@did) end @@ -56,10 +56,10 @@ defmodule Latch.ClientTest do :ok = Latch.TestStore.put_session(@did, stale_session) expect(XRPC, :query, 2, fn - _config, ^stale_session, "app.bsky.actor.getProfile", actor: @did -> + _config, ^stale_session, "app.bsky.actor.getProfile", params: [actor: @did] -> {:error, %XRPCError{status: 401, body: %{}}} - _config, ^refreshed_session, "app.bsky.actor.getProfile", actor: @did -> + _config, ^refreshed_session, "app.bsky.actor.getProfile", params: [actor: @did] -> {:ok, %{"did" => @did}} end) @@ -70,7 +70,7 @@ defmodule Latch.ClientTest do end) assert {:ok, %{"did" => @did}} = - Client.query(config, @did, "app.bsky.actor.getProfile", actor: @did) + Client.query(config, @did, "app.bsky.actor.getProfile", params: [actor: @did]) assert {:ok, ^refreshed_session} = Latch.TestStore.fetch_session(@did) end @@ -84,12 +84,16 @@ defmodule Latch.ClientTest do :ok = Latch.TestStore.put_session(@did, valid_session) - expect(XRPC, :procedure, fn _config, ^valid_session, "com.atproto.repo.putRecord", ^body -> + expect(XRPC, :procedure, fn _config, + ^valid_session, + "com.atproto.repo.putRecord", + ^body, + [] -> {:ok, %{"uri" => "at://#{@did}/app.bsky.feed.post/abc"}} end) assert {:ok, %{"uri" => _}} = - Client.procedure(config, @did, "com.atproto.repo.putRecord", body) + Client.procedure(config, @did, "com.atproto.repo.putRecord", body, []) end end @@ -100,11 +104,12 @@ defmodule Latch.ClientTest do :ok = Latch.TestStore.put_session(@did, valid_session) - expect(XRPC, :upload_blob, fn _config, ^valid_session, <<1, 2, 3>>, "image/png" -> + expect(XRPC, :upload_blob, fn _config, ^valid_session, <<1, 2, 3>>, "image/png", [] -> {:ok, %{"blob" => %{"ref" => "reffers"}}} end) - assert {:ok, %{"blob" => _}} = Client.upload_blob(config, @did, <<1, 2, 3>>, "image/png") + assert {:ok, %{"blob" => _}} = + Client.upload_blob(config, @did, <<1, 2, 3>>, "image/png", []) end end diff --git a/test/latch/xrpc_test.exs b/test/latch/xrpc_test.exs new file mode 100644 index 0000000..0bfe4b1 --- /dev/null +++ b/test/latch/xrpc_test.exs @@ -0,0 +1,156 @@ +defmodule Latch.XRPCTest do + use ExUnit.Case, async: true + use Mimic + + alias Latch.Config + alias Latch.DPoP + alias Latch.HTTP + alias Latch.Session + alias Latch.XRPC + + @did "did:plc:bvraa6gajy4tfr3eh2sisdkr" + + describe "query/4" do + test "passes all relevant information to HTTP call" do + config = make_config() + session = make_session() + method = "app.bsky.actor.getProfile" + params = [actor: @did] + service = "did:web:api.bsky.app#bsky_appview" + + start_link_supervised!( + {Latch.NonceCache, config: config, name: config.name, sweep_disabled: true} + ) + + expect(HTTP, :request, fn http_method, url, headers, body -> + assert http_method == "GET" + assert url == "#{session.pds_endpoint}/xrpc/#{method}?#{URI.encode_query(params)}" + assert {"atproto-proxy", service} in headers + assert {"authorization", "DPoP access-token"} in headers + assert {"dpop", _} = List.keyfind(headers, "dpop", 0) + + assert body == nil + + {:ok, %{status: 200, body: ~s|{"did": "string"}|, headers: %{}}} + end) + + assert {:ok, %{"did" => "string"}} = + XRPC.query(config, session, method, + params: params, + service: service + ) + end + + test "handles no params case" do + config = make_config() + session = make_session() + method = "com.atproto.server.getSession" + + start_link_supervised!( + {Latch.NonceCache, config: config, name: config.name, sweep_disabled: true} + ) + + expect(HTTP, :request, fn http_method, url, headers, body -> + assert http_method == "GET" + assert url == "#{session.pds_endpoint}/xrpc/#{method}" + assert {"authorization", "DPoP access-token"} in headers + assert {"dpop", _} = List.keyfind(headers, "dpop", 0) + + assert body == nil + + {:ok, %{status: 200, body: ~s|{"did": "string"}|, headers: %{}}} + end) + + assert {:ok, %{"did" => "string"}} = XRPC.query(config, session, method, []) + end + end + + describe "procedure/5" do + test "passes all relevant information to HTTP call" do + config = make_config() + session = make_session() + method = "app.bsky.actor.getProfile" + service = "did:web:api.bsky.app#bsky_appview" + expected_body = ~s|{"x": "y"}| + + start_link_supervised!( + {Latch.NonceCache, config: config, name: config.name, sweep_disabled: true} + ) + + expect(HTTP, :request, fn http_method, url, headers, body -> + assert http_method == "POST" + assert url == "#{session.pds_endpoint}/xrpc/#{method}" + assert {"atproto-proxy", service} in headers + assert {"authorization", "DPoP access-token"} in headers + assert {"dpop", _} = List.keyfind(headers, "dpop", 0) + + assert {:json, expected_body} == body + + {:ok, %{status: 200, body: ~s|{"did": "string"}|, headers: %{}}} + end) + + assert {:ok, %{"did" => "string"}} = + XRPC.procedure(config, session, method, expected_body, service: service) + end + end + + describe "upload_blob/5" do + test "passes all relevant information to HTTP call" do + config = make_config() + session = make_session() + method = "com.atproto.repo.uploadBlob" + service = "did:web:api.bsky.app#bsky_appview" + expected_body = "bytes" + + start_link_supervised!( + {Latch.NonceCache, config: config, name: config.name, sweep_disabled: true} + ) + + expect(HTTP, :request, fn http_method, url, headers, body -> + assert http_method == "POST" + assert url == "#{session.pds_endpoint}/xrpc/#{method}" + assert {"atproto-proxy", service} in headers + assert {"authorization", "DPoP access-token"} in headers + assert {"dpop", _} = List.keyfind(headers, "dpop", 0) + + assert {:raw, expected_body, "image/png"} == body + + {:ok, %{status: 200, body: ~s|{"did": "string"}|, headers: %{}}} + end) + + assert {:ok, %{"did" => "string"}} = + XRPC.upload_blob(config, session, expected_body, "image/png", service: service) + end + end + + defp make_config(overrides \\ []) do + defaults = [ + store: Latch.TestStore, + client_id: "https://client.example.com/oauth-client-metadata.json", + redirect_uri: "https://client.example.com/oauth/callback", + scope: "atproto", + signing_key: Jason.decode!(Jason.encode!(DPoP.generate_key())), + name: :"flow_test_#{inspect(self())}", + mode: :confidential + ] + + attrs = Keyword.merge(defaults, overrides) + struct!(Config, attrs) + end + + defp make_session(overrides \\ []) do + defaults = [ + did: @did, + access_token: "access-token", + refresh_token: "refresh-token", + dpop_key: DPoP.generate_key(), + scope: "atproto", + issuer: "https://issuer.example.com", + pds_endpoint: "https://pds.example.com", + expires_at: ~U[2026-01-01 00:00:00Z] + ] + + attrs = Keyword.merge(defaults, overrides) + struct!(Session, attrs) + end +end diff --git a/test/latch_test.exs b/test/latch_test.exs index ebd5679..77b41e9 100644 --- a/test/latch_test.exs +++ b/test/latch_test.exs @@ -156,12 +156,15 @@ defmodule LatchTest do mode: :confidential ) - expect(Client, :query, fn _config, @did, "app.bsky.actor.getProfile", actor: @did -> + expect(Client, :query, fn _config, + @did, + "app.bsky.actor.getProfile", + params: [actor: @did] -> {:ok, %{"did" => @did}} end) assert {:ok, %{"did" => @did}} = - Latch.query(pid, @did, "app.bsky.actor.getProfile", actor: @did) + Latch.query(pid, @did, "app.bsky.actor.getProfile", params: [actor: @did]) end end @@ -183,7 +186,7 @@ defmodule LatchTest do "record" => %{"text" => "hi"} } - expect(Client, :procedure, fn _config, @did, "com.atproto.repo.putRecord", ^body -> + expect(Client, :procedure, fn _config, @did, "com.atproto.repo.putRecord", ^body, [] -> {:ok, %{"uri" => "at://#{@did}/app.bsky.feed.post/abc"}} end) @@ -204,7 +207,7 @@ defmodule LatchTest do mode: :confidential ) - expect(Client, :upload_blob, fn _config, @did, <<1, 2, 3>>, "image/png" -> + expect(Client, :upload_blob, fn _config, @did, <<1, 2, 3>>, "image/png", [] -> {:ok, %{"blob" => %{"ref" => "reffers"}}} end) -- 2.51.2