diff --git a/appview/config/config.go b/appview/config/config.go
index a3c4f650..af14b971 100644
--- a/appview/config/config.go
+++ b/appview/config/config.go
@@ -69,6 +69,7 @@ type ResendConfig struct {
ApiKey string `env:"API_KEY"`
SentFrom string `env:"SENT_FROM, default=noreply@notifs.tangled.sh"`
NewsletterSegmentId string `env:"NEWSLETTER_SEGMENT_ID"`
+ AssetsURL string `env:"ASSETS_URL, default=https://assets.tangled.network/email/"`
}
type CamoConfig struct {
diff --git a/appview/notify/email/dispatcher.go b/appview/notify/email/dispatcher.go
new file mode 100644
index 00000000..a76391b8
--- /dev/null
+++ b/appview/notify/email/dispatcher.go
@@ -0,0 +1,381 @@
+package email
+
+import (
+ "bytes"
+ "context"
+ "fmt"
+ "html/template"
+ "log/slog"
+ "strings"
+ gotemplate "text/template"
+ "time"
+
+ "tangled.org/core/appview/config"
+ "tangled.org/core/appview/db"
+ "tangled.org/core/appview/email"
+ "tangled.org/core/appview/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}}
+---
+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}}
+
+
+ |
+
+
+
+ |
+
+
+
+
+`
+
+type digestGroup struct {
+ IconURL string
+ Header string
+ HeaderHTML template.HTML
+ EntityRef string
+ URL string
+}
+
+func notifHeader(n *models.NotificationWithEntity, 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.Issue != nil {
+ return actor + " mentioned you on an issue in " + repo
+ } else if n.Pull != nil {
+ return actor + " mentioned you on a pull request in " + repo
+ }
+ return actor + " mentioned you 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.NotificationWithEntity) string {
+ if n.Issue != nil {
+ return fmt.Sprintf("#%d %s", n.Issue.IssueId, n.Issue.Title)
+ }
+ if n.Pull != nil {
+ return fmt.Sprintf("#%d %s", n.Pull.PullId, n.Pull.Title)
+ }
+ return ""
+}
+
+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 digestData struct {
+ RecipientHandle string
+ Count int
+ Groups []digestGroup
+ SettingsURL string
+ AssetsURL string
+}
+
+// Dispatcher polls the notifications table and sends digest emails.
+type Dispatcher struct {
+ db *db.DB
+ resend config.ResendConfig
+ 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 *db.DB,
+ 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,
+ resend: resend,
+ 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)),
+ }
+}
+
+// Start runs the dispatcher ticker loop until ctx is cancelled.
+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 := db.GetPendingEmailDigestRecipients(d.db, cutoff)
+ if err != nil {
+ d.logger.Error("email dispatcher: failed to get recipients", "err", err)
+ return
+ }
+
+ d.logger.Debug("email dispatcher: processing recipients", "count", len(recipients))
+
+ for _, recipientDid := range recipients {
+ if err := d.sendDigest(ctx, recipientDid, cutoff); err != nil {
+ d.logger.Error("email dispatcher: failed to send digest", "did", recipientDid, "err", err)
+ }
+ }
+}
+
+func (d *Dispatcher) sendDigest(ctx context.Context, recipientDid string, cutoff time.Time) error {
+ em, err := db.GetPrimaryEmail(d.db, recipientDid)
+ if err != nil || !em.Verified {
+ return nil
+ }
+
+ notifs, err := db.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 new notifications created during sending
+ // aren't accidentally marked as emailed.
+ ids := make([]int64, len(notifs))
+ for i, n := range notifs {
+ ids[i] = n.ID
+ }
+
+ if err := email.SendEmail(email.Email{
+ APIKey: d.resend.ApiKey,
+ From: "Tangled <" + d.resend.SentFrom + ">",
+ To: em.Address,
+ Subject: subject,
+ Text: text,
+ Html: html,
+ }); err != nil {
+ return fmt.Errorf("send email: %w", err)
+ }
+
+ d.logger.Info("email dispatcher: digest sent", "did", recipientDid, "notifications", len(notifs))
+
+ if err := db.MarkNotificationsEmailed(d.db, ids); err != nil {
+ d.logger.Error("email dispatcher: failed to mark notifications emailed", "did", recipientDid, "err", err)
+ }
+
+ return nil
+}
+
+func (d *Dispatcher) renderDigest(ctx context.Context, recipientHandle string, notifs []*models.NotificationWithEntity) (subject, text, html string, err error) {
+ groups := make([]digestGroup, 0, len(notifs))
+
+ for _, n := range notifs {
+ actorHandle := n.ActorDid
+ if id, err2 := d.resolver.ResolveIdent(ctx, n.ActorDid); err2 == nil && !id.Handle.IsInvalidHandle() {
+ actorHandle = id.Handle.String()
+ }
+
+ repoStr := ""
+ if n.Repo != nil {
+ repoHandle := n.Repo.Did
+ if id, err2 := d.resolver.ResolveIdent(ctx, n.Repo.Did); err2 == nil && !id.Handle.IsInvalidHandle() {
+ repoHandle = id.Handle.String()
+ }
+ repoStr = repoHandle + "/" + n.Repo.Slug()
+ }
+
+ 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),
+ })
+ }
+
+ count := len(notifs)
+ data := digestData{
+ RecipientHandle: recipientHandle,
+ Count: count,
+ Groups: groups,
+ SettingsURL: d.baseURL + "/settings/notifications",
+ AssetsURL: d.assetsURL,
+ }
+
+ 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
+}