Documentation
¶
Overview ¶
Package notifyqueue is the durable delivery queue: the PostgreSQL implementation of notification.QueueStore over the notifications table, and the NOTIFY side of the pg_notify wakeup it fires on every enqueue.
The store and the listener are two halves of one mechanism — the store NOTIFYs on NotifyChannel, a listener LISTENs on it and wakes the send worker — so the channel name that binds them lives here. The listener itself is internal/pglisten, shared with every other queue that wants the same wakeup.
Index ¶
- Constants
- func NewListener(dsn string, notifiers ...pglisten.Notifier) *pglisten.Listener
- type KindState
- type PostgresStore
- func (s *PostgresStore) ClaimDigest(ctx context.Context, lease time.Duration, filter notification.TransportFilter) ([]notification.Notification, error)
- func (s *PostgresStore) ClaimImmediate(ctx context.Context, lease time.Duration, filter notification.TransportFilter) (*notification.Notification, error)
- func (s *PostgresStore) Count(ctx context.Context, filter notification.HistoryFilter) (int, error)
- func (s *PostgresStore) CountsByStatus(ctx context.Context, filter notification.HistoryFilter) (map[string]int, error)
- func (s *PostgresStore) Enqueue(ctx context.Context, n notification.Notification) error
- func (s *PostgresStore) Fail(ctx context.Context, ids []int64, sendErr string) error
- func (s *PostgresStore) List(ctx context.Context, filter notification.HistoryFilter) ([]notification.Notification, error)
- func (s *PostgresStore) MarkSent(ctx context.Context, ids []int64) error
- func (s *PostgresStore) PurgeOld(ctx context.Context, resolvedRetention, pendingTTL time.Duration) (int64, error)
- func (s *PostgresStore) QueueState(ctx context.Context) ([]KindState, error)
- func (s *PostgresStore) Retry(ctx context.Context, ids []int64, sendErr string, backoff time.Duration) error
Constants ¶
const ( KindChannel = "channel" KindEmail = "email" )
Transport kinds the queue state is split by: a row addressed to a channel destination, and every other row, which is mail to a person.
const NotifyChannel = "notifications"
NotifyChannel is the pg_notify channel that wakes the send worker when a row is enqueued. Producers fire it best-effort; the worker also polls.
Variables ¶
This section is empty.
Functions ¶
Types ¶
type KindState ¶ added in v1.141.0
type KindState struct {
Kind string
Pending int64
Sending int64
Waiting int64
OldestDue time.Duration
}
KindState is the open rows of one transport kind, read on a metrics scrape (#1897): due and unclaimed, claimed and sending, scheduled for later (a digest window, a retry's backoff), and how long the oldest due row has waited. Every replica reads the same rows, so a reader takes the max.
type PostgresStore ¶
type PostgresStore struct {
// contains filtered or unexported fields
}
PostgresStore implements notification.QueueStore backed by the notifications table.
func NewPostgresStore ¶
func NewPostgresStore(db *sql.DB) *PostgresStore
NewPostgresStore creates a PostgreSQL-backed notification queue.
func (*PostgresStore) ClaimDigest ¶
func (s *PostgresStore) ClaimDigest(ctx context.Context, lease time.Duration, filter notification.TransportFilter) ([]notification.Notification, error)
ClaimDigest claims all due digest rows for the recipient with the oldest due digest row. Concurrent workers racing for the same recipient are safe: the loser's UPDATE matches zero rows and reports notification.ErrNoWork.
func (*PostgresStore) ClaimImmediate ¶
func (s *PostgresStore) ClaimImmediate(ctx context.Context, lease time.Duration, filter notification.TransportFilter) (*notification.Notification, error)
ClaimImmediate claims the next due non-digest row.
func (*PostgresStore) Count ¶
func (s *PostgresStore) Count(ctx context.Context, filter notification.HistoryFilter) (int, error)
Count returns how many rows match the filter, ignoring its paging fields.
func (*PostgresStore) CountsByStatus ¶
func (s *PostgresStore) CountsByStatus(ctx context.Context, filter notification.HistoryFilter) (map[string]int, error)
CountsByStatus returns the per-status row counts for the filter. Statuses with no rows are absent from the map; callers render a zero for them.
func (*PostgresStore) Enqueue ¶
func (s *PostgresStore) Enqueue(ctx context.Context, n notification.Notification) error
Enqueue inserts a pending notification row and fires a best-effort pg_notify so a listening worker wakes without waiting for the next poll.
func (*PostgresStore) List ¶
func (s *PostgresStore) List(ctx context.Context, filter notification.HistoryFilter) ([]notification.Notification, error)
List returns one page of notification history, newest first.
func (*PostgresStore) MarkSent ¶
func (s *PostgresStore) MarkSent(ctx context.Context, ids []int64) error
MarkSent transitions claimed rows to sent.
func (*PostgresStore) PurgeOld ¶
func (s *PostgresStore) PurgeOld(ctx context.Context, resolvedRetention, pendingTTL time.Duration) (int64, error)
PurgeOld deletes resolved rows past retention and unresolved rows past their delivery-relevance window.
func (*PostgresStore) QueueState ¶ added in v1.141.0
func (s *PostgresStore) QueueState(ctx context.Context) ([]KindState, error)
QueueState reads the open rows by transport kind. A kind with no open row is reported with zeros, so its gauges fall to 0 rather than keeping the last value a scrape saw.