mirror of
https://github.com/coder/coder.git
synced 2026-09-24 15:04:27 +08:00
feat: add notification preferences database & audit support (#14100)
This commit is contained in:
@@ -3,6 +3,7 @@ package notifications
|
||||
import (
|
||||
"context"
|
||||
"encoding/json"
|
||||
"strings"
|
||||
"text/template"
|
||||
|
||||
"github.com/google/uuid"
|
||||
@@ -16,14 +17,13 @@ import (
|
||||
"github.com/coder/coder/v2/codersdk"
|
||||
)
|
||||
|
||||
var ErrCannotEnqueueDisabledNotification = xerrors.New("user has disabled this notification")
|
||||
|
||||
type StoreEnqueuer struct {
|
||||
store Store
|
||||
log slog.Logger
|
||||
|
||||
// TODO: expand this to allow for each notification to have custom delivery methods, or multiple, or none.
|
||||
// For example, Larry might want email notifications for "workspace deleted" notifications, but Harry wants
|
||||
// Slack notifications, and Mary doesn't want any.
|
||||
method database.NotificationMethod
|
||||
defaultMethod database.NotificationMethod
|
||||
// helpers holds a map of template funcs which are used when rendering templates. These need to be passed in because
|
||||
// the template funcs will return values which are inappropriately encapsulated in this struct.
|
||||
helpers template.FuncMap
|
||||
@@ -37,17 +37,31 @@ func NewStoreEnqueuer(cfg codersdk.NotificationsConfig, store Store, helpers tem
|
||||
}
|
||||
|
||||
return &StoreEnqueuer{
|
||||
store: store,
|
||||
log: log,
|
||||
method: method,
|
||||
helpers: helpers,
|
||||
store: store,
|
||||
log: log,
|
||||
defaultMethod: method,
|
||||
helpers: helpers,
|
||||
}, nil
|
||||
}
|
||||
|
||||
// Enqueue queues a notification message for later delivery.
|
||||
// Messages will be dequeued by a notifier later and dispatched.
|
||||
func (s *StoreEnqueuer) Enqueue(ctx context.Context, userID, templateID uuid.UUID, labels map[string]string, createdBy string, targets ...uuid.UUID) (*uuid.UUID, error) {
|
||||
payload, err := s.buildPayload(ctx, userID, templateID, labels)
|
||||
metadata, err := s.store.FetchNewMessageMetadata(ctx, database.FetchNewMessageMetadataParams{
|
||||
UserID: userID,
|
||||
NotificationTemplateID: templateID,
|
||||
})
|
||||
if err != nil {
|
||||
s.log.Warn(ctx, "failed to fetch message metadata", slog.F("template_id", templateID), slog.F("user_id", userID), slog.Error(err))
|
||||
return nil, xerrors.Errorf("new message metadata: %w", err)
|
||||
}
|
||||
|
||||
dispatchMethod := s.defaultMethod
|
||||
if metadata.CustomMethod.Valid {
|
||||
dispatchMethod = metadata.CustomMethod.NotificationMethod
|
||||
}
|
||||
|
||||
payload, err := s.buildPayload(metadata, labels)
|
||||
if err != nil {
|
||||
s.log.Warn(ctx, "failed to build payload", slog.F("template_id", templateID), slog.F("user_id", userID), slog.Error(err))
|
||||
return nil, xerrors.Errorf("enqueue notification (payload build): %w", err)
|
||||
@@ -63,12 +77,21 @@ func (s *StoreEnqueuer) Enqueue(ctx context.Context, userID, templateID uuid.UUI
|
||||
ID: id,
|
||||
UserID: userID,
|
||||
NotificationTemplateID: templateID,
|
||||
Method: s.method,
|
||||
Method: dispatchMethod,
|
||||
Payload: input,
|
||||
Targets: targets,
|
||||
CreatedBy: createdBy,
|
||||
})
|
||||
if err != nil {
|
||||
// We have a trigger on the notification_messages table named `inhibit_enqueue_if_disabled` which prevents messages
|
||||
// from being enqueued if the user has disabled them via notification_preferences. The trigger will fail the insertion
|
||||
// with the message "cannot enqueue message: user has disabled this notification".
|
||||
//
|
||||
// This is more efficient than fetching the user's preferences for each enqueue, and centralizes the business logic.
|
||||
if strings.Contains(err.Error(), ErrCannotEnqueueDisabledNotification.Error()) {
|
||||
return nil, ErrCannotEnqueueDisabledNotification
|
||||
}
|
||||
|
||||
s.log.Warn(ctx, "failed to enqueue notification", slog.F("template_id", templateID), slog.F("input", input), slog.Error(err))
|
||||
return nil, xerrors.Errorf("enqueue notification: %w", err)
|
||||
}
|
||||
@@ -80,15 +103,7 @@ func (s *StoreEnqueuer) Enqueue(ctx context.Context, userID, templateID uuid.UUI
|
||||
// buildPayload creates the payload that the notification will for variable substitution and/or routing.
|
||||
// The payload contains information about the recipient, the event that triggered the notification, and any subsequent
|
||||
// actions which can be taken by the recipient.
|
||||
func (s *StoreEnqueuer) buildPayload(ctx context.Context, userID, templateID uuid.UUID, labels map[string]string) (*types.MessagePayload, error) {
|
||||
metadata, err := s.store.FetchNewMessageMetadata(ctx, database.FetchNewMessageMetadataParams{
|
||||
UserID: userID,
|
||||
NotificationTemplateID: templateID,
|
||||
})
|
||||
if err != nil {
|
||||
return nil, xerrors.Errorf("new message metadata: %w", err)
|
||||
}
|
||||
|
||||
func (s *StoreEnqueuer) buildPayload(metadata database.FetchNewMessageMetadataRow, labels map[string]string) (*types.MessagePayload, error) {
|
||||
payload := types.MessagePayload{
|
||||
Version: "1.0",
|
||||
|
||||
|
||||
@@ -149,7 +149,7 @@ func (m *Manager) loop(ctx context.Context) error {
|
||||
var eg errgroup.Group
|
||||
|
||||
// Create a notifier to run concurrently, which will handle dequeueing and dispatching notifications.
|
||||
m.notifier = newNotifier(m.cfg, uuid.New(), m.log, m.store, m.handlers, m.method, m.metrics)
|
||||
m.notifier = newNotifier(m.cfg, uuid.New(), m.log, m.store, m.handlers, m.metrics)
|
||||
eg.Go(func() error {
|
||||
return m.notifier.run(ctx, m.success, m.failure)
|
||||
})
|
||||
@@ -249,15 +249,24 @@ func (m *Manager) syncUpdates(ctx context.Context) {
|
||||
for i := 0; i < nFailure; i++ {
|
||||
res := <-m.failure
|
||||
|
||||
status := database.NotificationMessageStatusPermanentFailure
|
||||
if res.retryable {
|
||||
var (
|
||||
reason string
|
||||
status database.NotificationMessageStatus
|
||||
)
|
||||
|
||||
switch {
|
||||
case res.retryable:
|
||||
status = database.NotificationMessageStatusTemporaryFailure
|
||||
case res.inhibited:
|
||||
status = database.NotificationMessageStatusInhibited
|
||||
reason = "disabled by user"
|
||||
default:
|
||||
status = database.NotificationMessageStatusPermanentFailure
|
||||
}
|
||||
|
||||
failureParams.IDs = append(failureParams.IDs, res.msg)
|
||||
failureParams.FailedAts = append(failureParams.FailedAts, res.ts)
|
||||
failureParams.Statuses = append(failureParams.Statuses, status)
|
||||
var reason string
|
||||
if res.err != nil {
|
||||
reason = res.err.Error()
|
||||
}
|
||||
@@ -367,4 +376,5 @@ type dispatchResult struct {
|
||||
ts time.Time
|
||||
err error
|
||||
retryable bool
|
||||
inhibited bool
|
||||
}
|
||||
|
||||
@@ -339,6 +339,81 @@ func TestInflightDispatchesMetric(t *testing.T) {
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
}
|
||||
|
||||
func TestCustomMethodMetricCollection(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// SETUP
|
||||
if !dbtestutil.WillUsePostgres() {
|
||||
// UpdateNotificationTemplateMethodByID only makes sense with a real database.
|
||||
t.Skip("This test requires postgres; it relies on business-logic only implemented in the database")
|
||||
}
|
||||
ctx, logger, store := setup(t)
|
||||
|
||||
var (
|
||||
reg = prometheus.NewRegistry()
|
||||
metrics = notifications.NewMetrics(reg)
|
||||
template = notifications.TemplateWorkspaceDeleted
|
||||
anotherTemplate = notifications.TemplateWorkspaceDormant
|
||||
)
|
||||
|
||||
const (
|
||||
customMethod = database.NotificationMethodWebhook
|
||||
defaultMethod = database.NotificationMethodSmtp
|
||||
)
|
||||
|
||||
// GIVEN: a template whose notification method differs from the default.
|
||||
out, err := store.UpdateNotificationTemplateMethodByID(ctx, database.UpdateNotificationTemplateMethodByIDParams{
|
||||
ID: template,
|
||||
Method: database.NullNotificationMethod{NotificationMethod: customMethod, Valid: true},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, customMethod, out.Method.NotificationMethod)
|
||||
|
||||
// WHEN: two notifications (each with different templates) are enqueued.
|
||||
cfg := defaultNotificationsConfig(defaultMethod)
|
||||
mgr, err := notifications.NewManager(cfg, store, metrics, logger.Named("manager"))
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
assert.NoError(t, mgr.Stop(ctx))
|
||||
})
|
||||
|
||||
smtpHandler := &fakeHandler{}
|
||||
webhookHandler := &fakeHandler{}
|
||||
mgr.WithHandlers(map[database.NotificationMethod]notifications.Handler{
|
||||
defaultMethod: smtpHandler,
|
||||
customMethod: webhookHandler,
|
||||
})
|
||||
|
||||
enq, err := notifications.NewStoreEnqueuer(cfg, store, defaultHelpers(), logger.Named("enqueuer"))
|
||||
require.NoError(t, err)
|
||||
|
||||
user := createSampleUser(t, store)
|
||||
|
||||
_, err = enq.Enqueue(ctx, user.ID, template, map[string]string{"type": "success"}, "test")
|
||||
require.NoError(t, err)
|
||||
_, err = enq.Enqueue(ctx, user.ID, anotherTemplate, map[string]string{"type": "success"}, "test")
|
||||
require.NoError(t, err)
|
||||
|
||||
mgr.Run(ctx)
|
||||
|
||||
// THEN: the fake handlers to "dispatch" the notifications.
|
||||
require.Eventually(t, func() bool {
|
||||
smtpHandler.mu.RLock()
|
||||
webhookHandler.mu.RLock()
|
||||
defer smtpHandler.mu.RUnlock()
|
||||
defer webhookHandler.mu.RUnlock()
|
||||
|
||||
return len(smtpHandler.succeeded) == 1 && len(smtpHandler.failed) == 0 &&
|
||||
len(webhookHandler.succeeded) == 1 && len(webhookHandler.failed) == 0
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
|
||||
// THEN: we should have metric series for both the default and custom notification methods.
|
||||
require.Eventually(t, func() bool {
|
||||
return promtest.ToFloat64(metrics.DispatchAttempts.WithLabelValues(string(defaultMethod), anotherTemplate.String(), notifications.ResultSuccess)) > 0 &&
|
||||
promtest.ToFloat64(metrics.DispatchAttempts.WithLabelValues(string(customMethod), template.String(), notifications.ResultSuccess)) > 0
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
}
|
||||
|
||||
// hasMatchingFingerprint checks if the given metric's series fingerprint matches the reference fingerprint.
|
||||
func hasMatchingFingerprint(metric *dto.Metric, fp model.Fingerprint) bool {
|
||||
return fingerprintLabelPairs(metric.Label) == fp
|
||||
|
||||
@@ -604,7 +604,7 @@ func TestNotifierPaused(t *testing.T) {
|
||||
}, testutil.WaitShort, testutil.IntervalFast)
|
||||
}
|
||||
|
||||
func TestNotifcationTemplatesBody(t *testing.T) {
|
||||
func TestNotificationTemplatesBody(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
if !dbtestutil.WillUsePostgres() {
|
||||
@@ -705,6 +705,194 @@ func TestNotifcationTemplatesBody(t *testing.T) {
|
||||
}
|
||||
}
|
||||
|
||||
// TestDisabledBeforeEnqueue ensures that notifications cannot be enqueued once a user has disabled that notification template
|
||||
func TestDisabledBeforeEnqueue(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// SETUP
|
||||
if !dbtestutil.WillUsePostgres() {
|
||||
t.Skip("This test requires postgres; it is testing business-logic implemented in the database")
|
||||
}
|
||||
|
||||
ctx, logger, db := setup(t)
|
||||
|
||||
// GIVEN: an enqueuer & a sample user
|
||||
cfg := defaultNotificationsConfig(database.NotificationMethodSmtp)
|
||||
enq, err := notifications.NewStoreEnqueuer(cfg, db, defaultHelpers(), logger.Named("enqueuer"))
|
||||
require.NoError(t, err)
|
||||
user := createSampleUser(t, db)
|
||||
|
||||
// WHEN: the user has a preference set to not receive the "workspace deleted" notification
|
||||
templateID := notifications.TemplateWorkspaceDeleted
|
||||
n, err := db.UpdateUserNotificationPreferences(ctx, database.UpdateUserNotificationPreferencesParams{
|
||||
UserID: user.ID,
|
||||
NotificationTemplateIds: []uuid.UUID{templateID},
|
||||
Disableds: []bool{true},
|
||||
})
|
||||
require.NoError(t, err, "failed to set preferences")
|
||||
require.EqualValues(t, 1, n, "unexpected number of affected rows")
|
||||
|
||||
// THEN: enqueuing the "workspace deleted" notification should fail with an error
|
||||
_, err = enq.Enqueue(ctx, user.ID, templateID, map[string]string{}, "test")
|
||||
require.ErrorIs(t, err, notifications.ErrCannotEnqueueDisabledNotification, "enqueueing did not fail with expected error")
|
||||
}
|
||||
|
||||
// TestDisabledAfterEnqueue ensures that notifications enqueued before a notification template was disabled will not be
|
||||
// sent, and will instead be marked as "inhibited".
|
||||
func TestDisabledAfterEnqueue(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// SETUP
|
||||
if !dbtestutil.WillUsePostgres() {
|
||||
t.Skip("This test requires postgres; it is testing business-logic implemented in the database")
|
||||
}
|
||||
|
||||
ctx, logger, db := setup(t)
|
||||
|
||||
method := database.NotificationMethodSmtp
|
||||
cfg := defaultNotificationsConfig(method)
|
||||
|
||||
mgr, err := notifications.NewManager(cfg, db, createMetrics(), logger.Named("manager"))
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
assert.NoError(t, mgr.Stop(ctx))
|
||||
})
|
||||
|
||||
enq, err := notifications.NewStoreEnqueuer(cfg, db, defaultHelpers(), logger.Named("enqueuer"))
|
||||
require.NoError(t, err)
|
||||
user := createSampleUser(t, db)
|
||||
|
||||
// GIVEN: a notification is enqueued which has not (yet) been disabled
|
||||
templateID := notifications.TemplateWorkspaceDeleted
|
||||
msgID, err := enq.Enqueue(ctx, user.ID, templateID, map[string]string{}, "test")
|
||||
require.NoError(t, err)
|
||||
|
||||
// Disable the notification template.
|
||||
n, err := db.UpdateUserNotificationPreferences(ctx, database.UpdateUserNotificationPreferencesParams{
|
||||
UserID: user.ID,
|
||||
NotificationTemplateIds: []uuid.UUID{templateID},
|
||||
Disableds: []bool{true},
|
||||
})
|
||||
require.NoError(t, err, "failed to set preferences")
|
||||
require.EqualValues(t, 1, n, "unexpected number of affected rows")
|
||||
|
||||
// WHEN: running the manager to trigger dequeueing of (now-disabled) messages
|
||||
mgr.Run(ctx)
|
||||
|
||||
// THEN: the message should not be sent, and must be set to "inhibited"
|
||||
require.EventuallyWithT(t, func(ct *assert.CollectT) {
|
||||
m, err := db.GetNotificationMessagesByStatus(ctx, database.GetNotificationMessagesByStatusParams{
|
||||
Status: database.NotificationMessageStatusInhibited,
|
||||
Limit: 10,
|
||||
})
|
||||
assert.NoError(ct, err)
|
||||
if assert.Equal(ct, len(m), 1) {
|
||||
assert.Equal(ct, m[0].ID.String(), msgID.String())
|
||||
assert.Contains(ct, m[0].StatusReason.String, "disabled by user")
|
||||
}
|
||||
}, testutil.WaitLong, testutil.IntervalFast, "did not find the expected inhibited message")
|
||||
}
|
||||
|
||||
func TestCustomNotificationMethod(t *testing.T) {
|
||||
t.Parallel()
|
||||
|
||||
// SETUP
|
||||
if !dbtestutil.WillUsePostgres() {
|
||||
t.Skip("This test requires postgres; it relies on business-logic only implemented in the database")
|
||||
}
|
||||
|
||||
ctx, logger, db := setup(t)
|
||||
|
||||
received := make(chan uuid.UUID, 1)
|
||||
|
||||
// SETUP:
|
||||
// Start mock server to simulate webhook endpoint.
|
||||
mockWebhookSrv := httptest.NewServer(http.HandlerFunc(func(w http.ResponseWriter, r *http.Request) {
|
||||
var payload dispatch.WebhookPayload
|
||||
err := json.NewDecoder(r.Body).Decode(&payload)
|
||||
assert.NoError(t, err)
|
||||
|
||||
received <- payload.MsgID
|
||||
close(received)
|
||||
|
||||
w.WriteHeader(http.StatusOK)
|
||||
_, err = w.Write([]byte("noted."))
|
||||
require.NoError(t, err)
|
||||
}))
|
||||
defer mockWebhookSrv.Close()
|
||||
|
||||
// Start mock SMTP server.
|
||||
mockSMTPSrv := smtpmock.New(smtpmock.ConfigurationAttr{
|
||||
LogToStdout: false,
|
||||
LogServerActivity: true,
|
||||
})
|
||||
require.NoError(t, mockSMTPSrv.Start())
|
||||
t.Cleanup(func() {
|
||||
assert.NoError(t, mockSMTPSrv.Stop())
|
||||
})
|
||||
|
||||
endpoint, err := url.Parse(mockWebhookSrv.URL)
|
||||
require.NoError(t, err)
|
||||
|
||||
// GIVEN: a notification template which has a method explicitly set
|
||||
var (
|
||||
template = notifications.TemplateWorkspaceDormant
|
||||
defaultMethod = database.NotificationMethodSmtp
|
||||
customMethod = database.NotificationMethodWebhook
|
||||
)
|
||||
out, err := db.UpdateNotificationTemplateMethodByID(ctx, database.UpdateNotificationTemplateMethodByIDParams{
|
||||
ID: template,
|
||||
Method: database.NullNotificationMethod{NotificationMethod: customMethod, Valid: true},
|
||||
})
|
||||
require.NoError(t, err)
|
||||
require.Equal(t, customMethod, out.Method.NotificationMethod)
|
||||
|
||||
// GIVEN: a manager configured with multiple dispatch methods
|
||||
cfg := defaultNotificationsConfig(defaultMethod)
|
||||
cfg.SMTP = codersdk.NotificationsEmailConfig{
|
||||
From: "danny@coder.com",
|
||||
Hello: "localhost",
|
||||
Smarthost: serpent.HostPort{Host: "localhost", Port: fmt.Sprintf("%d", mockSMTPSrv.PortNumber())},
|
||||
}
|
||||
cfg.Webhook = codersdk.NotificationsWebhookConfig{
|
||||
Endpoint: *serpent.URLOf(endpoint),
|
||||
}
|
||||
|
||||
mgr, err := notifications.NewManager(cfg, db, createMetrics(), logger.Named("manager"))
|
||||
require.NoError(t, err)
|
||||
t.Cleanup(func() {
|
||||
_ = mgr.Stop(ctx)
|
||||
})
|
||||
|
||||
enq, err := notifications.NewStoreEnqueuer(cfg, db, defaultHelpers(), logger)
|
||||
require.NoError(t, err)
|
||||
|
||||
// WHEN: a notification of that template is enqueued, it should be delivered with the configured method - not the default.
|
||||
user := createSampleUser(t, db)
|
||||
msgID, err := enq.Enqueue(ctx, user.ID, template, map[string]string{}, "test")
|
||||
require.NoError(t, err)
|
||||
|
||||
// THEN: the notification should be received by the custom dispatch method
|
||||
mgr.Run(ctx)
|
||||
|
||||
receivedMsgID := testutil.RequireRecvCtx(ctx, t, received)
|
||||
require.Equal(t, msgID.String(), receivedMsgID.String())
|
||||
|
||||
// Ensure no messages received by default method (SMTP):
|
||||
msgs := mockSMTPSrv.MessagesAndPurge()
|
||||
require.Len(t, msgs, 0)
|
||||
|
||||
// Enqueue a notification which does not have a custom method set to ensure default works correctly.
|
||||
msgID, err = enq.Enqueue(ctx, user.ID, notifications.TemplateWorkspaceDeleted, map[string]string{}, "test")
|
||||
require.NoError(t, err)
|
||||
require.EventuallyWithT(t, func(ct *assert.CollectT) {
|
||||
msgs := mockSMTPSrv.MessagesAndPurge()
|
||||
if assert.Len(ct, msgs, 1) {
|
||||
assert.Contains(ct, msgs[0].MsgRequest(), fmt.Sprintf("Message-Id: %s", msgID))
|
||||
}
|
||||
}, testutil.WaitLong, testutil.IntervalFast)
|
||||
}
|
||||
|
||||
type fakeHandler struct {
|
||||
mu sync.RWMutex
|
||||
succeeded, failed []string
|
||||
|
||||
@@ -10,6 +10,7 @@ import (
|
||||
"golang.org/x/sync/errgroup"
|
||||
"golang.org/x/xerrors"
|
||||
|
||||
"github.com/coder/coder/v2/coderd/database/dbtime"
|
||||
"github.com/coder/coder/v2/coderd/notifications/dispatch"
|
||||
"github.com/coder/coder/v2/coderd/notifications/render"
|
||||
"github.com/coder/coder/v2/coderd/notifications/types"
|
||||
@@ -33,12 +34,11 @@ type notifier struct {
|
||||
quit chan any
|
||||
done chan any
|
||||
|
||||
method database.NotificationMethod
|
||||
handlers map[database.NotificationMethod]Handler
|
||||
metrics *Metrics
|
||||
}
|
||||
|
||||
func newNotifier(cfg codersdk.NotificationsConfig, id uuid.UUID, log slog.Logger, db Store, hr map[database.NotificationMethod]Handler, method database.NotificationMethod, metrics *Metrics) *notifier {
|
||||
func newNotifier(cfg codersdk.NotificationsConfig, id uuid.UUID, log slog.Logger, db Store, hr map[database.NotificationMethod]Handler, metrics *Metrics) *notifier {
|
||||
return ¬ifier{
|
||||
id: id,
|
||||
cfg: cfg,
|
||||
@@ -48,7 +48,6 @@ func newNotifier(cfg codersdk.NotificationsConfig, id uuid.UUID, log slog.Logger
|
||||
tick: time.NewTicker(cfg.FetchInterval.Value()),
|
||||
store: db,
|
||||
handlers: hr,
|
||||
method: method,
|
||||
metrics: metrics,
|
||||
}
|
||||
}
|
||||
@@ -144,6 +143,12 @@ func (n *notifier) process(ctx context.Context, success chan<- dispatchResult, f
|
||||
|
||||
var eg errgroup.Group
|
||||
for _, msg := range msgs {
|
||||
// If a notification template has been disabled by the user after a notification was enqueued, mark it as inhibited
|
||||
if msg.Disabled {
|
||||
failure <- n.newInhibitedDispatch(msg)
|
||||
continue
|
||||
}
|
||||
|
||||
// A message failing to be prepared correctly should not affect other messages.
|
||||
deliverFn, err := n.prepare(ctx, msg)
|
||||
if err != nil {
|
||||
@@ -234,17 +239,17 @@ func (n *notifier) deliver(ctx context.Context, msg database.AcquireNotification
|
||||
logger := n.log.With(slog.F("msg_id", msg.ID), slog.F("method", msg.Method), slog.F("attempt", msg.AttemptCount+1))
|
||||
|
||||
if msg.AttemptCount > 0 {
|
||||
n.metrics.RetryCount.WithLabelValues(string(n.method), msg.TemplateID.String()).Inc()
|
||||
n.metrics.RetryCount.WithLabelValues(string(msg.Method), msg.TemplateID.String()).Inc()
|
||||
}
|
||||
|
||||
n.metrics.InflightDispatches.WithLabelValues(string(n.method), msg.TemplateID.String()).Inc()
|
||||
n.metrics.QueuedSeconds.WithLabelValues(string(n.method)).Observe(msg.QueuedSeconds)
|
||||
n.metrics.InflightDispatches.WithLabelValues(string(msg.Method), msg.TemplateID.String()).Inc()
|
||||
n.metrics.QueuedSeconds.WithLabelValues(string(msg.Method)).Observe(msg.QueuedSeconds)
|
||||
|
||||
start := time.Now()
|
||||
retryable, err := deliver(ctx, msg.ID)
|
||||
|
||||
n.metrics.DispatcherSendSeconds.WithLabelValues(string(n.method)).Observe(time.Since(start).Seconds())
|
||||
n.metrics.InflightDispatches.WithLabelValues(string(n.method), msg.TemplateID.String()).Dec()
|
||||
n.metrics.DispatcherSendSeconds.WithLabelValues(string(msg.Method)).Observe(time.Since(start).Seconds())
|
||||
n.metrics.InflightDispatches.WithLabelValues(string(msg.Method), msg.TemplateID.String()).Dec()
|
||||
|
||||
if err != nil {
|
||||
// Don't try to accumulate message responses if the context has been canceled.
|
||||
@@ -281,12 +286,12 @@ func (n *notifier) deliver(ctx context.Context, msg database.AcquireNotification
|
||||
}
|
||||
|
||||
func (n *notifier) newSuccessfulDispatch(msg database.AcquireNotificationMessagesRow) dispatchResult {
|
||||
n.metrics.DispatchAttempts.WithLabelValues(string(n.method), msg.TemplateID.String(), ResultSuccess).Inc()
|
||||
n.metrics.DispatchAttempts.WithLabelValues(string(msg.Method), msg.TemplateID.String(), ResultSuccess).Inc()
|
||||
|
||||
return dispatchResult{
|
||||
notifier: n.id,
|
||||
msg: msg.ID,
|
||||
ts: time.Now(),
|
||||
ts: dbtime.Now(),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -301,17 +306,27 @@ func (n *notifier) newFailedDispatch(msg database.AcquireNotificationMessagesRow
|
||||
result = ResultPermFail
|
||||
}
|
||||
|
||||
n.metrics.DispatchAttempts.WithLabelValues(string(n.method), msg.TemplateID.String(), result).Inc()
|
||||
n.metrics.DispatchAttempts.WithLabelValues(string(msg.Method), msg.TemplateID.String(), result).Inc()
|
||||
|
||||
return dispatchResult{
|
||||
notifier: n.id,
|
||||
msg: msg.ID,
|
||||
ts: time.Now(),
|
||||
ts: dbtime.Now(),
|
||||
err: err,
|
||||
retryable: retryable,
|
||||
}
|
||||
}
|
||||
|
||||
func (n *notifier) newInhibitedDispatch(msg database.AcquireNotificationMessagesRow) dispatchResult {
|
||||
return dispatchResult{
|
||||
notifier: n.id,
|
||||
msg: msg.ID,
|
||||
ts: dbtime.Now(),
|
||||
retryable: false,
|
||||
inhibited: true,
|
||||
}
|
||||
}
|
||||
|
||||
// stop stops the notifier from processing any new notifications.
|
||||
// This is a graceful stop, so any in-flight notifications will be completed before the notifier stops.
|
||||
// Once a notifier has stopped, it cannot be restarted.
|
||||
|
||||
Reference in New Issue
Block a user