Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
116 changes: 116 additions & 0 deletions lib/sequin/monitors/sink_monitor.ex
Original file line number Diff line number Diff line change
@@ -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
70 changes: 70 additions & 0 deletions lib/sequin/monitors/sink_monitor_alert.ex
Original file line number Diff line number Diff line change
@@ -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
50 changes: 50 additions & 0 deletions lib/sequin/monitors/sink_monitor_metric_sample.ex
Original file line number Diff line number Diff line change
@@ -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
96 changes: 96 additions & 0 deletions lib/sequin/monitors/sink_monitor_notification_channel.ex
Original file line number Diff line number Diff line change
@@ -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
Loading