package deliberi import ( "bytes" "context" "fmt" "html/template" "log/slog" "strings" gotemplate "text/template" "time" "github.com/bluesky-social/indigo/atproto/syntax" "tangled.org/core/deliberi/config" deldb "tangled.org/core/deliberi/db" "tangled.org/core/deliberi/mailer" "tangled.org/core/deliberi/models" "tangled.org/core/idresolver" ) const digestTextTmpl = `Hi {{.RecipientHandle}}, You have {{.Count}} new notification(s) on Tangled: {{range .Groups}}- {{wrap 70 " " .Header}}{{if .EntityRef}} {{wrap 70 " " .EntityRef}}{{end}} {{.URL}} {{end}}{{if .HasMore}} View more notifications: {{.NotificationsURL}} {{end}}--- Manage notifications: {{.SettingsURL}} ` const digestHTMLTmpl = ` Tangled notifications
Tangled

Hi {{.RecipientHandle}},

You have {{.Count}} new notification{{if gt .Count 1}}s{{end}}:

{{range .Groups}} {{end}}

{{.HeaderHTML}}

{{if .EntityRef}}

{{.EntityRef}}

{{end}}
{{if .HasMore}}

View more notifications

{{end}}

Manage notification settings

Tangled Labs Oy. © 2026 All rights reserved.

tangled.org

` const digestMaxGroups = 10 type digestGroup struct { IconURL string Header string HeaderHTML template.HTML EntityRef string URL string } type digestData struct { RecipientHandle string Count int Groups []digestGroup HasMore bool NotificationsURL string SettingsURL string AssetsURL string } func notifHeader(n *models.Notification, actor, repo string) string { switch n.Type { case models.NotificationTypeIssueCreated: return actor + " opened an issue on " + repo case models.NotificationTypeIssueCommented: return actor + " commented on an issue on " + repo case models.NotificationTypeIssueClosed: return actor + " closed an issue on " + repo case models.NotificationTypeIssueReopen: return actor + " reopened an issue on " + repo case models.NotificationTypePullCreated: return actor + " created a PR on " + repo case models.NotificationTypePullCommented: return actor + " commented on a PR on " + repo case models.NotificationTypePullMerged: return actor + " merged a PR on " + repo case models.NotificationTypePullClosed: return actor + " closed a PR on " + repo case models.NotificationTypePullReopen: return actor + " reopened a PR on " + repo case models.NotificationTypeUserMentioned: if n.EntityAt != "" && syntax.ATURI(n.EntityAt).Collection().String() == "sh.tangled.repo.pull" { return actor + " mentioned you on a pull request in " + repo } return actor + " mentioned you on an issue in " + repo case models.NotificationTypeIssueAssigned: return actor + " assigned you to an issue on " + repo case models.NotificationTypeIssueUnassigned: return actor + " unassigned you from an issue on " + repo case models.NotificationTypePullAssigned: return actor + " assigned you to a PR on " + repo case models.NotificationTypePullUnassigned: return actor + " unassigned you from a PR on " + repo default: return actor + " updated " + repo } } func notifEntityRef(n *models.Notification) string { return n.EntityTitle } func wordwrap(width int, indent, text string) string { words := strings.Fields(text) if len(words) == 0 { return text } var b strings.Builder col := 0 for i, w := range words { if i == 0 { b.WriteString(w) col = len(w) continue } if col+1+len(w) > width { b.WriteString("\n" + indent) b.WriteString(w) col = len(indent) + len(w) } else { b.WriteByte(' ') b.WriteString(w) col += 1 + len(w) } } return b.String() } type Dispatcher struct { db *deldb.DB sender *mailer.Sender baseURL string assetsURL string resolver *idresolver.Resolver logger *slog.Logger batchWait time.Duration interval time.Duration textTmpl *gotemplate.Template htmlTmpl *template.Template } func NewDispatcher(database *deldb.DB, sender *mailer.Sender, resend config.ResendConfig, baseURL string, resolver *idresolver.Resolver, logger *slog.Logger, dev bool) *Dispatcher { batchWait := 10 * time.Minute interval := 5 * time.Minute if dev { batchWait = 30 * time.Second interval = 15 * time.Second } return &Dispatcher{ db: database, sender: sender, baseURL: strings.TrimRight(baseURL, "/"), assetsURL: strings.TrimRight(resend.AssetsURL, "/") + "/", resolver: resolver, logger: logger, batchWait: batchWait, interval: interval, textTmpl: gotemplate.Must(gotemplate.New("digest-text").Funcs(gotemplate.FuncMap{"wrap": wordwrap}).Parse(digestTextTmpl)), htmlTmpl: template.Must(template.New("digest-html").Parse(digestHTMLTmpl)), } } func (d *Dispatcher) Start(ctx context.Context) { d.logger.Info("email dispatcher started", "interval", d.interval, "batchWait", d.batchWait) ticker := time.NewTicker(d.interval) defer ticker.Stop() for { select { case <-ticker.C: d.dispatch(ctx) case <-ctx.Done(): d.logger.Info("email dispatcher stopped") return } } } func (d *Dispatcher) dispatch(ctx context.Context) { cutoff := time.Now().Add(-d.batchWait) recipients, err := deldb.GetPendingEmailDigestRecipients(d.db, cutoff) if err != nil { d.logger.Error("email dispatcher: failed to get recipients", "err", err) return } for _, did := range recipients { if err := d.sendDigest(ctx, did, cutoff); err != nil { d.logger.Error("email dispatcher: failed to send digest", "did", did, "err", err) } } } func (d *Dispatcher) sendDigest(ctx context.Context, recipientDid string, cutoff time.Time) error { em, err := deldb.GetPrimaryEmail(d.db, recipientDid) if err != nil || !em.Verified { return nil } notifs, err := deldb.GetPendingNotificationsForEmailDigest(d.db, recipientDid, cutoff) if err != nil { return fmt.Errorf("get pending notifications: %w", err) } if len(notifs) == 0 { return nil } handle := recipientDid if id, err := d.resolver.ResolveIdent(ctx, recipientDid); err == nil && !id.Handle.IsInvalidHandle() { handle = id.Handle.String() } subject, text, html, err := d.renderDigest(ctx, handle, notifs) if err != nil { return fmt.Errorf("render digest: %w", err) } // collect ids before sending so notifications arriving mid-send are not // marked emailed. ids := make([]int64, len(notifs)) for i, n := range notifs { ids[i] = n.ID } if err := d.sender.Send(em.Address, subject, text, html); err != nil { return fmt.Errorf("send email: %w", err) } d.logger.Info("email dispatcher: digest sent", "did", recipientDid, "notifications", len(notifs)) if err := deldb.MarkEmailed(d.db, ids); err != nil { d.logger.Error("email dispatcher: failed to mark emailed", "did", recipientDid, "err", err) } return nil } func (d *Dispatcher) renderDigest(ctx context.Context, recipientHandle string, notifs []*models.Notification) (subject, text, html string, err error) { count := len(notifs) shown := notifs if len(shown) > digestMaxGroups { shown = shown[:digestMaxGroups] } groups := make([]digestGroup, 0, len(shown)) for _, n := range shown { actorHandle := n.ActorDid if id, err2 := d.resolver.ResolveIdent(ctx, n.ActorDid); err2 == nil && !id.Handle.IsInvalidHandle() { actorHandle = id.Handle.String() } repoStr := "" if n.RepoDid != "" { repoHandle := n.RepoDid if id, err2 := d.resolver.ResolveIdent(ctx, n.RepoDid); err2 == nil && !id.Handle.IsInvalidHandle() { repoHandle = id.Handle.String() } if n.RepoName != "" { repoStr = repoHandle + "/" + n.RepoName } else { repoStr = repoHandle } } header := notifHeader(n, actorHandle, repoStr) headerHTML := template.HTML(strings.Replace(header, actorHandle, ""+actorHandle+"", 1)) groups = append(groups, digestGroup{ IconURL: d.assetsURL + n.Icon() + ".png", Header: header, HeaderHTML: headerHTML, EntityRef: notifEntityRef(n), URL: d.baseURL + n.URL(d.resolver), }) } data := digestData{ RecipientHandle: recipientHandle, Count: count, Groups: groups, HasMore: count > digestMaxGroups, NotificationsURL: d.baseURL + "/notifications", SettingsURL: d.baseURL + "/settings/notifications", AssetsURL: d.assetsURL, } if count > digestMaxGroups { subject = fmt.Sprintf("[%s] %d+ notifications", recipientHandle, digestMaxGroups) } else { subject = fmt.Sprintf("[%s] %d notification(s)", recipientHandle, count) } var textBuf bytes.Buffer if err = d.textTmpl.Execute(&textBuf, data); err != nil { return } text = textBuf.String() var htmlBuf bytes.Buffer if err = d.htmlTmpl.Execute(&htmlBuf, data); err != nil { return } html = htmlBuf.String() return }