towerops/lib/towerops/admin.ex
graham ce7a8d196f fix: replace Repo.rollback with proper error handling for Oban compatibility (#201)
- Replace manual Repo.rollback/1 calls with raise/return patterns
- Fix typo: Repo.transact -> Repo.transaction
- Add unwrap_transaction_result/1 helper to reduce nesting depth
- Extract helper functions to satisfy Credo nesting depth requirements
- Fixes 'operation :rollback is rolling back unexpectedly' error

When Oban Pro's Smart engine uses Ecto.Multi for nested transactions,
manual Repo.rollback/1 calls break the transaction stack. This fix uses
proper error handling patterns:

1. Raise exceptions to trigger automatic rollback (admin.ex)
2. Return error tuples and unwrap with helper function (accounts.ex, mobile_sessions.ex)

Files changed:
- lib/towerops/admin.ex (3 functions: delete_user, delete_organization, update_billing_overrides)
- lib/towerops/accounts.ex (5 functions with helpers: update_user_email, reset_user_password, update_user_and_delete_all_tokens, do_update_user_and_delete_tokens, delete_account)
- lib/towerops/mobile_sessions.ex (2 functions with helper: complete_qr_login, do_complete_qr_login)

All tests passing (277 tests total).

Reviewed-on: graham/towerops-web#201
2026-03-28 09:11:21 -05:00

538 lines
15 KiB
Elixir

defmodule Towerops.Admin do
@moduledoc """
The Admin context for superuser operations.
Provides functions for managing users, organizations, and audit logging.
"""
import Ecto.Query
alias Towerops.Accounts.User
alias Towerops.Admin.AuditLog
alias Towerops.Organizations.Organization
alias Towerops.Repo
require Logger
## User Management
@doc """
Lists all users in the system with their organizations and device counts.
## Options
* `:limit` - Maximum number of users to return (default: 100)
* `:offset` - Number of users to skip (default: 0)
## Examples
iex> list_all_users()
[%User{}, ...]
iex> list_all_users(limit: 50, offset: 100)
[%User{}, ...]
"""
@spec list_all_users(Keyword.t()) :: [User.t()]
def list_all_users(opts \\ []) do
limit = Keyword.get(opts, :limit, 100)
offset = Keyword.get(opts, :offset, 0)
# Subquery to count devices across all user's organizations
device_count_query =
from d in "devices",
join: m in "organization_memberships",
on: d.organization_id == m.organization_id,
where: m.user_id == parent_as(:user).id,
select: count(d.id)
User
|> from(as: :user)
|> select_merge([u], %{device_count: subquery(device_count_query)})
|> order_by([u], desc: u.inserted_at)
|> limit(^limit)
|> offset(^offset)
|> preload(:organizations)
|> Repo.all()
end
@doc """
Returns the total count of users in the system.
## Examples
iex> count_users()
42
"""
@spec count_users() :: integer()
def count_users do
Repo.aggregate(User, :count)
end
@doc """
Deletes a user and creates an audit log entry.
## Examples
iex> delete_user(user_id, superuser_id, "127.0.0.1")
{:ok, %User{}}
iex> delete_user(invalid_id, superuser_id, "127.0.0.1")
** (Ecto.NoResultsError)
"""
@spec delete_user(String.t(), String.t(), String.t()) :: {:ok, User.t()} | {:error, any()}
def delete_user(user_id, superuser_id, ip_address) do
with {:ok, user} <- fetch_user(user_id) do
do_delete_user(user, superuser_id, ip_address)
end
end
defp fetch_user(user_id) do
case Repo.get(User, user_id) do
nil -> {:error, :not_found}
user -> {:ok, user}
end
end
defp do_delete_user(user, superuser_id, ip_address) do
Repo.transaction(fn ->
create_audit_log(%{
action: "user_delete",
superuser_id: superuser_id,
target_user_id: user.id,
metadata: %{email: user.email},
ip_address: ip_address
})
Towerops.Accounts.anonymize_user_login_history(user.id)
Towerops.Accounts.anonymize_user_browser_sessions(user.id)
case Repo.delete(user) do
{:ok, deleted_user} ->
deleted_user
{:error, changeset} ->
Logger.error("Failed to delete user: #{inspect(changeset)}")
# Raise to trigger rollback - never call Repo.rollback/1 directly
raise Ecto.InvalidChangesetError, action: :delete, changeset: changeset
end
end)
end
## Organization Management
@doc """
Lists all organizations in the system with their memberships.
## Options
* `:limit` - Maximum number of organizations to return (default: 100)
* `:offset` - Number of organizations to skip (default: 0)
## Examples
iex> list_all_organizations()
[%Organization{}, ...]
iex> list_all_organizations(limit: 50, offset: 100)
[%Organization{}, ...]
"""
def list_all_organizations(opts \\ []) do
limit = Keyword.get(opts, :limit, 100)
offset = Keyword.get(opts, :offset, 0)
device_count_query =
from(d in "devices",
where: d.organization_id == parent_as(:org).id,
select: count(d.id)
)
Organization
|> from(as: :org)
|> select_merge([o], %{device_count: subquery(device_count_query)})
|> order_by([o], desc: o.inserted_at)
|> limit(^limit)
|> offset(^offset)
|> preload(memberships: :user)
|> Repo.all()
end
@doc """
Returns the total count of organizations in the system.
## Examples
iex> count_organizations()
15
"""
def count_organizations do
Repo.aggregate(Organization, :count)
end
@doc """
Deletes an organization and creates an audit log entry.
## Examples
iex> delete_organization(org_id, superuser_id, "127.0.0.1")
{:ok, %Organization{}}
iex> delete_organization(invalid_id, superuser_id, "127.0.0.1")
** (Ecto.NoResultsError)
"""
def delete_organization(org_id, superuser_id, ip_address) do
with {:ok, org} <- fetch_organization(org_id) do
do_delete_organization(org, superuser_id, ip_address)
end
end
defp fetch_organization(org_id) do
case Repo.get(Organization, org_id) do
nil -> {:error, :not_found}
org -> {:ok, org}
end
end
defp do_delete_organization(org, superuser_id, ip_address) do
Repo.transaction(fn ->
create_audit_log(%{
action: "org_delete",
superuser_id: superuser_id,
target_user_id: nil,
metadata: %{name: org.name, slug: org.slug},
ip_address: ip_address
})
case Repo.delete(org) do
{:ok, deleted_org} ->
deleted_org
{:error, changeset} ->
Logger.error("Failed to delete organization: #{inspect(changeset)}")
# Raise to trigger rollback - never call Repo.rollback/1 directly
raise Ecto.InvalidChangesetError, action: :delete, changeset: changeset
end
end)
end
@doc """
Updates billing overrides for an organization with audit logging.
Only superusers should call this function. Creates an audit log entry
recording the change.
"""
def update_billing_overrides(org_id, attrs, superuser_id, ip_address) do
with {:ok, org} <- fetch_organization(org_id),
changeset = Organization.billing_override_changeset(org, attrs),
true <- changeset.valid? do
do_update_billing_overrides(changeset, org, superuser_id, ip_address)
else
{:error, _} = error -> error
false -> {:error, Organization.billing_override_changeset(%Organization{}, attrs)}
end
end
defp do_update_billing_overrides(changeset, org, superuser_id, ip_address) do
Repo.transaction(fn ->
create_audit_log(%{
action: "org_billing_override_updated",
superuser_id: superuser_id,
metadata: %{
organization_id: org.id,
organization_name: org.name,
changes: changeset.changes
},
ip_address: ip_address
})
case Repo.update(changeset) do
{:ok, updated_org} ->
updated_org
{:error, changeset} ->
Logger.error("Failed to update billing overrides: #{inspect(changeset.errors)}")
# Raise to trigger rollback - never call Repo.rollback/1 directly
raise Ecto.InvalidChangesetError, action: :update, changeset: changeset
end
end)
end
## Audit Logging
@doc """
Creates an audit log entry.
## Examples
iex> create_audit_log(%{action: "impersonate_start", superuser_id: superuser_id, target_user_id: user_id})
{:ok, %AuditLog{}}
iex> create_audit_log(%{action: "invalid"})
{:error, %Ecto.Changeset{}}
"""
def create_audit_log(attrs) do
%AuditLog{}
|> AuditLog.changeset(attrs)
|> Repo.insert()
end
@doc """
Lists recent audit logs with associated users.
## Options
* `:limit` - Maximum number of logs to return (default: 100)
## Examples
iex> list_audit_logs()
[%AuditLog{}, ...]
iex> list_audit_logs(limit: 50)
[%AuditLog{}, ...]
"""
def list_audit_logs(opts \\ []) do
limit = Keyword.get(opts, :limit, 100)
AuditLog
|> order_by([a], desc: a.inserted_at)
|> limit(^limit)
|> preload([:superuser, :target_user])
|> Repo.all()
end
@doc """
Returns audit logs related to a specific user (GDPR data access).
"""
def list_audit_logs_for_user(user_id, opts \\ []) do
limit = Keyword.get(opts, :limit, 100)
AuditLog
|> where([a], a.target_user_id == ^user_id or a.superuser_id == ^user_id)
|> order_by([a], desc: a.inserted_at)
|> limit(^limit)
|> preload([:superuser, :target_user])
|> Repo.all()
end
@doc """
Lists audit logs with optional filtering.
## Options
* `:limit` - Maximum number of logs to return (default: 100)
* `:offset` - Number of logs to skip (default: 0)
* `:email` - Filter by actor or target user email (partial match)
* `:action` - Filter by specific action
* `:date_from` - Filter logs from this date (ISO 8601 format)
* `:date_to` - Filter logs up to this date (ISO 8601 format)
## Examples
iex> list_audit_logs_with_filters(limit: 50, email: "user@example.com")
[%AuditLog{}, ...]
iex> list_audit_logs_with_filters(action: "impersonate_start", date_from: "2024-01-01")
[%AuditLog{}, ...]
"""
def list_audit_logs_with_filters(opts \\ []) do
limit = Keyword.get(opts, :limit, 100)
offset = Keyword.get(opts, :offset, 0)
email = Keyword.get(opts, :email, "")
action = Keyword.get(opts, :action, "")
date_from = Keyword.get(opts, :date_from, "")
date_to = Keyword.get(opts, :date_to, "")
query =
AuditLog
|> preload([:superuser, :target_user])
|> order_by([a], desc: a.inserted_at)
query = filter_by_email(query, email)
query = filter_by_action(query, action)
query = filter_by_date_range(query, date_from, date_to)
if limit == :all do
Repo.all(query)
else
query
|> limit(^limit)
|> offset(^offset)
|> Repo.all()
end
end
defp filter_by_email(query, ""), do: query
defp filter_by_email(query, email) do
from(a in query,
left_join: su in assoc(a, :superuser),
left_join: tu in assoc(a, :target_user),
where: ilike(su.email, ^"%#{email}%") or ilike(tu.email, ^"%#{email}%")
)
end
defp filter_by_action(query, ""), do: query
defp filter_by_action(query, action), do: where(query, [a], a.action == ^action)
defp filter_by_date_range(query, "", ""), do: query
defp filter_by_date_range(query, from_date, to_date) do
query =
if from_date == "" do
query
else
{:ok, from_datetime} = NaiveDateTime.new(Date.from_iso8601!(from_date), ~T[00:00:00])
where(query, [a], a.inserted_at >= ^from_datetime)
end
if to_date == "" do
query
else
{:ok, to_datetime} = NaiveDateTime.new(Date.from_iso8601!(to_date), ~T[23:59:59])
where(query, [a], a.inserted_at <= ^to_datetime)
end
end
## Global Pricing Management
@doc """
Update global pricing settings.
Creates new Stripe Price if price changed, migrates all active subscriptions,
updates ApplicationSettings, and creates audit log.
Returns {:ok, migration_report} or {:error, reason}.
"""
def update_global_pricing(attrs, superuser_id, ip_address) do
alias Towerops.Billing
alias Towerops.Billing.StripeClient
alias Towerops.Settings
alias Towerops.Settings.ApplicationSetting
with {:ok, validated} <- validate_pricing_attrs(attrs),
old_values = get_current_pricing(),
{:ok, new_price_id, migration_report} <- maybe_update_stripe_price(validated, old_values),
:ok <- update_pricing_settings(validated, new_price_id),
:ok <- create_pricing_audit_log(old_values, validated, migration_report, superuser_id, ip_address) do
{:ok, migration_report}
end
end
defp validate_pricing_attrs(attrs) do
free_devices = Map.get(attrs, "default_free_devices")
price = Map.get(attrs, "default_price_per_device")
with {:ok, free} <- validate_free_devices(free_devices),
{:ok, price} <- validate_price(price) do
{:ok, %{free_devices: free, price: price}}
end
end
defp validate_free_devices(nil), do: {:ok, nil}
defp validate_free_devices(value) when is_binary(value) do
case Integer.parse(value) do
{int, _} when int > 0 and int < 10_000 -> {:ok, int}
_ -> {:error, :invalid_free_devices}
end
end
defp validate_price(nil), do: {:ok, nil}
defp validate_price(value) when is_binary(value) do
case Decimal.parse(value) do
{decimal, _} ->
if Decimal.compare(decimal, "0.01") in [:gt, :eq] and
Decimal.compare(decimal, "999.99") in [:lt, :eq] do
{:ok, value}
else
{:error, :invalid_price}
end
:error ->
{:error, :invalid_price}
end
end
defp get_current_pricing do
alias Towerops.Settings
%{
free_devices: Settings.get_setting("default_free_devices"),
price: Settings.get_setting("default_price_per_device")
}
end
defp maybe_update_stripe_price(validated, old_values) do
alias Towerops.Billing
alias Towerops.Billing.StripeClient
price_changed = validated.price && validated.price != old_values.price
if price_changed do
case StripeClient.create_price(validated.price) do
{:ok, %{"id" => price_id}} ->
migration_report = Billing.migrate_all_subscriptions_to_price(price_id)
{:ok, price_id, migration_report}
{:error, reason} ->
{:error, {:stripe_error, reason}}
end
else
{:ok, nil, %{total: 0, succeeded: 0, failed: 0, failures: []}}
end
end
defp update_pricing_settings(validated, new_price_id) do
alias Towerops.Settings
alias Towerops.Settings.ApplicationSetting
if validated.free_devices do
Settings.update_setting("default_free_devices", validated.free_devices)
end
if validated.price do
Settings.update_setting("default_price_per_device", validated.price)
end
if new_price_id do
# Create setting if doesn't exist
case Settings.get_setting_record("stripe_price_id") do
nil ->
%ApplicationSetting{}
|> ApplicationSetting.changeset(%{
key: "stripe_price_id",
value: new_price_id,
value_type: "string",
description: "Current Stripe Price ID for metered billing"
})
|> Repo.insert()
_setting ->
Settings.update_setting("stripe_price_id", new_price_id)
end
end
:ok
end
defp create_pricing_audit_log(old_values, validated, migration_report, superuser_id, ip_address) do
metadata = %{
old_free_devices: old_values.free_devices,
new_free_devices: validated.free_devices,
old_price: old_values.price,
new_price: validated.price,
migration_total: migration_report.total,
migration_succeeded: migration_report.succeeded,
migration_failed: migration_report.failed
}
create_audit_log(%{
action: "global_pricing_updated",
superuser_id: superuser_id,
ip_address: ip_address,
metadata: metadata
})
:ok
end
end