diff --git a/lib/sequin/monitors/sink_monitor.ex b/lib/sequin/monitors/sink_monitor.ex new file mode 100644 index 000000000..ab685976d --- /dev/null +++ b/lib/sequin/monitors/sink_monitor.ex @@ -0,0 +1,116 @@ +defmodule Sequin.Monitors.SinkMonitor do + @moduledoc """ + Schema for sink health monitors. A monitor watches a specific metric + (pending or failed messages) on a sink consumer and fires alerts when + conditions are met. + """ + use Sequin.ConfigSchema + + import Ecto.Changeset + import Ecto.Query + + alias Sequin.Consumers.SinkConsumer + alias Sequin.Monitors.SinkMonitor + alias Sequin.Monitors.SinkMonitorAlert + alias Sequin.Monitors.SinkMonitorNotificationChannel + + @metrics [:pending_messages, :failed_messages] + @monitor_types [:threshold, :anomaly] + @anomaly_sensitivities [:low, :medium, :high] + + typed_schema "sink_monitors" do + field :name, :string + field :enabled, :boolean, default: true + field :metric, Ecto.Enum, values: @metrics + field :monitor_type, Ecto.Enum, values: @monitor_types + field :threshold_value, :integer + field :threshold_window_seconds, :integer, default: 60 + field :anomaly_sensitivity, Ecto.Enum, values: @anomaly_sensitivities + field :notify_on_recovery, :boolean, default: true + field :cooldown_seconds, :integer, default: 300 + + belongs_to :sink_consumer, SinkConsumer + belongs_to :account, Sequin.Accounts.Account + + has_many :notification_channels, SinkMonitorNotificationChannel + has_many :alerts, SinkMonitorAlert + + timestamps() + end + + def metrics, do: @metrics + def monitor_types, do: @monitor_types + def anomaly_sensitivities, do: @anomaly_sensitivities + + def changeset(%SinkMonitor{} = monitor, attrs) do + monitor + |> cast(attrs, [ + :name, + :enabled, + :metric, + :monitor_type, + :threshold_value, + :threshold_window_seconds, + :anomaly_sensitivity, + :notify_on_recovery, + :cooldown_seconds, + :sink_consumer_id, + :account_id + ]) + |> validate_required([:name, :metric, :monitor_type, :sink_consumer_id, :account_id]) + |> validate_inclusion(:metric, @metrics) + |> validate_inclusion(:monitor_type, @monitor_types) + |> validate_number(:cooldown_seconds, greater_than: 0) + |> validate_threshold_fields() + |> validate_anomaly_fields() + |> unique_constraint([:sink_consumer_id, :name], + name: :sink_monitors_consumer_name_unique, + error_key: :name, + message: "monitor name must be unique per sink" + ) + |> foreign_key_constraint(:sink_consumer_id) + |> foreign_key_constraint(:account_id) + end + + defp validate_threshold_fields(changeset) do + if get_field(changeset, :monitor_type) == :threshold do + changeset + |> validate_required([:threshold_value, :threshold_window_seconds]) + |> validate_number(:threshold_value, greater_than: 0) + |> validate_number(:threshold_window_seconds, greater_than_or_equal_to: 0) + else + changeset + |> put_change(:threshold_value, nil) + |> put_change(:threshold_window_seconds, nil) + end + end + + defp validate_anomaly_fields(changeset) do + if get_field(changeset, :monitor_type) == :anomaly do + changeset + |> validate_required([:anomaly_sensitivity]) + |> validate_inclusion(:anomaly_sensitivity, @anomaly_sensitivities) + else + changeset + |> put_change(:anomaly_sensitivity, nil) + end + end + + # Queries + + def where_account_id(query \\ base_query(), account_id) do + from(m in query, where: m.account_id == ^account_id) + end + + def where_sink_consumer_id(query \\ base_query(), sink_consumer_id) do + from(m in query, where: m.sink_consumer_id == ^sink_consumer_id) + end + + def where_enabled(query \\ base_query()) do + from(m in query, where: m.enabled == true) + end + + defp base_query do + from(m in __MODULE__, as: :sink_monitor) + end +end diff --git a/lib/sequin/monitors/sink_monitor_alert.ex b/lib/sequin/monitors/sink_monitor_alert.ex new file mode 100644 index 000000000..a3791c7fc --- /dev/null +++ b/lib/sequin/monitors/sink_monitor_alert.ex @@ -0,0 +1,70 @@ +defmodule Sequin.Monitors.SinkMonitorAlert do + @moduledoc """ + Schema for alert records. Each alert represents a triggered or resolved + incident for a sink monitor. + """ + use Sequin.ConfigSchema + + import Ecto.Changeset + import Ecto.Query + + alias Sequin.Monitors.SinkMonitor + alias Sequin.Monitors.SinkMonitorAlert + + @statuses [:triggered, :resolved] + + typed_schema "sink_monitor_alerts" do + field :status, Ecto.Enum, values: @statuses + field :metric_value, :integer + field :message, :string + field :triggered_at, :utc_datetime_usec + field :resolved_at, :utc_datetime_usec + + belongs_to :sink_monitor, SinkMonitor + + timestamps() + end + + def statuses, do: @statuses + + def changeset(%SinkMonitorAlert{} = alert, attrs) do + alert + |> cast(attrs, [:status, :metric_value, :message, :triggered_at, :resolved_at, :sink_monitor_id]) + |> validate_required([:status, :metric_value, :triggered_at, :sink_monitor_id]) + |> validate_inclusion(:status, @statuses) + |> foreign_key_constraint(:sink_monitor_id) + end + + def resolve_changeset(%SinkMonitorAlert{} = alert) do + now = DateTime.utc_now() + + alert + |> change(%{status: :resolved, resolved_at: now}) + end + + # Queries + + def where_sink_monitor_id(query \\ base_query(), sink_monitor_id) do + from(a in query, where: a.sink_monitor_id == ^sink_monitor_id) + end + + def where_status(query \\ base_query(), status) do + from(a in query, where: a.status == ^status) + end + + def where_triggered(query \\ base_query()) do + where_status(query, :triggered) + end + + def order_by_triggered_at(query \\ base_query(), direction \\ :desc) do + from(a in query, order_by: [{^direction, a.triggered_at}]) + end + + def where_older_than(query \\ base_query(), datetime) do + from(a in query, where: a.triggered_at < ^datetime) + end + + defp base_query do + from(a in __MODULE__, as: :sink_monitor_alert) + end +end diff --git a/lib/sequin/monitors/sink_monitor_metric_sample.ex b/lib/sequin/monitors/sink_monitor_metric_sample.ex new file mode 100644 index 000000000..82dcf96af --- /dev/null +++ b/lib/sequin/monitors/sink_monitor_metric_sample.ex @@ -0,0 +1,50 @@ +defmodule Sequin.Monitors.SinkMonitorMetricSample do + @moduledoc """ + Schema for metric samples used by anomaly detection. + Stores periodic snapshots of metric values for each monitor + to compute rolling statistics. + """ + use Sequin.ConfigSchema + + import Ecto.Changeset + import Ecto.Query + + alias Sequin.Monitors.SinkMonitor + alias Sequin.Monitors.SinkMonitorMetricSample + + typed_schema "sink_monitor_metric_samples" do + field :value, :integer + field :sampled_at, :utc_datetime_usec + + belongs_to :sink_monitor, SinkMonitor + end + + def changeset(%SinkMonitorMetricSample{} = sample, attrs) do + sample + |> cast(attrs, [:value, :sampled_at, :sink_monitor_id]) + |> validate_required([:value, :sampled_at, :sink_monitor_id]) + |> foreign_key_constraint(:sink_monitor_id) + end + + # Queries + + def where_sink_monitor_id(query \\ base_query(), sink_monitor_id) do + from(s in query, where: s.sink_monitor_id == ^sink_monitor_id) + end + + def where_after(query \\ base_query(), datetime) do + from(s in query, where: s.sampled_at >= ^datetime) + end + + def where_before(query \\ base_query(), datetime) do + from(s in query, where: s.sampled_at < ^datetime) + end + + def order_by_sampled_at(query \\ base_query(), direction \\ :asc) do + from(s in query, order_by: [{^direction, s.sampled_at}]) + end + + defp base_query do + from(s in __MODULE__, as: :sink_monitor_metric_sample) + end +end diff --git a/lib/sequin/monitors/sink_monitor_notification_channel.ex b/lib/sequin/monitors/sink_monitor_notification_channel.ex new file mode 100644 index 000000000..a71adef21 --- /dev/null +++ b/lib/sequin/monitors/sink_monitor_notification_channel.ex @@ -0,0 +1,96 @@ +defmodule Sequin.Monitors.SinkMonitorNotificationChannel do + @moduledoc """ + Schema for notification channels attached to a sink monitor. + Each channel defines a destination (webhook, Slack, Discord, etc.) + where alerts are sent. + """ + use Sequin.ConfigSchema + + import Ecto.Changeset + + alias Sequin.Monitors.SinkMonitor + alias Sequin.Monitors.SinkMonitorNotificationChannel + + @channel_types [:webhook, :slack, :discord, :pagerduty, :incident_io] + + typed_schema "sink_monitor_notification_channels" do + field :channel_type, Ecto.Enum, values: @channel_types + field :config, Sequin.Encrypted.Map + + belongs_to :sink_monitor, SinkMonitor + + timestamps() + end + + def channel_types, do: @channel_types + + def changeset(%SinkMonitorNotificationChannel{} = channel, attrs) do + channel + |> cast(attrs, [:channel_type, :config, :sink_monitor_id]) + |> validate_required([:channel_type, :config, :sink_monitor_id]) + |> validate_inclusion(:channel_type, @channel_types) + |> validate_channel_config() + |> foreign_key_constraint(:sink_monitor_id) + end + + defp validate_channel_config(changeset) do + channel_type = get_field(changeset, :channel_type) + config = get_field(changeset, :config) + + if channel_type && config do + case validate_config_for_type(channel_type, config) do + :ok -> changeset + {:error, message} -> add_error(changeset, :config, message) + end + else + changeset + end + end + + defp validate_config_for_type(:webhook, config) do + if is_binary(config["url"]) and config["url"] != "" do + :ok + else + {:error, "webhook config requires a 'url' field"} + end + end + + defp validate_config_for_type(:slack, config) do + if is_binary(config["webhook_url"]) and config["webhook_url"] != "" do + :ok + else + {:error, "slack config requires a 'webhook_url' field"} + end + end + + defp validate_config_for_type(:discord, config) do + if is_binary(config["webhook_url"]) and config["webhook_url"] != "" do + :ok + else + {:error, "discord config requires a 'webhook_url' field"} + end + end + + defp validate_config_for_type(:pagerduty, config) do + if is_binary(config["routing_key"]) and config["routing_key"] != "" do + :ok + else + {:error, "pagerduty config requires a 'routing_key' field"} + end + end + + defp validate_config_for_type(:incident_io, config) do + cond do + not (is_binary(config["api_key"]) and config["api_key"] != "") -> + {:error, "incident_io config requires an 'api_key' field"} + + not (config["severity"] in ["minor", "major", "critical"]) -> + {:error, "incident_io config requires 'severity' to be 'minor', 'major', or 'critical'"} + + true -> + :ok + end + end + + defp validate_config_for_type(_, _config), do: :ok +end diff --git a/priv/repo/migrations/20260226150000_create_sink_monitors.exs b/priv/repo/migrations/20260226150000_create_sink_monitors.exs new file mode 100644 index 000000000..91e4aae20 --- /dev/null +++ b/priv/repo/migrations/20260226150000_create_sink_monitors.exs @@ -0,0 +1,109 @@ +defmodule Sequin.Repo.Migrations.CreateSinkMonitors do + use Ecto.Migration + + @config_schema Application.compile_env(:sequin, [Sequin.Repo, :config_schema_prefix]) + + def change do + # Create enum types + execute( + "CREATE TYPE #{@config_schema}.sink_monitor_metric AS ENUM ('pending_messages', 'failed_messages')", + "DROP TYPE #{@config_schema}.sink_monitor_metric" + ) + + execute( + "CREATE TYPE #{@config_schema}.sink_monitor_type AS ENUM ('threshold', 'anomaly')", + "DROP TYPE #{@config_schema}.sink_monitor_type" + ) + + execute( + "CREATE TYPE #{@config_schema}.sink_monitor_anomaly_sensitivity AS ENUM ('low', 'medium', 'high')", + "DROP TYPE #{@config_schema}.sink_monitor_anomaly_sensitivity" + ) + + execute( + "CREATE TYPE #{@config_schema}.sink_monitor_channel_type AS ENUM ('webhook', 'slack', 'discord', 'pagerduty', 'incident_io')", + "DROP TYPE #{@config_schema}.sink_monitor_channel_type" + ) + + execute( + "CREATE TYPE #{@config_schema}.sink_monitor_alert_status AS ENUM ('triggered', 'resolved')", + "DROP TYPE #{@config_schema}.sink_monitor_alert_status" + ) + + # sink_monitors table + create table(:sink_monitors, prefix: @config_schema) do + add :sink_consumer_id, + references(:sink_consumers, on_delete: :delete_all, prefix: @config_schema), + null: false + + add :account_id, + references(:accounts, on_delete: :delete_all, prefix: @config_schema), + null: false + + add :name, :string, null: false + add :enabled, :boolean, null: false, default: true + add :metric, :"#{@config_schema}.sink_monitor_metric", null: false + add :monitor_type, :"#{@config_schema}.sink_monitor_type", null: false + add :threshold_value, :integer + add :threshold_window_seconds, :integer, default: 60 + add :anomaly_sensitivity, :"#{@config_schema}.sink_monitor_anomaly_sensitivity" + add :notify_on_recovery, :boolean, null: false, default: true + add :cooldown_seconds, :integer, null: false, default: 300 + + timestamps() + end + + create index(:sink_monitors, [:sink_consumer_id], prefix: @config_schema) + create index(:sink_monitors, [:account_id], prefix: @config_schema) + + create unique_index(:sink_monitors, [:sink_consumer_id, :name], + prefix: @config_schema, + name: :sink_monitors_consumer_name_unique + ) + + # sink_monitor_notification_channels table + create table(:sink_monitor_notification_channels, prefix: @config_schema) do + add :sink_monitor_id, + references(:sink_monitors, on_delete: :delete_all, prefix: @config_schema), + null: false + + add :channel_type, :"#{@config_schema}.sink_monitor_channel_type", null: false + add :config, :binary, null: false + + timestamps() + end + + create index(:sink_monitor_notification_channels, [:sink_monitor_id], prefix: @config_schema) + + # sink_monitor_alerts table + create table(:sink_monitor_alerts, prefix: @config_schema) do + add :sink_monitor_id, + references(:sink_monitors, on_delete: :delete_all, prefix: @config_schema), + null: false + + add :status, :"#{@config_schema}.sink_monitor_alert_status", null: false + add :metric_value, :integer, null: false + add :message, :text + add :triggered_at, :utc_datetime_usec, null: false + add :resolved_at, :utc_datetime_usec + + timestamps() + end + + create index(:sink_monitor_alerts, [:sink_monitor_id], prefix: @config_schema) + create index(:sink_monitor_alerts, [:triggered_at], prefix: @config_schema) + + # sink_monitor_metric_samples table + create table(:sink_monitor_metric_samples, prefix: @config_schema) do + add :sink_monitor_id, + references(:sink_monitors, on_delete: :delete_all, prefix: @config_schema), + null: false + + add :value, :integer, null: false + add :sampled_at, :utc_datetime_usec, null: false + end + + create index(:sink_monitor_metric_samples, [:sink_monitor_id], prefix: @config_schema) + create index(:sink_monitor_metric_samples, [:sampled_at], prefix: @config_schema) + end +end