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
|
Hi {{.RecipientHandle}},
You have {{.Count}} new notification{{if gt .Count 1}}s{{end}}:
{{range .Groups}}
|
{{.HeaderHTML}}
{{if .EntityRef}}{{.EntityRef}} {{end}}
|
|
{{end}}
{{if .HasMore}}
View more notifications
{{end}}
|
|
`
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
}