diff --git a/CHANGELOG.md b/CHANGELOG.md --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -31,6 +31,9 @@ - `Atex.OAuth.session_keys_name/0` and `Atex.OAuth.session_active_session_name/0` expose Plug session key atoms - `Atex.IdentityResolver` now has full module and function documentation - `Atex.XRPC.LoginClient` now has a `@moduledoc` +- Optional `:telemetry` instrumentation via `Atex.Telemetry`. Add `{:telemetry, "~> 1.0"}` to + your deps to receive events from XRPC requests, identity resolution, OAuth flows, and service + auth validation. See `Atex.Telemetry` for the full event catalogue. ### Fixed diff --git a/lib/atex/identity_resolver.ex b/lib/atex/identity_resolver.ex --- a/lib/atex/identity_resolver.ex +++ b/lib/atex/identity_resolver.ex @@ -47,15 +47,33 @@ @spec resolve(String.t(), list(options())) :: {:ok, Identity.t()} | {:error, any()} def resolve(identifier, opts \\ []) do opts = Keyword.validate!(opts, skip_cache: false) skip_cache = Keyword.get(opts, :skip_cache) + identifier_type = if String.starts_with?(identifier, "did:"), do: :did, else: :handle - cache_result = if skip_cache, do: {:error, :not_found}, else: Cache.get(identifier) + Atex.Telemetry.span( + [:atex, :identity_resolver, :resolve], + %{identifier: identifier, identifier_type: identifier_type}, + fn -> + cache_result = if skip_cache, do: {:error, :not_found}, else: Cache.get(identifier) - # If cache fetch succeeds, then the ok tuple will be retuned by the default `with` behaviour - with {:error, :not_found} <- cache_result, - {:ok, identity} <- do_resolve(identifier), - identity <- Cache.insert(identity) do - {:ok, identity} - end + cache_event = if match?({:ok, _}, cache_result), do: :hit, else: :miss + + Atex.Telemetry.execute( + [:atex, :identity_resolver, :cache, cache_event], + %{system_time: System.system_time()}, + %{identifier: identifier} + ) + + # If cache fetch succeeds, then the ok tuple will be retuned by the default `with` behaviour + result = + with {:error, :not_found} <- cache_result, + {:ok, identity} <- do_resolve(identifier), + identity <- Cache.insert(identity) do + {:ok, identity} + end + + {result, %{}} + end + ) end @spec do_resolve(identity :: String.t()) :: diff --git a/lib/atex/oauth/flow.ex b/lib/atex/oauth/flow.ex --- a/lib/atex/oauth/flow.ex +++ b/lib/atex/oauth/flow.ex @@ -149,40 +149,48 @@ String.t(), list(create_authorization_url_option()) ) :: {:ok, String.t()} | {:error, any()} def create_authorization_url(authz_metadata, state, code_verifier, login_hint, opts \\ []) do - opts = Keyword.validate!(opts, [:key, :client_id, :redirect_uri, :scopes]) + Atex.Telemetry.span( + [:atex, :oauth, :authorization_url], + %{issuer: Map.get(authz_metadata, :issuer)}, + fn -> + opts = Keyword.validate!(opts, [:key, :client_id, :redirect_uri, :scopes]) + key = Keyword.get_lazy(opts, :key, &Config.get_key/0) + client_id = Keyword.get_lazy(opts, :client_id, &Config.client_id/0) + redirect_uri = Keyword.get_lazy(opts, :redirect_uri, &Config.redirect_uri/0) + scopes = Keyword.get_lazy(opts, :scopes, &Config.scopes/0) - key = Keyword.get_lazy(opts, :key, &Config.get_key/0) - client_id = Keyword.get_lazy(opts, :client_id, &Config.client_id/0) - redirect_uri = Keyword.get_lazy(opts, :redirect_uri, &Config.redirect_uri/0) - scopes = Keyword.get_lazy(opts, :scopes, &Config.scopes/0) + code_challenge = :crypto.hash(:sha256, code_verifier) |> Base.url_encode64(padding: false) + client_assertion = create_client_assertion(key, client_id, authz_metadata.issuer) - code_challenge = :crypto.hash(:sha256, code_verifier) |> Base.url_encode64(padding: false) - client_assertion = create_client_assertion(key, client_id, authz_metadata.issuer) + body = %{ + response_type: "code", + client_id: client_id, + redirect_uri: redirect_uri, + state: state, + code_challenge_method: "S256", + code_challenge: code_challenge, + scope: scopes, + client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer", + client_assertion: client_assertion, + login_hint: login_hint + } - body = %{ - response_type: "code", - client_id: client_id, - redirect_uri: redirect_uri, - state: state, - code_challenge_method: "S256", - code_challenge: code_challenge, - scope: scopes, - client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer", - client_assertion: client_assertion, - login_hint: login_hint - } + result = + case Req.post(authz_metadata.par_endpoint, form: body) do + {:ok, %{body: %{"request_uri" => request_uri}}} -> + query = %{client_id: client_id, request_uri: request_uri} |> URI.encode_query() + {:ok, "#{authz_metadata.authorization_endpoint}?#{query}"} - case Req.post(authz_metadata.par_endpoint, form: body) do - {:ok, %{body: %{"request_uri" => request_uri}}} -> - query = %{client_id: client_id, request_uri: request_uri} |> URI.encode_query() - {:ok, "#{authz_metadata.authorization_endpoint}?#{query}"} + {:ok, _} -> + {:error, :invalid_par_response} - {:ok, _} -> - {:error, :invalid_par_response} + err -> + err + end - err -> - err - end + {result, %{}} + end + ) end @doc """ @@ -213,54 +221,62 @@ String.t(), list(validate_authorization_code_option()) ) :: {:ok, tokens(), String.t() | nil} | {:error, any()} def validate_authorization_code(authz_metadata, dpop_key, code, code_verifier, opts \\ []) do - opts = Keyword.validate!(opts, [:key, :client_id, :redirect_uri, :scopes]) + Atex.Telemetry.span( + [:atex, :oauth, :code_exchange], + %{issuer: Map.get(authz_metadata, :issuer)}, + fn -> + opts = Keyword.validate!(opts, [:key, :client_id, :redirect_uri, :scopes]) + key = Keyword.get_lazy(opts, :key, &Config.get_key/0) + client_id = Keyword.get_lazy(opts, :client_id, &Config.client_id/0) + redirect_uri = Keyword.get_lazy(opts, :redirect_uri, &Config.redirect_uri/0) - key = Keyword.get_lazy(opts, :key, &Config.get_key/0) - client_id = Keyword.get_lazy(opts, :client_id, &Config.client_id/0) - redirect_uri = Keyword.get_lazy(opts, :redirect_uri, &Config.redirect_uri/0) + client_assertion = create_client_assertion(key, client_id, authz_metadata.issuer) - client_assertion = create_client_assertion(key, client_id, authz_metadata.issuer) + body = %{ + grant_type: "authorization_code", + client_id: client_id, + redirect_uri: redirect_uri, + code: code, + code_verifier: code_verifier + } - body = %{ - grant_type: "authorization_code", - client_id: client_id, - redirect_uri: redirect_uri, - code: code, - code_verifier: code_verifier - } + body = + if Config.localhost?(), + do: body, + else: + Map.merge(body, %{ + client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer", + client_assertion: client_assertion + }) - body = - if Config.localhost?(), - do: body, - else: - Map.merge(body, %{ - client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer", - client_assertion: client_assertion - }) + result = + Req.new(method: :post, url: authz_metadata.token_endpoint, form: body) + |> DPoP.send_oauth_dpop_request(dpop_key) + |> case do + {:ok, + %{ + "access_token" => access_token, + "refresh_token" => refresh_token, + "expires_in" => expires_in, + "sub" => did + }, nonce} -> + expires_at = NaiveDateTime.utc_now() |> NaiveDateTime.add(expires_in, :second) - Req.new(method: :post, url: authz_metadata.token_endpoint, form: body) - |> DPoP.send_oauth_dpop_request(dpop_key) - |> case do - {:ok, - %{ - "access_token" => access_token, - "refresh_token" => refresh_token, - "expires_in" => expires_in, - "sub" => did - }, nonce} -> - expires_at = NaiveDateTime.utc_now() |> NaiveDateTime.add(expires_in, :second) + {:ok, + %{ + access_token: access_token, + refresh_token: refresh_token, + did: did, + expires_at: expires_at + }, nonce} - {:ok, - %{ - access_token: access_token, - refresh_token: refresh_token, - did: did, - expires_at: expires_at - }, nonce} + {:error, reason, _nonce} -> + {:error, reason} + end - {:error, reason, _nonce} -> - {:error, reason} - end + {result, %{}} + end + ) end @doc """ @@ -285,44 +301,52 @@ String.t(), list(refresh_token_option()) ) :: {:ok, tokens(), String.t() | nil} | {:error, any()} def refresh_token(refresh_token, dpop_key, issuer, token_endpoint, opts \\ []) do - opts = Keyword.validate!(opts, [:key, :client_id]) + Atex.Telemetry.span( + [:atex, :oauth, :token_refresh], + %{issuer: issuer}, + fn -> + opts = Keyword.validate!(opts, [:key, :client_id]) + key = Keyword.get_lazy(opts, :key, &Config.get_key/0) + client_id = Keyword.get_lazy(opts, :client_id, &Config.client_id/0) - key = Keyword.get_lazy(opts, :key, &Config.get_key/0) - client_id = Keyword.get_lazy(opts, :client_id, &Config.client_id/0) + client_assertion = create_client_assertion(key, client_id, issuer) - client_assertion = create_client_assertion(key, client_id, issuer) + body = %{ + grant_type: "refresh_token", + refresh_token: refresh_token, + client_id: client_id, + client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer", + client_assertion: client_assertion + } - body = %{ - grant_type: "refresh_token", - refresh_token: refresh_token, - client_id: client_id, - client_assertion_type: "urn:ietf:params:oauth:client-assertion-type:jwt-bearer", - client_assertion: client_assertion - } + result = + Req.new(method: :post, url: token_endpoint, form: body) + |> DPoP.send_oauth_dpop_request(dpop_key) + |> case do + {:ok, + %{ + "access_token" => access_token, + "refresh_token" => refresh_token, + "expires_in" => expires_in, + "sub" => did + }, nonce} -> + expires_at = NaiveDateTime.utc_now() |> NaiveDateTime.add(expires_in, :second) - Req.new(method: :post, url: token_endpoint, form: body) - |> DPoP.send_oauth_dpop_request(dpop_key) - |> case do - {:ok, - %{ - "access_token" => access_token, - "refresh_token" => refresh_token, - "expires_in" => expires_in, - "sub" => did - }, nonce} -> - expires_at = NaiveDateTime.utc_now() |> NaiveDateTime.add(expires_in, :second) + {:ok, + %{ + access_token: access_token, + refresh_token: refresh_token, + did: did, + expires_at: expires_at + }, nonce} - {:ok, - %{ - access_token: access_token, - refresh_token: refresh_token, - did: did, - expires_at: expires_at - }, nonce} + {:error, reason, _nonce} -> + {:error, reason} + end - {:error, reason, _nonce} -> - {:error, reason} - end + {result, %{}} + end + ) end @doc """ @@ -343,30 +367,39 @@ - `:ok` - Tokens revoked (or revocation endpoint unreachable - logged, not raised) """ @spec revoke_tokens(Session.t(), authorization_metadata()) :: :ok def revoke_tokens(%Session{} = session, authz_metadata) do - client_id = Config.client_id() + Atex.Telemetry.span( + [:atex, :oauth, :token_revocation], + %{issuer: Map.get(authz_metadata, :issuer)}, + fn -> + client_id = Config.client_id() - body = %{ - client_id: client_id, - token: session.refresh_token, - token_type_hint: "refresh_token" - } + body = %{ + client_id: client_id, + token: session.refresh_token, + token_type_hint: "refresh_token" + } - case Req.post(authz_metadata.revocation_endpoint, form: body) do - {:ok, %{status: status}} when status in [200, 204] -> - :ok + result = + case Req.post(authz_metadata.revocation_endpoint, form: body) do + {:ok, %{status: status}} when status in [200, 204] -> + :ok - {:ok, %{body: %{"error" => error}}} -> - Logger.warning("Token revocation failed: #{error}") - :ok + {:ok, %{body: %{"error" => error}}} -> + Logger.warning("Token revocation failed: #{error}") + :ok - {:error, reason} -> - Logger.warning("Token revocation request failed: #{inspect(reason)}") - :ok + {:error, reason} -> + Logger.warning("Token revocation request failed: #{inspect(reason)}") + :ok + + unexpected -> + Logger.warning("Unexpected token revocation response: #{inspect(unexpected)}") + :ok + end - unexpected -> - Logger.warning("Unexpected token revocation response: #{inspect(unexpected)}") - :ok - end + {result, %{}} + end + ) end @spec random_b64(integer()) :: String.t() diff --git a/lib/atex/service_auth.ex b/lib/atex/service_auth.ex --- a/lib/atex/service_auth.ex +++ b/lib/atex/service_auth.ex @@ -105,49 +105,75 @@ {expected_aud, expected_lxm} = options(opts) peek_result = try do - peeked = JOSE.JWT.peek(jwt) - {:ok, peeked} + {:ok, JOSE.JWT.peek(jwt)} rescue _ -> {:error, :invalid_jwt} end - case peek_result do - {:error, _} = err -> - err + {span_iss, span_lxm} = + case peek_result do + {:ok, %{fields: fields}} -> {Map.get(fields, "iss"), Map.get(fields, "lxm")} + _ -> {nil, nil} + end - {:ok, - %{ - fields: - %{ - "aud" => target_aud, - "iat" => iat, - "exp" => exp, - "iss" => issuing_did, - "jti" => nonce - } = fields - }} -> - target_lxm = Map.get(fields, "lxm") + Atex.Telemetry.span( + [:atex, :service_auth, :validate], + %{iss: span_iss, lxm: span_lxm}, + fn -> + result = do_validate_jwt(jwt, peek_result, expected_aud, expected_lxm) + {result, %{}} + end + ) + end - with :ok <- validate_aud(expected_aud, target_aud), - :ok <- validate_lxm(expected_lxm, target_lxm), - :ok <- validate_token_times(iat, exp), - # Resolve JWT's issuer to: a) make sure it's a real identity, b) get - # the signing key from their DID document to verify the token - {:ok, identity} <- Atex.IdentityResolver.resolve(issuing_did), - user_jwk when not is_nil(user_jwk) <- - Atex.DID.Document.get_atproto_signing_key(identity.document), - {true, %JOSE.JWT{} = jwt_struct, _jws} <- JOSE.JWT.verify(user_jwk, jwt), - # Record the nonce atomically after successful verification. insert_new - # is used under the hood so this returns :seen if the jti was already - # consumed, preventing replay attacks. - :ok <- Atex.ServiceAuth.JTICache.put(nonce, exp) do - {:ok, jwt_struct} - else - :seen -> {:error, :replayed_token} - err -> err - end + @spec do_validate_jwt( + String.t(), + {:ok, JOSE.JWT.t()} | {:error, atom()}, + String.t(), + String.t() | nil + ) :: {:ok, JOSE.JWT.t()} | {:error, atom()} + defp do_validate_jwt(_jwt, {:error, _} = err, _expected_aud, _expected_lxm), do: err + + defp do_validate_jwt( + jwt, + {:ok, + %{ + fields: + %{ + "aud" => target_aud, + "iat" => iat, + "exp" => exp, + "iss" => issuing_did, + "jti" => nonce + } = fields + }}, + expected_aud, + expected_lxm + ) do + target_lxm = Map.get(fields, "lxm") + + with :ok <- validate_aud(expected_aud, target_aud), + :ok <- validate_lxm(expected_lxm, target_lxm), + :ok <- validate_token_times(iat, exp), + # Resolve JWT's issuer to: a) make sure it's a real identity, b) get + # the signing key from their DID document to verify the token + {:ok, identity} <- Atex.IdentityResolver.resolve(issuing_did), + user_jwk when not is_nil(user_jwk) <- + Atex.DID.Document.get_atproto_signing_key(identity.document), + {true, %JOSE.JWT{} = jwt_struct, _jws} <- JOSE.JWT.verify(user_jwk, jwt), + # Record the nonce atomically after successful verification. insert_new + # is used under the hood so this returns :seen if the jti was already + # consumed, preventing replay attacks. + :ok <- Atex.ServiceAuth.JTICache.put(nonce, exp) do + {:ok, jwt_struct} + else + :seen -> {:error, :replayed_token} + err -> err end end + + defp do_validate_jwt(_jwt, {:ok, _unmatched}, _expected_aud, _expected_lxm), + do: {:error, :invalid_jwt} @spec validate_token_times(integer(), integer()) :: :ok | {:error, reason :: atom()} defp validate_token_times(iat, exp) do diff --git a/lib/atex/telemetry.ex b/lib/atex/telemetry.ex new file mode 100644 --- /dev/null +++ b/lib/atex/telemetry.ex @@ -0,0 +1,282 @@ +defmodule Atex.Telemetry do + @moduledoc """ + Telemetry instrumentation for Atex. + + Atex emits `:telemetry` events throughout its subsystems. To receive events, + attach handlers using `:telemetry.attach/4` or `:telemetry.attach_many/4`. + + `:telemetry` is an **optional dependency**. If it is not present in your + application's deps, all instrumentation calls compile to no-ops with zero + runtime overhead. Add it to your `mix.exs` to enable instrumentation: + + {:telemetry, "~> 1.0"} + + ## Event catalogue + + ### XRPC + + #### `[:atex, :xrpc, :request, :start | :stop | :exception]` + + Emitted for every outgoing XRPC HTTP request (all client types). + + - **start measurements:** `%{system_time: integer()}` + - **stop measurements:** `%{duration: integer()}` + - **exception measurements:** `%{duration: integer()}` + - **metadata (all events):** `%{method: :get | :post, resource: String.t(), endpoint: String.t(), client_type: :login | :oauth | :service_auth | :unauthed}` + - **stop additional metadata:** `%{status: integer()}` + - **exception additional metadata:** `%{kind: :error, reason: term(), stacktrace: list()}` + + #### `[:atex, :xrpc, :token_refresh, :start | :stop | :exception]` + + Emitted when a client performs a token refresh (`LoginClient` or `OAuthClient`). + + - **start measurements:** `%{system_time: integer()}` + - **stop measurements:** `%{duration: integer()}` + - **metadata:** `%{client_type: :login | :oauth}` + + ### Identity Resolver + + #### `[:atex, :identity_resolver, :resolve, :start | :stop | :exception]` + + Emitted for every call to `Atex.IdentityResolver.resolve/2`. + + - **start measurements:** `%{system_time: integer()}` + - **stop measurements:** `%{duration: integer()}` + - **metadata:** `%{identifier: String.t(), identifier_type: :did | :handle}` + + #### `[:atex, :identity_resolver, :cache, :hit | :miss]` + + Emitted at the cache check branch inside `resolve/2`. `:hit` means the result + was served from cache; `:miss` means a fresh resolution was performed (including + when `skip_cache: true` is passed). + + - **measurements:** `%{system_time: integer()}` + - **metadata:** `%{identifier: String.t()}` + + ### OAuth + + All OAuth spans share: + + - **start measurements:** `%{system_time: integer()}` + - **stop measurements:** `%{duration: integer()}` + - **metadata:** `%{issuer: String.t() | nil}` + + #### `[:atex, :oauth, :authorization_url, :start | :stop | :exception]` + + Wraps `Atex.OAuth.Flow.create_authorization_url/5` (PAR request + URL construction). + + #### `[:atex, :oauth, :code_exchange, :start | :stop | :exception]` + + Wraps `Atex.OAuth.Flow.validate_authorization_code/5`. + + #### `[:atex, :oauth, :token_refresh, :start | :stop | :exception]` + + Wraps `Atex.OAuth.Flow.refresh_token/5`. + + #### `[:atex, :oauth, :token_revocation, :start | :stop | :exception]` + + Wraps `Atex.OAuth.Flow.revoke_tokens/2`. + + ### Service Auth + + #### `[:atex, :service_auth, :validate, :start | :stop | :exception]` + + Wraps `Atex.ServiceAuth.validate_jwt/2`. + + - **start measurements:** `%{system_time: integer()}` + - **stop measurements:** `%{duration: integer()}` + - **metadata:** `%{iss: String.t() | nil, lxm: String.t() | nil}` + + ## Example handler + + :telemetry.attach_many( + "my-app-atex-handler", + [ + [:atex, :xrpc, :request, :stop], + [:atex, :identity_resolver, :resolve, :stop] + ], + fn event, measurements, metadata, _config -> + Logger.debug( + "Atex event: \#{inspect(event)} " <> + "duration=\#{measurements.duration} " <> + "metadata=\#{inspect(metadata)}" + ) + end, + nil + ) + """ + + @telemetry_available Code.ensure_loaded?(:telemetry) + + if @telemetry_available do + @doc """ + Execute a telemetry event. + + Delegates to `:telemetry.execute/3`. No-op when `:telemetry` is not loaded. + """ + @spec execute(list(atom()), map(), map()) :: :ok + def execute(event, measurements, metadata), + do: :telemetry.execute(event, measurements, metadata) + + @doc """ + Span a block with telemetry start/stop/exception events. + + Emits `event_prefix ++ [:start]` before calling `fun` and `event_prefix ++ [:stop]` + after it returns. The `:start` event carries `%{system_time: System.system_time()}` as + measurements and `start_metadata` as metadata. The `:stop` event carries + `%{duration: duration}` (in native time units) and the result of merging + `start_metadata` with the extra metadata returned by `fun`. + + `fun` must return `{result, extra_stop_metadata}`. This function returns `result`. + + No-op (calls `fun` and returns result) when `:telemetry` is not loaded. + """ + @spec span(list(atom()), map(), (-> {result, map()})) :: result when result: any() + def span(event_prefix, start_metadata, fun) do + start_time = System.monotonic_time() + + :telemetry.execute( + event_prefix ++ [:start], + %{system_time: System.system_time()}, + start_metadata + ) + + try do + {result, extra_stop_metadata} = fun.() + duration = System.monotonic_time() - start_time + + :telemetry.execute( + event_prefix ++ [:stop], + %{duration: duration}, + Map.merge(start_metadata, extra_stop_metadata) + ) + + result + rescue + exception -> + duration = System.monotonic_time() - start_time + + :telemetry.execute( + event_prefix ++ [:exception], + %{duration: duration}, + Map.merge(start_metadata, %{ + kind: :error, + reason: exception, + stacktrace: __STACKTRACE__ + }) + ) + + reraise exception, __STACKTRACE__ + catch + kind, reason -> + duration = System.monotonic_time() - start_time + + :telemetry.execute( + event_prefix ++ [:exception], + %{duration: duration}, + Map.merge(start_metadata, %{ + kind: kind, + reason: reason, + stacktrace: __STACKTRACE__ + }) + ) + + :erlang.raise(kind, reason, __STACKTRACE__) + end + end + + @doc """ + Attach telemetry instrumentation to a `Req.Request`. + + Adds request and response steps that emit `[:atex, :xrpc, :request, ...]` + events. Pass `client_type:` to identify the XRPC client variant. + + No-op when `:telemetry` is not loaded. + + ## Options + + - `:client_type` — one of `:login`, `:oauth`, `:service_auth`, `:unauthed` + (default: `:unknown`) + + ## Example + + Req.new(method: :get, url: url) + |> Atex.Telemetry.attach_req_plugin(client_type: :login) + |> Req.request() + """ + @spec attach_req_plugin(Req.Request.t(), keyword()) :: Req.Request.t() + def attach_req_plugin(req, opts \\ []) do + client_type = Keyword.get(opts, :client_type, :unknown) + + req + |> Req.Request.append_request_steps(atex_telemetry_start: &run_start_step(&1, client_type)) + |> Req.Request.prepend_response_steps(atex_telemetry_stop: &run_stop_step/1) + |> Req.Request.prepend_error_steps(atex_telemetry_stop: &run_stop_step/1) + end + + defp run_start_step(req, client_type) do + start_time = System.monotonic_time() + path = req.url.path || "" + resource = String.replace_prefix(path, "/xrpc/", "") + + endpoint = + URI.to_string(%URI{scheme: req.url.scheme, host: req.url.host, port: req.url.port}) + + :telemetry.execute( + [:atex, :xrpc, :request, :start], + %{system_time: System.system_time()}, + %{method: req.method, resource: resource, endpoint: endpoint, client_type: client_type} + ) + + req + |> Req.Request.put_private(:atex_start_time, start_time) + |> Req.Request.put_private(:atex_metadata, %{ + method: req.method, + resource: resource, + endpoint: endpoint, + client_type: client_type + }) + end + + defp run_stop_step({req, response}) do + start_time = Req.Request.get_private(req, :atex_start_time) + base_metadata = Req.Request.get_private(req, :atex_metadata) || %{} + + if start_time do + duration = System.monotonic_time() - start_time + + if match?(%Req.Response{}, response) do + :telemetry.execute( + [:atex, :xrpc, :request, :stop], + %{duration: duration}, + Map.put(base_metadata, :status, response.status) + ) + else + :telemetry.execute( + [:atex, :xrpc, :request, :exception], + %{duration: duration}, + # stacktrace unavailable in Req error steps — transport errors don't have one + Map.merge(base_metadata, %{kind: :error, reason: response, stacktrace: []}) + ) + end + end + + {req, response} + end + else + @doc false + @spec execute(list(atom()), map(), map()) :: :ok + def execute(_event, _measurements, _metadata), do: :ok + + @doc false + @spec span(list(atom()), map(), (-> {result, map()})) :: result when result: any() + def span(_event_prefix, _start_metadata, fun) do + {result, _meta} = fun.() + result + end + + @doc false + @spec attach_req_plugin(Req.Request.t(), keyword()) :: Req.Request.t() + def attach_req_plugin(req, _opts \\ []), do: req + end +end diff --git a/lib/atex/xrpc.ex b/lib/atex/xrpc.ex --- a/lib/atex/xrpc.ex +++ b/lib/atex/xrpc.ex @@ -162,7 +162,10 @@ """ @spec unauthed_get(String.t(), String.t(), keyword()) :: {:ok, Req.Response.t()} | {:error, any()} def unauthed_get(endpoint, name, opts \\ []) do - Req.get(url(endpoint, name), opts) + (opts ++ [method: :get, url: url(endpoint, name)]) + |> Req.new() + |> Atex.Telemetry.attach_req_plugin(client_type: :unauthed) + |> Req.request() end @doc """ @@ -171,7 +174,10 @@ """ @spec unauthed_post(String.t(), String.t(), keyword()) :: {:ok, Req.Response.t()} | {:error, any()} def unauthed_post(endpoint, name, opts \\ []) do - Req.post(url(endpoint, name), opts) + (opts ++ [method: :post, url: url(endpoint, name)]) + |> Req.new() + |> Atex.Telemetry.attach_req_plugin(client_type: :unauthed) + |> Req.request() end @doc """ diff --git a/lib/atex/xrpc/login_client.ex b/lib/atex/xrpc/login_client.ex --- a/lib/atex/xrpc/login_client.ex +++ b/lib/atex/xrpc/login_client.ex @@ -80,20 +80,29 @@ Request a new `refresh_token` for the given client. """ @spec refresh(t()) :: {:ok, t()} | {:error, any()} def refresh(%__MODULE__{endpoint: endpoint, refresh_token: refresh_token} = client) do - request = - Req.new(method: :post, url: XRPC.url(endpoint, "com.atproto.server.refreshSession")) - |> put_auth(refresh_token) + Atex.Telemetry.span( + [:atex, :xrpc, :token_refresh], + %{client_type: :login}, + fn -> + request = + Req.new(method: :post, url: XRPC.url(endpoint, "com.atproto.server.refreshSession")) + |> put_auth(refresh_token) + + result = + case Req.request(request) do + {:ok, %{body: %{"accessJwt" => access_token, "refreshJwt" => refresh_token}}} -> + {:ok, %{client | access_token: access_token, refresh_token: refresh_token}} - case Req.request(request) do - {:ok, %{body: %{"accessJwt" => access_token, "refreshJwt" => refresh_token}}} -> - {:ok, %{client | access_token: access_token, refresh_token: refresh_token}} + {:ok, response} -> + {:error, response} - {:ok, response} -> - {:error, response} + err -> + err + end - err -> - err - end + {result, %{}} + end + ) end @impl true @@ -109,7 +118,11 @@ @spec request(t(), keyword()) :: {:ok, Req.Response.t(), t()} | {:error, any()} defp request(client, opts) do with {:ok, client} <- validate_client(client) do - request = opts |> Req.new() |> put_auth(client.access_token) + request = + opts + |> Req.new() + |> put_auth(client.access_token) + |> Atex.Telemetry.attach_req_plugin(client_type: :login) case Req.request(request) do {:ok, %{status: 200} = response} -> diff --git a/lib/atex/xrpc/oauth_client.ex b/lib/atex/xrpc/oauth_client.ex --- a/lib/atex/xrpc/oauth_client.ex +++ b/lib/atex/xrpc/oauth_client.ex @@ -142,7 +142,18 @@ end) end @spec do_refresh(t()) :: {:ok, OAuth.Session.t()} | {:error, any()} - defp do_refresh(%__MODULE__{session_key: session_key}) do + defp do_refresh(%__MODULE__{} = client) do + Atex.Telemetry.span( + [:atex, :xrpc, :token_refresh], + %{client_type: :oauth}, + fn -> + {do_refresh_impl(client), %{}} + end + ) + end + + @spec do_refresh_impl(t()) :: {:ok, OAuth.Session.t()} | {:error, any()} + defp do_refresh_impl(%__MODULE__{session_key: session_key}) do with {:ok, session} <- OAuth.SessionStore.get(session_key), {:ok, authz_server} <- Discovery.get_authorization_server(session.aud), {:ok, %{token_endpoint: token_endpoint}} <- @@ -227,6 +238,7 @@ opts |> Keyword.put(:url, url) |> Req.new() |> Req.Request.put_header("authorization", "DPoP #{session.access_token}") + |> Atex.Telemetry.attach_req_plugin(client_type: :oauth) case DPoP.request_protected_dpop_resource( request, diff --git a/lib/atex/xrpc/service_auth_client.ex b/lib/atex/xrpc/service_auth_client.ex --- a/lib/atex/xrpc/service_auth_client.ex +++ b/lib/atex/xrpc/service_auth_client.ex @@ -51,7 +51,11 @@ end @spec request(t(), keyword()) :: {:ok, Req.Response.t(), t()} | {:error, any(), t()} defp request(client, opts) do - req = opts |> Req.new() |> put_auth(client.token) + req = + opts + |> Req.new() + |> put_auth(client.token) + |> Atex.Telemetry.attach_req_plugin(client_type: :service_auth) case Req.request(req) do {:ok, response} -> {:ok, response, client} diff --git a/lib/atex/xrpc/unauthed_client.ex b/lib/atex/xrpc/unauthed_client.ex --- a/lib/atex/xrpc/unauthed_client.ex +++ b/lib/atex/xrpc/unauthed_client.ex @@ -27,6 +27,7 @@ @impl true def get(%__MODULE__{endpoint: endpoint} = client, resource, opts \\ []) do (opts ++ [method: :get, url: Atex.XRPC.url(endpoint, resource)]) |> Req.new() + |> Atex.Telemetry.attach_req_plugin(client_type: :unauthed) |> Req.request() |> case do {:ok, response} -> {:ok, response, client} @@ -38,6 +39,7 @@ @impl true def post(%__MODULE__{endpoint: endpoint} = client, resource, opts \\ []) do (opts ++ [method: :post, url: Atex.XRPC.url(endpoint, resource)]) |> Req.new() + |> Atex.Telemetry.attach_req_plugin(client_type: :unauthed) |> Req.request() |> case do {:ok, response} -> {:ok, response, client} diff --git a/mix.exs b/mix.exs --- a/mix.exs +++ b/mix.exs @@ -46,6 +46,7 @@ {:jose, "~> 1.11"}, {:bandit, "~> 1.0", only: [:dev, :test]}, {:con_cache, "~> 1.1"}, {:mutex, "~> 3.0"}, + {:telemetry, "~> 1.0", optional: true}, {:dasl, "~> 0.1"}, {:mst, "~> 0.1"}, {:dialyxir, "~> 1.4", only: [:dev, :test], runtime: false}, diff --git a/test/atex/identity_resolver_telemetry_test.exs b/test/atex/identity_resolver_telemetry_test.exs new file mode 100644 --- /dev/null +++ b/test/atex/identity_resolver_telemetry_test.exs @@ -0,0 +1,86 @@ +defmodule Atex.IdentityResolverTelemetryTest do + use ExUnit.Case, async: true + + describe "resolve/2 telemetry" do + test "emits resolve start/stop and cache miss on first call" do + ref = make_ref() + identifier = "did:plc:test-#{inspect(ref)}" + + :telemetry.attach_many( + "test-resolver-#{inspect(ref)}", + [ + [:atex, :identity_resolver, :resolve, :start], + [:atex, :identity_resolver, :resolve, :stop], + [:atex, :identity_resolver, :cache, :miss], + [:atex, :identity_resolver, :cache, :hit] + ], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-resolver-#{inspect(ref)}") end) + + # resolve will fail because it's not a real DID, but telemetry still fires + Atex.IdentityResolver.resolve(identifier) + + assert_receive {:telemetry, [:atex, :identity_resolver, :resolve, :start], + %{system_time: _}, %{identifier: ^identifier, identifier_type: :did}} + + assert_receive {:telemetry, [:atex, :identity_resolver, :cache, :miss], %{system_time: _}, + %{identifier: ^identifier}} + + assert_receive {:telemetry, [:atex, :identity_resolver, :resolve, :stop], %{duration: _}, _} + end + + test "emits cache hit event on repeated call for same identifier" do + ref = make_ref() + # Use a DID-style identifier that we can pre-populate in the cache + identifier = "did:plc:cached-#{inspect(ref)}" + + :telemetry.attach( + "test-resolver-cache-#{inspect(ref)}", + [:atex, :identity_resolver, :cache, :hit], + fn _event, _measurements, metadata, _ -> + send(self(), {:cache_hit, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-resolver-cache-#{inspect(ref)}") end) + + # First call — cache miss, resolution fails (not a real DID) + Atex.IdentityResolver.resolve(identifier) + + # Pre-populate the cache directly so the second call gets a hit + identity = %Atex.IdentityResolver.Identity{did: identifier, handle: nil, document: nil} + Atex.IdentityResolver.Cache.insert(identity) + + # Second call — cache hit + Atex.IdentityResolver.resolve(identifier) + + assert_receive {:cache_hit, %{identifier: ^identifier}} + end + + test "identifier_type is :handle for non-did identifiers" do + ref = make_ref() + identifier = "user.bsky.social" + + :telemetry.attach( + "test-resolver-handle-#{inspect(ref)}", + [:atex, :identity_resolver, :resolve, :start], + fn _event, _measurements, metadata, _ -> + send(self(), {:start, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-resolver-handle-#{inspect(ref)}") end) + + Atex.IdentityResolver.resolve(identifier) + + assert_receive {:start, %{identifier_type: :handle}} + end + end +end diff --git a/test/atex/oauth/flow_test.exs b/test/atex/oauth/flow_test.exs --- a/test/atex/oauth/flow_test.exs +++ b/test/atex/oauth/flow_test.exs @@ -81,4 +81,99 @@ {true, %JOSE.JWT{}, _} = JOSE.JWT.verify(JOSE.JWK.to_public(key), token) end end + + describe "telemetry" do + setup do + key = JOSE.JWK.generate_key({:ec, "P-256"}) + key = %{key | fields: Map.put(key.fields, "kid", "test-kid")} + %{key: key} + end + + test "create_authorization_url/5 emits authorization_url start/stop events", %{key: key} do + ref = make_ref() + + :telemetry.attach_many( + "test-oauth-authz-url-#{inspect(ref)}", + [ + [:atex, :oauth, :authorization_url, :start], + [:atex, :oauth, :authorization_url, :stop] + ], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-oauth-authz-url-#{inspect(ref)}") end) + + authz_metadata = %{ + issuer: "https://bsky.social", + par_endpoint: "https://bsky.social/oauth/par", + token_endpoint: "https://bsky.social/oauth/token", + authorization_endpoint: "https://bsky.social/oauth/authorize", + revocation_endpoint: "https://bsky.social/oauth/revoke" + } + + # This will fail (no real PAR server) but telemetry still fires + Flow.create_authorization_url(authz_metadata, "state", "verifier", "user.bsky.social", + key: key, + client_id: "https://example.com/client", + redirect_uri: "https://example.com/callback", + scopes: "atproto" + ) + + assert_receive {:telemetry, [:atex, :oauth, :authorization_url, :start], %{system_time: _}, + %{issuer: "https://bsky.social"}} + + assert_receive {:telemetry, [:atex, :oauth, :authorization_url, :stop], %{duration: _}, + %{issuer: "https://bsky.social"}} + end + + test "revoke_tokens/2 emits token_revocation start/stop events" do + ref = make_ref() + + :telemetry.attach_many( + "test-oauth-revoke-#{inspect(ref)}", + [ + [:atex, :oauth, :token_revocation, :start], + [:atex, :oauth, :token_revocation, :stop] + ], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-oauth-revoke-#{inspect(ref)}") end) + + session = %Atex.OAuth.Session{ + iss: "https://bsky.social", + aud: "https://bsky.social", + sub: "did:plc:abc", + nonce: "nonce", + access_token: "token", + refresh_token: "refresh", + expires_at: NaiveDateTime.utc_now(), + dpop_key: JOSE.JWK.generate_key({:ec, "P-256"}), + dpop_nonce: nil + } + + authz_metadata = %{ + issuer: "https://bsky.social", + par_endpoint: "https://bsky.social/oauth/par", + token_endpoint: "https://bsky.social/oauth/token", + authorization_endpoint: "https://bsky.social/oauth/authorize", + revocation_endpoint: "https://bsky.social/oauth/revoke" + } + + # Will fail (no real server) but telemetry fires + Flow.revoke_tokens(session, authz_metadata) + + assert_receive {:telemetry, [:atex, :oauth, :token_revocation, :start], %{system_time: _}, + %{issuer: "https://bsky.social"}} + + assert_receive {:telemetry, [:atex, :oauth, :token_revocation, :stop], %{duration: _}, + %{issuer: "https://bsky.social"}} + end + end end diff --git a/test/atex/service_auth_test.exs b/test/atex/service_auth_test.exs --- a/test/atex/service_auth_test.exs +++ b/test/atex/service_auth_test.exs @@ -12,4 +12,32 @@ assert {:error, :invalid_jwt} = Atex.ServiceAuth.validate_jwt("", aud: "did:web:example.com") end end + + describe "validate_jwt/2 telemetry" do + test "emits validate start and stop events even for invalid JWT" do + ref = make_ref() + + :telemetry.attach_many( + "test-service-auth-#{inspect(ref)}", + [ + [:atex, :service_auth, :validate, :start], + [:atex, :service_auth, :validate, :stop] + ], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-service-auth-#{inspect(ref)}") end) + + Atex.ServiceAuth.validate_jwt("not.a.valid.jwt", aud: "did:web:example.com") + + assert_receive {:telemetry, [:atex, :service_auth, :validate, :start], %{system_time: _}, + %{iss: nil, lxm: nil}} + + assert_receive {:telemetry, [:atex, :service_auth, :validate, :stop], %{duration: _}, + %{iss: nil, lxm: nil}} + end + end end diff --git a/test/atex/telemetry_test.exs b/test/atex/telemetry_test.exs new file mode 100644 --- /dev/null +++ b/test/atex/telemetry_test.exs @@ -0,0 +1,269 @@ +defmodule Atex.TelemetryTest do + use ExUnit.Case, async: true + + describe "execute/3" do + test "emits a telemetry event" do + ref = make_ref() + + :telemetry.attach( + "test-execute-#{inspect(ref)}", + [:atex, :test, :event], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-execute-#{inspect(ref)}") end) + + Atex.Telemetry.execute([:atex, :test, :event], %{count: 1}, %{key: "val"}) + + assert_receive {:telemetry, [:atex, :test, :event], %{count: 1}, %{key: "val"}} + end + end + + describe "span/3" do + test "emits start and stop events and returns the result" do + ref = make_ref() + + :telemetry.attach_many( + "test-span-#{inspect(ref)}", + [[:atex, :test, :span, :start], [:atex, :test, :span, :stop]], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-span-#{inspect(ref)}") end) + + result = + Atex.Telemetry.span([:atex, :test, :span], %{key: "start_val"}, fn -> + {{:ok, :returned_value}, %{extra: "stop_val"}} + end) + + assert result == {:ok, :returned_value} + + assert_receive {:telemetry, [:atex, :test, :span, :start], %{system_time: _}, + %{key: "start_val"}} + + assert_receive {:telemetry, [:atex, :test, :span, :stop], %{duration: _}, + %{key: "start_val", extra: "stop_val"}} + end + end + + defmodule TransportErrorPlug do + @moduledoc false + def init(opts), do: opts + def call(conn, _opts), do: Req.Test.transport_error(conn, :closed) + end + + defmodule OkPlug do + @moduledoc false + import Plug.Conn + def init(opts), do: opts + + def call(conn, _opts) do + conn |> send_resp(200, Jason.encode!(%{ok: true})) + end + end + + defmodule ErrorPlug do + @moduledoc false + import Plug.Conn + def init(opts), do: opts + + def call(conn, _opts) do + conn |> send_resp(500, Jason.encode!(%{error: "ServerError"})) + end + end + + describe "attach_req_plugin/2" do + test "emits start and stop events on success" do + ref = make_ref() + + :telemetry.attach_many( + "test-plugin-success-#{inspect(ref)}", + [[:atex, :xrpc, :request, :start], [:atex, :xrpc, :request, :stop]], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-plugin-success-#{inspect(ref)}") end) + + req = + Req.new( + method: :get, + url: "http://bsky.social/xrpc/app.bsky.actor.getProfile", + plug: OkPlug + ) + |> Atex.Telemetry.attach_req_plugin(client_type: :login) + + {:ok, _response} = Req.request(req) + + assert_receive {:telemetry, [:atex, :xrpc, :request, :start], %{system_time: _}, + %{ + method: :get, + resource: "app.bsky.actor.getProfile", + endpoint: "http://bsky.social", + client_type: :login + }} + + assert_receive {:telemetry, [:atex, :xrpc, :request, :stop], %{duration: _}, + %{ + status: 200, + method: :get, + resource: "app.bsky.actor.getProfile", + endpoint: "http://bsky.social", + client_type: :login + }} + end + + test "includes status code from non-200 responses in stop event" do + ref = make_ref() + + :telemetry.attach( + "test-plugin-error-#{inspect(ref)}", + [:atex, :xrpc, :request, :stop], + fn _event, _measurements, metadata, _ -> + send(self(), {:stop, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-plugin-error-#{inspect(ref)}") end) + + req = + Req.new( + method: :get, + url: "http://bsky.social/xrpc/app.bsky.actor.getProfile", + plug: ErrorPlug, + retry: false + ) + |> Atex.Telemetry.attach_req_plugin(client_type: :login) + + {:ok, _response} = Req.request(req) + + assert_receive {:stop, + %{status: 500, resource: "app.bsky.actor.getProfile", client_type: :login}} + end + + test "emits exception event on transport error" do + ref = make_ref() + + :telemetry.attach( + "test-plugin-exception-#{inspect(ref)}", + [:atex, :xrpc, :request, :exception], + fn _event, _measurements, metadata, _ -> + send(self(), {:exception, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-plugin-exception-#{inspect(ref)}") end) + + req = + Req.new( + method: :get, + url: "http://bsky.social/xrpc/app.bsky.actor.getProfile", + plug: TransportErrorPlug, + retry: false + ) + |> Atex.Telemetry.attach_req_plugin(client_type: :login) + + {:error, _reason} = Req.request(req) + + assert_receive {:exception, + %{ + kind: :error, + reason: %Req.TransportError{reason: :closed}, + resource: "app.bsky.actor.getProfile", + client_type: :login + }} + end + + test "no-op when telemetry not attached — request still succeeds" do + req = + Req.new( + method: :get, + url: "http://bsky.social/xrpc/app.bsky.actor.getProfile", + plug: OkPlug + ) + |> Atex.Telemetry.attach_req_plugin(client_type: :login) + + assert {:ok, %{status: 200}} = Req.request(req) + end + end + + describe "XRPC client instrumentation" do + defmodule XRPCPlug do + @moduledoc false + import Plug.Conn + def init(opts), do: opts + + def call(conn, _opts) do + conn |> send_resp(200, Jason.encode!(%{did: "did:plc:abc", handle: "user.bsky.social"})) + end + end + + test "LoginClient emits request start/stop events" do + ref = make_ref() + + :telemetry.attach_many( + "test-login-client-#{inspect(ref)}", + [[:atex, :xrpc, :request, :start], [:atex, :xrpc, :request, :stop]], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-login-client-#{inspect(ref)}") end) + + client = Atex.XRPC.LoginClient.new("http://bsky.social", "fake-token", nil) + + Atex.XRPC.get(client, "com.atproto.identity.resolveHandle", + plug: XRPCPlug, + params: [handle: "user.bsky.social"] + ) + + assert_receive {:telemetry, [:atex, :xrpc, :request, :start], %{system_time: _}, + %{ + method: :get, + resource: "com.atproto.identity.resolveHandle", + client_type: :login + }} + + assert_receive {:telemetry, [:atex, :xrpc, :request, :stop], %{duration: _}, + %{status: 200, client_type: :login}} + end + + test "unauthed_get emits request start/stop events" do + ref = make_ref() + + :telemetry.attach_many( + "test-unauthed-#{inspect(ref)}", + [[:atex, :xrpc, :request, :start], [:atex, :xrpc, :request, :stop]], + fn event, measurements, metadata, _ -> + send(self(), {:telemetry, event, measurements, metadata}) + end, + nil + ) + + on_exit(fn -> :telemetry.detach("test-unauthed-#{inspect(ref)}") end) + + Atex.XRPC.unauthed_get("http://bsky.social", "com.atproto.identity.resolveHandle", + plug: XRPCPlug, + params: [handle: "user.bsky.social"] + ) + + assert_receive {:telemetry, [:atex, :xrpc, :request, :start], %{system_time: _}, + %{resource: "com.atproto.identity.resolveHandle", client_type: :unauthed}} + + assert_receive {:telemetry, [:atex, :xrpc, :request, :stop], %{duration: _}, + %{status: 200, client_type: :unauthed}} + end + end +end