379 lines
14 KiB
Go
379 lines
14 KiB
Go
package store
|
|
|
|
import (
|
|
"context"
|
|
"encoding/json"
|
|
"fmt"
|
|
"strings"
|
|
"time"
|
|
|
|
"github.com/ponzischeme89/memby/server/internal/notify"
|
|
)
|
|
|
|
// The notification log: what Memby sent, to whom, over which channel, and what became of
|
|
// it. Written by internal/notify — nothing else writes this table, which is the whole
|
|
// point of it — and read only by the console.
|
|
//
|
|
// It is deliberately a separate table from user_notifications rather than a set of extra
|
|
// columns on it. That table is *state*: one viewer's undismissed list, which they empty.
|
|
// This is *history*: it keeps a row for a notification that was dismissed, for one that
|
|
// was never delivered, and for a broadcast that belongs to no viewer at all — none of
|
|
// which the other table can represent.
|
|
|
|
// NotificationRetention is how far back the log goes. Ninety days is long enough that a
|
|
// question about "the summary I never got last month" is still answerable, and short
|
|
// enough that the table cannot outgrow the database on a household gateway. The
|
|
// housekeeping task prunes to it; the console derives its widest window from it, so the
|
|
// page can never offer a range the data does not cover.
|
|
const NotificationRetention = 90 * 24 * time.Hour
|
|
|
|
// notificationTextLimit bounds a stored string. A notification body is a sentence or two
|
|
// by construction, and this is only here so a bug in a producer cannot write a megabyte
|
|
// per row into the audit trail.
|
|
const notificationTextLimit = 2000
|
|
|
|
// NotificationLogEntry is one delivered — or refused — notification as the console reads
|
|
// it.
|
|
type NotificationLogEntry struct {
|
|
ID int64 `json:"id"`
|
|
OccurredAt time.Time `json:"occurredAt"`
|
|
Channel string `json:"channel"`
|
|
Kind string `json:"kind"`
|
|
Source string `json:"source"`
|
|
UserID string `json:"userId,omitempty"`
|
|
Username string `json:"username,omitempty"`
|
|
Title string `json:"title"`
|
|
Body string `json:"body,omitempty"`
|
|
ItemID string `json:"itemId,omitempty"`
|
|
Target string `json:"target,omitempty"`
|
|
SourceKey string `json:"sourceKey,omitempty"`
|
|
Status string `json:"status"`
|
|
Detail string `json:"detail,omitempty"`
|
|
DurationMS int64 `json:"durationMs"`
|
|
EventAt *time.Time `json:"eventAt,omitempty"`
|
|
Metadata json.RawMessage `json:"metadata,omitempty"`
|
|
}
|
|
|
|
// NotificationLogFilter is the console's question. Every field is optional and they
|
|
// combine with AND, which is what makes the filter bar above the table read the way it
|
|
// behaves.
|
|
type NotificationLogFilter struct {
|
|
UserID string
|
|
Kinds []string
|
|
Channels []string
|
|
Statuses []string
|
|
Sources []string
|
|
// Query searches the title, the body, the failure detail and the recipient's name. One
|
|
// box rather than four, because an operator arriving here is looking for a *thing* they
|
|
// half remember and does not yet know which column it is in.
|
|
Query string
|
|
From time.Time
|
|
To time.Time
|
|
Limit int
|
|
Offset int
|
|
}
|
|
|
|
// NotificationLogPage is a window on the log plus the counts the page heads itself with.
|
|
type NotificationLogPage struct {
|
|
Entries []NotificationLogEntry `json:"entries"`
|
|
Total int `json:"total"`
|
|
Limit int `json:"limit"`
|
|
Offset int `json:"offset"`
|
|
}
|
|
|
|
// NotificationLogTotals summarises the filtered window. Counted by its own query rather
|
|
// than tallied from the page, for the reason SearchTotals is: the page is capped, so
|
|
// adding it up would report the first hundred rows' totals as the window's.
|
|
type NotificationLogTotals struct {
|
|
Total int `json:"total"`
|
|
Sent int `json:"sent"`
|
|
Delivered int `json:"delivered"`
|
|
Failed int `json:"failed"`
|
|
Pending int `json:"pending"`
|
|
Skipped int `json:"skipped"`
|
|
Users int `json:"users"`
|
|
}
|
|
|
|
// NotificationFacet is one value of a filterable column and how many rows carry it. The
|
|
// console builds its dropdowns from these rather than from a list of constants, so the
|
|
// filter can neither offer a type that matches nothing nor miss one a feature added after
|
|
// the page was written — the stance the activity feed's type filter takes.
|
|
type NotificationFacet struct {
|
|
Value string `json:"value"`
|
|
Count int `json:"count"`
|
|
}
|
|
|
|
// NotificationFacets is every dropdown on the page.
|
|
type NotificationFacets struct {
|
|
Kinds []NotificationFacet `json:"kinds"`
|
|
Channels []NotificationFacet `json:"channels"`
|
|
Statuses []NotificationFacet `json:"statuses"`
|
|
Sources []NotificationFacet `json:"sources"`
|
|
}
|
|
|
|
// RecordNotification writes one row of the audit trail.
|
|
//
|
|
// It implements notify.Recorder, which is the only thing that calls it. Text is clamped
|
|
// here rather than at the caller so one careless producer cannot be the reason the console
|
|
// takes a second to draw.
|
|
func (s *Store) RecordNotification(ctx context.Context, record notify.Record) error {
|
|
if s == nil || s.pool == nil {
|
|
return nil
|
|
}
|
|
occurred := record.OccurredAt
|
|
if occurred.IsZero() {
|
|
occurred = time.Now().UTC()
|
|
}
|
|
var metadata any
|
|
if len(record.Metadata) > 0 {
|
|
metadata = []byte(record.Metadata)
|
|
}
|
|
_, err := s.pool.Exec(ctx, `
|
|
INSERT INTO notification_log
|
|
(occurred_at, channel, kind, source, emby_user_id, username, title, body,
|
|
item_id, target, source_key, status, detail, duration_ms, event_at, metadata)
|
|
VALUES ($1, $2, $3, $4, $5, $6, $7, $8, $9, $10, $11, $12, $13, $14, $15, $16)`,
|
|
occurred, string(record.Channel), record.Kind, record.Source,
|
|
record.UserID, record.Username,
|
|
clampText(record.Title), clampText(record.Body),
|
|
record.ItemID, record.Target, record.SourceKey,
|
|
string(record.Status), clampText(record.Detail),
|
|
record.DurationMS, record.EventAt, metadata)
|
|
if err != nil {
|
|
return fmt.Errorf("store: record notification: %w", err)
|
|
}
|
|
return nil
|
|
}
|
|
|
|
func clampText(value string) string {
|
|
runes := []rune(value)
|
|
if len(runes) <= notificationTextLimit {
|
|
return value
|
|
}
|
|
return string(runes[:notificationTextLimit]) + "…"
|
|
}
|
|
|
|
// notificationWhere builds the shared predicate. The log, the totals and the facets all
|
|
// answer for the *same* filtered window, so they must be filtered identically — writing
|
|
// the clause three times is how a page comes to show a total that disagrees with its own
|
|
// table.
|
|
func notificationWhere(filter NotificationLogFilter) (string, []any) {
|
|
clauses := []string{"TRUE"}
|
|
args := []any{}
|
|
add := func(clause string, value any) {
|
|
args = append(args, value)
|
|
clauses = append(clauses, fmt.Sprintf(clause, len(args)))
|
|
}
|
|
if filter.UserID != "" {
|
|
add("emby_user_id = $%d", filter.UserID)
|
|
}
|
|
if len(filter.Kinds) > 0 {
|
|
add("kind = ANY($%d)", filter.Kinds)
|
|
}
|
|
if len(filter.Channels) > 0 {
|
|
add("channel = ANY($%d)", filter.Channels)
|
|
}
|
|
if len(filter.Statuses) > 0 {
|
|
add("status = ANY($%d)", filter.Statuses)
|
|
}
|
|
if len(filter.Sources) > 0 {
|
|
add("source = ANY($%d)", filter.Sources)
|
|
}
|
|
if !filter.From.IsZero() {
|
|
add("occurred_at >= $%d", filter.From)
|
|
}
|
|
if !filter.To.IsZero() {
|
|
add("occurred_at < $%d", filter.To)
|
|
}
|
|
if query := strings.TrimSpace(filter.Query); query != "" {
|
|
// ILIKE over four columns rather than a tsvector: this table is a few tens of
|
|
// thousands of rows on a household gateway, always read with a date bound, and the
|
|
// operator is looking for a substring of a title or an error message — which is
|
|
// exactly what full-text search is worst at.
|
|
add("(title ILIKE $%[1]d OR body ILIKE $%[1]d OR detail ILIKE $%[1]d OR username ILIKE $%[1]d)",
|
|
"%"+query+"%")
|
|
}
|
|
return strings.Join(clauses, " AND "), args
|
|
}
|
|
|
|
// NotificationLog reads the filtered window, newest first.
|
|
func (s *Store) NotificationLog(
|
|
ctx context.Context, filter NotificationLogFilter,
|
|
) (NotificationLogPage, error) {
|
|
limit := filter.Limit
|
|
if limit <= 0 || limit > 500 {
|
|
limit = 100
|
|
}
|
|
offset := filter.Offset
|
|
if offset < 0 {
|
|
offset = 0
|
|
}
|
|
where, args := notificationWhere(filter)
|
|
|
|
page := NotificationLogPage{Entries: []NotificationLogEntry{}, Limit: limit, Offset: offset}
|
|
if err := s.pool.QueryRow(ctx,
|
|
`SELECT count(*) FROM notification_log WHERE `+where, args...,
|
|
).Scan(&page.Total); err != nil {
|
|
return page, fmt.Errorf("store: count notification log: %w", err)
|
|
}
|
|
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT id, occurred_at, channel, kind, source, emby_user_id, username, title, body,
|
|
item_id, target, source_key, status, detail, duration_ms, event_at, metadata
|
|
FROM notification_log
|
|
WHERE `+where+`
|
|
ORDER BY occurred_at DESC, id DESC
|
|
LIMIT $`+fmt.Sprint(len(args)+1)+` OFFSET $`+fmt.Sprint(len(args)+2),
|
|
append(args, limit, offset)...)
|
|
if err != nil {
|
|
return page, fmt.Errorf("store: read notification log: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
for rows.Next() {
|
|
var entry NotificationLogEntry
|
|
var metadata []byte
|
|
if err := rows.Scan(
|
|
&entry.ID, &entry.OccurredAt, &entry.Channel, &entry.Kind, &entry.Source,
|
|
&entry.UserID, &entry.Username, &entry.Title, &entry.Body,
|
|
&entry.ItemID, &entry.Target, &entry.SourceKey, &entry.Status, &entry.Detail,
|
|
&entry.DurationMS, &entry.EventAt, &metadata,
|
|
); err != nil {
|
|
return page, fmt.Errorf("store: scan notification log: %w", err)
|
|
}
|
|
if len(metadata) > 0 {
|
|
entry.Metadata = json.RawMessage(metadata)
|
|
}
|
|
page.Entries = append(page.Entries, entry)
|
|
}
|
|
if err := rows.Err(); err != nil {
|
|
return page, fmt.Errorf("store: read notification log: %w", err)
|
|
}
|
|
return page, nil
|
|
}
|
|
|
|
// NotificationLogTotals counts the same window the log is read with.
|
|
func (s *Store) NotificationLogTotals(
|
|
ctx context.Context, filter NotificationLogFilter,
|
|
) (NotificationLogTotals, error) {
|
|
where, args := notificationWhere(filter)
|
|
var totals NotificationLogTotals
|
|
err := s.pool.QueryRow(ctx, `
|
|
SELECT count(*),
|
|
count(*) FILTER (WHERE status = 'sent'),
|
|
count(*) FILTER (WHERE status = 'delivered'),
|
|
count(*) FILTER (WHERE status = 'failed'),
|
|
count(*) FILTER (WHERE status = 'pending'),
|
|
count(*) FILTER (WHERE status = 'skipped'),
|
|
count(DISTINCT emby_user_id) FILTER (WHERE emby_user_id <> '')
|
|
FROM notification_log
|
|
WHERE `+where, args...).Scan(
|
|
&totals.Total, &totals.Sent, &totals.Delivered, &totals.Failed,
|
|
&totals.Pending, &totals.Skipped, &totals.Users)
|
|
if err != nil {
|
|
return totals, fmt.Errorf("store: notification totals: %w", err)
|
|
}
|
|
return totals, nil
|
|
}
|
|
|
|
// NotificationLogFacets lists what the filters may offer.
|
|
//
|
|
// Deliberately computed over the retention window rather than over the operator's current
|
|
// filter: a dropdown whose options disappear as you narrow the table is one you cannot use
|
|
// to widen the question again.
|
|
func (s *Store) NotificationLogFacets(
|
|
ctx context.Context, since time.Time,
|
|
) (NotificationFacets, error) {
|
|
facets := NotificationFacets{
|
|
Kinds: []NotificationFacet{}, Channels: []NotificationFacet{},
|
|
Statuses: []NotificationFacet{}, Sources: []NotificationFacet{},
|
|
}
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT 'kind', kind, count(*) FROM notification_log
|
|
WHERE occurred_at >= $1 AND kind <> '' GROUP BY kind
|
|
UNION ALL
|
|
SELECT 'channel', channel, count(*) FROM notification_log
|
|
WHERE occurred_at >= $1 AND channel <> '' GROUP BY channel
|
|
UNION ALL
|
|
SELECT 'status', status, count(*) FROM notification_log
|
|
WHERE occurred_at >= $1 AND status <> '' GROUP BY status
|
|
UNION ALL
|
|
SELECT 'source', source, count(*) FROM notification_log
|
|
WHERE occurred_at >= $1 AND source <> '' GROUP BY source
|
|
ORDER BY 3 DESC, 2`, since)
|
|
if err != nil {
|
|
return facets, fmt.Errorf("store: notification facets: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
for rows.Next() {
|
|
var group string
|
|
var facet NotificationFacet
|
|
if err := rows.Scan(&group, &facet.Value, &facet.Count); err != nil {
|
|
return facets, fmt.Errorf("store: scan notification facet: %w", err)
|
|
}
|
|
switch group {
|
|
case "kind":
|
|
facets.Kinds = append(facets.Kinds, facet)
|
|
case "channel":
|
|
facets.Channels = append(facets.Channels, facet)
|
|
case "status":
|
|
facets.Statuses = append(facets.Statuses, facet)
|
|
case "source":
|
|
facets.Sources = append(facets.Sources, facet)
|
|
}
|
|
}
|
|
return facets, rows.Err()
|
|
}
|
|
|
|
// NotificationLogDays is the daily shape of the filtered window, for the chart above the
|
|
// table. Grouped in the database's own timezone, the stance the sign-in history takes, so
|
|
// an evening notification stays on the day it happened.
|
|
type NotificationLogDay struct {
|
|
Day string `json:"day"`
|
|
Sent int `json:"sent"`
|
|
Failed int `json:"failed"`
|
|
Skipped int `json:"skipped"`
|
|
Delivered int `json:"delivered"`
|
|
}
|
|
|
|
func (s *Store) NotificationLogDays(
|
|
ctx context.Context, filter NotificationLogFilter,
|
|
) ([]NotificationLogDay, error) {
|
|
where, args := notificationWhere(filter)
|
|
rows, err := s.pool.Query(ctx, `
|
|
SELECT to_char(date_trunc('day', occurred_at), 'YYYY-MM-DD'),
|
|
count(*) FILTER (WHERE status = 'sent'),
|
|
count(*) FILTER (WHERE status = 'failed'),
|
|
count(*) FILTER (WHERE status = 'skipped'),
|
|
count(*) FILTER (WHERE status = 'delivered')
|
|
FROM notification_log
|
|
WHERE `+where+`
|
|
GROUP BY 1 ORDER BY 1`, args...)
|
|
if err != nil {
|
|
return nil, fmt.Errorf("store: notification days: %w", err)
|
|
}
|
|
defer rows.Close()
|
|
days := []NotificationLogDay{}
|
|
for rows.Next() {
|
|
var day NotificationLogDay
|
|
if err := rows.Scan(&day.Day, &day.Sent, &day.Failed, &day.Skipped, &day.Delivered); err != nil {
|
|
return nil, fmt.Errorf("store: scan notification day: %w", err)
|
|
}
|
|
days = append(days, day)
|
|
}
|
|
return days, rows.Err()
|
|
}
|
|
|
|
// PruneNotificationLog is the retention policy, run by the housekeeping scheduler.
|
|
func (s *Store) PruneNotificationLog(ctx context.Context, older time.Duration) (int64, error) {
|
|
if older <= 0 {
|
|
return 0, nil
|
|
}
|
|
tag, err := s.pool.Exec(ctx,
|
|
`DELETE FROM notification_log WHERE occurred_at < now() - $1::interval`,
|
|
older.String())
|
|
if err != nil {
|
|
return 0, fmt.Errorf("store: prune notification log: %w", err)
|
|
}
|
|
return tag.RowsAffected(), nil
|
|
}
|