towerops/lib/towerops/activity_feed.ex

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