defmodule Gemini.Client.HTTP do @moduledoc """ Unified HTTP client for both Gemini and Vertex AI APIs using Req. Supports multiple authentication strategies and provides both regular and streaming request capabilities. ## Rate Limiting All requests are automatically routed through the rate limiter unless `disable_rate_limiter: true` is passed in options. The rate limiter: - Enforces concurrency limits per model - Honors 429 RetryInfo delays from the API - Retries transient failures with backoff - Tracks token usage for budget estimation See `Gemini.RateLimiter` for configuration options. """ alias Gemini.Auth alias Gemini.Config alias Gemini.Error alias Gemini.RateLimiter alias Gemini.Telemetry @doc """ Make a GET request using the configured authentication. """ def get(path, opts \\ []) do auth_config = Config.auth_config() request(:get, path, nil, auth_config, opts) end @doc """ Make a POST request using the configured authentication. """ def post(path, body, opts \\ []) do auth_config = Config.auth_config() request(:post, path, body, auth_config, opts) end @doc """ Make a PATCH request using the configured authentication. """ def patch(path, body, opts \\ []) do auth_config = Config.auth_config() request(:patch, path, body, auth_config, opts) end @doc """ Make a DELETE request using the configured authentication. """ def delete(path, opts \\ []) do auth_config = Config.auth_config() request(:delete, path, nil, auth_config, opts) end @doc """ Make an authenticated HTTP request. ## Options In addition to standard request options, supports rate limiter options: - `:disable_rate_limiter` - Bypass rate limiting (default: false) - `:non_blocking` - Return immediately if rate limited (default: false) - `:max_concurrency_per_model` - Override concurrency limit """ def request(method, path, body, auth_config, opts \\ []) do Config.validate!() case auth_config do nil -> {:error, Error.config_error("No authentication configured")} %{type: auth_type, credentials: credentials} -> execute_authenticated_request(method, path, body, auth_type, credentials, opts) end end # Execute the actual HTTP request with telemetry defp execute_request(method, url, headers, body, opts) do start_time = System.monotonic_time() metadata = Telemetry.build_request_metadata(url, method, opts) measurements = %{system_time: System.system_time()} Telemetry.execute([:gemini, :request, :start], measurements, metadata) timeout = Keyword.get(opts, :timeout, Config.timeout()) req_opts = [ method: method, url: url, headers: headers, receive_timeout: timeout, json: body ] try do result = Req.request(req_opts) |> handle_response() case result do {:ok, _response} -> duration = Telemetry.calculate_duration(start_time) stop_measurements = %{ duration: duration, status: 200 } Telemetry.execute([:gemini, :request, :stop], stop_measurements, metadata) {:error, error} -> Telemetry.execute( [:gemini, :request, :exception], measurements, Map.put(metadata, :reason, error) ) end result rescue exception -> Telemetry.execute( [:gemini, :request, :exception], measurements, Map.put(metadata, :reason, exception) ) reraise exception, __STACKTRACE__ end end @doc """ Stream a POST request for Server-Sent Events using configured authentication. """ def stream_post(path, body, opts \\ []) do auth_config = Config.auth_config() stream_post_with_auth(path, body, auth_config, opts) end @doc """ Stream a POST request with specific authentication configuration. """ def stream_post_with_auth(path, body, auth_config, opts \\ []) do Config.validate!() start_time = System.monotonic_time() case auth_config do nil -> {:error, Error.config_error("No authentication configured")} %{type: auth_type, credentials: credentials} -> stream_with_auth(path, body, auth_type, credentials, opts, start_time) end end @doc """ Raw streaming POST with full URL (used by streaming manager). """ def stream_post_raw(url, body, headers, opts \\ []) do timeout = Keyword.get(opts, :timeout, Config.timeout()) req_opts = [ url: url, headers: headers, receive_timeout: timeout, json: body, into: :self ] case Req.post(req_opts) do {:ok, %Req.Response{status: status, body: body}} when status in 200..299 -> case parse_sse_stream(body) do {:ok, events} -> {:ok, events} {:error, parse_error} -> {:error, parse_error} end {:ok, %Req.Response{status: status}} -> {:error, Error.http_error(status, "Stream request failed")} {:error, reason} -> {:error, Error.network_error(reason)} end end # Private functions defp execute_authenticated_request(method, path, body, auth_type, credentials, opts) do url = build_authenticated_url(auth_type, path, credentials) case Auth.build_headers(auth_type, credentials) do {:ok, headers} -> model = extract_model_from_path(path) request_fn = fn -> execute_request(method, url, headers, body, opts) end maybe_rate_limited_request(request_fn, model, opts) {:error, reason} -> {:error, Error.auth_error(reason)} end end defp maybe_rate_limited_request(request_fn, model, opts) do if Keyword.get(opts, :disable_rate_limiter, false) do request_fn.() else RateLimiter.execute_with_usage_tracking(request_fn, model, opts) end end defp stream_with_auth(path, body, auth_type, credentials, opts, start_time) do url = build_authenticated_url(auth_type, path, credentials) case Auth.build_headers(auth_type, credentials) do {:ok, headers} -> sse_url = build_sse_url(url, opts) stream_id = Telemetry.generate_stream_id() metadata = Telemetry.build_stream_metadata(sse_url, :post, stream_id, opts) measurements = %{system_time: System.system_time()} Telemetry.execute([:gemini, :stream, :start], measurements, metadata) req_opts = stream_request_opts(sse_url, headers, body, opts) try do Req.post(req_opts) |> handle_stream_response(start_time, metadata, measurements) rescue exception -> emit_stream_exception(metadata, measurements, exception) reraise exception, __STACKTRACE__ end {:error, reason} -> {:error, Error.auth_error(reason)} end end defp stream_request_opts(url, headers, body, opts) do timeout = Keyword.get(opts, :timeout, Config.timeout()) [ url: url, headers: headers, receive_timeout: timeout, json: body, into: :self ] end defp build_sse_url(url, opts) do if Keyword.get(opts, :add_sse_params, true) do cond do String.contains?(url, "alt=sse") -> url String.contains?(url, "?") -> "#{url}&alt=sse" true -> "#{url}?alt=sse" end else url end end defp handle_stream_response( {:ok, %Req.Response{status: status, body: body}}, start_time, metadata, measurements ) when status in 200..299 do handle_stream_body(body, start_time, metadata, measurements) end defp handle_stream_response( {:ok, %Req.Response{status: status}}, _start_time, metadata, measurements ) do error = {:http_error, status, "Stream request failed"} emit_stream_exception(metadata, measurements, error) {:error, Error.http_error(status, "Stream request failed")} end defp handle_stream_response({:error, reason}, _start_time, metadata, measurements) do emit_stream_exception(metadata, measurements, reason) {:error, Error.network_error(reason)} end defp handle_stream_body(body, start_time, metadata, measurements) do case parse_sse_stream(body) do {:ok, events} -> duration = Telemetry.calculate_duration(start_time) stop_measurements = %{ total_duration: duration, total_chunks: length(events) } Telemetry.execute([:gemini, :stream, :stop], stop_measurements, metadata) {:ok, events} {:error, parse_error} -> emit_stream_exception(metadata, measurements, parse_error) {:error, parse_error} end end defp emit_stream_exception(metadata, measurements, reason) do Telemetry.execute( [:gemini, :stream, :exception], measurements, Map.put(metadata, :reason, reason) ) end defp build_authenticated_url(auth_type, path, credentials) do base_url = Auth.get_base_url(auth_type, credentials) # Check if this is a model-specific endpoint (contains ":" separator) # or a general endpoint like "models" for listing if String.contains?(path, ":") do # Model-specific endpoint, use the auth strategy to build the path full_path = Auth.build_path( auth_type, extract_model_from_path(path), extract_endpoint_from_path(path), credentials ) "#{base_url}/#{full_path}" else # General endpoint (like "models"), use path directly "#{base_url}/#{path}" end end defp extract_model_from_path(path) do # Extract model from paths like "models/gemini-2.0-flash:generateContent" [model_path | _rest] = String.split(path, ":") model_path |> String.replace_prefix("models/", "") |> String.trim_leading("/") end defp extract_endpoint_from_path(path) do # Extract endpoint from paths like "models/gemini-2.0-flash:generateContent" path |> String.split(":") |> List.last() |> String.split("?") |> hd() end defp handle_response({:ok, %Req.Response{status: status, body: body}}) when status in 200..299 do case body do decoded when is_map(decoded) -> {:ok, decoded} json_string when is_binary(json_string) -> case Jason.decode(json_string) do {:ok, decoded} -> {:ok, decoded} {:error, _} -> {:error, Error.invalid_response("Invalid JSON response")} end _ -> {:error, Error.invalid_response("Invalid response format")} end end defp handle_response({:ok, %Req.Response{status: status, body: body}}) do {error_info, error_details} = parse_error_body(body, status) {:error, Error.api_error(status, error_info, error_details)} end defp handle_response({:error, reason}) do {:error, Error.network_error(reason)} end defp build_default_error(status) do message = %{"message" => "HTTP #{status}"} {message, %{"error" => message}} end defp parse_error_body(%{"error" => error} = decoded, _status), do: {error, decoded} defp parse_error_body(body, status) when is_binary(body) do case Jason.decode(body) do {:ok, decoded} -> parse_error_body(decoded, status) _ -> build_default_error(status) end end defp parse_error_body(decoded, _status) when is_map(decoded) do {decoded, %{"error" => decoded}} end defp parse_error_body(_body, status), do: build_default_error(status) # Parse Server-Sent Events format defp parse_sse_stream(data) when is_binary(data) do events = data |> String.split("\n\n") |> Enum.filter(&(String.trim(&1) != "")) |> Enum.map(&parse_sse_event/1) |> Enum.filter(&(&1 != nil)) {:ok, events} rescue exception -> {:error, Error.invalid_response("Failed to parse SSE stream: #{Exception.message(exception)}")} end defp parse_sse_stream(_), do: {:error, Error.invalid_response("Invalid SSE stream payload")} defp parse_sse_event(event_data) do event_data |> String.split("\n") |> Enum.reduce(%{}, &parse_sse_line/2) |> extract_sse_data() end defp parse_sse_line(line, acc) do case String.split(line, ": ", parts: 2) do ["data", json_data] -> maybe_put_sse_data(acc, json_data) [field, value] -> Map.put(acc, String.to_atom(field), value) _ -> acc end end defp maybe_put_sse_data(acc, json_data) do case Jason.decode(json_data) do {:ok, decoded} -> Map.put(acc, :data, decoded) _ -> acc end end defp extract_sse_data(%{data: data}), do: data defp extract_sse_data(_), do: nil end