fix/security-and-quality-issues (#136)
Reviewed-on: graham/towerops-web#136
This commit is contained in:
parent
872f717dba
commit
12a737e1f5
13 changed files with 416 additions and 190 deletions
|
|
@ -655,7 +655,7 @@ defmodule SnmpKit.SnmpLib.PDU.V3Encoder do
|
|||
Creates a discovery message for engine ID discovery.
|
||||
"""
|
||||
@spec create_discovery_message(non_neg_integer()) :: v3_message()
|
||||
def create_discovery_message(msg_id \\ :rand.uniform(2_147_483_647)) do
|
||||
def create_discovery_message(msg_id \\ 4 |> :crypto.strong_rand_bytes() |> :binary.decode_unsigned()) do
|
||||
%{
|
||||
version: 3,
|
||||
msg_id: msg_id,
|
||||
|
|
|
|||
|
|
@ -136,8 +136,8 @@ defmodule SnmpKit.SnmpLib.Security.USM do
|
|||
Logger.debug("Starting engine discovery for host: #{host}")
|
||||
|
||||
try do
|
||||
# Create discovery message
|
||||
msg_id = :rand.uniform(2_147_483_647)
|
||||
# Create discovery message with cryptographically secure message ID
|
||||
msg_id = 4 |> :crypto.strong_rand_bytes() |> :binary.decode_unsigned()
|
||||
discovery_message = V3Encoder.create_discovery_message(msg_id)
|
||||
|
||||
# Encode discovery message (no security)
|
||||
|
|
@ -186,7 +186,8 @@ defmodule SnmpKit.SnmpLib.Security.USM do
|
|||
|
||||
try do
|
||||
# Create time synchronization message (authenticated but not encrypted)
|
||||
msg_id = :rand.uniform(2_147_483_647)
|
||||
# Use cryptographically secure random for message ID
|
||||
msg_id = 4 |> :crypto.strong_rand_bytes() |> :binary.decode_unsigned()
|
||||
|
||||
time_sync_message = %{
|
||||
version: 3,
|
||||
|
|
@ -388,7 +389,7 @@ defmodule SnmpKit.SnmpLib.Security.USM do
|
|||
|
||||
%{
|
||||
version: 3,
|
||||
msg_id: :rand.uniform(2_147_483_647),
|
||||
msg_id: 4 |> :crypto.strong_rand_bytes() |> :binary.decode_unsigned(),
|
||||
msg_max_size: Constants.default_max_message_size(),
|
||||
msg_flags: flags,
|
||||
msg_security_model: Constants.usm_security_model(),
|
||||
|
|
|
|||
|
|
@ -18,6 +18,8 @@ defmodule Towerops.Accounts do
|
|||
alias Towerops.GeoIP
|
||||
alias Towerops.Repo
|
||||
|
||||
require Logger
|
||||
|
||||
## Database getters
|
||||
|
||||
@doc """
|
||||
|
|
@ -844,8 +846,17 @@ defmodule Towerops.Accounts do
|
|||
def deliver_user_reset_password_instructions(%User{} = user, reset_password_url_fun)
|
||||
when is_function(reset_password_url_fun, 1) do
|
||||
{encoded_token, user_token} = UserToken.build_email_token(user, "reset_password")
|
||||
Repo.insert!(user_token)
|
||||
UserNotifier.deliver_reset_password_instructions(user, reset_password_url_fun.(encoded_token))
|
||||
|
||||
case Repo.insert(user_token) do
|
||||
{:ok, _token} ->
|
||||
UserNotifier.deliver_reset_password_instructions(user, reset_password_url_fun.(encoded_token))
|
||||
|
||||
{:error, changeset} ->
|
||||
require Logger
|
||||
|
||||
Logger.error("Failed to insert reset password token: #{inspect(changeset.errors)}")
|
||||
{:error, changeset}
|
||||
end
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -900,8 +911,17 @@ defmodule Towerops.Accounts do
|
|||
"""
|
||||
def generate_user_session_token(user) do
|
||||
{token, user_token} = UserToken.build_session_token(user)
|
||||
Repo.insert!(user_token)
|
||||
token
|
||||
|
||||
case Repo.insert(user_token) do
|
||||
{:ok, _} ->
|
||||
token
|
||||
|
||||
{:error, changeset} ->
|
||||
require Logger
|
||||
|
||||
Logger.error("Failed to insert session token: #{inspect(changeset.errors)}")
|
||||
raise Ecto.InvalidChangesetError, changeset: changeset
|
||||
end
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -911,8 +931,17 @@ defmodule Towerops.Accounts do
|
|||
"""
|
||||
def generate_user_session_token_with_record(user) do
|
||||
{token, user_token_changeset} = UserToken.build_session_token(user)
|
||||
user_token = Repo.insert!(user_token_changeset)
|
||||
{token, user_token}
|
||||
|
||||
case Repo.insert(user_token_changeset) do
|
||||
{:ok, user_token} ->
|
||||
{token, user_token}
|
||||
|
||||
{:error, changeset} ->
|
||||
require Logger
|
||||
|
||||
Logger.error("Failed to insert session token with record: #{inspect(changeset.errors)}")
|
||||
raise Ecto.InvalidChangesetError, changeset: changeset
|
||||
end
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -967,34 +996,39 @@ defmodule Towerops.Accounts do
|
|||
`mix help phx.gen.auth`.
|
||||
"""
|
||||
def login_user_by_magic_link(token) do
|
||||
case UserToken.verify_magic_link_token_query(token) do
|
||||
{:ok, query} ->
|
||||
case Repo.one(query) do
|
||||
# Prevent session fixation attacks by disallowing magic links for unconfirmed users with password
|
||||
{%User{confirmed_at: nil, hashed_password: hash}, _token} when not is_nil(hash) ->
|
||||
raise """
|
||||
magic link log in is not allowed for unconfirmed users with a password set!
|
||||
with {:ok, query} <- UserToken.verify_magic_link_token_query(token),
|
||||
result when not is_nil(result) <- Repo.one(query) do
|
||||
handle_magic_link_login(result)
|
||||
else
|
||||
nil -> {:error, :not_found}
|
||||
:error -> {:error, :not_found}
|
||||
end
|
||||
end
|
||||
|
||||
This cannot happen with the default implementation, which indicates that you
|
||||
might have adapted the code to a different use case. Please make sure to read the
|
||||
"Mixing magic link and password registration" section of `mix help phx.gen.auth`.
|
||||
"""
|
||||
defp handle_magic_link_login({%User{confirmed_at: nil, hashed_password: hash}, _token}) when not is_nil(hash) do
|
||||
raise """
|
||||
magic link log in is not allowed for unconfirmed users with a password set!
|
||||
|
||||
{%User{confirmed_at: nil} = user, _token} ->
|
||||
user
|
||||
|> User.confirm_changeset()
|
||||
|> update_user_and_delete_all_tokens()
|
||||
This cannot happen with the default implementation, which indicates that you
|
||||
might have adapted the code to a different use case. Please make sure to read the
|
||||
"Mixing magic link and password registration" section of `mix help phx.gen.auth`.
|
||||
"""
|
||||
end
|
||||
|
||||
{user, token} ->
|
||||
Repo.delete!(token)
|
||||
{:ok, {user, []}}
|
||||
defp handle_magic_link_login({%User{confirmed_at: nil} = user, _token}) do
|
||||
user
|
||||
|> User.confirm_changeset()
|
||||
|> update_user_and_delete_all_tokens()
|
||||
end
|
||||
|
||||
nil ->
|
||||
{:error, :not_found}
|
||||
end
|
||||
defp handle_magic_link_login({user, token}) do
|
||||
case Repo.delete(token) do
|
||||
{:ok, _} ->
|
||||
{:ok, {user, []}}
|
||||
|
||||
:error ->
|
||||
{:error, :not_found}
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to delete magic link token: #{inspect(changeset)}")
|
||||
{:ok, {user, []}}
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -1011,8 +1045,16 @@ defmodule Towerops.Accounts do
|
|||
when is_function(update_email_url_fun, 1) do
|
||||
{encoded_token, user_token} = UserToken.build_email_token(user, "change:#{current_email}")
|
||||
|
||||
Repo.insert!(user_token)
|
||||
UserNotifier.deliver_update_email_instructions(user, update_email_url_fun.(encoded_token))
|
||||
case Repo.insert(user_token) do
|
||||
{:ok, _} ->
|
||||
UserNotifier.deliver_update_email_instructions(user, update_email_url_fun.(encoded_token))
|
||||
|
||||
{:error, changeset} ->
|
||||
require Logger
|
||||
|
||||
Logger.error("Failed to insert update email token: #{inspect(changeset.errors)}")
|
||||
{:error, changeset}
|
||||
end
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -1020,8 +1062,17 @@ defmodule Towerops.Accounts do
|
|||
"""
|
||||
def deliver_login_instructions(%User{} = user, magic_link_url_fun) when is_function(magic_link_url_fun, 1) do
|
||||
{encoded_token, user_token} = UserToken.build_email_token(user, "login")
|
||||
Repo.insert!(user_token)
|
||||
UserNotifier.deliver_login_instructions(user, magic_link_url_fun.(encoded_token))
|
||||
|
||||
case Repo.insert(user_token) do
|
||||
{:ok, _} ->
|
||||
UserNotifier.deliver_login_instructions(user, magic_link_url_fun.(encoded_token))
|
||||
|
||||
{:error, changeset} ->
|
||||
require Logger
|
||||
|
||||
Logger.error("Failed to insert login token: #{inspect(changeset.errors)}")
|
||||
{:error, changeset}
|
||||
end
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -1666,8 +1717,17 @@ defmodule Towerops.Accounts do
|
|||
"""
|
||||
def deliver_user_confirmation_instructions(%User{} = user, confirmation_url_fun) do
|
||||
{encoded_token, user_token} = UserToken.build_email_token(user, "confirm")
|
||||
Repo.insert!(user_token)
|
||||
UserNotifier.deliver_confirmation_email(user, confirmation_url_fun.(encoded_token))
|
||||
|
||||
case Repo.insert(user_token) do
|
||||
{:ok, _} ->
|
||||
UserNotifier.deliver_confirmation_email(user, confirmation_url_fun.(encoded_token))
|
||||
|
||||
{:error, changeset} ->
|
||||
require Logger
|
||||
|
||||
Logger.error("Failed to insert confirmation token: #{inspect(changeset.errors)}")
|
||||
{:error, changeset}
|
||||
end
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -1754,25 +1814,37 @@ defmodule Towerops.Accounts do
|
|||
end
|
||||
|
||||
defp do_delete_account(user) do
|
||||
case sole_owner_organizations(user.id) do
|
||||
[] ->
|
||||
Repo.transaction(fn ->
|
||||
{:ok, _} =
|
||||
Towerops.Admin.create_audit_log(%{
|
||||
action: "account_deletion_requested",
|
||||
superuser_id: user.id,
|
||||
target_user_id: user.id,
|
||||
metadata: %{email: user.email, self_service: true}
|
||||
})
|
||||
sole_owner_orgs = sole_owner_organizations(user.id)
|
||||
|
||||
anonymize_user_login_history(user.id)
|
||||
anonymize_user_browser_sessions(user.id)
|
||||
Repo.delete!(user)
|
||||
end)
|
||||
|
||||
sole_owner_orgs ->
|
||||
org_names = Enum.map(sole_owner_orgs, & &1.name)
|
||||
{:error, :sole_owner, org_names}
|
||||
if sole_owner_orgs == [] do
|
||||
delete_account_transaction(user)
|
||||
else
|
||||
org_names = Enum.map(sole_owner_orgs, & &1.name)
|
||||
{:error, :sole_owner, org_names}
|
||||
end
|
||||
end
|
||||
|
||||
defp delete_account_transaction(user) do
|
||||
Repo.transaction(fn ->
|
||||
{:ok, _} =
|
||||
Towerops.Admin.create_audit_log(%{
|
||||
action: "account_deletion_requested",
|
||||
superuser_id: user.id,
|
||||
target_user_id: user.id,
|
||||
metadata: %{email: user.email, self_service: true}
|
||||
})
|
||||
|
||||
anonymize_user_login_history(user.id)
|
||||
anonymize_user_browser_sessions(user.id)
|
||||
|
||||
case Repo.delete(user) do
|
||||
{:ok, deleted_user} ->
|
||||
deleted_user
|
||||
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to delete user account: #{inspect(changeset)}")
|
||||
Repo.rollback(changeset)
|
||||
end
|
||||
end)
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -11,6 +11,8 @@ defmodule Towerops.Admin do
|
|||
alias Towerops.Organizations.Organization
|
||||
alias Towerops.Repo
|
||||
|
||||
require Logger
|
||||
|
||||
## User Management
|
||||
|
||||
@doc """
|
||||
|
|
@ -78,27 +80,39 @@ defmodule Towerops.Admin do
|
|||
"""
|
||||
@spec delete_user(String.t(), String.t(), String.t()) :: {:ok, User.t()} | {:error, any()}
|
||||
def delete_user(user_id, superuser_id, ip_address) do
|
||||
user = Repo.get!(User, user_id)
|
||||
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 before deletion
|
||||
create_audit_log(%{
|
||||
action: "user_delete",
|
||||
superuser_id: superuser_id,
|
||||
target_user_id: user_id,
|
||||
target_user_id: user.id,
|
||||
metadata: %{email: user.email},
|
||||
ip_address: ip_address
|
||||
})
|
||||
|
||||
# Anonymize login history (GDPR compliance)
|
||||
# Preserves IP/location for security analysis, removes user link
|
||||
Towerops.Accounts.anonymize_user_login_history(user_id)
|
||||
Towerops.Accounts.anonymize_user_login_history(user.id)
|
||||
Towerops.Accounts.anonymize_user_browser_sessions(user.id)
|
||||
|
||||
# Anonymize browser sessions (GDPR compliance)
|
||||
Towerops.Accounts.anonymize_user_browser_sessions(user_id)
|
||||
case Repo.delete(user) do
|
||||
{:ok, deleted_user} ->
|
||||
deleted_user
|
||||
|
||||
# Delete user (cascades to memberships and tokens)
|
||||
Repo.delete!(user)
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to delete user: #{inspect(changeset)}")
|
||||
Repo.rollback(changeset)
|
||||
end
|
||||
end)
|
||||
end
|
||||
|
||||
|
|
@ -164,10 +178,20 @@ defmodule Towerops.Admin do
|
|||
** (Ecto.NoResultsError)
|
||||
"""
|
||||
def delete_organization(org_id, superuser_id, ip_address) do
|
||||
org = Repo.get!(Organization, org_id)
|
||||
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 before deletion
|
||||
create_audit_log(%{
|
||||
action: "org_delete",
|
||||
superuser_id: superuser_id,
|
||||
|
|
@ -176,8 +200,14 @@ defmodule Towerops.Admin do
|
|||
ip_address: ip_address
|
||||
})
|
||||
|
||||
# Delete organization (cascades to memberships, sites, device)
|
||||
Repo.delete!(org)
|
||||
case Repo.delete(org) do
|
||||
{:ok, deleted_org} ->
|
||||
deleted_org
|
||||
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to delete organization: #{inspect(changeset)}")
|
||||
Repo.rollback(changeset)
|
||||
end
|
||||
end)
|
||||
end
|
||||
|
||||
|
|
@ -188,30 +218,40 @@ defmodule Towerops.Admin do
|
|||
recording the change.
|
||||
"""
|
||||
def update_billing_overrides(org_id, attrs, superuser_id, ip_address) do
|
||||
org = Repo.get!(Organization, org_id)
|
||||
|
||||
changeset = Organization.billing_override_changeset(org, attrs)
|
||||
|
||||
if changeset.valid? 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
|
||||
})
|
||||
|
||||
Repo.update!(changeset)
|
||||
end)
|
||||
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, changeset}
|
||||
{: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)}")
|
||||
Repo.rollback(changeset)
|
||||
end
|
||||
end)
|
||||
end
|
||||
|
||||
## Audit Logging
|
||||
|
||||
@doc """
|
||||
|
|
|
|||
|
|
@ -66,7 +66,7 @@ defmodule Towerops.Alerts.StormDetector do
|
|||
"""
|
||||
@spec storm_mode?() :: boolean()
|
||||
def storm_mode? do
|
||||
GenServer.call(__MODULE__, :storm_mode?)
|
||||
GenServer.call(__MODULE__, :storm_mode?, 5000)
|
||||
end
|
||||
|
||||
@doc """
|
||||
|
|
@ -74,7 +74,7 @@ defmodule Towerops.Alerts.StormDetector do
|
|||
"""
|
||||
@spec stats() :: map()
|
||||
def stats do
|
||||
GenServer.call(__MODULE__, :stats)
|
||||
GenServer.call(__MODULE__, :stats, 5000)
|
||||
end
|
||||
|
||||
# --- GenServer Implementation ---
|
||||
|
|
|
|||
|
|
@ -59,34 +59,31 @@ defmodule Towerops.Monitoring.Executors.PingExecutor do
|
|||
# The ping command itself has -W/-w timeout flags.
|
||||
# Add a safety timeout in case the process hangs.
|
||||
safety_timeout = timeout_ms + count * 1000 + 2000
|
||||
caller = self()
|
||||
ref = make_ref()
|
||||
|
||||
pid =
|
||||
spawn(fn ->
|
||||
task =
|
||||
Task.async(fn ->
|
||||
try do
|
||||
result = System.cmd("ping", args, stderr_to_stdout: true)
|
||||
send(caller, {ref, {:ok, result}})
|
||||
System.cmd("ping", args, stderr_to_stdout: true)
|
||||
rescue
|
||||
e -> send(caller, {ref, {:error, Exception.message(e)}})
|
||||
e -> {:error, Exception.message(e)}
|
||||
end
|
||||
end)
|
||||
|
||||
receive do
|
||||
{^ref, {:ok, {output, 0}}} ->
|
||||
case Task.yield(task, safety_timeout) do
|
||||
{:ok, {:error, reason}} ->
|
||||
{:error, reason}
|
||||
|
||||
{:ok, {output, 0}} ->
|
||||
parse_output(output)
|
||||
|
||||
{^ref, {:ok, {output, _exit_code}}} ->
|
||||
{:ok, {output, _exit_code}} ->
|
||||
case parse_output(output) do
|
||||
{:ok, _, _} = success -> success
|
||||
{:error, _} = error -> error
|
||||
end
|
||||
|
||||
{^ref, {:error, reason}} ->
|
||||
{:error, reason}
|
||||
after
|
||||
safety_timeout ->
|
||||
Process.exit(pid, :kill)
|
||||
nil ->
|
||||
Task.shutdown(task, :brutal_kill)
|
||||
{:error, "Ping command timed out"}
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -67,7 +67,7 @@ defmodule Towerops.Profiles.YamlProfiles do
|
|||
Reloads all profiles from YAML files.
|
||||
"""
|
||||
def reload do
|
||||
GenServer.call(__MODULE__, :reload)
|
||||
GenServer.call(__MODULE__, :reload, 30_000)
|
||||
end
|
||||
|
||||
# GenServer callbacks
|
||||
|
|
|
|||
|
|
@ -854,25 +854,30 @@ defmodule Towerops.Snmp.Discovery do
|
|||
Repo.transaction(
|
||||
fn ->
|
||||
# Upsert Device
|
||||
snmp_device = upsert_device(device, sanitized_device_info)
|
||||
case upsert_device(device, sanitized_device_info) do
|
||||
{:ok, snmp_device} ->
|
||||
# Sync interfaces, sensors, and VLANs (preserving historical data)
|
||||
synced_interfaces = sync_interfaces(snmp_device, interfaces)
|
||||
_ = sync_sensors(snmp_device, sensors)
|
||||
_ = sync_vlans(snmp_device, vlans)
|
||||
|
||||
# Sync interfaces, sensors, and VLANs (preserving historical data)
|
||||
synced_interfaces = sync_interfaces(snmp_device, interfaces)
|
||||
_ = sync_sensors(snmp_device, sensors)
|
||||
_ = sync_vlans(snmp_device, vlans)
|
||||
# Sync IP addresses and processors inside transaction for atomicity
|
||||
device_with_interfaces = %{
|
||||
snmp_device
|
||||
| interfaces: Enum.map(synced_interfaces, &Map.put(&1, :device_id, snmp_device.device_id))
|
||||
}
|
||||
|
||||
# Sync IP addresses and processors inside transaction for atomicity
|
||||
device_with_interfaces = %{
|
||||
snmp_device
|
||||
| interfaces: Enum.map(synced_interfaces, &Map.put(&1, :device_id, snmp_device.device_id))
|
||||
}
|
||||
_ = sync_ip_addresses(device_with_interfaces, ip_addresses)
|
||||
_ = sync_processors(device_with_interfaces, processors)
|
||||
|
||||
_ = sync_ip_addresses(device_with_interfaces, ip_addresses)
|
||||
_ = sync_processors(device_with_interfaces, processors)
|
||||
# Return device with interfaces loaded and device_id added
|
||||
# Use snmp_device.device_id which references the Equipment table
|
||||
device_with_interfaces
|
||||
|
||||
# Return device with interfaces loaded and device_id added
|
||||
# Use snmp_device.device_id which references the Equipment table
|
||||
device_with_interfaces
|
||||
{:error, changeset} ->
|
||||
Logger.error("Failed to upsert SNMP device: #{inspect(changeset.errors)}")
|
||||
Repo.rollback({:upsert_failed, changeset})
|
||||
end
|
||||
end,
|
||||
timeout: 60_000
|
||||
)
|
||||
|
|
@ -914,47 +919,39 @@ defmodule Towerops.Snmp.Discovery do
|
|||
|> Map.put(:last_discovery_at, last_discovery_at)
|
||||
end
|
||||
|
||||
@spec upsert_device(DeviceSchema.t(), map()) :: Device.t()
|
||||
@spec upsert_device(DeviceSchema.t(), map()) :: {:ok, Device.t()} | {:error, Ecto.Changeset.t()}
|
||||
defp upsert_device(device, device_attrs) do
|
||||
Logger.info("Upserting SNMP device for #{device.name}")
|
||||
|
||||
case Repo.get_by(Device, device_id: device.id) do
|
||||
nil ->
|
||||
snmp_device =
|
||||
%Device{}
|
||||
|> Device.changeset(device_attrs)
|
||||
|> Repo.insert!()
|
||||
existing_device = Repo.get_by(Device, device_id: device.id)
|
||||
old_version = existing_device && existing_device.firmware_version
|
||||
new_version = device_attrs[:firmware_version]
|
||||
|
||||
# Log initial firmware version (no old version)
|
||||
if snmp_device.firmware_version do
|
||||
Firmware.log_firmware_change(
|
||||
snmp_device.id,
|
||||
nil,
|
||||
snmp_device.firmware_version
|
||||
)
|
||||
changeset = Device.changeset(existing_device || %Device{}, device_attrs)
|
||||
|
||||
result =
|
||||
Repo.insert_or_update(changeset,
|
||||
on_conflict: {:replace_all_except, [:id, :inserted_at]},
|
||||
conflict_target: :device_id
|
||||
)
|
||||
|
||||
case result do
|
||||
{:ok, snmp_device} ->
|
||||
cond do
|
||||
is_nil(existing_device) && snmp_device.firmware_version ->
|
||||
Firmware.log_firmware_change(snmp_device.id, nil, snmp_device.firmware_version)
|
||||
|
||||
version_changed?(old_version, new_version) ->
|
||||
Firmware.log_firmware_change(snmp_device.id, old_version, new_version)
|
||||
|
||||
true ->
|
||||
:ok
|
||||
end
|
||||
|
||||
snmp_device
|
||||
{:ok, snmp_device}
|
||||
|
||||
existing_device ->
|
||||
old_version = existing_device.firmware_version
|
||||
new_version = device_attrs[:firmware_version]
|
||||
|
||||
snmp_device =
|
||||
existing_device
|
||||
|> Device.changeset(device_attrs)
|
||||
|> Repo.update!()
|
||||
|
||||
# Detect and log firmware version changes
|
||||
if version_changed?(old_version, new_version) do
|
||||
Firmware.log_firmware_change(
|
||||
snmp_device.id,
|
||||
old_version,
|
||||
new_version
|
||||
)
|
||||
end
|
||||
|
||||
snmp_device
|
||||
{:error, _changeset} = error ->
|
||||
error
|
||||
end
|
||||
end
|
||||
|
||||
|
|
|
|||
|
|
@ -4,6 +4,14 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookController do
|
|||
|
||||
Called by CI after publishing a new GitHub Release. Broadcasts an update
|
||||
command to all connected agents via their existing WebSocket channels.
|
||||
|
||||
## Authentication
|
||||
Uses Bearer token authentication via the WebhookAuth plug.
|
||||
The token must match the configured `agent_webhook_secret`.
|
||||
|
||||
## Optional Signature Verification
|
||||
For enhanced security, also supports HMAC-SHA256 signature verification
|
||||
via X-Agent-Webhook-Signature header (format: t=<timestamp>,v1=<signature>).
|
||||
"""
|
||||
|
||||
use ToweropsWeb, :controller
|
||||
|
|
@ -12,7 +20,23 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookController do
|
|||
|
||||
require Logger
|
||||
|
||||
@max_age_seconds 300
|
||||
|
||||
def create(conn, _params) do
|
||||
case verify_optional_signature(conn) do
|
||||
:ok ->
|
||||
handle_update(conn)
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.warning("Agent release webhook signature verification failed: #{inspect(reason)}")
|
||||
|
||||
conn
|
||||
|> put_status(:unauthorized)
|
||||
|> json(%{status: "error", error: "Signature verification failed"})
|
||||
end
|
||||
end
|
||||
|
||||
defp handle_update(conn) do
|
||||
case Agents.broadcast_mass_update() do
|
||||
{:ok, result} ->
|
||||
json(conn, %{
|
||||
|
|
@ -30,4 +54,78 @@ defmodule ToweropsWeb.Api.V1.AgentReleaseWebhookController do
|
|||
|> json(%{status: "error", error: "Broadcast failed"})
|
||||
end
|
||||
end
|
||||
|
||||
defp verify_optional_signature(conn) do
|
||||
sig_headers = Plug.Conn.get_req_header(conn, "x-agent-webhook-signature")
|
||||
|
||||
case sig_headers do
|
||||
[] ->
|
||||
:ok
|
||||
|
||||
[header] ->
|
||||
secret = Application.get_env(:towerops, :agent_webhook_secret)
|
||||
raw_body = conn.private[:raw_body] || ""
|
||||
verify_hmac_signature(header, raw_body, secret)
|
||||
end
|
||||
end
|
||||
|
||||
defp verify_hmac_signature(header, raw_body, secret) do
|
||||
with {:ok, timestamp, v1_signature} <- parse_signature_header(header),
|
||||
:ok <- check_timestamp(timestamp) do
|
||||
signed_payload = timestamp <> "." <> raw_body
|
||||
|
||||
expected =
|
||||
:hmac
|
||||
|> :crypto.mac(:sha256, secret, signed_payload)
|
||||
|> Base.encode16(case: :lower)
|
||||
|
||||
if Plug.Crypto.secure_compare(expected, v1_signature) do
|
||||
:ok
|
||||
else
|
||||
{:error, :invalid_signature}
|
||||
end
|
||||
end
|
||||
end
|
||||
|
||||
defp parse_signature_header(header) do
|
||||
parts = String.split(header, ",")
|
||||
|
||||
with {:ok, timestamp} <- extract_part(parts, "t="),
|
||||
{:ok, signature} <- extract_part(parts, "v1=") do
|
||||
{:ok, timestamp, signature}
|
||||
else
|
||||
_ -> {:error, :malformed_signature}
|
||||
end
|
||||
end
|
||||
|
||||
defp extract_part(parts, prefix) do
|
||||
parts
|
||||
|> Enum.find_value(fn part ->
|
||||
case String.split(part, prefix, parts: 2) do
|
||||
["", value] -> {:ok, value}
|
||||
_ -> nil
|
||||
end
|
||||
end)
|
||||
|> case do
|
||||
{:ok, value} -> {:ok, value}
|
||||
nil -> {:error, :missing_part}
|
||||
end
|
||||
end
|
||||
|
||||
defp check_timestamp(timestamp_str) do
|
||||
case Integer.parse(timestamp_str) do
|
||||
{timestamp, ""} ->
|
||||
now = System.system_time(:second)
|
||||
age = now - timestamp
|
||||
|
||||
cond do
|
||||
age < 0 -> {:error, :timestamp_in_future}
|
||||
age > @max_age_seconds -> {:error, :timestamp_expired}
|
||||
true -> :ok
|
||||
end
|
||||
|
||||
_ ->
|
||||
{:error, :invalid_timestamp}
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -196,7 +196,8 @@ defmodule ToweropsWeb.Api.V1.MibController do
|
|||
|
||||
defp extract_tarball(conn, upload, vendor_dir, vendor) do
|
||||
# Extract to temporary directory first to validate contents
|
||||
temp_dir = Path.join(System.tmp_dir!(), "mib_extract_#{:rand.uniform(999_999_999)}")
|
||||
random_suffix = 16 |> :crypto.strong_rand_bytes() |> Base.encode16(case: :lower)
|
||||
temp_dir = Path.join(System.tmp_dir!(), "mib_extract_#{random_suffix}")
|
||||
File.mkdir_p!(temp_dir)
|
||||
|
||||
try do
|
||||
|
|
@ -242,7 +243,8 @@ defmodule ToweropsWeb.Api.V1.MibController do
|
|||
|
||||
defp extract_zip(conn, upload, vendor_dir, vendor) do
|
||||
# Extract to temporary directory first to validate contents
|
||||
temp_dir = Path.join(System.tmp_dir!(), "mib_extract_#{:rand.uniform(999_999_999)}")
|
||||
random_suffix = 16 |> :crypto.strong_rand_bytes() |> Base.encode16(case: :lower)
|
||||
temp_dir = Path.join(System.tmp_dir!(), "mib_extract_#{random_suffix}")
|
||||
File.mkdir_p!(temp_dir)
|
||||
|
||||
try do
|
||||
|
|
|
|||
|
|
@ -16,6 +16,8 @@ defmodule ToweropsWeb.DeviceLive.Show do
|
|||
alias ToweropsWeb.CheckLive.FormComponent
|
||||
alias ToweropsWeb.Live.Helpers.AccessControl
|
||||
|
||||
require Logger
|
||||
|
||||
# SNMP interface type to category mappings
|
||||
@interface_type_categories %{
|
||||
6 => "Ethernet",
|
||||
|
|
@ -1053,7 +1055,7 @@ defmodule ToweropsWeb.DeviceLive.Show do
|
|||
if Enum.empty?(datasets) do
|
||||
nil
|
||||
else
|
||||
Jason.encode!(%{datasets: datasets})
|
||||
safe_encode_chart_data(%{datasets: datasets})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -1096,7 +1098,7 @@ defmodule ToweropsWeb.DeviceLive.Show do
|
|||
if Enum.empty?(datasets) do
|
||||
nil
|
||||
else
|
||||
Jason.encode!(%{datasets: datasets})
|
||||
safe_encode_chart_data(%{datasets: datasets})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -1139,7 +1141,7 @@ defmodule ToweropsWeb.DeviceLive.Show do
|
|||
else
|
||||
twenty_four_hours_ago = DateTime.add(DateTime.utc_now(), -24, :hour)
|
||||
datasets = Enum.map(sensors, &sensor_to_chart_dataset(&1, twenty_four_hours_ago))
|
||||
Jason.encode!(%{datasets: datasets})
|
||||
safe_encode_chart_data(%{datasets: datasets})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -1178,7 +1180,7 @@ defmodule ToweropsWeb.DeviceLive.Show do
|
|||
data: Enum.map(checks, &latency_check_to_data_point/1)
|
||||
}
|
||||
|
||||
Jason.encode!(%{datasets: [dataset]})
|
||||
safe_encode_chart_data(%{datasets: [dataset]})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -2096,4 +2098,15 @@ defmodule ToweropsWeb.DeviceLive.Show do
|
|||
end
|
||||
|
||||
defp safe_to_integer(_, default), do: default
|
||||
|
||||
defp safe_encode_chart_data(data) do
|
||||
case Jason.encode(data) do
|
||||
{:ok, encoded} ->
|
||||
encoded
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Failed to encode chart data: #{inspect(reason)}")
|
||||
Jason.encode!(%{datasets: []})
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -162,7 +162,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
end)
|
||||
}
|
||||
|
||||
chart_data = Jason.encode!(%{datasets: [response_time_dataset, status_dataset]})
|
||||
chart_data = safe_encode_chart_data(%{datasets: [response_time_dataset, status_dataset]})
|
||||
{chart_data, "ms", true, true}
|
||||
end
|
||||
end
|
||||
|
|
@ -186,7 +186,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
end)
|
||||
}
|
||||
|
||||
Jason.encode!(%{datasets: [dataset]})
|
||||
safe_encode_chart_data(%{datasets: [dataset]})
|
||||
end
|
||||
|
||||
{chart_data, unit, auto_scale, false}
|
||||
|
|
@ -340,7 +340,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
|> assign(:device, device)
|
||||
|> assign(:page_title, "#{device.name} - #{chart_title} (Live)")
|
||||
|> assign(:chart_title, chart_title)
|
||||
|> assign(:chart_data, Jason.encode!(%{datasets: []}))
|
||||
|> assign(:chart_data, safe_encode_chart_data(%{datasets: []}))
|
||||
|> assign(:unit, unit)
|
||||
|> assign(:auto_scale, auto_scale)
|
||||
|> assign(:dual_axis, false)
|
||||
|
|
@ -441,7 +441,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
if Enum.empty?(dataset.data) do
|
||||
nil
|
||||
else
|
||||
Jason.encode!(%{datasets: [dataset]})
|
||||
safe_encode_chart_data(%{datasets: [dataset]})
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
@ -462,7 +462,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
since = get_datetime_from_range(range)
|
||||
limit = get_limit_for_range(range)
|
||||
dataset = storage_to_dataset(storage, since, limit)
|
||||
Jason.encode!(%{datasets: [dataset]})
|
||||
safe_encode_chart_data(%{datasets: [dataset]})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -504,7 +504,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
since = get_datetime_from_range(range)
|
||||
limit = get_limit_for_range(range)
|
||||
datasets = Enum.map(sensors, &sensor_to_dataset(&1, since, limit))
|
||||
Jason.encode!(%{datasets: datasets})
|
||||
safe_encode_chart_data(%{datasets: datasets})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -647,7 +647,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
datasets
|
||||
end
|
||||
|
||||
if Enum.empty?(datasets), do: nil, else: Jason.encode!(%{datasets: datasets})
|
||||
if Enum.empty?(datasets), do: nil, else: safe_encode_chart_data(%{datasets: datasets})
|
||||
end
|
||||
|
||||
defp build_per_interface_datasets(interfaces, since, limit) do
|
||||
|
|
@ -787,7 +787,7 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
%{label: "Inbound", data: in_data}
|
||||
]
|
||||
|
||||
Jason.encode!(%{datasets: datasets})
|
||||
safe_encode_chart_data(%{datasets: datasets})
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -1373,4 +1373,15 @@ defmodule ToweropsWeb.GraphLive.Show do
|
|||
cancel_live_polling_timer(socket)
|
||||
:ok
|
||||
end
|
||||
|
||||
defp safe_encode_chart_data(data) do
|
||||
case Jason.encode(data) do
|
||||
{:ok, encoded} ->
|
||||
encoded
|
||||
|
||||
{:error, reason} ->
|
||||
Logger.error("Failed to encode chart data: #{inspect(reason)}")
|
||||
Jason.encode!(%{datasets: []})
|
||||
end
|
||||
end
|
||||
end
|
||||
|
|
|
|||
|
|
@ -284,13 +284,11 @@ defmodule Towerops.AdminTest do
|
|||
assert log.metadata["email"] == target_user_email
|
||||
end
|
||||
|
||||
test "raises exception when user does not exist" do
|
||||
test "returns error when user does not exist" do
|
||||
superuser = user_fixture(%{superuser: true})
|
||||
non_existent_id = Ecto.UUID.generate()
|
||||
|
||||
assert_raise Ecto.NoResultsError, fn ->
|
||||
Admin.delete_user(non_existent_id, superuser.id, "127.0.0.1")
|
||||
end
|
||||
assert {:error, :not_found} = Admin.delete_user(non_existent_id, superuser.id, "127.0.0.1")
|
||||
end
|
||||
|
||||
test "anonymizes login history before deletion" do
|
||||
|
|
@ -456,13 +454,11 @@ defmodule Towerops.AdminTest do
|
|||
assert log.metadata["slug"] == org.slug
|
||||
end
|
||||
|
||||
test "raises exception when organization does not exist" do
|
||||
test "returns error when organization does not exist" do
|
||||
superuser = user_fixture(%{superuser: true})
|
||||
non_existent_id = Ecto.UUID.generate()
|
||||
|
||||
assert_raise Ecto.NoResultsError, fn ->
|
||||
Admin.delete_organization(non_existent_id, superuser.id, "127.0.0.1")
|
||||
end
|
||||
assert {:error, :not_found} = Admin.delete_organization(non_existent_id, superuser.id, "127.0.0.1")
|
||||
end
|
||||
end
|
||||
|
||||
|
|
@ -757,15 +753,14 @@ defmodule Towerops.AdminTest do
|
|||
)
|
||||
end
|
||||
|
||||
test "raises for non-existent organization", %{superuser: superuser} do
|
||||
assert_raise Ecto.NoResultsError, fn ->
|
||||
Admin.update_billing_overrides(
|
||||
Ecto.UUID.generate(),
|
||||
%{custom_free_device_limit: 50},
|
||||
superuser.id,
|
||||
"192.168.1.1"
|
||||
)
|
||||
end
|
||||
test "returns error for non-existent organization", %{superuser: superuser} do
|
||||
assert {:error, :not_found} =
|
||||
Admin.update_billing_overrides(
|
||||
Ecto.UUID.generate(),
|
||||
%{custom_free_device_limit: 50},
|
||||
superuser.id,
|
||||
"192.168.1.1"
|
||||
)
|
||||
end
|
||||
end
|
||||
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue