448 lines
14 KiB
Elixir
448 lines
14 KiB
Elixir
defmodule Towerops.ActivityFeed do
|
|
@moduledoc """
|
|
Aggregates events from multiple sources into a unified, chronologically sorted
|
|
activity feed for an organization.
|
|
|
|
Sources:
|
|
- Config changes (config_change_events)
|
|
- Alerts fired/resolved (alerts)
|
|
- Device events (device_events)
|
|
- Integration syncs (preseem_sync_logs)
|
|
- Device additions (devices)
|
|
"""
|
|
|
|
import Ecto.Query
|
|
|
|
alias Towerops.Alerts.Alert
|
|
alias Towerops.ConfigChanges.ConfigChangeEvent
|
|
alias Towerops.Devices.Device
|
|
alias Towerops.Devices.Event, as: DeviceEvent
|
|
alias Towerops.Preseem.SyncLog
|
|
alias Towerops.Repo
|
|
alias Towerops.Sites.Site
|
|
|
|
@type activity_item :: %{
|
|
type: atom(),
|
|
timestamp: DateTime.t(),
|
|
device_name: String.t() | nil,
|
|
site_name: String.t() | nil,
|
|
summary: String.t(),
|
|
detail: String.t() | nil,
|
|
severity: :info | :warning | :critical,
|
|
icon: String.t(),
|
|
link: String.t() | nil,
|
|
id: String.t()
|
|
}
|
|
|
|
@doc """
|
|
Lists unified activity for an organization, sorted by timestamp descending.
|
|
|
|
## Options
|
|
|
|
* `:limit` - max items to return (default 50)
|
|
* `:after` - only items after this datetime
|
|
* `:before` - only items before this datetime
|
|
* `:types` - list of atom type filters (e.g. `[:config_change, :alert_fired]`)
|
|
|
|
"""
|
|
@spec list_org_activity(String.t(), keyword()) :: [activity_item()]
|
|
def list_org_activity(organization_id, opts \\ []) do
|
|
limit = Keyword.get(opts, :limit, 50)
|
|
types = Keyword.get(opts, :types, nil)
|
|
after_dt = Keyword.get(opts, :after, nil)
|
|
before_dt = Keyword.get(opts, :before, nil)
|
|
search = Keyword.get(opts, :search, nil)
|
|
|
|
sources = active_sources(types)
|
|
|
|
items =
|
|
Enum.flat_map(sources, fn source ->
|
|
fetch_source(source, organization_id, after_dt, before_dt, limit)
|
|
end)
|
|
|
|
items
|
|
|> maybe_filter_search(search)
|
|
|> Enum.sort_by(& &1.timestamp, {:desc, DateTime})
|
|
|> Enum.take(limit)
|
|
end
|
|
|
|
@doc """
|
|
Returns counts of items per type for an organization (for filter chip badges).
|
|
"""
|
|
def count_by_type(organization_id) do
|
|
# Count recent items (last 24h) per type for badge display
|
|
since = DateTime.add(DateTime.utc_now(), -86_400, :second)
|
|
|
|
nil
|
|
|> active_sources()
|
|
|> Map.new(fn source ->
|
|
items = fetch_source(source, organization_id, since, nil, 1000)
|
|
{source, length(items)}
|
|
end)
|
|
end
|
|
|
|
defp maybe_filter_search(items, nil), do: items
|
|
defp maybe_filter_search(items, ""), do: items
|
|
|
|
defp maybe_filter_search(items, search) do
|
|
search_lower = String.downcase(search)
|
|
|
|
Enum.filter(items, fn item ->
|
|
matches_search?(item, search_lower)
|
|
end)
|
|
end
|
|
|
|
defp matches_search?(item, search) do
|
|
searchable =
|
|
[
|
|
item.device_name,
|
|
item.site_name,
|
|
item.summary,
|
|
item.detail
|
|
]
|
|
|> Enum.reject(&is_nil/1)
|
|
|> Enum.join(" ")
|
|
|> String.downcase()
|
|
|
|
String.contains?(searchable, search)
|
|
end
|
|
|
|
defp active_sources(nil), do: [:config_change, :alert_fired, :alert_resolved, :device_event, :sync, :device_added]
|
|
|
|
defp active_sources(types) when is_list(types), do: types
|
|
|
|
# --- Config Changes ---
|
|
|
|
defp fetch_source(:config_change, org_id, after_dt, before_dt, limit) do
|
|
query =
|
|
from(c in ConfigChangeEvent,
|
|
join: d in Device,
|
|
on: d.id == c.device_id,
|
|
left_join: s in Site,
|
|
on: s.id == d.site_id,
|
|
where: c.organization_id == ^org_id,
|
|
select: %{
|
|
id: c.id,
|
|
timestamp: c.changed_at,
|
|
device_name: d.name,
|
|
site_name: s.name,
|
|
device_id: c.device_id,
|
|
sections_changed: c.sections_changed,
|
|
change_size: c.change_size
|
|
},
|
|
order_by: [desc: c.changed_at],
|
|
limit: ^limit
|
|
)
|
|
|
|
query
|
|
|> maybe_after(:changed_at, after_dt)
|
|
|> maybe_before(:changed_at, before_dt)
|
|
|> Repo.all()
|
|
|> Enum.map(fn row ->
|
|
sections = row.sections_changed || []
|
|
sections_str = if sections == [], do: "", else: Enum.join(sections, ", ")
|
|
|
|
%{
|
|
type: :config_change,
|
|
id: "cc-#{row.id}",
|
|
timestamp: row.timestamp,
|
|
device_name: row.device_name,
|
|
site_name: row.site_name,
|
|
summary: "Config change on #{row.device_name}",
|
|
detail: "Sections: #{sections_str} \u00b7 #{row.change_size || 0} lines changed",
|
|
severity: severity_for_change_size(row.change_size),
|
|
icon: "hero-wrench-screwdriver",
|
|
link: "/devices/#{row.device_id}/config-timeline"
|
|
}
|
|
end)
|
|
end
|
|
|
|
# --- Alerts Fired ---
|
|
|
|
defp fetch_source(:alert_fired, org_id, after_dt, before_dt, limit) do
|
|
query =
|
|
from(a in Alert,
|
|
join: d in Device,
|
|
on: d.id == a.device_id,
|
|
left_join: s in Site,
|
|
on: s.id == d.site_id,
|
|
where: d.organization_id == ^org_id,
|
|
where: a.alert_type == "device_down",
|
|
select: %{
|
|
id: a.id,
|
|
timestamp: a.triggered_at,
|
|
device_name: d.name,
|
|
site_name: s.name,
|
|
device_id: a.device_id,
|
|
message: a.message,
|
|
gaiia_impact: a.gaiia_impact
|
|
},
|
|
order_by: [desc: a.triggered_at],
|
|
limit: ^limit
|
|
)
|
|
|
|
query
|
|
|> maybe_after(:triggered_at, after_dt)
|
|
|> maybe_before(:triggered_at, before_dt)
|
|
|> Repo.all()
|
|
|> Enum.map(fn row ->
|
|
impact_detail = format_gaiia_impact(row.gaiia_impact)
|
|
|
|
%{
|
|
type: :alert_fired,
|
|
id: "af-#{row.id}",
|
|
timestamp: row.timestamp,
|
|
device_name: row.device_name,
|
|
site_name: row.site_name,
|
|
summary: "Alert: #{row.device_name} unreachable",
|
|
detail: impact_detail || row.message,
|
|
severity: :critical,
|
|
icon: "hero-exclamation-triangle",
|
|
link: "/devices/#{row.device_id}"
|
|
}
|
|
end)
|
|
end
|
|
|
|
# --- Alerts Resolved ---
|
|
|
|
defp fetch_source(:alert_resolved, org_id, after_dt, before_dt, limit) do
|
|
query =
|
|
from(a in Alert,
|
|
join: d in Device,
|
|
on: d.id == a.device_id,
|
|
left_join: s in Site,
|
|
on: s.id == d.site_id,
|
|
where: d.organization_id == ^org_id,
|
|
where: not is_nil(a.resolved_at),
|
|
select: %{
|
|
id: a.id,
|
|
timestamp: a.resolved_at,
|
|
triggered_at: a.triggered_at,
|
|
resolved_at: a.resolved_at,
|
|
device_name: d.name,
|
|
site_name: s.name,
|
|
device_id: a.device_id
|
|
},
|
|
order_by: [desc: a.resolved_at],
|
|
limit: ^limit
|
|
)
|
|
|
|
query
|
|
|> maybe_after(:resolved_at, after_dt)
|
|
|> maybe_before(:resolved_at, before_dt)
|
|
|> Repo.all()
|
|
|> Enum.map(fn row ->
|
|
duration = format_alert_duration(row.triggered_at, row.resolved_at)
|
|
|
|
%{
|
|
type: :alert_resolved,
|
|
id: "ar-#{row.id}",
|
|
timestamp: row.timestamp,
|
|
device_name: row.device_name,
|
|
site_name: row.site_name,
|
|
summary: "Alert resolved: #{row.device_name} back online",
|
|
detail: "Duration: #{duration}",
|
|
severity: :info,
|
|
icon: "hero-check-circle",
|
|
link: "/devices/#{row.device_id}"
|
|
}
|
|
end)
|
|
end
|
|
|
|
# --- Device Events ---
|
|
|
|
defp fetch_source(:device_event, org_id, after_dt, before_dt, limit) do
|
|
query =
|
|
from(e in DeviceEvent,
|
|
join: d in Device,
|
|
on: d.id == e.device_id,
|
|
left_join: s in Site,
|
|
on: s.id == d.site_id,
|
|
where: d.organization_id == ^org_id,
|
|
select: %{
|
|
id: e.id,
|
|
timestamp: e.occurred_at,
|
|
device_name: d.name,
|
|
site_name: s.name,
|
|
device_id: e.device_id,
|
|
event_type: e.event_type,
|
|
severity: e.severity,
|
|
message: e.message
|
|
},
|
|
order_by: [desc: e.occurred_at],
|
|
limit: ^limit
|
|
)
|
|
|
|
query
|
|
|> maybe_after(:occurred_at, after_dt)
|
|
|> maybe_before(:occurred_at, before_dt)
|
|
|> Repo.all()
|
|
|> Enum.map(fn row ->
|
|
%{
|
|
type: :device_event,
|
|
id: "de-#{row.id}",
|
|
timestamp: row.timestamp,
|
|
device_name: row.device_name,
|
|
site_name: row.site_name,
|
|
summary: humanize_event_type(row.event_type) <> " on #{row.device_name}",
|
|
detail: row.message,
|
|
severity: parse_severity(row.severity),
|
|
icon: icon_for_event_type(row.event_type),
|
|
link: "/devices/#{row.device_id}"
|
|
}
|
|
end)
|
|
end
|
|
|
|
# --- Preseem Syncs ---
|
|
|
|
defp fetch_source(:sync, org_id, after_dt, before_dt, limit) do
|
|
# Group by minute to show only one sync entry per minute (deduplicate rapid syncs)
|
|
query =
|
|
from(sl in SyncLog,
|
|
where: sl.organization_id == ^org_id,
|
|
select: %{
|
|
id: fragment("(array_agg(?::text ORDER BY ? DESC))[1]", sl.id, sl.inserted_at),
|
|
timestamp: fragment("date_trunc('minute', ?)", sl.inserted_at),
|
|
status: fragment("(array_agg(?::text ORDER BY ? DESC))[1]", sl.status, sl.inserted_at),
|
|
records_synced: sum(sl.records_synced),
|
|
duration_ms: avg(sl.duration_ms)
|
|
},
|
|
group_by: fragment("date_trunc('minute', ?)", sl.inserted_at),
|
|
order_by: [desc: fragment("date_trunc('minute', ?)", sl.inserted_at)],
|
|
limit: ^limit
|
|
)
|
|
|
|
query
|
|
|> maybe_after(:inserted_at, after_dt)
|
|
|> maybe_before(:inserted_at, before_dt)
|
|
|> Repo.all()
|
|
|> Enum.map(fn row ->
|
|
status_label = if row.status == "success", do: "completed", else: row.status
|
|
|
|
%{
|
|
type: :sync,
|
|
id: "sy-#{row.id}",
|
|
timestamp: DateTime.from_naive!(row.timestamp, "Etc/UTC"),
|
|
device_name: nil,
|
|
site_name: nil,
|
|
summary: "Preseem sync #{status_label}",
|
|
detail:
|
|
"#{row.records_synced || 0} records synced#{if row.duration_ms, do: " in #{Decimal.to_integer(Decimal.round(row.duration_ms, 0))}ms", else: ""}",
|
|
severity: sync_severity(row.status),
|
|
icon: "hero-arrow-path",
|
|
link: nil
|
|
}
|
|
end)
|
|
end
|
|
|
|
# --- Device Additions ---
|
|
|
|
defp fetch_source(:device_added, org_id, after_dt, before_dt, limit) do
|
|
query =
|
|
from(d in Device,
|
|
left_join: s in Site,
|
|
on: s.id == d.site_id,
|
|
where: d.organization_id == ^org_id,
|
|
select: %{
|
|
id: d.id,
|
|
timestamp: d.inserted_at,
|
|
device_name: d.name,
|
|
site_name: s.name,
|
|
ip_address: d.ip_address
|
|
},
|
|
order_by: [desc: d.inserted_at],
|
|
limit: ^limit
|
|
)
|
|
|
|
query
|
|
|> maybe_after(:inserted_at, after_dt)
|
|
|> maybe_before(:inserted_at, before_dt)
|
|
|> Repo.all()
|
|
|> Enum.map(fn row ->
|
|
%{
|
|
type: :device_added,
|
|
id: "da-#{row.id}",
|
|
timestamp: row.timestamp,
|
|
device_name: row.device_name,
|
|
site_name: row.site_name,
|
|
summary: "Device added: #{row.device_name}",
|
|
detail:
|
|
[row.site_name && "Site: #{row.site_name}", row.ip_address]
|
|
|> Enum.reject(&is_nil/1)
|
|
|> Enum.join(" \u00b7 "),
|
|
severity: :info,
|
|
icon: "hero-plus-circle",
|
|
link: "/devices/#{row.id}"
|
|
}
|
|
end)
|
|
end
|
|
|
|
defp fetch_source(_, _org_id, _after_dt, _before_dt, _limit), do: []
|
|
|
|
# --- Helpers ---
|
|
|
|
defp maybe_after(query, _field, nil), do: query
|
|
defp maybe_after(query, :changed_at, dt), do: where(query, [c], c.changed_at >= ^dt)
|
|
defp maybe_after(query, :triggered_at, dt), do: where(query, [a], a.triggered_at >= ^dt)
|
|
defp maybe_after(query, :resolved_at, dt), do: where(query, [a], a.resolved_at >= ^dt)
|
|
defp maybe_after(query, :occurred_at, dt), do: where(query, [e], e.occurred_at >= ^dt)
|
|
defp maybe_after(query, :inserted_at, dt), do: where(query, [r], r.inserted_at >= ^dt)
|
|
|
|
defp maybe_before(query, _field, nil), do: query
|
|
defp maybe_before(query, :changed_at, dt), do: where(query, [c], c.changed_at <= ^dt)
|
|
defp maybe_before(query, :triggered_at, dt), do: where(query, [a], a.triggered_at <= ^dt)
|
|
defp maybe_before(query, :resolved_at, dt), do: where(query, [a], a.resolved_at <= ^dt)
|
|
defp maybe_before(query, :occurred_at, dt), do: where(query, [e], e.occurred_at <= ^dt)
|
|
defp maybe_before(query, :inserted_at, dt), do: where(query, [r], r.inserted_at <= ^dt)
|
|
|
|
defp severity_for_change_size(nil), do: :info
|
|
defp severity_for_change_size(size) when size > 50, do: :warning
|
|
defp severity_for_change_size(_), do: :info
|
|
|
|
defp sync_severity("failed"), do: :critical
|
|
defp sync_severity("partial"), do: :warning
|
|
defp sync_severity(_), do: :info
|
|
|
|
defp parse_severity("critical"), do: :critical
|
|
defp parse_severity("warning"), do: :warning
|
|
defp parse_severity(_), do: :info
|
|
|
|
defp format_gaiia_impact(nil), do: nil
|
|
|
|
defp format_gaiia_impact(%{"total_subscribers" => subs, "total_mrr" => mrr}) when is_integer(subs) and subs > 0 do
|
|
"Affects #{subs} subscribers, $#{mrr}/mo MRR at risk"
|
|
end
|
|
|
|
defp format_gaiia_impact(_), do: nil
|
|
|
|
defp format_alert_duration(nil, _), do: "unknown"
|
|
defp format_alert_duration(_, nil), do: "unknown"
|
|
|
|
defp format_alert_duration(triggered, resolved) do
|
|
diff = DateTime.diff(resolved, triggered, :second)
|
|
|
|
cond do
|
|
diff < 60 -> "#{diff}s"
|
|
diff < 3600 -> "#{div(diff, 60)}m"
|
|
diff < 86_400 -> "#{div(diff, 3600)}h #{div(rem(diff, 3600), 60)}m"
|
|
true -> "#{div(diff, 86_400)}d #{div(rem(diff, 86_400), 3600)}h"
|
|
end
|
|
end
|
|
|
|
defp humanize_event_type(type) when is_binary(type) do
|
|
type
|
|
|> String.replace("_", " ")
|
|
|> String.capitalize()
|
|
end
|
|
|
|
defp humanize_event_type(_), do: "Device event"
|
|
|
|
defp icon_for_event_type("interface_down"), do: "hero-arrow-down-circle"
|
|
defp icon_for_event_type("interface_up"), do: "hero-arrow-up-circle"
|
|
defp icon_for_event_type("device_discovered"), do: "hero-magnifying-glass"
|
|
defp icon_for_event_type("device_rediscovered"), do: "hero-magnifying-glass"
|
|
defp icon_for_event_type("sensor_threshold_critical"), do: "hero-fire"
|
|
defp icon_for_event_type("sensor_threshold_warning"), do: "hero-exclamation-circle"
|
|
defp icon_for_event_type("sensor_threshold_normal"), do: "hero-check-circle"
|
|
defp icon_for_event_type(_), do: "hero-bolt"
|
|
end
|