notifyfeed

package
v0.8.4 Latest Latest
Warning

This package is not in the latest version of its module.

Go to latest
Published: Oct 9, 2026 License: Apache-2.0 Imports: 8 Imported by: 0

Documentation

Overview

Package notifyfeed is the durable, per-user in-app notification feed — the data layer behind the bell. It records change-driven notifications (fanned out from the alert engine; in later slices, the transaction log), tracks per-user read state, and serves the unread count + list the UI renders.

This is distinct from internal/notification, which delivers alerts OUTWARD to Slack/email/webhook channels. notifyfeed is the INWARD, in-app surface.

Design: docs/engineering/notifications_design.md. Spec: system-notifications.

Index

Constants

This section is empty.

Variables

View Source
var ErrNotFound = errors.New("notifyfeed: notification not found")

ErrNotFound is returned when a mark-read targets a notification that does not exist for the calling user (wrong id, or another user's row).

Functions

This section is empty.

Types

type Channel

type Channel struct {
	// contains filtered or unexported fields
}

Channel is an alertrouter.Channel that fans each alert into the in-app feed — one durable row per active recipient. Registered at boot alongside the stdout/Slack/email channels, so every alert the engine already classifies (host_unreachable, host_recovered, drift_major/minor/improvement) lights up the bell for free. Later slices add a transaction-log projector for rule-level regressions. Spec: system-notifications.

func NewChannel

func NewChannel(store *Store) *Channel

NewChannel returns the in-app notification channel.

func (*Channel) Name

func (c *Channel) Name() string

Name identifies the channel in router metrics + logs.

func (*Channel) Send

func (c *Channel) Send(ctx context.Context, alert alertrouter.Alert) error

Send fans the alert into one in-app notification per active recipient, in a single statement (RecordFanout) so an alert is one DB round-trip rather than one per user.

type GovernanceProjector

type GovernanceProjector struct {
	// contains filtered or unexported fields
}

GovernanceProjector turns governance + remediation lifecycle events into RBAC-scoped in-app notifications — the bell's "action queue" (Slice 3). Unlike the alert Channel and the regression Projector (which fan to every active user), governance items reach only the users who can act on them: an exception pending approval reaches approvers, a decision reaches the requester, a failed remediation reaches the operators who can re-run it.

It implements the Notifier interface the exception service and the GovernanceNotifier the remediation worker hold, so neither producer imports this package. Spec: system-notifications (Slice 3).

func NewGovernanceProjector

func NewGovernanceProjector(store *Store) *GovernanceProjector

NewGovernanceProjector returns a governance projector over the feed store.

func (*GovernanceProjector) ExceptionDecided

func (g *GovernanceProjector) ExceptionDecided(ctx context.Context, exceptionID, requestedBy uuid.UUID, ruleID string, approved bool) error

ExceptionDecided records the outcome notification for the requester (only), closing the loop for the person who asked. Grouped per exception.

func (*GovernanceProjector) ExceptionExpired

func (g *GovernanceProjector) ExceptionExpired(ctx context.Context, exceptionID, hostID uuid.UUID, ruleID string) error

ExceptionExpired notifies approvers that an exception has lapsed and its rules are back in scope. Fires once per exception (the sweep flips each to expired exactly once), so the standard fan-out is used. Grouped per exception.

func (*GovernanceProjector) ExceptionExpiringSoon

func (g *GovernanceProjector) ExceptionExpiringSoon(ctx context.Context, exceptionID, hostID uuid.UUID, ruleID string) error

ExceptionExpiringSoon warns approvers that an approved exception is about to lapse (after which its rules re-enter scope). Uses the QUIET fan-out: the expiry sweep re-evaluates hourly, so a non-quiet record would re-surface the same warning unread every hour. Grouped per exception.

func (*GovernanceProjector) ExceptionRequested

func (g *GovernanceProjector) ExceptionRequested(ctx context.Context, exceptionID, hostID uuid.UUID, ruleID string) error

ExceptionRequested records an "exception pending approval" notification for every user whose role can approve it (exception:approve → security_admin, admin). Grouped per exception so a re-surfaced request collapses onto one row. Best-effort: an error is returned for the caller to log, never to fail the request.

func (*GovernanceProjector) PasswordExpiring added in v0.3.0

func (g *GovernanceProjector) PasswordExpiring(ctx context.Context, hostID uuid.UUID, username string, daysLeft int, expired bool) error

PasswordExpiring warns host operators that a host user account's password is about to expire (or has expired). Reaches everyone who can view the host (host:read). Uses the QUIET fan-out: the daily sweep re-evaluates every 24h, so a non-quiet record would re-surface the same warning unread each day. Grouped per (host, user) so the daily re-sweep collapses onto one row and a later "expired" simply updates the same item. Best-effort.

func (*GovernanceProjector) RemediationFailed

func (g *GovernanceProjector) RemediationFailed(ctx context.Context, hostID uuid.UUID, ruleID, action, finalStatus string) error

RemediationFailed records a "remediation failed / rolled back" notification for the operators who can act on it (remediation:execute → ops_lead, security_admin, admin). Grouped per (host, rule). finalStatus is the terminal remediation status ("failed" | "rolled_back"); action is "execute" | "rollback".

type Notification

type Notification struct {
	ID         uuid.UUID
	UserID     uuid.UUID
	Kind       string
	Severity   string
	Title      string
	Body       string
	HostID     *uuid.UUID
	Link       string
	GroupKey   string
	OccurredAt time.Time
	ReadAt     *time.Time
	CreatedAt  time.Time
}

Notification is one in-app notification row for one recipient.

type Projector

type Projector struct {
	// contains filtered or unexported fields
}

Projector turns the transaction-log changes a completed scan wrote into a single, grouped, per-host in-app notification — the headline "a rule that was passing is now failing" surface (notifications_design.md §3, Slice 2).

It is the second producer of the bell, alongside the alertrouter Channel (host_unreachable/recovered, drift_*). The Channel covers fleet/host-level alerts; the Projector covers RULE-level regressions, which the alert engine does not classify. Both write through the same notifyfeed.Store, so read state, grouping, and fan-out are uniform.

Runs in the scan worker (where transactionlog.Writer.Apply commits), called best-effort after the outcomes persist. Spec: system-notifications.

func NewProjector

func NewProjector(store *Store) *Projector

NewProjector returns a regression projector over the given feed store.

func (*Projector) ProjectScan

func (p *Projector) ProjectScan(ctx context.Context, scanID, hostID uuid.UUID) error

ProjectScan reads the changes the given scan wrote for the given host and, if any qualify as a regression, records ONE grouped "rule_regression" notification fanned to every active recipient.

What qualifies (notifications_design.md §3 "Compliance"):

  • state_changed -> fail: a rule that was passing (or skipped/errored) now fails. The unambiguous regression; always counted.
  • first_seen -> fail, severity critical: a brand-new critical finding. Only counted when the host has prior scan history — on a host's FIRST scan every rule is first_seen, which is a baseline, not a regression, and must not flood the bell with "N rules regressed".

Grouping (design §8): one notification per host (group_key "rule_regression:<host>"), so a burst of N rules in one scan is one row, and a later scan's regressions collapse onto (and re-surface) the same row rather than piling up. The notification severity is the highest among the regressed rules; the title summarizes the counts.

Best-effort: a nil/empty result is a no-op returning nil. The caller treats any error as non-fatal (the scan already persisted).

type Store

type Store struct {
	// contains filtered or unexported fields
}

Store is the PostgreSQL-backed notification feed.

func NewStore

func NewStore(pool *pgxpool.Pool) *Store

NewStore returns a feed store over the given pool.

func (*Store) List

func (s *Store) List(ctx context.Context, userID uuid.UUID, unreadOnly bool, limit int) ([]Notification, error)

List returns a user's notifications newest-first (by occurred_at). When unreadOnly is true, only unread rows are returned. limit caps the result (defaulted/clamped to a sane page size).

func (*Store) MarkAllRead

func (s *Store) MarkAllRead(ctx context.Context, userID uuid.UUID) (int, error)

MarkAllRead marks every unread notification for the user read; returns the number affected.

func (*Store) MarkRead

func (s *Store) MarkRead(ctx context.Context, userID, id uuid.UUID) error

MarkRead marks one notification read, scoped to the owning user (so a user can never read or mutate another user's row). Idempotent on an already-read row; returns ErrNotFound when no such row belongs to the user.

func (*Store) Record

func (s *Store) Record(ctx context.Context, n Notification) error

Record upserts one notification for one user. A repeat of the same change (same user_id + group_key) collapses onto the existing row, refreshing its content + occurred_at and re-surfacing it as UNREAD — so a recurring problem re-pings the bell without creating a second entry. id/created_at on the passed Notification are ignored (the store assigns them).

func (*Store) RecordFanout

func (s *Store) RecordFanout(ctx context.Context, n Notification) error

RecordFanout records one notification for EVERY active (non-deleted) user in a single statement (one row per recipient, same upsert/collapse semantics as Record). This replaces a per-user Record loop, so a fleet alert is one DB round-trip rather than N — important on large user bases. The UserID on the template is ignored; recipients come from the users table.

func (*Store) RecordForRoles

func (s *Store) RecordForRoles(ctx context.Context, roleIDs []string, n Notification) error

RecordForRoles records one notification for every active (non-deleted) user who holds at least one of the given built-in roles, in a single statement (same upsert/collapse semantics as RecordFanout). This is the RBAC-scoped fan-out behind governance notifications — e.g. an exception pending approval reaches only users whose role grants exception:approve, not the whole fleet. An empty roleIDs slice matches no one (a no-op). The UserID on the template is ignored; recipients come from users joined to user_roles. EXISTS (not JOIN) keeps a user holding two matching roles to one row, avoiding an ON CONFLICT double-hit within the statement.

func (*Store) RecordForRolesQuiet

func (s *Store) RecordForRolesQuiet(ctx context.Context, roleIDs []string, n Notification) error

RecordForRolesQuiet is RecordForRoles with at-most-once-per-recipient semantics: ON CONFLICT it does NOTHING (the existing row is left untouched — content, occurred_at, AND read state preserved). This is for repeatedly-swept warnings (e.g. "exception expiring soon", re-evaluated every hour): a recipient is notified once and the row does not re-surface unread on every sweep, which would be noise. A recipient who joins later still gets it on the next sweep (their INSERT). Empty roleIDs is a no-op.

func (*Store) UnreadCount

func (s *Store) UnreadCount(ctx context.Context, userID uuid.UUID) (int, error)

UnreadCount returns how many unread notifications a user has.

Jump to

Keyboard shortcuts

? : This menu
/ : Search site
f or F : Jump to
y or Y : Canonical URL