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