defmodule PhoenixKit.Modules.AI do @moduledoc """ Main context for PhoenixKit AI system. Provides AI endpoint management and usage tracking for AI API requests. ## Architecture Each **Endpoint** is a unified configuration that combines: - Provider credentials (api_key, base_url, provider_settings) - Model selection (single model per endpoint) - Generation parameters (temperature, max_tokens, etc.) Users create as many endpoints as needed, each representing one complete AI configuration ready for making API requests. ## Core Functions ### System Management - `enabled?/0` - Check if AI module is enabled - `enable_system/0` - Enable the AI module - `disable_system/0` - Disable the AI module - `get_config/0` - Get module configuration with statistics ### Endpoint CRUD - `list_endpoints/1` - List all endpoints with filters - `get_endpoint!/1` - Get endpoint by UUID (raises) - `get_endpoint/1` - Get endpoint by UUID - `create_endpoint/1` - Create new endpoint - `update_endpoint/2` - Update existing endpoint - `delete_endpoint/1` - Delete endpoint ### Completion API - `ask/3` - Simple single-turn completion - `complete/3` - Multi-turn chat completion - `embed/3` - Generate embeddings ### Usage Tracking - `list_requests/1` - List requests with pagination/filters - `create_request/1` - Log a new request - `get_usage_stats/1` - Get aggregated statistics - `get_dashboard_stats/0` - Get stats for dashboard display ## Usage Examples # Enable the module PhoenixKit.Modules.AI.enable_system() # Create an endpoint {:ok, endpoint} = PhoenixKit.Modules.AI.create_endpoint(%{ name: "Claude Fast", provider: "openrouter", api_key: "sk-or-v1-...", model: "anthropic/claude-3-haiku", temperature: 0.7 }) # Use the endpoint {:ok, response} = PhoenixKit.Modules.AI.ask(endpoint.uuid, "Hello!") # Extract the response text {:ok, text} = PhoenixKit.Modules.AI.extract_content(response) """ use PhoenixKit.Module import Ecto.Query, warn: false require Logger alias PhoenixKit.Dashboard.Tab alias PhoenixKit.Modules.AI.Endpoint alias PhoenixKit.Modules.AI.Prompt alias PhoenixKit.Modules.AI.Request alias PhoenixKit.PubSub.Manager, as: PubSub alias PhoenixKit.Settings alias PhoenixKit.Utils.Date, as: UtilsDate alias PhoenixKit.Utils.UUID, as: UUIDUtils # =========================================== # PUBSUB TOPICS # =========================================== @endpoints_topic "phoenix_kit:ai:endpoints" @prompts_topic "phoenix_kit:ai:prompts" @requests_topic "phoenix_kit:ai:requests" @doc """ Returns the PubSub topic for AI endpoints. Subscribe to this topic to receive real-time updates. """ def endpoints_topic, do: @endpoints_topic @doc """ Returns the PubSub topic for AI prompts. """ def prompts_topic, do: @prompts_topic @doc """ Returns the PubSub topic for AI requests/usage. """ def requests_topic, do: @requests_topic @doc """ Subscribes the current process to AI endpoint changes. """ def subscribe_endpoints do PubSub.subscribe(@endpoints_topic) end @doc """ Subscribes the current process to AI prompt changes. """ def subscribe_prompts do PubSub.subscribe(@prompts_topic) end @doc """ Subscribes the current process to AI request/usage changes. """ def subscribe_requests do PubSub.subscribe(@requests_topic) end # =========================================== # HELPERS # =========================================== defp repo do PhoenixKit.RepoHelper.repo() end defp broadcast_endpoint_change(result, event) do case result do {:ok, endpoint} -> PubSub.broadcast(@endpoints_topic, {event, endpoint}) {:ok, endpoint} error -> error end end defp broadcast_prompt_change(result, event) do case result do {:ok, prompt} -> PubSub.broadcast(@prompts_topic, {event, prompt}) {:ok, prompt} error -> error end end defp broadcast_request_change(result, event) do case result do {:ok, request} -> PubSub.broadcast(@requests_topic, {event, request}) {:ok, request} error -> error end end # =========================================== # SYSTEM MANAGEMENT # =========================================== @doc """ Checks if the AI module is enabled. """ @impl PhoenixKit.Module def enabled? do Settings.get_boolean_setting("ai_enabled", false) end @doc """ Enables the AI module. """ @impl PhoenixKit.Module def enable_system do Settings.update_boolean_setting_with_module("ai_enabled", true, "ai") end @doc """ Disables the AI module. """ @impl PhoenixKit.Module def disable_system do Settings.update_boolean_setting_with_module("ai_enabled", false, "ai") end @doc """ Gets the AI module configuration with statistics. """ @impl PhoenixKit.Module def get_config do %{ enabled: enabled?(), endpoints_count: count_endpoints(), total_requests: count_requests(), total_tokens: sum_tokens() } end # =========================================== # MODULE BEHAVIOUR CALLBACKS # =========================================== @impl PhoenixKit.Module def module_key, do: "ai" @impl PhoenixKit.Module def module_name, do: "AI" @impl PhoenixKit.Module def permission_metadata do %{ key: "ai", label: "AI", icon: "hero-sparkles", description: "AI endpoints, prompts, and usage tracking" } end @impl PhoenixKit.Module def admin_tabs do [ %Tab{ id: :admin_ai, label: "AI", icon: "hero-cpu-chip", path: "ai", priority: 550, level: :admin, permission: "ai", match: :prefix, group: :admin_modules, subtab_display: :when_active, highlight_with_subtabs: false }, %Tab{ id: :admin_ai_endpoints, label: "Endpoints", icon: "hero-server-stack", path: "ai/endpoints", priority: 551, level: :admin, permission: "ai", parent: :admin_ai }, %Tab{ id: :admin_ai_prompts, label: "Prompts", icon: "hero-document-text", path: "ai/prompts", priority: 552, level: :admin, permission: "ai", parent: :admin_ai }, %Tab{ id: :admin_ai_usage, label: "Usage", icon: "hero-chart-bar", path: "ai/usage", priority: 553, level: :admin, permission: "ai", parent: :admin_ai } ] end # =========================================== # ENDPOINT CRUD # =========================================== @doc """ Lists all AI endpoints. ## Options - `:provider` - Filter by provider type - `:enabled` - Filter by enabled status - `:preload` - Associations to preload ## Examples PhoenixKit.Modules.AI.list_endpoints() PhoenixKit.Modules.AI.list_endpoints(provider: "openrouter", enabled: true) """ def list_endpoints(opts \\ []) do sort_by = Keyword.get(opts, :sort_by, :sort_order) sort_dir = Keyword.get(opts, :sort_dir, :asc) # Always paginate, default to page 1 and ensure it's > 0 page = Keyword.get(opts, :page, 1) |> max(1) page_size = Keyword.get(opts, :page_size, 20) # Build base query with filters (no sorting yet) base_query = from(e in Endpoint) base_query = case Keyword.get(opts, :provider) do nil -> base_query provider -> where(base_query, [e], e.provider == ^provider) end base_query = case Keyword.get(opts, :enabled) do nil -> base_query enabled -> where(base_query, [e], e.enabled == ^enabled) end # Count on base query BEFORE applying sorting (which may add group_by) total = repo().aggregate(base_query, :count) # Now apply sorting (may add group_by for usage/tokens/cost/last_used) query = apply_endpoint_sorting(base_query, sort_by, sort_dir) query = case Keyword.get(opts, :preload) do nil -> query preloads -> preload(query, ^preloads) end offset = (page - 1) * page_size endpoints = query |> limit(^page_size) |> offset(^offset) |> repo().all() # Always return the same shape {endpoints, total} end defp apply_endpoint_sorting(query, :usage, dir) do # Sort by total request count using a subquery to avoid GROUP BY issues stats_subquery = from(r in Request, where: not is_nil(r.endpoint_uuid), group_by: r.endpoint_uuid, select: %{endpoint_uuid: r.endpoint_uuid, count: count()} ) from(e in query, left_join: s in subquery(stats_subquery), on: s.endpoint_uuid == e.uuid, order_by: [{^dir, coalesce(s.count, 0)}] ) end defp apply_endpoint_sorting(query, :tokens, dir) do # Sort by total tokens used stats_subquery = from(r in Request, where: not is_nil(r.endpoint_uuid), group_by: r.endpoint_uuid, select: %{endpoint_uuid: r.endpoint_uuid, total: coalesce(sum(r.total_tokens), 0)} ) from(e in query, left_join: s in subquery(stats_subquery), on: s.endpoint_uuid == e.uuid, order_by: [{^dir, coalesce(s.total, 0)}] ) end defp apply_endpoint_sorting(query, :cost, dir) do # Sort by total cost stats_subquery = from(r in Request, where: not is_nil(r.endpoint_uuid), group_by: r.endpoint_uuid, select: %{endpoint_uuid: r.endpoint_uuid, total: coalesce(sum(r.cost_cents), 0)} ) from(e in query, left_join: s in subquery(stats_subquery), on: s.endpoint_uuid == e.uuid, order_by: [{^dir, coalesce(s.total, 0)}] ) end defp apply_endpoint_sorting(query, :last_used, dir) do # Sort by most recent request time stats_subquery = from(r in Request, where: not is_nil(r.endpoint_uuid), group_by: r.endpoint_uuid, select: %{endpoint_uuid: r.endpoint_uuid, last_used: max(r.inserted_at)} ) from(e in query, left_join: s in subquery(stats_subquery), on: s.endpoint_uuid == e.uuid, order_by: [{^dir, s.last_used}] ) end defp apply_endpoint_sorting(query, field, dir) when field in [:name, :enabled, :model, :sort_order] do order_by(query, [e], [{^dir, field(e, ^field)}]) end defp apply_endpoint_sorting(query, _field, _dir) do # Default sorting order_by(query, [e], asc: e.sort_order, desc: e.inserted_at) end @doc """ Returns usage statistics for each endpoint. Returns a map of endpoint_uuid => %{request_count, total_tokens, total_cost, last_used_at} """ def get_endpoint_usage_stats do query = from(r in Request, where: not is_nil(r.endpoint_uuid), group_by: r.endpoint_uuid, select: { r.endpoint_uuid, %{ request_count: count(), total_tokens: coalesce(sum(r.total_tokens), 0), total_cost: coalesce(sum(r.cost_cents), 0), last_used_at: max(r.inserted_at) } } ) query |> repo().all() |> Map.new() end @doc """ Gets a single endpoint by UUID. Raises `Ecto.NoResultsError` if the endpoint does not exist. """ def get_endpoint!(id) do case get_endpoint(id) do nil -> raise Ecto.NoResultsError, queryable: Endpoint endpoint -> endpoint end end @doc """ Gets a single endpoint by UUID. Accepts a UUID string (e.g., "550e8400-e29b-41d4-a716-446655440000"). Returns `nil` if the endpoint does not exist. """ def get_endpoint(id) when is_binary(id) do if UUIDUtils.valid?(id) do repo().get_by(Endpoint, uuid: id) else nil end end def get_endpoint(_), do: nil @doc """ Resolves an endpoint from an ID (UUID string) or Endpoint struct. ## Examples {:ok, endpoint} = PhoenixKit.Modules.AI.resolve_endpoint("019abc12-3456-7def-8901-234567890abc") {:ok, endpoint} = PhoenixKit.Modules.AI.resolve_endpoint(endpoint) """ def resolve_endpoint(id) when is_binary(id) do case get_endpoint(id) do nil -> {:error, "Endpoint not found"} endpoint -> {:ok, endpoint} end end def resolve_endpoint(%Endpoint{} = endpoint), do: {:ok, endpoint} def resolve_endpoint(_), do: {:error, "Invalid endpoint identifier"} @doc """ Creates a new AI endpoint. ## Examples {:ok, endpoint} = PhoenixKit.Modules.AI.create_endpoint(%{ name: "Claude Fast", provider: "openrouter", api_key: "sk-or-v1-...", model: "anthropic/claude-3-haiku", temperature: 0.7 }) """ def create_endpoint(attrs) do %Endpoint{} |> Endpoint.changeset(attrs) |> repo().insert() |> broadcast_endpoint_change(:endpoint_created) end @doc """ Updates an existing AI endpoint. """ def update_endpoint(%Endpoint{} = endpoint, attrs) do endpoint |> Endpoint.changeset(attrs) |> repo().update() |> broadcast_endpoint_change(:endpoint_updated) end @doc """ Deletes an AI endpoint. """ def delete_endpoint(%Endpoint{} = endpoint) do repo().delete(endpoint) |> broadcast_endpoint_change(:endpoint_deleted) end @doc """ Returns an endpoint changeset for use in forms. """ def change_endpoint(%Endpoint{} = endpoint, attrs \\ %{}) do Endpoint.changeset(endpoint, attrs) end @doc """ Marks an endpoint as validated by updating its last_validated_at timestamp. """ def mark_endpoint_validated(%Endpoint{} = endpoint) do endpoint |> Endpoint.validation_changeset() |> repo().update() end @doc """ Counts the total number of endpoints. """ def count_endpoints do repo().aggregate(Endpoint, :count) end @doc """ Counts the number of enabled endpoints. """ def count_enabled_endpoints do query = from(e in Endpoint, where: e.enabled == true) repo().aggregate(query, :count) end # =========================================== # PROMPT CRUD # =========================================== @doc """ Lists all AI prompts. ## Options - `:sort_by` - Field to sort by (default: :sort_order) - `:sort_dir` - Sort direction, :asc or :desc (default: :asc) - `:enabled` - Filter by enabled status ## Examples PhoenixKit.Modules.AI.list_prompts() PhoenixKit.Modules.AI.list_prompts(sort_by: :name, sort_dir: :asc) PhoenixKit.Modules.AI.list_prompts(enabled: true) """ def list_prompts(opts \\ []) do sort_by = Keyword.get(opts, :sort_by, :sort_order) sort_dir = Keyword.get(opts, :sort_dir, :asc) page = Keyword.get(opts, :page) page_size = Keyword.get(opts, :page_size, 20) query = from(p in Prompt) query = case Keyword.get(opts, :enabled) do nil -> query enabled -> where(query, [p], p.enabled == ^enabled) end query = order_by(query, [p], [{^sort_dir, field(p, ^sort_by)}]) # If page is provided, return paginated results with total count if page do total = repo().aggregate(query, :count) offset = (page - 1) * page_size prompts = query |> limit(^page_size) |> offset(^offset) |> repo().all() {prompts, total} else # No pagination - return all (backwards compatible) repo().all(query) end end @doc """ Lists only enabled prompts. Convenience wrapper for `list_prompts(enabled: true)`. ## Examples PhoenixKit.Modules.AI.list_enabled_prompts() """ def list_enabled_prompts do list_prompts(enabled: true) end @doc """ Gets a single prompt by UUID. Raises `Ecto.NoResultsError` if the prompt does not exist. """ def get_prompt!(id) do case get_prompt(id) do nil -> raise Ecto.NoResultsError, queryable: Prompt prompt -> prompt end end @doc """ Gets a single prompt by UUID. Accepts a UUID string (e.g., "550e8400-e29b-41d4-a716-446655440000"). Returns `nil` if the prompt does not exist. """ def get_prompt(id) when is_binary(id) do if UUIDUtils.valid?(id) do repo().get_by(Prompt, uuid: id) else nil end end def get_prompt(_), do: nil @doc """ Gets a prompt by slug. Returns `nil` if the prompt does not exist. """ def get_prompt_by_slug(slug) when is_binary(slug) do repo().get_by(Prompt, slug: slug) end @doc """ Creates a new AI prompt. ## Examples {:ok, prompt} = PhoenixKit.Modules.AI.create_prompt(%{ name: "Translator", content: "Translate the following text to {{Language}}:\\n\\n{{Text}}" }) """ def create_prompt(attrs) do %Prompt{} |> Prompt.changeset(attrs) |> repo().insert() |> broadcast_prompt_change(:prompt_created) end @doc """ Updates an existing AI prompt. """ def update_prompt(%Prompt{} = prompt, attrs) do prompt |> Prompt.changeset(attrs) |> repo().update() |> broadcast_prompt_change(:prompt_updated) end @doc """ Deletes an AI prompt. """ def delete_prompt(%Prompt{} = prompt) do repo().delete(prompt) |> broadcast_prompt_change(:prompt_deleted) end @doc """ Returns a prompt changeset for use in forms. """ def change_prompt(%Prompt{} = prompt, attrs \\ %{}) do Prompt.changeset(prompt, attrs) end @doc """ Increments the usage count for a prompt and updates last_used_at. """ def record_prompt_usage(%Prompt{} = prompt) do prompt |> Prompt.usage_changeset() |> repo().update() end @doc """ Counts the total number of prompts. """ def count_prompts do repo().aggregate(Prompt, :count) end @doc """ Counts the number of enabled prompts. """ def count_enabled_prompts do query = from(p in Prompt, where: p.enabled == true) repo().aggregate(query, :count) end @doc """ Resolves a prompt from various input types. Accepts: - UUID string (e.g., "019abc12-3456-7def-8901-234567890abc") - String slug (e.g., "my-prompt") - Prompt struct (returned as-is) Returns `{:ok, prompt}` or `{:error, reason}`. """ def resolve_prompt(%Prompt{} = prompt), do: {:ok, prompt} def resolve_prompt(id_or_slug) when is_binary(id_or_slug) do if UUIDUtils.valid?(id_or_slug) do # It's a UUID case get_prompt(id_or_slug) do nil -> {:error, "Prompt not found"} prompt -> {:ok, prompt} end else # It's a slug case get_prompt_by_slug(id_or_slug) do nil -> {:error, "Prompt not found"} prompt -> {:ok, prompt} end end end def resolve_prompt(_), do: {:error, "Invalid prompt identifier"} @doc """ Renders a prompt by replacing variables with provided values. Returns `{:ok, rendered_text}` or `{:error, reason}`. """ def render_prompt(prompt_uuid, variables \\ %{}) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do Prompt.render(prompt, variables) end end @doc """ Increments the usage count for a prompt and updates last_used_at. """ def increment_prompt_usage(prompt_uuid) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do record_prompt_usage(prompt) end end @doc """ Makes an AI completion using a prompt template. The prompt content is rendered with the provided variables and sent as the user message. """ def ask_with_prompt(endpoint_uuid, prompt_uuid, variables \\ %{}, opts \\ []) do with {:ok, prompt} <- resolve_prompt(prompt_uuid), {:ok, _} <- validate_prompt(prompt), {:ok, rendered} <- Prompt.render(prompt, variables) do # Pass prompt info to ask for request logging opts_with_prompt = opts |> Keyword.put(:prompt_uuid, prompt.uuid) |> Keyword.put(:prompt_name, prompt.name) case ask(endpoint_uuid, rendered, opts_with_prompt) do {:ok, response} -> # Only increment usage on successful completion increment_prompt_usage(prompt_uuid) {:ok, response} error -> error end end end @doc """ Makes an AI completion with a prompt template as the system message. The prompt is rendered and used as the system message, with the user_message as the user message. """ def complete_with_system_prompt(endpoint_uuid, prompt_uuid, variables, user_message, opts \\ []) do with {:ok, prompt} <- resolve_prompt(prompt_uuid), {:ok, _} <- validate_prompt(prompt), {:ok, system_prompt} <- Prompt.render(prompt, variables) do # Build messages with system prompt messages = [ %{role: "system", content: system_prompt}, %{role: "user", content: user_message} ] # Pass prompt info to complete for request logging opts_with_prompt = opts |> Keyword.put(:prompt_uuid, prompt.uuid) |> Keyword.put(:prompt_name, prompt.name) case complete(endpoint_uuid, messages, opts_with_prompt) do {:ok, response} -> # Only increment usage on successful completion increment_prompt_usage(prompt_uuid) {:ok, response} error -> error end end end @doc """ Validates that a prompt is ready for use. Returns `{:ok, prompt}` if valid, or `{:error, reason}` if not. """ def validate_prompt(prompt) do cond do prompt.content == nil or prompt.content == "" -> {:error, "Prompt has no content"} prompt.enabled == false -> {:error, "Prompt is disabled"} true -> {:ok, prompt} end end @doc """ Duplicates a prompt with a new name. """ def duplicate_prompt(prompt_uuid, new_name) when is_binary(new_name) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do create_prompt(%{ name: new_name, description: prompt.description, content: prompt.content, enabled: prompt.enabled, sort_order: prompt.sort_order, metadata: prompt.metadata }) end end @doc """ Enables a prompt. """ def enable_prompt(prompt_uuid) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do update_prompt(prompt, %{enabled: true}) end end @doc """ Disables a prompt. """ def disable_prompt(prompt_uuid) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do update_prompt(prompt, %{enabled: false}) end end @doc """ Gets the variables defined in a prompt. """ def get_prompt_variables(prompt_uuid) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do {:ok, prompt.variables || []} end end @doc """ Previews a rendered prompt without making an AI call. """ def preview_prompt(prompt_uuid, variables \\ %{}) do render_prompt(prompt_uuid, variables) end @doc """ Validates that all required variables are provided for a prompt. """ def validate_prompt_variables(prompt_uuid, variables) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do Prompt.validate_variables(prompt, variables) end end @doc """ Searches prompts by name, description, or content. """ def search_prompts(query, opts \\ []) when is_binary(query) do pattern = "%#{query}%" limit = Keyword.get(opts, :limit, 50) base_query = from(p in Prompt, where: ilike(p.name, ^pattern) or ilike(p.description, ^pattern) or ilike(p.content, ^pattern), order_by: [asc: p.sort_order, desc: p.inserted_at], limit: ^limit ) base_query = case Keyword.get(opts, :enabled) do nil -> base_query enabled -> where(base_query, [p], p.enabled == ^enabled) end repo().all(base_query) end @doc """ Finds all prompts that use a specific variable. """ def get_prompts_with_variable(variable_name) when is_binary(variable_name) do query = from(p in Prompt, where: ^variable_name in p.variables, order_by: [asc: p.sort_order, desc: p.inserted_at] ) repo().all(query) end @doc """ Validates that the content has valid variable syntax. """ def validate_prompt_content(content) when is_binary(content) do all_patterns = Regex.scan(~r/\{\{([^}]+)\}\}/, content) invalid = all_patterns |> Enum.map(fn [_full, inner] -> inner end) |> Enum.reject(fn inner -> Regex.match?(~r/^\w+$/, inner) end) if Enum.empty?(invalid) do :ok else {:error, invalid} end end def validate_prompt_content(_), do: {:error, "Content must be a string"} @doc """ Gets usage statistics for all prompts. """ def get_prompt_usage_stats(opts \\ []) do query = from(p in Prompt, select: %{ prompt: p, usage_count: p.usage_count, last_used_at: p.last_used_at }, order_by: [desc: p.usage_count, desc: p.last_used_at] ) query = case Keyword.get(opts, :enabled) do nil -> query enabled -> where(query, [p], p.enabled == ^enabled) end query = case Keyword.get(opts, :limit) do nil -> query limit -> limit(query, ^limit) end repo().all(query) end @doc """ Resets the usage statistics for a prompt. """ def reset_prompt_usage(prompt_uuid) do with {:ok, prompt} <- resolve_prompt(prompt_uuid) do prompt |> Ecto.Changeset.change(%{usage_count: 0, last_used_at: nil}) |> repo().update() end end @doc """ Updates the sort order for multiple prompts. Accepts prompt UUIDs. """ def reorder_prompts(order_list) when is_list(order_list) do repo().transaction(fn -> Enum.each(order_list, fn {id, sort_order} -> build_prompt_uuid_query(id) |> repo().update_all(set: [sort_order: sort_order]) end) end) :ok end defp build_prompt_uuid_query(id) when is_binary(id) do if UUIDUtils.valid?(id) do from(p in Prompt, where: p.uuid == ^id) else from(p in Prompt, where: false) end end defp build_prompt_uuid_query(_), do: from(p in Prompt, where: false) # =========================================== # USAGE TRACKING (REQUESTS) # =========================================== @doc """ Lists AI requests with pagination and filters. ## Options - `:page` - Page number (default: 1) - `:page_size` - Results per page (default: 20) - `:endpoint_uuid` - Filter by endpoint - `:user_uuid` - Filter by user - `:status` - Filter by status - `:model` - Filter by model - `:source` - Filter by source (from metadata) - `:since` - Filter by date (requests after this date) - `:preload` - Associations to preload ## Returns `{requests, total_count}` """ def list_requests(opts \\ []) do page = Keyword.get(opts, :page, 1) page_size = Keyword.get(opts, :page_size, 20) offset = (page - 1) * page_size sort_by = Keyword.get(opts, :sort_by, :inserted_at) sort_dir = Keyword.get(opts, :sort_dir, :desc) base_query = from(r in Request) base_query = apply_request_filters(base_query, opts) base_query = apply_request_sorting(base_query, sort_by, sort_dir) total = repo().aggregate(base_query, :count) query = base_query |> limit(^page_size) |> offset(^offset) query = case Keyword.get(opts, :preload) do nil -> query preloads -> preload(query, ^preloads) end requests = repo().all(query) {requests, total} end defp apply_request_sorting(query, field, dir) when field in [ :inserted_at, :model, :total_tokens, :latency_ms, :cost_cents, :status, :endpoint_name ] do order_by(query, [r], [{^dir, field(r, ^field)}]) end defp apply_request_sorting(query, _field, _dir) do order_by(query, [r], desc: r.inserted_at) end @doc """ Gets a single request by UUID. """ def get_request!(id) do case get_request(id) do nil -> raise Ecto.NoResultsError, queryable: Request request -> request end end @doc """ Gets a single request by UUID. Accepts a UUID string (e.g., "550e8400-e29b-41d4-a716-446655440000"). Returns `nil` if the request does not exist. """ def get_request(id) when is_binary(id) do if UUIDUtils.valid?(id) do repo().get_by(Request, uuid: id) else nil end end def get_request(_), do: nil @doc """ Creates a new AI request record. Used to log every AI API call for tracking and statistics. """ def create_request(attrs) do %Request{} |> Request.changeset(attrs) |> repo().insert() |> broadcast_request_change(:request_created) end @doc """ Counts the total number of requests. """ def count_requests do repo().aggregate(Request, :count) end @doc """ Sums the total tokens used across all requests. """ def sum_tokens do repo().aggregate(Request, :sum, :total_tokens) || 0 end defp apply_request_filters(query, opts) do query |> maybe_filter_by(:endpoint_uuid, Keyword.get(opts, :endpoint_uuid)) |> maybe_filter_by(:user_uuid, Keyword.get(opts, :user_uuid)) |> maybe_filter_by(:status, Keyword.get(opts, :status)) |> maybe_filter_by(:model, Keyword.get(opts, :model)) |> maybe_filter_by(:source, Keyword.get(opts, :source)) |> maybe_filter_since(Keyword.get(opts, :since)) end defp maybe_filter_by(query, _field, nil), do: query defp maybe_filter_by(query, :endpoint_uuid, uuid) when is_binary(uuid) do where(query, [r], r.endpoint_uuid == ^uuid) end defp maybe_filter_by(query, :user_uuid, uuid) when is_binary(uuid) do where(query, [r], r.user_uuid == ^uuid) end defp maybe_filter_by(query, :status, status), do: where(query, [r], r.status == ^status) defp maybe_filter_by(query, :model, model), do: where(query, [r], r.model == ^model) defp maybe_filter_by(query, :source, source), do: where(query, [r], fragment("?->>'source' = ?", r.metadata, ^source)) defp maybe_filter_since(query, nil), do: query defp maybe_filter_since(query, date), do: where(query, [r], r.inserted_at >= ^date) @doc """ Returns filter options for requests (distinct endpoints, models, and sources). """ def get_request_filter_options do endpoints_query = from(r in Request, where: not is_nil(r.endpoint_uuid) and not is_nil(r.endpoint_name), distinct: true, select: {r.endpoint_uuid, r.endpoint_name}, order_by: r.endpoint_name ) models_query = from(r in Request, where: not is_nil(r.model), distinct: r.model, select: r.model, order_by: r.model ) # Query unique sources from metadata JSONB field sources_query = from(r in Request, where: not is_nil(fragment("?->>'source'", r.metadata)), distinct: fragment("?->>'source'", r.metadata), select: fragment("?->>'source'", r.metadata), order_by: fragment("?->>'source'", r.metadata) ) %{ endpoints: repo().all(endpoints_query), models: repo().all(models_query), statuses: Request.valid_statuses(), sources: repo().all(sources_query) } end # =========================================== # STATISTICS # =========================================== @doc """ Gets aggregated usage statistics. ## Options - `:since` - Start date for statistics - `:until` - End date for statistics - `:endpoint_uuid` - Filter by endpoint ## Returns Map with statistics including total_requests, total_tokens, success_rate, etc. """ def get_usage_stats(opts \\ []) do base_query = from(r in Request) base_query = apply_request_filters(base_query, opts) total_requests = repo().aggregate(base_query, :count) total_tokens = repo().aggregate(base_query, :sum, :total_tokens) || 0 total_cost = repo().aggregate(base_query, :sum, :cost_cents) || 0 avg_latency = repo().aggregate(base_query, :avg, :latency_ms) success_query = where(base_query, [r], r.status == "success") success_count = repo().aggregate(success_query, :count) success_rate = if total_requests > 0 do Float.round(success_count / total_requests * 100, 1) else 0.0 end %{ total_requests: total_requests, total_tokens: total_tokens, total_cost_cents: total_cost, success_count: success_count, error_count: total_requests - success_count, success_rate: success_rate, avg_latency_ms: decimal_to_int(avg_latency) } end # Convert Decimal or number to integer, handling nil defp decimal_to_int(nil), do: nil defp decimal_to_int(%Decimal{} = d), do: d |> Decimal.round() |> Decimal.to_integer() defp decimal_to_int(n) when is_float(n), do: round(n) defp decimal_to_int(n) when is_integer(n), do: n @doc """ Gets dashboard statistics for display. Returns stats for the last 30 days plus all-time totals. """ def get_dashboard_stats do thirty_days_ago = UtilsDate.utc_now() |> DateTime.add(-30, :day) today_start = Date.utc_today() |> DateTime.new!(~T[00:00:00], "Etc/UTC") all_time = get_usage_stats() last_30_days = get_usage_stats(since: thirty_days_ago) today = get_usage_stats(since: today_start) tokens_by_model = get_tokens_by_model(since: thirty_days_ago) requests_by_day = get_requests_by_day(since: thirty_days_ago) %{ all_time: all_time, last_30_days: last_30_days, today: today, tokens_by_model: tokens_by_model, requests_by_day: requests_by_day } end @doc """ Gets token usage grouped by model. """ def get_tokens_by_model(opts \\ []) do base_query = from(r in Request) base_query = apply_request_filters(base_query, opts) query = from r in subquery(base_query), where: not is_nil(r.model) and r.model != "", group_by: r.model, select: %{ model: r.model, total_tokens: sum(r.total_tokens), request_count: count() }, order_by: [desc: sum(r.total_tokens)] repo().all(query) end @doc """ Gets request counts grouped by day. """ def get_requests_by_day(opts \\ []) do base_query = from(r in Request) base_query = apply_request_filters(base_query, opts) query = from r in subquery(base_query), group_by: fragment("DATE(?)", r.inserted_at), select: %{ date: fragment("DATE(?)", r.inserted_at), count: count(), tokens: sum(r.total_tokens) }, order_by: [asc: fragment("DATE(?)", r.inserted_at)] repo().all(query) end # =========================================== # COMPLETION API # =========================================== alias PhoenixKit.Modules.AI.Completion @doc """ Makes a chat completion request using a configured endpoint. ## Parameters - `endpoint_uuid` - Endpoint UUID string or Endpoint struct - `messages` - List of message maps with `:role` and `:content` - `opts` - Optional parameter overrides ## Options All standard completion parameters plus: - `:source` - Override auto-detected source for request tracking ## Examples {:ok, response} = PhoenixKit.Modules.AI.complete(endpoint_uuid, [ %{role: "user", content: "Hello!"} ]) # With system message {:ok, response} = PhoenixKit.Modules.AI.complete(endpoint_uuid, [ %{role: "system", content: "You are a helpful assistant."}, %{role: "user", content: "What is 2+2?"} ]) # With parameter overrides {:ok, response} = PhoenixKit.Modules.AI.complete(endpoint_uuid, messages, temperature: 0.5, max_tokens: 500 ) # With custom source for tracking {:ok, response} = PhoenixKit.Modules.AI.complete(endpoint_uuid, messages, source: "MyModule" ) ## Returns - `{:ok, response}` - Full API response including usage stats - `{:error, reason}` - Error with reason string """ def complete(endpoint_uuid, messages, opts \\ []) do with {:ok, endpoint} <- resolve_endpoint(endpoint_uuid), {:ok, _} <- validate_endpoint(endpoint) do # Capture caller info (source + stacktrace + context) {auto_source, stacktrace, caller_context} = capture_caller_info() # Allow manual override of source, but all debug info is always captured source = Keyword.get(opts, :source) || auto_source # Extract prompt info if present (from ask_with_prompt, complete_with_system_prompt) prompt_info = %{ prompt_uuid: Keyword.get(opts, :prompt_uuid), prompt_name: Keyword.get(opts, :prompt_name) } merged_opts = merge_endpoint_opts(endpoint, opts) case Completion.chat_completion(endpoint, messages, merged_opts) do {:ok, response} -> log_request( endpoint, messages, response, source, stacktrace, caller_context, prompt_info ) {:ok, response} {:error, reason} -> log_failed_request( endpoint, messages, reason, source, stacktrace, caller_context, prompt_info ) {:error, reason} end end end @doc """ Simple helper for single-turn chat completion. ## Parameters - `endpoint_uuid` - Endpoint UUID string or Endpoint struct - `prompt` - User prompt string - `opts` - Optional parameter overrides and system message ## Options All options from `complete/3` plus: - `:system` - System message string - `:source` - Override auto-detected source for request tracking ## Examples # Simple question {:ok, response} = PhoenixKit.Modules.AI.ask(endpoint_uuid, "What is the capital of France?") # With system message {:ok, response} = PhoenixKit.Modules.AI.ask(endpoint_uuid, "Translate: Hello", system: "You are a translator. Translate to French." ) # With custom source for tracking {:ok, response} = PhoenixKit.Modules.AI.ask(endpoint_uuid, "Hello!", source: "Languages" ) # Extract just the text content {:ok, response} = PhoenixKit.Modules.AI.ask(endpoint_uuid, "Hello!") {:ok, text} = PhoenixKit.Modules.AI.extract_content(response) ## Returns Same as `complete/3` """ def ask(endpoint_uuid, prompt, opts \\ []) when is_binary(prompt) do {system, opts} = Keyword.pop(opts, :system) messages = case system do nil -> [%{role: "user", content: prompt}] sys -> [%{role: "system", content: sys}, %{role: "user", content: prompt}] end complete(endpoint_uuid, messages, opts) end @doc """ Makes an embeddings request using a configured endpoint. ## Parameters - `endpoint_uuid` - Endpoint UUID string or Endpoint struct - `input` - Text or list of texts to embed - `opts` - Optional parameter overrides ## Options - `:dimensions` - Override embedding dimensions - `:source` - Override auto-detected source for request tracking ## Examples # Single text {:ok, response} = PhoenixKit.Modules.AI.embed(endpoint_uuid, "Hello, world!") # Multiple texts {:ok, response} = PhoenixKit.Modules.AI.embed(endpoint_uuid, ["Hello", "World"]) # With dimension override {:ok, response} = PhoenixKit.Modules.AI.embed(endpoint_uuid, "Hello", dimensions: 512) # With custom source for tracking {:ok, response} = PhoenixKit.Modules.AI.embed(endpoint_uuid, "Hello", source: "SemanticSearch" ) ## Returns - `{:ok, response}` - Response with embeddings - `{:error, reason}` - Error with reason """ def embed(endpoint_uuid, input, opts \\ []) do with {:ok, endpoint} <- resolve_endpoint(endpoint_uuid), {:ok, _} <- validate_endpoint(endpoint) do # Capture caller info (source + stacktrace + context) {auto_source, stacktrace, caller_context} = capture_caller_info() # Allow manual override of source, but all debug info is always captured source = Keyword.get(opts, :source) || auto_source merged_opts = merge_embedding_opts(endpoint, opts) case Completion.embeddings(endpoint, input, merged_opts) do {:ok, response} -> log_embedding_request(endpoint, input, response, source, stacktrace, caller_context) {:ok, response} {:error, reason} -> log_failed_embedding_request(endpoint, reason, source, stacktrace, caller_context) {:error, reason} end end end @doc """ Extracts the text content from a completion response. ## Examples {:ok, response} = PhoenixKit.Modules.AI.ask(endpoint_uuid, "Hello!") {:ok, text} = PhoenixKit.Modules.AI.extract_content(response) # => "Hello! How can I help you today?" """ defdelegate extract_content(response), to: Completion @doc """ Extracts usage information from a response. ## Examples {:ok, response} = PhoenixKit.Modules.AI.complete(endpoint_uuid, messages) usage = PhoenixKit.Modules.AI.extract_usage(response) # => %{prompt_tokens: 10, completion_tokens: 15, total_tokens: 25} """ defdelegate extract_usage(response), to: Completion # Private helpers for completion API defp validate_endpoint(endpoint) do cond do endpoint.model == nil or endpoint.model == "" -> {:error, "Endpoint has no model configured"} endpoint.api_key == nil or endpoint.api_key == "" -> {:error, "Endpoint has no API key configured"} endpoint.enabled == false -> {:error, "Endpoint is disabled"} true -> {:ok, endpoint} end end defp merge_endpoint_opts(endpoint, opts) do # Endpoint defaults, then user overrides base_opts = [ temperature: endpoint.temperature, max_tokens: endpoint.max_tokens, top_p: endpoint.top_p, top_k: endpoint.top_k, frequency_penalty: endpoint.frequency_penalty, presence_penalty: endpoint.presence_penalty, repetition_penalty: endpoint.repetition_penalty, stop: endpoint.stop, seed: endpoint.seed, # Reasoning/thinking parameters (for models like DeepSeek R1, Qwen QwQ) reasoning_enabled: endpoint.reasoning_enabled, reasoning_effort: endpoint.reasoning_effort, reasoning_max_tokens: endpoint.reasoning_max_tokens, reasoning_exclude: endpoint.reasoning_exclude ] # Filter out nil values and merge with user opts base_opts |> Enum.reject(fn {_k, v} -> is_nil(v) end) |> Keyword.merge(opts) end defp merge_embedding_opts(endpoint, opts) do base_opts = [dimensions: endpoint.dimensions] base_opts |> Enum.reject(fn {_k, v} -> is_nil(v) end) |> Keyword.merge(opts) end # =========================================== # CALLER INFO CAPTURE (for source tracking & debugging) # =========================================== @doc false # Captures full debug context: source, stacktrace, and caller context defp capture_caller_info do # Get process info (stacktrace + memory in one call) [{:current_stacktrace, stack}, {:memory, memory}] = Process.info(self(), [:current_stacktrace, :memory]) # Format stacktrace for storage formatted_stack = format_stacktrace(stack) # Extract clean source from first non-internal caller source = extract_source(stack) # Build caller context with additional debug info caller_context = build_caller_context(memory) {source, formatted_stack, caller_context} end defp format_stacktrace(stack) do stack # Limit depth for storage |> Enum.take(20) |> Enum.map(fn {mod, fun, arity, location} -> mod_str = Atom.to_string(mod) |> String.replace_prefix("Elixir.", "") file = Keyword.get(location, :file, ~c"unknown") |> to_string() line = Keyword.get(location, :line, 0) "#{mod_str}.#{fun}/#{arity} (#{file}:#{line})" end) end defp extract_source(stack) do # Modules to skip (PhoenixKit.Modules.AI internals, Elixir/Erlang core) skip_prefixes = ["PhoenixKit.Modules.AI", "Elixir.PhoenixKit.Modules.AI"] skip_modules = [Process, :proc_lib, :gen_server, :gen, :elixir, :erl_eval] caller = Enum.find(stack, fn {mod, _fun, _arity, _loc} -> mod_str = Atom.to_string(mod) not Enum.any?(skip_prefixes, &String.starts_with?(mod_str, &1)) and mod not in skip_modules end) case caller do {mod, fun, _arity, _loc} -> mod_str = Atom.to_string(mod) |> String.replace_prefix("Elixir.", "") "#{mod_str}.#{fun}" nil -> nil end end defp build_caller_context(memory) do # Get Phoenix request_id from Logger metadata (if in request context) logger_meta = Logger.metadata() %{ request_id: Keyword.get(logger_meta, :request_id), node: node() |> Atom.to_string(), pid: self() |> inspect(), memory_bytes: memory } end # =========================================== # REQUEST LOGGING # =========================================== defp log_request(endpoint, messages, response, source, stacktrace, caller_context, prompt_info) do usage = Completion.extract_usage(response) # Extract response content response_content = case Completion.extract_content(response) do {:ok, content} -> content _ -> nil end # Build the request payload we sent (for debugging) request_payload = %{ model: endpoint.model, messages: normalize_messages(messages), temperature: endpoint.temperature, max_tokens: endpoint.max_tokens } create_request(%{ endpoint_uuid: endpoint.uuid, endpoint_name: endpoint.name, prompt_uuid: prompt_info[:prompt_uuid], prompt_name: prompt_info[:prompt_name], model: endpoint.model, request_type: "chat", input_tokens: usage.prompt_tokens, output_tokens: usage.completion_tokens, total_tokens: usage.total_tokens, cost_cents: usage.cost_cents, latency_ms: response["latency_ms"], status: "success", metadata: %{ temperature: endpoint.temperature, max_tokens: endpoint.max_tokens, messages: normalize_messages(messages), response: response_content, request_payload: request_payload, raw_response: response, # Debug context (source tracking) source: source, stacktrace: stacktrace, caller_context: caller_context } }) end defp log_failed_request( endpoint, messages, reason, source, stacktrace, caller_context, prompt_info ) do create_request(%{ endpoint_uuid: endpoint.uuid, endpoint_name: endpoint.name, prompt_uuid: prompt_info[:prompt_uuid], prompt_name: prompt_info[:prompt_name], model: endpoint.model, request_type: "chat", status: "error", error_message: reason, metadata: %{ messages: normalize_messages(messages), # Debug context (source tracking) source: source, stacktrace: stacktrace, caller_context: caller_context } }) end # Normalize messages to ensure consistent format for storage defp normalize_messages(messages) do Enum.map(messages, fn msg -> %{ "role" => to_string(msg[:role] || msg["role"]), "content" => msg[:content] || msg["content"] } end) end defp log_embedding_request(endpoint, input, response, source, stacktrace, caller_context) do usage = Completion.extract_usage(response) input_count = if is_list(input), do: length(input), else: 1 create_request(%{ endpoint_uuid: endpoint.uuid, endpoint_name: endpoint.name, model: endpoint.model, request_type: "embedding", input_tokens: usage.prompt_tokens, total_tokens: usage.total_tokens, cost_cents: usage.cost_cents, latency_ms: response["latency_ms"], status: "success", metadata: %{ input_count: input_count, dimensions: endpoint.dimensions, # Debug context (source tracking) source: source, stacktrace: stacktrace, caller_context: caller_context } }) end defp log_failed_embedding_request(endpoint, reason, source, stacktrace, caller_context) do create_request(%{ endpoint_uuid: endpoint.uuid, endpoint_name: endpoint.name, model: endpoint.model, request_type: "embedding", status: "error", error_message: reason, metadata: %{ # Debug context (source tracking) source: source, stacktrace: stacktrace, caller_context: caller_context } }) end end